CB-306: spawn-readiness gate — ClaudeCodeLauncher blocks until the worker is injectable or throws PeerUnreachableException

The launcher now polls AgentControl.status(paneId) after starting the pane.
It returns the handle only once the worker reports an injectable state
(IDLE/BLOCKED/DONE). If the timeout elapses while still UNKNOWN, the
pane is self-reaped and a PeerUnreachableException is thrown — no orphan
left behind. The gate is disabled when spawnReadyTimeoutMs == 0 (legacy
non-blocking spawn, the default for the 6-arg constructor).

Key changes:
- PeerUnreachableException (new) in dev.ltms.bridged.peer
- BridgedConfig: spawnReadyTimeoutMs (default 20000), spawnReadyPollMs (default 300)
- ClaudeCodeLauncher: 3 constructor overloads:
  (a) 6-arg backward-compat: gate disabled (timeout=0)
  (b) 8-arg production: gate with config knobs + real clock/sleep
  (c) 10-arg testability: full seam (LongSupplier clock + Runnable sleeper)
- waitUntilInjectableOrThrow() loop in spawn(SpawnRequest)
- sleepUninterruptibly() helper for the production sleeper
- BridgeMcp.spawn + BridgedApp.spawnWorker catch PeerUnreachableException
  → clean tool error / 502 response (not an uncaught 500)
- SessionManager.acquire inherently registers nothing on throw (both
  worktree and non-worktree paths) — confirmed by new test

Tests:
- ClaudeCodeLauncherTest: 4 new tests
  - unknown→idle: returns handle, no pane.close
  - always-unknown: throws PeerUnreachableException, pane closed,
    clock advanced past timeout
  - timeout=0 (6-arg ctor): no agent.get calls, returns handle
  - timeout=0 (10-arg ctor): no orphan pane close
- SessionManagerTest: 1 new test
  - acquire → PeerUnreachableException: roster remains empty
