#164 WIP: do not report an empty or crashed scrape as a successful reply

Rescued from an abandoned worker worktree, uncommitted since 2026-08-28.
Not reviewed and not built at this commit. The branch is 42 commits behind main.
This commit is contained in:
Dai Ha
2026-08-31 10:31:23 +07:00
parent 7c4170ff6d
commit 851ebcacaf
6 changed files with 187 additions and 43 deletions
@@ -351,6 +351,12 @@ public final class Fleetd {
sessions.onTurnComplete(target); sessions.onTurnComplete(target);
} }
@Override
public void onTurnComplete(String target, long elapsedNanos) {
completion.onTurnComplete(target, elapsedNanos);
sessions.onTurnComplete(target);
}
@Override @Override
public boolean hasPostTurnAction(String target) { public boolean hasPostTurnAction(String target) {
return sessions.hasPostTurnAction(target); return sessions.hasPostTurnAction(target);
@@ -362,6 +368,12 @@ public final class Fleetd {
return sessions.onTurnCompleteWithPostAction(target); return sessions.onTurnCompleteWithPostAction(target);
} }
@Override
public boolean onTurnCompleteWithPostAction(String target, long elapsedNanos) {
completion.resolveBeforePostAction(target, elapsedNanos);
return sessions.onTurnCompleteWithPostAction(target);
}
@Override @Override
public void onDelivered(String target, dev.ltms.fleet.msg.TurnToken token) { public void onDelivered(String target, dev.ltms.fleet.msg.TurnToken token) {
completion.onDelivered(target, token); completion.onDelivered(target, token);
@@ -12,6 +12,7 @@ import java.util.Set;
import java.util.TreeSet; import java.util.TreeSet;
import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.regex.Pattern; import java.util.regex.Pattern;
/** /**
@@ -61,6 +62,15 @@ public final class CompletionResolver implements TurnListener {
/** Cap the scraped tail so a long transcript can't return an unbounded blob. */ /** Cap the scraped tail so a long transcript can't return an unbounded blob. */
static final int MAX_SCRAPE_CHARS = 4000; static final int MAX_SCRAPE_CHARS = 4000;
/**
* A backend rejection can return the pane to idle in about one second. Two seconds is above the
* measured 1.03s crash while staying low enough not to reject ordinary short model answers.
*/
static final long MIN_COMPLETED_TURN_NANOS = TimeUnit.SECONDS.toNanos(2);
/** One stable, explicit backend-failure marker; broader error lists would be brittle. */
private static final Pattern BACKEND_ERROR = Pattern.compile("(?i)\\bAPI Error\\s*:");
private static final String CLIPPED_PANE_TAIL_MARKER = private static final String CLIPPED_PANE_TAIL_MARKER =
"[Pane tail clipped: member did not call fleet_reply.]"; "[Pane tail clipped: member did not call fleet_reply.]";
@@ -140,10 +150,15 @@ public final class CompletionResolver implements TurnListener {
@Override @Override
public void onTurnComplete(String target) { public void onTurnComplete(String target) {
onTurnComplete(target, Long.MAX_VALUE);
}
@Override
public void onTurnComplete(String target, long elapsedNanos) {
// Read the in-flight turn on the poller thread — before any next-turn delivery can overwrite // Read the in-flight turn on the poller thread — before any next-turn delivery can overwrite
// it — then off-load the scrape (a herdr round-trip we must not block polling on) to a vthread. // it — then off-load the scrape (a herdr round-trip we must not block polling on) to a vthread.
InFlight turn = inFlight.get(target); InFlight turn = inFlight.get(target);
Thread.ofVirtual().name("completion-" + target).start(() -> resolve(target, turn)); Thread.ofVirtual().name("completion-" + target).start(() -> resolve(target, turn, elapsedNanos));
} }
/** /**
@@ -152,7 +167,12 @@ public final class CompletionResolver implements TurnListener {
* path remains off-loaded so polling is not blocked by a scrape. * path remains off-loaded so polling is not blocked by a scrape.
*/ */
public void resolveBeforePostAction(String target) { public void resolveBeforePostAction(String target) {
resolve(target, inFlight.get(target)); resolveBeforePostAction(target, Long.MAX_VALUE);
}
/** Timed form used when adapter housekeeping must start immediately after the turn. */
public void resolveBeforePostAction(String target, long elapsedNanos) {
resolve(target, inFlight.get(target), elapsedNanos);
} }
@Override @Override
@@ -169,6 +189,11 @@ public final class CompletionResolver implements TurnListener {
/** Synchronous resolve (the unit-testable core of {@link #onTurnComplete}). */ /** Synchronous resolve (the unit-testable core of {@link #onTurnComplete}). */
void resolve(String target, InFlight turn) { void resolve(String target, InFlight turn) {
resolve(target, turn, Long.MAX_VALUE);
}
/** Synchronous timed resolve used by the real injector path. */
void resolve(String target, InFlight turn, long elapsedNanos) {
CompletableFuture<Rendezvous.Resolution> waiter = turn == null ? null : turn.waiter(); CompletableFuture<Rendezvous.Resolution> waiter = turn == null ? null : turn.waiter();
if (waiter == null || waiter.isDone()) { if (waiter == null || waiter.isDone()) {
// Nobody is blocked on THIS turn (it had no send, or its fleet_reply already won). Skip // Nobody is blocked on THIS turn (it had no send, or its fleet_reply already won). Skip
@@ -180,50 +205,65 @@ public final class CompletionResolver implements TurnListener {
String assistantBlock = null; String assistantBlock = null;
int originalLength = 0; int originalLength = 0;
boolean clipped = false; boolean clipped = false;
boolean scrapeFailed = false; String raw;
try { try {
assistantBlock = lastAssistantBlock(agents.read(target, SCRAPE_SOURCE)); raw = agents.read(target, SCRAPE_SOURCE);
assistantBlock = lastAssistantBlock(raw);
originalLength = assistantBlock.strip().length(); originalLength = assistantBlock.strip().length();
clipped = originalLength > MAX_SCRAPE_CHARS; clipped = originalLength > MAX_SCRAPE_CHARS;
tail = clip(assistantBlock); tail = clip(assistantBlock);
} catch (RuntimeException e) { } catch (RuntimeException e) {
// The worker finished but we couldn't read its screen — still resolve the send so the resolveFailure(target, turn, "member " + target + " ended its turn, but its completion "
// caller unblocks; an empty tail beats hanging until the caller's timeout. + "screen could not be read: " + e.getMessage());
log.warn("completion scrape for {} failed; resolving with an empty tail: {}", return;
target, e.getMessage()); }
tail = ""; // CB-578 stage A: a turn that ended with no fleet_reply AND whose scrape matches the
scrapeFailed = true; // backend's configured usage-limit pattern is a refusal, not an answer. Classify it as
// BACKEND_EXHAUSTED rather than handing the caller a scrape that reads like a real reply.
// When TUI chrome hides the assistant block, inspect the raw screen so an on-screen refusal
// still reaches the caller as an error.
String visibleTurn = assistantBlock.isBlank() ? raw : assistantBlock;
Pattern exhausted = exhaustedPatterns.patternFor(target);
String matchedLine = exhausted == null ? null : firstMatchingLine(visibleTurn, exhausted);
if (matchedLine != null) {
String reason = "backend exhausted (usage limit): " + matchedLine;
if (rendezvous.resolveExhausted(waiter, reason)) {
inFlight.remove(target, turn);
log.warn("completion for {} classified BACKEND_EXHAUSTED (no fleet_reply; scrape "
+ "matched the profile's exhausted pattern): {}", target, reason);
// CB-578 stage B: only on the resolution that actually won the race — a late
// duplicate must never quarantine a credential twice for one refusal.
exhaustionSink.onExhausted(target, reason);
}
return;
}
String backendError = firstMatchingLine(visibleTurn, BACKEND_ERROR);
if (backendError != null) {
resolveFailure(target, turn, "member " + target + " ended on a backend error: " + backendError);
return;
}
if (tail.isBlank()) {
resolveFailure(target, turn, "member " + target
+ " ended its turn without a reply or any readable completion output");
return;
}
if (elapsedNanos < MIN_COMPLETED_TURN_NANOS) {
resolveFailure(target, turn, "member " + target + " returned to idle after "
+ TimeUnit.NANOSECONDS.toMillis(elapsedNanos)
+ "ms, below the minimum credible turn duration of "
+ TimeUnit.NANOSECONDS.toMillis(MIN_COMPLETED_TURN_NANOS) + "ms");
return;
} }
// Misattribution guard (CB-115): if the scrape is byte-identical to the pane content at // Misattribution guard (CB-115): if the scrape is byte-identical to the pane content at
// delivery, this turn produced no new output — the boundary belongs to the previous turn's // delivery, this turn produced no new output — the boundary belongs to the previous turn's
// wind-down (common on rapid back-to-back sends). Suppress rather than resolve the send with // wind-down (common on rapid back-to-back sends). Suppress rather than resolve the send with
// a stale answer; the real fleet_reply (or a later genuine completion) resolves it instead. // a stale answer; the real fleet_reply (or a later genuine completion) resolves it instead.
// A scrape that failed to read is exempt — an empty tail there is "couldn't see", not "no change".
String baseline = turn.baseline(); String baseline = turn.baseline();
if (!scrapeFailed && baseline != null && baseline.equals(tail)) { if (baseline != null && baseline.equals(tail)) {
log.debug("suppressing misattributed completion for {} (no output change since delivery)", log.debug("suppressing misattributed completion for {} (no output change since delivery)",
target); target);
return; // keep the in-flight record: a later genuine completion still needs it return; // keep the in-flight record: a later genuine completion still needs it
} }
// CB-578 stage A: a turn that ended with no fleet_reply AND whose scrape matches the
// backend's configured usage-limit pattern is a refusal, not an answer. Classify it as
// BACKEND_EXHAUSTED rather than handing the caller a scrape that reads like a real reply.
if (!scrapeFailed) {
Pattern exhausted = exhaustedPatterns.patternFor(target);
String matchedLine = exhausted == null ? null : firstMatchingLine(assistantBlock, exhausted);
if (matchedLine != null) {
String reason = "backend exhausted (usage limit): " + matchedLine;
if (rendezvous.resolveExhausted(waiter, reason)) {
inFlight.remove(target, turn);
log.warn("completion for {} classified BACKEND_EXHAUSTED (no fleet_reply; scrape "
+ "matched the profile's exhausted pattern): {}", target, reason);
// CB-578 stage B: only on the resolution that actually won the race — a late
// duplicate must never quarantine a credential twice for one refusal.
exhaustionSink.onExhausted(target, reason);
}
return;
}
}
String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail; String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail;
if (rendezvous.resolveCompletion(waiter, completion)) { if (rendezvous.resolveCompletion(waiter, completion)) {
inFlight.remove(target, turn); inFlight.remove(target, turn);
@@ -272,6 +312,14 @@ public final class CompletionResolver implements TurnListener {
} }
} }
/** Resolve a captured completion waiter as {@code WORKER_FAILED}, never as reply content. */
private void resolveFailure(String target, InFlight turn, String reason) {
if (rendezvous.resolveFailure(turn.waiter(), reason)) {
inFlight.remove(target, turn);
log.warn("failing send to {} via completion fallback: {}", target, reason);
}
}
/** /**
* The first line of {@code text} matching {@code pattern}, stripped — the CB-578 stage A * The first line of {@code text} matching {@code pattern}, stripped — the CB-578 stage A
* evidence carried in a {@code BACKEND_EXHAUSTED} reason so the operator sees the real refusal * evidence carried in a {@code BACKEND_EXHAUSTED} reason so the operator sees the real refusal
@@ -14,6 +14,7 @@ import java.util.Set;
import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.function.Consumer; import java.util.function.Consumer;
import java.util.function.LongSupplier;
import java.util.function.Predicate; import java.util.function.Predicate;
import java.util.stream.Collectors; import java.util.stream.Collectors;
@@ -92,6 +93,7 @@ public final class Injector {
private final TurnListener turnListener; private final TurnListener turnListener;
private final Predicate<String> ready; // CB-113: a target is deliverable only when available private final Predicate<String> ready; // CB-113: a target is deliverable only when available
private final Consumer<String> forget; // CB-114: clear a gone worker's readiness/presence private final Consumer<String> forget; // CB-114: clear a gone worker's readiness/presence
private final LongSupplier nowNanos;
private final ConcurrentHashMap<String, Target> targets = new ConcurrentHashMap<>(); private final ConcurrentHashMap<String, Target> targets = new ConcurrentHashMap<>();
/** Delivery only; completion signalling is a no-op and every target is treated as available. */ /** Delivery only; completion signalling is a no-op and every target is treated as available. */
@@ -123,10 +125,17 @@ public final class Injector {
*/ */
public Injector(AgentControl agents, TurnListener turnListener, Predicate<String> ready, public Injector(AgentControl agents, TurnListener turnListener, Predicate<String> ready,
Consumer<String> forget) { Consumer<String> forget) {
this(agents, turnListener, ready, forget, System::nanoTime);
}
/** Full constructor with an injectable monotonic clock for turn-duration tests. */
public Injector(AgentControl agents, TurnListener turnListener, Predicate<String> ready,
Consumer<String> forget, LongSupplier nowNanos) {
this.agents = agents; this.agents = agents;
this.turnListener = turnListener; this.turnListener = turnListener;
this.ready = ready; this.ready = ready;
this.forget = forget; this.forget = forget;
this.nowNanos = nowNanos;
} }
/** A pending message and the future that completes when it has been delivered. */ /** A pending message and the future that completes when it has been delivered. */
@@ -140,6 +149,7 @@ public final class Injector {
int injectableSincePickup; // consecutive injectable samples while awaitingPickup int injectableSincePickup; // consecutive injectable samples while awaitingPickup
boolean awaitingCompletion; // a delivered message's turn is not yet known-complete boolean awaitingCompletion; // a delivered message's turn is not yet known-complete
boolean turnObserved; // saw a real `working` sample since that delivery (turn ran) boolean turnObserved; // saw a real `working` sample since that delivery (turn ran)
long turnStartedAtNanos; // injected monotonic time of the first `working` sample
int unknownSinceTurn; // consecutive `unknown` samples while a delegation is outstanding (CB-109) int unknownSinceTurn; // consecutive `unknown` samples while a delegation is outstanding (CB-109)
int notReadySincePoll; // consecutive injectable samples a queued message waited on the readiness gate (CB-114) int notReadySincePoll; // consecutive injectable samples a queued message waited on the readiness gate (CB-114)
boolean postTurnPending; // completion observed; adapter housekeeping has not started yet boolean postTurnPending; // completion observed; adapter housekeeping has not started yet
@@ -185,6 +195,7 @@ public final class Injector {
Pending sent = null; Pending sent = null;
RuntimeException sendError = null; RuntimeException sendError = null;
boolean turnCompleted = false; boolean turnCompleted = false;
long turnElapsedNanos = 0;
boolean turnFailed = false; boolean turnFailed = false;
boolean resubmit = false; boolean resubmit = false;
boolean startPostTurn = false; boolean startPostTurn = false;
@@ -202,7 +213,10 @@ public final class Injector {
t.injectableSincePickup = 0; t.injectableSincePickup = 0;
t.unknownSinceTurn = 0; t.unknownSinceTurn = 0;
t.notReadySincePoll = 0; t.notReadySincePoll = 0;
if (t.awaitingCompletion) t.turnObserved = true; if (t.awaitingCompletion && !t.turnObserved) {
t.turnObserved = true;
t.turnStartedAtNanos = nowNanos.getAsLong();
}
} else if (status.injectable()) { // IDLE or BLOCKED } else if (status.injectable()) { // IDLE or BLOCKED
t.unknownSinceTurn = 0; t.unknownSinceTurn = 0;
if (t.awaitingPostTurnPickup) { if (t.awaitingPostTurnPickup) {
@@ -238,6 +252,7 @@ public final class Injector {
if (t.awaitingCompletion && t.turnObserved) { if (t.awaitingCompletion && t.turnObserved) {
t.awaitingCompletion = false; t.awaitingCompletion = false;
t.turnObserved = false; t.turnObserved = false;
turnElapsedNanos = Math.max(0, nowNanos.getAsLong() - t.turnStartedAtNanos);
turnCompleted = true; turnCompleted = true;
if (turnListener.hasPostTurnAction(target)) { if (turnListener.hasPostTurnAction(target)) {
t.postTurnPending = true; t.postTurnPending = true;
@@ -332,7 +347,7 @@ public final class Injector {
} }
if (turnCompleted) { if (turnCompleted) {
if (startPostTurn) { if (startPostTurn) {
boolean started = turnListener.onTurnCompleteWithPostAction(target); boolean started = turnListener.onTurnCompleteWithPostAction(target, turnElapsedNanos);
synchronized (t) { synchronized (t) {
t.postTurnPending = false; t.postTurnPending = false;
if (started) { if (started) {
@@ -344,7 +359,7 @@ public final class Injector {
} }
} }
} else { } else {
turnListener.onTurnComplete(target); turnListener.onTurnComplete(target, turnElapsedNanos);
} }
} }
if (turnFailed) { if (turnFailed) {
@@ -15,6 +15,14 @@ public interface TurnListener {
/** A worker's delegated turn finished (worker returned to idle after visibly working). */ /** A worker's delegated turn finished (worker returned to idle after visibly working). */
void onTurnComplete(String target); void onTurnComplete(String target);
/**
* A worker's delegated turn finished after {@code elapsedNanos} of observed work. The default
* keeps listeners that do not use timing source-compatible with the original callback.
*/
default void onTurnComplete(String target, long elapsedNanos) {
onTurnComplete(target);
}
/** /**
* Whether completion must pause ordinary delivery while adapter-specific housekeeping starts. * Whether completion must pause ordinary delivery while adapter-specific housekeeping starts.
* This is queried before the injector considers the next queued message, closing the same-tick * This is queried before the injector considers the next queued message, closing the same-tick
@@ -34,6 +42,11 @@ public interface TurnListener {
return false; return false;
} }
/** Timed form of {@link #onTurnCompleteWithPostAction(String)}. */
default boolean onTurnCompleteWithPostAction(String target, long elapsedNanos) {
return onTurnCompleteWithPostAction(target);
}
/** /**
* A worker that visibly ran a delegated turn then wedged in a non-idle, non-working state * 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 * (CB-109) — e.g. an error screen herdr classifies as {@code unknown} — so no
@@ -249,12 +249,9 @@ class CompletionResolverTest {
} }
@Test @Test
void resolvesWhenTheScrapeItselfFailsEvenWithABaselinePresent() { void failsWhenTheScrapeItselfFailsEvenWithABaselinePresent() {
// The most important branch of the CB-115 guard: a failed read means the resolver could not // A failed read means the resolver could not see the screen. It must unblock the caller as a
// SEE the screen — "couldn't see", not "no change". It must still resolve the send (an empty // failure, never turn that missing evidence into a successful empty completion.
// tail beats hanging until the caller's timeout), even though a baseline was captured. The
// baseline here is "" (an empty pane at delivery), so without the !scrapeFailed clause the
// byte-identical guard would wrongly match the empty tail and suppress.
FakeHerdr herdr = new FakeHerdr().healthy(false); // agent.read throws HerdrException FakeHerdr herdr = new FakeHerdr().healthy(false); // agent.read throws HerdrException
Rendezvous rendezvous = new Rendezvous(); Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
@@ -265,8 +262,8 @@ class CompletionResolverTest {
assertTrue(waiter.isDone(), assertTrue(waiter.isDone(),
"a failed scrape must still resolve the send, not hang until the caller's timeout"); "a failed scrape must still resolve the send, not hang until the caller's timeout");
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind()); assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind());
assertEquals("", waiter.getNow(null).text(), "the tail is empty because the screen was unreadable"); assertTrue(waiter.getNow(null).text().contains("term_a"), "the failure names the member");
} }
// --- CB-115/CB-116 fail guard: an already-done or absent waiter is left alone --------- // --- CB-115/CB-116 fail guard: an already-done or absent waiter is left alone ---------
@@ -14,6 +14,7 @@ import org.junit.jupiter.api.Test;
import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertFalse;
@@ -35,7 +36,8 @@ class MessageServiceTest {
private final Rendezvous rendezvous = new Rendezvous(); private final Rendezvous rendezvous = new Rendezvous();
private final CompletionResolver completion = private final CompletionResolver completion =
new CompletionResolver(agents, rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none()); new CompletionResolver(agents, rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
private final Injector injector = new Injector(agents, completion); private final AtomicLong nowNanos = new AtomicLong();
private final Injector injector = new Injector(agents, completion, _ -> true, _ -> { }, nowNanos::get);
private final InMemoryReplyInbox inbox = new InMemoryReplyInbox(); private final InMemoryReplyInbox inbox = new InMemoryReplyInbox();
private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox); private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
@@ -68,6 +70,7 @@ class MessageServiceTest {
injector.onStatus(T, AgentStatus.IDLE); // deliver the task (baselines the pre-turn content) injector.onStatus(T, AgentStatus.IDLE); // deliver the task (baselines the pre-turn content)
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up and works injector.onStatus(T, AgentStatus.WORKING); // worker picks it up and works
herdr.readText("BUILD GREEN: 391 files"); // the worker's turn produced new output herdr.readText("BUILD GREEN: 391 files"); // the worker's turn produced new output
nowNanos.addAndGet(TimeUnit.SECONDS.toNanos(3));
injector.onStatus(T, AgentStatus.IDLE); // working → idle: turn complete, no fleet_reply injector.onStatus(T, AgentStatus.IDLE); // working → idle: turn complete, no fleet_reply
MessageService.Reply reply = send.get(5, TimeUnit.SECONDS); MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);
@@ -77,6 +80,62 @@ class MessageServiceTest {
assertTrue(reply.completed(), "a scraped completion still counts as completed"); assertTrue(reply.completed(), "a scraped completion still counts as completed");
} }
@Test
void emptyCompletionThroughMessageServiceFailsAndNamesTheMember() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
herdr.readText("$ prompt");
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
nowNanos.addAndGet(TimeUnit.SECONDS.toNanos(3));
herdr.readText("");
injector.onStatus(T, AgentStatus.IDLE);
MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.WORKER_FAILED, reply.outcome(),
"an empty scrape must use the caller's failure outcome, not an empty completion");
assertFalse(reply.completed());
assertTrue(reply.text().contains(T), "the failure must name the member: " + reply.text());
}
@Test
void shortCompletionThroughMessageServiceFailsEvenWithContent() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
herdr.readText("$ prompt");
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
nowNanos.addAndGet(TimeUnit.SECONDS.toNanos(1));
herdr.readText("looks like an answer but the backend returned immediately");
injector.onStatus(T, AgentStatus.IDLE);
MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.WORKER_FAILED, reply.outcome(),
"a turn below the duration floor is a crash signature, not a completion");
assertTrue(reply.text().contains("1000ms"), "the reason carries the observed duration: " + reply.text());
}
@Test
void backendErrorScrapeThroughMessageServiceFailsInsteadOfBecomingReplyText() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
herdr.readText("$ prompt");
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
nowNanos.addAndGet(TimeUnit.SECONDS.toNanos(3));
herdr.readText("● API Error: 400 invalid request body");
injector.onStatus(T, AgentStatus.IDLE);
MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.WORKER_FAILED, reply.outcome());
assertFalse(reply.completed(), "plain backend errors are never fallback reply content");
assertTrue(reply.text().contains("API Error: 400"),
"the visible backend error is carried as the failure reason: " + reply.text());
}
@Test @Test
void explicitFleetReplyResolvesAsReplied() throws Exception { void explicitFleetReplyResolvesAsReplied() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync(); CompletableFuture<MessageService.Reply> send = sendAsync();