Merge CB-568: a dropped send reports the real cause
CI / build (push) Successful in 55s
CI / contract (push) Successful in 1m0s

Injector.drop knew the precise cause (herdr agent_not_found) but the
sender was told only 'worker unreachable or stuck', so a lead could not
tell a dead pane from a stalled model.

TurnListener.onTurnFailed gains a reason, defaulting to the old one-arg
form. CompletionResolver prefers that reason, then the pane scrape, then
the old fixed text.

drop now fires onTurnFailed unconditionally. That is the substantive
fix: the sender blocks on the rendezvous waiter, not on the delivered
future, so failing delivered() alone never woke it and a queued send sat
until its timeout.
This commit is contained in:
Dai Ha
2026-08-15 06:08:41 +02:00
6 changed files with 95 additions and 14 deletions
@@ -285,6 +285,12 @@ public final class Bridged {
completion.onTurnFailed(target);
sessions.onTurnFailed(target);
}
@Override
public void onTurnFailed(String target, String reason) {
completion.onTurnFailed(target, reason);
sessions.onTurnFailed(target);
}
};
Injector injector = new Injector(agents, turnListener, deliverableTo(presence, leads),
presence::forget);
@@ -137,7 +137,13 @@ public final class CompletionResolver implements TurnListener {
@Override
public void onTurnFailed(String target) {
InFlight turn = inFlight.get(target);
Thread.ofVirtual().name("turn-failed-" + target).start(() -> fail(target, turn));
Thread.ofVirtual().name("turn-failed-" + target).start(() -> fail(target, turn, null));
}
@Override
public void onTurnFailed(String target, String reason) {
InFlight turn = inFlight.get(target);
Thread.ofVirtual().name("turn-failed-" + target).start(() -> fail(target, turn, reason));
}
/** Synchronous resolve (the unit-testable core of {@link #onTurnComplete}). */
@@ -192,6 +198,11 @@ public final class CompletionResolver implements TurnListener {
/** Synchronous fail (the unit-testable core of {@link #onTurnFailed}). */
void fail(String target, InFlight turn) {
fail(target, turn, null);
}
/** Synchronous fail with an optional reason supplied by a dropped worker queue. */
void fail(String target, InFlight turn, String explicitReason) {
// A never-delivered readiness failure has no in-flight record but still has a blocked send;
// fall back to the currently-registered waiter (unambiguous — that send never completed, so
// no next turn exists to confuse it with).
@@ -201,16 +212,18 @@ public final class CompletionResolver implements TurnListener {
inFlight.remove(target, turn); // nobody blocked on this worker — nothing to fail
return;
}
String reason;
try {
reason = clip(agents.read(target, SCRAPE_SOURCE));
} catch (RuntimeException e) {
reason = "";
}
if (reason.isBlank()) {
// No screen to scrape — either the worker is stuck (CB-109) or gone (CB-110).
reason = "worker did not reply; its turn ended in an unrecoverable state "
+ "(worker unreachable or stuck)";
String reason = explicitReason;
if (reason == null || reason.isBlank()) {
try {
reason = clip(agents.read(target, SCRAPE_SOURCE));
} catch (RuntimeException e) {
reason = "";
}
if (reason.isBlank()) {
// No screen to scrape — either the worker is stuck (CB-109) or gone (CB-110).
reason = "worker did not reply; its turn ended in an unrecoverable state "
+ "(worker unreachable or stuck)";
}
}
if (rendezvous.resolveFailure(waiter, reason)) {
inFlight.remove(target, turn);
@@ -410,8 +410,8 @@ public final class Injector {
for (Pending p : pending) {
p.delivered().completeExceptionally(cause);
}
if (hadDeliveredTurn) {
turnListener.onTurnFailed(target);
}
// A queued send has no in-flight record, while a delivered turn does. CompletionResolver
// handles both forms and resolves its waiter at most once.
turnListener.onTurnFailed(target, cause.getMessage());
}
}
@@ -41,6 +41,14 @@ public interface TurnListener {
default void onTurnFailed(String target) {
}
/**
* As {@link #onTurnFailed(String)}, carrying the reason a worker became unreachable. The default
* keeps existing listeners working while allowing the completion resolver to report a useful cause.
*/
default void onTurnFailed(String target, String reason) {
onTurnFailed(target);
}
/**
* A message was just delivered into {@code target}'s pane (CB-115). Fired so the completion
* resolver can snapshot the pane's pre-turn content: a later {@link #onTurnComplete} whose
@@ -259,6 +259,7 @@ class InjectorTest {
private static final class Captor implements TurnListener {
final List<String> completed = new ArrayList<>();
final List<String> failed = new ArrayList<>();
final List<String> failureReasons = new ArrayList<>();
@Override
public void onTurnComplete(String target) {
@@ -269,6 +270,12 @@ class InjectorTest {
public void onTurnFailed(String target) {
failed.add(target);
}
@Override
public void onTurnFailed(String target, String reason) {
failed.add(target);
failureReasons.add(reason);
}
}
// ~30s of unknown at the 250ms prod poll interval; enough onStatus samples to trip the stall.
@@ -322,6 +329,23 @@ class InjectorTest {
assertTrue(f.isCompletedExceptionally(), "queued waiters unblock when the worker vanishes");
}
@Test
void dropPassesTheRealCauseForQueuedAndDeliveredWork() {
Captor cap = new Captor();
Injector inj = new Injector(new AgentControl(herdr), cap);
CompletableFuture<Void> delivered = inj.enqueue(T, "delivered");
CompletableFuture<Void> queued = inj.enqueue(T, "queued");
inj.onStatus(T, AgentStatus.IDLE); // deliver the first message
inj.onStatus(T, AgentStatus.WORKING); // its turn is now in flight; one remains queued
inj.drop(T, new HerdrException("agent target sol not found", "agent_not_found", null));
assertEquals(List.of(T), cap.failed, "drop signals one turn failure for both affected states");
assertEquals(List.of("agent target sol not found"), cap.failureReasons);
assertTrue(delivered.isDone(), "the delivered future has already completed");
assertTrue(queued.isCompletedExceptionally(), "the queued future fails with the drop cause");
}
@Test
void dropFailsTheTurnOfADeliveredMessageWhenTheWorkerVanishes() {
// CB-110: the message was delivered (no longer queued), so failing queued waiters alone would
@@ -120,6 +120,36 @@ class MessageServiceTest {
assertFalse(reply.completed());
}
@Test
void droppedQueuedAndDeliveredTurnsExposeTheRealCauseExactlyOnce() throws Exception {
CompletableFuture<MessageService.Reply> first = sendAsync();
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // first delivery
injector.onStatus(T, AgentStatus.WORKING); // first turn in flight
CompletableFuture<Void> queued = injector.enqueue(T, "second task");
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(T);
injector.drop(T, new HerdrException("agent target sol not found", "agent_not_found", null));
MessageService.Reply reply = first.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.WORKER_FAILED, reply.outcome());
assertEquals("agent target sol not found", reply.text());
assertTrue(queued.isCompletedExceptionally(), "the queued delivery future also fails");
assertFalse(rendezvous.resolveFailure(waiter, "second failure"), "the waiter fails exactly once");
}
@Test
void noDropReasonKeepsTheExistingFallbackText() throws Exception {
herdr.readText("");
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.open(T);
completion.onTurnFailed(T);
Rendezvous.Resolution resolution = waiter.get(2, TimeUnit.SECONDS);
assertEquals("worker did not reply; its turn ended in an unrecoverable state "
+ "(worker unreachable or stuck)", resolution.text());
}
// --- bridge_ask reverse rendezvous (CB-205) ------------------------------------------------
@Test