From 2fb46f670f7a3c2c939f8a67c86ffcdf21a9cfea Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 16 Jul 2026 05:01:58 +0200 Subject: [PATCH] CB-114: readiness-gate timeout + presence cleanup (delegated review findings) An off-sub worker's review of CB-113 (delegated through the bridge) surfaced two real gaps in the readiness gate: - A worker that herdr reports idle but whose Claude never connects the bridge MCP (crashed during boot, or wedged on a startup prompt) left its message queued forever: ready.test() never passed, the target was polled indefinitely, and the caller's future never completed (async waiter hung for the full 30-min window). Injector now counts injectable-but-not-ready samples and, after a ~60s grace (READINESS_GRACE_POLLS, deliberately longer than the UNKNOWN stall grace since a first boot is slower than an in-turn blip), fails the queued messages, fires onTurnFailed so blocking/async waiters resolve WORKER_FAILED, and reclaims the target. Mirrors the CB-109 UNKNOWN-stall path. - WorkerPresence.forget had no caller, so a worker's readiness lingered past its life. Injector now clears presence via a forget callback on drop() (pane crash) and on the readiness timeout. +3 InjectorTest cases (never-ready failure, ready-within-grace delivery, drop clears presence). 105 tests green. --- .../main/java/dev/ltms/bridged/Bridged.java | 2 +- .../dev/ltms/bridged/inject/Injector.java | 54 +++++++++++++++++++ .../dev/ltms/bridged/inject/InjectorTest.java | 52 ++++++++++++++++++ 3 files changed, 107 insertions(+), 1 deletion(-) diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 63ebf48..808f6fe 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -63,7 +63,7 @@ public final class Bridged { CompletionResolver completion = new CompletionResolver(agents, rendezvous); // 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); + Injector injector = new Injector(agents, completion, presence::isPresent, presence::forget); StatusPoller poller = new StatusPoller(agents, injector, INJECT_POLL_MILLIS); poller.start(); Runtime.getRuntime().addShutdownHook(new Thread(poller::stop)); diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java b/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java index 0f68730..09b4056 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java @@ -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.Consumer; import java.util.function.Predicate; import java.util.stream.Collectors; @@ -64,9 +65,23 @@ public final class Injector { */ private static final int TURN_STALL_GRACE_POLLS = 120; + /** + * How many consecutive injectable samples a queued-but-undelivered message may wait on the + * {@link #ready} gate before we give up and fail it (CB-114). The gate holds a message out of a + * worker's boot window (herdr reports {@code idle} while its Claude is still starting), but a + * worker whose Claude crashes during boot — or never connects the bridge MCP — stays "idle and + * not ready" forever: {@link #ready} never accepts it, the message is never delivered, and the + * target would be polled indefinitely with its caller's future never completing. After this + * grace the queued messages are failed and the target released. At the 250ms poll interval this + * is ~60s — deliberately longer than {@link #TURN_STALL_GRACE_POLLS}, since a first boot (spawn + * + model load + MCP connect) legitimately takes longer than an in-turn detection blip. + */ + private static final int READINESS_GRACE_POLLS = 240; + private final AgentControl agents; private final TurnListener turnListener; private final Predicate ready; // CB-113: a target is deliverable only when available + private final Consumer forget; // CB-114: clear a gone worker's readiness/presence private final ConcurrentHashMap targets = new ConcurrentHashMap<>(); /** Delivery only; completion signalling is a no-op and every target is treated as available. */ @@ -86,9 +101,22 @@ public final class Injector { * {@code idle} but the TUI would drop an injected paste. */ public Injector(AgentControl agents, TurnListener turnListener, Predicate ready) { + this(agents, turnListener, ready, _ -> { + }); + } + + /** + * Delivery, completion signalling (CB-106), a readiness gate (CB-113), and readiness cleanup + * (CB-114): {@code forget} is invoked with a target when its worker is gone — dropped + * (pane crash) or timed out on the readiness gate — so its stale presence/readiness is cleared + * and does not linger past the worker's life. + */ + public Injector(AgentControl agents, TurnListener turnListener, Predicate ready, + Consumer forget) { this.agents = agents; this.turnListener = turnListener; this.ready = ready; + this.forget = forget; } /** A pending message and the future that completes when it has been delivered. */ @@ -103,6 +131,7 @@ public final class Injector { boolean awaitingCompletion; // a delivered message's turn is not yet known-complete boolean turnObserved; // saw a real `working` sample since that delivery (turn ran) int unknownSinceTurn; // consecutive `unknown` samples while a delegation is outstanding (CB-109) + int notReadySincePoll; // consecutive injectable samples a queued message waited on the readiness gate (CB-114) synchronized void add(Pending p) { queue.add(p); @@ -144,6 +173,7 @@ public final class Injector { boolean turnCompleted = false; boolean turnFailed = false; boolean resubmit = false; + List notReady = null; // queued messages failed because the worker never became ready synchronized (t) { if (status == AgentStatus.WORKING) { // Definitive pickup: the worker is busy on our last message, and (if a delivery is @@ -151,6 +181,7 @@ public final class Injector { t.awaitingPickup = false; t.injectableSincePickup = 0; t.unknownSinceTurn = 0; + t.notReadySincePoll = 0; if (t.awaitingCompletion) t.turnObserved = true; } else if (status.injectable()) { // IDLE or BLOCKED t.unknownSinceTurn = 0; @@ -185,6 +216,7 @@ public final class Injector { if (!t.awaitingCompletion) { Pending p = t.queue.peek(); if (p != null && ready.test(target)) { + t.notReadySincePoll = 0; try { agents.send(target, p.text()); t.queue.poll(); @@ -200,6 +232,15 @@ public final class Injector { sent = p; sendError = e; } + } else if (p != null && ++t.notReadySincePoll >= READINESS_GRACE_POLLS) { + // The worker has been idle-but-not-ready for the whole grace: its Claude + // never connected the bridge MCP (crashed during boot, or wedged on a + // startup prompt). The readiness gate would hold this message forever, so + // fail every queued message and release the target (CB-114) instead of + // polling it indefinitely with the caller's future never completing. + notReady = new ArrayList<>(t.queue); + t.queue.clear(); + t.notReadySincePoll = 0; } } } @@ -236,6 +277,18 @@ public final class Injector { log.debug("resubmit to {} failed (will retry next poll): {}", target, e.getMessage()); } } + if (notReady != null) { + // Worker never became available: forget its (never-set) readiness, unblock every queued + // caller, and route the awaiting send through the same failure path as a stalled turn so + // a blocking or async waiter resolves WORKER_FAILED rather than riding out the timeout. + forget.accept(target); + RuntimeException cause = new IllegalStateException( + target + " never became available (no bridge MCP connection within the boot window)"); + for (Pending p : notReady) { + p.delivered().completeExceptionally(cause); + } + turnListener.onTurnFailed(target); + } if (turnCompleted) { turnListener.onTurnComplete(target); } @@ -288,6 +341,7 @@ public final class Injector { t.awaitingCompletion = false; t.awaitingPickup = false; } + forget.accept(target); // the worker is gone — clear its readiness/presence too (CB-114) for (Pending p : pending) { p.delivered().completeExceptionally(cause); } diff --git a/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java b/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java index f886abd..2e81012 100644 --- a/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java @@ -308,6 +308,58 @@ class InjectorTest { assertEquals(List.of(T), cap.failed, "a delivery that vanishes before pickup still fails its send"); } + // ~60s of idle-but-not-ready at the 250ms prod poll interval; enough to trip the readiness grace. + private static final int READINESS_SAMPLES = 245; + + @Test + void failsAQueuedMessageWhoseWorkerNeverBecomesReady() { + // CB-114: herdr keeps reporting the worker idle, but its Claude never connects the bridge MCP, + // so the readiness gate never opens. The message must not be held (and the target polled) + // forever — after the grace it fails, the caller unblocks via the worker-failure path, the + // target is reclaimed, and the never-set presence is cleared. + Captor cap = new Captor(); + List forgotten = new ArrayList<>(); + Injector inj = new Injector(new AgentControl(herdr), cap, _ -> false, forgotten::add); + CompletableFuture f = inj.enqueue(T, "task"); + + for (int i = 0; i < READINESS_SAMPLES; i++) inj.onStatus(T, AgentStatus.IDLE); + + assertEquals(List.of(), sent(), "a never-ready worker is never delivered to"); + assertTrue(f.isCompletedExceptionally(), "the caller's future fails instead of hanging forever"); + assertEquals(List.of(T), cap.failed, "the awaiting send resolves through the worker-failure path"); + assertEquals(List.of(), cap.completed, "a never-ready worker is a failure, not a completion"); + assertEquals(List.of(T), forgotten, "the never-ready worker's presence is cleared"); + assertTrue(inj.activeTargets().isEmpty(), "the target is reclaimed, not polled forever"); + } + + @Test + void aWorkerThatBecomesReadyWithinTheGraceIsDeliveredNormally() { + // The readiness grace must not fail a worker that is merely slow to boot: once it becomes + // available before the grace elapses, delivery proceeds as usual (the counter resets). + Set ready = new java.util.HashSet<>(); + Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, ready::contains, _ -> { + }); + inj.enqueue(T, "task"); + + for (int i = 0; i < 100; i++) inj.onStatus(T, AgentStatus.IDLE); // still booting, well under grace + assertEquals(List.of(), sent()); + + ready.add(T); // MCP connects before the grace elapses + inj.onStatus(T, AgentStatus.IDLE); + assertEquals(List.of("task"), sent(), "a worker that connects within the grace is delivered to"); + } + + @Test + void dropClearsWorkerPresence() { + // CB-114 (finding #1): a vanished worker's readiness must be forgotten so a stale entry cannot + // linger past the worker's life (WorkerPresence.forget had no caller before this). + List forgotten = new ArrayList<>(); + Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> true, forgotten::add); + inj.enqueue(T, "orphan"); + inj.drop(T, new HerdrException("worker gone", "pane_not_found", null)); + assertEquals(List.of(T), forgotten, "drop clears the gone worker's presence"); + } + @Test void pollerDeliversToAnIdleWorker() throws Exception { // End-to-end through the poller: idle worker → message delivered without manual onStatus.