CB-110: fail an in-flight delegation when its worker vanishes

Companion to CB-109. When a worker disappears mid-turn (pane crash → herdr
*_not_found), the status poller drops the target, which failed only *queued*
messages — a message already DELIVERED is out of the queue, so its send's
rendezvous was left hanging until the 30-min async timeout. Injector.drop now
fires onTurnFailed for a delivered-but-unresolved turn (awaitingCompletion), so
the send resolves as WORKER_FAILED. Reuses the CB-109 resolver path; the failure
reason is neutral to cover both wedge (stuck) and vanish (gone).
This commit is contained in:
Dai Ha
2026-07-15 15:23:17 +02:00
parent 2052929768
commit a629a7ee73
4 changed files with 47 additions and 4 deletions
@@ -89,7 +89,9 @@ public final class CompletionResolver implements TurnListener {
reason = "";
}
if (reason.isBlank()) {
reason = "worker turn ended in an unrecoverable state (herdr status stuck at unknown)";
// 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(target, reason)) {
log.debug("failed send to {} via turn-stall fallback", target);
@@ -240,20 +240,30 @@ public final class Injector {
}
/**
* Forget a target whose worker is gone, failing every still-queued message so awaiting
* callers unblock instead of hanging forever. Futures are completed after the monitor is
* released.
* Forget a target whose worker is gone, failing every still-queued message so awaiting callers
* unblock instead of hanging forever. If a message had already been <em>delivered</em> but its
* turn was not yet resolved (CB-110 — the worker vanished mid-turn, e.g. its pane crashed), fire
* {@link TurnListener#onTurnFailed} for it: a delivered message is no longer in the queue, so
* failing queued waiters alone would leave that send's rendezvous hanging until the async
* timeout. Futures and listeners are completed after the monitor is released.
*/
public void drop(String target, Throwable cause) {
Target t = targets.remove(target);
if (t == null) return;
List<Pending> pending;
boolean hadDeliveredTurn;
synchronized (t) {
pending = new ArrayList<>(t.queue);
t.queue.clear();
hadDeliveredTurn = t.awaitingCompletion;
t.awaitingCompletion = false;
t.awaitingPickup = false;
}
for (Pending p : pending) {
p.delivered().completeExceptionally(cause);
}
if (hadDeliveredTurn) {
turnListener.onTurnFailed(target);
}
}
}
@@ -239,6 +239,20 @@ class InjectorTest {
assertTrue(f.isCompletedExceptionally(), "queued waiters unblock when the worker vanishes");
}
@Test
void dropFailsTheTurnOfADeliveredMessageWhenTheWorkerVanishes() {
// CB-110: the message was delivered (no longer queued), so failing queued waiters alone would
// leave its send hanging. A vanished worker must fail that in-flight turn too.
Captor cap = new Captor();
Injector inj = new Injector(new AgentControl(herdr), cap);
inj.enqueue(T, "task");
inj.onStatus(T, AgentStatus.IDLE); // deliver
inj.onStatus(T, AgentStatus.WORKING); // turn running
inj.drop(T, new HerdrException("worker gone", "pane_not_found", null));
assertEquals(List.of(T), cap.failed, "a worker that vanishes mid-turn fails its in-flight send");
}
@Test
void pollerDeliversToAnIdleWorker() throws Exception {
// End-to-end through the poller: idle worker → message delivered without manual onStatus.
@@ -3,6 +3,7 @@ package dev.ltms.bridged.msg;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.inject.CompletionResolver;
import dev.ltms.bridged.inject.Injector;
import org.junit.jupiter.api.Test;
@@ -88,4 +89,20 @@ class MessageServiceTest {
assertFalse(reply.completed(), "a wedge is terminal but not a successful completion");
assertTrue(reply.text().contains("ENOTFOUND"), "the error screen is carried as the failure reason");
}
@Test
void aWorkerThatVanishesMidTurnResolvesTheSendAsFailed() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker starts the turn
// The worker's pane crashes — the poller sees a *_not_found and drops it (CB-110).
injector.drop(T, new HerdrException("worker gone", "pane_not_found", null));
MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.WORKER_FAILED, reply.outcome(),
"a delivered send whose worker vanishes fails instead of hanging to the timeout");
assertFalse(reply.completed());
}
}