From 16e17b32adb9a9825f6dfb3da6f2f39be10e827c Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 06:06:07 +0200 Subject: [PATCH] CB-568: preserve dropped turn causes --- .../main/java/dev/ltms/bridged/Bridged.java | 6 ++++ .../bridged/inject/CompletionResolver.java | 35 +++++++++++++------ .../dev/ltms/bridged/inject/Injector.java | 6 ++-- .../dev/ltms/bridged/inject/TurnListener.java | 8 +++++ .../dev/ltms/bridged/inject/InjectorTest.java | 24 +++++++++++++ .../ltms/bridged/msg/MessageServiceTest.java | 30 ++++++++++++++++ 6 files changed, 95 insertions(+), 14 deletions(-) diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 85d057c..f0e2a85 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -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); diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java index 1d1be03..e504b96 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java @@ -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); 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 b6848e2..6144c6b 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java @@ -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()); } } diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java b/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java index 6ac5062..68e377d 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java @@ -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 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 990b7a0..61a10b9 100644 --- a/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java @@ -259,6 +259,7 @@ class InjectorTest { private static final class Captor implements TurnListener { final List completed = new ArrayList<>(); final List failed = new ArrayList<>(); + final List 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 delivered = inj.enqueue(T, "delivered"); + CompletableFuture 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 diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java index 5668890..5743816 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -120,6 +120,36 @@ class MessageServiceTest { assertFalse(reply.completed()); } + @Test + void droppedQueuedAndDeliveredTurnsExposeTheRealCauseExactlyOnce() throws Exception { + CompletableFuture first = sendAsync(); + awaitWaiting(); + injector.onStatus(T, AgentStatus.IDLE); // first delivery + injector.onStatus(T, AgentStatus.WORKING); // first turn in flight + + CompletableFuture queued = injector.enqueue(T, "second task"); + CompletableFuture 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 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