diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index d57c2e7..9fdaaa2 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -390,7 +390,12 @@ public final class Fleetd { // now that `sessions` exists to resolve target -> session -> profile. exhaustionSinkRef.set(exhaustionSink); AgentControl agents = router.memberAgents(); - CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns, exhaustionSink); + CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns, exhaustionSink, + target -> sessions.roster().stream() + .filter(session -> target.equals(session.terminalId())) + .findFirst() + .map(session -> new CompletionResolver.WorktreeBranch(session.worktree(), session.branch())) + .orElse(null)); // CB-113: deliver only to an available worker (its MCP is connected), never its boot window. // CB-301: the manager's presence bridge records availability and drives SPAWNING → READY. MemberPresence presence = sessions.asPresence(); 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 f289783..8d21ba2 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/inject/CompletionResolver.java +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/CompletionResolver.java @@ -14,6 +14,7 @@ import java.util.TreeSet; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.function.LongSupplier; +import java.util.function.Function; import java.util.regex.Pattern; /** @@ -85,6 +86,14 @@ public final class CompletionResolver implements TurnListener { private static final String CLIPPED_PANE_TAIL_MARKER = "[Pane tail clipped: member did not call fleet_reply.]"; + /** A pane echo must be this large before it can replace a completion report. */ + static final int ECHO_MIN_CHARS = 400; + + /** The explicit result returned instead of a lead's echoed injected brief. */ + public static final String NO_REPORT_PREFIX = "[no report — the member ended its turn without fleet_reply, " + + "and the pane still shows the injected brief. Nothing was produced on the pane. Check the " + + "member's worktree and branch for committed work before re-delegating."; + private final AgentControl agents; private final Rendezvous rendezvous; private final ExhaustedPatternLookup exhaustedPatterns; @@ -92,6 +101,11 @@ public final class CompletionResolver implements TurnListener { private final BackendErrorPatternLookup backendErrorPatterns; private final BackendErrorSink backendErrorSink; private final LongSupplier nowNanos; + private final Function worktreeBranches; + + /** Known member location, used only to guide a lead after an echoed brief. */ + public record WorktreeBranch(String worktree, String branch) { + } /** * Per-target record of the turn currently in flight: the exact {@link Rendezvous} waiter its @@ -108,7 +122,8 @@ public final class CompletionResolver implements TurnListener { * the other half of the {@link #MIN_TURN_NANOS} floor check, compared against a fresh reading at * resolution time. */ - record InFlight(CompletableFuture waiter, String baseline, long deliveredAtNanos) { + record InFlight(CompletableFuture waiter, String baseline, long deliveredAtNanos, + String injectedText) { /** * Convenience for tests exercising scrape/suppression logic that don't care about turn @@ -116,7 +131,11 @@ public final class CompletionResolver implements TurnListener { * Not used by production code — {@link #captureBaseline} always records a real reading. */ InFlight(CompletableFuture waiter, String baseline) { - this(waiter, baseline, Long.MIN_VALUE / 2); + this(waiter, baseline, Long.MIN_VALUE / 2, null); + } + + InFlight(CompletableFuture waiter, String baseline, long deliveredAtNanos) { + this(waiter, baseline, deliveredAtNanos, null); } } @@ -139,11 +158,18 @@ public final class CompletionResolver implements TurnListener { * (fleetd#201 Unit 5). */ public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns, - ExhaustionSink exhaustionSink) { + ExhaustionSink exhaustionSink) { this(agents, rendezvous, exhaustedPatterns, exhaustionSink, BackendErrorPatternLookup.legacy(), BackendErrorSink.none(), System::nanoTime); } + /** Production constructor with a lookup for the member worktree and branch. */ + public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns, + ExhaustionSink exhaustionSink, Function worktreeBranches) { + this(agents, rendezvous, exhaustedPatterns, exhaustionSink, + BackendErrorPatternLookup.legacy(), BackendErrorSink.none(), System::nanoTime, worktreeBranches); + } + /** * Transition constructor (fleetd#201 Unit 1): same legacy backend-error defaults as the 4-arg * constructor above, but with the injectable clock. Kept so existing fleetd#164 timing tests @@ -172,7 +198,7 @@ public final class CompletionResolver implements TurnListener { ExhaustionSink exhaustionSink, BackendErrorPatternLookup backendErrorPatterns, BackendErrorSink backendErrorSink) { this(agents, rendezvous, exhaustedPatterns, exhaustionSink, backendErrorPatterns, backendErrorSink, - System::nanoTime); + System::nanoTime, _ -> null); } /** @@ -185,8 +211,17 @@ public final class CompletionResolver implements TurnListener { * package. */ public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns, - ExhaustionSink exhaustionSink, BackendErrorPatternLookup backendErrorPatterns, - BackendErrorSink backendErrorSink, LongSupplier nowNanos) { + ExhaustionSink exhaustionSink, BackendErrorPatternLookup backendErrorPatterns, + BackendErrorSink backendErrorSink, LongSupplier nowNanos) { + this(agents, rendezvous, exhaustedPatterns, exhaustionSink, backendErrorPatterns, backendErrorSink, + nowNanos, _ -> null); + } + + /** Full constructor with injectable clock and member location lookup. */ + public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns, + ExhaustionSink exhaustionSink, BackendErrorPatternLookup backendErrorPatterns, + BackendErrorSink backendErrorSink, LongSupplier nowNanos, + Function worktreeBranches) { this.agents = agents; this.rendezvous = rendezvous; this.exhaustedPatterns = Objects.requireNonNull(exhaustedPatterns, "exhaustedPatterns"); @@ -194,6 +229,7 @@ public final class CompletionResolver implements TurnListener { this.backendErrorPatterns = Objects.requireNonNull(backendErrorPatterns, "backendErrorPatterns"); this.backendErrorSink = Objects.requireNonNull(backendErrorSink, "backendErrorSink"); this.nowNanos = Objects.requireNonNull(nowNanos, "nowNanos"); + this.worktreeBranches = Objects.requireNonNull(worktreeBranches, "worktreeBranches"); } @Override @@ -223,7 +259,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, nowNanos.getAsLong())); + inFlight.put(target, new InFlight(waiter, baseline, nowNanos.getAsLong(), token.injectedText())); } /** The turn currently baselined for {@code target}, or {@code null} — a test hook for the captureBaseline path. */ @@ -369,6 +405,9 @@ public final class CompletionResolver implements TurnListener { return; } String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail; + if (echoesInjectedBrief(tail, turn.injectedText())) { + completion = noReportMessage(target) + (clipped ? "\n" + CLIPPED_PANE_TAIL_MARKER : ""); + } if (rendezvous.resolveCompletion(waiter, completion)) { inFlight.remove(target, turn); if (clipped) { @@ -381,6 +420,36 @@ public final class CompletionResolver implements TurnListener { } } + /** + * A full echoed brief is at least 400 normalised characters and one normalised value contains the + * other. This accepts harmless TUI whitespace and punctuation changes, but preserves a real report + * that quotes only one part of the brief. + */ + static boolean echoesInjectedBrief(String scrape, String injectedText) { + String normalScrape = normalize(scrape); + String normalInjected = normalize(injectedText); + if (normalScrape.length() < ECHO_MIN_CHARS || normalInjected.length() < ECHO_MIN_CHARS) { + return false; + } + return normalInjected.contains(normalScrape) || normalScrape.contains(normalInjected); + } + + private static String normalize(String text) { + return text == null ? "" : text.toLowerCase().replaceAll("[^a-z0-9]+", ""); + } + + private String noReportMessage(String target) { + WorktreeBranch location = worktreeBranches.apply(target); + if (location == null || (location.worktree() == null && location.branch() == null)) { + return NO_REPORT_PREFIX + "]"; + } + String locationText = location.worktree() == null ? "" : " worktree=" + location.worktree(); + if (location.branch() != null) { + locationText += " branch=" + location.branch(); + } + return NO_REPORT_PREFIX + locationText + "]"; + } + /** * fleetd#211: the raw-scrape fallback classification, run only when {@link #lastAssistantBlock} * found nothing usable (see the call site in {@link #resolve}). Mirrors the two classifications diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index 4b664e3..7620b5b 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -10,6 +10,7 @@ import dev.ltms.fleet.herdr.Agent; import dev.ltms.fleet.metrics.FleetMetrics; import dev.ltms.fleet.metrics.Metrics; import dev.ltms.fleet.inject.MemberPresence; +import dev.ltms.fleet.inject.CompletionResolver; import dev.ltms.fleet.herdr.HerdrException; import dev.ltms.fleet.msg.LeadChannel; import dev.ltms.fleet.msg.LeadMessage; @@ -543,8 +544,9 @@ public final class FleetMcp { case REPLIED -> text(r.text()); // The worker's turn finished but it never called fleet_reply — hand back the scraped // transcript tail, flagged so the primary knows it isn't a structured reply. - case COMPLETED_UNREPLIED -> text( - "[worker finished without a structured fleet_reply — transcript tail follows]\n" + r.text()); + case COMPLETED_UNREPLIED -> text(r.text().startsWith(CompletionResolver.NO_REPORT_PREFIX) + ? r.text() + : "[worker finished without a structured fleet_reply — transcript tail follows]\n" + r.text()); // The worker ran the turn then wedged (CB-109) — surface the error context. case WORKER_FAILED -> text("[worker failed — turn ended in an unrecoverable state]\n" + r.text()); // The backend refused on a subscription usage limit (CB-578 stage A) — the worker's @@ -659,6 +661,7 @@ public final class FleetMcp { } return switch (v.phase()) { case DONE -> text(v.replySource() != null && v.replySource().equals("transcript") + && !v.reply().startsWith(CompletionResolver.NO_REPORT_PREFIX) ? "[done — worker finished without a structured fleet_reply; transcript tail follows]\n" + v.reply() : v.reply()); case PENDING -> text("[pending — " + v.detail() + "]"); diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java index b7baab1..7138991 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -746,7 +746,7 @@ public final class MessageService { if (task != null) { asyncTasksByWaiter.put(reply, task); } - TurnToken token = new TurnToken(target, reply); + TurnToken token = new TurnToken(target, reply, content); // The send has won the lock; the accepted-delivery hook records delegator ownership // here (CB-548). It runs BEFORE enqueue so a throwing hook — onAccepted is now a // public callback — fails the send without queuing a message that would orphan. diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/TurnToken.java b/fleetd/src/main/java/dev/ltms/fleet/msg/TurnToken.java index bd1d839..574cbff 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/TurnToken.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/TurnToken.java @@ -9,12 +9,19 @@ import java.util.concurrent.CompletableFuture; public final class TurnToken { private final String target; private final CompletableFuture waiter; + private final String injectedText; public TurnToken(String target, CompletableFuture waiter) { + this(target, waiter, null); + } + + public TurnToken(String target, CompletableFuture waiter, String injectedText) { this.target = target; this.waiter = waiter; + this.injectedText = injectedText; } public String target() { return target; } public CompletableFuture waiter() { return waiter; } + public String injectedText() { return injectedText; } } 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 95f1e92..f62684c 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -56,7 +56,11 @@ class MessageServiceTest { /** Run {@code send} on a background thread; the current thread drives the worker's turn. */ private CompletableFuture sendAsync() { - return CompletableFuture.supplyAsync(() -> messages.send(T, "do the task", 5000)); + return sendAsync("do the task"); + } + + private CompletableFuture sendAsync(String content) { + return CompletableFuture.supplyAsync(() -> messages.send(T, content, 5000)); } private void awaitWaiting() throws InterruptedException { @@ -86,6 +90,94 @@ class MessageServiceTest { assertTrue(reply.completed(), "a scraped completion still counts as completed"); } + @Test + void completionFallbackReplacesAnEchoedInjectedBriefWithNoReportOutcome() throws Exception { + String brief = "Implement the requested change. ".repeat(20); + CompletableFuture send = sendAsync(brief); + awaitWaiting(); + + herdr.readText("$ prompt"); + injector.onStatus(T, AgentStatus.IDLE); + injector.onStatus(T, AgentStatus.WORKING); + herdr.readText("⏺ " + brief + "\n❯ "); + injector.onStatus(T, AgentStatus.IDLE); + + MessageService.Reply reply = send.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.COMPLETED_UNREPLIED, reply.outcome()); + assertEquals(CompletionResolver.NO_REPORT_PREFIX + "]", reply.text(), + "the real injector -> completion fallback path must not return the lead's brief"); + } + + @Test + void completionFallbackKeepsARealReportThatQuotesPartOfTheBrief() throws Exception { + String quotedAcceptance = "Acceptance: the new test passes and the build is green. ".repeat(4); + String brief = quotedAcceptance + "Implement the requested change. ".repeat(20); + String report = "I changed the fallback and added tests. " + quotedAcceptance + + "The final build passed."; + CompletableFuture send = sendAsync(brief); + awaitWaiting(); + + herdr.readText("$ prompt"); + injector.onStatus(T, AgentStatus.IDLE); + injector.onStatus(T, AgentStatus.WORKING); + herdr.readText("⏺ " + report + "\n❯ "); + injector.onStatus(T, AgentStatus.IDLE); + + MessageService.Reply reply = send.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.COMPLETED_UNREPLIED, reply.outcome()); + assertEquals(report, reply.text(), "a report that quotes part of the brief must survive unchanged"); + } + + @Test + void completionFallbackNamesTheKnownWorktreeAndBranchForAnEchoedBrief() throws Exception { + FakeHerdr localHerdr = new FakeHerdr(); + Rendezvous localRendezvous = new Rendezvous(); + java.util.concurrent.atomic.AtomicLong localClock = new java.util.concurrent.atomic.AtomicLong(); + CompletionResolver localCompletion = new CompletionResolver(new AgentControl(localHerdr), localRendezvous, + ExhaustedPatternLookup.none(), ExhaustionSink.none(), + dev.ltms.fleet.inject.BackendErrorPatternLookup.legacy(), + dev.ltms.fleet.inject.BackendErrorSink.none(), + () -> localClock.addAndGet(CompletionResolver.MIN_TURN_NANOS + 1), + target -> new CompletionResolver.WorktreeBranch("/tmp/member-worktree", "worker/cb241")); + Injector localInjector = new Injector(new AgentControl(localHerdr), localCompletion); + MessageService localMessages = new MessageService(new AgentControl(localHerdr), localInjector, + localRendezvous, new InMemoryReplyInbox()); + String brief = "Implement the requested change. ".repeat(20); + CompletableFuture send = + CompletableFuture.supplyAsync(() -> localMessages.send(T, brief, 5000)); + long deadline = System.currentTimeMillis() + 2000; + while (!localRendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + } + assertTrue(localRendezvous.isWaiting(T)); + + localHerdr.readText("$ prompt"); + localInjector.onStatus(T, AgentStatus.IDLE); + localInjector.onStatus(T, AgentStatus.WORKING); + localHerdr.readText("⏺ " + brief + "\n❯ "); + localInjector.onStatus(T, AgentStatus.IDLE); + + assertEquals(CompletionResolver.NO_REPORT_PREFIX + " worktree=/tmp/member-worktree " + + "branch=worker/cb241]", send.get(5, TimeUnit.SECONDS).text()); + } + + @Test + void completionFallbackKeepsTheClippedMarkerWhenAnEchoedBriefIsTooLong() throws Exception { + String brief = "a".repeat(4_001); // CompletionResolver's 4,000-character scrape cap + CompletableFuture send = sendAsync(brief); + awaitWaiting(); + + herdr.readText("$ prompt"); + injector.onStatus(T, AgentStatus.IDLE); + injector.onStatus(T, AgentStatus.WORKING); + herdr.readText("⏺ " + brief + "\n❯ "); + injector.onStatus(T, AgentStatus.IDLE); + + assertEquals(CompletionResolver.NO_REPORT_PREFIX + "]\n" + + "[Pane tail clipped: member did not call fleet_reply.]", + send.get(5, TimeUnit.SECONDS).text()); + } + @Test void backendErrorScrapeThroughMessageServiceFailsInsteadOfBecomingReplyText() throws Exception { // fleetd#164 (part 2 addendum): a scrape that reads cleanly but is only the backend's own