From ba04b2359bd0c3a19d859708fb4b58a0a8a334c0 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 3 Sep 2026 11:45:11 +0700 Subject: [PATCH] fleetd #201/#227 unit 5: wire errorPattern, cool-off spawn gate, and fleet views Wires the already-merged units into production: - Per-profile errorPattern config (beside exhaustedPattern), compiled once at startup; falls back to the legacy (?i)\bAPI Error\s*: pattern when unset. Startup logs configured-vs-legacy coverage, same as exhaustedPattern. - One production BackendErrorSink in Fleetd.java: mark backend_error on the session, resolve the profile's credential fail-loud (never Optional.ifPresent), record it in BackendOutagePolicy, and push a lead nudge on a new incident. - CompositePeerLauncher's explicit and automatic spawn paths both refuse a cooling-off credential; exhaustion quarantine wins when both are active. PlacementContext gets a separate coolingOff set so refusal text says "cooling off", never "exhausted". - fleet_list/fleet_profiles report coolingOffForSeconds as an independent fact from quarantinedForSeconds; both can appear together. - fleetd.example.yaml documents errorPattern and the 2/60/60 cool-off policy; CLAUDE.md tells leads how to read the two independent outage states. Also: FixedPlacementPolicy.java, not in the original file list, needed the same coolingOff filtering as PlacementPolicyUtil (it does its own inline candidate filtering rather than delegating). 1190 tests, 0 failures (up from the 1163 baseline); BUILD SUCCESS. --- CLAUDE.md | 12 + fleetd/fleetd.example.yaml | 46 ++- .../src/main/java/dev/ltms/fleet/Fleetd.java | 91 +++++- .../java/dev/ltms/fleet/config/ConfigRef.java | 10 +- .../dev/ltms/fleet/config/FleetConfig.java | 99 +++++- .../java/dev/ltms/fleet/mcp/FleetMcp.java | 115 ++++++- .../fleet/member/CompositePeerLauncher.java | 115 ++++++- .../fleet/placement/FixedPlacementPolicy.java | 51 +-- .../fleet/placement/PlacementContext.java | 11 +- .../fleet/placement/PlacementPolicyUtil.java | 28 +- .../dev/ltms/fleet/config/ConfigRefTest.java | 48 +++ .../ltms/fleet/config/FleetConfigTest.java | 82 +++++ .../fleet/inject/BackendOutageFlowTest.java | 295 ++++++++++++++++++ .../java/dev/ltms/fleet/mcp/FleetMcpTest.java | 95 ++++++ .../member/CompositePeerLauncherTest.java | 182 +++++++++++ .../fleet/placement/PlacementPolicyTest.java | 82 ++++- 16 files changed, 1295 insertions(+), 67 deletions(-) create mode 100644 fleetd/src/test/java/dev/ltms/fleet/inject/BackendOutageFlowTest.java diff --git a/CLAUDE.md b/CLAUDE.md index 202e993..5962a63 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -188,6 +188,18 @@ must obey belongs in the charter, not here. - **This repo is the bridge.** The daemon is `fleetd`, its MCP mount is `http://127.0.0.1:8765/mcp`, and the code behind the rules above is `mcp/FleetMcp` (tools), `auth/Authz` (the role table), `mcp/ConnectionIdentity` (connection→role), and `worker/*Launcher` (`REPLY_CHARTER`). +- **`fleet_profiles`/`fleet_list` report two separate outage states, and they are not the same + thing.** *Quarantined* (CB-578) means the backend told us it is out of capacity — a long, + 1800s-default cooldown. *Cooling off* (fleetd #201/#227) means a profile's credential threw two + distinct backend errors (a non-exhaustion failure such as an HTTP 5xx) within 60 seconds — a + short, fixed 60s cooldown, not configurable per profile. Each check runs independently, so a + profile can show both at once. In the JSON: a cooling profile carries `credentialId` and + `coolingOffForSeconds`; a quarantined profile carries `quarantinedForSeconds`; a profile hit by + both carries all three fields, and either state alone already sets that profile's `free` to `0`. + A `fleet_spawn` naming a cooling-off profile is refused before it ever reaches the backend + adapter, with a message naming the credential and the remaining seconds ("cooling off after + repeated backend errors") — distinct wording from a quarantine refusal, so don't conflate the + two when reading a spawn failure. - **Skills available to delegate:** `implementer` (worktree → commit → push → own PR) and `reviewer` (scoped review → one structured finding). Name one in every delegation. - **Primary-side skills** (not delegation playbooks — a worker cannot use them): diff --git a/fleetd/fleetd.example.yaml b/fleetd/fleetd.example.yaml index 80a9707..f33f622 100644 --- a/fleetd/fleetd.example.yaml +++ b/fleetd/fleetd.example.yaml @@ -206,6 +206,39 @@ herdrSocket: ~/.config/herdr/herdr.sock # 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. +# errorPattern → fleetd #201 / #227: regex matched against a completion-fallback scrape to +# classify a turn that ended with no fleet_reply as a BACKEND ERROR — a +# credential outage or a provider 5xx — rather than a real answer or a +# usage-limit exhaustion (exhaustedPattern above always wins when a line +# matches both). Opt-in. Omit it and this profile falls back to fleetd's +# built-in legacy pattern `(?i)\bAPI Error\s*:` — classification still +# happens, just without a profile-specific match; every backend words its +# failure differently, so a hardcoded sentence would only ever match one +# of them. +# DEFERRED: compiled once into a startup pattern map, same as exhaustedPattern +# — editing it needs a daemon restart. +# # errorPattern: "503 Service Unavailable" # opt-in: classify a backend outage +# +# What happens once a match fires (BackendOutagePolicy, credentialId-keyed, +# SEPARATE from the CB-578 stage B quarantine above and never merged with it): +# - threshold 2 — TWO DISTINCT TARGETS (never raw events) on the same +# effective credential inside a 60-second window start an "incident" and a +# 60-second cool-off for that credential. One member repeating the same +# classified line twice never cools anything off — a real outage hits +# every target on that credential, so requiring a second, independent +# target loses nothing against the case this guards against, while +# protecting against a heuristic misfire on one flaky member. +# - a fresh error while a credential is already cooling off is ignored +# outright: it neither extends the 60s deadline nor starts a new incident. +# - `fleet_list`/`fleet_profiles` report a cooling credential with +# `coolingOffForSeconds` (never `quarantinedForSeconds`, unless CB-578 +# exhaustion quarantine is ALSO independently active for the same +# credential — the two checks can both fire at once). A spawn onto a +# cooling profile is refused with a message naming the credential and +# remaining seconds — "cooling off", never "exhausted", so an operator can +# tell a short transient fault from a spent subscription at a glance. +# - the lead gets ONE nudge per incident (not one per affected target), via +# the same push loop that already delivers ticket/question reminders. # env → extra environment for this profile's workers, as a literal key/value map # (CB-511). Use it to give workers a toolchain. # @@ -273,6 +306,7 @@ profiles: # 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) + # errorPattern: "503 Service Unavailable" # opt-in: classify a backend outage (fleetd #201/#227) — see the key doc above # 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: [".env", ".envrc"] # the default; never add .mcp.json or .claude/settings.local.json — see above @@ -369,6 +403,12 @@ placement: weighted # 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. +# +# This does NOT govern the fleetd #201 / #227 backend-error cool-off documented under errorPattern +# above — that mechanism is a separate, shorter-lived, NOT-configurable policy (threshold 2 distinct +# targets, 60-second window, 60-second cool-off), on purpose: it exists to survive a brief transient +# fault, not to replace this 30-minute exhaustion quarantine. Do not conflate the two when reading +# fleet_list/fleet_profiles — coolingOffForSeconds and quarantinedForSeconds are independent facts. # quarantineCooldownSeconds: 1800 # Re-read this file without restarting the daemon (CB-559). Off unless you add this block, so an @@ -395,8 +435,10 @@ placement: weighted # 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 +# configDir, mcpUrl, tabLabel, exhaustedPattern, errorPattern (fleetd #201 / #227 — +# compiled once into a startup pattern map the same way exhaustedPattern is). 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 diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index d57c2e7..98b8178 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -13,6 +13,8 @@ import dev.ltms.fleet.lead.LeadLauncher; import dev.ltms.fleet.herdr.PaneLocator; import dev.ltms.fleet.herdr.UnixSocketHerdrClient; import dev.ltms.fleet.herdr.WorkspaceControl; +import dev.ltms.fleet.inject.BackendErrorPatternLookup; +import dev.ltms.fleet.inject.BackendErrorSink; import dev.ltms.fleet.inject.CompletionResolver; import dev.ltms.fleet.inject.ExhaustedPatternLookup; import dev.ltms.fleet.inject.ExhaustionSink; @@ -50,6 +52,7 @@ import dev.ltms.fleet.member.ClaudeCodeLauncher; import dev.ltms.fleet.member.CompositePeerLauncher; import dev.ltms.fleet.member.HerdrPeerLauncher; import dev.ltms.fleet.member.OpenCodeLauncher; +import dev.ltms.fleet.placement.BackendOutagePolicy; import dev.ltms.fleet.placement.BackendQuarantine; import io.javalin.Javalin; import org.slf4j.Logger; @@ -61,6 +64,7 @@ import java.util.ArrayList; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.Set; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; @@ -205,12 +209,18 @@ public final class Fleetd { // 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())); + // fleetd #201 Unit 5: one outage-cool-off tracker for the whole daemon, shared between the + // launcher (checked at spawn, like `quarantine` above) and the backend-error sink wired in + // below (written on a classified backend error). A SEPARATE, shorter-lived mechanism from + // `quarantine` — see BackendOutagePolicy's class doc — never merged with it. + BackendOutagePolicy outagePolicy = new BackendOutagePolicy(System::nanoTime); PeerLauncher workers = new CompositePeerLauncher( adapters, cfg.effectiveDefaultProfile(), config, profileName -> liveCountRef.get().apply(profileName), - quarantine); + quarantine, + outagePolicy); // CB-504: under supervision (launchd/systemd) fleetd 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 @@ -342,6 +352,26 @@ public final class Fleetd { .orElse(null); log.info("backend-exhausted classification (CB-578 stage A): {}", CompletionResolver.coverage(cfg.profiles().keySet(), exhaustedPatternsByProfile.keySet())); + // fleetd #201 Unit 5: classify a completion-fallback scrape that matches a profile's + // configured backend-error refusal (a credential outage, a provider 5xx) as a backend error + // rather than handing it back as a real answer. Compiled once at startup, keyed by profile + // name, mirroring exhaustedPatternsByProfile above — a profile with no configured + // errorPattern is simply absent here, so CompletionResolver falls back to its built-in + // narrow {@code (?i)\bAPI Error\s*:} compatibility pattern for that profile's targets + // (BackendErrorPatternLookup#legacy's contract — see backendErrorPatterns below). + Map errorPatternsByProfile = new LinkedHashMap<>(); + cfg.profiles().forEach((name, profile) -> { + if (profile.hasErrorPattern()) { + errorPatternsByProfile.put(name, Pattern.compile(profile.errorPattern())); + } + }); + BackendErrorPatternLookup backendErrorPatterns = target -> sessions.roster().stream() + .filter(session -> target.equals(session.terminalId())) + .findFirst() + .map(session -> errorPatternsByProfile.get(session.profile())) + .orElse(null); + log.info("backend-error classification (fleetd #201 Unit 5): {}", + CompletionResolver.coverage(cfg.profiles().keySet(), errorPatternsByProfile.keySet())); // 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 @@ -389,8 +419,56 @@ public final class Fleetd { // fleetd #175: point the forwarding sink handed to OpenCodeLauncher above at the real one, // now that `sessions` exists to resolve target -> session -> profile. exhaustionSinkRef.set(exhaustionSink); + // fleetd #201 Unit 5: the production BackendErrorSink needs `pushLoop` (built further below, + // after `sessions`) to tell a lead about an incident or an unmapped target — the same + // construction-order cycle `exhaustionSinkRef` breaks above, broken the same way: a mutable + // holder set once `pushLoop` exists, read lazily from inside the lambda built here. + AtomicReference pushLoopRef = new AtomicReference<>(); + // Order: (1) mark the member BACKEND_ERROR; (2) resolve profile/credential through the + // roster — fail loud (never Optional.ifPresent, the fleetd #234 lesson applied to this new + // sink) and notify the lead via onBackendTargetUnmapped when it cannot be resolved; (3) + // record the error in BackendOutagePolicy; (4) on a NEW incident (the record() call that + // actually crosses the threshold), tell the lead via onBackendIncident. + BackendErrorSink backendErrorSink = (target, matchedLine, reason) -> { + sessions.onBackendError(target, reason); + + String profileName = sessions.roster().stream() + .filter(session -> target.equals(session.terminalId())) + .findFirst() + .map(MemberSession::profile) + .orElse(null); + FleetConfig.Profile profile = profileName == null ? null : config.get().profiles().get(profileName); + if (profile == null) { + log.error("backend error on target '{}' ({}) but no profile could be resolved — the " + + "target is not (yet) in the roster — no cool-off applied (fleetd #201 Unit 5)", + target, reason); + ReplyPushLoop loop = pushLoopRef.get(); + if (loop != null) { + loop.onBackendTargetUnmapped(target, reason); + } + return; + } + String credentialId = profile.effectiveCredentialId(); + Optional incident = outagePolicy.record(credentialId, target, reason); + incident.ifPresent(inc -> { + List affectedProfiles = config.get().profiles().values().stream() + .filter(p -> credentialId.equals(p.effectiveCredentialId())) + .map(FleetConfig.Profile::profile) + .sorted() + .toList(); + log.warn("credential '{}' cooling off for {}s after backend errors on {} distinct " + + "target(s) (profile '{}'): {}", credentialId, + inc.remainingCoolOffSeconds(), inc.evidenceCount(), profile.profile(), reason); + ReplyPushLoop loop = pushLoopRef.get(); + if (loop != null) { + loop.onBackendIncident(inc.id(), inc.targets(), credentialId, affectedProfiles, + (int) inc.remainingCoolOffSeconds()); + } + }); + }; AgentControl agents = router.memberAgents(); - CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns, exhaustionSink); + CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns, + exhaustionSink, backendErrorPatterns, backendErrorSink); // 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(); @@ -472,6 +550,9 @@ public final class Fleetd { Metrics metrics = FleetMetrics.create(sessions, replyInbox); var pushLoop = new ReplyPushLoop(primaryRegistry, router.leadAgents(), replyInbox, pushScheduler, maxReminders, backoffMs, metrics); + // fleetd #201 Unit 5: point the forwarding holder captured by the backendErrorSink lambda + // above at the real push loop, now that it exists. + pushLoopRef.set(pushLoop); // CB-551: idle-lead heartbeat. Opt-in; absent `leadHeartbeat:` this is never constructed, so // an upgraded daemon cannot silently start spending subscription on nudging an idle lead. // It has its own single-thread scheduler and holds its own scheduler shutdown via close(). @@ -579,7 +660,11 @@ public final class Fleetd { var configured = config.get().profiles().get(profile); return configured == null ? null : configured.effectiveCredentialId(); }, quarantine), - leadMailbox); + leadMailbox, + new FleetMcp.OutageSource(profile -> { + var configured = config.get().profiles().get(profile); + return configured == null ? null : configured.effectiveCredentialId(); + }, outagePolicy)); // CB-637: the receive half. Only constructed when a lead mailbox actually opened — with no // coordinator (or an unreachable one) there is nothing to deliver, so no scheduler is diff --git a/fleetd/src/main/java/dev/ltms/fleet/config/ConfigRef.java b/fleetd/src/main/java/dev/ltms/fleet/config/ConfigRef.java index ce13fc1..419253a 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/config/ConfigRef.java +++ b/fleetd/src/main/java/dev/ltms/fleet/config/ConfigRef.java @@ -43,7 +43,9 @@ import java.util.function.Supplier; * which is constructed once), and an existing profile's launch settings — * {@code model}, {@code baseUrl}, {@code argv}, {@code env}, {@code mcpUrl}, * {@code exhaustedPattern} (CB-578 stage A — compiled once into {@code Fleetd.main}'s - * pattern map at startup), and the rest. {@code credentialId} (CB-578 stage B) is NOT on + * pattern map at startup), {@code errorPattern} (fleetd #201 Unit 5 — compiled once into + * {@code Fleetd.main}'s backend-error pattern map at startup, the same way), 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 @@ -292,6 +294,10 @@ public final class ConfigRef implements Supplier { // CB-578 stage B: exhaustedPattern is compiled once into Fleetd.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()); + && Objects.equals(a.exhaustedPattern(), b.exhaustedPattern()) + // fleetd #201 Unit 5: errorPattern is compiled once into Fleetd.main's backend-error + // pattern map at startup (see BackendErrorPatternLookup wiring), the same way + // exhaustedPattern is — a reload never re-reads it either. + && Objects.equals(a.errorPattern(), b.errorPattern()); } } diff --git a/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java b/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java index 04a98a3..3fbe5ff 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java +++ b/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java @@ -24,6 +24,8 @@ import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Set; +import java.util.regex.Pattern; +import java.util.regex.PatternSyntaxException; /** * {@code fleetd} configuration, loaded from a YAML file (see @@ -356,6 +358,16 @@ public record FleetConfig( * {@code limit.context} instead, which bounds the window opencode compacts * within, and only when {@code model} resolves to a * {@code provider/model} pair. + * @param errorPattern regex matched against a completion-fallback scrape (fleetd #201 Unit 5) + * to classify a turn that ended with no {@code fleet_reply} as a backend + * error (credential outage, provider 5xx) rather than a real answer or an + * exhausted usage limit. {@code null}/blank ⇒ this profile relies on + * {@link dev.ltms.fleet.inject.CompletionResolver}'s built-in + * {@code (?i)\bAPI Error\s*:} compatibility pattern instead — classification + * still happens, just without a profile-specific match (reported separately + * at startup as legacy-default coverage, never full coverage). Compiled once + * at startup ({@code Fleetd.main}), like {@code exhaustedPattern}, so a + * reload of this key is DEFERRED (see {@code ConfigRef}). */ @JsonIgnoreProperties(ignoreUnknown = true) public record Profile(String profile, String baseUrl, String model, @@ -374,7 +386,8 @@ public record FleetConfig( String ideMcpUrl, String ideProjectDir, String ideOpenCommand, - Integer autoCompactWindow) { + Integer autoCompactWindow, + String errorPattern) { /** Peer kind spawned by {@link dev.ltms.fleet.member.ClaudeCodeLauncher} (the default). */ public static final String KIND_CLAUDE_CODE = "claude-code"; @@ -448,6 +461,32 @@ public record FleetConfig( // there is no blank-string form to normalize (it's an Integer), and the [100000, 1000000] // range is enforced eagerly at config load (rejectAutoCompactWindowOutOfRange), naming the // profile, rather than silently clamped here. A profile that never sets it keeps null. + // errorPattern stays null when unset/blank (opt-in) — same rule as exhaustedPattern: no + // defaulting, no vendor wording. A profile that never sets it relies on + // CompletionResolver's built-in compatibility pattern instead (never "off"). + errorPattern = (errorPattern == null || errorPattern.isBlank()) ? null : errorPattern; + } + + /** + * Backward-compatible constructor without the fleetd #201 {@code errorPattern} field — the + * profile relies on {@code CompletionResolver}'s built-in compatibility pattern instead + * (legacy-default coverage). This is the shape the canonical constructor had before the + * field was added; every pre-existing Java call site (and any YAML that omits the key) keeps + * compiling and behaving identically. Jackson binds the canonical (longest) constructor, so + * YAML omitting {@code errorPattern:} still lands here as {@code null} via that path too. + */ + 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, String credentialId, + String ideMcpUrl, String ideProjectDir, String ideOpenCommand, + Integer autoCompactWindow) { + this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel, + mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, + subscription, exhaustedPattern, credentialId, ideMcpUrl, ideProjectDir, ideOpenCommand, + autoCompactWindow, null); } /** @@ -492,7 +531,8 @@ public record FleetConfig( 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, credentialId, ideMcpUrl, ideProjectDir, ideOpenCommand, autoCompactWindow); + exhaustedPattern, credentialId, ideMcpUrl, ideProjectDir, ideOpenCommand, autoCompactWindow, + errorPattern); } /** @@ -589,6 +629,15 @@ public record FleetConfig( return exhaustedPattern != null; } + /** + * True when this profile configures its own fleetd #201 Unit 5 backend-error classification + * pattern. {@code false} means this profile relies on {@code CompletionResolver}'s built-in + * {@code (?i)\bAPI Error\s*:} compatibility pattern instead (legacy-default coverage). + */ + public boolean hasErrorPattern() { + return errorPattern != 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 @@ -1389,6 +1438,7 @@ public record FleetConfig( rejectDuplicateMemberSlots(yaml); rejectNegativeMaxLoad(yaml); rejectAutoCompactWindowOutOfRange(yaml); + rejectMalformedErrorPattern(yaml); rejectUnknownKind(yaml); rejectUnknownAuthMode(yaml); rejectUnknownPlacement(yaml); @@ -1737,6 +1787,51 @@ public record FleetConfig( } } + /** + * Reject a profile whose {@code errorPattern} (fleetd #201 Unit 5) is not a valid Java regex, + * naming the profile, the key, and the parser's own message. + * + *

