fleetd #201/#227 Unit 5: wire the backend-error cool-off into config, placement and the MCP surface
A profile's credential that throws two distinct backend errors inside 60
seconds now cools off for 60 seconds. Automatic placement skips it,
an explicit fleet_spawn naming it is refused before the adapter is
called, and fleet_list/fleet_profiles report it as coolingOffForSeconds
next to the separate CB-578 quarantinedForSeconds.
Verified by the lead: see the merge check below. The worker ran 7
mutations, all killed with 0 compile errors; M7 was NOT killed on the
first pass (the assertion only checked .contains("quarantined"), which
is true of both the correct message and the mutated fallback), and the
worker strengthened it to assertEquals on the exact literal and kept
that change. That is the right call and it is reported honestly.
Two deviations, both justified in the PR:
- FixedPlacementPolicy needed the same coolingOff filter because it
filters candidates inline instead of using PlacementPolicyUtil.
- BackendOutageFlowTest sits in dev.ltms.fleet.inject because
CompletionResolver.InFlight is package-private there.
The startup coverage log line is a code-reading claim, not a captured
line from a live daemon. The worker said so rather than overclaiming.
PR #246, branch worker/cb201-unit5-wiring-6c12e6-8
This commit is contained in:
@@ -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):
|
||||
|
||||
@@ -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());
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user