Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha 16de9df000 CB-601: make the recovery-race test's head start deterministic, not a sleep
CI / contract (pull_request) Successful in 1m8s
CI / build (pull_request) Successful in 1m9s
2026-08-16 18:02:57 +02:00
8 changed files with 44 additions and 209 deletions
+6 -35
View File
@@ -88,28 +88,15 @@ bind:
# backoffMs: 60000
# quietNudgeCap: 3
# Fleet health detection is dormant unless enabled (CB-573). It reads one whole-fleet agent list
# per tick.
# intervalSeconds → how often a tick runs (default 30). ENFORCED floor of 15: the code computes
# Math.max(15, intervalSeconds), so a lower value is silently raised, not
# rejected.
# workingSuspectAfterSeconds, paneProbeIntervalSeconds → accepted and parsed, but NOT YET READ by
# anything — the dormant monitor only consumes intervalSeconds today (CB-573
# shipped ahead of the evidence publishers these two knobs are for). Setting
# them changes nothing right now, and no minimum is enforced on either, because
# nothing reads them to enforce one. They exist so a later build can start
# honouring them without another config-shape change.
# notifications.mode → "webhook" flips what bridge_list REPORTS (healthCoverage: "full" instead
# of "detection-only") — it does NOT make bridged send any webhook call; no
# delivery mechanism is implemented yet. Any other value, or omitting the
# block, reports "detection-only".
# Fleet health detection is dormant unless enabled. It reads one whole-fleet agent list per tick.
# It can run without a webhook; bridge_list then reports healthCoverage: detection-only.
# health:
# enabled: true
# intervalSeconds: 30
# workingSuspectAfterSeconds: 600
# paneProbeIntervalSeconds: 60
# intervalSeconds: 30 # minimum 15
# workingSuspectAfterSeconds: 600 # minimum 300
# paneProbeIntervalSeconds: 60 # minimum 60
# notifications:
# mode: disabled
# mode: disabled # disabled (default) or webhook
# herdr Unix socket. Omit to use the client default
# (${HERDR_SOCKET_PATH:-~/.config/herdr/herdr.sock}).
@@ -299,11 +286,6 @@ placement: weighted
# / credentialId. Those are hot because the placement policy (and, for credentialId,
# the CB-578 stage B quarantine check) reads them through a supplier — being config is
# not by itself enough to make a key hot.
# EXCEPT `fleet.leaders`: Bridged.main reads it once at startup to build the lead tab
# scanner and launcher, and neither is rebuilt on reload. A changed/added/removed
# `fleet.leaders` entry is silently accepted — the reload reports "config reloaded"
# with nothing in the deferred list — but has NO effect until you restart. Treat it
# as deferred in practice, even though today's reload output does not say so.
# 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
@@ -386,12 +368,6 @@ fleet:
# An auto-launched lead is NOT a member: it gets no worker reply charter, is never registered with
# the session lifecycle (the idle reaper would kill your orchestrator), and stays on the
# subscription — ANTHROPIC_BASE_URL/AUTH_TOKEN are stripped from its env whatever the profile says.
#
# GET THE `tab:` VALUE RIGHT. A pane that does not match any configured `tab:` (a typo, a renamed
# tab, a pane no entry names at all) is not recognised as a lead — it resolves as an ordinary
# WORKER instead, silently, and every orchestration call it makes (spawn/stop/send/drain) is
# refused. There is no error at startup for this: an unmatched pane is simply not a lead. If your
# primary suddenly can't spawn or send, check this section first.
# leaders:
# opus-5.0:
# profile: opus # omit to never create this lead, only recognise it
@@ -447,15 +423,10 @@ guard:
# idleTtlSeconds → reap READY/DONE sessions idle longer than this (never BUSY/SPAWNING)
# contextCap → force-release a session after this many delegated turns
# drainTimeoutSeconds → seconds to wait for BUSY sessions on shutdown before forced teardown
# clearAfterTurn → whether a reusable worker discards its conversation context after every
# completed delegated turn (default false). Works for claude-code workers
# only — any other peer kind (e.g. opencode) logs "context reset is
# unsupported for peer kind …" once and the reset is a no-op.
# lifecycle:
# idleTtlSeconds: 300
# contextCap: 10
# drainTimeoutSeconds: 5
# clearAfterTurn: false
# Durable reply delivery (CB-307 Stage 2). OMIT this block entirely to keep the default
# in-memory, soft-state reply inbox (late worker replies are held only until a daemon bounce).
@@ -967,12 +967,8 @@ public record BridgedConfig(
/**
* Top-level keys this version understands. Used only to warn about the rest — see
* {@link #warnUnknownTopLevelKeys}. Keep in step with the record components.
*
* <p>Package-private (not {@code private}) so a test can assert every key here is documented in
* {@code bridged.example.yaml} — the only committed description of the config schema, since
* {@code bridged.yaml} itself is gitignored.
*/
static final Set<String> KNOWN_TOP_LEVEL_KEYS = Set.of(
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");
@@ -15,7 +15,6 @@ 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.placement.PlacementException;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.bridged.session.WorktreeRequest;
@@ -693,10 +692,6 @@ public final class BridgeMcp {
return text(json(memberView(member)));
} catch (GuardException e) {
return error("subscription boundary: " + e.getMessage());
} catch (PlacementException e) {
// CB-599: no candidate had capacity (maxLoad, quarantine, or all-exhausted) — distinct
// from "profile does not exist" below.
return error("no capacity: " + e.getMessage());
} catch (IllegalArgumentException e) {
return error(e.getMessage()); // unknown / no-default profile, or a refused resumeSessionId
} catch (PeerUnreachableException e) {
@@ -13,7 +13,6 @@ import dev.ltms.bridged.herdr.HerdrClient;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.inject.MemberPresence;
import dev.ltms.bridged.peer.PeerUnreachableException;
import dev.ltms.bridged.placement.PlacementException;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.peer.MemberRole;
@@ -293,11 +292,6 @@ public final class BridgedApp {
ctx.status(201).json(view(member));
} catch (GuardException e) {
ctx.status(403).json(Map.of("error", "subscription_boundary", "detail", e.getMessage()));
} catch (PlacementException e) {
// CB-599: no candidate had capacity (maxLoad, quarantine, or all-exhausted) — a benign,
// likely-transient refusal, distinct from "profile does not exist" below. 503: the
// request was valid and will likely succeed later.
ctx.status(503).json(Map.of("error", "no_capacity", "detail", e.getMessage()));
} catch (IllegalArgumentException e) {
ctx.status(400).json(Map.of("error", "unknown_profile", "detail", e.getMessage()));
} catch (PeerUnreachableException e) {
@@ -10,7 +10,6 @@ import java.nio.file.Path;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.regex.Pattern;
import static org.junit.jupiter.api.Assertions.*;
@@ -1190,79 +1189,6 @@ class BridgedConfigTest {
assertEquals(5, cfg.leadHeartbeat().quietNudgeCap());
}
/**
* A top-level key {@code BridgedConfig} reads but that appears nowhere in
* {@code bridged.example.yaml} — live or commented — is invisible drift: {@code bridged.yaml}
* is gitignored, so the example is the ONLY committed description of the config schema, and
* neither {@link #shippedExampleConfigParses} (example → code: does the example still parse)
* nor {@link #everyOptionalKnobDocumentedInTheExampleBinds} (a hand-maintained list of keys
* that must bind) can catch a brand-new key nobody added to either.
*
* <p>This test compares the OTHER direction: every key in {@link BridgedConfig#KNOWN_TOP_LEVEL_KEYS}
* (the parser's own accepted set, which backs the unknown-key WARN) must appear as a top-level
* key in the example text, live or commented-out — see {@link #topLevelKeyDocumented}.
*/
@Test
void everyKnownTopLevelKeyIsDocumentedInTheExample() throws Exception {
Path example = Path.of("bridged.example.yaml");
assertTrue(Files.exists(example), "bridged.example.yaml must ship next to the pom");
String text = Files.readString(example);
List<String> undocumented = BridgedConfig.KNOWN_TOP_LEVEL_KEYS.stream()
.filter(key -> !topLevelKeyDocumented(text, key))
.sorted()
.toList();
assertTrue(undocumented.isEmpty(), () -> "key(s) " + undocumented
+ " are read by BridgedConfig but appear nowhere in bridged.example.yaml — "
+ "document each one there, commented out if optional. bridged.yaml is "
+ "gitignored, so this file is the only committed description of the config "
+ "schema an operator or a worker can see.");
}
/**
* Most of {@code bridged.example.yaml} is deliberately commented out — optional sections are
* documented as commented blocks so the shipped file stays a working minimal config. A key
* documented ONLY as a comment must still count as documented; parsing the file as YAML and
* reading its live key set (as an earlier attempt at this guard did) gets this wrong, because
* every commented section then looks entirely absent.
*/
@Test
void commentedOnlyTopLevelKeyCountsAsDocumented() {
String yaml = """
bind:
port: 8765
# broker:
# uri: amqp://guest:guest@127.0.0.1:5672
""";
assertTrue(topLevelKeyDocumented(yaml, "broker"),
"a key documented only inside a commented-out block must still count as documented");
}
/** A key that appears in neither a live nor a commented top-level line must NOT count. */
@Test
void absentTopLevelKeyIsNotDocumented() {
String yaml = """
bind:
port: 8765
""";
assertFalse(topLevelKeyDocumented(yaml, "broker"),
"a key mentioned nowhere in the example must not be reported as documented");
}
/**
* True when {@code key} appears as a top-level YAML key in {@code yaml} — either live
* ({@code key:} at column 0) or commented out ({@code # key:}, also at column 0, with only
* whitespace between the {@code #} and the key). Anchoring on column 0 is what keeps this a
* top-level check: an indented occurrence (a nested field, or prose inside a comment that
* happens to end in a colon) never matches, because {@code ^} requires the key's own first
* character — or the sole leading {@code #} — to sit at the very start of the line.
*/
private static boolean topLevelKeyDocumented(String yaml, String key) {
Pattern p = Pattern.compile("(?m)^(?:#\\s*)?" + Pattern.quote(key) + ":");
return p.matcher(yaml).find();
}
@Test
void placementDefaultsToFixedForExistingConfigs(@TempDir Path dir) throws Exception {
Path f = dir.resolve("no-placement.yaml");
@@ -16,15 +16,12 @@ 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.member.CompositePeerLauncher;
import dev.ltms.bridged.placement.BackendQuarantine;
import dev.ltms.bridged.placement.PlacementPolicies;
import io.modelcontextprotocol.spec.McpSchema;
import dev.ltms.bridged.msg.InMemoryReplyInbox;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
@@ -426,35 +423,6 @@ class BridgeMcpTest {
assertTrue(textOf(res).contains("unknown worker profile"), textOf(res));
}
/**
* CB-599: a profile at its {@code maxLoad} cap must surface a readable reason on the MCP
* surface too, not merely flip {@code isError} with an opaque or absent message.
*/
@Test
void spawnAtMaxLoadSurfacesTheCapacityReason() {
FakeHerdr h = new FakeHerdr();
BridgedConfig.Profile wcfg = new BridgedConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN", null,
"tab", "bridged-workers", "worker: {profile} #{n}", null,
null, null, null, null, null, null, null, 0, null, null, null);
Map<String, BridgedConfig.Profile> profiles = Map.of(wcfg.profile(), wcfg);
ClaudeCodeLauncher delegate = new ClaudeCodeLauncher(
new AgentControl(h), new WorkspaceControl(h), new SubscriptionGuard(Set.of("gx00.gw")),
profiles, wcfg.profile(), k -> "BRIDGED_WORKER_TOKEN".equals(k) ? "tok" : null);
CompositePeerLauncher composite = new CompositePeerLauncher(
List.of(delegate), wcfg.profile(), profiles, PlacementPolicies.fixed(), _ -> 0);
SessionManager sm = new SessionManager(composite);
McpSchema.CallToolResult res = BridgeMcp.spawn(sm, "ltms-local");
assertTrue(res.isError());
String text = textOf(res);
assertTrue(text.contains("no capacity"), "surfaces a capacity reason, not a bare error: " + text);
assertTrue(text.contains("ltms-local"), "names the profile: " + text);
assertTrue(text.contains("maxLoad"), "explains the refusal: " + text);
assertFalse(h.called("agent.start"), "at cap, the spawn is refused before any herdr call");
}
@Test
void spawnPassesTheRequestedCwdToTheWorker() {
FakeHerdr h = new FakeHerdr();
@@ -35,9 +35,12 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
*/
class AmqpReplyInboxRecoveryRaceTest {
/** Large enough that the (unfixed) unsynchronized sweep's iteration is a real, observable window
* a concurrently-started publish can land in — not just a best case, single-entry sprint. */
private static final int STALE_PUBLISHES = 100_000;
/** Large enough that thousands of entries are still unprocessed by the time the very first one
* is observed as failed (see {@code sweepIsHoldingTheLock} below) — that gap is what makes the
* head start deterministic instead of a coin flip. 100,000 gave the same guarantee but made the
* test far more expensive than the guarantee needs; the ordering no longer depends on a timing
* window sized to the full backlog; just to the tail of it. */
private static final int STALE_PUBLISHES = 2_000;
@Test
@Timeout(30)
@@ -58,6 +61,12 @@ class AmqpReplyInboxRecoveryRaceTest {
// pendingByMsgId exactly like publishes whose confirm never arrived before a connection drop.
// Virtual threads make this many concurrent blocking publish() calls cheap.
CountDownLatch staleStarted = new CountDownLatch(STALE_PUBLISHES);
// Counted down by the FIRST stale publish thread to observe its own failure. That can only
// happen from inside failPendingPublishesOnRecovery() — nothing else in this test ever
// completes a stale Pending exceptionally (no nack/return is simulated for any "stale-*"
// msgId) — so seeing it fire is direct, observable proof the sweep is inside its loop, not a
// timing guess. It replaces the old fixed Thread.sleep(5) head start.
CountDownLatch sweepIsHoldingTheLock = new CountDownLatch(1);
for (int i = 0; i < STALE_PUBLISHES; i++) {
String msgId = "stale-" + i;
Thread.ofVirtual().start(() -> {
@@ -65,7 +74,7 @@ class AmqpReplyInboxRecoveryRaceTest {
try {
inbox.publish("worker-stale", msgId, "x");
} catch (IllegalStateException expected) {
// resolved (failed by the sweep) — that is exactly what this thread is here for
sweepIsHoldingTheLock.countDown();
}
});
}
@@ -83,13 +92,28 @@ class AmqpReplyInboxRecoveryRaceTest {
}
}, "recovery-sweep");
sweepThread.start();
// A short, deliberate head start: with STALE_PUBLISHES this large, the (unfixed) sweep's own
// iteration takes several milliseconds, so this guarantees the sweep has already begun —
// and, once guarded, is already holding publishChannelLock — before "fresh" attempts to
// register. Without this head start, "fresh" sometimes wins the race for the lock and
// registers before the sweep even starts, which is the accepted "already in flight when
// recovery fires" case (correctly failed either way) rather than the bug under test.
Thread.sleep(5);
// Deterministic head start: block until the sweep has actually failed one of the stale
// publishes. failPendingPublishesOnRecovery() (once guarded, as it is on main) holds
// publishChannelLock for its ENTIRE loop, not just per entry — so this failure proves the
// sweep is, at this instant, still holding that lock. With STALE_PUBLISHES this large, the
// remaining ~1,999 entries give an enormous margin between "first failure observed" and "sweep
// releases the lock": there is no window left for "fresh" to slip in before the sweep starts,
// or to win the lock ahead of it — see the case-2 note below. This also means Case 1 (the sweep
// is already inside its loop, holding the lock, when "fresh" tries to register) is now
// guaranteed by construction rather than merely likely under a fixed sleep.
assertTrue(sweepIsHoldingTheLock.await(20, TimeUnit.SECONDS),
"the sweep never failed a single stale publish — it may not have started");
// Case 2 ("fresh" wins publishChannelLock before the sweep even starts, so it genuinely
// published on the stale channel and the sweep correctly fails it) is impossible by
// construction in this test: freshThread.start() below is reached only after
// sweepIsHoldingTheLock has counted down, which can only happen once
// failPendingPublishesOnRecovery() is already running and has already failed a stale entry.
// There is no code path that lets "fresh" start before the sweep starts. That case is real
// and correct production behaviour (see AmqpReplyInbox#failPendingPublishesOnRecovery's
// javadoc), it is just not reachable from this deterministic ordering, so it does not need a
// separate assertion here.
// This is the exact interleaving CB-528's follow-up describes: "the still-running recovery
// sweep" racing a publish that registers while it is mid-flight.
@@ -109,8 +133,8 @@ class AmqpReplyInboxRecoveryRaceTest {
// Simulate the broker's real confirm for "fresh" now that the sweep is done, so a correct
// implementation's publish() returns normally instead of idling out CONFIRM_TIMEOUT_MS. Poll
// for the registration rather than checking once: freshThread may still be contending for
// publishChannelLock (behind the 20,000 stale threads' own lock acquisitions) even though the
// sweep itself has already finished.
// publishChannelLock (behind the sweep's own hold on it, and possibly other stale threads
// still unwinding) even though the sweep itself has already finished.
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(9);
int idx = -1;
while (idx < 0 && System.nanoTime() < deadline) {
@@ -18,8 +18,6 @@ import dev.ltms.bridged.session.GitWorktrees;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.session.Worktrees;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
import dev.ltms.bridged.member.CompositePeerLauncher;
import dev.ltms.bridged.placement.PlacementPolicies;
import io.javalin.Javalin;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
@@ -242,43 +240,6 @@ class BridgedAppTest {
assertFalse(herdr.called("agent.start"), "an unknown profile must not spawn anything");
}
/**
* CB-599: a profile at its {@code maxLoad} cap must not surface as a bare 500 — the caller
* needs a structured, readable reason, distinct from "unknown_profile".
*/
@Test
void spawnAtMaxLoadIs503WithTheCapacityReasonNotABare500() throws Exception {
FakeHerdr herdr = new FakeHerdr();
BridgedConfig.Profile wcfg = new BridgedConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN", null,
"tab", "bridged-workers", "worker: {profile} #{n}", null,
null, null, null, null, null, null, null, 0, null, null, null);
Map<String, BridgedConfig.Profile> profiles = Map.of(wcfg.profile(), wcfg);
ClaudeCodeLauncher delegate = new ClaudeCodeLauncher(
new AgentControl(herdr), new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
profiles, wcfg.profile(), k -> "BRIDGED_WORKER_TOKEN".equals(k) ? "tok-abc" : null);
CompositePeerLauncher workers = new CompositePeerLauncher(
List.of(delegate), wcfg.profile(), profiles, PlacementPolicies.fixed(), _ -> 0);
SessionManager sessions = new SessionManager(workers, new GitWorktrees());
this.presence = sessions.asPresence();
Injector injector = new Injector(new AgentControl(herdr));
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
sessions.onAcquire(inbox::own);
MessageService messages = new MessageService(new AgentControl(herdr), injector, new Rendezvous(), inbox);
app = new BridgedApp(herdr, workers, sessions, messages, this.presence, null).build().start("127.0.0.1", 0);
int port = app.port();
HttpResponse<String> res = req(port, "POST", "/members?profile=ltms-local");
assertEquals(503, res.statusCode(), res.body());
JsonNode body = mapper.readTree(res.body());
assertEquals("no_capacity", body.get("error").asText());
String detail = body.get("detail").asText();
assertTrue(detail.contains("ltms-local"), "detail names the profile: " + detail);
assertTrue(detail.contains("maxLoad"), "detail explains the refusal: " + detail);
assertFalse(herdr.called("agent.start"), "at cap, the spawn is refused before any herdr call");
}
@Test
void spawnWorkerReusesExistingWorkerSpace() throws Exception {
// A space labelled "bridged-workers" already exists → no second workspace.create.