Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 8fd2d7e5e7 |
@@ -144,17 +144,6 @@ 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.
|
||||
#
|
||||
@@ -189,7 +178,6 @@ 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
|
||||
@@ -255,15 +243,6 @@ 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
|
||||
@@ -273,20 +252,17 @@ 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
|
||||
# / 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
|
||||
# `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
|
||||
# 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`, `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.
|
||||
# `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.
|
||||
# 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,7 +14,6 @@ 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;
|
||||
@@ -38,7 +37,6 @@ 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;
|
||||
@@ -46,7 +44,6 @@ 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;
|
||||
@@ -144,18 +141,11 @@ 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),
|
||||
quarantine);
|
||||
profileName -> liveCountRef.get().apply(profileName));
|
||||
// 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
|
||||
@@ -282,23 +272,7 @@ public final class Bridged {
|
||||
.orElse(null);
|
||||
log.info("backend-exhausted classification (CB-578 stage A): {}",
|
||||
CompletionResolver.coverage(cfg.profiles().keySet(), exhaustedPatternsByProfile.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
|
||||
// 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);
|
||||
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns);
|
||||
// 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();
|
||||
@@ -466,11 +440,7 @@ 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,12 +58,6 @@ 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(
|
||||
@@ -82,11 +76,7 @@ public record BridgedConfig(
|
||||
Health health,
|
||||
String placement,
|
||||
Auth auth,
|
||||
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;
|
||||
ConfigReload configReload) {
|
||||
|
||||
/** 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,
|
||||
@@ -94,7 +84,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, null);
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, null);
|
||||
}
|
||||
|
||||
/** Back-compat form before the optional {@code health:} block was added. */
|
||||
@@ -103,18 +93,7 @@ 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, 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);
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, configReload);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -208,17 +187,6 @@ 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,
|
||||
@@ -232,8 +200,7 @@ public record BridgedConfig(
|
||||
Float weight,
|
||||
Integer maxLoad,
|
||||
Boolean subscription,
|
||||
String exhaustedPattern,
|
||||
String credentialId) {
|
||||
String exhaustedPattern) {
|
||||
|
||||
/** Peer kind spawned by {@link dev.ltms.bridged.member.ClaudeCodeLauncher} (the default). */
|
||||
public static final String KIND_CLAUDE_CODE = "claude-code";
|
||||
@@ -279,10 +246,6 @@ 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;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -327,7 +290,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, credentialId);
|
||||
exhaustedPattern);
|
||||
}
|
||||
|
||||
/** True when this profile is served by the Claude Code adapter (the default kind). */
|
||||
@@ -363,38 +326,11 @@ 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();
|
||||
@@ -914,7 +850,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", "quarantineCooldownSeconds");
|
||||
"leadHeartbeat", "health", "placement", "auth", "configReload");
|
||||
|
||||
/** Load and validate config from {@code path}. */
|
||||
public static BridgedConfig load(Path path) {
|
||||
@@ -1221,15 +1157,8 @@ 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,
|
||||
quarantineCooldown);
|
||||
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -37,15 +37,10 @@ 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 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,
|
||||
* {@code spawnReadyTimeoutMs} / {@code spawnReadyPollMs}, {@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},
|
||||
* {@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 model}, {@code baseUrl}, {@code argv}, {@code env}, {@code mcpUrl} and the rest.
|
||||
* {@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>
|
||||
@@ -220,12 +215,6 @@ 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 =
|
||||
@@ -261,10 +250,8 @@ 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}, {@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.
|
||||
* 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.
|
||||
*/
|
||||
private static boolean sameLaunchSettings(BridgedConfig.Profile a, BridgedConfig.Profile b) {
|
||||
return Objects.equals(a.baseUrl(), b.baseUrl())
|
||||
@@ -282,10 +269,6 @@ 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())
|
||||
// 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());
|
||||
&& Objects.equals(a.subscription(), b.subscription());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -67,7 +67,6 @@ 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
|
||||
@@ -90,17 +89,11 @@ 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,
|
||||
ExhaustionSink exhaustionSink) {
|
||||
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns) {
|
||||
this.agents = agents;
|
||||
this.rendezvous = rendezvous;
|
||||
this.exhaustedPatterns = Objects.requireNonNull(exhaustedPatterns, "exhaustedPatterns");
|
||||
this.exhaustionSink = Objects.requireNonNull(exhaustionSink, "exhaustionSink");
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -217,9 +210,6 @@ 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;
|
||||
}
|
||||
|
||||
@@ -1,31 +0,0 @@
|
||||
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,7 +14,6 @@ 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;
|
||||
@@ -34,7 +33,6 @@ 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;
|
||||
@@ -80,7 +78,6 @@ 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,
|
||||
@@ -93,30 +90,17 @@ 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,
|
||||
QuarantineSource quarantine) {
|
||||
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage) {
|
||||
this.capacity = capacity;
|
||||
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
|
||||
this.healthCoverage = healthCoverage;
|
||||
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
|
||||
this.transport = HttpServletStreamableServerTransportProvider.builder()
|
||||
@@ -244,7 +228,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, quarantine);
|
||||
return profiles(workers);
|
||||
})
|
||||
.toolCall(whoamiTool(), (exchange, _) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
||||
@@ -722,33 +706,11 @@ public final class BridgeMcp {
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* {@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));
|
||||
/** {@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())));
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -978,10 +940,7 @@ 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. 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.",
|
||||
"List the configured worker profiles (backends) and which one bridge_spawn uses by default.",
|
||||
objectSchema(Map.of(), List.of()));
|
||||
}
|
||||
|
||||
|
||||
@@ -8,7 +8,6 @@ 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;
|
||||
@@ -24,12 +23,10 @@ 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
|
||||
@@ -81,9 +78,6 @@ 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
|
||||
@@ -105,10 +99,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
}
|
||||
|
||||
/**
|
||||
* 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}.
|
||||
* Production constructor with a placement policy and live-worker counter.
|
||||
*
|
||||
* @param delegates one adapter per configured peer kind; must be non-empty and declare
|
||||
* disjoint profile-name sets
|
||||
@@ -122,14 +113,13 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
Map<String, BridgedConfig.Profile> profileConfigs,
|
||||
PlacementPolicy placementPolicy,
|
||||
Function<String, Integer> liveCount) {
|
||||
this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, null, BackendQuarantine.none());
|
||||
this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* 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. Quarantine is off for this
|
||||
* constructor too, for the same reason as the 5-arg overload above.
|
||||
* reviewer backend and never on, say, the architect-only one.
|
||||
*
|
||||
* @param fleet the configured role pools; {@code null} ⇒ every profile is a candidate for every
|
||||
* role, which is the pre-CB-557 behaviour
|
||||
@@ -140,50 +130,29 @@ 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), quarantine);
|
||||
constant(placementPolicy), liveCount, constant(fleet));
|
||||
}
|
||||
|
||||
/**
|
||||
* 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 quarantine required — CB-578 stage B; pass {@link BackendQuarantine#none()} to opt out
|
||||
* @param config the live configuration — read at every spawn, never captured
|
||||
*/
|
||||
public CompositePeerLauncher(List<HerdrPeerLauncher> delegates,
|
||||
String defaultProfile,
|
||||
Supplier<BridgedConfig> config,
|
||||
Function<String, Integer> liveCount,
|
||||
BackendQuarantine quarantine) {
|
||||
Function<String, Integer> liveCount) {
|
||||
this(delegates, defaultProfile,
|
||||
() -> config.get().profiles(),
|
||||
() -> PlacementPolicies.fromName(config.get().placement()),
|
||||
liveCount,
|
||||
() -> config.get().fleet(),
|
||||
quarantine);
|
||||
() -> config.get().fleet());
|
||||
}
|
||||
|
||||
/** The all-suppliers form every other constructor funnels into. */
|
||||
@@ -192,10 +161,8 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
Supplier<Map<String, BridgedConfig.Profile>> profileConfigs,
|
||||
Supplier<PlacementPolicy> placementPolicy,
|
||||
Function<String, Integer> liveCount,
|
||||
Supplier<BridgedConfig.Fleet> fleet,
|
||||
BackendQuarantine quarantine) {
|
||||
Supplier<BridgedConfig.Fleet> fleet) {
|
||||
this.fleet = fleet;
|
||||
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
|
||||
if (delegates.isEmpty()) {
|
||||
throw new IllegalArgumentException("at least one peer adapter must be configured");
|
||||
}
|
||||
@@ -258,7 +225,6 @@ 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);
|
||||
@@ -272,10 +238,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
List<PlacementCandidate> candidates = candidates(req.role());
|
||||
String roleDefault = defaultProfileFor(req.role());
|
||||
Set<String> unreachable = new HashSet<>();
|
||||
// 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);
|
||||
PlacementContext ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable);
|
||||
|
||||
int maxAttempts = candidates.isEmpty() ? 1 : candidates.size();
|
||||
for (int attempt = 0; attempt < maxAttempts; attempt++) {
|
||||
@@ -305,7 +268,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, quarantined);
|
||||
ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -335,42 +298,6 @@ 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.
|
||||
|
||||
@@ -1,99 +0,0 @@
|
||||
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,28 +4,18 @@ 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() && !ctx.quarantined().contains(d)) {
|
||||
if (d != null && !d.isBlank()) {
|
||||
return new PlacementCandidate(d, null, 1.0f, null);
|
||||
}
|
||||
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");
|
||||
if (!ctx.candidates().isEmpty()) {
|
||||
PlacementCandidate first = ctx.candidates().getFirst();
|
||||
return new PlacementCandidate(first.profile(), null, first.weight(), first.maxLoad());
|
||||
}
|
||||
throw new PlacementException("no worker profiles configured");
|
||||
}
|
||||
|
||||
@@ -11,13 +11,9 @@ 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> quarantined) {
|
||||
Set<String> unreachable) {
|
||||
}
|
||||
|
||||
@@ -12,13 +12,13 @@ final class PlacementPolicyUtil {
|
||||
}
|
||||
|
||||
/**
|
||||
* Candidates that are not known-unreachable, not quarantined (CB-578 stage B), and have not
|
||||
* reached their maxLoad. A {@code null} maxLoad means unlimited.
|
||||
* Candidates that are not known-unreachable 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()) || ctx.quarantined().contains(c.profile())) {
|
||||
if (ctx.unreachable().contains(c.profile())) {
|
||||
continue;
|
||||
}
|
||||
Integer cap = c.maxLoad();
|
||||
@@ -35,17 +35,14 @@ final class PlacementPolicyUtil {
|
||||
|
||||
/**
|
||||
* Build a clear exception describing why every candidate was dropped: all at capacity,
|
||||
* all unreachable, all quarantined, or a mix.
|
||||
* all unreachable, 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.quarantined().contains(c.profile())) {
|
||||
quarantined++;
|
||||
} else if (ctx.unreachable().contains(c.profile())) {
|
||||
if (ctx.unreachable().contains(c.profile())) {
|
||||
unreachable++;
|
||||
} else if (cap != null && ctx.liveCount().apply(c.profile()) >= cap) {
|
||||
atCap++;
|
||||
@@ -56,9 +53,6 @@ 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");
|
||||
}
|
||||
@@ -66,7 +60,6 @@ final class PlacementPolicyUtil {
|
||||
return new PlacementException("all worker profiles are unreachable");
|
||||
}
|
||||
return new PlacementException("no worker profile available: " + atCap + " at maxLoad, "
|
||||
+ unreachable + " unreachable, " + quarantined + " quarantined, "
|
||||
+ (total - atCap - unreachable - quarantined) + " remaining");
|
||||
+ unreachable + " unreachable, " + (total - atCap - unreachable) + " remaining");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -197,27 +197,44 @@ public final class SessionManager implements TurnListener {
|
||||
MemberSession removed = registry.remove(paneId);
|
||||
boolean preserveWorktree = cause == ReleaseCause.SHUTDOWN;
|
||||
if (removed != null) {
|
||||
memberLifecycle.released(removed.terminalId());
|
||||
log.debug("releasing session pane={} terminal={} state={} cause={}",
|
||||
removed.paneId(), removed.terminalId(), removed.state(), cause);
|
||||
if (preserveWorktree && removed.worktree() != null) {
|
||||
logPreservedForShutdown(removed);
|
||||
} else if (removed.worktree() != null && worktrees.hasUncommitted(removed.worktree())) {
|
||||
// CB-576: a release that would otherwise remove the worktree finds it holding
|
||||
// uncommitted work the bridge cannot see. A worker that ends a turn without
|
||||
// committing (normally because it stopped to ask a question or refused the turn)
|
||||
// has its only copy of that work in the worktree. Remove would --force-delete it,
|
||||
// so preserve the directory and tell an operator where to find it.
|
||||
try {
|
||||
memberLifecycle.released(removed.terminalId());
|
||||
log.debug("releasing session pane={} terminal={} state={} cause={}",
|
||||
removed.paneId(), removed.terminalId(), removed.state(), cause);
|
||||
if (preserveWorktree && removed.worktree() != null) {
|
||||
logPreservedForShutdown(removed);
|
||||
} else if (removed.worktree() != null && worktrees.hasUncommitted(removed.worktree())) {
|
||||
// CB-576: a release that would otherwise remove the worktree finds it holding
|
||||
// uncommitted work the bridge cannot see. A worker that ends a turn without
|
||||
// committing (normally because it stopped to ask a question or refused the turn)
|
||||
// has its only copy of that work in the worktree. Remove would --force-delete it,
|
||||
// so preserve the directory and tell an operator where to find it.
|
||||
preserveWorktree = true;
|
||||
log.warn("release {} preserves dirty worktree {} for pane={} terminal={}: "
|
||||
+ "the worktree holds uncommitted changes that --force remove would destroy",
|
||||
cause, removed.worktree(), removed.paneId(), removed.terminalId());
|
||||
}
|
||||
} catch (RuntimeException e) {
|
||||
// CB-581: hasUncommitted shells out to `git status` and can throw on a non-zero
|
||||
// exit. We can no longer tell whether the worktree holds uncommitted work, so fail
|
||||
// toward the safe answer and preserve it — deleting on a guess can destroy work
|
||||
// that has no other copy (CB-576), while keeping it on a false alarm only costs
|
||||
// disk. The exception must not propagate: the pane still has to stop below.
|
||||
preserveWorktree = true;
|
||||
log.warn("release {} preserves dirty worktree {} for pane={} terminal={}: "
|
||||
+ "the worktree holds uncommitted changes that --force remove would destroy",
|
||||
cause, removed.worktree(), removed.paneId(), removed.terminalId());
|
||||
log.warn("release {} could not tell whether worktree {} for pane={} terminal={} has "
|
||||
+ "uncommitted changes; preserving it rather than risk destroying unsaved work: {}",
|
||||
cause, removed.worktree(), removed.paneId(), removed.terminalId(), e.toString());
|
||||
} finally {
|
||||
// CB-516/CB-581: a send still waiting on this worker can never be answered now, no
|
||||
// matter what happened above. Tell the listener BEFORE the pane is torn down, so a
|
||||
// blocked caller fails fast with a real reason instead of sitting on a rendezvous
|
||||
// nothing will ever resolve.
|
||||
notifyReleased(removed.terminalId());
|
||||
}
|
||||
// CB-516: a send still waiting on this worker can never be answered now. Tell the
|
||||
// listener BEFORE the pane is torn down, so a blocked caller fails fast with a real
|
||||
// reason instead of sitting on a rendezvous nothing will ever resolve.
|
||||
notifyReleased(removed.terminalId());
|
||||
}
|
||||
// CB-581: the pane must always stop, even if the dirty check above threw. A session removed
|
||||
// from the registry with no pane stop is an orphaned pane — a live terminal burning a fleet
|
||||
// slot that no longer appears in the roster and can never be reclaimed.
|
||||
launcher.stop(paneId);
|
||||
if (removed != null && !preserveWorktree && removed.worktree() != null) {
|
||||
worktrees.remove(worktrees.repoRoot(removed.cwd()), removed.worktree());
|
||||
@@ -539,8 +556,15 @@ public final class SessionManager implements TurnListener {
|
||||
log.debug("reaping idle session terminal={} pane={}: idle {}s exceeds the {}s ttl",
|
||||
s.terminalId(), s.paneId(), TimeUnit.NANOSECONDS.toSeconds(idleNanos),
|
||||
TimeUnit.NANOSECONDS.toSeconds(idleTtlNanos));
|
||||
release(s.paneId());
|
||||
reaped++;
|
||||
// CB-581: one session that fails to release must not abort the whole reaping pass —
|
||||
// match drainAll's per-session try/catch so the rest of the roster still gets reaped.
|
||||
try {
|
||||
release(s.paneId());
|
||||
reaped++;
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("reap failed for pane={} terminal={} worktree={}; continuing with "
|
||||
+ "remaining sessions", s.paneId(), s.terminalId(), s.worktree(), e);
|
||||
}
|
||||
}
|
||||
}
|
||||
return reaped;
|
||||
|
||||
@@ -338,54 +338,6 @@ 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(), ExhaustionSink.none());
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.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(), ExhaustionSink.none());
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.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(), ExhaustionSink.none());
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.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(), ExhaustionSink.none());
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.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(), ExhaustionSink.none());
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.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(), ExhaustionSink.none());
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.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(), ExhaustionSink.none());
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.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(), ExhaustionSink.none());
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.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(), ExhaustionSink.none());
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.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(), ExhaustionSink.none());
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.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(), ExhaustionSink.none());
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.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(), ExhaustionSink.none());
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.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(), ExhaustionSink.none());
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.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(), ExhaustionSink.none());
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.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(), ExhaustionSink.none());
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.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, ExhaustionSink.none());
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns);
|
||||
|
||||
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, ExhaustionSink.none());
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
@@ -407,53 +407,12 @@ 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, ExhaustionSink.none());
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
@@ -471,7 +430,7 @@ class CompletionResolverTest {
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver =
|
||||
new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
@@ -72,8 +72,7 @@ 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"),
|
||||
BridgeMcp.QuarantineSource.none());
|
||||
metrics, BridgeMcp.CapacitySource.none(), new BridgeMcp.HealthCoverageSource(() -> "off"));
|
||||
return mcp;
|
||||
}
|
||||
|
||||
|
||||
@@ -16,7 +16,6 @@ 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;
|
||||
@@ -438,29 +437,11 @@ class BridgeMcpTest {
|
||||
@Test
|
||||
void profilesListsConfiguredProfilesAndDefault() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
McpSchema.CallToolResult res = BridgeMcp.profiles(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), BridgeMcp.QuarantineSource.none());
|
||||
McpSchema.CallToolResult res = BridgeMcp.profiles(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
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,7 +16,6 @@ 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;
|
||||
@@ -28,8 +27,6 @@ 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.*;
|
||||
|
||||
@@ -131,13 +128,6 @@ 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 —
|
||||
@@ -564,105 +554,4 @@ 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,7 +6,6 @@ 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;
|
||||
@@ -34,7 +33,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(), ExhaustionSink.none());
|
||||
new CompletionResolver(agents, rendezvous, ExhaustedPatternLookup.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);
|
||||
|
||||
@@ -1,108 +0,0 @@
|
||||
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,13 +24,7 @@ class PlacementPolicyTest {
|
||||
private static PlacementContext ctx(List<PlacementCandidate> candidates,
|
||||
Function<String, Integer> liveCount,
|
||||
Set<String> 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);
|
||||
return new PlacementContext("b", candidates, liveCount, unreachable);
|
||||
}
|
||||
|
||||
private static PlacementContext ctx(List<PlacementCandidate> candidates,
|
||||
@@ -52,37 +46,17 @@ 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());
|
||||
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());
|
||||
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();
|
||||
@@ -108,18 +82,6 @@ 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();
|
||||
@@ -188,29 +150,6 @@ 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();
|
||||
|
||||
@@ -46,6 +46,82 @@ class SessionManagerTest {
|
||||
return sessionManager(herdr, clock, 0);
|
||||
}
|
||||
|
||||
private SessionManager sessionManager(FakeHerdr herdr, Worktrees worktrees) {
|
||||
return sessionManager(herdr, worktrees, System::nanoTime);
|
||||
}
|
||||
|
||||
private SessionManager sessionManager(FakeHerdr herdr, Worktrees worktrees, LongSupplier clock) {
|
||||
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN",
|
||||
List.of("ccs", "ltms-local"), "tab", "bridged-workers",
|
||||
"worker: {profile} #{n}", null, null, null);
|
||||
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
|
||||
return new SessionManager(workers, worktrees, clock);
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-581: a {@link Worktrees} test double whose {@code hasUncommitted} and {@code remove} can
|
||||
* be told to throw, so {@link SessionManager#release} can be exercised against exactly the
|
||||
* failure {@code GitWorktrees} produces when {@code git status}/{@code git worktree remove}
|
||||
* exits non-zero.
|
||||
*/
|
||||
private static final class RecordingWorktrees implements Worktrees {
|
||||
private final List<String> removeCalls = new java.util.ArrayList<>();
|
||||
private final java.util.Set<String> failRemoveFor = new java.util.HashSet<>();
|
||||
private volatile boolean dirty = false;
|
||||
private volatile RuntimeException hasUncommittedFailure;
|
||||
|
||||
RecordingWorktrees dirty(boolean dirty) {
|
||||
this.dirty = dirty;
|
||||
return this;
|
||||
}
|
||||
|
||||
RecordingWorktrees failHasUncommittedWith(RuntimeException e) {
|
||||
this.hasUncommittedFailure = e;
|
||||
return this;
|
||||
}
|
||||
|
||||
RecordingWorktrees failRemoveFor(String worktreePath) {
|
||||
failRemoveFor.add(worktreePath);
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String add(String repoRoot, String branch, String baseRef) {
|
||||
return "/wt/" + branch.replace('/', '_');
|
||||
}
|
||||
|
||||
@Override
|
||||
public void remove(String repoRoot, String worktreePath) {
|
||||
if (failRemoveFor.contains(worktreePath)) {
|
||||
throw new WorktreeException("simulated remove failure for " + worktreePath);
|
||||
}
|
||||
removeCalls.add(worktreePath);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean hasUncommitted(String worktreePath) {
|
||||
if (hasUncommittedFailure != null) {
|
||||
throw hasUncommittedFailure;
|
||||
}
|
||||
return dirty;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void overlayParity(String repoRoot, String worktreePath, List<String> overlay) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public String repoRoot(String cwd) {
|
||||
return "/repo";
|
||||
}
|
||||
|
||||
List<String> removeCalls() {
|
||||
return List.copyOf(removeCalls);
|
||||
}
|
||||
}
|
||||
|
||||
private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock, int contextCap) {
|
||||
return sessionManager(herdr, clock, contextCap, false);
|
||||
}
|
||||
@@ -525,4 +601,172 @@ class SessionManagerTest {
|
||||
"a listener failure must never prevent the teardown it is reacting to");
|
||||
assertTrue(sessions.get(s.paneId()).isEmpty(), "and the session is still deregistered");
|
||||
}
|
||||
|
||||
// --- CB-581: a throw inside release() must not orphan the pane or abort reapIdle -----------
|
||||
|
||||
@Test
|
||||
void releasePreservesWorktreeWhenDirtyCheckThrows() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees();
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-581a", null));
|
||||
worktrees.failHasUncommittedWith(new WorktreeException("git status exited 128"));
|
||||
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
sessionLog.addAppender(appender);
|
||||
sessionLog.setLevel(Level.WARN);
|
||||
try {
|
||||
assertDoesNotThrow(() -> sessions.release(s.paneId()),
|
||||
"a throwing dirty check must not abort the release");
|
||||
|
||||
assertTrue(worktrees.removeCalls().isEmpty(),
|
||||
"the worktree is preserved when its dirty state cannot be determined");
|
||||
String warn = appender.list.stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.filter(m -> m.contains(s.worktree()))
|
||||
.findFirst()
|
||||
.orElse("no warn logged naming the worktree");
|
||||
assertTrue(warn.contains(s.paneId()), "the WARN names the pane: " + warn);
|
||||
assertTrue(warn.contains(s.terminalId()), "the WARN names the terminal: " + warn);
|
||||
} finally {
|
||||
sessionLog.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void releaseStillStopsThePaneWhenDirtyCheckThrows() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees();
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-581b", null));
|
||||
worktrees.failHasUncommittedWith(new WorktreeException("git status exited 128"));
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertEquals(1, paneCloseCallsFor(herdr, "w9:pRoot_1"),
|
||||
"the pane is stopped exactly once even though the dirty check threw");
|
||||
}
|
||||
|
||||
@Test
|
||||
void releaseStillNotifiesTheListenerWhenDirtyCheckThrows() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees();
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
java.util.List<String> released = new java.util.concurrent.CopyOnWriteArrayList<>();
|
||||
sessions.onRelease(released::add);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-581c", null));
|
||||
worktrees.failHasUncommittedWith(new WorktreeException("git status exited 128"));
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertEquals(java.util.List.of(s.terminalId()), released,
|
||||
"a blocked caller must still be told the terminal was released, even though the "
|
||||
+ "dirty check threw");
|
||||
}
|
||||
|
||||
@Test
|
||||
void reapIdleSurvivesOneSessionThatFailsToRelease() {
|
||||
long[] clock = {0};
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees();
|
||||
SessionManager sessions = sessionManager(herdr, worktrees, () -> clock[0]);
|
||||
MemberSession a = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-581d", null));
|
||||
MemberSession b = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-581e", null));
|
||||
MemberSession c = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-581f", null));
|
||||
sessions.asPresence().markPresent(a.terminalId());
|
||||
sessions.asPresence().markPresent(b.terminalId());
|
||||
sessions.asPresence().markPresent(c.terminalId());
|
||||
// The middle session's worktree removal fails — release() propagates that, so this is the
|
||||
// one call reapIdle's per-session guard must survive without skipping the rest of the pass.
|
||||
worktrees.failRemoveFor(b.worktree());
|
||||
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
sessionLog.addAppender(appender);
|
||||
sessionLog.setLevel(Level.WARN);
|
||||
int reaped;
|
||||
try {
|
||||
clock[0] = 100;
|
||||
reaped = sessions.reapIdle(10);
|
||||
|
||||
String warn = appender.list.stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.filter(m -> m.contains(b.paneId()))
|
||||
.findFirst()
|
||||
.orElse("no reap-failure WARN logged");
|
||||
assertTrue(warn.contains(b.terminalId()), "the WARN names the failed session's terminal: " + warn);
|
||||
assertTrue(warn.contains(b.worktree()), "the WARN names the failed session's worktree: " + warn);
|
||||
} finally {
|
||||
sessionLog.detachAppender(appender);
|
||||
}
|
||||
|
||||
assertEquals(2, reaped, "the middle session's failure is logged, not counted as reaped");
|
||||
assertTrue(sessions.get(a.paneId()).isEmpty(), "the first session is still released");
|
||||
assertTrue(sessions.get(c.paneId()).isEmpty(), "the third session is still released");
|
||||
assertTrue(sessions.get(b.paneId()).isEmpty(),
|
||||
"the middle session is still deregistered even though its worktree removal threw");
|
||||
assertEquals(1, paneCloseCallsFor(herdr, "w9:pRoot_1"), "the first pane is stopped");
|
||||
assertEquals(1, paneCloseCallsFor(herdr, "w9:pRoot_2"),
|
||||
"the middle pane is still stopped even though its worktree removal failed");
|
||||
assertEquals(1, paneCloseCallsFor(herdr, "w9:pRoot_3"), "the third pane is stopped");
|
||||
}
|
||||
|
||||
@Test
|
||||
void unchangedRegressionCleanCompletedReleaseStillRemovesTheWorktree() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees();
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-581g", null));
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertEquals(List.of(s.worktree()), worktrees.removeCalls(),
|
||||
"COMPLETED release of a clean worktree still removes it");
|
||||
}
|
||||
|
||||
@Test
|
||||
void unchangedRegressionDirtyCompletedReleaseStillPreservesTheWorktree() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees().dirty(true);
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-581h", null));
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertTrue(worktrees.removeCalls().isEmpty(),
|
||||
"COMPLETED release of a dirty worktree still preserves it");
|
||||
}
|
||||
|
||||
@Test
|
||||
void unchangedRegressionShutdownDrainStillPreservesTheWorktree() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees();
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-581i", null));
|
||||
sessions.asPresence().markPresent(s.terminalId());
|
||||
|
||||
sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(100));
|
||||
|
||||
assertTrue(worktrees.removeCalls().isEmpty(), "SHUTDOWN drain still preserves the worktree");
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user