fleetd #201/#227 unit 5: wire errorPattern, cool-off spawn gate, and fleet views
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Successful in 1m52s

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.
This commit is contained in:
Dai Ha
2026-09-03 11:45:11 +07:00
parent 26bafe824b
commit ba04b2359b
16 changed files with 1295 additions and 67 deletions
+12
View File
@@ -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):
+44 -2
View File
@@ -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
@@ -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<String, Pattern> 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<ReplyPushLoop> 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<BackendOutagePolicy.Incident> incident = outagePolicy.record(credentialId, target, reason);
incident.ifPresent(inc -> {
List<String> 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
@@ -43,7 +43,9 @@ import java.util.function.Supplier;
* which is constructed once), <em>and an existing profile's launch settings</em> —
* {@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<FleetConfig> {
// 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());
}
}
@@ -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
* <em>within</em>, 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<String> argv,
String placement, String workspace, String tabLabel, String mcpUrl,
String cwd, List<String> parityOverlay, String gitTokenEnv, String gitHostEnv,
String kind, Map<String, String> 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.
*
* <p>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<String> 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<String> KNOWN_KINDS = Set.of(Profile.KIND_CLAUDE_CODE, Profile.KIND_OPENCODE);
@@ -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<String, String> 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<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> 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<String, Object> result = new LinkedHashMap<>();
result.put("profiles", workers.profiles());
result.put("default", workers.defaultProfile() == null ? "" : workers.defaultProfile());
Map<String, Object> quarantined = new LinkedHashMap<>();
Map<String, Object> 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<String, Object> 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<String, Object> row = new LinkedHashMap<>();
row.put("credentialId", outageCredentialId);
row.put("coolingOffForSeconds", remaining);
coolingOff.put(profile, row);
});
}
quarantine.quarantine().remainingSeconds(credentialId).ifPresent(remaining -> {
Map<String, Object> 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<String, String> 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<String, String> 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<String, String> leads, String selfTerm, String selfCoordId) {
try {
Map<String, Agent> 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.
*
* <p>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<String, Object> capacityView(String profile, Function<String, Integer> liveCount,
Function<String, Integer> maxLoad, List<MemberSession> 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;
}
@@ -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<String, FleetConfig.Profile> profileConfigs,
PlacementPolicy placementPolicy,
Function<String, Integer> 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<String, Integer> 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<String, Integer> 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<HerdrPeerLauncher> delegates,
String defaultProfile,
Map<String, FleetConfig.Profile> profileConfigs,
PlacementPolicy placementPolicy,
Function<String, Integer> 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<FleetConfig> config,
Function<String, Integer> 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<HerdrPeerLauncher> delegates,
String defaultProfile,
Supplier<FleetConfig> config,
Function<String, Integer> 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> placementPolicy,
Function<String, Integer> liveCount,
Supplier<FleetConfig.Fleet> 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<String> 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<String> 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<String> coolingOffProfiles(List<PlacementCandidate> 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
@@ -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.
*
* <p>Two exceptions walk past the default instead of returning it unconditionally:
* <p>Three exceptions walk past the default instead of returning it unconditionally:
* <ul>
* <li>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.
* <li>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.
* <li>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.
* </ul>
* 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<String> 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");
}
@@ -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<PlacementCandidate> candidates,
Function<String, Integer> liveCount,
Set<String> unreachable,
Set<String> quarantined) {
Set<String> quarantined,
Set<String> coolingOff) {
}
@@ -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<PlacementCandidate> available(PlacementContext ctx) {
List<PlacementCandidate> 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");
}
}
@@ -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,
@@ -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());
@@ -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.
*
* <p>{@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<ScheduledExecutorService> 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<Object> 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<String, FleetConfig.Profile> 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<ReplyPushLoop> 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<BackendOutagePolicy.Incident> incident = outagePolicy.record(credentialId, target, reason);
incident.ifPresent(inc -> {
List<String> 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<String, FleetConfig.Profile> orderedProfiles() {
Map<String, FleetConfig.Profile> 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");
}
}
@@ -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();
@@ -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<String, FleetConfig.Profile> 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<String, FleetConfig.Profile> 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<String, FleetConfig.Profile> 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<String, FleetConfig.Profile> 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<String, FleetConfig.Profile> 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<String, FleetConfig.Profile> 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<String, FleetConfig.Profile> 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)));
}
}
@@ -30,7 +30,14 @@ class PlacementPolicyTest {
private static PlacementContext ctx(List<PlacementCandidate> candidates,
Function<String, Integer> liveCount,
Set<String> unreachable, Set<String> quarantined) {
return new PlacementContext("b", candidates, liveCount, unreachable, quarantined);
return ctx(candidates, liveCount, unreachable, quarantined, Set.of());
}
private static PlacementContext ctx(List<PlacementCandidate> candidates,
Function<String, Integer> liveCount,
Set<String> unreachable, Set<String> quarantined,
Set<String> coolingOff) {
return new PlacementContext("b", candidates, liveCount, unreachable, quarantined, coolingOff);
}
private static PlacementContext ctx(List<PlacementCandidate> 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<PlacementCandidate> 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<PlacementCandidate> 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<PlacementCandidate> 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<PlacementCandidate> 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());
}