CB-109: fail a delegation whose worker wedges in an unknown state
Dogfood found the gap: a turn that dies into an error screen herdr reports as 'unknown' (e.g. the worker hitting API ENOTFOUND) never produces a working->idle boundary, so CB-106 never fires and the async send rides its full 30-min timeout. The injector now counts consecutive 'unknown' samples while a delegation is outstanding; any working/idle sample resets the streak, so only a genuine wedge (~30s continuous unknown) trips it. It then fires TurnListener.onTurnFailed; CompletionResolver scrapes the error screen and resolves the send via Rendezvous.resolveFailure (Kind.FAILED -> Outcome.WORKER_FAILED), surfaced as async phase=failed / REST status=failed / a [worker failed] MCP note, with the error context as the reason. This also frees a delivery that wedged before pickup, which the injectable-only pickup grace could never release.
This commit is contained in:
@@ -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.
|
||||
*
|
||||
* <p>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.
|
||||
* <p>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.
|
||||
*
|
||||
* <p>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();
|
||||
|
||||
@@ -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<String, Target> 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());
|
||||
|
||||
@@ -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 = _ -> {
|
||||
};
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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<Rendezvous.Resolution> 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). */
|
||||
|
||||
@@ -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<Resolution> waiter = waiters.get(session);
|
||||
return waiter != null && waiter.complete(resolution);
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<String> completed = new ArrayList<>();
|
||||
final List<String> 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");
|
||||
|
||||
@@ -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<MessageService.Reply> 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");
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user