Compare commits

...

2 Commits

Author SHA1 Message Date
Dai Ha 3437d6313d fleetd #241: bound the echo match so a real report is never swallowed
CI / contract (pull_request) Successful in 1m9s
CI / build (pull_request) Successful in 1m25s
Round 1 used plain bidirectional containment. The direction that catches the real bug --
the pane holds the brief plus a status bar, so the scrape contains the brief -- also fires
when a member restates the whole brief and then writes a genuine report under it. That
threw the report away and told the lead nothing was produced, which is worse than the bug
being fixed: it destroys a delivery instead of merely obscuring one.

The safe direction (the scrape is a fragment of the brief) stays unbounded, because a
fragment of the brief is by definition not a report. The dangerous direction now requires
the scrape to add at most MAX_ECHO_EXCESS_CHARS beyond the brief, which is the amount of
TUI chrome a real echo carries.

Work by the cb241 worker, committed by the lead: its backend stopped answering after the
fix was written, so two turns ended with no commit and no reply. Verified by the lead:
1169 tests, 0 failures; removing the bound turns pinsTheMaximumTuiChromeExcess and
completionFallbackKeepsARealReportThatRestatesTheWholeBrief red with 0 compile errors.
2026-09-03 11:39:11 +07:00
Dai Ha 321d8dcbb5 fleetd #241: suppress echoed fallback briefs
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Successful in 1m20s
2026-09-03 11:10:48 +07:00
7 changed files with 215 additions and 12 deletions
@@ -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();
@@ -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,17 @@ 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;
/** Normalised TUI chrome may add this many characters to an otherwise echoed brief. */
static final int MAX_ECHO_EXCESS_CHARS = 160;
/** 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 +104,11 @@ public final class CompletionResolver implements TurnListener {
private final BackendErrorPatternLookup backendErrorPatterns;
private final BackendErrorSink backendErrorSink;
private final LongSupplier nowNanos;
private final Function<String, WorktreeBranch> 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 +125,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<Rendezvous.Resolution> waiter, String baseline, long deliveredAtNanos) {
record InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline, long deliveredAtNanos,
String injectedText) {
/**
* Convenience for tests exercising scrape/suppression logic that don't care about turn
@@ -116,7 +134,11 @@ public final class CompletionResolver implements TurnListener {
* 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);
this(waiter, baseline, Long.MIN_VALUE / 2, null);
}
InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline, long deliveredAtNanos) {
this(waiter, baseline, deliveredAtNanos, null);
}
}
@@ -139,11 +161,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<String, WorktreeBranch> 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 +201,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 +214,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<String, WorktreeBranch> worktreeBranches) {
this.agents = agents;
this.rendezvous = rendezvous;
this.exhaustedPatterns = Objects.requireNonNull(exhaustedPatterns, "exhaustedPatterns");
@@ -194,6 +232,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 +262,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 +408,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 +423,40 @@ public final class CompletionResolver implements TurnListener {
}
}
/**
* A full echoed brief is at least 400 normalised characters. A scrape that contains the brief may
* add no more than 160 normalised characters of TUI chrome. This accepts harmless status text, but
* preserves a real report that restates the full brief before adding substantive content.
*/
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;
}
if (normalInjected.contains(normalScrape)) {
return true;
}
return normalScrape.contains(normalInjected)
&& normalScrape.length() - normalInjected.length() <= MAX_ECHO_EXCESS_CHARS;
}
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
@@ -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() + "]");
@@ -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.
@@ -9,12 +9,19 @@ import java.util.concurrent.CompletableFuture;
public final class TurnToken {
private final String target;
private final CompletableFuture<Rendezvous.Resolution> waiter;
private final String injectedText;
public TurnToken(String target, CompletableFuture<Rendezvous.Resolution> waiter) {
this(target, waiter, null);
}
public TurnToken(String target, CompletableFuture<Rendezvous.Resolution> waiter, String injectedText) {
this.target = target;
this.waiter = waiter;
this.injectedText = injectedText;
}
public String target() { return target; }
public CompletableFuture<Rendezvous.Resolution> waiter() { return waiter; }
public String injectedText() { return injectedText; }
}
@@ -196,6 +196,26 @@ class CompletionResolverTest {
assertEquals("complete report", waiter.getNow(null).text());
}
@Test
void suppressesABareEchoWithOnlyTuiChrome() {
String injected = "Load the implementer skill. You own fleetd #999. ".repeat(12);
String scrape = injected + "\nDev auto - GPT-5.6 Terra OpenAI";
assertTrue(CompletionResolver.echoesInjectedBrief(scrape, injected));
}
@Test
void pinsTheMaximumTuiChromeExcess() {
String injected = "a".repeat(CompletionResolver.ECHO_MIN_CHARS);
String underMargin = injected + "b".repeat(CompletionResolver.MAX_ECHO_EXCESS_CHARS);
String overMargin = injected + "b".repeat(CompletionResolver.MAX_ECHO_EXCESS_CHARS + 1);
assertTrue(CompletionResolver.echoesInjectedBrief(underMargin, injected),
"the configured excess itself remains an echoed brief");
assertFalse(CompletionResolver.echoesInjectedBrief(overMargin, injected),
"one character beyond the excess must preserve the scrape as a real report");
}
@Test
void resolvesSynchronouslyBeforePostTurnContextClearing() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ previous answer\n❯ ");
@@ -56,7 +56,11 @@ class MessageServiceTest {
/** Run {@code send} on a background thread; the current thread drives the worker's turn. */
private CompletableFuture<MessageService.Reply> sendAsync() {
return CompletableFuture.supplyAsync(() -> messages.send(T, "do the task", 5000));
return sendAsync("do the task");
}
private CompletableFuture<MessageService.Reply> 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<MessageService.Reply> 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 completionFallbackKeepsARealReportThatRestatesTheWholeBrief() throws Exception {
String brief = "Load the implementer skill. You own fleetd #999. ".repeat(12);
String report = brief + "\n\n## Report\n"
+ ("I implemented the fix in GitWorktrees.java, added five tests, ran mvn clean install "
+ "and got 1168 tests with 0 failures. Commit 321d8dc pushed. ").repeat(8);
CompletableFuture<MessageService.Reply> 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.strip(), reply.text(), "a report that restates the whole 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<MessageService.Reply> 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<MessageService.Reply> 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