CB-113: reliable worker readiness gate + submission nudge
Delegating right after spawn failed: herdr reports 'idle' during the worker's
boot, so the injector delivered into a not-ready TUI (paste lost) and wedged the
worker. Two fixes:
- Readiness gate: a worker is 'available' only once its Claude connects the bridge
MCP (the daemon observes it via peer-PID->terminal). WorkerPresence tracks it; the
injector holds the first delivery until present, so it never pastes into the boot
window. Exposed as 'ready' on GET /sessions/{id}/status.
- Submission nudge: the Enter accompanying a delivery can race the paste (esp. right
as the TUI becomes ready), leaving text unsubmitted. While a delivered message
stays idle (not picked up), the injector re-sends Enter each poll until the worker
starts (WORKING) or the grace expires.
Validated live: spawn + immediate delegate now holds during boot (ready=false),
delivers on MCP-connect, re-nudges Enter, worker replies. 102 tests green.
This commit is contained in:
@@ -9,6 +9,7 @@ import dev.ltms.bridged.herdr.WorkspaceControl;
|
||||
import dev.ltms.bridged.inject.CompletionResolver;
|
||||
import dev.ltms.bridged.inject.Injector;
|
||||
import dev.ltms.bridged.inject.StatusPoller;
|
||||
import dev.ltms.bridged.inject.WorkerPresence;
|
||||
import dev.ltms.bridged.mcp.BridgeMcp;
|
||||
import dev.ltms.bridged.mcp.ConnectionIdentity;
|
||||
import dev.ltms.bridged.mcp.LsofPeerPidLookup;
|
||||
@@ -60,7 +61,9 @@ public final class Bridged {
|
||||
// CB-106: a confirmed turn completion resolves a blocked send whose worker never replied.
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver completion = new CompletionResolver(agents, rendezvous);
|
||||
Injector injector = new Injector(agents, completion);
|
||||
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
|
||||
WorkerPresence presence = new WorkerPresence();
|
||||
Injector injector = new Injector(agents, completion, presence::isPresent);
|
||||
StatusPoller poller = new StatusPoller(agents, injector, INJECT_POLL_MILLIS);
|
||||
poller.start();
|
||||
Runtime.getRuntime().addShutdownHook(new Thread(poller::stop));
|
||||
@@ -72,9 +75,9 @@ public final class Bridged {
|
||||
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
|
||||
ConnectionIdentity identity = new ConnectionIdentity(
|
||||
new PaneLocator(herdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
|
||||
BridgeMcp mcp = new BridgeMcp(messages, rendezvous, workers, identity);
|
||||
BridgeMcp mcp = new BridgeMcp(messages, rendezvous, workers, identity, presence);
|
||||
Runtime.getRuntime().addShutdownHook(new Thread(mcp::close));
|
||||
Javalin app = new BridgedApp(herdr, workers, messages, rendezvous, mcp.servlet()).build();
|
||||
Javalin app = new BridgedApp(herdr, workers, messages, rendezvous, presence, mcp.servlet()).build();
|
||||
app.start(cfg.bind().host(), cfg.bind().port());
|
||||
log.info("bridged listening on {}:{}, herdr socket {}",
|
||||
cfg.bind().host(), cfg.bind().port(), socket);
|
||||
|
||||
@@ -86,6 +86,15 @@ public final class AgentControl {
|
||||
herdr.call("agent.send", Map.of("target", target, "text", SUBMIT_KEY));
|
||||
}
|
||||
|
||||
/**
|
||||
* Re-send the submit keystroke (Enter) to {@code target}. The Enter that accompanies a delivery
|
||||
* can race the paste — especially right as the worker's TUI becomes interactive — leaving the
|
||||
* text unsubmitted; the injector nudges it with this until the worker actually picks up (CB-113).
|
||||
*/
|
||||
public void submit(String target) {
|
||||
herdr.call("agent.send", Map.of("target", target, "text", SUBMIT_KEY));
|
||||
}
|
||||
|
||||
/**
|
||||
* Read an agent's terminal.
|
||||
*
|
||||
|
||||
@@ -12,6 +12,7 @@ import java.util.List;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.function.Predicate;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
/**
|
||||
@@ -65,17 +66,29 @@ public final class Injector {
|
||||
|
||||
private final AgentControl agents;
|
||||
private final TurnListener turnListener;
|
||||
private final Predicate<String> ready; // CB-113: a target is deliverable only when available
|
||||
private final ConcurrentHashMap<String, Target> targets = new ConcurrentHashMap<>();
|
||||
|
||||
/** Delivery only; completion signalling is a no-op. */
|
||||
/** Delivery only; completion signalling is a no-op and every target is treated as available. */
|
||||
public Injector(AgentControl agents) {
|
||||
this(agents, TurnListener.NOOP);
|
||||
}
|
||||
|
||||
/** Delivery plus turn-completion signalling to {@code turnListener} (CB-106). */
|
||||
/** Delivery plus turn-completion signalling (CB-106); every target is treated as available. */
|
||||
public Injector(AgentControl agents, TurnListener turnListener) {
|
||||
this(agents, turnListener, _ -> true);
|
||||
}
|
||||
|
||||
/**
|
||||
* Delivery, completion signalling (CB-106), and a readiness gate (CB-113): a message is delivered
|
||||
* only when {@code ready} accepts the target — i.e. the worker's Claude has connected the bridge
|
||||
* MCP. This holds the first delivery out of the worker's boot window, where herdr already reports
|
||||
* {@code idle} but the TUI would drop an injected paste.
|
||||
*/
|
||||
public Injector(AgentControl agents, TurnListener turnListener, Predicate<String> ready) {
|
||||
this.agents = agents;
|
||||
this.turnListener = turnListener;
|
||||
this.ready = ready;
|
||||
}
|
||||
|
||||
/** A pending message and the future that completes when it has been delivered. */
|
||||
@@ -130,6 +143,7 @@ public final class Injector {
|
||||
RuntimeException sendError = null;
|
||||
boolean turnCompleted = false;
|
||||
boolean turnFailed = false;
|
||||
boolean resubmit = false;
|
||||
synchronized (t) {
|
||||
if (status == AgentStatus.WORKING) {
|
||||
// Definitive pickup: the worker is busy on our last message, and (if a delivery is
|
||||
@@ -140,15 +154,22 @@ public final class Injector {
|
||||
if (t.awaitingCompletion) t.turnObserved = true;
|
||||
} else if (status.injectable()) { // IDLE or BLOCKED
|
||||
t.unknownSinceTurn = 0;
|
||||
if (t.awaitingPickup && ++t.injectableSincePickup >= PICKUP_GRACE_POLLS) {
|
||||
// Pickup edge was never sampled (turn faster than the poll, or status lag).
|
||||
// Release the latch rather than wedge — and give up on synthesizing a completion
|
||||
// for this message, since without a confirmed `working` we cannot trust that a
|
||||
// task-processing turn actually ran.
|
||||
t.awaitingPickup = false;
|
||||
t.injectableSincePickup = 0;
|
||||
t.awaitingCompletion = false;
|
||||
t.turnObserved = false;
|
||||
if (t.awaitingPickup) {
|
||||
if (++t.injectableSincePickup >= PICKUP_GRACE_POLLS) {
|
||||
// Pickup edge was never sampled (turn faster than the poll, or status lag).
|
||||
// Release the latch rather than wedge — and give up on synthesizing a
|
||||
// completion for this message, since without a confirmed `working` we cannot
|
||||
// trust that a task-processing turn actually ran.
|
||||
t.awaitingPickup = false;
|
||||
t.injectableSincePickup = 0;
|
||||
t.awaitingCompletion = false;
|
||||
t.turnObserved = false;
|
||||
} else {
|
||||
// Delivered but still idle → the worker hasn't picked it up; the submit
|
||||
// keystroke likely raced the paste (esp. right as the TUI became ready).
|
||||
// Re-nudge Enter (CB-113) until the worker starts (WORKING) or the grace ends.
|
||||
resubmit = true;
|
||||
}
|
||||
}
|
||||
if (!t.awaitingPickup) {
|
||||
// A confirmed turn (a `working` sample was seen) that has now returned to idle is
|
||||
@@ -159,10 +180,11 @@ public final class Injector {
|
||||
turnCompleted = true;
|
||||
}
|
||||
// Deliver the next queued message only once the prior turn is fully settled, so a
|
||||
// completion is never confused with the pickup of the following message.
|
||||
// completion is never confused with the pickup of the following message — and only
|
||||
// once the worker is available (CB-113), so we never paste into its boot window.
|
||||
if (!t.awaitingCompletion) {
|
||||
Pending p = t.queue.peek();
|
||||
if (p != null) {
|
||||
if (p != null && ready.test(target)) {
|
||||
try {
|
||||
agents.send(target, p.text());
|
||||
t.queue.poll();
|
||||
@@ -205,8 +227,15 @@ public final class Injector {
|
||||
}
|
||||
}
|
||||
|
||||
// Fire listeners after releasing the monitor so a continuation (which may call herdr) never
|
||||
// runs on the poller thread while it holds the target lock.
|
||||
// Fire listeners / herdr calls after releasing the monitor so nothing runs on the poller
|
||||
// thread while it holds the target lock.
|
||||
if (resubmit) {
|
||||
try {
|
||||
agents.submit(target); // nudge a raced Enter so the pending paste submits
|
||||
} catch (RuntimeException e) {
|
||||
log.debug("resubmit to {} failed (will retry next poll): {}", target, e.getMessage());
|
||||
}
|
||||
}
|
||||
if (turnCompleted) {
|
||||
turnListener.onTurnComplete(target);
|
||||
}
|
||||
|
||||
@@ -0,0 +1,38 @@
|
||||
package dev.ltms.bridged.inject;
|
||||
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.Set;
|
||||
|
||||
/**
|
||||
* Tracks which workers are <em>available</em> — their Claude has booted and connected its MCP client
|
||||
* to the bridge (CB-113). This is the reliable readiness signal, unlike herdr's {@code agent_status},
|
||||
* which reports {@code idle} for a worker whose Claude is still booting. Delivering into that boot
|
||||
* window pastes into a not-yet-ready TUI (the text is lost) and wedges the worker's delivery state,
|
||||
* so the {@link Injector} holds the first delivery until the worker is present here.
|
||||
*
|
||||
* <p>Populated from the MCP transport: any MCP request whose connection resolves to a worker terminal
|
||||
* marks that worker present (its {@code initialize} is the first such contact). A worker that never
|
||||
* mounts the bridge MCP is never marked present — its sends stay queued until they time out, which is
|
||||
* correct (it could not have replied anyway).
|
||||
*/
|
||||
public final class WorkerPresence {
|
||||
|
||||
private final Set<String> present = ConcurrentHashMap.newKeySet();
|
||||
|
||||
/** Record that {@code terminal}'s worker has connected its MCP client (is available). */
|
||||
public void markPresent(String terminal) {
|
||||
if (terminal != null && !terminal.isBlank()) {
|
||||
present.add(terminal);
|
||||
}
|
||||
}
|
||||
|
||||
/** Whether {@code terminal}'s worker is available (has been seen on the bridge MCP). */
|
||||
public boolean isPresent(String terminal) {
|
||||
return present.contains(terminal);
|
||||
}
|
||||
|
||||
/** Forget a torn-down worker so its terminal id does not linger as "present". */
|
||||
public void forget(String terminal) {
|
||||
present.remove(terminal);
|
||||
}
|
||||
}
|
||||
@@ -2,6 +2,7 @@ package dev.ltms.bridged.mcp;
|
||||
|
||||
import dev.ltms.bridged.guard.GuardException;
|
||||
import dev.ltms.bridged.herdr.Agent;
|
||||
import dev.ltms.bridged.inject.WorkerPresence;
|
||||
import dev.ltms.bridged.herdr.HerdrException;
|
||||
import dev.ltms.bridged.msg.MessageService;
|
||||
import dev.ltms.bridged.msg.Rendezvous;
|
||||
@@ -50,16 +51,18 @@ public final class BridgeMcp {
|
||||
private final McpSyncServer server;
|
||||
|
||||
public BridgeMcp(MessageService messages, Rendezvous rendezvous, WorkerService workers,
|
||||
ConnectionIdentity identity) {
|
||||
ConnectionIdentity identity, WorkerPresence presence) {
|
||||
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
|
||||
this.transport = HttpServletStreamableServerTransportProvider.builder()
|
||||
.jsonMapper(json)
|
||||
.mcpEndpoint("/mcp")
|
||||
// Resolve the caller from the connection (peer PID → herdr pane) in one lookup: the
|
||||
// worker terminal for bridge_reply (no spoofable arg), and the PID so bridge_spawn can
|
||||
// inherit the primary's cwd (CB-112).
|
||||
// inherit the primary's cwd (CB-112). Any contact from a worker marks it available
|
||||
// (CB-113) — its MCP initialize is the reliable "the agent is up" signal.
|
||||
.contextExtractor(req -> {
|
||||
ConnectionIdentity.Caller c = identity.resolve(req.getRemoteAddr(), req.getRemotePort());
|
||||
presence.markPresent(c.terminal()); // no-op for the primary (null terminal)
|
||||
return McpTransportContext.create(Map.of(
|
||||
CALLER_TERMINAL, orEmpty(c.terminal()),
|
||||
CALLER_PID, Long.toString(c.pid())));
|
||||
|
||||
@@ -6,6 +6,7 @@ import dev.ltms.bridged.guard.GuardException;
|
||||
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.msg.MessageService;
|
||||
import dev.ltms.bridged.msg.Rendezvous;
|
||||
import dev.ltms.bridged.worker.WorkerService;
|
||||
@@ -38,15 +39,17 @@ public final class BridgedApp {
|
||||
private final WorkerService workers;
|
||||
private final MessageService messages;
|
||||
private final Rendezvous rendezvous;
|
||||
private final WorkerPresence presence; // CB-113: which workers are MCP-connected (available)
|
||||
private final HttpServlet mcpServlet; // MCP Streamable-HTTP endpoint, mounted at /mcp (nullable)
|
||||
private final ObjectMapper mapper = new ObjectMapper();
|
||||
|
||||
public BridgedApp(HerdrClient herdr, WorkerService workers, MessageService messages,
|
||||
Rendezvous rendezvous, HttpServlet mcpServlet) {
|
||||
Rendezvous rendezvous, WorkerPresence presence, HttpServlet mcpServlet) {
|
||||
this.herdr = herdr;
|
||||
this.workers = workers;
|
||||
this.messages = messages;
|
||||
this.rendezvous = rendezvous;
|
||||
this.presence = presence;
|
||||
this.mcpServlet = mcpServlet;
|
||||
}
|
||||
|
||||
@@ -240,13 +243,19 @@ public final class BridgedApp {
|
||||
}
|
||||
}
|
||||
|
||||
/** Live lifecycle status of a worker (MCP `bridge_status` wraps this in CB-105). */
|
||||
/**
|
||||
* Live lifecycle status of a worker (MCP `bridge_status` wraps this in CB-105), plus its
|
||||
* <em>readiness</em> (CB-113): {@code ready} is true once the worker's Claude has connected the
|
||||
* bridge MCP — the reliable "available to receive a task" signal, unlike bare {@code idle}, which
|
||||
* is also true during boot.
|
||||
*/
|
||||
private void sessionStatus(Context ctx) {
|
||||
String id = ctx.pathParam("id");
|
||||
try {
|
||||
ctx.status(200).json(Map.of(
|
||||
"sessionId", id,
|
||||
"status", messages.status(id).name().toLowerCase()));
|
||||
"status", messages.status(id).name().toLowerCase(),
|
||||
"ready", presence.isPresent(id)));
|
||||
} catch (HerdrException e) {
|
||||
herdrError(ctx, e);
|
||||
}
|
||||
|
||||
@@ -51,6 +51,47 @@ class InjectorTest {
|
||||
assertEquals(List.of("hello"), sent());
|
||||
}
|
||||
|
||||
@Test
|
||||
void holdsDeliveryUntilTheWorkerIsAvailable() {
|
||||
// CB-113: idle alone is not enough — hold until the worker's MCP is connected (ready).
|
||||
java.util.Set<String> ready = new java.util.HashSet<>();
|
||||
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, ready::contains);
|
||||
inj.enqueue(T, "task");
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE); // idle but not yet available → held out of the boot window
|
||||
assertEquals(List.of(), sent(), "must not deliver into a not-yet-available worker");
|
||||
|
||||
ready.add(T); // the worker's Claude connects the bridge MCP
|
||||
inj.onStatus(T, AgentStatus.IDLE);
|
||||
assertEquals(List.of("task"), sent(), "delivers once the worker is available");
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private long enterKeystrokes() {
|
||||
return herdr.calls.stream()
|
||||
.filter(c -> c.method().equals("agent.send"))
|
||||
.filter(c -> "\r".equals(((Map<String, Object>) c.params()).get("text")))
|
||||
.count();
|
||||
}
|
||||
|
||||
@Test
|
||||
void resubmitsEnterWhenADeliveredMessageIsNotPickedUp() {
|
||||
// CB-113: the Enter at delivery can race the paste; while the worker stays idle (not picked
|
||||
// up), the injector re-nudges Enter so the pending paste submits.
|
||||
injector.enqueue(T, "task");
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver: paste + one Enter
|
||||
long afterDeliver = enterKeystrokes();
|
||||
|
||||
injector.onStatus(T, AgentStatus.IDLE); // still idle → re-nudge Enter
|
||||
injector.onStatus(T, AgentStatus.IDLE); // and again
|
||||
assertTrue(enterKeystrokes() > afterDeliver, "an unpicked-up delivery re-nudges Enter");
|
||||
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker finally starts
|
||||
long atPickup = enterKeystrokes();
|
||||
injector.onStatus(T, AgentStatus.WORKING);
|
||||
assertEquals(atPickup, enterKeystrokes(), "no more nudges once the worker has picked up");
|
||||
}
|
||||
|
||||
@Test
|
||||
void holdsWhileWorkingThenDeliversOnIdle() {
|
||||
injector.enqueue(T, "later");
|
||||
|
||||
@@ -0,0 +1,33 @@
|
||||
package dev.ltms.bridged.inject;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/** The CB-113 worker-availability registry. */
|
||||
class WorkerPresenceTest {
|
||||
|
||||
@Test
|
||||
void tracksPresenceAndForgets() {
|
||||
WorkerPresence p = new WorkerPresence();
|
||||
assertFalse(p.isPresent("term_a"), "unseen worker is not available");
|
||||
|
||||
p.markPresent("term_a");
|
||||
assertTrue(p.isPresent("term_a"), "a worker seen on the MCP is available");
|
||||
|
||||
p.forget("term_a");
|
||||
assertFalse(p.isPresent("term_a"), "a torn-down worker is no longer available");
|
||||
}
|
||||
|
||||
@Test
|
||||
void nullOrBlankMarkIsANoOp() {
|
||||
WorkerPresence p = new WorkerPresence();
|
||||
assertDoesNotThrow(() -> {
|
||||
p.markPresent(null);
|
||||
p.markPresent(" ");
|
||||
});
|
||||
assertFalse(p.isPresent(""), "blank/null contacts (the primary) are never present");
|
||||
}
|
||||
}
|
||||
@@ -9,6 +9,7 @@ import dev.ltms.bridged.herdr.FakeHerdr;
|
||||
import dev.ltms.bridged.herdr.WorkspaceControl;
|
||||
import dev.ltms.bridged.inject.Injector;
|
||||
import dev.ltms.bridged.inject.StatusPoller;
|
||||
import dev.ltms.bridged.inject.WorkerPresence;
|
||||
import dev.ltms.bridged.msg.MessageService;
|
||||
import dev.ltms.bridged.msg.Rendezvous;
|
||||
import dev.ltms.bridged.worker.WorkerService;
|
||||
@@ -36,6 +37,7 @@ class BridgedAppTest {
|
||||
|
||||
private final ObjectMapper mapper = new ObjectMapper();
|
||||
private final HttpClient http = HttpClient.newHttpClient();
|
||||
private final WorkerPresence presence = new WorkerPresence();
|
||||
private Javalin app;
|
||||
private StatusPoller poller;
|
||||
|
||||
@@ -63,7 +65,8 @@ class BridgedAppTest {
|
||||
poller.start();
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
MessageService messages = new MessageService(agents, injector, rendezvous);
|
||||
app = new BridgedApp(herdr, workers, messages, rendezvous, null).build().start("127.0.0.1", 0);
|
||||
app = new BridgedApp(herdr, workers, messages, rendezvous, this.presence, null)
|
||||
.build().start("127.0.0.1", 0);
|
||||
return app.port();
|
||||
}
|
||||
|
||||
@@ -380,6 +383,16 @@ class BridgedAppTest {
|
||||
assertEquals("blocked", mapper.readTree(res.body()).get("status").asText());
|
||||
}
|
||||
|
||||
@Test
|
||||
void sessionStatusReportsReadinessFromMcpPresence() throws Exception {
|
||||
int port = start(new FakeHerdr(), "http://gx00.gw:8000", Set.of("gx00.gw"));
|
||||
// Not yet seen on the bridge MCP → not ready.
|
||||
assertFalse(mapper.readTree(req(port, "GET", "/sessions/term_a/status").body()).get("ready").asBoolean());
|
||||
// Worker connects its MCP client → available.
|
||||
presence.markPresent("term_a");
|
||||
assertTrue(mapper.readTree(req(port, "GET", "/sessions/term_a/status").body()).get("ready").asBoolean());
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopWorkerInPanePlacementClosesOnlyThePane() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
Reference in New Issue
Block a user