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.
This commit is contained in:
@@ -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));
|
||||
|
||||
@@ -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<String> ready; // CB-113: a target is deliverable only when available
|
||||
private final Consumer<String> forget; // CB-114: clear a gone worker's readiness/presence
|
||||
private final ConcurrentHashMap<String, Target> 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<String> 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<String> ready,
|
||||
Consumer<String> 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<Pending> 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);
|
||||
}
|
||||
|
||||
@@ -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<String> forgotten = new ArrayList<>();
|
||||
Injector inj = new Injector(new AgentControl(herdr), cap, _ -> false, forgotten::add);
|
||||
CompletableFuture<Void> 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<String> 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<String> 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.
|
||||
|
||||
Reference in New Issue
Block a user