fleetd#164: an empty or suspiciously fast scrape must fail, never resolve as a success
CI / contract (pull_request) Successful in 1m3s
CI / build (pull_request) Successful in 1m20s

CompletionResolver.resolve() used to hand the caller a successful "" reply whenever a
turn's scrape came back empty (whether the read failed, or genuinely produced nothing),
making a lost turn indistinguishable from a real empty answer. It also had no way to
tell a crashed backend's near-instant BUSY -> DONE transition apart from a genuine
completion.

Add MIN_TURN_NANOS (2s), a named floor below which a completed turn is treated as a
crash signature and failed rather than resolved as a reply. Fail on any empty scrape
(read failure or a clean-but-empty read) instead of resolving with "". Both failures
name the member and carry whatever is on the pane for context.

Thread an injectable LongSupplier clock through CompletionResolver (matching the
SessionManager/MessageService nowNanos pattern) so the floor is testable without a
real sleep.
This commit is contained in:
Dai Ha
2026-08-28 06:01:00 +07:00
parent 7c4170ff6d
commit 3bfa82839b
3 changed files with 211 additions and 33 deletions
@@ -6,12 +6,14 @@ import dev.ltms.fleet.msg.TurnToken;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.time.Duration;
import java.util.List;
import java.util.Objects;
import java.util.Set;
import java.util.TreeSet;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.LongSupplier;
import java.util.regex.Pattern;
/**
@@ -61,6 +63,17 @@ public final class CompletionResolver implements TurnListener {
/** Cap the scraped tail so a long transcript can't return an unbounded blob. */
static final int MAX_SCRAPE_CHARS = 4000;
/**
* fleetd#164: the floor below which a {@code BUSY -> DONE} transition cannot be real work. A
* backend that rejects a turn outright (e.g. an HTTP 400 from the model, before the worker read
* a single file or produced a token) drives the exact same confirmed {@code working -> idle}
* transition a genuine completion does — just in about a second instead of the many seconds a
* real turn costs. {@link #onTurnComplete} cannot tell those two cases apart from the transition
* alone, so a turn that settles inside this floor is treated as a crash signature and resolved
* as a failure, never as a (possibly empty) success.
*/
public static final long MIN_TURN_NANOS = Duration.ofSeconds(2).toNanos();
private static final String CLIPPED_PANE_TAIL_MARKER =
"[Pane tail clipped: member did not call fleet_reply.]";
@@ -68,6 +81,7 @@ public final class CompletionResolver implements TurnListener {
private final Rendezvous rendezvous;
private final ExhaustedPatternLookup exhaustedPatterns;
private final ExhaustionSink exhaustionSink;
private final LongSupplier nowNanos;
/**
* Per-target record of the turn currently in flight: the exact {@link Rendezvous} waiter its
@@ -79,8 +93,21 @@ public final class CompletionResolver implements TurnListener {
* reference: a completion scrape equal to it means the worker produced no new output (the previous
* turn's wind-down sampled as this boundary), so it is suppressed. Overwritten on each delivery;
* cleared when the turn resolves. Package-private so tests can capture and replay a specific turn.
*
* <p>{@code deliveredAtNanos} (fleetd#164) is the {@link #nowNanos} reading taken at delivery —
* the other half of the {@link #MIN_TURN_NANOS} floor check, compared against a fresh reading at
* resolution time.
*/
record InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline) {
record InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline, long deliveredAtNanos) {
/**
* Convenience for tests exercising scrape/suppression logic that don't care about turn
* timing: back-dates the delivery far enough that {@link #MIN_TURN_NANOS} can never fire.
* Not used by production code — {@link #captureBaseline} always records a real reading.
*/
InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline) {
this(waiter, baseline, Long.MIN_VALUE / 2);
}
}
private final ConcurrentHashMap<String, InFlight> inFlight = new ConcurrentHashMap<>();
@@ -97,10 +124,25 @@ public final class CompletionResolver implements TurnListener {
*/
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
ExhaustionSink exhaustionSink) {
this(agents, rendezvous, exhaustedPatterns, exhaustionSink, System::nanoTime);
}
/**
* Test constructor with an injectable clock (fleetd#164), matching the {@code LongSupplier}
* pattern {@link dev.ltms.fleet.session.SessionManager} and {@link dev.ltms.fleet.msg.MessageService}
* already use: lets a test place a turn's delivery and its resolution at an exact, controllable
* distance apart around the {@link #MIN_TURN_NANOS} floor, without a real sleep. Public (rather
* than package-private like those two) because callers that wire a full {@code MessageService}
* fixture — e.g. {@code MessageServiceTest} — construct this resolver directly from another
* package.
*/
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
ExhaustionSink exhaustionSink, LongSupplier nowNanos) {
this.agents = agents;
this.rendezvous = rendezvous;
this.exhaustedPatterns = Objects.requireNonNull(exhaustedPatterns, "exhaustedPatterns");
this.exhaustionSink = Objects.requireNonNull(exhaustionSink, "exhaustionSink");
this.nowNanos = Objects.requireNonNull(nowNanos, "nowNanos");
}
@Override
@@ -130,7 +172,7 @@ public final class CompletionResolver implements TurnListener {
baseline = null; // fail open: no baseline ⇒ no suppression
log.debug("delivery baseline for {} failed: {}", target, e.getMessage());
}
inFlight.put(target, new InFlight(waiter, baseline));
inFlight.put(target, new InFlight(waiter, baseline, nowNanos.getAsLong()));
}
/** The turn currently baselined for {@code target}, or {@code null} — a test hook for the captureBaseline path. */
@@ -176,6 +218,15 @@ public final class CompletionResolver implements TurnListener {
inFlight.remove(target, turn);
return;
}
// fleetd#164: a BUSY -> DONE transition inside the floor cannot be real work — it's a crash
// signature (e.g. a backend HTTP 400 before the worker did anything), not a fast answer. Fail
// it before spending a scrape on the ordinary path; the reason still carries whatever is on
// screen, since that is usually the backend's own error.
long elapsedNanos = nowNanos.getAsLong() - turn.deliveredAtNanos();
if (elapsedNanos < MIN_TURN_NANOS) {
fail(target, turn, tooFastReason(target, elapsedNanos));
return;
}
String tail;
String assistantBlock = null;
int originalLength = 0;
@@ -187,20 +238,25 @@ public final class CompletionResolver implements TurnListener {
clipped = originalLength > MAX_SCRAPE_CHARS;
tail = clip(assistantBlock);
} catch (RuntimeException e) {
// The worker finished but we couldn't read its screen — still resolve the send so the
// caller unblocks; an empty tail beats hanging until the caller's timeout.
log.warn("completion scrape for {} failed; resolving with an empty tail: {}",
target, e.getMessage());
log.warn("completion scrape for {} failed: {}", target, e.getMessage());
tail = "";
scrapeFailed = true;
}
// fleetd#164: a scrape nobody could read, and a scrape that read cleanly but produced nothing,
// both used to resolve the send as a SUCCESS carrying "" — indistinguishable from a worker that
// genuinely finished with nothing to say. That is the defect: fail loudly instead, naming the
// member, so a caller (including a lead deciding whether to delegate again) can tell a lost
// turn from a real empty answer.
if (scrapeFailed || tail.isEmpty()) {
fail(target, turn, emptyScrapeReason(target, scrapeFailed));
return;
}
// 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
// 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 scrape that failed to read is exempt — an empty tail there is "couldn't see", not "no change".
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)",
target);
return; // keep the in-flight record: a later genuine completion still needs it
@@ -208,21 +264,19 @@ public final class CompletionResolver implements TurnListener {
// 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;
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;
if (rendezvous.resolveCompletion(waiter, completion)) {
@@ -272,6 +326,36 @@ public final class CompletionResolver implements TurnListener {
}
}
/**
* fleetd#164: the failure reason for a turn that settled inside {@link #MIN_TURN_NANOS} — names
* the member and both timings, and appends whatever the pane shows (usually the backend's own
* error) so the caller sees the cause, not just "it failed".
*/
private String tooFastReason(String target, long elapsedNanos) {
String scrape;
try {
scrape = clip(agents.read(target, SCRAPE_SOURCE));
} catch (RuntimeException e) {
scrape = "";
}
String reason = String.format(
"member %s went BUSY -> DONE in %dms (floor %dms) — too fast to be real work, most "
+ "likely a backend error before any work started",
target, elapsedNanos / 1_000_000, MIN_TURN_NANOS / 1_000_000);
return scrape.isBlank() ? reason : reason + ": " + scrape;
}
/**
* fleetd#164: the failure reason for a scrape that produced zero characters — names the member
* and says plainly that the turn produced nothing, so a caller (a lead deciding whether to
* delegate again included) never mistakes a lost turn for a genuinely empty reply.
*/
private static String emptyScrapeReason(String target, boolean scrapeFailed) {
return "member " + target + " turn completed with an empty scrape (0 chars) — "
+ (scrapeFailed ? "its pane could not be read; " : "")
+ "treating as a lost turn, not a real answer";
}
/**
* 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
@@ -197,10 +197,15 @@ class CompletionResolverTest {
void resolvesSynchronouslyBeforePostTurnContextClearing() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ previous answer\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
// fleetd#164: an injectable clock, so this real (non-crash) turn lands outside MIN_TURN_NANOS
// — captureBaseline and resolveBeforePostAction below run back-to-back with no real delay.
long[] clock = {1_000_000_000L};
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
var waiter = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter));
herdr.readText("⏺ answer that /clear would erase\n❯ ");
clock[0] += CompletionResolver.MIN_TURN_NANOS + 1; // this turn took longer than the floor
resolver.resolveBeforePostAction("term_a");
@@ -219,7 +224,11 @@ class CompletionResolverTest {
String longBlock = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 500) + "\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(longBlock);
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
// fleetd#164: an injectable clock so captureBaseline and resolve (back-to-back, no real
// delay) don't trip the too-fast-turn floor — this test is about the suppression guard, not timing.
long[] clock = {1_000_000_000L};
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
var waiter = rendezvous.open("term_a"); // a send is blocked on this turn
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter)); // baseline is the clipped >cap block
@@ -227,6 +236,7 @@ class CompletionResolverTest {
assertEquals(CompletionResolver.MAX_SCRAPE_CHARS, turn.baseline().length(),
"the delivery baseline is clipped to the same cap resolve() applies to the tail");
clock[0] += CompletionResolver.MIN_TURN_NANOS + 1; // outside the floor
resolver.resolve("term_a", turn); // scrape unchanged → clipped tail == baseline → suppress
assertFalse(waiter.isDone(),
@@ -249,12 +259,13 @@ class CompletionResolverTest {
}
@Test
void resolvesWhenTheScrapeItselfFailsEvenWithABaselinePresent() {
// The most important branch of the CB-115 guard: a failed read means the resolver could not
// SEE the screen — "couldn't see", not "no change". It must still resolve the send (an empty
// 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.
void aFailedScrapeResolvesAsAFailureEvenWithABaselinePresent() {
// fleetd#164: this test used to assert that a failed read resolved the send as a SUCCESS
// carrying an empty string ("an empty tail beats hanging until the caller's timeout") — that
// was the bug this ticket fixes: a lost turn and a genuine empty answer looked identical to
// every caller. This test encoded the bug and is changed here: a failed read must fail the
// send instead, naming the member, whether or not a baseline was captured (the baseline here
// is "", an empty pane at delivery — proof this isn't the CB-115 misattribution path either).
FakeHerdr herdr = new FakeHerdr().healthy(false); // agent.read throws HerdrException
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
@@ -265,8 +276,82 @@ class CompletionResolverTest {
assertTrue(waiter.isDone(),
"a failed scrape must still resolve the send, not hang until the caller's timeout");
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind());
assertEquals("", waiter.getNow(null).text(), "the tail is empty because the screen was unreadable");
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
"a failed read is a lost turn, not a successful empty reply");
assertTrue(waiter.getNow(null).text().contains("term_a"),
"the failure names the member: " + waiter.getNow(null).text());
assertTrue(waiter.getNow(null).text().contains("could not be read"),
"the failure explains the scrape could not be read: " + waiter.getNow(null).text());
}
// --- fleetd#164: an empty (but readable) scrape must never resolve as a success --------------
@Test
void anEmptyScrapeResolvesAsAFailureNamingTheMember() {
// The core defect: a scrape that read CLEANLY but produced zero characters used to resolve
// the send as a SUCCESS carrying "" — indistinguishable, to every caller, from a worker that
// genuinely finished with nothing to say. A lost turn must never look like a real empty reply.
FakeHerdr herdr = new FakeHerdr().readText("");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertTrue(waiter.isDone(), "an empty scrape must still resolve the send, not hang");
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
"an empty scrape is a lost turn, not a successful empty reply");
assertTrue(waiter.getNow(null).text().contains("term_a"),
"the failure names the member: " + waiter.getNow(null).text());
assertTrue(waiter.getNow(null).text().toLowerCase().contains("empty"),
"the failure says the scrape was empty: " + waiter.getNow(null).text());
}
// --- fleetd#164: a BUSY -> DONE transition inside the floor is a crash, not a fast answer ------
@Test
void aBusyToDoneTransitionInsideTheFloorResolvesAsAFailure() {
// The exact fleetd#164 scenario: the backend returned an HTTP 400 before the worker did
// anything, and the member went BUSY -> DONE in ~1s. That transition alone is indistinguishable
// from a genuine (if unusually fast) completion, so the resolver leans on the floor to catch it.
FakeHerdr herdr = new FakeHerdr().readText("⏺ HTTP 400: invalid request\n❯ ");
Rendezvous rendezvous = new Rendezvous();
long[] clock = {10_000_000_000L};
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
var waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, null, clock[0]); // delivered "now"
clock[0] += CompletionResolver.MIN_TURN_NANOS - 1; // 1ns inside the floor
resolver.resolve("term_a", turn);
assertTrue(waiter.isDone(), "a suspiciously fast turn must still resolve (as a failure)");
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind());
assertTrue(waiter.getNow(null).text().contains("term_a"),
"the failure names the member: " + waiter.getNow(null).text());
assertTrue(waiter.getNow(null).text().contains("HTTP 400"),
"the failure carries whatever was on screen: " + waiter.getNow(null).text());
}
@Test
void aBusyToDoneTransitionJustOutsideTheFloorResolvesNormally() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ a real, if quick, answer\n❯ ");
Rendezvous rendezvous = new Rendezvous();
long[] clock = {10_000_000_000L};
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
var waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, null, clock[0]); // delivered "now"
clock[0] += CompletionResolver.MIN_TURN_NANOS + 1; // 1ns outside the floor
resolver.resolve("term_a", turn);
assertTrue(waiter.isDone());
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind(),
"a turn that took longer than the floor resolves normally");
assertEquals("a real, if quick, answer", waiter.getNow(null).text());
}
// --- CB-115/CB-116 fail guard: an already-done or absent waiter is left alone ---------
@@ -33,8 +33,17 @@ class MessageServiceTest {
private final FakeHerdr herdr = new FakeHerdr().readText("BUILD GREEN: 391 files");
private final AgentControl agents = new AgentControl(herdr);
private final Rendezvous rendezvous = new Rendezvous();
/**
* fleetd#164: this fixture drives a delivery and its completion back-to-back with no real time
* between them, so the real clock would trip {@link CompletionResolver#MIN_TURN_NANOS} on every
* completion-fallback test here. An ever-advancing fake clock stands in for the model-latency and
* herdr round-trips a real turn would spend, so each delivery-then-resolve pair still lands
* outside the floor.
*/
private final java.util.concurrent.atomic.AtomicLong resolverClock = new java.util.concurrent.atomic.AtomicLong();
private final CompletionResolver completion =
new CompletionResolver(agents, rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
new CompletionResolver(agents, rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none(),
() -> resolverClock.addAndGet(CompletionResolver.MIN_TURN_NANOS + 1));
private final Injector injector = new Injector(agents, completion);
private final InMemoryReplyInbox inbox = new InMemoryReplyInbox();
private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox);