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 2ae1d50..a444704 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java @@ -16,8 +16,12 @@ import org.slf4j.LoggerFactory; * — a fleet worker's own turns, or a send that already timed out, cost no herdr traffic. An explicit * {@code bridge_reply} that raced in first wins; {@link Rendezvous#resolveCompletion} is then a no-op. * - *

Wired as the {@link Injector}'s {@link TurnListener}; {@link #onTurnComplete} hands off to a - * virtual thread so the scrape's herdr round-trip never stalls the status poller. + *

It also handles the CB-109 stall signal ({@link #onTurnFailed}): a worker that ran a turn then + * wedged in an {@code unknown} state resolves the send as a failure (with the error screen as + * context) rather than leaving it to time out. + * + *

Wired as the {@link Injector}'s {@link TurnListener}; both handlers hand off to a virtual thread + * so the scrape's herdr round-trip never stalls the status poller. */ public final class CompletionResolver implements TurnListener { @@ -47,6 +51,11 @@ public final class CompletionResolver implements TurnListener { Thread.ofVirtual().name("completion-" + target).start(() -> resolve(target)); } + @Override + public void onTurnFailed(String target) { + Thread.ofVirtual().name("turn-failed-" + target).start(() -> fail(target)); + } + /** Synchronous resolve (the unit-testable core of {@link #onTurnComplete}). */ void resolve(String target) { if (!rendezvous.isWaiting(target)) { @@ -68,6 +77,25 @@ public final class CompletionResolver implements TurnListener { } } + /** Synchronous fail (the unit-testable core of {@link #onTurnFailed}). */ + void fail(String target) { + if (!rendezvous.isWaiting(target)) { + return; // nobody blocked on this worker — nothing to fail + } + String reason; + try { + reason = clip(agents.read(target, SCRAPE_SOURCE)); + } catch (RuntimeException e) { + reason = ""; + } + if (reason.isBlank()) { + reason = "worker turn ended in an unrecoverable state (herdr status stuck at unknown)"; + } + if (rendezvous.resolveFailure(target, reason)) { + log.debug("failed send to {} via turn-stall fallback", target); + } + } + private static String clip(String s) { if (s == null) return ""; String trimmed = s.strip(); 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 d0af9bc..ba91500 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java @@ -53,6 +53,16 @@ public final class Injector { */ private static final int PICKUP_GRACE_POLLS = 8; + /** + * How many consecutive {@code unknown} samples while a delegation is outstanding before we + * declare it stalled and fire {@link TurnListener#onTurnFailed} (CB-109). A worker wedged in a + * state herdr can't classify (e.g. an API-error screen) stays {@code unknown} indefinitely and + * would otherwise never resolve; any {@code working}/{@code idle} sample resets the streak, so a + * transient detection glitch cannot trip it. At the 250ms poll interval this is ~30s — far longer + * than any real detection blip, and still vastly better than the async send's timeout. + */ + private static final int TURN_STALL_GRACE_POLLS = 120; + private final AgentControl agents; private final TurnListener turnListener; private final ConcurrentHashMap targets = new ConcurrentHashMap<>(); @@ -79,6 +89,7 @@ public final class Injector { int injectableSincePickup; // consecutive injectable samples while awaitingPickup 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) synchronized void add(Pending p) { queue.add(p); @@ -118,14 +129,17 @@ public final class Injector { Pending sent = null; RuntimeException sendError = null; boolean turnCompleted = false; + boolean turnFailed = false; synchronized (t) { if (status == AgentStatus.WORKING) { // Definitive pickup: the worker is busy on our last message, and (if a delivery is // outstanding) a real turn is now confirmed to be running. t.awaitingPickup = false; t.injectableSincePickup = 0; + t.unknownSinceTurn = 0; if (t.awaitingCompletion) t.turnObserved = true; } else if (status.injectable()) { // IDLE or BLOCKED + t.unknownSinceTurn = 0; if (t.awaitingPickup && ++t.injectableSincePickup >= PICKUP_GRACE_POLLS) { // Pickup edge was never sampled (turn faster than the poll, or status lag). // Release the latch rather than wedge — and give up on synthesizing a completion @@ -167,9 +181,22 @@ public final class Injector { } } } + } else { + // UNKNOWN (or any other non-injectable, non-working): not a safe window nor a + // reliable pickup signal, so we never deliver or release the pickup latch here. But + // an outstanding delegation whose worker has gone unresponsive — stuck in a state + // herdr can't classify (CB-109) — will never yield a working→idle boundary. After a + // sustained streak, declare it failed so the awaiting send resolves rather than + // riding out the async timeout. (This also frees a delivery that wedged before it + // was ever picked up, which the injectable-only pickup grace could never release.) + if (t.awaitingCompletion && ++t.unknownSinceTurn >= TURN_STALL_GRACE_POLLS) { + t.awaitingPickup = false; + t.awaitingCompletion = false; + t.turnObserved = false; + t.unknownSinceTurn = 0; + turnFailed = true; + } } - // UNKNOWN (and any other non-injectable, non-working): do nothing — neither a safe - // window nor a reliable pickup signal, so we must not deliver or release the latch. // Reclaim the entry once the worker is fully quiescent (nothing queued, no pickup or // completion awaited), so the map cannot grow without bound across short-lived workers. @@ -183,6 +210,9 @@ public final class Injector { if (turnCompleted) { turnListener.onTurnComplete(target); } + if (turnFailed) { + turnListener.onTurnFailed(target); + } if (sent != null) { if (sendError != null) { log.warn("inject to {} failed, dropped message: {}", target, sendError.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 d3a0f04..e09478b 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java @@ -13,6 +13,15 @@ public interface TurnListener { /** A worker's delegated turn finished (worker returned to idle after visibly working). */ void onTurnComplete(String target); + /** + * A worker that visibly ran a delegated turn then wedged in a non-idle, non-working state + * (CB-109) — e.g. an error screen herdr classifies as {@code unknown} — so no + * {@code working → idle} completion boundary will ever arrive. A default no-op keeps this a + * functional interface; the completion resolver overrides it to fail the awaiting send. + */ + default void onTurnFailed(String target) { + } + /** No-op default for callers that only need delivery, not completion signalling. */ TurnListener NOOP = _ -> { }; diff --git a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java index fb0548e..714242b 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java @@ -122,6 +122,10 @@ public final class BridgeMcp { return text("[worker finished without a structured bridge_reply — transcript tail follows]\n" + r.text()); } + if (r.outcome() == MessageService.Outcome.WORKER_FAILED) { + // The worker ran the turn then wedged (CB-109) — surface the error context. + return text("[worker failed — turn ended in an unrecoverable state]\n" + r.text()); + } return text("[no reply within " + timeout + "ms — worker " + r.outcome().name().toLowerCase().replace("timed_out_", "") + "; retry or poll status]"); } catch (HerdrException e) { diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java index 94ae3eb..e1fa187 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -60,6 +60,11 @@ public final class MessageService { * {@code text} is the scraped transcript tail rather than a structured answer. */ COMPLETED_UNREPLIED, + /** + * The worker ran the turn then wedged in an unrecoverable state (CB-109); {@code text} is the + * failure context (e.g. the error screen). Terminal, but not a successful completion. + */ + WORKER_FAILED, /** Timed out after the message was delivered — the worker is still working. */ TIMED_OUT_WORKING, /** Timed out before delivery — the message is still queued for the worker. */ @@ -142,9 +147,11 @@ public final class MessageService { CompletableFuture reply = rendezvous.open(target); try { Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS); - Outcome outcome = r.kind() == Rendezvous.Kind.REPLY - ? Outcome.REPLIED - : Outcome.COMPLETED_UNREPLIED; + Outcome outcome = switch (r.kind()) { + case REPLY -> Outcome.REPLIED; + case COMPLETION -> Outcome.COMPLETED_UNREPLIED; + case FAILED -> Outcome.WORKER_FAILED; + }; return new Reply(outcome, r.text()); } catch (TimeoutException e) { boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally(); @@ -206,7 +213,12 @@ public final class MessageService { String source = r.outcome() == Outcome.REPLIED ? "reply" : "transcript"; return new TaskView(ticket, Phase.DONE, r.text(), source, null); } - return new TaskView(ticket, Phase.FAILED, null, null, "no reply — " + r.outcome().name().toLowerCase()); + // A wedged worker (CB-109) carries the error context as its reason; the timeout/busy + // outcomes carry none, so fall back to the outcome name. + String detail = r.outcome() == Outcome.WORKER_FAILED && r.text() != null + ? r.text() + : "no reply — " + r.outcome().name().toLowerCase(); + return new TaskView(ticket, Phase.FAILED, null, null, detail); } /** Best-effort live worker status for a pending poll; never throws (a lookup error is just noise). */ diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java b/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java index 5346ca8..f7db678 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java @@ -20,7 +20,9 @@ public final class Rendezvous { /** The worker called {@code bridge_reply} with a structured answer. */ REPLY, /** The worker's turn finished without a {@code bridge_reply}; {@code text} is a scrape. */ - COMPLETION + COMPLETION, + /** The worker ran the turn then wedged (CB-109); {@code text} is the failure context. */ + FAILED } /** The resolved outcome of a send: its {@link Kind} and the associated text. */ @@ -67,6 +69,17 @@ public final class Rendezvous { return complete(session, new Resolution(Kind.COMPLETION, text)); } + /** + * Resolve the send awaiting on {@code session} as a failure — the worker ran the turn but wedged + * in an unrecoverable state (CB-109); {@code reason} is the failure context (e.g. the error + * screen). A no-op if the worker had already raced a reply/completion in — first resolution wins. + * + * @return {@code true} if a waiter was resolved, {@code false} if none was waiting + */ + public boolean resolveFailure(String session, String reason) { + return complete(session, new Resolution(Kind.FAILED, reason)); + } + private boolean complete(String session, Resolution resolution) { CompletableFuture waiter = waiters.get(session); return waiter != null && waiter.complete(resolution); diff --git a/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java b/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java index ae87208..c224d22 100644 --- a/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java +++ b/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java @@ -173,9 +173,12 @@ public final class BridgedApp { case TIMED_OUT_WORKING -> "working"; case TIMED_OUT_QUEUED -> "queued"; case BUSY -> "busy"; + case WORKER_FAILED -> "failed"; case REPLIED, COMPLETED_UNREPLIED -> "done"; // unreachable (completed() branch) }, - "detail", "no reply within " + timeout + "ms; poll status or retry")); + "detail", reply.outcome() == MessageService.Outcome.WORKER_FAILED && reply.text() != null + ? reply.text() + : "no reply within " + timeout + "ms; poll status or retry")); } } catch (HerdrException e) { herdrError(ctx, e); diff --git a/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java b/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java index 4fec392..0d48606 100644 --- a/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java @@ -21,4 +21,16 @@ class CompletionResolverTest { assertFalse(herdr.called("agent.read"), "a turn nobody is blocked on must not cost a transcript scrape"); } + + @Test + void failSkipsTheScrapeWhenNoSendIsWaiting() { + FakeHerdr herdr = new FakeHerdr(); + Rendezvous rendezvous = new Rendezvous(); // no waiter opened + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + + resolver.fail("term_a"); + + assertFalse(herdr.called("agent.read"), + "a wedge nobody is blocked on must not cost a transcript scrape"); + } } 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 5f46da6..ff69286 100644 --- a/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java @@ -172,6 +172,55 @@ class InjectorTest { assertEquals(List.of(), completed, "no completion is synthesized from an unconfirmed turn"); } + /** Captures both turn-lifecycle callbacks so the CB-109 stall path can be asserted. */ + private static final class Captor implements TurnListener { + final List completed = new ArrayList<>(); + final List failed = new ArrayList<>(); + + @Override + public void onTurnComplete(String target) { + completed.add(target); + } + + @Override + public void onTurnFailed(String target) { + failed.add(target); + } + } + + // ~30s of unknown at the 250ms prod poll interval; enough onStatus samples to trip the stall. + private static final int STALL_SAMPLES = 130; + + @Test + void failsAnOutstandingDelegationWhoseWorkerWedgesInUnknown() { + 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); // worker starts the turn + for (int i = 0; i < STALL_SAMPLES; i++) inj.onStatus(T, AgentStatus.UNKNOWN); // then wedges + + assertEquals(List.of(T), cap.failed, "a sustained unknown streak fails the outstanding send"); + assertEquals(List.of(), cap.completed, "a wedge is a failure, not a completion"); + assertTrue(inj.activeTargets().isEmpty(), "the wedged target is reclaimed, not polled forever"); + } + + @Test + void aTransientUnknownGlitchNeitherFailsNorBlocksCompletion() { + 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); // confirmed turn + for (int i = 0; i < 10; i++) inj.onStatus(T, AgentStatus.UNKNOWN); // brief glitch, well under grace + inj.onStatus(T, AgentStatus.IDLE); // working → idle: the real completion + + assertEquals(List.of(), cap.failed, "a short unknown blip must not fail the turn"); + assertEquals(List.of(T), cap.completed, "the streak reset, so the turn still completes"); + } + @Test void sendFailureDropsMessageAndFailsItsFuture() { FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed"); 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 948acdc..7382c8d 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -11,6 +11,7 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertTrue; /** @@ -71,4 +72,20 @@ class MessageServiceTest { assertEquals(MessageService.Outcome.REPLIED, reply.outcome()); assertEquals("LGTM ship it", reply.text()); } + + @Test + void aWedgedWorkerResolvesTheSendAsFailedWithTheErrorContext() throws Exception { + herdr.readText("API Error: Unable to connect to API (ENOTFOUND)"); + CompletableFuture send = sendAsync(); + awaitWaiting(); + + injector.onStatus(T, AgentStatus.IDLE); // deliver + injector.onStatus(T, AgentStatus.WORKING); // worker starts the turn + for (int i = 0; i < 130; i++) injector.onStatus(T, AgentStatus.UNKNOWN); // then wedges (CB-109) + + MessageService.Reply reply = send.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.WORKER_FAILED, reply.outcome()); + 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"); + } }