Compare commits

...

1 Commits

Author SHA1 Message Date
Dai Ha e501d39988 CB-578 stage B: quarantine the exhausted credential, not the profile
CI / build (pull_request) Successful in 55s
CI / contract (pull_request) Successful in 1m21s
A BACKEND_EXHAUSTED classification (stage A) now puts that profile's
credential into a BackendQuarantine for a configurable cooldown. A spawn
onto a quarantined profile is refused naming the credential and roughly
when it lifts; weighted/round-robin/fixed placement skip a quarantined
candidate; the quarantine lifts itself on the injected clock; and it is
visible on bridge_profiles.

Keyed by credential, not by profile name, via the new Profile.credentialId
(profiles sharing one credential quarantine together — e.g. two models on
one account) and effectiveCredentialId() (unset ⇒ quarantines alone,
today's behaviour unchanged). Fixed a related gap along the way: a reload
changing exhaustedPattern was silently reported "applied" even though it's
deferred — sameLaunchSettings() now catches it too.
2026-08-15 10:29:54 +02:00
20 changed files with 886 additions and 79 deletions
+32 -8
View File
@@ -144,6 +144,17 @@ herdrSocket: ~/.config/herdr/herdr.sock
# omit and this profile's completion fallback behaves exactly as before.
# Every backend words its refusal differently, so this is config, never a
# vendor string baked into bridged itself.
# DEFERRED: compiled once into a startup pattern map — editing it needs a
# daemon restart, same as this profile's model/baseUrl/argv.
# credentialId → CB-578 stage B: the credential this profile quarantines WITH when a
# BACKEND_EXHAUSTED classification fires. Two profiles that set the SAME
# credentialId share one quarantine — the case this exists for is two models
# on one account (e.g. sol and terra both billing one OpenAI credential): an
# exhaustion on either one must lock out both, or the fleet just walks onto
# the same dead account under the sibling's name. Opt-in — omit and this
# profile quarantines alone, under its own name, exactly as if the field did
# not exist. Cooldown length is the top-level quarantineCooldownSeconds below.
# HOT: read live at every spawn/exhaustion check — no restart needed.
# env → extra environment for this profile's workers, as a literal key/value map
# (CB-511). Use it to give workers a toolchain.
#
@@ -178,6 +189,7 @@ profiles:
# gitTokenEnv: GITEA_TOKEN # opt-in: let this profile's workers open their own PR (CB-302)
# gitHostEnv: GITEA_HOST # defaults to GITEA_HOST; injected only with gitTokenEnv
# exhaustedPattern: "usage limit has been reached" # opt-in: classify a usage-limit refusal (CB-578)
# credentialId: shared-openai # opt-in: quarantine together with every other profile sharing this id (CB-578)
# configDir: /Users/me/.ccs/instances/gx10 # CLAUDE_CONFIG_DIR — inherit that profile's skills/MCP
# cwd: /Users/me/src/myrepo # pin the working dir; omit to inherit the primary's
# parityOverlay: [".claude/settings.local.json", ".env", ".envrc"] # never add .mcp.json — see above
@@ -243,6 +255,15 @@ profiles:
# round-robin, or weighted. Omitting this key is a strict no-op for existing configs.
placement: weighted
# How long a credential sits out after a BACKEND_EXHAUSTED classification (CB-578 stage B), in
# seconds, before a spawn may land on it again. Applies to every profile's effective credential
# (its own name, or its credentialId if set above) — there is no per-profile override. Default
# 1800 (30 minutes) when omitted or non-positive.
# DEFERRED: baked once into the BackendQuarantine built at startup — a running quarantine keeps
# its original cooldown regardless; a new value only applies to a quarantine that starts after a
# restart. Editing this needs a daemon restart to take effect.
# quarantineCooldownSeconds: 1800
# Re-read this file without restarting the daemon (CB-559). Off unless you add this block, so an
# upgraded bridged keeps the old behaviour: the file is read once at boot and never again.
# enabled → turn the watch on. bridged checks the file's modified time on a timer and
@@ -252,17 +273,20 @@ placement: weighted
# Not every key can move under a running daemon, and the difference is about what already exists
# when the reload happens — not about how important the key is:
# HOT → takes effect on the next spawn: the whole `fleet:` block (every role pool,
# `charters`, and `tabLabel`), `placement:`, and an existing profile's weight / maxLoad. Those are
# hot because the placement policy reads them through a supplier — being config is
# `charters`, and `tabLabel`), `placement:`, and an existing profile's weight / maxLoad
# / credentialId. Those are hot because the placement policy (and, for credentialId,
# the CB-578 stage B quarantine check) reads them through a supplier — being config is
# not by itself enough to make a key hot.
# DEFERRED → accepted into the new config, but the wiring built at startup keeps the old value
# until you restart: `lifecycle:`, `leadHeartbeat:`, `guard:`, `worktreeRoot:`,
# `spawnReadyTimeoutMs` / `spawnReadyPollMs`, ADDING or REMOVING a profile (a new
# backend needs its own launcher, and launchers are built once), AND an existing
# profile's launch settings — model, baseUrl, argv, env, configDir, mcpUrl, tabLabel.
# The launcher takes a copy of `profiles:` at startup and resolves every spawn out of
# that copy, so those never reach a launch until you restart. The reload logs them by
# name rather than pretending they applied.
# `spawnReadyTimeoutMs` / `spawnReadyPollMs`, `quarantineCooldownSeconds` (CB-578
# stage B — baked once into the quarantine tracker built at startup), ADDING or
# REMOVING a profile (a new backend needs its own launcher, and launchers are built
# once), AND an existing profile's launch settings — model, baseUrl, argv, env,
# configDir, mcpUrl, tabLabel, exhaustedPattern. The launcher takes a copy of
# `profiles:` at startup and resolves every spawn out of that copy, so those never
# reach a launch until you restart. The reload logs them by name rather than
# pretending they applied.
# COLD → cannot change at all: `bind:`, `herdrSocket:`, `broker:` and `auth:`. The socket is
# bound, the broker connection is open, and the auth mode decides who may reach the
# port that is already listening.
@@ -14,6 +14,7 @@ import dev.ltms.bridged.herdr.UnixSocketHerdrClient;
import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.inject.CompletionResolver;
import dev.ltms.bridged.inject.ExhaustedPatternLookup;
import dev.ltms.bridged.inject.ExhaustionSink;
import dev.ltms.bridged.inject.Injector;
import dev.ltms.bridged.inject.StatusPoller;
import dev.ltms.bridged.inject.TurnListener;
@@ -37,6 +38,7 @@ import dev.ltms.bridged.msg.LeadHeartbeatLoop;
import dev.ltms.bridged.msg.ReplyPushLoop;
import dev.ltms.bridged.rest.BridgedApp;
import dev.ltms.bridged.session.GitWorktrees;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.peer.PeerLauncher;
import dev.ltms.bridged.session.SessionReaper;
@@ -44,6 +46,7 @@ import dev.ltms.bridged.member.ClaudeCodeLauncher;
import dev.ltms.bridged.member.CompositePeerLauncher;
import dev.ltms.bridged.member.HerdrPeerLauncher;
import dev.ltms.bridged.member.OpenCodeLauncher;
import dev.ltms.bridged.placement.BackendQuarantine;
import io.javalin.Javalin;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -141,11 +144,18 @@ public final class Bridged {
() -> config.get().fleet()));
}
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
// CB-578 stage B: one quarantine tracker for the whole daemon, shared between the launcher
// (checked at spawn) and the exhaustion sink wired in below (written on BACKEND_EXHAUSTED).
// The cooldown is deferred (see BridgedConfig#quarantineCooldownSeconds): it is read once
// here, at startup, and a config reload only changes it for a daemon restart.
BackendQuarantine quarantine = new BackendQuarantine(System::nanoTime,
TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));
PeerLauncher workers = new CompositePeerLauncher(
adapters,
cfg.effectiveDefaultProfile(),
config,
profileName -> liveCountRef.get().apply(profileName));
profileName -> liveCountRef.get().apply(profileName),
quarantine);
// CB-504: under supervision (launchd/systemd) bridged can start before herdr's socket
// exists. The client itself is lazy — it connects per call — but the orphan reap below is
// the first thing that actually talks to herdr, so without this wait a boot-order race
@@ -272,7 +282,23 @@ public final class Bridged {
.orElse(null);
log.info("backend-exhausted classification (CB-578 stage A): {}",
CompletionResolver.coverage(cfg.profiles().keySet(), exhaustedPatternsByProfile.keySet()));
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns);
// CB-578 stage B: on a classification that actually wins, quarantine the exhausted profile's
// CREDENTIAL — not the profile name — so a profile sharing that credential (e.g. two models
// on one OpenAI account) is refused too, not just the one that happened to report it. Reads
// the profile config live off `config`, so a credentialId edit is hot: no restart needed.
ExhaustionSink exhaustionSink = (target, reason) -> sessions.roster().stream()
.filter(session -> target.equals(session.terminalId()))
.findFirst()
.map(MemberSession::profile)
.map(profileName -> config.get().profiles().get(profileName))
.ifPresent(profile -> {
String credentialId = profile.effectiveCredentialId();
quarantine.quarantine(credentialId);
log.warn("credential '{}' quarantined for {}s (profile '{}' classified "
+ "BACKEND_EXHAUSTED): {}", credentialId,
cfg.quarantineCooldownSeconds(), profile.profile(), reason);
});
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns, exhaustionSink);
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
MemberPresence presence = sessions.asPresence();
@@ -440,7 +466,11 @@ public final class Bridged {
var health = config.get().health();
return FleetHealthMonitor.coverage(health != null && health.isEnabled(),
health != null && health.notifications() != null && health.notifications().configured());
}));
}),
new BridgeMcp.QuarantineSource(profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.effectiveCredentialId();
}, quarantine));
// CB-559: opt-in config reload. With no `configReload:` block nothing is constructed, so an
// upgraded daemon behaves exactly as before — the file is read once at boot and never again.
@@ -58,6 +58,12 @@ import java.util.Set;
* {@code fixed} (default), {@code round-robin}, or {@code weighted}
* @param auth API authentication mode ({@code null} → {@code loopback-trust}, the
* historical behaviour), CB-501
* @param quarantineCooldownSeconds how long a credential stays quarantined after a
* {@code BACKEND_EXHAUSTED} classification (CB-578 stage B); {@code null}/{@code
* <=0} → {@link #DEFAULT_QUARANTINE_COOLDOWN_SECONDS}. Baked once into the
* {@code BackendQuarantine} built at startup, so it is DEFERRED: changing it
* needs a restart, and a quarantine already running keeps whatever cooldown was
* live when it started.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record BridgedConfig(
@@ -76,7 +82,11 @@ public record BridgedConfig(
Health health,
String placement,
Auth auth,
ConfigReload configReload) {
ConfigReload configReload,
Integer quarantineCooldownSeconds) {
/** Default cooldown (CB-578 stage B) when {@code quarantineCooldownSeconds} is absent/non-positive. */
public static final int DEFAULT_QUARANTINE_COOLDOWN_SECONDS = 1800;
/** Back-compat 14-arg form — no {@code configReload:} block, so file watching stays off. */
public BridgedConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
@@ -84,7 +94,7 @@ public record BridgedConfig(
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, String placement, Auth auth) {
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, null);
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, null, null);
}
/** Back-compat form before the optional {@code health:} block was added. */
@@ -93,7 +103,18 @@ public record BridgedConfig(
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, String placement, Auth auth, ConfigReload configReload) {
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, configReload);
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, configReload, null);
}
/** Back-compat form before the CB-578 stage B {@code quarantineCooldownSeconds} field was added. */
public BridgedConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
ConfigReload configReload) {
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
configReload, null);
}
/**
@@ -187,6 +208,17 @@ public record BridgedConfig(
* {@code null}/blank ⇒ the classification never fires for this profile and
* today's completion-fallback behaviour is unchanged. Every backend words
* its refusal differently, so this is config, never a vendor string in code.
* @param credentialId the shared account this profile authenticates as (CB-578 stage B). Two
* or more profiles setting the <em>same</em> non-blank value are quarantined
* together by one {@code exhaustedPattern} classification on any one of
* them — the case a single OpenAI (or any other) credential backing several
* profiles ({@code sol}, {@code terra}, ...) needs, so a spawn does not walk
* straight onto the other profile sharing the same exhausted account.
* {@code null}/blank ⇒ this profile's {@link #effectiveCredentialId()} is
* its own name, so it quarantines alone — today's behaviour for every
* profile that does not opt in. Read live off the current config, so it is
* HOT: a change takes effect on the next exhaustion classification / spawn,
* no restart needed.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Profile(String profile, String baseUrl, String model,
@@ -200,7 +232,8 @@ public record BridgedConfig(
Float weight,
Integer maxLoad,
Boolean subscription,
String exhaustedPattern) {
String exhaustedPattern,
String credentialId) {
/** Peer kind spawned by {@link dev.ltms.bridged.member.ClaudeCodeLauncher} (the default). */
public static final String KIND_CLAUDE_CODE = "claude-code";
@@ -246,6 +279,10 @@ public record BridgedConfig(
// exhaustedPattern stays null when unset/blank (opt-in) — no defaulting, no vendor
// wording: an unconfigured profile keeps today's completion-fallback behaviour exactly.
exhaustedPattern = (exhaustedPattern == null || exhaustedPattern.isBlank()) ? null : exhaustedPattern;
// credentialId stays null when unset/blank — effectiveCredentialId() is where the
// "quarantines alone" fallback actually lives, so today's behaviour needs no defaulting
// here at all.
credentialId = (credentialId == null || credentialId.isBlank()) ? null : credentialId;
}
/**
@@ -290,7 +327,7 @@ public record BridgedConfig(
public Profile withProfile(String p) {
return new Profile(p, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, subscription,
exhaustedPattern);
exhaustedPattern, credentialId);
}
/** True when this profile is served by the Claude Code adapter (the default kind). */
@@ -326,11 +363,38 @@ public record BridgedConfig(
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, null, null);
}
/**
* Backward-compatible constructor without the CB-578 stage B {@code credentialId} — the
* profile quarantines alone under its own name (see {@link #effectiveCredentialId()}). Keeps
* pre-stage-B call sites (and any YAML that omits the key) compiling and behaving identically.
*/
public Profile(String profile, String baseUrl, String model,
String configDir, String tokenEnv, List<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) {
this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad,
subscription, exhaustedPattern, null);
}
/** True when this profile's CB-578 stage A backend-exhausted classification is configured. */
public boolean hasExhaustedPattern() {
return exhaustedPattern != null;
}
/**
* The credential group this profile quarantines with (CB-578 stage B): the configured
* {@link #credentialId} when set, else this profile's own name — so an unconfigured profile
* quarantines alone, exactly as it did before this field existed. Two profiles that set the
* same non-blank {@code credentialId} share one quarantine: a {@code BACKEND_EXHAUSTED}
* classification on either one quarantines both.
*/
public String effectiveCredentialId() {
return (credentialId == null || credentialId.isBlank()) ? profile : credentialId;
}
/** True when this profile's workers are granted a forge token to open their own PR (CB-302). */
public boolean hasGitToken() {
return gitTokenEnv != null && !gitTokenEnv.isBlank();
@@ -850,7 +914,7 @@ public record BridgedConfig(
private static final Set<String> KNOWN_TOP_LEVEL_KEYS = Set.of(
"bind", "herdrSocket", "profiles", "guard", "worktreeRoot",
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
"leadHeartbeat", "health", "placement", "auth", "configReload");
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds");
/** Load and validate config from {@code path}. */
public static BridgedConfig load(Path path) {
@@ -1157,8 +1221,15 @@ public record BridgedConfig(
// configReload is left as-is: null is "off", and ConfigReload's own compact constructor
// defaults the fields of a block that IS present. Defaulting it here would start watching
// the file for every config that never asked to be watched.
// quarantineCooldownSeconds IS defaulted, unlike leadHeartbeat/configReload above: it has no
// separate on/off switch of its own (BackendQuarantine only ever quarantines a credential
// after a BACKEND_EXHAUSTED classification, which stays opt-in via exhaustedPattern), so a
// config that never mentions it should still get a sane cooldown rather than a null one.
Integer quarantineCooldown = (quarantineCooldownSeconds != null && quarantineCooldownSeconds > 0)
? quarantineCooldownSeconds : DEFAULT_QUARANTINE_COOLDOWN_SECONDS;
return new BridgedConfig(b, herdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload);
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
quarantineCooldown);
}
/**
@@ -37,10 +37,15 @@ import java.util.function.Supplier;
* {@code fleet.leaders} needs a restart, the same as any deferred key below.</li>
* <li><strong>Deferred</strong> — accepted into the new snapshot, but the wiring built at startup
* keeps the old value until a restart: {@code lifecycle:}, {@code leadHeartbeat:},
* {@code spawnReadyTimeoutMs} / {@code spawnReadyPollMs}, {@code guard:},
* {@code worktreeRoot:}, adding or removing a profile (a new backend needs its own launcher,
* {@code spawnReadyTimeoutMs} / {@code spawnReadyPollMs}, {@code quarantineCooldownSeconds}
* (CB-578 stage B — baked once into the {@code BackendQuarantine} built at startup),
* {@code guard:}, {@code worktreeRoot:}, adding or removing a profile (a new backend needs its own launcher,
* which is constructed once), <em>and an existing profile's launch settings</em> —
* {@code model}, {@code baseUrl}, {@code argv}, {@code env}, {@code mcpUrl} and the rest.
* {@code model}, {@code baseUrl}, {@code argv}, {@code env}, {@code mcpUrl},
* {@code exhaustedPattern} (CB-578 stage A — compiled once into {@code Bridged.main}'s
* pattern map at startup), and the rest. {@code credentialId} (CB-578 stage B) is NOT on
* this list — it is read live off the config supplier at every quarantine check and
* exhaustion event, exactly like {@code weight} / {@code maxLoad}, so it is hot instead.
* {@code HerdrPeerLauncher} takes {@code Map.copyOf(profiles)} at construction and resolves
* each spawn out of that copy, so those never reach a launch until the daemon restarts. A
* reload logs these rather than pretending they applied.</li>
@@ -215,6 +220,12 @@ public final class ConfigRef implements Supplier<BridgedConfig> {
|| !Objects.equals(old.spawnReadyPollMs(), fresh.spawnReadyPollMs())) {
changed.add("spawnReady*");
}
// CB-578 stage B: baked once into the BackendQuarantine built at startup — a running
// quarantine keeps its original cooldown regardless, and a new cooldown only applies to a
// quarantine that starts after a restart.
if (!Objects.equals(old.quarantineCooldownSeconds(), fresh.quarantineCooldownSeconds())) {
changed.add("quarantineCooldownSeconds");
}
Map<String, BridgedConfig.Profile> before =
old.profiles() == null ? Map.of() : old.profiles();
Map<String, BridgedConfig.Profile> after =
@@ -250,8 +261,10 @@ public final class ConfigRef implements Supplier<BridgedConfig> {
/**
* Whether two versions of a profile would launch a peer identically. Compares every component
* the launcher reads at spawn; {@code weight} and {@code maxLoad} are excluded because those are
* read live by the placement policy and really do take effect on the next spawn.
* the launcher reads at spawn; {@code weight}, {@code maxLoad} and {@code credentialId} are
* excluded because those are read live (by the placement policy and, for credentialId, by
* {@code CompositePeerLauncher}/the CB-578 stage B exhaustion sink) and really do take effect on
* the next spawn.
*/
private static boolean sameLaunchSettings(BridgedConfig.Profile a, BridgedConfig.Profile b) {
return Objects.equals(a.baseUrl(), b.baseUrl())
@@ -269,6 +282,10 @@ public final class ConfigRef implements Supplier<BridgedConfig> {
&& Objects.equals(a.gitHostEnv(), b.gitHostEnv())
&& Objects.equals(a.kind(), b.kind())
&& Objects.equals(a.env(), b.env())
&& Objects.equals(a.subscription(), b.subscription());
&& Objects.equals(a.subscription(), b.subscription())
// CB-578 stage B: exhaustedPattern is compiled once into Bridged.main's pattern map
// at startup (see ExhaustedPatternLookup wiring) — a reload never re-reads it, so a
// changed pattern must be reported as deferred, exactly like model/baseUrl/argv.
&& Objects.equals(a.exhaustedPattern(), b.exhaustedPattern());
}
}
@@ -67,6 +67,7 @@ public final class CompletionResolver implements TurnListener {
private final AgentControl agents;
private final Rendezvous rendezvous;
private final ExhaustedPatternLookup exhaustedPatterns;
private final ExhaustionSink exhaustionSink;
/**
* Per-target record of the turn currently in flight: the exact {@link Rendezvous} waiter its
@@ -89,11 +90,17 @@ public final class CompletionResolver implements TurnListener {
* usage-limit refusal pattern. Required — there is deliberately no
* defaulting overload; a caller that does not want the classification
* must pass an explicit inert value ({@link ExhaustedPatternLookup#none()}).
* @param exhaustionSink CB-578 stage B: notified when a {@code BACKEND_EXHAUSTED}
* classification actually resolves a waiter. Required for the same
* reason as {@code exhaustedPatterns} — pass {@link ExhaustionSink#none()}
* to opt out.
*/
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns) {
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
ExhaustionSink exhaustionSink) {
this.agents = agents;
this.rendezvous = rendezvous;
this.exhaustedPatterns = Objects.requireNonNull(exhaustedPatterns, "exhaustedPatterns");
this.exhaustionSink = Objects.requireNonNull(exhaustionSink, "exhaustionSink");
}
@Override
@@ -210,6 +217,9 @@ public final class CompletionResolver implements TurnListener {
inFlight.remove(target, turn);
log.warn("completion for {} classified BACKEND_EXHAUSTED (no bridge_reply; scrape "
+ "matched the profile's exhausted pattern): {}", target, reason);
// CB-578 stage B: only on the resolution that actually won the race — a late
// duplicate must never quarantine a credential twice for one refusal.
exhaustionSink.onExhausted(target, reason);
}
return;
}
@@ -0,0 +1,31 @@
package dev.ltms.bridged.inject;
/**
* Notified when {@link CompletionResolver} actually delivers a {@code BACKEND_EXHAUSTED}
* classification to a waiting send (CB-578 stage B) — never on a race that lost (see
* {@link CompletionResolver#resolve}, which only calls this after
* {@code Rendezvous.resolveExhausted} returns {@code true}).
*
* <p>{@link CompletionResolver} knows only {@code target} (a herdr terminal id); it has no notion of
* profiles or credentials, so mapping {@code target} to whatever should be quarantined is entirely
* the sink's job — see {@code Bridged.main}'s wiring, which resolves target → session → profile →
* {@code effectiveCredentialId()} and calls {@code BackendQuarantine.quarantine} on it.
*/
@FunctionalInterface
public interface ExhaustionSink {
/**
* @param target the herdr terminal id whose turn was classified {@code BACKEND_EXHAUSTED}
* @param reason the matched-line reason carried by the classification
*/
void onExhausted(String target, String reason);
/**
* Inert sink — nothing happens on exhaustion. The explicit stand-in a caller (or a test not
* exercising this feature) passes instead of a defaulting overload, exactly like
* {@link ExhaustedPatternLookup#none()}.
*/
static ExhaustionSink none() {
return (target, reason) -> { };
}
}
@@ -14,6 +14,7 @@ import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.peer.PeerUnreachableException;
import dev.ltms.bridged.placement.BackendQuarantine;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.bridged.session.WorktreeRequest;
@@ -33,6 +34,7 @@ import jakarta.servlet.http.HttpServlet;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.function.Function;
import java.util.function.LongSupplier;
@@ -78,6 +80,7 @@ public final class BridgeMcp {
private final Metrics metrics; // CB-502: null → auth failures not counted
private final CapacitySource capacity;
private final HealthCoverageSource healthCoverage;
private final QuarantineSource quarantine;
/** Capacity facts used by {@code bridge_list}; production must supply the placement live count. */
public record CapacitySource(Function<String, Integer> liveCount, Function<String, Integer> maxLoad,
@@ -90,17 +93,30 @@ public final class BridgeMcp {
/** Coverage is supplied by the health wiring, not inferred from a missing dependency. */
public record HealthCoverageSource(Supplier<String> value) { }
/**
* CB-578 stage B quarantine facts used by {@code bridge_profiles}: a profile → credential id
* lookup, plus the shared {@link BackendQuarantine} to read remaining cooldowns off.
*/
public record QuarantineSource(Function<String, String> credentialIdFor, BackendQuarantine quarantine) {
/** Inert source — no profile is ever reported quarantined. Explicit stand-in, not a default. */
public static QuarantineSource none() { return new QuarantineSource(_ -> null, BackendQuarantine.none()); }
}
/**
* @param callers resolves each call's {@link Principal}; {@code null} disables authorization.
* This surface needs its own enforcement: {@code /mcp} is a raw servlet on
* Jetty's context handler and never passes through Javalin's {@code before}
* filter, so the REST guard does not cover it.
* @param metrics registry for auth-failure counting; may be {@code null}
* @param quarantine CB-578 stage B facts for {@code bridge_profiles}; required — pass
* {@link QuarantineSource#none()} for a caller that does not want the feature
*/
public BridgeMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage) {
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine) {
this.capacity = capacity;
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
this.healthCoverage = healthCoverage;
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
this.transport = HttpServletStreamableServerTransportProvider.builder()
@@ -228,7 +244,7 @@ public final class BridgeMcp {
.toolCall(profilesTool(), (exchange, _) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied;
return profiles(workers);
return profiles(workers, quarantine);
})
.toolCall(whoamiTool(), (exchange, _) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
@@ -706,11 +722,33 @@ public final class BridgeMcp {
return null;
}
/** {@code bridge_profiles}: the configured worker profiles and the default. */
static McpSchema.CallToolResult profiles(PeerLauncher workers) {
return text(json(Map.of(
"profiles", workers.profiles(),
"default", workers.defaultProfile() == null ? "" : workers.defaultProfile())));
/**
* {@code bridge_profiles}: the configured worker profiles, the default, and — CB-578 stage B —
* which of them are currently quarantined (backend exhausted) and for how much longer. The
* {@code quarantined} key is present only when at least one profile is, so a fleet where nothing
* has ever been quarantined gets exactly the pre-stage-B shape.
*/
static McpSchema.CallToolResult profiles(PeerLauncher workers, QuarantineSource quarantine) {
Map<String, Object> result = new LinkedHashMap<>();
result.put("profiles", workers.profiles());
result.put("default", workers.defaultProfile() == null ? "" : workers.defaultProfile());
Map<String, Object> quarantined = new LinkedHashMap<>();
for (String profile : workers.profiles()) {
String credentialId = quarantine.credentialIdFor().apply(profile);
if (credentialId == null) {
continue;
}
quarantine.quarantine().remainingSeconds(credentialId).ifPresent(remaining -> {
Map<String, Object> row = new LinkedHashMap<>();
row.put("credentialId", credentialId);
row.put("quarantinedForSeconds", remaining);
quarantined.put(profile, row);
});
}
if (!quarantined.isEmpty()) {
result.put("quarantined", quarantined);
}
return text(json(result));
}
/**
@@ -940,7 +978,10 @@ public final class BridgeMcp {
private static McpSchema.Tool profilesTool() {
return tool("bridge_profiles",
"List the configured worker profiles (backends) and which one bridge_spawn uses by default.",
"List the configured worker profiles (backends) and which one bridge_spawn uses by "
+ "default. A 'quarantined' map is present when a backend-exhausted refusal put "
+ "a profile's credential on cooldown — bridge_spawn onto it is refused until "
+ "quarantinedForSeconds elapses; a profile sharing that credential is listed too.",
objectSchema(Map.of(), List.of()));
}
@@ -8,6 +8,7 @@ import dev.ltms.bridged.peer.PeerHandle;
import dev.ltms.bridged.peer.PeerLauncher;
import dev.ltms.bridged.peer.PeerUnreachableException;
import dev.ltms.bridged.peer.SpawnRequest;
import dev.ltms.bridged.placement.BackendQuarantine;
import dev.ltms.bridged.placement.PlacementCandidate;
import dev.ltms.bridged.placement.PlacementContext;
import dev.ltms.bridged.placement.PlacementException;
@@ -23,10 +24,12 @@ import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.Function;
import java.util.function.Supplier;
import java.util.stream.Collectors;
/**
* The {@link PeerLauncher} the core actually holds when more than one adapter is configured — a thin
@@ -78,6 +81,9 @@ public final class CompositePeerLauncher implements PeerLauncher {
private final Supplier<Map<String, BridgedConfig.Profile>> profileConfigs;
private final Supplier<PlacementPolicy> placementPolicy;
/** CB-578 stage B: credential cooldown, checked before an explicit spawn and filtered into placement. */
private final BackendQuarantine quarantine;
/**
* CB-557: the role pools an unqualified spawn draws its candidates from. A supplier that yields
* {@code null}, and an empty pool for a role, both fall back to every configured profile — the
@@ -99,7 +105,10 @@ public final class CompositePeerLauncher implements PeerLauncher {
}
/**
* Production constructor with a placement policy and live-worker counter.
* Production constructor with a placement policy and live-worker counter. Quarantine (CB-578
* stage B) is off for this constructor — {@link BackendQuarantine#none()} — since it predates
* the feature and existing callers of this exact overload never exercised it; use the 7-arg
* overload below to wire a real {@link BackendQuarantine}.
*
* @param delegates one adapter per configured peer kind; must be non-empty and declare
* disjoint profile-name sets
@@ -113,13 +122,14 @@ public final class CompositePeerLauncher implements PeerLauncher {
Map<String, BridgedConfig.Profile> profileConfigs,
PlacementPolicy placementPolicy,
Function<String, Integer> liveCount) {
this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, null);
this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, null, BackendQuarantine.none());
}
/**
* Production constructor with role pools (CB-557). An unqualified spawn draws its candidates from
* {@code fleet.<role>} instead of from every configured profile, so a reviewer is placed on a
* reviewer backend and never on, say, the architect-only one.
* reviewer backend and never on, say, the architect-only one. Quarantine is off for this
* constructor too, for the same reason as the 5-arg overload above.
*
* @param fleet the configured role pools; {@code null} ⇒ every profile is a candidate for every
* role, which is the pre-CB-557 behaviour
@@ -130,29 +140,50 @@ public final class CompositePeerLauncher implements PeerLauncher {
PlacementPolicy placementPolicy,
Function<String, Integer> liveCount,
BridgedConfig.Fleet fleet) {
this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, fleet, BackendQuarantine.none());
}
/**
* Production constructor with role pools and quarantine (CB-578 stage B). The full-featured
* non-reloading form; {@link #CompositePeerLauncher(List, String, Supplier, Function, BackendQuarantine)}
* is what {@code Bridged.main} actually wires up.
*
* @param quarantine required — pass {@link BackendQuarantine#none()} for a caller that does not
* want the feature, never a defaulting overload (CB-578 stage B's own rule).
*/
public CompositePeerLauncher(List<HerdrPeerLauncher> delegates,
String defaultProfile,
Map<String, BridgedConfig.Profile> profileConfigs,
PlacementPolicy placementPolicy,
Function<String, Integer> liveCount,
BridgedConfig.Fleet fleet,
BackendQuarantine quarantine) {
// LinkedHashMap, not Map.copyOf: candidates() promises definition order and the weighted
// policy breaks exact-weight ties on it, so a salted iteration order would make placement
// differ from one JVM run to the next.
this(delegates, defaultProfile,
constant(Collections.unmodifiableMap(new LinkedHashMap<>(profileConfigs))),
constant(placementPolicy), liveCount, constant(fleet));
constant(placementPolicy), liveCount, constant(fleet), quarantine);
}
/**
* Production constructor that re-reads its placement inputs per spawn (CB-559), so a config
* reload retargets the next member without a restart.
*
* @param config the live configuration — read at every spawn, never captured
* @param config the live configuration — read at every spawn, never captured
* @param quarantine required — CB-578 stage B; pass {@link BackendQuarantine#none()} to opt out
*/
public CompositePeerLauncher(List<HerdrPeerLauncher> delegates,
String defaultProfile,
Supplier<BridgedConfig> config,
Function<String, Integer> liveCount) {
Function<String, Integer> liveCount,
BackendQuarantine quarantine) {
this(delegates, defaultProfile,
() -> config.get().profiles(),
() -> PlacementPolicies.fromName(config.get().placement()),
liveCount,
() -> config.get().fleet());
() -> config.get().fleet(),
quarantine);
}
/** The all-suppliers form every other constructor funnels into. */
@@ -161,8 +192,10 @@ public final class CompositePeerLauncher implements PeerLauncher {
Supplier<Map<String, BridgedConfig.Profile>> profileConfigs,
Supplier<PlacementPolicy> placementPolicy,
Function<String, Integer> liveCount,
Supplier<BridgedConfig.Fleet> fleet) {
Supplier<BridgedConfig.Fleet> fleet,
BackendQuarantine quarantine) {
this.fleet = fleet;
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
if (delegates.isEmpty()) {
throw new IllegalArgumentException("at least one peer adapter must be configured");
}
@@ -225,6 +258,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
// the charter makes explicit-profile spawns the normal path — so skipping the check
// here would leave the cap dead config in real operation.
HerdrPeerLauncher d = route(requestedProfile);
enforceNotQuarantined(requestedProfile);
enforceMaxLoad(requestedProfile);
PeerHandle handle = d.spawn(req);
spawnedBy.put(handle.id(), d);
@@ -238,7 +272,10 @@ public final class CompositePeerLauncher implements PeerLauncher {
List<PlacementCandidate> candidates = candidates(req.role());
String roleDefault = defaultProfileFor(req.role());
Set<String> unreachable = new HashSet<>();
PlacementContext ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable);
// CB-578 stage B: computed once up front — a quarantine's expiry cannot pass within one spawn
// call, so re-deriving it per retry would only cost work, never change the answer.
Set<String> quarantined = quarantinedProfiles(candidates);
PlacementContext ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable, quarantined);
int maxAttempts = candidates.isEmpty() ? 1 : candidates.size();
for (int attempt = 0; attempt < maxAttempts; attempt++) {
@@ -268,7 +305,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
chosen.profile(), e.getMessage());
unreachable.add(chosen.profile());
// Update the context for the next selection so the policy excludes this profile.
ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable);
ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable, quarantined);
}
}
@@ -298,6 +335,42 @@ public final class CompositePeerLauncher implements PeerLauncher {
* @param profile the profile the caller explicitly named
* @throws PlacementException when the profile is at capacity
*/
/**
* Refuse an explicit-profile spawn whose credential is quarantined (CB-578 stage B): a prior
* {@code BACKEND_EXHAUSTED} classification on this profile, or on another profile sharing its
* {@code credentialId}, is still on cooldown.
*
* <p>Checked before {@link #enforceMaxLoad}, and for the same reason that check exists: an
* explicit profile bypasses the placement policy's filtering entirely, so without this the cap
* (there, quarantine here) would be dead config on the very path the charter calls normal.
* Deliberately no fallback to another profile, matching {@link #enforceMaxLoad}'s own reasoning —
* the caller named this profile for a cost/model reason.
*
* @throws PlacementException naming the profile, its credential, and the remaining cooldown
*/
private void enforceNotQuarantined(String profile) {
String credentialId = credentialIdFor(profile);
quarantine.remainingSeconds(credentialId).ifPresent(remaining -> {
throw new PlacementException("worker profile '" + profile + "' is quarantined "
+ "(credential '" + credentialId + "' exhausted; ~" + remaining
+ "s remaining) — refusing spawn");
});
}
/** {@code profile}'s credential group (CB-578 stage B), or the profile's own name if unconfigured. */
private String credentialIdFor(String profile) {
BridgedConfig.Profile cfg = profiles0().get(profile);
return cfg == null ? profile : cfg.effectiveCredentialId();
}
/** The subset of {@code candidates} whose credential is currently quarantined (CB-578 stage B). */
private Set<String> quarantinedProfiles(List<PlacementCandidate> candidates) {
return candidates.stream()
.map(PlacementCandidate::profile)
.filter(p -> quarantine.isQuarantined(credentialIdFor(p)))
.collect(Collectors.toSet());
}
private void enforceMaxLoad(String profile) {
// Absent config, or a config whose maxLoad normalized to null (non-positive ⇒ unlimited at
// load), means no cap — never cap what wasn't configured.
@@ -0,0 +1,99 @@
package dev.ltms.bridged.placement;
import java.util.LinkedHashMap;
import java.util.Map;
import java.util.Objects;
import java.util.OptionalLong;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.LongSupplier;
/**
* Where a credential (not a profile — see {@code BridgedConfig.Profile#effectiveCredentialId()})
* sits out a cooldown after a {@code BACKEND_EXHAUSTED} classification (CB-578 stage B), so a fresh
* spawn does not walk straight back onto the account that just refused on a usage limit.
*
* <p>Keyed by credential id, never by profile name: two profiles sharing one credential (e.g. two
* models on the same OpenAI account) share one quarantine — {@link #quarantine} one credential id
* and every profile whose {@code effectiveCredentialId()} equals it is quarantined too, without this
* class knowing anything about profiles at all. That mapping is the caller's job (see
* {@code CompositePeerLauncher} and {@code dev.ltms.bridged.inject.ExhaustionSink}).
*
* <p>The clock is injected ({@link LongSupplier}, conventionally {@code System::nanoTime} like
* {@code FleetHealthMonitor}), never read inline, so a quarantine's expiry is testable without a
* real sleep.
*/
public final class BackendQuarantine {
private final ConcurrentHashMap<String, Long> quarantinedUntilNanos = new ConcurrentHashMap<>();
private final LongSupplier nowNanos;
private final long cooldownNanos;
/**
* @param nowNanos monotonic clock, injected for testability
* @param cooldownNanos how long a fresh {@link #quarantine} call blocks the credential for;
* must be positive
*/
public BackendQuarantine(LongSupplier nowNanos, long cooldownNanos) {
this.nowNanos = Objects.requireNonNull(nowNanos, "nowNanos");
if (cooldownNanos <= 0) {
throw new IllegalArgumentException("cooldownNanos must be positive: " + cooldownNanos);
}
this.cooldownNanos = cooldownNanos;
}
/**
* Inert quarantine — nothing is ever quarantined unless {@link #quarantine} is actually called on
* this instance. The explicit stand-in a caller (or a test not exercising this feature) passes
* instead of a defaulting overload, exactly like {@code ExhaustedPatternLookup.none()}.
*/
public static BackendQuarantine none() {
return new BackendQuarantine(() -> 0L, 1);
}
/**
* Quarantine {@code credentialId} for the configured cooldown, starting now. A repeat call while
* already quarantined restarts the cooldown at full length — a fresh refusal is fresh evidence the
* account is still exhausted, not a reason to let an earlier, shorter wait stand.
*/
public void quarantine(String credentialId) {
Objects.requireNonNull(credentialId, "credentialId");
quarantinedUntilNanos.put(credentialId, nowNanos.getAsLong() + cooldownNanos);
}
/** Whether {@code credentialId} is quarantined right now. */
public boolean isQuarantined(String credentialId) {
return remainingNanos(credentialId) > 0;
}
/** Seconds left on {@code credentialId}'s quarantine, or empty when it is not quarantined. */
public OptionalLong remainingSeconds(String credentialId) {
long remaining = remainingNanos(credentialId);
return remaining > 0 ? OptionalLong.of(toSecondsRoundedUp(remaining)) : OptionalLong.empty();
}
/**
* Every currently-quarantined credential id and its remaining seconds (CB-578 stage B fleet
* reporting) — expired entries are never included. Not pruned from the backing map here: it stays
* small (bounded by the number of distinct credentials ever exhausted) and a lazily-stale entry is
* harmless, since every read already checks the deadline.
*/
public Map<String, Long> activeRemainingSeconds() {
Map<String, Long> out = new LinkedHashMap<>();
quarantinedUntilNanos.forEach((credentialId, deadline) -> {
long remaining = deadline - nowNanos.getAsLong();
if (remaining > 0) {
out.put(credentialId, toSecondsRoundedUp(remaining));
}
});
return out;
}
private long remainingNanos(String credentialId) {
Long deadline = quarantinedUntilNanos.get(credentialId);
return deadline == null ? 0L : deadline - nowNanos.getAsLong();
}
private static long toSecondsRoundedUp(long nanos) {
return (nanos + 999_999_999L) / 1_000_000_000L;
}
}
@@ -4,18 +4,28 @@ package dev.ltms.bridged.placement;
* Backward-compatible placement: an unqualified spawn always resolves to the configured default
* profile, exactly as {@code CompositePeerLauncher} did before CB-518. This ignores caps and
* reachability so that a pre-existing config behaves identically after upgrade.
*
* <p>Quarantine (CB-578 stage B) is the one exception: a quarantined default is a credential that
* just refused on a usage limit, not a transient capacity or reachability concern, so {@code fixed}
* steps to the first non-quarantined candidate instead of walking straight back onto it. A fleet
* where nothing is ever quarantined never exercises this path, so today's behaviour is unchanged.
*/
final class FixedPlacementPolicy implements PlacementPolicy {
@Override
public PlacementCandidate select(PlacementContext ctx) {
String d = ctx.defaultProfile();
if (d != null && !d.isBlank()) {
if (d != null && !d.isBlank() && !ctx.quarantined().contains(d)) {
return new PlacementCandidate(d, null, 1.0f, null);
}
if (!ctx.candidates().isEmpty()) {
PlacementCandidate first = ctx.candidates().getFirst();
return new PlacementCandidate(first.profile(), null, first.weight(), first.maxLoad());
for (PlacementCandidate c : ctx.candidates()) {
if (!ctx.quarantined().contains(c.profile())) {
return new PlacementCandidate(c.profile(), null, c.weight(), c.maxLoad());
}
}
if (d != null && !d.isBlank()) {
throw new PlacementException("worker profile '" + d + "' is quarantined (backend "
+ "exhausted) and no un-quarantined candidate is available");
}
throw new PlacementException("no worker profiles configured");
}
@@ -11,9 +11,13 @@ import java.util.function.Function;
* @param candidates every configured candidate; the policy filters out those at cap or unreachable
* @param liveCount current live worker count per profile (from the session registry)
* @param unreachable profiles already known to have failed in this spawn attempt
* @param quarantined profiles whose credential is currently quarantined (CB-578 stage B) — a
* {@code BACKEND_EXHAUSTED} classification put it, or a profile it shares a
* credential with, on cooldown. Filtered the same way as {@code unreachable}.
*/
public record PlacementContext(String defaultProfile,
List<PlacementCandidate> candidates,
Function<String, Integer> liveCount,
Set<String> unreachable) {
Set<String> unreachable,
Set<String> quarantined) {
}
@@ -12,13 +12,13 @@ final class PlacementPolicyUtil {
}
/**
* Candidates that are not known-unreachable and have not reached their maxLoad.
* A {@code null} maxLoad means unlimited.
* Candidates that are not known-unreachable, not quarantined (CB-578 stage B), and have not
* reached their maxLoad. A {@code null} maxLoad means unlimited.
*/
static List<PlacementCandidate> available(PlacementContext ctx) {
List<PlacementCandidate> out = new ArrayList<>();
for (PlacementCandidate c : ctx.candidates()) {
if (ctx.unreachable().contains(c.profile())) {
if (ctx.unreachable().contains(c.profile()) || ctx.quarantined().contains(c.profile())) {
continue;
}
Integer cap = c.maxLoad();
@@ -35,14 +35,17 @@ final class PlacementPolicyUtil {
/**
* Build a clear exception describing why every candidate was dropped: all at capacity,
* all unreachable, or a mix.
* all unreachable, all quarantined, or a mix.
*/
static PlacementException emptyException(PlacementContext ctx) {
int atCap = 0;
int unreachable = 0;
int quarantined = 0;
for (PlacementCandidate c : ctx.candidates()) {
Integer cap = c.maxLoad();
if (ctx.unreachable().contains(c.profile())) {
if (ctx.quarantined().contains(c.profile())) {
quarantined++;
} else if (ctx.unreachable().contains(c.profile())) {
unreachable++;
} else if (cap != null && ctx.liveCount().apply(c.profile()) >= cap) {
atCap++;
@@ -53,6 +56,9 @@ final class PlacementPolicyUtil {
if (total == 0) {
return new PlacementException("no worker profiles configured");
}
if (quarantined == total) {
return new PlacementException("all worker profiles are quarantined (backend exhausted)");
}
if (atCap == total) {
return new PlacementException("all worker profiles are at maxLoad");
}
@@ -60,6 +66,7 @@ final class PlacementPolicyUtil {
return new PlacementException("all worker profiles are unreachable");
}
return new PlacementException("no worker profile available: " + atCap + " at maxLoad, "
+ unreachable + " unreachable, " + (total - atCap - unreachable) + " remaining");
+ unreachable + " unreachable, " + quarantined + " quarantined, "
+ (total - atCap - unreachable - quarantined) + " remaining");
}
}
@@ -338,6 +338,54 @@ class ConfigRefTest {
assertEquals("deepseek-v4-flash", ref.get().profiles().get("sonnet").model());
}
/**
* CB-578 stage B: exhaustedPattern is compiled once into Bridged.main's pattern map at startup
* (see ExhaustedPatternLookup), so a reload never re-reads it — a changed pattern must be
* reported deferred exactly like model/baseUrl, not silently claimed as applied.
*/
@Test
void changingAProfilesExhaustedPatternIsReportedAsDeferred(@TempDir Path dir) throws Exception {
Path f = dir.resolve("bridged.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
sonnet:
baseUrl: http://gx00.gw:8000
model: sonnet
exhaustedPattern: "usage limit has been reached"
guard:
offSubscriptionHosts:
- gx00.gw
""");
ConfigRef ref = refFor(f);
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
sonnet:
baseUrl: http://gx00.gw:8000
model: sonnet
exhaustedPattern: "rate limit exceeded"
guard:
offSubscriptionHosts:
- gx00.gw
""");
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertEquals(1, out.deferred().size(), out.deferred().toString());
assertTrue(out.deferred().getFirst().contains("sonnet"), out.deferred().toString());
assertTrue(out.deferred().getFirst().contains("launch settings"), out.deferred().toString());
// The snapshot still carries the new value — a restart is what makes it take effect.
assertEquals("rate limit exceeded", ref.get().profiles().get("sonnet").exhaustedPattern());
}
@Test
void aFixedRefHasNoFileAndRefusesToReload() {
BridgedConfig cfg = new BridgedConfig(null, null, null, null, null, null,
@@ -26,7 +26,7 @@ class CompletionResolverTest {
void skipsTheScrapeWhenNoSendIsWaiting() {
FakeHerdr herdr = new FakeHerdr();
Rendezvous rendezvous = new Rendezvous(); // no waiter opened
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
resolver.resolve("term_a", null); // no in-flight turn captured for this target
@@ -38,7 +38,7 @@ class CompletionResolverTest {
void failSkipsTheScrapeWhenNoSendIsWaiting() {
FakeHerdr herdr = new FakeHerdr();
Rendezvous rendezvous = new Rendezvous(); // no waiter opened
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
resolver.fail("term_a", null); // no in-flight turn, and no registered waiter to fall back to
@@ -50,7 +50,7 @@ class CompletionResolverTest {
void captureBaselineSkipsTheReadWhenNoSendIsWaiting() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ X\n❯ ");
Rendezvous rendezvous = new Rendezvous(); // no waiter opened
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
resolver.captureBaseline("term_a", TestTurnTokens.inert("term_a")); // no send to attribute a later completion to
@@ -139,7 +139,7 @@ class CompletionResolverTest {
// send must NOT be resolved with the stale answer.
FakeHerdr herdr = new FakeHerdr().readText("⏺ 391\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiter = rendezvous.open("term_a"); // a send is blocked on this turn
// The turn as captured at delivery: its waiter, and the previous turn's answer still on screen.
@@ -154,7 +154,7 @@ class CompletionResolverTest {
void resolvesACompletionWhoseScrapeChangedSinceDelivery() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ No, 391 = 17 × 23.\n❯ "); // the worker's real answer
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiter = rendezvous.open("term_a");
// Delivery baseline was the previous turn's "391"; the scrape now differs → resolve.
@@ -171,7 +171,7 @@ class CompletionResolverTest {
String block = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 1) + "\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
@@ -185,7 +185,7 @@ class CompletionResolverTest {
void leavesAnUnclippedCompletionPaneTailUnmarked() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ complete report\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
@@ -197,7 +197,7 @@ class CompletionResolverTest {
void resolvesSynchronouslyBeforePostTurnContextClearing() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ previous answer\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiter = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter));
herdr.readText("⏺ answer that /clear would erase\n❯ ");
@@ -219,7 +219,7 @@ class CompletionResolverTest {
String longBlock = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 500) + "\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(longBlock);
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiter = rendezvous.open("term_a"); // a send is blocked on this turn
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter)); // baseline is the clipped >cap block
@@ -239,7 +239,7 @@ class CompletionResolverTest {
// No delivery baseline (e.g. the pre-turn read failed) ⇒ never suppress; the completion resolves.
FakeHerdr herdr = new FakeHerdr().readText("⏺ hello\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
@@ -257,7 +257,7 @@ class CompletionResolverTest {
// byte-identical guard would wrongly match the empty tail and suppress.
FakeHerdr herdr = new FakeHerdr().healthy(false); // agent.read throws HerdrException
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, ""); // empty pane baselined at delivery
@@ -277,7 +277,7 @@ class CompletionResolverTest {
// fail must not overwrite that value, and must not even scrape the worker — nobody needs it.
FakeHerdr herdr = new FakeHerdr().readText("an error screen");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, null);
@@ -298,7 +298,7 @@ class CompletionResolverTest {
// fail falls back to the waiter currently registered on the Rendezvous and fails it.
FakeHerdr herdr = new FakeHerdr().readText("stuck on an error screen");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiter = rendezvous.open("term_a"); // send registered, but no captureBaseline ever ran
resolver.fail("term_a", null); // no in-flight turn → fall back to the registered waiter
@@ -324,7 +324,7 @@ class CompletionResolverTest {
try {
FakeHerdr herdr = new FakeHerdr().readText("stuck on an error screen");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiter = rendezvous.open("term_a");
resolver.fail("term_a", null);
@@ -352,7 +352,7 @@ class CompletionResolverTest {
// scrape to turn N+1; targeting turn N's captured waiter makes the late completion a no-op.
FakeHerdr herdr = new FakeHerdr().readText("⏺ turn N answer\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiterN = rendezvous.open("term_a"); // turn N's send
// The turn as the injector captured it at delivery (waiter + pre-turn baseline).
@@ -382,7 +382,7 @@ class CompletionResolverTest {
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, ExhaustionSink.none());
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
@@ -398,7 +398,7 @@ class CompletionResolverTest {
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, ExhaustionSink.none());
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
@@ -407,12 +407,53 @@ class CompletionResolverTest {
waiter.getNow(null).text(), "the reason names the real cause and carries the matched line");
}
@Test
void aWinningBackendExhaustedClassificationNotifiesTheExhaustionSink() {
String block = "⏺ Working on it...\nThe usage limit has been reached. Try again later.\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
java.util.List<String> notified = new java.util.ArrayList<>();
ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals(1, notified.size(), "the sink is notified exactly once for the winning classification");
assertTrue(notified.get(0).startsWith("term_a: "), "the sink is told which target exhausted");
assertTrue(notified.get(0).contains("The usage limit has been reached"),
"the sink is told the matched reason: " + notified.get(0));
}
@Test
void aLosingBackendExhaustedClassificationNeverNotifiesTheExhaustionSink() {
// The waiter was already resolved (e.g. by the worker's own reply) before this scrape landed —
// resolveExhausted loses the race and must return false, so the sink must not fire either.
String block = "⏺ Working on it...\nThe usage limit has been reached. Try again later.\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
java.util.List<String> notified = new java.util.ArrayList<>();
ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
var waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, null);
assertTrue(rendezvous.resolveCompletion(waiter, "already replied"));
resolver.resolve("term_a", turn);
assertTrue(notified.isEmpty(), "a classification that loses the race must not quarantine anything");
assertEquals("already replied", waiter.getNow(null).text(), "the earlier resolution stands untouched");
}
@Test
void aNonMatchingScrapeResolvesAsAnOrdinaryCompletion() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ complete report\n❯ ");
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, ExhaustionSink.none());
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
@@ -430,7 +471,7 @@ class CompletionResolverTest {
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver =
new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
@@ -72,7 +72,8 @@ class BridgeMcpAuthzTest {
new PrimaryRegistry(null),
enforce ? CallerResolver.withLeadsAndMembers(identity, false, null,
Map::of, new MemberRegistry(null)) : null,
metrics, BridgeMcp.CapacitySource.none(), new BridgeMcp.HealthCoverageSource(() -> "off"));
metrics, BridgeMcp.CapacitySource.none(), new BridgeMcp.HealthCoverageSource(() -> "off"),
BridgeMcp.QuarantineSource.none());
return mcp;
}
@@ -16,6 +16,7 @@ import dev.ltms.bridged.inject.MemberPresence;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.bridged.session.WorktreeRequest;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
import dev.ltms.bridged.placement.BackendQuarantine;
import io.modelcontextprotocol.spec.McpSchema;
import dev.ltms.bridged.msg.InMemoryReplyInbox;
import org.junit.jupiter.api.BeforeEach;
@@ -437,11 +438,29 @@ class BridgeMcpTest {
@Test
void profilesListsConfiguredProfilesAndDefault() {
FakeHerdr h = new FakeHerdr();
McpSchema.CallToolResult res = BridgeMcp.profiles(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
McpSchema.CallToolResult res = BridgeMcp.profiles(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), BridgeMcp.QuarantineSource.none());
assertNotEquals(Boolean.TRUE, res.isError());
String out = textOf(res);
assertTrue(out.contains("ltms-local"), out);
assertTrue(out.contains("\"default\":\"ltms-local\""), out);
assertFalse(out.contains("quarantined"), "no profile is quarantined, so the key is omitted: " + out);
}
@Test
void profilesReportsAQuarantinedCredential() {
FakeHerdr h = new FakeHerdr();
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
quarantine.quarantine("shared-openai");
BridgeMcp.QuarantineSource source = new BridgeMcp.QuarantineSource(
profile -> "ltms-local".equals(profile) ? "shared-openai" : null, quarantine);
McpSchema.CallToolResult res = BridgeMcp.profiles(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), source);
assertNotEquals(Boolean.TRUE, res.isError());
String out = textOf(res);
assertTrue(out.contains("\"quarantined\""), out);
assertTrue(out.contains("shared-openai"), out);
assertTrue(out.contains("\"quarantinedForSeconds\":1800"), out);
}
@Test
@@ -16,6 +16,7 @@ import dev.ltms.bridged.peer.PeerHandle;
import dev.ltms.bridged.peer.PeerLauncher;
import dev.ltms.bridged.peer.PeerUnreachableException;
import dev.ltms.bridged.peer.SpawnRequest;
import dev.ltms.bridged.placement.BackendQuarantine;
import dev.ltms.bridged.placement.PlacementException;
import dev.ltms.bridged.placement.PlacementPolicies;
import org.junit.jupiter.api.Test;
@@ -27,6 +28,8 @@ import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.*;
@@ -128,6 +131,13 @@ class CompositePeerLauncherTest {
weight, maxLoad);
}
private static BridgedConfig.Profile stubWorker(String profile, String credentialId) {
return new BridgedConfig.Profile(profile, "http://gx00.gw:8000", "coder",
null, "BRIDGED_WORKER_TOKEN", List.of("claude"), "tab", "bridged-workers",
"w #{n}", null, null, null, null, null, null, null, null, null,
null, null, credentialId);
}
/**
* An <em>order-preserving</em> profile map. Never {@code Map.of} here: its iteration order is
* salted per JVM run, and the weighted policy breaks an exact-weight tie on candidate order —
@@ -554,4 +564,105 @@ class CompositePeerLauncherTest {
new SpawnRequest(null, null, null, null, null, MemberRole.DEV)));
assertTrue(e.getMessage().contains("maxLoad"), e.getMessage());
}
// ── CB-578 stage B: a BACKEND_EXHAUSTED classification quarantines the credential ──────────
@Test
void explicitSpawnOntoAQuarantinedProfileIsRefusedNamingTheCredentialAndRemainingTime() {
FakeHerdr herdr = new FakeHerdr();
Map<String, BridgedConfig.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");
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
PlacementPolicies.fixed(), _ -> 0, null, quarantine);
PlacementException e = assertThrows(PlacementException.class,
() -> composite.spawn(new SpawnRequest("sol", null, null)));
assertTrue(e.getMessage().contains("sol"), "message names the profile: " + e.getMessage());
assertTrue(e.getMessage().contains("shared-openai"), "message names the credential: " + e.getMessage());
assertTrue(e.getMessage().contains("1800"), "message names roughly when it lifts: " + e.getMessage());
assertEquals(0, adapter.spawnCount("sol"), "the quarantined profile is never delegated to");
}
/**
* The part CB-578 stage B calls out as easy to get wrong: sol and terra are two different
* profiles sharing one OpenAI credential. Quarantining because of an exhaustion classified on
* ONE of them must lock out the other too, or the fleet just walks onto the same dead account
* under the sibling's name.
*/
@Test
void twoProfilesSharingACredentialAreBothQuarantinedByOneExhaustionEvent() {
FakeHerdr herdr = new FakeHerdr();
Map<String, BridgedConfig.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));
// Only "sol" was classified BACKEND_EXHAUSTED — but the two profiles share one credential.
quarantine.quarantine("shared-openai");
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
PlacementPolicies.fixed(), _ -> 0, null, quarantine);
assertThrows(PlacementException.class, () -> composite.spawn(new SpawnRequest("sol", null, null)),
"sol was the one classified exhausted");
assertThrows(PlacementException.class, () -> composite.spawn(new SpawnRequest("terra", null, null)),
"terra shares sol's credential, so it must be locked out too");
assertEquals(0, adapter.spawnCount("sol"));
assertEquals(0, adapter.spawnCount("terra"));
}
@Test
void placementSkipsAQuarantinedProfileAndRoutesToAnUnquarantinedOne() {
FakeHerdr herdr = new FakeHerdr();
Map<String, BridgedConfig.Profile> profiles = ordered(
"sol", stubWorker("sol", "shared-openai"),
"b", stubWorker("b"));
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "sol", Set.of());
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
quarantine.quarantine("shared-openai");
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
PlacementPolicies.weighted(), _ -> 0, null, quarantine);
PeerHandle h = composite.spawn(new SpawnRequest(null, null, null));
assertEquals("b", h.profile(), "sol is quarantined, so an unqualified spawn must land on b");
assertEquals(0, adapter.spawnCount("sol"));
}
@Test
void aQuarantineLiftsOnTheInjectedClockAndTheProfileBecomesSpawnableAgain() {
FakeHerdr herdr = new FakeHerdr();
Map<String, BridgedConfig.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);
BackendQuarantine quarantine = new BackendQuarantine(nowNanos::get, TimeUnit.MINUTES.toNanos(30));
quarantine.quarantine("shared-openai");
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
PlacementPolicies.fixed(), _ -> 0, null, quarantine);
assertThrows(PlacementException.class, () -> composite.spawn(new SpawnRequest("sol", null, null)),
"still inside the cooldown");
nowNanos.set(TimeUnit.MINUTES.toNanos(31));
PeerHandle h = composite.spawn(new SpawnRequest("sol", null, null));
assertEquals("sol", h.profile(), "the cooldown expired on the injected clock — sol is spawnable again");
assertEquals(1, adapter.spawnCount("sol"));
}
@Test
void aFleetWithNoExhaustedPatternAnywhereBehavesExactlyAsBeforeQuarantineExisted() {
// BackendQuarantine.none() is the inert stand-in every constructor already defaults to when
// no quarantine is wired — the 2/5/6-arg constructors used throughout this file all exercise
// it. This test pins that an explicit .none() also never refuses a spawn, for any profile.
FakeHerdr herdr = new FakeHerdr();
PeerLauncher composite = composite(herdr);
assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("claude", null, null)));
assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("gemini", null, null)));
}
}
@@ -6,6 +6,7 @@ import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.inject.CompletionResolver;
import dev.ltms.bridged.inject.ExhaustedPatternLookup;
import dev.ltms.bridged.inject.ExhaustionSink;
import dev.ltms.bridged.mcp.PrimaryRegistry;
import dev.ltms.bridged.inject.Injector;
import org.junit.jupiter.api.BeforeEach;
@@ -33,7 +34,7 @@ class MessageServiceTest {
private final AgentControl agents = new AgentControl(herdr);
private final Rendezvous rendezvous = new Rendezvous();
private final CompletionResolver completion =
new CompletionResolver(agents, rendezvous, ExhaustedPatternLookup.none());
new CompletionResolver(agents, rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
private final Injector injector = new Injector(agents, completion);
private final InMemoryReplyInbox inbox = new InMemoryReplyInbox();
private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
@@ -0,0 +1,108 @@
package dev.ltms.bridged.placement;
import org.junit.jupiter.api.Test;
import java.util.Map;
import java.util.OptionalLong;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* CB-578 stage B: the credential-keyed quarantine tracker itself, isolated from placement/spawn
* wiring (that's {@code CompositePeerLauncherTest}). The clock is a plain {@link AtomicLong} of
* nanos so expiry is exercised without a real sleep.
*/
class BackendQuarantineTest {
@Test
void aFreshCredentialIsNotQuarantined() {
BackendQuarantine q = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
assertFalse(q.isQuarantined("shared-openai"));
assertEquals(OptionalLong.empty(), q.remainingSeconds("shared-openai"));
}
@Test
void quarantineBlocksTheCredentialForTheFullCooldown() {
BackendQuarantine q = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
q.quarantine("shared-openai");
assertTrue(q.isQuarantined("shared-openai"));
assertEquals(OptionalLong.of(1800L), q.remainingSeconds("shared-openai"));
}
@Test
void onlyTheQuarantinedCredentialIsAffected() {
BackendQuarantine q = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
q.quarantine("shared-openai");
assertFalse(q.isQuarantined("some-other-credential"),
"an unrelated credential must not be swept into the quarantine");
}
@Test
void expiresOnTheInjectedClock() {
AtomicLong now = new AtomicLong(0L);
BackendQuarantine q = new BackendQuarantine(now::get, TimeUnit.MINUTES.toNanos(30));
q.quarantine("shared-openai");
assertTrue(q.isQuarantined("shared-openai"));
now.set(TimeUnit.MINUTES.toNanos(29));
assertTrue(q.isQuarantined("shared-openai"), "still inside the cooldown");
now.set(TimeUnit.MINUTES.toNanos(31));
assertFalse(q.isQuarantined("shared-openai"), "the cooldown has elapsed on the injected clock");
assertEquals(OptionalLong.empty(), q.remainingSeconds("shared-openai"));
}
@Test
void aRepeatQuarantineCallRestartsTheCooldownAtFullLength() {
AtomicLong now = new AtomicLong(0L);
BackendQuarantine q = new BackendQuarantine(now::get, TimeUnit.MINUTES.toNanos(30));
q.quarantine("shared-openai");
now.set(TimeUnit.MINUTES.toNanos(20));
q.quarantine("shared-openai");
now.set(TimeUnit.MINUTES.toNanos(45)); // 25 min after the second call, 45 after the first
assertTrue(q.isQuarantined("shared-openai"),
"a fresh exhaustion resets the cooldown to full length, not the earlier shorter wait");
}
@Test
void activeRemainingSecondsListsOnlyStillQuarantinedCredentials() {
AtomicLong now = new AtomicLong(0L);
BackendQuarantine q = new BackendQuarantine(now::get, TimeUnit.MINUTES.toNanos(30));
q.quarantine("shared-openai");
q.quarantine("another-credential");
now.set(TimeUnit.MINUTES.toNanos(31));
q.quarantine("shared-openai"); // re-quarantined after the first one expired
Map<String, Long> active = q.activeRemainingSeconds();
assertEquals(Map.of("shared-openai", 1800L), active,
"the expired credential is dropped; the re-quarantined one is reported");
}
@Test
void noneReportsNothingQuarantinedWhenNeverToldTo() {
// .none() is the stand-in for a caller whose code path never calls #quarantine at all (e.g.
// the 2/5/6-arg CompositePeerLauncher constructors) — not a guarantee that a call to
// #quarantine on it is a no-op. Left alone, as those call sites leave it, nothing is ever
// quarantined.
BackendQuarantine q = BackendQuarantine.none();
assertFalse(q.isQuarantined("anything"));
assertTrue(q.activeRemainingSeconds().isEmpty());
}
@Test
void aNonPositiveCooldownIsRejected() {
assertThrows(IllegalArgumentException.class, () -> new BackendQuarantine(() -> 0L, 0L));
assertThrows(IllegalArgumentException.class, () -> new BackendQuarantine(() -> 0L, -1L));
}
}
@@ -24,7 +24,13 @@ class PlacementPolicyTest {
private static PlacementContext ctx(List<PlacementCandidate> candidates,
Function<String, Integer> liveCount,
Set<String> unreachable) {
return new PlacementContext("b", candidates, liveCount, unreachable);
return ctx(candidates, liveCount, unreachable, Set.of());
}
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);
}
private static PlacementContext ctx(List<PlacementCandidate> candidates,
@@ -46,17 +52,37 @@ class PlacementPolicyTest {
PlacementPolicy policy = PlacementPolicies.fixed();
PlacementContext ctx = new PlacementContext(null,
List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")),
noSessions(), Set.of());
noSessions(), Set.of(), Set.of());
assertEquals("a", policy.select(ctx).profile());
}
@Test
void fixedThrowsWhenNoProfilesAndNoDefault() {
PlacementPolicy policy = PlacementPolicies.fixed();
PlacementContext ctx = new PlacementContext(null, List.of(), noSessions(), Set.of());
PlacementContext ctx = new PlacementContext(null, List.of(), noSessions(), Set.of(), Set.of());
assertThrows(PlacementException.class, () -> policy.select(ctx));
}
@Test
void fixedSkipsQuarantinedDefault() {
PlacementPolicy policy = PlacementPolicies.fixed();
PlacementContext ctx = new PlacementContext("b",
List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")),
noSessions(), Set.of(), Set.of("b"));
assertEquals("a", policy.select(ctx).profile(),
"the default 'b' is quarantined, so fixed falls through to the first un-quarantined candidate");
}
@Test
void fixedThrowsWhenDefaultAndEveryCandidateQuarantined() {
PlacementPolicy policy = PlacementPolicies.fixed();
PlacementContext ctx = new PlacementContext("b",
List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")),
noSessions(), Set.of(), Set.of("a", "b"));
PlacementException e = assertThrows(PlacementException.class, () -> policy.select(ctx));
assertTrue(e.getMessage().contains("quarantined"), e.getMessage());
}
@Test
void roundRobinCyclesThroughAvailableProfiles() {
PlacementPolicy policy = PlacementPolicies.roundRobin();
@@ -82,6 +108,18 @@ class PlacementPolicyTest {
}
}
@Test
void roundRobinSkipsQuarantinedProfiles() {
PlacementPolicy policy = PlacementPolicies.roundRobin();
List<PlacementCandidate> candidates = List.of(
PlacementCandidate.profile("a"),
PlacementCandidate.profile("b"));
PlacementContext ctx = ctx(candidates, noSessions(), Set.of(), Set.of("a"));
for (int i = 0; i < 5; i++) {
assertEquals("b", policy.select(ctx).profile(), "a is quarantined, so every pick lands on b");
}
}
@Test
void roundRobinThrowsWhenAllAtMaxLoad() {
PlacementPolicy policy = PlacementPolicies.roundRobin();
@@ -150,6 +188,29 @@ class PlacementPolicyTest {
assertTrue(e.getMessage().contains("maxLoad"), e.getMessage());
}
@Test
void weightedSkipsQuarantinedProfile() {
PlacementPolicy policy = PlacementPolicies.weighted();
List<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("a"));
for (int i = 0; i < 5; i++) {
assertEquals("b", policy.select(ctx).profile(), "a is quarantined, so every pick lands on b");
}
}
@Test
void weightedThrowsWhenAllQuarantined() {
PlacementPolicy policy = PlacementPolicies.weighted();
List<PlacementCandidate> candidates = List.of(
PlacementCandidate.profile("a"),
PlacementCandidate.profile("b"));
PlacementException e = assertThrows(PlacementException.class,
() -> policy.select(ctx(candidates, noSessions(), Set.of(), Set.of("a", "b"))));
assertTrue(e.getMessage().contains("quarantined"), e.getMessage());
}
@Test
void weightedThrowsWhenAllUnreachable() {
PlacementPolicy policy = PlacementPolicies.weighted();