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
5 changed files with 37 additions and 95 deletions
@@ -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) {
@@ -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.