diff --git a/bridged/bridged.example.yaml b/bridged/bridged.example.yaml index 3fcf852..e297fe9 100644 --- a/bridged/bridged.example.yaml +++ b/bridged/bridged.example.yaml @@ -144,6 +144,17 @@ herdrSocket: ~/.config/herdr/herdr.sock # omit and this profile's completion fallback behaves exactly as before. # Every backend words its refusal differently, so this is config, never a # vendor string baked into bridged itself. +# DEFERRED: compiled once into a startup pattern map — editing it needs a +# daemon restart, same as this profile's model/baseUrl/argv. +# credentialId → CB-578 stage B: the credential this profile quarantines WITH when a +# BACKEND_EXHAUSTED classification fires. Two profiles that set the SAME +# credentialId share one quarantine — the case this exists for is two models +# on one account (e.g. sol and terra both billing one OpenAI credential): an +# exhaustion on either one must lock out both, or the fleet just walks onto +# the same dead account under the sibling's name. Opt-in — omit and this +# profile quarantines alone, under its own name, exactly as if the field did +# not exist. Cooldown length is the top-level quarantineCooldownSeconds below. +# HOT: read live at every spawn/exhaustion check — no restart needed. # env → extra environment for this profile's workers, as a literal key/value map # (CB-511). Use it to give workers a toolchain. # @@ -178,6 +189,7 @@ profiles: # gitTokenEnv: GITEA_TOKEN # opt-in: let this profile's workers open their own PR (CB-302) # gitHostEnv: GITEA_HOST # defaults to GITEA_HOST; injected only with gitTokenEnv # exhaustedPattern: "usage limit has been reached" # opt-in: classify a usage-limit refusal (CB-578) + # credentialId: shared-openai # opt-in: quarantine together with every other profile sharing this id (CB-578) # configDir: /Users/me/.ccs/instances/gx10 # CLAUDE_CONFIG_DIR — inherit that profile's skills/MCP # cwd: /Users/me/src/myrepo # pin the working dir; omit to inherit the primary's # parityOverlay: [".claude/settings.local.json", ".env", ".envrc"] # never add .mcp.json — see above @@ -243,6 +255,15 @@ profiles: # round-robin, or weighted. Omitting this key is a strict no-op for existing configs. placement: weighted +# How long a credential sits out after a BACKEND_EXHAUSTED classification (CB-578 stage B), in +# seconds, before a spawn may land on it again. Applies to every profile's effective credential +# (its own name, or its credentialId if set above) — there is no per-profile override. Default +# 1800 (30 minutes) when omitted or non-positive. +# DEFERRED: baked once into the BackendQuarantine built at startup — a running quarantine keeps +# its original cooldown regardless; a new value only applies to a quarantine that starts after a +# restart. Editing this needs a daemon restart to take effect. +# quarantineCooldownSeconds: 1800 + # Re-read this file without restarting the daemon (CB-559). Off unless you add this block, so an # upgraded bridged keeps the old behaviour: the file is read once at boot and never again. # enabled → turn the watch on. bridged checks the file's modified time on a timer and @@ -252,17 +273,20 @@ placement: weighted # Not every key can move under a running daemon, and the difference is about what already exists # when the reload happens — not about how important the key is: # HOT → takes effect on the next spawn: the whole `fleet:` block (every role pool, -# `charters`, and `tabLabel`), `placement:`, and an existing profile's weight / maxLoad. Those are -# hot because the placement policy reads them through a supplier — being config is +# `charters`, and `tabLabel`), `placement:`, and an existing profile's weight / maxLoad +# / credentialId. Those are hot because the placement policy (and, for credentialId, +# the CB-578 stage B quarantine check) reads them through a supplier — being config is # not by itself enough to make a key hot. # DEFERRED → accepted into the new config, but the wiring built at startup keeps the old value # until you restart: `lifecycle:`, `leadHeartbeat:`, `guard:`, `worktreeRoot:`, -# `spawnReadyTimeoutMs` / `spawnReadyPollMs`, ADDING or REMOVING a profile (a new -# backend needs its own launcher, and launchers are built once), AND an existing -# profile's launch settings — model, baseUrl, argv, env, configDir, mcpUrl, tabLabel. -# The launcher takes a copy of `profiles:` at startup and resolves every spawn out of -# that copy, so those never reach a launch until you restart. The reload logs them by -# name rather than pretending they applied. +# `spawnReadyTimeoutMs` / `spawnReadyPollMs`, `quarantineCooldownSeconds` (CB-578 +# stage B — baked once into the quarantine tracker built at startup), ADDING or +# REMOVING a profile (a new backend needs its own launcher, and launchers are built +# once), AND an existing profile's launch settings — model, baseUrl, argv, env, +# configDir, mcpUrl, tabLabel, exhaustedPattern. The launcher takes a copy of +# `profiles:` at startup and resolves every spawn out of that copy, so those never +# reach a launch until you restart. The reload logs them by name rather than +# pretending they applied. # COLD → cannot change at all: `bind:`, `herdrSocket:`, `broker:` and `auth:`. The socket is # bound, the broker connection is open, and the auth mode decides who may reach the # port that is already listening. diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 97211d2..797ad3b 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -14,6 +14,7 @@ import dev.ltms.bridged.herdr.UnixSocketHerdrClient; import dev.ltms.bridged.herdr.WorkspaceControl; import dev.ltms.bridged.inject.CompletionResolver; import dev.ltms.bridged.inject.ExhaustedPatternLookup; +import dev.ltms.bridged.inject.ExhaustionSink; import dev.ltms.bridged.inject.Injector; import dev.ltms.bridged.inject.StatusPoller; import dev.ltms.bridged.inject.TurnListener; @@ -37,6 +38,7 @@ import dev.ltms.bridged.msg.LeadHeartbeatLoop; import dev.ltms.bridged.msg.ReplyPushLoop; import dev.ltms.bridged.rest.BridgedApp; import dev.ltms.bridged.session.GitWorktrees; +import dev.ltms.bridged.session.MemberSession; import dev.ltms.bridged.session.SessionManager; import dev.ltms.bridged.peer.PeerLauncher; import dev.ltms.bridged.session.SessionReaper; @@ -44,6 +46,7 @@ import dev.ltms.bridged.member.ClaudeCodeLauncher; import dev.ltms.bridged.member.CompositePeerLauncher; import dev.ltms.bridged.member.HerdrPeerLauncher; import dev.ltms.bridged.member.OpenCodeLauncher; +import dev.ltms.bridged.placement.BackendQuarantine; import io.javalin.Javalin; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -141,11 +144,18 @@ public final class Bridged { () -> config.get().fleet())); } AtomicReference> liveCountRef = new AtomicReference<>(_ -> 0); + // CB-578 stage B: one quarantine tracker for the whole daemon, shared between the launcher + // (checked at spawn) and the exhaustion sink wired in below (written on BACKEND_EXHAUSTED). + // The cooldown is deferred (see BridgedConfig#quarantineCooldownSeconds): it is read once + // here, at startup, and a config reload only changes it for a daemon restart. + BackendQuarantine quarantine = new BackendQuarantine(System::nanoTime, + TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds())); PeerLauncher workers = new CompositePeerLauncher( adapters, cfg.effectiveDefaultProfile(), config, - profileName -> liveCountRef.get().apply(profileName)); + profileName -> liveCountRef.get().apply(profileName), + quarantine); // CB-504: under supervision (launchd/systemd) bridged can start before herdr's socket // exists. The client itself is lazy — it connects per call — but the orphan reap below is // the first thing that actually talks to herdr, so without this wait a boot-order race @@ -272,7 +282,23 @@ public final class Bridged { .orElse(null); log.info("backend-exhausted classification (CB-578 stage A): {}", CompletionResolver.coverage(cfg.profiles().keySet(), exhaustedPatternsByProfile.keySet())); - CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns); + // CB-578 stage B: on a classification that actually wins, quarantine the exhausted profile's + // CREDENTIAL — not the profile name — so a profile sharing that credential (e.g. two models + // on one OpenAI account) is refused too, not just the one that happened to report it. Reads + // the profile config live off `config`, so a credentialId edit is hot: no restart needed. + ExhaustionSink exhaustionSink = (target, reason) -> sessions.roster().stream() + .filter(session -> target.equals(session.terminalId())) + .findFirst() + .map(MemberSession::profile) + .map(profileName -> config.get().profiles().get(profileName)) + .ifPresent(profile -> { + String credentialId = profile.effectiveCredentialId(); + quarantine.quarantine(credentialId); + log.warn("credential '{}' quarantined for {}s (profile '{}' classified " + + "BACKEND_EXHAUSTED): {}", credentialId, + cfg.quarantineCooldownSeconds(), profile.profile(), reason); + }); + CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns, exhaustionSink); // CB-113: deliver only to an available worker (its MCP is connected), never its boot window. // CB-301: the manager's presence bridge records availability and drives SPAWNING → READY. MemberPresence presence = sessions.asPresence(); @@ -440,7 +466,11 @@ public final class Bridged { var health = config.get().health(); return FleetHealthMonitor.coverage(health != null && health.isEnabled(), health != null && health.notifications() != null && health.notifications().configured()); - })); + }), + new BridgeMcp.QuarantineSource(profile -> { + var configured = config.get().profiles().get(profile); + return configured == null ? null : configured.effectiveCredentialId(); + }, quarantine)); // CB-559: opt-in config reload. With no `configReload:` block nothing is constructed, so an // upgraded daemon behaves exactly as before — the file is read once at boot and never again. diff --git a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java index 4d257a9..6ff7fc5 100644 --- a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java +++ b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java @@ -58,6 +58,12 @@ import java.util.Set; * {@code fixed} (default), {@code round-robin}, or {@code weighted} * @param auth API authentication mode ({@code null} → {@code loopback-trust}, the * historical behaviour), CB-501 + * @param quarantineCooldownSeconds how long a credential stays quarantined after a + * {@code BACKEND_EXHAUSTED} classification (CB-578 stage B); {@code null}/{@code + * <=0} → {@link #DEFAULT_QUARANTINE_COOLDOWN_SECONDS}. Baked once into the + * {@code BackendQuarantine} built at startup, so it is DEFERRED: changing it + * needs a restart, and a quarantine already running keeps whatever cooldown was + * live when it started. */ @JsonIgnoreProperties(ignoreUnknown = true) public record BridgedConfig( @@ -76,7 +82,11 @@ public record BridgedConfig( Health health, String placement, Auth auth, - ConfigReload configReload) { + ConfigReload configReload, + Integer quarantineCooldownSeconds) { + + /** Default cooldown (CB-578 stage B) when {@code quarantineCooldownSeconds} is absent/non-positive. */ + public static final int DEFAULT_QUARANTINE_COOLDOWN_SECONDS = 1800; /** Back-compat 14-arg form — no {@code configReload:} block, so file watching stays off. */ public BridgedConfig(Bind bind, String herdrSocket, Map profiles, Guard guard, @@ -84,7 +94,7 @@ public record BridgedConfig( Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet, LeadHeartbeat leadHeartbeat, String placement, Auth auth) { this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs, - spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, null); + spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, null, null); } /** Back-compat form before the optional {@code health:} block was added. */ @@ -93,7 +103,18 @@ public record BridgedConfig( Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet, LeadHeartbeat leadHeartbeat, String placement, Auth auth, ConfigReload configReload) { this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs, - spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, configReload); + spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, configReload, null); + } + + /** Back-compat form before the CB-578 stage B {@code quarantineCooldownSeconds} field was added. */ + public BridgedConfig(Bind bind, String herdrSocket, Map profiles, Guard guard, + String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs, + Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet, + LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth, + ConfigReload configReload) { + this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs, + spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth, + configReload, null); } /** @@ -187,6 +208,17 @@ public record BridgedConfig( * {@code null}/blank ⇒ the classification never fires for this profile and * today's completion-fallback behaviour is unchanged. Every backend words * its refusal differently, so this is config, never a vendor string in code. + * @param credentialId the shared account this profile authenticates as (CB-578 stage B). Two + * or more profiles setting the same non-blank value are quarantined + * together by one {@code exhaustedPattern} classification on any one of + * them — the case a single OpenAI (or any other) credential backing several + * profiles ({@code sol}, {@code terra}, ...) needs, so a spawn does not walk + * straight onto the other profile sharing the same exhausted account. + * {@code null}/blank ⇒ this profile's {@link #effectiveCredentialId()} is + * its own name, so it quarantines alone — today's behaviour for every + * profile that does not opt in. Read live off the current config, so it is + * HOT: a change takes effect on the next exhaustion classification / spawn, + * no restart needed. */ @JsonIgnoreProperties(ignoreUnknown = true) public record Profile(String profile, String baseUrl, String model, @@ -200,7 +232,8 @@ public record BridgedConfig( Float weight, Integer maxLoad, Boolean subscription, - String exhaustedPattern) { + String exhaustedPattern, + String credentialId) { /** Peer kind spawned by {@link dev.ltms.bridged.member.ClaudeCodeLauncher} (the default). */ public static final String KIND_CLAUDE_CODE = "claude-code"; @@ -246,6 +279,10 @@ public record BridgedConfig( // exhaustedPattern stays null when unset/blank (opt-in) — no defaulting, no vendor // wording: an unconfigured profile keeps today's completion-fallback behaviour exactly. exhaustedPattern = (exhaustedPattern == null || exhaustedPattern.isBlank()) ? null : exhaustedPattern; + // credentialId stays null when unset/blank — effectiveCredentialId() is where the + // "quarantines alone" fallback actually lives, so today's behaviour needs no defaulting + // here at all. + credentialId = (credentialId == null || credentialId.isBlank()) ? null : credentialId; } /** @@ -290,7 +327,7 @@ public record BridgedConfig( public Profile withProfile(String p) { return new Profile(p, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel, mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, subscription, - exhaustedPattern); + exhaustedPattern, credentialId); } /** True when this profile is served by the Claude Code adapter (the default kind). */ @@ -326,11 +363,38 @@ public record BridgedConfig( mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, null, null); } + /** + * Backward-compatible constructor without the CB-578 stage B {@code credentialId} — the + * profile quarantines alone under its own name (see {@link #effectiveCredentialId()}). Keeps + * pre-stage-B call sites (and any YAML that omits the key) compiling and behaving identically. + */ + public Profile(String profile, String baseUrl, String model, + String configDir, String tokenEnv, List argv, + String placement, String workspace, String tabLabel, String mcpUrl, + String cwd, List parityOverlay, String gitTokenEnv, String gitHostEnv, + String kind, Map env, Float weight, Integer maxLoad, + Boolean subscription, String exhaustedPattern) { + this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel, + mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, + subscription, exhaustedPattern, null); + } + /** True when this profile's CB-578 stage A backend-exhausted classification is configured. */ public boolean hasExhaustedPattern() { return exhaustedPattern != null; } + /** + * The credential group this profile quarantines with (CB-578 stage B): the configured + * {@link #credentialId} when set, else this profile's own name — so an unconfigured profile + * quarantines alone, exactly as it did before this field existed. Two profiles that set the + * same non-blank {@code credentialId} share one quarantine: a {@code BACKEND_EXHAUSTED} + * classification on either one quarantines both. + */ + public String effectiveCredentialId() { + return (credentialId == null || credentialId.isBlank()) ? profile : credentialId; + } + /** True when this profile's workers are granted a forge token to open their own PR (CB-302). */ public boolean hasGitToken() { return gitTokenEnv != null && !gitTokenEnv.isBlank(); @@ -850,7 +914,7 @@ public record BridgedConfig( private static final Set KNOWN_TOP_LEVEL_KEYS = Set.of( "bind", "herdrSocket", "profiles", "guard", "worktreeRoot", "lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet", - "leadHeartbeat", "health", "placement", "auth", "configReload"); + "leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds"); /** Load and validate config from {@code path}. */ public static BridgedConfig load(Path path) { @@ -1157,8 +1221,15 @@ public record BridgedConfig( // configReload is left as-is: null is "off", and ConfigReload's own compact constructor // defaults the fields of a block that IS present. Defaulting it here would start watching // the file for every config that never asked to be watched. + // quarantineCooldownSeconds IS defaulted, unlike leadHeartbeat/configReload above: it has no + // separate on/off switch of its own (BackendQuarantine only ever quarantines a credential + // after a BACKEND_EXHAUSTED classification, which stays opt-in via exhaustedPattern), so a + // config that never mentions it should still get a sane cooldown rather than a null one. + Integer quarantineCooldown = (quarantineCooldownSeconds != null && quarantineCooldownSeconds > 0) + ? quarantineCooldownSeconds : DEFAULT_QUARANTINE_COOLDOWN_SECONDS; return new BridgedConfig(b, herdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs, - broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload); + broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload, + quarantineCooldown); } /** diff --git a/bridged/src/main/java/dev/ltms/bridged/config/ConfigRef.java b/bridged/src/main/java/dev/ltms/bridged/config/ConfigRef.java index d145555..f802e42 100644 --- a/bridged/src/main/java/dev/ltms/bridged/config/ConfigRef.java +++ b/bridged/src/main/java/dev/ltms/bridged/config/ConfigRef.java @@ -37,10 +37,15 @@ import java.util.function.Supplier; * {@code fleet.leaders} needs a restart, the same as any deferred key below. *
  • Deferred — accepted into the new snapshot, but the wiring built at startup * keeps the old value until a restart: {@code lifecycle:}, {@code leadHeartbeat:}, - * {@code spawnReadyTimeoutMs} / {@code spawnReadyPollMs}, {@code guard:}, - * {@code worktreeRoot:}, adding or removing a profile (a new backend needs its own launcher, + * {@code spawnReadyTimeoutMs} / {@code spawnReadyPollMs}, {@code quarantineCooldownSeconds} + * (CB-578 stage B — baked once into the {@code BackendQuarantine} built at startup), + * {@code guard:}, {@code worktreeRoot:}, adding or removing a profile (a new backend needs its own launcher, * which is constructed once), and an existing profile's launch settings — - * {@code model}, {@code baseUrl}, {@code argv}, {@code env}, {@code mcpUrl} and the rest. + * {@code model}, {@code baseUrl}, {@code argv}, {@code env}, {@code mcpUrl}, + * {@code exhaustedPattern} (CB-578 stage A — compiled once into {@code Bridged.main}'s + * pattern map at startup), and the rest. {@code credentialId} (CB-578 stage B) is NOT on + * this list — it is read live off the config supplier at every quarantine check and + * exhaustion event, exactly like {@code weight} / {@code maxLoad}, so it is hot instead. * {@code HerdrPeerLauncher} takes {@code Map.copyOf(profiles)} at construction and resolves * each spawn out of that copy, so those never reach a launch until the daemon restarts. A * reload logs these rather than pretending they applied.
  • @@ -215,6 +220,12 @@ public final class ConfigRef implements Supplier { || !Objects.equals(old.spawnReadyPollMs(), fresh.spawnReadyPollMs())) { changed.add("spawnReady*"); } + // CB-578 stage B: baked once into the BackendQuarantine built at startup — a running + // quarantine keeps its original cooldown regardless, and a new cooldown only applies to a + // quarantine that starts after a restart. + if (!Objects.equals(old.quarantineCooldownSeconds(), fresh.quarantineCooldownSeconds())) { + changed.add("quarantineCooldownSeconds"); + } Map before = old.profiles() == null ? Map.of() : old.profiles(); Map after = @@ -250,8 +261,10 @@ public final class ConfigRef implements Supplier { /** * Whether two versions of a profile would launch a peer identically. Compares every component - * the launcher reads at spawn; {@code weight} and {@code maxLoad} are excluded because those are - * read live by the placement policy and really do take effect on the next spawn. + * the launcher reads at spawn; {@code weight}, {@code maxLoad} and {@code credentialId} are + * excluded because those are read live (by the placement policy and, for credentialId, by + * {@code CompositePeerLauncher}/the CB-578 stage B exhaustion sink) and really do take effect on + * the next spawn. */ private static boolean sameLaunchSettings(BridgedConfig.Profile a, BridgedConfig.Profile b) { return Objects.equals(a.baseUrl(), b.baseUrl()) @@ -269,6 +282,10 @@ public final class ConfigRef implements Supplier { && Objects.equals(a.gitHostEnv(), b.gitHostEnv()) && Objects.equals(a.kind(), b.kind()) && Objects.equals(a.env(), b.env()) - && Objects.equals(a.subscription(), b.subscription()); + && Objects.equals(a.subscription(), b.subscription()) + // CB-578 stage B: exhaustedPattern is compiled once into Bridged.main's pattern map + // at startup (see ExhaustedPatternLookup wiring) — a reload never re-reads it, so a + // changed pattern must be reported as deferred, exactly like model/baseUrl/argv. + && Objects.equals(a.exhaustedPattern(), b.exhaustedPattern()); } } diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java index 82c0b6a..9a44702 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java @@ -67,6 +67,7 @@ public final class CompletionResolver implements TurnListener { private final AgentControl agents; private final Rendezvous rendezvous; private final ExhaustedPatternLookup exhaustedPatterns; + private final ExhaustionSink exhaustionSink; /** * Per-target record of the turn currently in flight: the exact {@link Rendezvous} waiter its @@ -89,11 +90,17 @@ public final class CompletionResolver implements TurnListener { * usage-limit refusal pattern. Required — there is deliberately no * defaulting overload; a caller that does not want the classification * must pass an explicit inert value ({@link ExhaustedPatternLookup#none()}). + * @param exhaustionSink CB-578 stage B: notified when a {@code BACKEND_EXHAUSTED} + * classification actually resolves a waiter. Required for the same + * reason as {@code exhaustedPatterns} — pass {@link ExhaustionSink#none()} + * to opt out. */ - public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns) { + public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns, + ExhaustionSink exhaustionSink) { this.agents = agents; this.rendezvous = rendezvous; this.exhaustedPatterns = Objects.requireNonNull(exhaustedPatterns, "exhaustedPatterns"); + this.exhaustionSink = Objects.requireNonNull(exhaustionSink, "exhaustionSink"); } @Override @@ -210,6 +217,9 @@ public final class CompletionResolver implements TurnListener { inFlight.remove(target, turn); log.warn("completion for {} classified BACKEND_EXHAUSTED (no bridge_reply; scrape " + "matched the profile's exhausted pattern): {}", target, reason); + // CB-578 stage B: only on the resolution that actually won the race — a late + // duplicate must never quarantine a credential twice for one refusal. + exhaustionSink.onExhausted(target, reason); } return; } diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/ExhaustionSink.java b/bridged/src/main/java/dev/ltms/bridged/inject/ExhaustionSink.java new file mode 100644 index 0000000..d7b35a0 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/inject/ExhaustionSink.java @@ -0,0 +1,31 @@ +package dev.ltms.bridged.inject; + +/** + * Notified when {@link CompletionResolver} actually delivers a {@code BACKEND_EXHAUSTED} + * classification to a waiting send (CB-578 stage B) — never on a race that lost (see + * {@link CompletionResolver#resolve}, which only calls this after + * {@code Rendezvous.resolveExhausted} returns {@code true}). + * + *

    {@link CompletionResolver} knows only {@code target} (a herdr terminal id); it has no notion of + * profiles or credentials, so mapping {@code target} to whatever should be quarantined is entirely + * the sink's job — see {@code Bridged.main}'s wiring, which resolves target → session → profile → + * {@code effectiveCredentialId()} and calls {@code BackendQuarantine.quarantine} on it. + */ +@FunctionalInterface +public interface ExhaustionSink { + + /** + * @param target the herdr terminal id whose turn was classified {@code BACKEND_EXHAUSTED} + * @param reason the matched-line reason carried by the classification + */ + void onExhausted(String target, String reason); + + /** + * Inert sink — nothing happens on exhaustion. The explicit stand-in a caller (or a test not + * exercising this feature) passes instead of a defaulting overload, exactly like + * {@link ExhaustedPatternLookup#none()}. + */ + static ExhaustionSink none() { + return (target, reason) -> { }; + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java index da37ea9..0527bc7 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java @@ -14,6 +14,7 @@ import dev.ltms.bridged.herdr.HerdrException; import dev.ltms.bridged.msg.MessageService; import dev.ltms.bridged.msg.Rendezvous; import dev.ltms.bridged.peer.PeerUnreachableException; +import dev.ltms.bridged.placement.BackendQuarantine; import dev.ltms.bridged.session.SessionManager; import dev.ltms.bridged.session.MemberSession; import dev.ltms.bridged.session.WorktreeRequest; @@ -33,6 +34,7 @@ import jakarta.servlet.http.HttpServlet; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.Objects; import java.util.Set; import java.util.function.Function; import java.util.function.LongSupplier; @@ -78,6 +80,7 @@ public final class BridgeMcp { private final Metrics metrics; // CB-502: null → auth failures not counted private final CapacitySource capacity; private final HealthCoverageSource healthCoverage; + private final QuarantineSource quarantine; /** Capacity facts used by {@code bridge_list}; production must supply the placement live count. */ public record CapacitySource(Function liveCount, Function maxLoad, @@ -90,17 +93,30 @@ public final class BridgeMcp { /** Coverage is supplied by the health wiring, not inferred from a missing dependency. */ public record HealthCoverageSource(Supplier value) { } + /** + * CB-578 stage B quarantine facts used by {@code bridge_profiles}: a profile → credential id + * lookup, plus the shared {@link BackendQuarantine} to read remaining cooldowns off. + */ + public record QuarantineSource(Function credentialIdFor, BackendQuarantine quarantine) { + /** Inert source — no profile is ever reported quarantined. Explicit stand-in, not a default. */ + public static QuarantineSource none() { return new QuarantineSource(_ -> null, BackendQuarantine.none()); } + } + /** * @param callers resolves each call's {@link Principal}; {@code null} disables authorization. * This surface needs its own enforcement: {@code /mcp} is a raw servlet on * Jetty's context handler and never passes through Javalin's {@code before} * filter, so the REST guard does not cover it. * @param metrics registry for auth-failure counting; may be {@code null} + * @param quarantine CB-578 stage B facts for {@code bridge_profiles}; required — pass + * {@link QuarantineSource#none()} for a caller that does not want the feature */ public BridgeMcp(MessageService messages, PeerLauncher workers, SessionManager sessions, ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry, - CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage) { + CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage, + QuarantineSource quarantine) { this.capacity = capacity; + this.quarantine = Objects.requireNonNull(quarantine, "quarantine"); this.healthCoverage = healthCoverage; McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get(); this.transport = HttpServletStreamableServerTransportProvider.builder() @@ -228,7 +244,7 @@ public final class BridgeMcp { .toolCall(profilesTool(), (exchange, _) -> { McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null); if (denied != null) return denied; - return profiles(workers); + return profiles(workers, quarantine); }) .toolCall(whoamiTool(), (exchange, _) -> { McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null); @@ -706,11 +722,33 @@ public final class BridgeMcp { return null; } - /** {@code bridge_profiles}: the configured worker profiles and the default. */ - static McpSchema.CallToolResult profiles(PeerLauncher workers) { - return text(json(Map.of( - "profiles", workers.profiles(), - "default", workers.defaultProfile() == null ? "" : workers.defaultProfile()))); + /** + * {@code bridge_profiles}: the configured worker profiles, the default, and — CB-578 stage B — + * which of them are currently quarantined (backend exhausted) and for how much longer. The + * {@code quarantined} key is present only when at least one profile is, so a fleet where nothing + * has ever been quarantined gets exactly the pre-stage-B shape. + */ + static McpSchema.CallToolResult profiles(PeerLauncher workers, QuarantineSource quarantine) { + Map result = new LinkedHashMap<>(); + result.put("profiles", workers.profiles()); + result.put("default", workers.defaultProfile() == null ? "" : workers.defaultProfile()); + Map quarantined = new LinkedHashMap<>(); + for (String profile : workers.profiles()) { + String credentialId = quarantine.credentialIdFor().apply(profile); + if (credentialId == null) { + continue; + } + quarantine.quarantine().remainingSeconds(credentialId).ifPresent(remaining -> { + Map row = new LinkedHashMap<>(); + row.put("credentialId", credentialId); + row.put("quarantinedForSeconds", remaining); + quarantined.put(profile, row); + }); + } + if (!quarantined.isEmpty()) { + result.put("quarantined", quarantined); + } + return text(json(result)); } /** @@ -940,7 +978,10 @@ public final class BridgeMcp { private static McpSchema.Tool profilesTool() { return tool("bridge_profiles", - "List the configured worker profiles (backends) and which one bridge_spawn uses by default.", + "List the configured worker profiles (backends) and which one bridge_spawn uses by " + + "default. A 'quarantined' map is present when a backend-exhausted refusal put " + + "a profile's credential on cooldown — bridge_spawn onto it is refused until " + + "quarantinedForSeconds elapses; a profile sharing that credential is listed too.", objectSchema(Map.of(), List.of())); } diff --git a/bridged/src/main/java/dev/ltms/bridged/member/CompositePeerLauncher.java b/bridged/src/main/java/dev/ltms/bridged/member/CompositePeerLauncher.java index 605e0c2..73e08b4 100644 --- a/bridged/src/main/java/dev/ltms/bridged/member/CompositePeerLauncher.java +++ b/bridged/src/main/java/dev/ltms/bridged/member/CompositePeerLauncher.java @@ -8,6 +8,7 @@ import dev.ltms.bridged.peer.PeerHandle; import dev.ltms.bridged.peer.PeerLauncher; import dev.ltms.bridged.peer.PeerUnreachableException; import dev.ltms.bridged.peer.SpawnRequest; +import dev.ltms.bridged.placement.BackendQuarantine; import dev.ltms.bridged.placement.PlacementCandidate; import dev.ltms.bridged.placement.PlacementContext; import dev.ltms.bridged.placement.PlacementException; @@ -23,10 +24,12 @@ import java.util.HashSet; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.Objects; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.function.Function; import java.util.function.Supplier; +import java.util.stream.Collectors; /** * The {@link PeerLauncher} the core actually holds when more than one adapter is configured — a thin @@ -78,6 +81,9 @@ public final class CompositePeerLauncher implements PeerLauncher { private final Supplier> profileConfigs; private final Supplier placementPolicy; + /** CB-578 stage B: credential cooldown, checked before an explicit spawn and filtered into placement. */ + private final BackendQuarantine quarantine; + /** * CB-557: the role pools an unqualified spawn draws its candidates from. A supplier that yields * {@code null}, and an empty pool for a role, both fall back to every configured profile — the @@ -99,7 +105,10 @@ public final class CompositePeerLauncher implements PeerLauncher { } /** - * Production constructor with a placement policy and live-worker counter. + * Production constructor with a placement policy and live-worker counter. Quarantine (CB-578 + * stage B) is off for this constructor — {@link BackendQuarantine#none()} — since it predates + * the feature and existing callers of this exact overload never exercised it; use the 7-arg + * overload below to wire a real {@link BackendQuarantine}. * * @param delegates one adapter per configured peer kind; must be non-empty and declare * disjoint profile-name sets @@ -113,13 +122,14 @@ public final class CompositePeerLauncher implements PeerLauncher { Map profileConfigs, PlacementPolicy placementPolicy, Function liveCount) { - this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, null); + this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, null, BackendQuarantine.none()); } /** * Production constructor with role pools (CB-557). An unqualified spawn draws its candidates from * {@code fleet.} instead of from every configured profile, so a reviewer is placed on a - * reviewer backend and never on, say, the architect-only one. + * reviewer backend and never on, say, the architect-only one. Quarantine is off for this + * constructor too, for the same reason as the 5-arg overload above. * * @param fleet the configured role pools; {@code null} ⇒ every profile is a candidate for every * role, which is the pre-CB-557 behaviour @@ -130,29 +140,50 @@ public final class CompositePeerLauncher implements PeerLauncher { PlacementPolicy placementPolicy, Function liveCount, BridgedConfig.Fleet fleet) { + this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, fleet, BackendQuarantine.none()); + } + + /** + * Production constructor with role pools and quarantine (CB-578 stage B). The full-featured + * non-reloading form; {@link #CompositePeerLauncher(List, String, Supplier, Function, BackendQuarantine)} + * is what {@code Bridged.main} actually wires up. + * + * @param quarantine required — pass {@link BackendQuarantine#none()} for a caller that does not + * want the feature, never a defaulting overload (CB-578 stage B's own rule). + */ + public CompositePeerLauncher(List delegates, + String defaultProfile, + Map profileConfigs, + PlacementPolicy placementPolicy, + Function liveCount, + BridgedConfig.Fleet fleet, + BackendQuarantine quarantine) { // LinkedHashMap, not Map.copyOf: candidates() promises definition order and the weighted // policy breaks exact-weight ties on it, so a salted iteration order would make placement // differ from one JVM run to the next. this(delegates, defaultProfile, constant(Collections.unmodifiableMap(new LinkedHashMap<>(profileConfigs))), - constant(placementPolicy), liveCount, constant(fleet)); + constant(placementPolicy), liveCount, constant(fleet), quarantine); } /** * Production constructor that re-reads its placement inputs per spawn (CB-559), so a config * reload retargets the next member without a restart. * - * @param config the live configuration — read at every spawn, never captured + * @param config the live configuration — read at every spawn, never captured + * @param quarantine required — CB-578 stage B; pass {@link BackendQuarantine#none()} to opt out */ public CompositePeerLauncher(List delegates, String defaultProfile, Supplier config, - Function liveCount) { + Function liveCount, + BackendQuarantine quarantine) { this(delegates, defaultProfile, () -> config.get().profiles(), () -> PlacementPolicies.fromName(config.get().placement()), liveCount, - () -> config.get().fleet()); + () -> config.get().fleet(), + quarantine); } /** The all-suppliers form every other constructor funnels into. */ @@ -161,8 +192,10 @@ public final class CompositePeerLauncher implements PeerLauncher { Supplier> profileConfigs, Supplier placementPolicy, Function liveCount, - Supplier fleet) { + Supplier fleet, + BackendQuarantine quarantine) { this.fleet = fleet; + this.quarantine = Objects.requireNonNull(quarantine, "quarantine"); if (delegates.isEmpty()) { throw new IllegalArgumentException("at least one peer adapter must be configured"); } @@ -225,6 +258,7 @@ public final class CompositePeerLauncher implements PeerLauncher { // the charter makes explicit-profile spawns the normal path — so skipping the check // here would leave the cap dead config in real operation. HerdrPeerLauncher d = route(requestedProfile); + enforceNotQuarantined(requestedProfile); enforceMaxLoad(requestedProfile); PeerHandle handle = d.spawn(req); spawnedBy.put(handle.id(), d); @@ -238,7 +272,10 @@ public final class CompositePeerLauncher implements PeerLauncher { List candidates = candidates(req.role()); String roleDefault = defaultProfileFor(req.role()); Set unreachable = new HashSet<>(); - PlacementContext ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable); + // CB-578 stage B: computed once up front — a quarantine's expiry cannot pass within one spawn + // call, so re-deriving it per retry would only cost work, never change the answer. + Set quarantined = quarantinedProfiles(candidates); + PlacementContext ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable, quarantined); int maxAttempts = candidates.isEmpty() ? 1 : candidates.size(); for (int attempt = 0; attempt < maxAttempts; attempt++) { @@ -268,7 +305,7 @@ public final class CompositePeerLauncher implements PeerLauncher { chosen.profile(), e.getMessage()); unreachable.add(chosen.profile()); // Update the context for the next selection so the policy excludes this profile. - ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable); + ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable, quarantined); } } @@ -298,6 +335,42 @@ public final class CompositePeerLauncher implements PeerLauncher { * @param profile the profile the caller explicitly named * @throws PlacementException when the profile is at capacity */ + /** + * Refuse an explicit-profile spawn whose credential is quarantined (CB-578 stage B): a prior + * {@code BACKEND_EXHAUSTED} classification on this profile, or on another profile sharing its + * {@code credentialId}, is still on cooldown. + * + *

    Checked before {@link #enforceMaxLoad}, and for the same reason that check exists: an + * explicit profile bypasses the placement policy's filtering entirely, so without this the cap + * (there, quarantine here) would be dead config on the very path the charter calls normal. + * Deliberately no fallback to another profile, matching {@link #enforceMaxLoad}'s own reasoning — + * the caller named this profile for a cost/model reason. + * + * @throws PlacementException naming the profile, its credential, and the remaining cooldown + */ + private void enforceNotQuarantined(String profile) { + String credentialId = credentialIdFor(profile); + quarantine.remainingSeconds(credentialId).ifPresent(remaining -> { + throw new PlacementException("worker profile '" + profile + "' is quarantined " + + "(credential '" + credentialId + "' exhausted; ~" + remaining + + "s remaining) — refusing spawn"); + }); + } + + /** {@code profile}'s credential group (CB-578 stage B), or the profile's own name if unconfigured. */ + private String credentialIdFor(String profile) { + BridgedConfig.Profile cfg = profiles0().get(profile); + return cfg == null ? profile : cfg.effectiveCredentialId(); + } + + /** The subset of {@code candidates} whose credential is currently quarantined (CB-578 stage B). */ + private Set quarantinedProfiles(List candidates) { + return candidates.stream() + .map(PlacementCandidate::profile) + .filter(p -> quarantine.isQuarantined(credentialIdFor(p))) + .collect(Collectors.toSet()); + } + private void enforceMaxLoad(String profile) { // Absent config, or a config whose maxLoad normalized to null (non-positive ⇒ unlimited at // load), means no cap — never cap what wasn't configured. diff --git a/bridged/src/main/java/dev/ltms/bridged/placement/BackendQuarantine.java b/bridged/src/main/java/dev/ltms/bridged/placement/BackendQuarantine.java new file mode 100644 index 0000000..1917a58 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/placement/BackendQuarantine.java @@ -0,0 +1,99 @@ +package dev.ltms.bridged.placement; + +import java.util.LinkedHashMap; +import java.util.Map; +import java.util.Objects; +import java.util.OptionalLong; +import java.util.concurrent.ConcurrentHashMap; +import java.util.function.LongSupplier; + +/** + * Where a credential (not a profile — see {@code BridgedConfig.Profile#effectiveCredentialId()}) + * sits out a cooldown after a {@code BACKEND_EXHAUSTED} classification (CB-578 stage B), so a fresh + * spawn does not walk straight back onto the account that just refused on a usage limit. + * + *

    Keyed by credential id, never by profile name: two profiles sharing one credential (e.g. two + * models on the same OpenAI account) share one quarantine — {@link #quarantine} one credential id + * and every profile whose {@code effectiveCredentialId()} equals it is quarantined too, without this + * class knowing anything about profiles at all. That mapping is the caller's job (see + * {@code CompositePeerLauncher} and {@code dev.ltms.bridged.inject.ExhaustionSink}). + * + *

    The clock is injected ({@link LongSupplier}, conventionally {@code System::nanoTime} like + * {@code FleetHealthMonitor}), never read inline, so a quarantine's expiry is testable without a + * real sleep. + */ +public final class BackendQuarantine { + + private final ConcurrentHashMap quarantinedUntilNanos = new ConcurrentHashMap<>(); + private final LongSupplier nowNanos; + private final long cooldownNanos; + + /** + * @param nowNanos monotonic clock, injected for testability + * @param cooldownNanos how long a fresh {@link #quarantine} call blocks the credential for; + * must be positive + */ + public BackendQuarantine(LongSupplier nowNanos, long cooldownNanos) { + this.nowNanos = Objects.requireNonNull(nowNanos, "nowNanos"); + if (cooldownNanos <= 0) { + throw new IllegalArgumentException("cooldownNanos must be positive: " + cooldownNanos); + } + this.cooldownNanos = cooldownNanos; + } + + /** + * Inert quarantine — nothing is ever quarantined unless {@link #quarantine} is actually called on + * this instance. The explicit stand-in a caller (or a test not exercising this feature) passes + * instead of a defaulting overload, exactly like {@code ExhaustedPatternLookup.none()}. + */ + public static BackendQuarantine none() { + return new BackendQuarantine(() -> 0L, 1); + } + + /** + * Quarantine {@code credentialId} for the configured cooldown, starting now. A repeat call while + * already quarantined restarts the cooldown at full length — a fresh refusal is fresh evidence the + * account is still exhausted, not a reason to let an earlier, shorter wait stand. + */ + public void quarantine(String credentialId) { + Objects.requireNonNull(credentialId, "credentialId"); + quarantinedUntilNanos.put(credentialId, nowNanos.getAsLong() + cooldownNanos); + } + + /** Whether {@code credentialId} is quarantined right now. */ + public boolean isQuarantined(String credentialId) { + return remainingNanos(credentialId) > 0; + } + + /** Seconds left on {@code credentialId}'s quarantine, or empty when it is not quarantined. */ + public OptionalLong remainingSeconds(String credentialId) { + long remaining = remainingNanos(credentialId); + return remaining > 0 ? OptionalLong.of(toSecondsRoundedUp(remaining)) : OptionalLong.empty(); + } + + /** + * Every currently-quarantined credential id and its remaining seconds (CB-578 stage B fleet + * reporting) — expired entries are never included. Not pruned from the backing map here: it stays + * small (bounded by the number of distinct credentials ever exhausted) and a lazily-stale entry is + * harmless, since every read already checks the deadline. + */ + public Map activeRemainingSeconds() { + Map out = new LinkedHashMap<>(); + quarantinedUntilNanos.forEach((credentialId, deadline) -> { + long remaining = deadline - nowNanos.getAsLong(); + if (remaining > 0) { + out.put(credentialId, toSecondsRoundedUp(remaining)); + } + }); + return out; + } + + private long remainingNanos(String credentialId) { + Long deadline = quarantinedUntilNanos.get(credentialId); + return deadline == null ? 0L : deadline - nowNanos.getAsLong(); + } + + private static long toSecondsRoundedUp(long nanos) { + return (nanos + 999_999_999L) / 1_000_000_000L; + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/placement/FixedPlacementPolicy.java b/bridged/src/main/java/dev/ltms/bridged/placement/FixedPlacementPolicy.java index e0896ce..995ca15 100644 --- a/bridged/src/main/java/dev/ltms/bridged/placement/FixedPlacementPolicy.java +++ b/bridged/src/main/java/dev/ltms/bridged/placement/FixedPlacementPolicy.java @@ -4,18 +4,28 @@ package dev.ltms.bridged.placement; * Backward-compatible placement: an unqualified spawn always resolves to the configured default * profile, exactly as {@code CompositePeerLauncher} did before CB-518. This ignores caps and * reachability so that a pre-existing config behaves identically after upgrade. + * + *

    Quarantine (CB-578 stage B) is the one exception: a quarantined default is a credential that + * just refused on a usage limit, not a transient capacity or reachability concern, so {@code fixed} + * steps to the first non-quarantined candidate instead of walking straight back onto it. A fleet + * where nothing is ever quarantined never exercises this path, so today's behaviour is unchanged. */ final class FixedPlacementPolicy implements PlacementPolicy { @Override public PlacementCandidate select(PlacementContext ctx) { String d = ctx.defaultProfile(); - if (d != null && !d.isBlank()) { + if (d != null && !d.isBlank() && !ctx.quarantined().contains(d)) { return new PlacementCandidate(d, null, 1.0f, null); } - if (!ctx.candidates().isEmpty()) { - PlacementCandidate first = ctx.candidates().getFirst(); - return new PlacementCandidate(first.profile(), null, first.weight(), first.maxLoad()); + for (PlacementCandidate c : ctx.candidates()) { + if (!ctx.quarantined().contains(c.profile())) { + return new PlacementCandidate(c.profile(), null, c.weight(), c.maxLoad()); + } + } + if (d != null && !d.isBlank()) { + throw new PlacementException("worker profile '" + d + "' is quarantined (backend " + + "exhausted) and no un-quarantined candidate is available"); } throw new PlacementException("no worker profiles configured"); } diff --git a/bridged/src/main/java/dev/ltms/bridged/placement/PlacementContext.java b/bridged/src/main/java/dev/ltms/bridged/placement/PlacementContext.java index 068f526..ad73c16 100644 --- a/bridged/src/main/java/dev/ltms/bridged/placement/PlacementContext.java +++ b/bridged/src/main/java/dev/ltms/bridged/placement/PlacementContext.java @@ -11,9 +11,13 @@ import java.util.function.Function; * @param candidates every configured candidate; the policy filters out those at cap or unreachable * @param liveCount current live worker count per profile (from the session registry) * @param unreachable profiles already known to have failed in this spawn attempt + * @param quarantined profiles whose credential is currently quarantined (CB-578 stage B) — a + * {@code BACKEND_EXHAUSTED} classification put it, or a profile it shares a + * credential with, on cooldown. Filtered the same way as {@code unreachable}. */ public record PlacementContext(String defaultProfile, List candidates, Function liveCount, - Set unreachable) { + Set unreachable, + Set quarantined) { } diff --git a/bridged/src/main/java/dev/ltms/bridged/placement/PlacementPolicyUtil.java b/bridged/src/main/java/dev/ltms/bridged/placement/PlacementPolicyUtil.java index 1e32f8b..66f0ca6 100644 --- a/bridged/src/main/java/dev/ltms/bridged/placement/PlacementPolicyUtil.java +++ b/bridged/src/main/java/dev/ltms/bridged/placement/PlacementPolicyUtil.java @@ -12,13 +12,13 @@ final class PlacementPolicyUtil { } /** - * Candidates that are not known-unreachable and have not reached their maxLoad. - * A {@code null} maxLoad means unlimited. + * Candidates that are not known-unreachable, not quarantined (CB-578 stage B), and have not + * reached their maxLoad. A {@code null} maxLoad means unlimited. */ static List available(PlacementContext ctx) { List out = new ArrayList<>(); for (PlacementCandidate c : ctx.candidates()) { - if (ctx.unreachable().contains(c.profile())) { + if (ctx.unreachable().contains(c.profile()) || ctx.quarantined().contains(c.profile())) { continue; } Integer cap = c.maxLoad(); @@ -35,14 +35,17 @@ final class PlacementPolicyUtil { /** * Build a clear exception describing why every candidate was dropped: all at capacity, - * all unreachable, or a mix. + * all unreachable, all quarantined, or a mix. */ static PlacementException emptyException(PlacementContext ctx) { int atCap = 0; int unreachable = 0; + int quarantined = 0; for (PlacementCandidate c : ctx.candidates()) { Integer cap = c.maxLoad(); - if (ctx.unreachable().contains(c.profile())) { + if (ctx.quarantined().contains(c.profile())) { + quarantined++; + } else if (ctx.unreachable().contains(c.profile())) { unreachable++; } else if (cap != null && ctx.liveCount().apply(c.profile()) >= cap) { atCap++; @@ -53,6 +56,9 @@ final class PlacementPolicyUtil { if (total == 0) { return new PlacementException("no worker profiles configured"); } + if (quarantined == total) { + return new PlacementException("all worker profiles are quarantined (backend exhausted)"); + } if (atCap == total) { return new PlacementException("all worker profiles are at maxLoad"); } @@ -60,6 +66,7 @@ final class PlacementPolicyUtil { return new PlacementException("all worker profiles are unreachable"); } return new PlacementException("no worker profile available: " + atCap + " at maxLoad, " - + unreachable + " unreachable, " + (total - atCap - unreachable) + " remaining"); + + unreachable + " unreachable, " + quarantined + " quarantined, " + + (total - atCap - unreachable - quarantined) + " remaining"); } } diff --git a/bridged/src/test/java/dev/ltms/bridged/config/ConfigRefTest.java b/bridged/src/test/java/dev/ltms/bridged/config/ConfigRefTest.java index 735cd96..7b33a21 100644 --- a/bridged/src/test/java/dev/ltms/bridged/config/ConfigRefTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/config/ConfigRefTest.java @@ -338,6 +338,54 @@ class ConfigRefTest { assertEquals("deepseek-v4-flash", ref.get().profiles().get("sonnet").model()); } + /** + * CB-578 stage B: exhaustedPattern is compiled once into Bridged.main's pattern map at startup + * (see ExhaustedPatternLookup), so a reload never re-reads it — a changed pattern must be + * reported deferred exactly like model/baseUrl, not silently claimed as applied. + */ + @Test + void changingAProfilesExhaustedPatternIsReportedAsDeferred(@TempDir Path dir) throws Exception { + Path f = dir.resolve("bridged.yaml"); + Files.writeString(f, """ + bind: + host: 127.0.0.1 + port: 8765 + herdrSocket: ~/.config/herdr/herdr.sock + profiles: + sonnet: + baseUrl: http://gx00.gw:8000 + model: sonnet + exhaustedPattern: "usage limit has been reached" + guard: + offSubscriptionHosts: + - gx00.gw + """); + ConfigRef ref = refFor(f); + + Files.writeString(f, """ + bind: + host: 127.0.0.1 + port: 8765 + herdrSocket: ~/.config/herdr/herdr.sock + profiles: + sonnet: + baseUrl: http://gx00.gw:8000 + model: sonnet + exhaustedPattern: "rate limit exceeded" + guard: + offSubscriptionHosts: + - gx00.gw + """); + ConfigRef.Outcome out = ref.reload(); + + assertTrue(out.applied()); + assertEquals(1, out.deferred().size(), out.deferred().toString()); + assertTrue(out.deferred().getFirst().contains("sonnet"), out.deferred().toString()); + assertTrue(out.deferred().getFirst().contains("launch settings"), out.deferred().toString()); + // The snapshot still carries the new value — a restart is what makes it take effect. + assertEquals("rate limit exceeded", ref.get().profiles().get("sonnet").exhaustedPattern()); + } + @Test void aFixedRefHasNoFileAndRefusesToReload() { BridgedConfig cfg = new BridgedConfig(null, null, null, null, null, null, diff --git a/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java b/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java index a778b00..7c78236 100644 --- a/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java @@ -26,7 +26,7 @@ class CompletionResolverTest { void skipsTheScrapeWhenNoSendIsWaiting() { FakeHerdr herdr = new FakeHerdr(); Rendezvous rendezvous = new Rendezvous(); // no waiter opened - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); resolver.resolve("term_a", null); // no in-flight turn captured for this target @@ -38,7 +38,7 @@ class CompletionResolverTest { void failSkipsTheScrapeWhenNoSendIsWaiting() { FakeHerdr herdr = new FakeHerdr(); Rendezvous rendezvous = new Rendezvous(); // no waiter opened - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); resolver.fail("term_a", null); // no in-flight turn, and no registered waiter to fall back to @@ -50,7 +50,7 @@ class CompletionResolverTest { void captureBaselineSkipsTheReadWhenNoSendIsWaiting() { FakeHerdr herdr = new FakeHerdr().readText("⏺ X\n❯ "); Rendezvous rendezvous = new Rendezvous(); // no waiter opened - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); resolver.captureBaseline("term_a", TestTurnTokens.inert("term_a")); // no send to attribute a later completion to @@ -139,7 +139,7 @@ class CompletionResolverTest { // send must NOT be resolved with the stale answer. FakeHerdr herdr = new FakeHerdr().readText("⏺ 391\n❯ "); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); var waiter = rendezvous.open("term_a"); // a send is blocked on this turn // The turn as captured at delivery: its waiter, and the previous turn's answer still on screen. @@ -154,7 +154,7 @@ class CompletionResolverTest { void resolvesACompletionWhoseScrapeChangedSinceDelivery() { FakeHerdr herdr = new FakeHerdr().readText("⏺ No, 391 = 17 × 23.\n❯ "); // the worker's real answer Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); var waiter = rendezvous.open("term_a"); // Delivery baseline was the previous turn's "391"; the scrape now differs → resolve. @@ -171,7 +171,7 @@ class CompletionResolverTest { String block = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 1) + "\n❯ "; FakeHerdr herdr = new FakeHerdr().readText(block); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); var waiter = rendezvous.open("term_a"); resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); @@ -185,7 +185,7 @@ class CompletionResolverTest { void leavesAnUnclippedCompletionPaneTailUnmarked() { FakeHerdr herdr = new FakeHerdr().readText("⏺ complete report\n❯ "); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); var waiter = rendezvous.open("term_a"); resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); @@ -197,7 +197,7 @@ class CompletionResolverTest { void resolvesSynchronouslyBeforePostTurnContextClearing() { FakeHerdr herdr = new FakeHerdr().readText("⏺ previous answer\n❯ "); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); var waiter = rendezvous.open("term_a"); resolver.captureBaseline("term_a", new TurnToken("term_a", waiter)); herdr.readText("⏺ answer that /clear would erase\n❯ "); @@ -219,7 +219,7 @@ class CompletionResolverTest { String longBlock = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 500) + "\n❯ "; FakeHerdr herdr = new FakeHerdr().readText(longBlock); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); var waiter = rendezvous.open("term_a"); // a send is blocked on this turn resolver.captureBaseline("term_a", new TurnToken("term_a", waiter)); // baseline is the clipped >cap block @@ -239,7 +239,7 @@ class CompletionResolverTest { // No delivery baseline (e.g. the pre-turn read failed) ⇒ never suppress; the completion resolves. FakeHerdr herdr = new FakeHerdr().readText("⏺ hello\n❯ "); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); var waiter = rendezvous.open("term_a"); resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); @@ -257,7 +257,7 @@ class CompletionResolverTest { // byte-identical guard would wrongly match the empty tail and suppress. FakeHerdr herdr = new FakeHerdr().healthy(false); // agent.read throws HerdrException Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); var waiter = rendezvous.open("term_a"); var turn = new CompletionResolver.InFlight(waiter, ""); // empty pane baselined at delivery @@ -277,7 +277,7 @@ class CompletionResolverTest { // fail must not overwrite that value, and must not even scrape the worker — nobody needs it. FakeHerdr herdr = new FakeHerdr().readText("an error screen"); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); var waiter = rendezvous.open("term_a"); var turn = new CompletionResolver.InFlight(waiter, null); @@ -298,7 +298,7 @@ class CompletionResolverTest { // fail falls back to the waiter currently registered on the Rendezvous and fails it. FakeHerdr herdr = new FakeHerdr().readText("stuck on an error screen"); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); var waiter = rendezvous.open("term_a"); // send registered, but no captureBaseline ever ran resolver.fail("term_a", null); // no in-flight turn → fall back to the registered waiter @@ -324,7 +324,7 @@ class CompletionResolverTest { try { FakeHerdr herdr = new FakeHerdr().readText("stuck on an error screen"); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); var waiter = rendezvous.open("term_a"); resolver.fail("term_a", null); @@ -352,7 +352,7 @@ class CompletionResolverTest { // scrape to turn N+1; targeting turn N's captured waiter makes the late completion a no-op. FakeHerdr herdr = new FakeHerdr().readText("⏺ turn N answer\n❯ "); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); var waiterN = rendezvous.open("term_a"); // turn N's send // The turn as the injector captured it at delivery (waiter + pre-turn baseline). @@ -382,7 +382,7 @@ class CompletionResolverTest { FakeHerdr herdr = new FakeHerdr().readText(block); Rendezvous rendezvous = new Rendezvous(); ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached"); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, ExhaustionSink.none()); var waiter = rendezvous.open("term_a"); resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); @@ -398,7 +398,7 @@ class CompletionResolverTest { FakeHerdr herdr = new FakeHerdr().readText(block); Rendezvous rendezvous = new Rendezvous(); ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached"); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, ExhaustionSink.none()); var waiter = rendezvous.open("term_a"); resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); @@ -407,12 +407,53 @@ class CompletionResolverTest { waiter.getNow(null).text(), "the reason names the real cause and carries the matched line"); } + @Test + void aWinningBackendExhaustedClassificationNotifiesTheExhaustionSink() { + String block = "⏺ Working on it...\nThe usage limit has been reached. Try again later.\n❯ "; + FakeHerdr herdr = new FakeHerdr().readText(block); + Rendezvous rendezvous = new Rendezvous(); + ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached"); + java.util.List notified = new java.util.ArrayList<>(); + ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink); + + var waiter = rendezvous.open("term_a"); + resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); + + assertEquals(1, notified.size(), "the sink is notified exactly once for the winning classification"); + assertTrue(notified.get(0).startsWith("term_a: "), "the sink is told which target exhausted"); + assertTrue(notified.get(0).contains("The usage limit has been reached"), + "the sink is told the matched reason: " + notified.get(0)); + } + + @Test + void aLosingBackendExhaustedClassificationNeverNotifiesTheExhaustionSink() { + // The waiter was already resolved (e.g. by the worker's own reply) before this scrape landed — + // resolveExhausted loses the race and must return false, so the sink must not fire either. + String block = "⏺ Working on it...\nThe usage limit has been reached. Try again later.\n❯ "; + FakeHerdr herdr = new FakeHerdr().readText(block); + Rendezvous rendezvous = new Rendezvous(); + ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached"); + java.util.List notified = new java.util.ArrayList<>(); + ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink); + + var waiter = rendezvous.open("term_a"); + var turn = new CompletionResolver.InFlight(waiter, null); + assertTrue(rendezvous.resolveCompletion(waiter, "already replied")); + + resolver.resolve("term_a", turn); + + assertTrue(notified.isEmpty(), "a classification that loses the race must not quarantine anything"); + assertEquals("already replied", waiter.getNow(null).text(), "the earlier resolution stands untouched"); + } + @Test void aNonMatchingScrapeResolvesAsAnOrdinaryCompletion() { FakeHerdr herdr = new FakeHerdr().readText("⏺ complete report\n❯ "); Rendezvous rendezvous = new Rendezvous(); ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached"); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, ExhaustionSink.none()); var waiter = rendezvous.open("term_a"); resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); @@ -430,7 +471,7 @@ class CompletionResolverTest { FakeHerdr herdr = new FakeHerdr().readText(block); Rendezvous rendezvous = new Rendezvous(); CompletionResolver resolver = - new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); + new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); var waiter = rendezvous.open("term_a"); resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); diff --git a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpAuthzTest.java b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpAuthzTest.java index 7b1af8b..523137d 100644 --- a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpAuthzTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpAuthzTest.java @@ -72,7 +72,8 @@ class BridgeMcpAuthzTest { new PrimaryRegistry(null), enforce ? CallerResolver.withLeadsAndMembers(identity, false, null, Map::of, new MemberRegistry(null)) : null, - metrics, BridgeMcp.CapacitySource.none(), new BridgeMcp.HealthCoverageSource(() -> "off")); + metrics, BridgeMcp.CapacitySource.none(), new BridgeMcp.HealthCoverageSource(() -> "off"), + BridgeMcp.QuarantineSource.none()); return mcp; } diff --git a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java index 43ba54d..7568244 100644 --- a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java @@ -16,6 +16,7 @@ import dev.ltms.bridged.inject.MemberPresence; import dev.ltms.bridged.session.MemberSession; import dev.ltms.bridged.session.WorktreeRequest; import dev.ltms.bridged.member.ClaudeCodeLauncher; +import dev.ltms.bridged.placement.BackendQuarantine; import io.modelcontextprotocol.spec.McpSchema; import dev.ltms.bridged.msg.InMemoryReplyInbox; import org.junit.jupiter.api.BeforeEach; @@ -437,11 +438,29 @@ class BridgeMcpTest { @Test void profilesListsConfiguredProfilesAndDefault() { FakeHerdr h = new FakeHerdr(); - McpSchema.CallToolResult res = BridgeMcp.profiles(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))); + McpSchema.CallToolResult res = BridgeMcp.profiles( + workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), BridgeMcp.QuarantineSource.none()); assertNotEquals(Boolean.TRUE, res.isError()); String out = textOf(res); assertTrue(out.contains("ltms-local"), out); assertTrue(out.contains("\"default\":\"ltms-local\""), out); + assertFalse(out.contains("quarantined"), "no profile is quarantined, so the key is omitted: " + out); + } + + @Test + void profilesReportsAQuarantinedCredential() { + FakeHerdr h = new FakeHerdr(); + BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30)); + quarantine.quarantine("shared-openai"); + BridgeMcp.QuarantineSource source = new BridgeMcp.QuarantineSource( + profile -> "ltms-local".equals(profile) ? "shared-openai" : null, quarantine); + McpSchema.CallToolResult res = BridgeMcp.profiles( + workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), source); + assertNotEquals(Boolean.TRUE, res.isError()); + String out = textOf(res); + assertTrue(out.contains("\"quarantined\""), out); + assertTrue(out.contains("shared-openai"), out); + assertTrue(out.contains("\"quarantinedForSeconds\":1800"), out); } @Test diff --git a/bridged/src/test/java/dev/ltms/bridged/member/CompositePeerLauncherTest.java b/bridged/src/test/java/dev/ltms/bridged/member/CompositePeerLauncherTest.java index b1f1994..cbc5078 100644 --- a/bridged/src/test/java/dev/ltms/bridged/member/CompositePeerLauncherTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/member/CompositePeerLauncherTest.java @@ -16,6 +16,7 @@ import dev.ltms.bridged.peer.PeerHandle; import dev.ltms.bridged.peer.PeerLauncher; import dev.ltms.bridged.peer.PeerUnreachableException; import dev.ltms.bridged.peer.SpawnRequest; +import dev.ltms.bridged.placement.BackendQuarantine; import dev.ltms.bridged.placement.PlacementException; import dev.ltms.bridged.placement.PlacementPolicies; import org.junit.jupiter.api.Test; @@ -27,6 +28,8 @@ import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Set; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; import static org.junit.jupiter.api.Assertions.*; @@ -128,6 +131,13 @@ class CompositePeerLauncherTest { weight, maxLoad); } + private static BridgedConfig.Profile stubWorker(String profile, String credentialId) { + return new BridgedConfig.Profile(profile, "http://gx00.gw:8000", "coder", + null, "BRIDGED_WORKER_TOKEN", List.of("claude"), "tab", "bridged-workers", + "w #{n}", null, null, null, null, null, null, null, null, null, + null, null, credentialId); + } + /** * An order-preserving profile map. Never {@code Map.of} here: its iteration order is * salted per JVM run, and the weighted policy breaks an exact-weight tie on candidate order — @@ -554,4 +564,105 @@ class CompositePeerLauncherTest { new SpawnRequest(null, null, null, null, null, MemberRole.DEV))); assertTrue(e.getMessage().contains("maxLoad"), e.getMessage()); } + + // ── CB-578 stage B: a BACKEND_EXHAUSTED classification quarantines the credential ────────── + + @Test + void explicitSpawnOntoAQuarantinedProfileIsRefusedNamingTheCredentialAndRemainingTime() { + FakeHerdr herdr = new FakeHerdr(); + Map profiles = ordered( + "sol", stubWorker("sol", "shared-openai"), + "terra", stubWorker("terra", "shared-openai")); + StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "sol", Set.of()); + BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30)); + quarantine.quarantine("shared-openai"); + CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles, + PlacementPolicies.fixed(), _ -> 0, null, quarantine); + + PlacementException e = assertThrows(PlacementException.class, + () -> composite.spawn(new SpawnRequest("sol", null, null))); + assertTrue(e.getMessage().contains("sol"), "message names the profile: " + e.getMessage()); + assertTrue(e.getMessage().contains("shared-openai"), "message names the credential: " + e.getMessage()); + assertTrue(e.getMessage().contains("1800"), "message names roughly when it lifts: " + e.getMessage()); + assertEquals(0, adapter.spawnCount("sol"), "the quarantined profile is never delegated to"); + } + + /** + * The part CB-578 stage B calls out as easy to get wrong: sol and terra are two different + * profiles sharing one OpenAI credential. Quarantining because of an exhaustion classified on + * ONE of them must lock out the other too, or the fleet just walks onto the same dead account + * under the sibling's name. + */ + @Test + void twoProfilesSharingACredentialAreBothQuarantinedByOneExhaustionEvent() { + FakeHerdr herdr = new FakeHerdr(); + Map profiles = ordered( + "sol", stubWorker("sol", "shared-openai"), + "terra", stubWorker("terra", "shared-openai")); + StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "sol", Set.of()); + BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30)); + // Only "sol" was classified BACKEND_EXHAUSTED — but the two profiles share one credential. + quarantine.quarantine("shared-openai"); + CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles, + PlacementPolicies.fixed(), _ -> 0, null, quarantine); + + assertThrows(PlacementException.class, () -> composite.spawn(new SpawnRequest("sol", null, null)), + "sol was the one classified exhausted"); + assertThrows(PlacementException.class, () -> composite.spawn(new SpawnRequest("terra", null, null)), + "terra shares sol's credential, so it must be locked out too"); + assertEquals(0, adapter.spawnCount("sol")); + assertEquals(0, adapter.spawnCount("terra")); + } + + @Test + void placementSkipsAQuarantinedProfileAndRoutesToAnUnquarantinedOne() { + FakeHerdr herdr = new FakeHerdr(); + Map profiles = ordered( + "sol", stubWorker("sol", "shared-openai"), + "b", stubWorker("b")); + StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "sol", Set.of()); + BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30)); + quarantine.quarantine("shared-openai"); + CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles, + PlacementPolicies.weighted(), _ -> 0, null, quarantine); + + PeerHandle h = composite.spawn(new SpawnRequest(null, null, null)); + assertEquals("b", h.profile(), "sol is quarantined, so an unqualified spawn must land on b"); + assertEquals(0, adapter.spawnCount("sol")); + } + + @Test + void aQuarantineLiftsOnTheInjectedClockAndTheProfileBecomesSpawnableAgain() { + FakeHerdr herdr = new FakeHerdr(); + Map profiles = ordered( + "sol", stubWorker("sol", "shared-openai"), + "b", stubWorker("b")); + StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "sol", Set.of()); + AtomicLong nowNanos = new AtomicLong(0L); + BackendQuarantine quarantine = new BackendQuarantine(nowNanos::get, TimeUnit.MINUTES.toNanos(30)); + quarantine.quarantine("shared-openai"); + CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles, + PlacementPolicies.fixed(), _ -> 0, null, quarantine); + + assertThrows(PlacementException.class, () -> composite.spawn(new SpawnRequest("sol", null, null)), + "still inside the cooldown"); + + nowNanos.set(TimeUnit.MINUTES.toNanos(31)); + + PeerHandle h = composite.spawn(new SpawnRequest("sol", null, null)); + assertEquals("sol", h.profile(), "the cooldown expired on the injected clock — sol is spawnable again"); + assertEquals(1, adapter.spawnCount("sol")); + } + + @Test + void aFleetWithNoExhaustedPatternAnywhereBehavesExactlyAsBeforeQuarantineExisted() { + // BackendQuarantine.none() is the inert stand-in every constructor already defaults to when + // no quarantine is wired — the 2/5/6-arg constructors used throughout this file all exercise + // it. This test pins that an explicit .none() also never refuses a spawn, for any profile. + FakeHerdr herdr = new FakeHerdr(); + PeerLauncher composite = composite(herdr); + + assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("claude", null, null))); + assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("gemini", null, null))); + } } diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java index eb5de30..0f4fe23 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -6,6 +6,7 @@ import dev.ltms.bridged.herdr.FakeHerdr; import dev.ltms.bridged.herdr.HerdrException; import dev.ltms.bridged.inject.CompletionResolver; import dev.ltms.bridged.inject.ExhaustedPatternLookup; +import dev.ltms.bridged.inject.ExhaustionSink; import dev.ltms.bridged.mcp.PrimaryRegistry; import dev.ltms.bridged.inject.Injector; import org.junit.jupiter.api.BeforeEach; @@ -33,7 +34,7 @@ class MessageServiceTest { private final AgentControl agents = new AgentControl(herdr); private final Rendezvous rendezvous = new Rendezvous(); private final CompletionResolver completion = - new CompletionResolver(agents, rendezvous, ExhaustedPatternLookup.none()); + new CompletionResolver(agents, rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); private final Injector injector = new Injector(agents, completion); private final InMemoryReplyInbox inbox = new InMemoryReplyInbox(); private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox); diff --git a/bridged/src/test/java/dev/ltms/bridged/placement/BackendQuarantineTest.java b/bridged/src/test/java/dev/ltms/bridged/placement/BackendQuarantineTest.java new file mode 100644 index 0000000..de10c0d --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/placement/BackendQuarantineTest.java @@ -0,0 +1,108 @@ +package dev.ltms.bridged.placement; + +import org.junit.jupiter.api.Test; + +import java.util.Map; +import java.util.OptionalLong; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * CB-578 stage B: the credential-keyed quarantine tracker itself, isolated from placement/spawn + * wiring (that's {@code CompositePeerLauncherTest}). The clock is a plain {@link AtomicLong} of + * nanos so expiry is exercised without a real sleep. + */ +class BackendQuarantineTest { + + @Test + void aFreshCredentialIsNotQuarantined() { + BackendQuarantine q = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30)); + assertFalse(q.isQuarantined("shared-openai")); + assertEquals(OptionalLong.empty(), q.remainingSeconds("shared-openai")); + } + + @Test + void quarantineBlocksTheCredentialForTheFullCooldown() { + BackendQuarantine q = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30)); + q.quarantine("shared-openai"); + + assertTrue(q.isQuarantined("shared-openai")); + assertEquals(OptionalLong.of(1800L), q.remainingSeconds("shared-openai")); + } + + @Test + void onlyTheQuarantinedCredentialIsAffected() { + BackendQuarantine q = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30)); + q.quarantine("shared-openai"); + + assertFalse(q.isQuarantined("some-other-credential"), + "an unrelated credential must not be swept into the quarantine"); + } + + @Test + void expiresOnTheInjectedClock() { + AtomicLong now = new AtomicLong(0L); + BackendQuarantine q = new BackendQuarantine(now::get, TimeUnit.MINUTES.toNanos(30)); + q.quarantine("shared-openai"); + assertTrue(q.isQuarantined("shared-openai")); + + now.set(TimeUnit.MINUTES.toNanos(29)); + assertTrue(q.isQuarantined("shared-openai"), "still inside the cooldown"); + + now.set(TimeUnit.MINUTES.toNanos(31)); + assertFalse(q.isQuarantined("shared-openai"), "the cooldown has elapsed on the injected clock"); + assertEquals(OptionalLong.empty(), q.remainingSeconds("shared-openai")); + } + + @Test + void aRepeatQuarantineCallRestartsTheCooldownAtFullLength() { + AtomicLong now = new AtomicLong(0L); + BackendQuarantine q = new BackendQuarantine(now::get, TimeUnit.MINUTES.toNanos(30)); + q.quarantine("shared-openai"); + + now.set(TimeUnit.MINUTES.toNanos(20)); + q.quarantine("shared-openai"); + + now.set(TimeUnit.MINUTES.toNanos(45)); // 25 min after the second call, 45 after the first + assertTrue(q.isQuarantined("shared-openai"), + "a fresh exhaustion resets the cooldown to full length, not the earlier shorter wait"); + } + + @Test + void activeRemainingSecondsListsOnlyStillQuarantinedCredentials() { + AtomicLong now = new AtomicLong(0L); + BackendQuarantine q = new BackendQuarantine(now::get, TimeUnit.MINUTES.toNanos(30)); + q.quarantine("shared-openai"); + q.quarantine("another-credential"); + + now.set(TimeUnit.MINUTES.toNanos(31)); + q.quarantine("shared-openai"); // re-quarantined after the first one expired + + Map active = q.activeRemainingSeconds(); + assertEquals(Map.of("shared-openai", 1800L), active, + "the expired credential is dropped; the re-quarantined one is reported"); + } + + @Test + void noneReportsNothingQuarantinedWhenNeverToldTo() { + // .none() is the stand-in for a caller whose code path never calls #quarantine at all (e.g. + // the 2/5/6-arg CompositePeerLauncher constructors) — not a guarantee that a call to + // #quarantine on it is a no-op. Left alone, as those call sites leave it, nothing is ever + // quarantined. + BackendQuarantine q = BackendQuarantine.none(); + + assertFalse(q.isQuarantined("anything")); + assertTrue(q.activeRemainingSeconds().isEmpty()); + } + + @Test + void aNonPositiveCooldownIsRejected() { + assertThrows(IllegalArgumentException.class, () -> new BackendQuarantine(() -> 0L, 0L)); + assertThrows(IllegalArgumentException.class, () -> new BackendQuarantine(() -> 0L, -1L)); + } +} diff --git a/bridged/src/test/java/dev/ltms/bridged/placement/PlacementPolicyTest.java b/bridged/src/test/java/dev/ltms/bridged/placement/PlacementPolicyTest.java index 105056b..9c8a1c3 100644 --- a/bridged/src/test/java/dev/ltms/bridged/placement/PlacementPolicyTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/placement/PlacementPolicyTest.java @@ -24,7 +24,13 @@ class PlacementPolicyTest { private static PlacementContext ctx(List candidates, Function liveCount, Set unreachable) { - return new PlacementContext("b", candidates, liveCount, unreachable); + return ctx(candidates, liveCount, unreachable, Set.of()); + } + + private static PlacementContext ctx(List candidates, + Function liveCount, + Set unreachable, Set quarantined) { + return new PlacementContext("b", candidates, liveCount, unreachable, quarantined); } private static PlacementContext ctx(List candidates, @@ -46,17 +52,37 @@ class PlacementPolicyTest { PlacementPolicy policy = PlacementPolicies.fixed(); PlacementContext ctx = new PlacementContext(null, List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")), - noSessions(), Set.of()); + noSessions(), Set.of(), Set.of()); assertEquals("a", policy.select(ctx).profile()); } @Test void fixedThrowsWhenNoProfilesAndNoDefault() { PlacementPolicy policy = PlacementPolicies.fixed(); - PlacementContext ctx = new PlacementContext(null, List.of(), noSessions(), Set.of()); + PlacementContext ctx = new PlacementContext(null, List.of(), noSessions(), Set.of(), Set.of()); assertThrows(PlacementException.class, () -> policy.select(ctx)); } + @Test + void fixedSkipsQuarantinedDefault() { + PlacementPolicy policy = PlacementPolicies.fixed(); + PlacementContext ctx = new PlacementContext("b", + List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")), + noSessions(), Set.of(), Set.of("b")); + assertEquals("a", policy.select(ctx).profile(), + "the default 'b' is quarantined, so fixed falls through to the first un-quarantined candidate"); + } + + @Test + void fixedThrowsWhenDefaultAndEveryCandidateQuarantined() { + PlacementPolicy policy = PlacementPolicies.fixed(); + PlacementContext ctx = new PlacementContext("b", + List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")), + noSessions(), Set.of(), Set.of("a", "b")); + PlacementException e = assertThrows(PlacementException.class, () -> policy.select(ctx)); + assertTrue(e.getMessage().contains("quarantined"), e.getMessage()); + } + @Test void roundRobinCyclesThroughAvailableProfiles() { PlacementPolicy policy = PlacementPolicies.roundRobin(); @@ -82,6 +108,18 @@ class PlacementPolicyTest { } } + @Test + void roundRobinSkipsQuarantinedProfiles() { + PlacementPolicy policy = PlacementPolicies.roundRobin(); + List candidates = List.of( + PlacementCandidate.profile("a"), + PlacementCandidate.profile("b")); + PlacementContext ctx = ctx(candidates, noSessions(), Set.of(), Set.of("a")); + for (int i = 0; i < 5; i++) { + assertEquals("b", policy.select(ctx).profile(), "a is quarantined, so every pick lands on b"); + } + } + @Test void roundRobinThrowsWhenAllAtMaxLoad() { PlacementPolicy policy = PlacementPolicies.roundRobin(); @@ -150,6 +188,29 @@ class PlacementPolicyTest { assertTrue(e.getMessage().contains("maxLoad"), e.getMessage()); } + @Test + void weightedSkipsQuarantinedProfile() { + PlacementPolicy policy = PlacementPolicies.weighted(); + List candidates = List.of( + PlacementCandidate.profile("a", 1.0f, null), + PlacementCandidate.profile("b", 1.0f, null)); + PlacementContext ctx = ctx(candidates, noSessions(), Set.of(), Set.of("a")); + for (int i = 0; i < 5; i++) { + assertEquals("b", policy.select(ctx).profile(), "a is quarantined, so every pick lands on b"); + } + } + + @Test + void weightedThrowsWhenAllQuarantined() { + PlacementPolicy policy = PlacementPolicies.weighted(); + List candidates = List.of( + PlacementCandidate.profile("a"), + PlacementCandidate.profile("b")); + PlacementException e = assertThrows(PlacementException.class, + () -> policy.select(ctx(candidates, noSessions(), Set.of(), Set.of("a", "b")))); + assertTrue(e.getMessage().contains("quarantined"), e.getMessage()); + } + @Test void weightedThrowsWhenAllUnreachable() { PlacementPolicy policy = PlacementPolicies.weighted();