Total: 188 tests, all pass (no existing test changed semantics)
This commit is contained in:
Dai Ha
2026-07-18 16:17:43 +02:00
parent 3a5cdc5108
commit 7dd6c46156
9 changed files with 267 additions and 5 deletions
+5
View File
@@ -54,6 +54,11 @@ guard:
- gx01.gw
- ollama.ltms.dev
# Spawn-readiness gate (CB-306). The launcher blocks until the worker's herdr status is
# injectable (IDLE/BLOCKED/DONE) or the timeout elapses. 0 disables the gate.
# spawn_ready_timeout_ms: 20000
# spawn_ready_poll_ms: 300
# Session lifecycle limits (CB-303). All knobs are opt-in; omit or set to null to keep
# the feature disabled. By default the daemon never reaps, caps, or drains sessions.
# idleTtlSeconds → reap READY/DONE sessions idle longer than this (never BUSY/SPAWNING)
@@ -58,7 +58,8 @@ public final class Bridged {
AgentControl agents = new AgentControl(herdr);
WorkspaceControl spaces = new WorkspaceControl(herdr);
PeerLauncher workers = new ClaudeCodeLauncher(agents, spaces, guard,
cfg.workerProfiles(), cfg.defaultProfile(), System::getenv);
cfg.workerProfiles(), cfg.defaultProfile(), System::getenv,
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs());
// CB-117: herdr keeps worker panes alive across a daemon restart, and their ids died with
// the previous process — reap those leaked orphans now, before we start serving.
workers.reapOrphanWorkers();
@@ -27,7 +27,10 @@ import java.util.Set;
* @param guard subscription-boundary allowlist
* @param worktreeRoot nullable root directory for provisioned worktrees; defaults to a sibling
* of the repo root
* @param lifecycle session lifecycle limits ({@code null} = all disabled)
* @param lifecycle session lifecycle limits ({@code null} = all disabled)
* @param spawnReadyTimeoutMs max ms to wait for a spawned worker to reach an injectable state
* ({@code null} / 0 disables the poll gate — legacy non-blocking behaviour)
* @param spawnReadyPollMs poll interval while waiting for the worker to become injectable
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record BridgedConfig(
@@ -38,7 +41,9 @@ public record BridgedConfig(
String defaultWorker,
Guard guard,
String worktreeRoot,
Lifecycle lifecycle) {
Lifecycle lifecycle,
Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs) {
@JsonIgnoreProperties(ignoreUnknown = true)
public record Bind(String host, int port) {
@@ -229,6 +234,8 @@ public record BridgedConfig(
Bind b = bind != null ? bind : new Bind(null, 0);
Guard g = guard != null ? guard : new Guard(List.of());
Lifecycle l = lifecycle != null ? lifecycle : new Lifecycle(null, null, null);
return new BridgedConfig(b, herdrSocket, worker, workers, defaultWorker, g, worktreeRoot, l);
Integer timeout = (spawnReadyTimeoutMs != null) ? spawnReadyTimeoutMs : 20000;
Integer pollMs = (spawnReadyPollMs != null) ? spawnReadyPollMs : 300;
return new BridgedConfig(b, herdrSocket, worker, workers, defaultWorker, g, worktreeRoot, l, timeout, pollMs);
}
}
@@ -6,6 +6,7 @@ import dev.ltms.bridged.inject.WorkerPresence;
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.session.SessionManager;
import dev.ltms.bridged.session.WorkerSession;
import dev.ltms.bridged.session.WorktreeRequest;
@@ -308,6 +309,8 @@ public final class BridgeMcp {
return error("subscription boundary: " + e.getMessage());
} catch (IllegalArgumentException e) {
return error(e.getMessage()); // unknown / no-default profile
} catch (PeerUnreachableException e) {
return error("spawn timed out — worker pane never reached injectable state: " + e.getMessage());
} catch (HerdrException e) {
return error("herdr error spawning worker: " + e.getMessage());
}
@@ -0,0 +1,18 @@
package dev.ltms.bridged.peer;
/**
* Thrown when a {@link PeerLauncher} starts a peer process but the peer
* does not reach an injectable (ready-to-receive) state within the configured
* timeout. The launcher MUST clean up any resources it created (pane, tab)
* before throwing — no orphaned peer or pane is left behind.
*
* <p>This is a spawn-time failure, distinct from a post-spawn disconnect.
* Callers treat this as a clean spawn error (the peer never materialized
* into a usable session), not a mid-life session fault.
*/
public final class PeerUnreachableException extends RuntimeException {
public PeerUnreachableException(String message) {
super(message);
}
}
@@ -7,6 +7,7 @@ import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.herdr.HerdrClient;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.inject.WorkerPresence;
import dev.ltms.bridged.peer.PeerUnreachableException;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.session.SessionManager;
@@ -178,6 +179,8 @@ public final class BridgedApp {
ctx.status(403).json(Map.of("error", "subscription_boundary", "detail", e.getMessage()));
} catch (IllegalArgumentException e) {
ctx.status(400).json(Map.of("error", "unknown_profile", "detail", e.getMessage()));
} catch (PeerUnreachableException e) {
ctx.status(502).json(Map.of("error", "spawn_timeout", "detail", e.getMessage()));
}
}
@@ -11,6 +11,7 @@ import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.peer.Capability;
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 org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -24,6 +25,7 @@ import java.util.Map;
import java.util.Set;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Function;
import java.util.function.LongSupplier;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
@@ -65,6 +67,11 @@ public final class ClaudeCodeLauncher implements PeerLauncher {
private final Function<String, String> env; // host env lookup (injectable for tests)
private final AtomicLong nameSeq = new AtomicLong(); // per-worker counter (also the tab #)
private final long spawnReadyTimeoutMs; // 0 = disable gate (legacy non-blocking spawn)
private final long spawnReadyPollMs;
private final LongSupplier nowMillis; // monotonic clock (injectable for tests)
private final Runnable sleeper; // sleep/wait hook (injectable for tests; never real-sleep in unit tests)
/**
* Standing instruction appended to the worker's system prompt so it returns its result via
* {@code bridge_reply}. Injected as a launch flag, so nothing is written to the worker's
@@ -86,15 +93,62 @@ public final class ClaudeCodeLauncher implements PeerLauncher {
// cannot collide with same-profile workers that outlived a restart. See startUniquelyNamed.
private final String nameNonce = String.format("%06x", new SecureRandom().nextInt(1 << 24));
/**
* Production constructor — disables the spawn-ready gate ({@code spawnReadyTimeoutMs == 0}) so
* existing deployments and tests keep the legacy non-blocking spawn semantics.
*/
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, BridgedConfig.Worker> profiles, String defaultProfile,
Function<String, String> env) {
this(agents, spaces, guard, profiles, defaultProfile, env, 0, 300,
System::currentTimeMillis, () -> sleepUninterruptibly(300));
}
/**
* Production constructor with spawn-ready gate enabled. The gate polls
* {@code agents.status()} until the pane reports an injectable state or
* {@code spawnReadyTimeoutMs} elapses.
*/
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, BridgedConfig.Worker> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs) {
this(agents, spaces, guard, profiles, defaultProfile, env,
spawnReadyTimeoutMs, spawnReadyPollMs,
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs));
}
/**
* Full testability constructor. Every injectable collaborator is explicit so unit tests
* supply fakes for the clock ({@code nowMillis}) and poll-loop wait ({@code sleeper}).
* The {@code sleeper} is never called when the gate is disabled ({@code spawnReadyTimeoutMs == 0}).
*
* @param agents herdr agent control (start, status, close)
* @param spaces workspace / tab control (ensure, create, close)
* @param guard subscription-boundary guard (checked before spawning)
* @param profiles configured worker profiles
* @param defaultProfile profile a no-argument spawn uses (nullable)
* @param env host env lookup (injectable for tests)
* @param spawnReadyTimeoutMs max ms to wait for injectable state (0 disables the gate)
* @param spawnReadyPollMs poll interval ms while waiting
* @param nowMillis monotonic clock source (e.g. {@code System::currentTimeMillis})
* @param sleeper sleep/wait hook (e.g. {@code () -> Thread.sleep(spawnReadyPollMs)})
*/
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, BridgedConfig.Worker> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs,
LongSupplier nowMillis, Runnable sleeper) {
this.agents = agents;
this.spaces = spaces;
this.guard = guard;
this.profiles = Map.copyOf(profiles);
this.defaultProfile = defaultProfile;
this.env = env;
this.spawnReadyTimeoutMs = spawnReadyTimeoutMs;
this.spawnReadyPollMs = spawnReadyPollMs > 0 ? spawnReadyPollMs : 300;
this.nowMillis = nowMillis;
this.sleeper = sleeper;
}
/** The configured worker profile names (what {@code spawn(profile)} accepts). */
@@ -459,11 +513,39 @@ public final class ClaudeCodeLauncher implements PeerLauncher {
* <p>Delegates to the three-arg {@link #spawn(String, String, String)} and wraps the
* resulting herdr {@link Agent} in a {@link WorkerHandle} whose {@link PeerHandle#id()}
* equals the agent's paneId.
*
* <p>When {@link #spawnReadyTimeoutMs} > 0, blocks until the worker's herdr status is
* {@link AgentStatus#injectable()} or the timeout elapses. On timeout the pane is closed
* (no orphan) and a {@link PeerUnreachableException} is thrown.
*/
@Override
public PeerHandle spawn(SpawnRequest req) {
Agent agent = spawn(req.profileName(), req.requestedCwd(), req.callerCwd());
return new WorkerHandle(agent.paneId(), agent.terminalId());
String paneId = agent.paneId();
if (spawnReadyTimeoutMs > 0) {
waitUntilInjectableOrThrow(paneId);
}
return new WorkerHandle(paneId, agent.terminalId());
}
/**
* Poll {@link AgentControl#status} until the pane reports an injectable state or the
* configured timeout elapses. On timeout, close the pane (self-reap) and throw.
*/
private void waitUntilInjectableOrThrow(String paneId) {
long deadline = nowMillis.getAsLong() + spawnReadyTimeoutMs;
while (nowMillis.getAsLong() < deadline) {
if (agents.status(paneId).injectable()) {
log.debug("worker pane={} reached injectable state", paneId);
return;
}
sleeper.run();
}
log.warn("worker pane={} did not become injectable within {}ms — closing", paneId, spawnReadyTimeoutMs);
stop(paneId);
throw new PeerUnreachableException(
"worker pane " + paneId + " did not reach injectable state within "
+ spawnReadyTimeoutMs + "ms");
}
/** A concrete {@link PeerHandle} wrapping herdr agent coordinates. */
@@ -485,4 +567,18 @@ public final class ClaudeCodeLauncher implements PeerLauncher {
private String resolveEnv(String name) {
return (name == null || name.isBlank()) ? null : env.apply(name);
}
/**
* Uninterruptible sleep — the production {@link #sleeper}. Tests supply their own
* no-op / fast-faking sleeper so they never real-sleep.
*/
private static void sleepUninterruptibly(long ms) {
try {
Thread.sleep(ms);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
// preserve the interrupt flag but continue — poll loops should not be
// aborted by an interrupt that was not meant for them.
}
}
}
@@ -6,6 +6,7 @@ import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.worker.ClaudeCodeLauncher;
import dev.ltms.bridged.peer.PeerUnreachableException;
import org.junit.jupiter.api.Test;
import java.util.List;
@@ -332,4 +333,33 @@ class SessionManagerTest {
.filter(c -> paneId.equals(((Map<?, ?>) c.params()).get("pane_id")))
.count();
}
// --- CB-306 spawn-readiness gate: no half-registered session on timeout ----------------
@Test
void acquireThrowsPeerUnreachableWhenGateTimesOutAndRegistersNoSession() {
FakeHerdr herdr = new FakeHerdr();
herdr.agentStatus("unknown"); // never becomes injectable
long[] clock = {0};
// Gate-enabled launcher (1 ms timeout + no-op sleeper that advances clock past deadline)
BridgedConfig.Worker cfg = new BridgedConfig.Worker(
"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,
1, 1, () -> clock[0], () -> clock[0] += 10);
SessionManager sessions = new SessionManager(workers, new GitWorktrees(), () -> 0L, 0);
assertThrows(PeerUnreachableException.class,
() -> sessions.acquire("ltms-local", null, "/caller", "term_primary"),
"acquire must throw PeerUnreachableException when spawn times out");
// No half-registered session — the error happened inside spawn, before
// SessionManager could put() anything into the registry.
assertTrue(sessions.roster().isEmpty(),
"no session is registered when spawn times out (roster empty)");
}
}
@@ -7,6 +7,7 @@ import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.peer.Capability;
import dev.ltms.bridged.peer.PeerHandle;
import dev.ltms.bridged.peer.PeerUnreachableException;
import dev.ltms.bridged.peer.SpawnRequest;
import org.junit.jupiter.api.Test;
@@ -318,4 +319,102 @@ class ClaudeCodeLauncherTest {
assertTrue(herdr.called("pane.close"), "stop via handle.id() must close the pane");
}
// --- CB-306 spawn-readiness gate -----------------------------------------------------------
private static Map<String, BridgedConfig.Worker> workerConfigMap(String profile, String mcpUrl) {
BridgedConfig.Worker cfg = new BridgedConfig.Worker(
profile, "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN",
List.of("ccs", profile), "tab", "bridged-workers",
"worker: {profile} #{n}", mcpUrl, null, null);
return Map.of(cfg.profile(), cfg);
}
@Test
void spawnWaitsUntilInjectableThenReturnsHandle() {
FakeHerdr herdr = new FakeHerdr();
herdr.agentStatus("unknown"); // first status call sees UNKNOWN
long[] clock = {0};
boolean[] firstSleep = {true};
// The sleeper: advance the fake clock, and on the first call flip the
// agent status to IDLE so the next poll succeeds.
Runnable sleeper = () -> {
clock[0] += 300;
if (firstSleep[0]) {
herdr.agentStatus("idle");
firstSleep[0] = false;
}
};
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(
new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")),
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
5000, 300, () -> clock[0], sleeper);
PeerHandle handle = svc.spawn(new SpawnRequest(null, null, null));
assertNotNull(handle, "spawn returns a handle when worker becomes injectable");
assertEquals("w9:pW_1", handle.id(), "handle id matches the started pane");
assertEquals(0, paneCloseCount(herdr, "w9:pW_1"),
"no pane.close when worker becomes injectable before timeout");
}
@Test
void spawnThrowsPeerUnreachableWhenNeverInjectableAndReapsPane() {
FakeHerdr herdr = new FakeHerdr();
herdr.agentStatus("unknown"); // always UNKNOWN
long[] clock = {0};
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(
new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")),
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
1000, 50, () -> clock[0], () -> clock[0] += 50);
PeerUnreachableException ex = assertThrows(
PeerUnreachableException.class,
() -> svc.spawn(new SpawnRequest(null, null, null)));
assertTrue(ex.getMessage().contains("w9:pW_1"),
"exception message references the paneId: " + ex.getMessage());
assertTrue(ex.getMessage().contains("1000"),
"exception message references the timeout: " + ex.getMessage());
assertTrue(clock[0] >= 1000, "fake clock advanced past the timeout: " + clock[0]);
assertEquals(1, paneCloseCount(herdr, "w9:pW_1"),
"pane was closed on timeout (no orphan left behind)");
}
@Test
void spawnReturnsImmediatelyWhenGateIsDisabled() {
FakeHerdr herdr = new FakeHerdr();
// The default 6-arg constructor has spawnReadyTimeoutMs=0 (gate disabled).
ClaudeCodeLauncher svc = service(herdr, List.of("ccs", "ltms-local"), null);
PeerHandle handle = svc.spawn(new SpawnRequest(null, null, null));
assertNotNull(handle, "spawn returns a handle when the gate is disabled");
assertFalse(herdr.called("agent.get"),
"agent.get is never called when the gate is disabled (no polling)");
}
@Test
void spawnGateRespectsZeroTimeoutEvenWithFullConstructor() {
FakeHerdr herdr = new FakeHerdr();
long[] clock = {0};
// Explicit zero timeout with the full testability constructor — should
// skip polling entirely, just like the legacy default path.
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(
new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")),
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
0, 300, () -> clock[0], () -> clock[0] += 1);
PeerHandle handle = svc.spawn(new SpawnRequest(null, null, null));
assertNotNull(handle, "spawn still succeeds with zero timeout");
assertEquals(0, paneCloseCount(herdr, handle.id()),
"no orphan pane close from the gate path");
}
}