diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index 9b66efc..0d327f8 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -351,6 +351,12 @@ public final class Fleetd { sessions.onTurnComplete(target); } + @Override + public void onTurnComplete(String target, long elapsedNanos) { + completion.onTurnComplete(target, elapsedNanos); + sessions.onTurnComplete(target); + } + @Override public boolean hasPostTurnAction(String target) { return sessions.hasPostTurnAction(target); @@ -362,6 +368,12 @@ public final class Fleetd { return sessions.onTurnCompleteWithPostAction(target); } + @Override + public boolean onTurnCompleteWithPostAction(String target, long elapsedNanos) { + completion.resolveBeforePostAction(target, elapsedNanos); + return sessions.onTurnCompleteWithPostAction(target); + } + @Override public void onDelivered(String target, dev.ltms.fleet.msg.TurnToken token) { completion.onDelivered(target, token); diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/CompletionResolver.java b/fleetd/src/main/java/dev/ltms/fleet/inject/CompletionResolver.java index 9d140c6..c3381e7 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/inject/CompletionResolver.java +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/CompletionResolver.java @@ -12,6 +12,7 @@ import java.util.Set; import java.util.TreeSet; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.TimeUnit; 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. */ 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 = "[Pane tail clipped: member did not call fleet_reply.]"; @@ -140,10 +150,15 @@ public final class CompletionResolver implements TurnListener { @Override 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 // it — then off-load the scrape (a herdr round-trip we must not block polling on) to a vthread. 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. */ 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 @@ -169,6 +189,11 @@ public final class CompletionResolver implements TurnListener { /** Synchronous resolve (the unit-testable core of {@link #onTurnComplete}). */ 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 waiter = turn == null ? null : turn.waiter(); if (waiter == null || waiter.isDone()) { // 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; int originalLength = 0; boolean clipped = false; - boolean scrapeFailed = false; + String raw; try { - assistantBlock = lastAssistantBlock(agents.read(target, SCRAPE_SOURCE)); + raw = agents.read(target, SCRAPE_SOURCE); + assistantBlock = lastAssistantBlock(raw); originalLength = assistantBlock.strip().length(); 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()); - tail = ""; - scrapeFailed = true; + resolveFailure(target, turn, "member " + target + " ended its turn, but its completion " + + "screen could not be read: " + e.getMessage()); + return; + } + // 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. + // 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 // 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 } - // 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; if (rendezvous.resolveCompletion(waiter, completion)) { 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 * evidence carried in a {@code BACKEND_EXHAUSTED} reason so the operator sees the real refusal diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java b/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java index 9bafae7..529efb9 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java @@ -14,6 +14,7 @@ import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.function.Consumer; +import java.util.function.LongSupplier; import java.util.function.Predicate; import java.util.stream.Collectors; @@ -92,6 +93,7 @@ public final class Injector { private final TurnListener turnListener; private final Predicate ready; // CB-113: a target is deliverable only when available private final Consumer forget; // CB-114: clear a gone worker's readiness/presence + private final LongSupplier nowNanos; private final ConcurrentHashMap targets = new ConcurrentHashMap<>(); /** 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 ready, Consumer 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 ready, + Consumer forget, LongSupplier nowNanos) { this.agents = agents; this.turnListener = turnListener; this.ready = ready; this.forget = forget; + this.nowNanos = nowNanos; } /** 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 boolean awaitingCompletion; // a delivered message's turn is not yet known-complete 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 notReadySincePoll; // consecutive injectable samples a queued message waited on the readiness gate (CB-114) boolean postTurnPending; // completion observed; adapter housekeeping has not started yet @@ -185,6 +195,7 @@ public final class Injector { Pending sent = null; RuntimeException sendError = null; boolean turnCompleted = false; + long turnElapsedNanos = 0; boolean turnFailed = false; boolean resubmit = false; boolean startPostTurn = false; @@ -202,7 +213,10 @@ public final class Injector { t.injectableSincePickup = 0; t.unknownSinceTurn = 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 t.unknownSinceTurn = 0; if (t.awaitingPostTurnPickup) { @@ -238,6 +252,7 @@ public final class Injector { if (t.awaitingCompletion && t.turnObserved) { t.awaitingCompletion = false; t.turnObserved = false; + turnElapsedNanos = Math.max(0, nowNanos.getAsLong() - t.turnStartedAtNanos); turnCompleted = true; if (turnListener.hasPostTurnAction(target)) { t.postTurnPending = true; @@ -332,7 +347,7 @@ public final class Injector { } if (turnCompleted) { if (startPostTurn) { - boolean started = turnListener.onTurnCompleteWithPostAction(target); + boolean started = turnListener.onTurnCompleteWithPostAction(target, turnElapsedNanos); synchronized (t) { t.postTurnPending = false; if (started) { @@ -344,7 +359,7 @@ public final class Injector { } } } else { - turnListener.onTurnComplete(target); + turnListener.onTurnComplete(target, turnElapsedNanos); } } if (turnFailed) { diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/TurnListener.java b/fleetd/src/main/java/dev/ltms/fleet/inject/TurnListener.java index 1ab1b35..c190e22 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/inject/TurnListener.java +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/TurnListener.java @@ -15,6 +15,14 @@ public interface TurnListener { /** A worker's delegated turn finished (worker returned to idle after visibly working). */ 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. * This is queried before the injector considers the next queued message, closing the same-tick @@ -34,6 +42,11 @@ public interface TurnListener { 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 * (CB-109) — e.g. an error screen herdr classifies as {@code unknown} — so no diff --git a/fleetd/src/test/java/dev/ltms/fleet/inject/CompletionResolverTest.java b/fleetd/src/test/java/dev/ltms/fleet/inject/CompletionResolverTest.java index 388a49b..a113f94 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/inject/CompletionResolverTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/inject/CompletionResolverTest.java @@ -249,12 +249,9 @@ 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 failsWhenTheScrapeItselfFailsEvenWithABaselinePresent() { + // A failed read means the resolver could not see the screen. It must unblock the caller as a + // failure, never turn that missing evidence into a successful empty completion. 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 +262,8 @@ 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()); + 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 --------- diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java index c0feb46..4d6bbf9 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -14,6 +14,7 @@ import org.junit.jupiter.api.Test; import java.util.concurrent.CompletableFuture; 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.assertFalse; @@ -35,7 +36,8 @@ class MessageServiceTest { private final Rendezvous rendezvous = new Rendezvous(); private final CompletionResolver completion = 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 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.WORKING); // worker picks it up and works 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 MessageService.Reply reply = send.get(5, TimeUnit.SECONDS); @@ -77,6 +80,62 @@ class MessageServiceTest { assertTrue(reply.completed(), "a scraped completion still counts as completed"); } + @Test + void emptyCompletionThroughMessageServiceFailsAndNamesTheMember() throws Exception { + CompletableFuture 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 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 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 void explicitFleetReplyResolvesAsReplied() throws Exception { CompletableFuture send = sendAsync();