Unset/{@code null} means "use {@code CompletionResolver}'s built-in {@code (?i)\bAPI + * Error\s*:} compatibility pattern" and passes silently. A profile that DOES set the key gets it + * compiled once at daemon startup ({@code Fleetd.main}, mirroring {@code exhaustedPattern}) — an + * uncaught {@link java.util.regex.PatternSyntaxException} there crashes startup without naming + * which profile or key is at fault. Validate eagerly here instead, at config load, the same + * "fail loud at load, not lazily later" reasoning as {@link #rejectAutoCompactWindowOutOfRange}. + * + * @param yaml the raw config text + * @throws IllegalStateException when any profile's {@code errorPattern} fails to compile + */ + static void rejectMalformedErrorPattern(String yaml) { + Map raw; + try { + raw = YAML.readValue(yaml, Map.class); + } catch (IOException | IllegalArgumentException e) { + return; // a malformed file is reported by the real parse, not here + } + if (raw == null || !(raw.get("profiles") instanceof Map profiles)) { + return; + } + List bad = new ArrayList<>(); + for (Map.Entry e : profiles.entrySet()) { + if (!(e.getValue() instanceof Map p)) { + continue; + } + if (!(p.get("errorPattern") instanceof String pattern) || pattern.isBlank()) { + continue; + } + try { + Pattern.compile(pattern); + } catch (PatternSyntaxException ex) { + bad.add("profiles." + e.getKey() + ".errorPattern (\"" + pattern + "\"): " + ex.getMessage()); + } + } + bad.sort(String::compareTo); + if (!bad.isEmpty()) { + throw new IllegalStateException("refusing to start: malformed errorPattern — " + + String.join("; ", bad)); + } + } + /** The peer kinds this build has an adapter for — {@link Profile#kind()}'s only valid values. */ private static final Set KNOWN_KINDS = Set.of(Profile.KIND_CLAUDE_CODE, Profile.KIND_OPENCODE); diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index 4b664e3..47de01a 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -16,6 +16,7 @@ import dev.ltms.fleet.msg.LeadMessage; import dev.ltms.fleet.msg.MessageService; import dev.ltms.fleet.msg.Rendezvous; import dev.ltms.fleet.peer.PeerUnreachableException; +import dev.ltms.fleet.placement.BackendOutagePolicy; import dev.ltms.fleet.placement.BackendQuarantine; import dev.ltms.fleet.placement.PlacementException; import dev.ltms.fleet.session.SessionManager; @@ -92,6 +93,8 @@ public final class FleetMcp { private final CapacitySource capacity; private final HealthCoverageSource healthCoverage; private final QuarantineSource quarantine; + /** fleetd #201 Unit 5: SEPARATE from {@link #quarantine} — see {@link OutageSource}'s doc. */ + private final OutageSource outage; /** CB-637: this daemon's lead-to-lead channel; {@code null} when no coordinator is configured. */ private final LeadChannel leadChannel; @@ -115,6 +118,22 @@ public final class FleetMcp { public static QuarantineSource none() { return new QuarantineSource(_ -> null, BackendQuarantine.none()); } } + /** + * fleetd #201 Unit 5 cool-off facts used by {@code fleet_list}/{@code fleet_profiles}: a + * profile → credential id lookup, plus the shared {@link BackendOutagePolicy} to read remaining + * cool-offs off. A SEPARATE source from {@link QuarantineSource} — never merged into it — so a + * credential that is cooling off after repeated backend errors is reported independently of + * whether it is ALSO under CB-578 stage B exhaustion quarantine; the two checks can both fire + * for the same profile, and when they do, {@code fleet_list}/{@code fleet_profiles} report both + * (see {@link #capacityView}/{@link #profiles}). + */ + public record OutageSource(Function credentialIdFor, BackendOutagePolicy outagePolicy) { + /** Inert source — no profile is ever reported cooling off. Explicit stand-in, not a default. */ + public static OutageSource none() { + return new OutageSource(_ -> null, new BackendOutagePolicy(System::nanoTime)); + } + } + /** * @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 @@ -129,7 +148,7 @@ public final class FleetMcp { CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage, QuarantineSource quarantine) { this(messages, workers, sessions, identity, presence, primaryRegistry, callers, metrics, capacity, - healthCoverage, quarantine, null); + healthCoverage, quarantine, null, OutageSource.none()); } /** @@ -142,9 +161,25 @@ public final class FleetMcp { ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry, CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage, QuarantineSource quarantine, LeadChannel leadChannel) { + this(messages, workers, sessions, identity, presence, primaryRegistry, callers, metrics, capacity, + healthCoverage, quarantine, leadChannel, OutageSource.none()); + } + + /** + * As above, with fleetd #201 Unit 5 cool-off facts for {@code fleet_list}/{@code fleet_profiles} + * (see {@link OutageSource}). This is what {@code Fleetd.main} actually wires up. + * + * @param outage required — pass {@link OutageSource#none()} for a caller that does not want the + * feature, never a defaulting overload (the same rule {@code quarantine} follows). + */ + public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions, + ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry, + CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage, + QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage) { this.leadChannel = leadChannel; this.capacity = capacity; this.quarantine = Objects.requireNonNull(quarantine, "quarantine"); + this.outage = Objects.requireNonNull(outage, "outage"); this.healthCoverage = healthCoverage; McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get(); this.transport = HttpServletStreamableServerTransportProvider.builder() @@ -273,7 +308,7 @@ public final class FleetMcp { (exchange, _) -> { McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null); if (denied != null) return denied; - return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, + return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage, callers == null ? Map.of() : callers.leads(), callerTerminal(exchange), leadChannel == null ? null : leadChannel.selfCoordId()); @@ -289,7 +324,7 @@ public final class FleetMcp { (exchange, _) -> { McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null); if (denied != null) return denied; - return profiles(workers, quarantine); + return profiles(workers, quarantine, outage); }; BiFunction whoamiHandler = (exchange, _) -> { @@ -872,25 +907,47 @@ public final class FleetMcp { * has ever been quarantined gets exactly the pre-stage-B shape. */ static McpSchema.CallToolResult profiles(PeerLauncher workers, QuarantineSource quarantine) { + return profiles(workers, quarantine, OutageSource.none()); + } + + /** + * As above, plus a SEPARATE {@code coolingOff} map (fleetd #201 Unit 5) — never merged into + * {@code quarantined} — for a profile whose credential is cooling off after repeated backend + * errors ({@link BackendOutagePolicy}). The two checks are independent: a profile can appear in + * both maps at once when it is both exhaustion-quarantined AND cooling off. + */ + static McpSchema.CallToolResult profiles(PeerLauncher workers, QuarantineSource quarantine, OutageSource outage) { Map result = new LinkedHashMap<>(); result.put("profiles", workers.profiles()); result.put("default", workers.defaultProfile() == null ? "" : workers.defaultProfile()); Map quarantined = new LinkedHashMap<>(); + Map coolingOff = new LinkedHashMap<>(); for (String profile : workers.profiles()) { String credentialId = quarantine.credentialIdFor().apply(profile); - if (credentialId == null) { - continue; + if (credentialId != null) { + quarantine.quarantine().remainingSeconds(credentialId).ifPresent(remaining -> { + Map row = new LinkedHashMap<>(); + row.put("credentialId", credentialId); + row.put("quarantinedForSeconds", remaining); + quarantined.put(profile, row); + }); + } + String outageCredentialId = outage.credentialIdFor().apply(profile); + if (outageCredentialId != null) { + outage.outagePolicy().remainingCoolOffSeconds(outageCredentialId).ifPresent(remaining -> { + Map row = new LinkedHashMap<>(); + row.put("credentialId", outageCredentialId); + row.put("coolingOffForSeconds", remaining); + coolingOff.put(profile, row); + }); } - 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); } + if (!coolingOff.isEmpty()) { + result.put("coolingOff", coolingOff); + } return text(json(result)); } @@ -929,6 +986,14 @@ public final class FleetMcp { return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, leads, selfTerm, null); } + /** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */ + static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages, + CapacitySource capacity, HealthCoverageSource healthCoverage, + QuarantineSource quarantine, OutageSource outage, + Map leads, String selfTerm) { + return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage, leads, selfTerm, null); + } + /** * As above, additionally reporting this daemon's own lead coordination id (CB-637) when one is * configured and its channel opened. There is no peer-discovery surface yet — a lead addresses a @@ -942,6 +1007,15 @@ public final class FleetMcp { CapacitySource capacity, HealthCoverageSource healthCoverage, QuarantineSource quarantine, Map leads, String selfTerm, String selfCoordId) { + return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, OutageSource.none(), + leads, selfTerm, selfCoordId); + } + + /** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */ + static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages, + CapacitySource capacity, HealthCoverageSource healthCoverage, + QuarantineSource quarantine, OutageSource outage, + Map leads, String selfTerm, String selfCoordId) { try { Map live = workers.list().stream() .map(Agent.class::cast) @@ -968,7 +1042,7 @@ public final class FleetMcp { } if (capacity.available()) result.put("capacity", profiles.stream() .map(profile -> capacityView(profile, capacity.liveCount(), capacity.maxLoad(), roster, messages, - capacity.clock().getAsLong(), quarantine)).toList()); + capacity.clock().getAsLong(), quarantine, outage)).toList()); return text(json(result)); } catch (HerdrException e) { return error("herdr error listing the fleet: " + e.getMessage()); @@ -1001,10 +1075,17 @@ public final class FleetMcp { * reports, reusing {@link QuarantineSource} rather than a second lookup. Both new keys are * added only when the profile is actually quarantined, so an ordinary fleet's rows are * byte-identical to before this change. + * + *

fleetd #201 Unit 5: a profile whose credential is cooling off (a SEPARATE, independent + * check from quarantine — see {@link OutageSource}) also forces {@code free} to 0 and carries + * {@code credentialId}/{@code coolingOffForSeconds}, but never {@code quarantinedForSeconds} — + * that key is added only when exhaustion quarantine is ALSO active for this profile, since the + * two checks are independent and either, both, or neither can be true. */ private static Map capacityView(String profile, Function liveCount, Function maxLoad, List roster, - MessageService messages, long nowNanos, QuarantineSource quarantine) { + MessageService messages, long nowNanos, QuarantineSource quarantine, + OutageSource outage) { Integer cap = maxLoad.apply(profile); int live = liveCount.apply(profile); int reclaimable = (int) roster.stream().filter(s -> profile.equals(s.profile())) @@ -1022,6 +1103,14 @@ public final class FleetMcp { row.put("quarantinedForSeconds", remaining); }); } + String outageCredentialId = outage.credentialIdFor().apply(profile); + if (outageCredentialId != null) { + outage.outagePolicy().remainingCoolOffSeconds(outageCredentialId).ifPresent(remaining -> { + row.put("free", 0); + row.put("credentialId", outageCredentialId); + row.put("coolingOffForSeconds", remaining); + }); + } return row; } diff --git a/fleetd/src/main/java/dev/ltms/fleet/member/CompositePeerLauncher.java b/fleetd/src/main/java/dev/ltms/fleet/member/CompositePeerLauncher.java index ee83b17..80d0237 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/member/CompositePeerLauncher.java +++ b/fleetd/src/main/java/dev/ltms/fleet/member/CompositePeerLauncher.java @@ -10,6 +10,7 @@ import dev.ltms.fleet.peer.PeerHandle; import dev.ltms.fleet.peer.PeerLauncher; import dev.ltms.fleet.peer.PeerUnreachableException; import dev.ltms.fleet.peer.SpawnRequest; +import dev.ltms.fleet.placement.BackendOutagePolicy; import dev.ltms.fleet.placement.BackendQuarantine; import dev.ltms.fleet.placement.PlacementCandidate; import dev.ltms.fleet.placement.PlacementContext; @@ -92,6 +93,18 @@ public final class CompositePeerLauncher implements PeerLauncher { /** CB-578 stage B: credential cooldown, checked before an explicit spawn and filtered into placement. */ private final BackendQuarantine quarantine; + /** + * fleetd #201 Unit 5: credential cool-off after repeated backend errors, checked before an + * explicit spawn and filtered into placement — a SEPARATE, shorter-lived source from + * {@link #quarantine}. A never-{@code record}-called instance is naturally inert (its + * {@code remainingCoolOffSeconds} always returns empty), so back-compat constructors that predate + * this feature share one fixed instance rather than needing a {@code none()} sentinel. + */ + private final BackendOutagePolicy outagePolicy; + + /** The shared inert instance back-compat constructors wire in — never {@code record}-called. */ + private static final BackendOutagePolicy NO_OUTAGE_POLICY = new BackendOutagePolicy(System::nanoTime); + /** * 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 @@ -130,7 +143,8 @@ public final class CompositePeerLauncher implements PeerLauncher { Map profileConfigs, PlacementPolicy placementPolicy, Function liveCount) { - this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, null, BackendQuarantine.none()); + this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, null, + BackendQuarantine.none(), NO_OUTAGE_POLICY); } /** @@ -148,13 +162,14 @@ public final class CompositePeerLauncher implements PeerLauncher { PlacementPolicy placementPolicy, Function liveCount, FleetConfig.Fleet fleet) { - this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, fleet, BackendQuarantine.none()); + this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, fleet, + BackendQuarantine.none(), NO_OUTAGE_POLICY); } /** - * 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 Fleetd.main} actually wires up. + * Production constructor with role pools and quarantine (CB-578 stage B), no cool-off (fleetd + * #201 Unit 5). Kept for callers that predate the cool-off feature; use the 8-arg overload below + * to wire a real {@link BackendOutagePolicy}. * * @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). @@ -166,17 +181,42 @@ public final class CompositePeerLauncher implements PeerLauncher { Function liveCount, FleetConfig.Fleet fleet, BackendQuarantine quarantine) { + this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, fleet, quarantine, + NO_OUTAGE_POLICY); + } + + /** + * Production constructor with role pools, quarantine (CB-578 stage B), and cool-off (fleetd + * #201 Unit 5). The full-featured non-reloading form; + * {@link #CompositePeerLauncher(List, String, Supplier, Function, BackendQuarantine, BackendOutagePolicy)} + * is what {@code Fleetd.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). + * @param outagePolicy required — pass a fresh, never-{@code record}-called {@link + * BackendOutagePolicy} for a caller that does not want the feature; same + * "explicit opt-out, never a silent default" rule as {@code quarantine}. + */ + public CompositePeerLauncher(List delegates, + String defaultProfile, + Map profileConfigs, + PlacementPolicy placementPolicy, + Function liveCount, + FleetConfig.Fleet fleet, + BackendQuarantine quarantine, + BackendOutagePolicy outagePolicy) { // 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), quarantine); + constant(placementPolicy), liveCount, constant(fleet), quarantine, outagePolicy); } /** * Production constructor that re-reads its placement inputs per spawn (CB-559), so a config - * reload retargets the next member without a restart. + * reload retargets the next member without a restart. No cool-off (fleetd #201 Unit 5); use the + * 6-arg overload below to wire a real {@link BackendOutagePolicy}. * * @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 @@ -186,12 +226,31 @@ public final class CompositePeerLauncher implements PeerLauncher { Supplier config, Function liveCount, BackendQuarantine quarantine) { + this(delegates, defaultProfile, config, liveCount, quarantine, NO_OUTAGE_POLICY); + } + + /** + * Production constructor that re-reads its placement inputs per spawn (CB-559), with cool-off + * (fleetd #201 Unit 5). This is what {@code Fleetd.main} actually wires up. + * + * @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 + * @param outagePolicy required — fleetd #201 Unit 5; pass a fresh, never-{@code record}-called + * {@link BackendOutagePolicy} to opt out + */ + public CompositePeerLauncher(List delegates, + String defaultProfile, + Supplier config, + Function liveCount, + BackendQuarantine quarantine, + BackendOutagePolicy outagePolicy) { this(delegates, defaultProfile, () -> config.get().profiles(), () -> PlacementPolicies.fromName(config.get().placement()), liveCount, () -> config.get().fleet(), - quarantine); + quarantine, + outagePolicy); } /** The all-suppliers form every other constructor funnels into. */ @@ -201,9 +260,11 @@ public final class CompositePeerLauncher implements PeerLauncher { Supplier placementPolicy, Function liveCount, Supplier fleet, - BackendQuarantine quarantine) { + BackendQuarantine quarantine, + BackendOutagePolicy outagePolicy) { this.fleet = fleet; this.quarantine = Objects.requireNonNull(quarantine, "quarantine"); + this.outagePolicy = Objects.requireNonNull(outagePolicy, "outagePolicy"); if (delegates.isEmpty()) { throw new IllegalArgumentException("at least one peer adapter must be configured"); } @@ -266,7 +327,10 @@ 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); + // Checked in this order so exhaustion quarantine wins when both are active: quarantine + // throws first and short-circuits before the cool-off check ever runs (fleetd #201 Unit 5). enforceNotQuarantined(requestedProfile); + enforceNotCoolingOff(requestedProfile); enforceMaxLoad(requestedProfile); PeerHandle handle = d.spawn(req); spawnedBy.put(handle.id(), d); @@ -283,7 +347,10 @@ public final class CompositePeerLauncher implements PeerLauncher { // 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); + // fleetd #201 Unit 5: a distinct set from quarantined — see PlacementContext.coolingOff. + Set coolingOff = coolingOffProfiles(candidates); + PlacementContext ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable, + quarantined, coolingOff); int maxAttempts = candidates.isEmpty() ? 1 : candidates.size(); for (int attempt = 0; attempt < maxAttempts; attempt++) { @@ -313,7 +380,8 @@ 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, quarantined); + ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable, + quarantined, coolingOff); } } @@ -379,6 +447,31 @@ public final class CompositePeerLauncher implements PeerLauncher { .collect(Collectors.toSet()); } + /** + * Refuse an explicit-profile spawn whose credential is cooling off after repeated backend errors + * (fleetd #201 Unit 5 — {@link BackendOutagePolicy}): a SEPARATE, shorter-lived source from + * {@link #enforceNotQuarantined}'s exhaustion quarantine. Checked after quarantine so exhaustion + * wins when both are active — see the call site in {@link #spawn}. + * + * @throws PlacementException naming the profile, its credential, and the remaining cool-off + */ + private void enforceNotCoolingOff(String profile) { + String credentialId = credentialIdFor(profile); + outagePolicy.remainingCoolOffSeconds(credentialId).ifPresent(remaining -> { + throw new PlacementException("worker profile '" + profile + "' credential '" + credentialId + + "' is cooling off after repeated backend errors; ~" + remaining + + "s remaining — refusing spawn"); + }); + } + + /** The subset of {@code candidates} whose credential is currently cooling off (fleetd #201 Unit 5). */ + private Set coolingOffProfiles(List candidates) { + return candidates.stream() + .map(PlacementCandidate::profile) + .filter(p -> outagePolicy.remainingCoolOffSeconds(credentialIdFor(p)).isPresent()) + .collect(Collectors.toSet()); + } + private void enforceMaxLoad(String profile) { // Absent config, or a config whose maxLoad normalized to null (ABSENT ⇒ unlimited at load), // means no cap — never cap what wasn't configured. Note "non-positive ⇒ unlimited" was true diff --git a/fleetd/src/main/java/dev/ltms/fleet/placement/FixedPlacementPolicy.java b/fleetd/src/main/java/dev/ltms/fleet/placement/FixedPlacementPolicy.java index 121463c..b56635d 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/placement/FixedPlacementPolicy.java +++ b/fleetd/src/main/java/dev/ltms/fleet/placement/FixedPlacementPolicy.java @@ -1,55 +1,68 @@ package dev.ltms.fleet.placement; +import java.util.ArrayList; +import java.util.List; + /** * 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. * - *

Two exceptions walk past the default instead of returning it unconditionally: + *

Three exceptions walk past the default instead of returning it unconditionally: *

    *
  • Quarantine (CB-578 stage B): a quarantined default is a credential that just refused on * a usage limit, not a transient capacity or reachability concern. + *
  • Cooling off (fleetd #201 Unit 5): a credential cooling off after repeated backend errors + * ({@code BackendOutagePolicy}) — a separate, shorter-lived source from quarantine. When a + * profile is both quarantined and cooling off, only the quarantine reason is reported + * (exhaustion takes priority), matching {@code CompositePeerLauncher}'s explicit-spawn order. *
  • Weight 0 (CB-554): {@code fixed} is still automatic selection, so a profile the operator * marked "never auto-select me" ({@code weight <= 0}) must be skipped here exactly as * {@code weighted}/{@code round-robin} skip it — an explicit {@code fleet_spawn} naming * the profile is unaffected, only this automatic fallback walk. *
- * A fleet where nothing is ever quarantined or weight-0 never exercises either path, so today's - * behaviour is unchanged. + * A fleet where nothing is ever quarantined, cooling off, or weight-0 never exercises any of these + * paths, 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() && !ctx.quarantined().contains(d) && !weightExcluded(ctx, d)) { + if (d != null && !d.isBlank() && !ctx.quarantined().contains(d) && !ctx.coolingOff().contains(d) + && !weightExcluded(ctx, d)) { return new PlacementCandidate(d, null, 1.0f, null); } for (PlacementCandidate c : ctx.candidates()) { - if (!ctx.quarantined().contains(c.profile()) && !c.excluded()) { + if (!ctx.quarantined().contains(c.profile()) && !ctx.coolingOff().contains(c.profile()) + && !c.excluded()) { return new PlacementCandidate(c.profile(), null, c.weight(), c.maxLoad()); } } if (d != null && !d.isBlank()) { boolean dQuarantined = ctx.quarantined().contains(d); + // Exhaustion quarantine takes priority: reported only when quarantine is absent, so the + // message never claims "cooling off" for a profile that is really backend-exhausted. + boolean dCoolingOff = !dQuarantined && ctx.coolingOff().contains(d); boolean dWeightExcluded = weightExcluded(ctx, d); - if (dQuarantined && dWeightExcluded) { - throw new PlacementException("worker profile '" + d + "' is quarantined (backend " - + "exhausted) and has weight 0 (excluded from automatic selection), and no " - + "available candidate remains"); - } - if (dWeightExcluded) { - throw new PlacementException("worker profile '" + d + "' has weight 0 (excluded " - + "from automatic selection) and no available candidate remains"); - } - if (dQuarantined) { - throw new PlacementException("worker profile '" + d + "' is quarantined (backend " - + "exhausted) and no un-quarantined candidate is available"); + if (dQuarantined || dCoolingOff || dWeightExcluded) { + List reasons = new ArrayList<>(); + if (dQuarantined) { + reasons.add("is quarantined (backend exhausted)"); + } + if (dCoolingOff) { + reasons.add("is cooling off after repeated backend errors"); + } + if (dWeightExcluded) { + reasons.add("has weight 0 (excluded from automatic selection)"); + } + throw new PlacementException("worker profile '" + d + "' " + + String.join(" and ", reasons) + ", and no available candidate remains"); } } if (!ctx.candidates().isEmpty()) { - throw new PlacementException( - "all worker profiles are excluded from automatic selection (quarantined or weight-0)"); + throw new PlacementException("all worker profiles are excluded from automatic " + + "selection (quarantined, cooling off, or weight-0)"); } throw new PlacementException("no worker profiles configured"); } diff --git a/fleetd/src/main/java/dev/ltms/fleet/placement/PlacementContext.java b/fleetd/src/main/java/dev/ltms/fleet/placement/PlacementContext.java index 0d83314..0665c75 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/placement/PlacementContext.java +++ b/fleetd/src/main/java/dev/ltms/fleet/placement/PlacementContext.java @@ -14,10 +14,19 @@ import java.util.function.Function; * @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}. + * @param coolingOff profiles whose credential is currently cooling off after repeated backend + * errors (fleetd #201 Unit 5 — {@code BackendOutagePolicy}), a SEPARATE, + * shorter-lived source from {@code quarantined}: a credential outage cools off + * even when no member was ever exhausted. Deliberately its own set rather than + * merged into {@code quarantined} — {@link PlacementPolicyUtil} needs to tell + * the two apart so its refusal message says "cooling off", not "exhausted", + * when only this one is active. A profile can be in both sets at once; when it + * is, exhaustion quarantine is reported (it takes priority). */ public record PlacementContext(String defaultProfile, List candidates, Function liveCount, Set unreachable, - Set quarantined) { + Set quarantined, + Set coolingOff) { } diff --git a/fleetd/src/main/java/dev/ltms/fleet/placement/PlacementPolicyUtil.java b/fleetd/src/main/java/dev/ltms/fleet/placement/PlacementPolicyUtil.java index f402479..5cc2e77 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/placement/PlacementPolicyUtil.java +++ b/fleetd/src/main/java/dev/ltms/fleet/placement/PlacementPolicyUtil.java @@ -14,14 +14,16 @@ final class PlacementPolicyUtil { /** * Candidates that are not weight-excluded (CB-554: explicit {@code weight <= 0}, checked * first because it is a static config choice rather than transient state), not - * known-unreachable, not quarantined (CB-578 stage B), and have not reached their maxLoad. - * A {@code null} maxLoad means unlimited. + * known-unreachable, not quarantined (CB-578 stage B), not cooling off after repeated backend + * errors (fleetd #201 Unit 5 — a separate, shorter-lived source from quarantine), 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 (c.excluded() || ctx.unreachable().contains(c.profile()) - || ctx.quarantined().contains(c.profile())) { + || ctx.quarantined().contains(c.profile()) + || ctx.coolingOff().contains(c.profile())) { continue; } Integer cap = c.maxLoad(); @@ -38,21 +40,27 @@ final class PlacementPolicyUtil { /** * Build a clear exception describing why every candidate was dropped: all weight-0, all - * quarantined, all at capacity, all unreachable, or a mix. Each candidate is counted into - * exactly one bucket (weight-excluded takes priority) so a candidate excluded for more than - * one reason is never double-counted. + * quarantined, all cooling off, all at capacity, all unreachable, or a mix. Each candidate is + * counted into exactly one bucket (weight-excluded first, then quarantined, then cooling off) + * so a candidate excluded for more than one reason is never double-counted — a candidate that is + * both quarantined (CB-578 stage B, backend exhausted) and cooling off (fleetd #201 Unit 5, + * repeated backend errors) counts only as quarantined, matching {@code CompositePeerLauncher}'s + * explicit-spawn ordering: exhaustion quarantine takes priority when both are active. */ static PlacementException emptyException(PlacementContext ctx) { int weightExcluded = 0; int atCap = 0; int unreachable = 0; int quarantined = 0; + int coolingOff = 0; for (PlacementCandidate c : ctx.candidates()) { Integer cap = c.maxLoad(); if (c.excluded()) { weightExcluded++; } else if (ctx.quarantined().contains(c.profile())) { quarantined++; + } else if (ctx.coolingOff().contains(c.profile())) { + coolingOff++; } else if (ctx.unreachable().contains(c.profile())) { unreachable++; } else if (cap != null && ctx.liveCount().apply(c.profile()) >= cap) { @@ -71,6 +79,10 @@ final class PlacementPolicyUtil { if (quarantined == total) { return new PlacementException("all worker profiles are quarantined (backend exhausted)"); } + if (coolingOff == total) { + return new PlacementException( + "all worker profiles are cooling off after repeated backend errors"); + } if (atCap == total) { return new PlacementException("all worker profiles are at maxLoad"); } @@ -79,7 +91,9 @@ final class PlacementPolicyUtil { } return new PlacementException("no worker profile available: " + atCap + " at maxLoad, " + unreachable + " unreachable, " + quarantined + " quarantined, " + + coolingOff + " cooling off, " + weightExcluded + " weight-0, " - + (total - atCap - unreachable - quarantined - weightExcluded) + " remaining"); + + (total - atCap - unreachable - quarantined - coolingOff - weightExcluded) + + " remaining"); } } diff --git a/fleetd/src/test/java/dev/ltms/fleet/config/ConfigRefTest.java b/fleetd/src/test/java/dev/ltms/fleet/config/ConfigRefTest.java index 44ea0fe..dceccd6 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/config/ConfigRefTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/config/ConfigRefTest.java @@ -386,6 +386,54 @@ class ConfigRefTest { assertEquals("rate limit exceeded", ref.get().profiles().get("sonnet").exhaustedPattern()); } + /** + * fleetd #201 Unit 5: errorPattern is compiled once into Fleetd.main's backend-error pattern + * map at startup (see BackendErrorPatternLookup), the same way exhaustedPattern is above — a + * reload never re-reads it either, so a changed value must be reported deferred. + */ + @Test + void changingAProfilesErrorPatternIsReportedAsDeferred(@TempDir Path dir) throws Exception { + Path f = dir.resolve("fleetd.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 + errorPattern: "credential outage" + 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 + errorPattern: "provider 5xx" + 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("provider 5xx", ref.get().profiles().get("sonnet").errorPattern()); + } + @Test void aFixedRefHasNoFileAndRefusesToReload() { FleetConfig cfg = new FleetConfig(null, null, null, null, null, null, diff --git a/fleetd/src/test/java/dev/ltms/fleet/config/FleetConfigTest.java b/fleetd/src/test/java/dev/ltms/fleet/config/FleetConfigTest.java index 7a4f09f..52de820 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/config/FleetConfigTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/config/FleetConfigTest.java @@ -86,6 +86,85 @@ class FleetConfigTest { "unset means off — today's behaviour, unchanged"); } + // ── fleetd #201 Unit 5: errorPattern ──────────────────────────────────────────────────────── + + @Test + void aProfileWithAnErrorPatternBindsAndHasErrorPatternIsTrue(@TempDir Path dir) throws Exception { + Path f = dir.resolve("error-pattern.yaml"); + Files.writeString(f, """ + profiles: + ltms-local: + baseUrl: http://gx00.gw:8000 + errorPattern: "(?i)\\\\bcredential outage\\\\b" + """); + + FleetConfig.Profile w = FleetConfig.load(f).profiles().get("ltms-local"); + assertTrue(w.hasErrorPattern(), "errorPattern binds and enables fleetd #201 Unit 5 classification"); + assertEquals("(?i)\\bcredential outage\\b", w.errorPattern()); + } + + @Test + void aProfileWithNoErrorPatternLeavesItNullAndHasErrorPatternIsFalse(@TempDir Path dir) throws Exception { + Path f = dir.resolve("no-error-pattern.yaml"); + Files.writeString(f, """ + profiles: + ltms-local: + baseUrl: http://gx00.gw:8000 + """); + + FleetConfig.Profile w = FleetConfig.load(f).profiles().get("ltms-local"); + assertNull(w.errorPattern(), "unset means 'use the built-in fallback' — never off"); + assertFalse(w.hasErrorPattern()); + } + + @Test + void aBlankErrorPatternNormalizesToNullJustLikeUnset(@TempDir Path dir) throws Exception { + Path f = dir.resolve("blank-error-pattern.yaml"); + Files.writeString(f, """ + profiles: + ltms-local: + baseUrl: http://gx00.gw:8000 + errorPattern: " " + """); + + FleetConfig.Profile w = FleetConfig.load(f).profiles().get("ltms-local"); + assertNull(w.errorPattern()); + assertFalse(w.hasErrorPattern()); + } + + @Test + void aProfileWithAMalformedErrorPatternIsRejectedAtLoadNamingTheProfileAndKey(@TempDir Path dir) + throws Exception { + Path f = dir.resolve("malformed-error-pattern.yaml"); + Files.writeString(f, """ + profiles: + ltms-local: + baseUrl: http://gx00.gw:8000 + errorPattern: "(unterminated[" + """); + + IllegalStateException e = assertThrows(IllegalStateException.class, () -> FleetConfig.load(f)); + assertTrue(e.getMessage().contains("ltms-local"), "the offending profile is named: " + e.getMessage()); + assertTrue(e.getMessage().contains("errorPattern"), "the offending key is named: " + e.getMessage()); + } + + @Test + void withProfileCarriesErrorPatternThrough(@TempDir Path dir) throws Exception { + Path f = dir.resolve("with-profile-error-pattern.yaml"); + Files.writeString(f, """ + profiles: + ltms-local: + baseUrl: http://gx00.gw:8000 + errorPattern: "credential outage" + """); + + // The compact constructor defaults `profile` to the map key via withProfile(name) — see + // loadsFullConfig above. errorPattern must survive that same rebuild. + FleetConfig.Profile w = FleetConfig.load(f).profiles().get("ltms-local"); + assertEquals("ltms-local", w.profile()); + assertEquals("credential outage", w.errorPattern()); + } + @Test void appliesDefaultsForMissingSections(@TempDir Path dir) throws Exception { Path f = dir.resolve("minimal.yaml"); @@ -1313,6 +1392,7 @@ class FleetConfigTest { weight: 0.5 maxLoad: 2 exhaustedPattern: "usage limit has been reached" + errorPattern: "credential outage" placement: weighted lifecycle: idleTtlSeconds: 300 @@ -1350,6 +1430,8 @@ class FleetConfigTest { assertEquals(2, w.maxLoad(), "maxLoad binds as an integer"); assertTrue(w.hasExhaustedPattern(), "exhaustedPattern binds and enables the CB-578 stage A classification"); assertEquals("usage limit has been reached", w.exhaustedPattern()); + assertTrue(w.hasErrorPattern(), "errorPattern binds and enables fleetd #201 Unit 5 classification"); + assertEquals("credential outage", w.errorPattern()); assertEquals("weighted", cfg.placement(), "placement binds at the top level"); assertEquals(300, cfg.lifecycle().idleTtlSeconds()); diff --git a/fleetd/src/test/java/dev/ltms/fleet/inject/BackendOutageFlowTest.java b/fleetd/src/test/java/dev/ltms/fleet/inject/BackendOutageFlowTest.java new file mode 100644 index 0000000..c82ed4e --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/inject/BackendOutageFlowTest.java @@ -0,0 +1,295 @@ +package dev.ltms.fleet.inject; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import dev.ltms.fleet.config.FleetConfig; +import dev.ltms.fleet.guard.SubscriptionGuard; +import dev.ltms.fleet.herdr.AgentControl; +import dev.ltms.fleet.herdr.FakeHerdr; +import dev.ltms.fleet.herdr.HerdrClient; +import dev.ltms.fleet.herdr.WorkspaceControl; +import dev.ltms.fleet.mcp.PrimaryRegistry; +import dev.ltms.fleet.member.ClaudeCodeLauncher; +import dev.ltms.fleet.member.CompositePeerLauncher; +import dev.ltms.fleet.msg.InMemoryReplyInbox; +import dev.ltms.fleet.msg.Rendezvous; +import dev.ltms.fleet.msg.ReplyPushLoop; +import dev.ltms.fleet.peer.SpawnRequest; +import dev.ltms.fleet.placement.BackendOutagePolicy; +import dev.ltms.fleet.placement.BackendQuarantine; +import dev.ltms.fleet.placement.PlacementException; +import dev.ltms.fleet.placement.PlacementPolicies; +import dev.ltms.fleet.session.MemberSession; +import dev.ltms.fleet.session.SessionManager; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; + +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.Set; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicReference; +import java.util.regex.Pattern; + +import static org.junit.jupiter.api.Assertions.*; + +/** + * fleetd #201 / #227 Unit 5: proves the WIRED production path — a real {@link CompletionResolver} + * pane-scrape classification, through the exact {@code BackendErrorSink} ordering {@code Fleetd.main} + * builds (mirrored here line for line since the lambda itself lives inline in {@code Fleetd.main}), + * into a real {@link BackendOutagePolicy}, a real {@link CompositePeerLauncher} spawn gate, and a + * real {@link ReplyPushLoop} lead nudge — never by calling {@link BackendOutagePolicy#record} or the + * sink directly, which is what the finer-grained unit tests in + * {@code dev.ltms.fleet.member.CompositePeerLauncherTest} and {@code ReplyPushLoopTest} do for their + * own concerns. This class exists to catch a wiring mistake those unit tests cannot see: each of them + * hands the collaborator a value it assumes was computed correctly one layer up. + * + *

{@code fleet_list}/{@code fleet_profiles} rendering of the resulting cool-off (criteria 11/12) + * is proven separately in {@code dev.ltms.fleet.mcp.FleetMcpTest} — {@code FleetMcp}'s view methods + * are package-private to {@code dev.ltms.fleet.mcp} and this class needs {@code CompletionResolver}'s + * package-private {@code InFlight}, so the two cannot share one test class. What matters for THIS + * test is only that the same {@link BackendOutagePolicy} object the real scrape path populates is the + * one {@code FleetMcp} reads from — proven by asserting directly against + * {@link BackendOutagePolicy#remainingCoolOffSeconds} after the real classification below. + */ +class BackendOutageFlowTest { + + private final List schedulers = new ArrayList<>(); + + @AfterEach + void tearDown() { + schedulers.forEach(ScheduledExecutorService::shutdownNow); + } + + private static FleetConfig.Profile stubWorker(String profile, String credentialId) { + return new FleetConfig.Profile(profile, "http://gx00.gw:8000", "coder", + null, "FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers", + "w #{n}", null, null, null, null, null, null, null, null, null, + null, null, credentialId, null); + } + + /** Minimal recording {@code HerdrClient} for the LEAD pane — mirrors ReplyPushLoopTest's own. */ + private static final class RecordingLeadClient implements HerdrClient { + private static final ObjectMapper MAPPER = new ObjectMapper(); + private final List prompts = new CopyOnWriteArrayList<>(); + volatile CountDownLatch sendLatch = new CountDownLatch(1); + + @Override + public JsonNode call(String method, Object params) { + if ("agent.get".equals(method)) { + return MAPPER.createObjectNode().set("agent", MAPPER.createObjectNode() + .put("terminal_id", "term_primary").put("agent_status", "idle")); + } + if ("agent.prompt".equals(method)) { + prompts.add(params); + sendLatch.countDown(); + } + return MAPPER.createObjectNode(); + } + + @Override + public void close() { + } + + int sendCount() { + return prompts.size(); + } + + String lastText() { + return ((Map) prompts.get(prompts.size() - 1)).get("text").toString(); + } + } + + /** Every collaborator the flow wires together, built exactly once per test. */ + private final class Flow { + final FakeHerdr herdr = new FakeHerdr() + .readText("⏺ 503 Service Unavailable: upstream credential rejected\n❯ "); + final Map profiles = orderedProfiles(); + final AtomicLong clockNanos = new AtomicLong(0L); + final BackendOutagePolicy outagePolicy = new BackendOutagePolicy(clockNanos::get); + final CompositePeerLauncher workers; + final SessionManager sessions; + final MemberSession s1; + final MemberSession s2; + final PrimaryRegistry registry = new PrimaryRegistry(null); + final InMemoryReplyInbox inbox = new InMemoryReplyInbox(); + final RecordingLeadClient leadClient = new RecordingLeadClient(); + final ReplyPushLoop pushLoop; + final CompletionResolver resolver; + final Rendezvous rendezvous = new Rendezvous(); + + Flow() { + ClaudeCodeLauncher adapter = new ClaudeCodeLauncher(new AgentControl(herdr), + new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")), + profiles, "terra", _ -> "tok"); + workers = new CompositePeerLauncher(List.of(adapter), "terra", profiles, + PlacementPolicies.weighted(), _ -> 0, null, BackendQuarantine.none(), outagePolicy); + sessions = new SessionManager(workers); + // Two REAL member sessions on "terra" — the roster lookup the production sink does. + s1 = sessions.acquire("terra", null, null, null); + s2 = sessions.acquire("terra", null, null, null); + registry.recordDelegation(s1.terminalId(), "term_primary"); + registry.recordDelegation(s2.terminalId(), "term_primary"); + inbox.own(s1.terminalId()); + inbox.own(s2.terminalId()); + + ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(); + schedulers.add(scheduler); + pushLoop = new ReplyPushLoop(registry, new AgentControl(leadClient), inbox, scheduler, 3, 50); + AtomicReference pushLoopRef = new AtomicReference<>(pushLoop); + + // --- mirrors Fleetd.main's backendErrorSink lambda EXACTLY: (1) mark BACKEND_ERROR, + // (2) resolve profile/credential via the roster, fail-loud + notify unmapped-target, + // (3) record in BackendOutagePolicy, (4) on a NEW incident, notify the lead. ------------ + BackendErrorSink backendErrorSink = (target, matchedLine, reason) -> { + sessions.onBackendError(target, reason); + + String profileName = sessions.roster().stream() + .filter(session -> target.equals(session.terminalId())) + .findFirst() + .map(MemberSession::profile) + .orElse(null); + FleetConfig.Profile profile = profileName == null ? null : profiles.get(profileName); + if (profile == null) { + ReplyPushLoop loop = pushLoopRef.get(); + if (loop != null) { + loop.onBackendTargetUnmapped(target, reason); + } + return; + } + String credentialId = profile.effectiveCredentialId(); + Optional incident = outagePolicy.record(credentialId, target, reason); + incident.ifPresent(inc -> { + List affectedProfiles = profiles.values().stream() + .filter(p -> credentialId.equals(p.effectiveCredentialId())) + .map(FleetConfig.Profile::profile) + .sorted() + .toList(); + ReplyPushLoop loop = pushLoopRef.get(); + if (loop != null) { + loop.onBackendIncident(inc.id(), inc.targets(), credentialId, affectedProfiles, + (int) inc.remainingCoolOffSeconds()); + } + }); + }; + BackendErrorPatternLookup patterns = target -> Pattern.compile("(?i)503 Service Unavailable"); + resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), + ExhaustionSink.none(), patterns, backendErrorSink); + } + + /** Drives one real pane-scrape classification for {@code target} through {@code resolver}. */ + void classifyBackendErrorOn(String target) { + var waiter = rendezvous.open(target); + resolver.resolve(target, new CompletionResolver.InFlight(waiter, null)); + assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(), + "a classified backend error still resolves the blocked send as a failure"); + } + } + + private static Map orderedProfiles() { + Map m = new LinkedHashMap<>(); + m.put("terra", stubWorker("terra", "shared-openai")); + m.put("sol", stubWorker("sol", "shared-openai")); + return m; + } + + private long agentStartCalls(FakeHerdr herdr) { + return herdr.calls.stream().filter(c -> c.method().equals("agent.start")).count(); + } + + // --- criteria 4/5: two DISTINCT targets start an incident; one target never does ------------ + + @Test + void aSingleClassifiedBackendErrorOnOneTargetNeverCoolsTheCredential() { + Flow flow = new Flow(); + + flow.classifyBackendErrorOn(flow.s1.terminalId()); + + assertTrue(flow.outagePolicy.remainingCoolOffSeconds("shared-openai").isEmpty(), + "one distinct target's classified error must not start a cool-off"); + assertEquals(0, flow.leadClient.sendCount(), "no incident yet, so no lead nudge"); + assertDoesNotThrow(() -> flow.workers.spawn(new SpawnRequest("terra", null, null)), + "terra must still be spawnable after only one target's error"); + } + + @Test + void twoDistinctTargetsClassifiedThroughTheRealScrapeStartAnIncidentAndCoolTheCredential() throws Exception { + Flow flow = new Flow(); + + flow.classifyBackendErrorOn(flow.s1.terminalId()); + assertTrue(flow.outagePolicy.remainingCoolOffSeconds("shared-openai").isEmpty()); + + flow.classifyBackendErrorOn(flow.s2.terminalId()); + + assertTrue(flow.leadClient.sendLatch.await(3, TimeUnit.SECONDS), + "the second distinct target must cross the threshold and nudge the lead"); + Thread.sleep(150); + var remaining = flow.outagePolicy.remainingCoolOffSeconds("shared-openai"); + assertTrue(remaining.isPresent(), "two distinct targets must start a cool-off"); + assertEquals(60L, remaining.getAsLong(), + "this is the exact fact FleetMcp.capacityView/profiles reads as coolingOffForSeconds " + + "(see FleetMcpTest for the rendering side)"); + } + + // --- criterion 13: the incident reaches the lead through ReplyPushLoop exactly once ---------- + + @Test + void theIncidentSendsExactlyOneNudgeToTheLeadNamingBothCredentialAndTargets() throws Exception { + Flow flow = new Flow(); + + flow.classifyBackendErrorOn(flow.s1.terminalId()); + flow.classifyBackendErrorOn(flow.s2.terminalId()); + + assertTrue(flow.leadClient.sendLatch.await(3, TimeUnit.SECONDS)); + Thread.sleep(200); + assertEquals(1, flow.leadClient.sendCount(), "one incident must send exactly one nudge"); + String nudge = flow.leadClient.lastText(); + assertTrue(nudge.contains("shared-openai"), nudge); + assertTrue(nudge.contains(flow.s1.terminalId()) && nudge.contains(flow.s2.terminalId()), nudge); + assertTrue(nudge.contains("60"), nudge); + } + + // --- criterion 6: explicit spawn is refused BEFORE the adapter is ever called ----------------- + + @Test + void explicitSpawnOnACooledCredentialIsRefusedBeforeReachingTheAdapter() throws Exception { + Flow flow = new Flow(); + flow.classifyBackendErrorOn(flow.s1.terminalId()); + flow.classifyBackendErrorOn(flow.s2.terminalId()); + assertTrue(flow.leadClient.sendLatch.await(3, TimeUnit.SECONDS)); + + long startsBefore = agentStartCalls(flow.herdr); + PlacementException e = assertThrows(PlacementException.class, + () -> flow.workers.spawn(new SpawnRequest("terra", null, null))); + + assertTrue(e.getMessage().contains("cooling off"), e.getMessage()); + assertTrue(e.getMessage().contains("shared-openai"), e.getMessage()); + assertEquals(startsBefore, agentStartCalls(flow.herdr), + "the refusal must happen before the adapter is ever called — no new agent.start"); + } + + // --- criteria 7/9: automatic placement skips every profile sharing the cooled credential ----- + + @Test + void automaticPlacementRefusesAndNamesCoolingOffWhenEveryCandidateSharesTheCredential() throws Exception { + Flow flow = new Flow(); + flow.classifyBackendErrorOn(flow.s1.terminalId()); + flow.classifyBackendErrorOn(flow.s2.terminalId()); + assertTrue(flow.leadClient.sendLatch.await(3, TimeUnit.SECONDS)); + + PlacementException e = assertThrows(PlacementException.class, + () -> flow.workers.spawn(new SpawnRequest(null, null, null)), + "terra and sol both share the cooled-off credential — nothing is available"); + assertTrue(e.getMessage().contains("cooling off"), e.getMessage()); + assertThrows(PlacementException.class, () -> flow.workers.spawn(new SpawnRequest("sol", null, null)), + "sol shares terra's credential, so it must be locked out too"); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java index 5f97c77..600e3e7 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java @@ -20,6 +20,7 @@ import dev.ltms.fleet.session.MemberSession; import dev.ltms.fleet.session.WorktreeRequest; import dev.ltms.fleet.member.ClaudeCodeLauncher; import dev.ltms.fleet.member.CompositePeerLauncher; +import dev.ltms.fleet.placement.BackendOutagePolicy; import dev.ltms.fleet.placement.BackendQuarantine; import dev.ltms.fleet.placement.PlacementPolicies; import io.modelcontextprotocol.spec.McpSchema; @@ -500,6 +501,25 @@ class FleetMcpTest { assertTrue(out.contains("\"quarantinedForSeconds\":1800"), out); } + /** fleetd #201 Unit 5: {@code coolingOff} is a SEPARATE map from {@code quarantined}. */ + @Test + void profilesReportsACoolingOffCredentialInASeparateMap() { + FakeHerdr h = new FakeHerdr(); + BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L); + outagePolicy.record("shared-openai", "t1", "API Error: rate limited"); + outagePolicy.record("shared-openai", "t2", "API Error: rate limited"); // 2nd distinct target starts the incident + FleetMcp.OutageSource outage = new FleetMcp.OutageSource( + profile -> "ltms-local".equals(profile) ? "shared-openai" : null, outagePolicy); + McpSchema.CallToolResult res = FleetMcp.profiles( + workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), FleetMcp.QuarantineSource.none(), outage); + assertNotEquals(Boolean.TRUE, res.isError()); + String out = textOf(res); + assertTrue(out.contains("\"coolingOff\""), out); + assertTrue(out.contains("shared-openai"), out); + assertTrue(out.contains("\"coolingOffForSeconds\":60"), out); + assertFalse(out.contains("\"quarantined\""), "nothing is exhaustion-quarantined here: " + out); + } + @Test void listReportsTrackedWorkers() { FakeHerdr h = new FakeHerdr(); @@ -631,6 +651,81 @@ class FleetMcpTest { assertEquals(2, out.split("\"free\":0", -1).length - 1, out); } + /** fleetd #201 Unit 5: cool-off forces {@code free:0} but never adds {@code quarantinedForSeconds}. */ + @Test + void coolingOffProfileReportsZeroFreeButNeverQuarantinedForSeconds() { + FakeHerdr h = new FakeHerdr(); + SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))); + BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L); + outagePolicy.record("shared-openai", "t1", "API Error: rate limited"); + outagePolicy.record("shared-openai", "t2", "API Error: rate limited"); + FleetMcp.OutageSource outage = new FleetMcp.OutageSource( + profile -> "terra".equals(profile) ? "shared-openai" : null, outagePolicy); + + String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), + sessions, null, new FleetMcp.CapacitySource(profile -> 0, profile -> 2, + () -> Set.of("terra"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"), + FleetMcp.QuarantineSource.none(), outage, Map.of(), "")); + + assertTrue(out.contains("\"profile\":\"terra\""), out); + assertTrue(out.contains("\"free\":0"), out); + assertTrue(out.contains("\"credentialId\":\"shared-openai\""), out); + assertTrue(out.contains("\"coolingOffForSeconds\":60"), out); + assertFalse(out.contains("quarantinedForSeconds"), + "exhaustion quarantine was never active — this key must not appear: " + out); + } + + /** + * The two checks are independent — a profile can be BOTH exhaustion-quarantined AND cooling off + * for the same credential at once, and both maps/fields report it simultaneously. + */ + @Test + void bothQuarantineAndCoolingOffCanFireForTheSameProfileAtOnce() { + FakeHerdr h = new FakeHerdr(); + SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))); + BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(20)); + quarantine.quarantine("shared-openai"); + FleetMcp.QuarantineSource quarantineSource = new FleetMcp.QuarantineSource( + profile -> "terra".equals(profile) ? "shared-openai" : null, quarantine); + BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L); + outagePolicy.record("shared-openai", "t1", "API Error: rate limited"); + outagePolicy.record("shared-openai", "t2", "API Error: rate limited"); + FleetMcp.OutageSource outage = new FleetMcp.OutageSource( + profile -> "terra".equals(profile) ? "shared-openai" : null, outagePolicy); + + String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), + sessions, null, new FleetMcp.CapacitySource(profile -> 0, profile -> 2, + () -> Set.of("terra"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"), + quarantineSource, outage, Map.of(), "")); + + assertTrue(out.contains("\"free\":0"), out); + assertTrue(out.contains("\"quarantinedForSeconds\":1200"), out); + assertTrue(out.contains("\"coolingOffForSeconds\":60"), out); + } + + @Test + void listAndProfilesAgreeOnWhatIsCoolingOff() { + FakeHerdr h = new FakeHerdr(); + SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))); + BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L); + outagePolicy.record("shared-openai", "t1", "API Error: rate limited"); + outagePolicy.record("shared-openai", "t2", "API Error: rate limited"); + FleetMcp.OutageSource outage = new FleetMcp.OutageSource( + profile -> "ltms-local".equals(profile) ? "shared-openai" : null, outagePolicy); + var workers = workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")); + + String listOut = textOf(FleetMcp.listFleet(workers, sessions, null, + new FleetMcp.CapacitySource(profile -> 0, profile -> 2, () -> Set.of("ltms-local"), () -> 0), + new FleetMcp.HealthCoverageSource(() -> "off"), FleetMcp.QuarantineSource.none(), outage, + Map.of(), "")); + String profilesOut = textOf(FleetMcp.profiles(workers, FleetMcp.QuarantineSource.none(), outage)); + + assertTrue(listOut.contains("\"free\":0"), listOut); + assertTrue(listOut.contains("\"credentialId\":\"shared-openai\""), listOut); + assertTrue(profilesOut.contains("\"coolingOff\""), profilesOut); + assertTrue(profilesOut.contains("shared-openai"), profilesOut); + } + @Test void listAndProfilesAgreeOnWhatIsQuarantined() { FakeHerdr h = new FakeHerdr(); diff --git a/fleetd/src/test/java/dev/ltms/fleet/member/CompositePeerLauncherTest.java b/fleetd/src/test/java/dev/ltms/fleet/member/CompositePeerLauncherTest.java index 1e791d5..b50d0cd 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/member/CompositePeerLauncherTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/member/CompositePeerLauncherTest.java @@ -17,6 +17,7 @@ import dev.ltms.fleet.peer.PeerHandle; import dev.ltms.fleet.peer.PeerLauncher; import dev.ltms.fleet.peer.PeerUnreachableException; import dev.ltms.fleet.peer.SpawnRequest; +import dev.ltms.fleet.placement.BackendOutagePolicy; import dev.ltms.fleet.placement.BackendQuarantine; import dev.ltms.fleet.placement.PlacementException; import dev.ltms.fleet.placement.PlacementPolicies; @@ -911,4 +912,185 @@ class CompositePeerLauncherTest { assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("claude", null, null))); assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("gemini", null, null))); } + + // ── fleetd #201 Unit 5: repeated backend errors cool a credential off — a SEPARATE, shorter- ── + // ── lived mechanism from CB-578 stage B quarantine above, never merged with it ─────────────── + + @Test + void explicitSpawnOntoACoolingOffProfileIsRefusedNamingTheCredentialAndRemainingTime() { + 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()); + BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L); + // Two DISTINCT targets on the same credential — evidenceCount() counts distinct targets, + // never raw events, so one target repeating a classified error can never mint an incident. + outagePolicy.record("shared-openai", "term-1", "API Error: 500"); + assertTrue(outagePolicy.record("shared-openai", "term-2", "API Error: 500").isPresent(), + "the second distinct target crosses the threshold and starts the cool-off"); + CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles, + PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(), outagePolicy); + + 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("cooling off"), "message says cooling off, not exhausted: " + e.getMessage()); + assertTrue(e.getMessage().contains("60"), "message names roughly when it lifts: " + e.getMessage()); + assertEquals(0, adapter.spawnCount("sol"), "the cooling-off profile is never delegated to"); + } + + /** + * A single error from one target must never cool a credential off — the threshold is on + * distinct targets. This proves the spawn gate really reads {@link BackendOutagePolicy}'s + * result rather than reacting to any recorded event. + */ + @Test + void explicitSpawnIsNotRefusedAfterOnlyOneBackendErrorOnOneTarget() { + 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()); + BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L); + assertTrue(outagePolicy.record("shared-openai", "term-1", "API Error: 500").isEmpty(), + "one distinct target is below THRESHOLD=2 — no incident yet"); + CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles, + PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(), outagePolicy); + + assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("sol", null, null)), + "one target's error never cools the credential off"); + assertEquals(1, adapter.spawnCount("sol")); + } + + /** + * sol and terra share one OpenAI credential; cooling off because of an outage classified via + * ONE of them must lock out the other too, exactly like CB-578 stage B quarantine does. + */ + @Test + void twoProfilesSharingACredentialAreBothCoolingOffByOneOutageIncident() { + 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()); + BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L); + outagePolicy.record("shared-openai", "term-1", "API Error: 500"); + outagePolicy.record("shared-openai", "term-2", "API Error: 500"); + CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles, + PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(), outagePolicy); + + assertThrows(PlacementException.class, () -> composite.spawn(new SpawnRequest("sol", null, null))); + 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 placementSkipsACoolingOffProfileAndRoutesToAnotherOne() { + 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()); + BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L); + outagePolicy.record("shared-openai", "term-1", "API Error: 500"); + outagePolicy.record("shared-openai", "term-2", "API Error: 500"); + CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles, + PlacementPolicies.weighted(), _ -> 0, null, BackendQuarantine.none(), outagePolicy); + + PeerHandle h = composite.spawn(new SpawnRequest(null, null, null)); + assertEquals("b", h.profile(), "sol is cooling off, so an unqualified spawn must land on b"); + assertEquals(0, adapter.spawnCount("sol")); + } + + @Test + void automaticPlacementNamesCoolingOffWhenEveryCandidateIsCoolingOff() { + 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()); + BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L); + outagePolicy.record("shared-openai", "term-1", "API Error: 500"); + outagePolicy.record("shared-openai", "term-2", "API Error: 500"); + CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles, + PlacementPolicies.weighted(), _ -> 0, null, BackendQuarantine.none(), outagePolicy); + + // PlacementPolicy.select throws uncaught out of spawn() when available() is empty — never + // wrapped in PeerUnreachableException, which is only for a delegate that actually refused. + PlacementException e = assertThrows(PlacementException.class, + () -> composite.spawn(new SpawnRequest(null, null, null)), + "both candidates share the cooling-off credential — nothing is available"); + assertTrue(e.getMessage().contains("cooling off"), + "message says cooling off, not exhausted, since nothing here is quarantined: " + e.getMessage()); + } + + @Test + void aCoolOffLiftsOnTheInjectedClockAndTheProfileBecomesSpawnableAgain() { + 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); + BackendOutagePolicy outagePolicy = new BackendOutagePolicy(nowNanos::get); + outagePolicy.record("shared-openai", "term-1", "API Error: 500"); + outagePolicy.record("shared-openai", "term-2", "API Error: 500"); + CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles, + PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(), outagePolicy); + + assertThrows(PlacementException.class, () -> composite.spawn(new SpawnRequest("sol", null, null)), + "still inside the 60s cool-off"); + + nowNanos.set(TimeUnit.SECONDS.toNanos(61)); + + PeerHandle h = composite.spawn(new SpawnRequest("sol", null, null)); + assertEquals("sol", h.profile(), "the cool-off expired on the injected clock — sol is spawnable again"); + assertEquals(1, adapter.spawnCount("sol")); + } + + /** + * A credential that is BOTH exhaustion-quarantined (CB-578 stage B) and cooling off (fleetd + * #201 Unit 5) at once — the two checks are independent and can both be true for one profile. + * Exhaustion quarantine must win: the refusal names quarantine, never cooling off, because + * {@code enforceNotQuarantined} is checked (and throws) before {@code enforceNotCoolingOff} ever + * runs. + */ + @Test + void quarantineWinsOverCoolingOffWhenBothAreActiveOnTheSameCredential() { + 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"); + BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L); + outagePolicy.record("shared-openai", "term-1", "API Error: 500"); + outagePolicy.record("shared-openai", "term-2", "API Error: 500"); + CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles, + PlacementPolicies.fixed(), _ -> 0, null, quarantine, outagePolicy); + + PlacementException e = assertThrows(PlacementException.class, + () -> composite.spawn(new SpawnRequest("sol", null, null))); + assertTrue(e.getMessage().contains("quarantined"), + "exhaustion quarantine takes priority in the message: " + e.getMessage()); + assertFalse(e.getMessage().contains("cooling off"), + "cooling off is never mentioned when quarantine already refused the spawn: " + e.getMessage()); + } + + @Test + void aFleetWithNoBackendErrorsEverRecordedBehavesExactlyAsBeforeCoolOffExisted() { + // A never-record()-called BackendOutagePolicy is naturally inert — the 2/5/6-arg + // constructors used throughout this file all wire this in implicitly (NO_OUTAGE_POLICY). + // This test pins that an explicit fresh instance 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/fleetd/src/test/java/dev/ltms/fleet/placement/PlacementPolicyTest.java b/fleetd/src/test/java/dev/ltms/fleet/placement/PlacementPolicyTest.java index b377ec7..69e5de4 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/placement/PlacementPolicyTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/placement/PlacementPolicyTest.java @@ -30,7 +30,14 @@ class PlacementPolicyTest { private static PlacementContext ctx(List candidates, Function liveCount, Set unreachable, Set quarantined) { - return new PlacementContext("b", candidates, liveCount, unreachable, quarantined); + return ctx(candidates, liveCount, unreachable, quarantined, Set.of()); + } + + private static PlacementContext ctx(List candidates, + Function liveCount, + Set unreachable, Set quarantined, + Set coolingOff) { + return new PlacementContext("b", candidates, liveCount, unreachable, quarantined, coolingOff); } private static PlacementContext ctx(List candidates, @@ -52,14 +59,14 @@ class PlacementPolicyTest { PlacementPolicy policy = PlacementPolicies.fixed(); PlacementContext ctx = new PlacementContext(null, List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")), - noSessions(), Set.of(), Set.of()); + noSessions(), Set.of(), 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(), Set.of()); + PlacementContext ctx = new PlacementContext(null, List.of(), noSessions(), Set.of(), Set.of(), Set.of()); assertThrows(PlacementException.class, () -> policy.select(ctx)); } @@ -68,7 +75,7 @@ class PlacementPolicyTest { PlacementPolicy policy = PlacementPolicies.fixed(); PlacementContext ctx = new PlacementContext("b", List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")), - noSessions(), Set.of(), Set.of("b")); + noSessions(), Set.of(), Set.of("b"), Set.of()); assertEquals("a", policy.select(ctx).profile(), "the default 'b' is quarantined, so fixed falls through to the first un-quarantined candidate"); } @@ -78,7 +85,7 @@ class PlacementPolicyTest { PlacementPolicy policy = PlacementPolicies.fixed(); PlacementContext ctx = new PlacementContext("b", List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")), - noSessions(), Set.of(), Set.of("a", "b")); + noSessions(), Set.of(), Set.of("a", "b"), Set.of()); PlacementException e = assertThrows(PlacementException.class, () -> policy.select(ctx)); assertTrue(e.getMessage().contains("quarantined"), e.getMessage()); } @@ -211,6 +218,67 @@ class PlacementPolicyTest { assertTrue(e.getMessage().contains("quarantined"), e.getMessage()); } + // --- fleetd #201 Unit 5: coolingOff excludes a candidate, SEPARATE from quarantined -------- + + @Test + void roundRobinSkipsCoolingOffProfiles() { + PlacementPolicy policy = PlacementPolicies.roundRobin(); + List candidates = List.of( + PlacementCandidate.profile("a"), + PlacementCandidate.profile("b")); + PlacementContext ctx = ctx(candidates, noSessions(), Set.of(), Set.of(), Set.of("a")); + for (int i = 0; i < 5; i++) { + assertEquals("b", policy.select(ctx).profile(), "a is cooling off, so every pick lands on b"); + } + } + + @Test + void weightedSkipsCoolingOffProfile() { + 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(), Set.of("a")); + for (int i = 0; i < 5; i++) { + assertEquals("b", policy.select(ctx).profile(), "a is cooling off, so every pick lands on b"); + } + } + + @Test + void weightedThrowsWhenAllCoolingOff() { + 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(), Set.of("a", "b")))); + assertTrue(e.getMessage().contains("cooling off"), e.getMessage()); + assertFalse(e.getMessage().contains("quarantined"), + "nothing here is quarantined — the message must say cooling off, not exhausted: " + e.getMessage()); + } + + /** + * A candidate that is BOTH quarantined (CB-578 stage B) and cooling off (fleetd #201 Unit 5) + * counts only as quarantined — exhaustion takes priority, matching + * {@code CompositePeerLauncher}'s explicit-spawn check order. + */ + @Test + void quarantinedAndCoolingOffIsNotDoubleCounted() { + PlacementPolicy policy = PlacementPolicies.weighted(); + List candidates = List.of( + PlacementCandidate.profile("a"), + PlacementCandidate.profile("b")); + // a is BOTH quarantined and cooling off; b is quarantined only. If a were double-bucketed + // as cooling-off instead of quarantined, quarantined would undercount to 1 of 2 candidates + // and this would fall through to the generic "no worker profile available: N quarantined, M + // cooling off, ..." message instead — which also happens to contain the substring + // "quarantined", so a loose contains() check here would pass either way. Pinning the exact + // "all ... quarantined" message is what actually proves the two are not double-counted. + PlacementContext ctx = ctx(candidates, noSessions(), Set.of(), Set.of("a", "b"), Set.of("a")); + PlacementException e = assertThrows(PlacementException.class, () -> policy.select(ctx)); + assertEquals("all worker profiles are quarantined (backend exhausted)", e.getMessage()); + } + @Test void weightedThrowsWhenAllUnreachable() { PlacementPolicy policy = PlacementPolicies.weighted(); @@ -320,7 +388,7 @@ class PlacementPolicyTest { PlacementContext ctx = new PlacementContext("b", List.of(PlacementCandidate.profile("a", 1.0f, null), PlacementCandidate.profile("b", 0.0f, null)), - noSessions(), Set.of(), Set.of()); + noSessions(), Set.of(), Set.of(), Set.of()); assertEquals("a", policy.select(ctx).profile(), "the default 'b' has weight 0, so fixed falls through to the first available candidate"); } @@ -331,7 +399,7 @@ class PlacementPolicyTest { PlacementContext ctx = new PlacementContext("b", List.of(PlacementCandidate.profile("a", 0.0f, null), PlacementCandidate.profile("b", 0.0f, null)), - noSessions(), Set.of(), Set.of()); + noSessions(), Set.of(), Set.of(), Set.of()); PlacementException e = assertThrows(PlacementException.class, () -> policy.select(ctx)); assertTrue(e.getMessage().contains("weight 0"), e.getMessage()); }