From 31e34d177b1af1585d2a544a83d693e5b6c5cec3 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 16 Jul 2026 08:27:38 +0200 Subject: [PATCH] =?UTF-8?q?CB-115/CB-116:=20reliable=20turn=20completion?= =?UTF-8?q?=20=E2=80=94=20status=20refinement,=20clean=20scrape,=20waiter?= =?UTF-8?q?=20identity?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A 5-turn primary↔worker conversation test (see the e2e harness) surfaced three delegation-channel gaps; this closes them. CB-115 — status + scrape correctness: - AgentStatus gains DONE (herdr's explicit turn-complete marker) so a finished turn is no longer misread as UNKNOWN and left to wedge or false-fail. - StatusRefiner reclassifies content-bearing UNKNOWN samples (StatusPoller wired to it), and the Injector baselines pane content on delivery (TurnListener gains onDelivered) to guard completion against previous-turn misattribution. - CompletionResolver.lastAssistantBlock stops at the first hard TUI boundary, so a scrape returns only the assistant answer — no input box, prompt echo, spinner, tips or warnings. CB-116 — waiter identity (the cross-turn stale reply): - The completion/failure fallback ran on a virtual thread and resolved whichever waiter was currently registered for the session. Since the rendezvous holds one waiter per session and sends serialize, turn N's late completion could land on turn N+1's waiter and deliver turn N's stale scrape as turn N+1's answer. The baseline guard missed it because turn N was resolved by bridge_reply, which never updates the completion baseline. - Fix: capture the exact waiter (and pre-turn baseline) when a turn is delivered, on the poller thread before any next-turn delivery can overwrite it, and resolve THAT waiter — a no-op if it was already resolved. Rendezvous.resolveCompletion/ resolveFailure now take the captured CompletableFuture; currentWaiter exposes the registered one for capture. A late completion for turn N can no longer touch turn N+1's send. Verified: 128 unit tests green (incl. a CB-116 regression asserting a late completion never resolves the next turn's waiter); a re-run of the conversation test passes with turn 5 resolving to its own reply rather than turn 4's text. --- .../dev/ltms/bridged/herdr/AgentStatus.java | 18 +- .../bridged/inject/CompletionResolver.java | 167 +++++++++++++++-- .../dev/ltms/bridged/inject/Injector.java | 3 + .../dev/ltms/bridged/inject/StatusPoller.java | 11 +- .../ltms/bridged/inject/StatusRefiner.java | 89 +++++++++ .../dev/ltms/bridged/inject/TurnListener.java | 11 ++ .../java/dev/ltms/bridged/msg/Rendezvous.java | 51 ++++-- .../ltms/bridged/worker/WorkerService.java | 14 +- .../ltms/bridged/herdr/AgentStatusTest.java | 43 +++++ .../inject/CompletionResolverTest.java | 170 +++++++++++++++++- .../bridged/inject/StatusRefinerTest.java | 88 +++++++++ .../ltms/bridged/msg/MessageServiceTest.java | 4 +- 12 files changed, 629 insertions(+), 40 deletions(-) create mode 100644 bridged/src/main/java/dev/ltms/bridged/inject/StatusRefiner.java create mode 100644 bridged/src/test/java/dev/ltms/bridged/herdr/AgentStatusTest.java create mode 100644 bridged/src/test/java/dev/ltms/bridged/inject/StatusRefinerTest.java diff --git a/bridged/src/main/java/dev/ltms/bridged/herdr/AgentStatus.java b/bridged/src/main/java/dev/ltms/bridged/herdr/AgentStatus.java index 36e712a..e04dbc5 100644 --- a/bridged/src/main/java/dev/ltms/bridged/herdr/AgentStatus.java +++ b/bridged/src/main/java/dev/ltms/bridged/herdr/AgentStatus.java @@ -2,28 +2,38 @@ package dev.ltms.bridged.herdr; /** * A herdr agent's lifecycle state, as reported by {@code agent_status}. Drives the - * status-gated injector: a worker is safe to inject into only when {@link #IDLE} or - * {@link #BLOCKED}, never mid-turn ({@link #WORKING}). + * status-gated injector: a worker is safe to inject into only when {@link #IDLE}, + * {@link #BLOCKED}, or {@link #DONE}, never mid-turn ({@link #WORKING}). */ public enum AgentStatus { IDLE, WORKING, BLOCKED, + /** + * The worker has finished its turn and is settled at an idle prompt. herdr emits this + * (observed live alongside {@code idle}) as a turn-complete marker; earlier code mapped the + * unrecognized string to {@link #UNKNOWN}, which both wedged delivery (not {@link #injectable}) + * and mis-fired the CB-109 stall-failure on a worker that had actually answered. It is a + * turn-boundary equivalent to {@link #IDLE}: injectable, and a {@code working → done} edge is a + * real completion. + */ + DONE, UNKNOWN; - /** Map herdr's wire string ({@code idle|working|blocked|unknown}) to the enum. */ + /** Map herdr's wire string ({@code idle|working|blocked|done|unknown}) to the enum. */ public static AgentStatus fromWire(String s) { if (s == null) return UNKNOWN; return switch (s.toLowerCase()) { case "idle" -> IDLE; case "working" -> WORKING; case "blocked" -> BLOCKED; + case "done" -> DONE; default -> UNKNOWN; }; } /** Whether {@code bridged} may inject a message now without stepping on a live turn. */ public boolean injectable() { - return this == IDLE || this == BLOCKED; + return this == IDLE || this == BLOCKED || this == DONE; } } diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java index 3f8dc2f..8c3e0b1 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java @@ -5,6 +5,9 @@ import dev.ltms.bridged.msg.Rendezvous; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentHashMap; + /** * The CB-106 completion fallback: bridges the {@link Injector}'s turn-completion signal to the * {@link Rendezvous} so a blocking {@code bridge_send} resolves even when the worker finishes its @@ -20,8 +23,23 @@ import org.slf4j.LoggerFactory; * wedged in an {@code unknown} state resolves the send as a failure (with the error screen as * context) rather than leaving it to time out. * - *

Wired as the {@link Injector}'s {@link TurnListener}; both handlers hand off to a virtual thread - * so the scrape's herdr round-trip never stalls the status poller. + *

The scrape is cleaned to the last {@code ⏺} assistant block (stripping TUI chrome) and guarded + * against misattribution (CB-115): the pane content is baselined on delivery ({@link #onDelivered}), + * and a completion whose scrape is unchanged from that baseline — the previous turn's wind-down + * sampled as this turn's boundary on a rapid back-to-back send — is suppressed rather than resolving + * the send with a stale answer. + * + *

Waiter-specific resolution (CB-116). On delivery we also capture the exact + * {@link Rendezvous} waiter this turn belongs to, and the completion/failure fallbacks resolve + * that waiter — never "whatever send is waiting now". A completion fallback runs on a virtual + * thread and can land after the worker's {@code bridge_reply} already resolved the turn and the + * next send opened its own waiter on the same session; resolving the current waiter would + * then deliver turn N's stale scrape as turn N+1's answer. Targeting the captured waiter makes a late + * completion a harmless no-op (its waiter is already done) instead of a cross-turn stale reply. + * + *

Wired as the {@link Injector}'s {@link TurnListener}; the handlers hand off to a virtual thread + * so the scrape's herdr round-trip never stalls the status poller. The captured waiter is read on the + * poller thread (before any next-turn delivery can overwrite it) and passed into the virtual thread. */ public final class CompletionResolver implements TurnListener { @@ -40,47 +58,116 @@ public final class CompletionResolver implements TurnListener { private final AgentControl agents; private final Rendezvous rendezvous; + /** + * Per-target record of the turn currently in flight: the exact {@link Rendezvous} waiter its + * delivering send opened, plus the assistant block present when it was delivered. + * + *

The {@code waiter} is what makes a late fallback safe (CB-116): we resolve it, not "whoever + * is waiting now", so a completion that fires after the next send has opened its own waiter is a + * no-op rather than a cross-turn stale reply. The {@code baseline} is the CB-115 staleness + * 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. + */ + record InFlight(CompletableFuture waiter, String baseline) { + } + + private final ConcurrentHashMap inFlight = new ConcurrentHashMap<>(); + public CompletionResolver(AgentControl agents, Rendezvous rendezvous) { this.agents = agents; this.rendezvous = rendezvous; } + @Override + public void onDelivered(String target) { + // Capture the exact waiter this turn belongs to (CB-116) and snapshot the pane's pre-turn + // content — what it shows *before* the just-delivered turn produces output — as the staleness + // reference (CB-115). Done synchronously (like the delivering send itself) so both are in + // place before this turn's completion can fire. + captureBaseline(target); + } + + /** Capture the in-flight turn: its waiter and pre-turn baseline (the testable core of {@link #onDelivered}). */ + void captureBaseline(String target) { + CompletableFuture waiter = rendezvous.currentWaiter(target); + if (waiter == null) { + inFlight.remove(target); // no send is waiting on this delivery — nothing to resolve later + return; + } + String baseline; + try { + baseline = lastAssistantBlock(agents.read(target, SCRAPE_SOURCE)); + } catch (RuntimeException e) { + baseline = null; // fail open: no baseline ⇒ no suppression + log.debug("delivery baseline for {} failed: {}", target, e.getMessage()); + } + inFlight.put(target, new InFlight(waiter, baseline)); + } + @Override public void onTurnComplete(String target) { - // Off the poller thread: the scrape is a herdr round-trip we must not block polling on. - Thread.ofVirtual().name("completion-" + target).start(() -> resolve(target)); + // 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)); } @Override public void onTurnFailed(String target) { - Thread.ofVirtual().name("turn-failed-" + target).start(() -> fail(target)); + InFlight turn = inFlight.get(target); + Thread.ofVirtual().name("turn-failed-" + target).start(() -> fail(target, turn)); } /** Synchronous resolve (the unit-testable core of {@link #onTurnComplete}). */ - void resolve(String target) { - if (!rendezvous.isWaiting(target)) { - return; // nobody is blocked on this worker's turn — nothing to resolve, skip the scrape + void resolve(String target, InFlight turn) { + CompletableFuture waiter = turn == null ? null : turn.waiter(); + if (waiter == null || waiter.isDone()) { + // Nobody is blocked on THIS turn (it had no send, or its bridge_reply already won). Skip + // the scrape; resolving the current waiter here would be the CB-116 cross-turn stale reply. + inFlight.remove(target, turn); + return; } String tail; + boolean scrapeFailed = false; try { - tail = clip(agents.read(target, SCRAPE_SOURCE)); + tail = clip(lastAssistantBlock(agents.read(target, SCRAPE_SOURCE))); } 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; } - if (rendezvous.resolveCompletion(target, tail)) { + // 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 bridge_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)) { + 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 + } + if (rendezvous.resolveCompletion(waiter, tail)) { + inFlight.remove(target, turn); log.debug("resolved send to {} via turn-completion fallback ({} chars scraped)", target, tail.length()); } } /** Synchronous fail (the unit-testable core of {@link #onTurnFailed}). */ - void fail(String target) { - if (!rendezvous.isWaiting(target)) { - return; // nobody blocked on this worker — nothing to fail + void fail(String target, InFlight turn) { + // A never-delivered readiness failure has no in-flight record but still has a blocked send; + // fall back to the currently-registered waiter (unambiguous — that send never completed, so + // no next turn exists to confuse it with). + CompletableFuture waiter = + turn != null ? turn.waiter() : rendezvous.currentWaiter(target); + if (waiter == null || waiter.isDone()) { + inFlight.remove(target, turn); // nobody blocked on this worker — nothing to fail + return; } String reason; try { @@ -93,7 +180,8 @@ public final class CompletionResolver implements TurnListener { reason = "worker did not reply; its turn ended in an unrecoverable state " + "(worker unreachable or stuck)"; } - if (rendezvous.resolveFailure(target, reason)) { + if (rendezvous.resolveFailure(waiter, reason)) { + inFlight.remove(target, turn); log.debug("failed send to {} via turn-stall fallback", target); } } @@ -105,4 +193,55 @@ public final class CompletionResolver implements TurnListener { ? trimmed : trimmed.substring(trimmed.length() - MAX_SCRAPE_CHARS); } + + /** + * Extract the last assistant message from a raw Claude Code pane scrape (CB-115). Claude Code + * prefixes each assistant turn with {@code ⏺}; the delegator wants that answer, not the TUI + * chrome around it. Take everything from the final {@code ⏺} onward and stop at the first + * hard interface boundary below it — the spinner/status line, input box, {@code ❯} prompt (which + * may echo the next turn's text), footer, or tips/warnings. Stopping at the first + * boundary (rather than trimming only trailing chrome) is what keeps a following turn's echoed + * prompt out of this reply. Blank lines are not boundaries, so a multi-paragraph answer survives; + * trailing blanks are trimmed at the end. With no {@code ⏺} marker (an unusual render) the whole + * text is scanned the same way, so we never lose the reply. + * + *

Package-private and pure so it is unit-testable without herdr. + */ + static String lastAssistantBlock(String raw) { + if (raw == null || raw.isBlank()) return ""; + int marker = raw.lastIndexOf('⏺'); + String block = marker >= 0 ? raw.substring(marker + 1) : raw; + StringBuilder out = new StringBuilder(); + int kept = 0; + for (String line : block.split("\n", -1)) { + if (isBoundary(line)) break; // first TUI boundary ends the assistant message + if (kept++ > 0) out.append('\n'); + out.append(line); + } + return out.toString().strip(); + } + + /** + * A hard TUI boundary line that marks the end of an assistant message and the start of interface + * chrome (input box, prompt, spinner, footer, tips/warnings). Blank lines are not + * boundaries — an answer may contain them — so they are kept and trimmed only if trailing. + */ + private static boolean isBoundary(String line) { + String t = line.strip(); + if (t.isEmpty()) return false; + // A horizontal rule / all box-drawing separators (e.g. "──────"). + if (t.chars().allMatch(c -> c == '─' || c == '—' || c == '━' || c == '═' || c == '-')) { + return true; + } + String lower = t.toLowerCase(); + return t.startsWith("╭") || t.startsWith("│") || t.startsWith("╰") || t.startsWith("┌") + || t.startsWith("└") || t.startsWith("❯") || t.startsWith("⏵") + || t.startsWith("⎿") || t.startsWith("⚠") + // Status/spinner lines Claude Code renders below a settled or in-flight turn, + // e.g. "✻ Baked for 21s", "✶ Forming…". + || t.startsWith("✻") || t.startsWith("✳") || t.startsWith("✽") || t.startsWith("·") + || t.startsWith("●") || t.startsWith("◐") || t.startsWith("✢") || t.startsWith("✶") + || lower.contains("auto mode") || lower.contains("for shortcuts") + || lower.contains("esc to interrupt") || lower.contains("bypass permissions"); + } } diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java b/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java index 09b4056..20e61c6 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java @@ -300,6 +300,9 @@ public final class Injector { log.warn("inject to {} failed, dropped message: {}", target, sendError.getMessage()); sent.delivered().completeExceptionally(sendError); } else { + // Baseline the pane's pre-turn content so a misattributed completion (no new output) + // can't resolve this send with the previous turn's stale answer (CB-115). + turnListener.onDelivered(target); sent.delivered().complete(null); } } diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/StatusPoller.java b/bridged/src/main/java/dev/ltms/bridged/inject/StatusPoller.java index 43dfa64..d3592a6 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/StatusPoller.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/StatusPoller.java @@ -23,13 +23,20 @@ public final class StatusPoller { private final AgentControl agents; private final Injector injector; + private final StatusRefiner refiner; private final long intervalMillis; private volatile boolean running; private Thread thread; public StatusPoller(AgentControl agents, Injector injector, long intervalMillis) { + this(agents, injector, new StatusRefiner(agents), intervalMillis); + } + + public StatusPoller(AgentControl agents, Injector injector, StatusRefiner refiner, + long intervalMillis) { this.agents = agents; this.injector = injector; + this.refiner = refiner; this.intervalMillis = intervalMillis; } @@ -47,7 +54,9 @@ public final class StatusPoller { for (String target : active) { if (!running) return; try { - AgentStatus status = agents.status(target); + // herdr's agent_status can misreport a settled worker as `unknown`; refine it + // against the pane content before it drives delivery/completion (CB-115). + AgentStatus status = refiner.refine(target, agents.status(target)); injector.onStatus(target, status); } catch (HerdrException e) { // The worker's agent is gone — stop trying and unblock its waiters. diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/StatusRefiner.java b/bridged/src/main/java/dev/ltms/bridged/inject/StatusRefiner.java new file mode 100644 index 0000000..bc7551f --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/inject/StatusRefiner.java @@ -0,0 +1,89 @@ +package dev.ltms.bridged.inject; + +import dev.ltms.bridged.herdr.AgentControl; +import dev.ltms.bridged.herdr.AgentStatus; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Refines an unreliable {@link AgentStatus#UNKNOWN} into a real state by reading the worker's + * terminal content (CB-115). + * + *

Some workers' panes are misclassified by herdr as {@code unknown} even when the worker is + * plainly settled at an idle prompt (empty {@code ❯}, "auto mode on" footer, a completed + * {@code ⏺} answer above). Left as {@code UNKNOWN} that both wedges delivery — the + * status-gated {@link Injector} only injects into an {@link AgentStatus#injectable} worker — and + * mis-fires the CB-109 stall failure on a worker that has actually answered. herdr's + * {@code agent_status} is a heuristic; the pane content is the ground truth. + * + *

The refinement only ever runs on a raw {@code UNKNOWN} sample (every other status is trusted + * as-is), so a healthy worker adds zero extra herdr traffic; a persistently-{@code unknown} worker + * costs one extra {@code agent.read} per poll while it has work outstanding. Classification is + * deliberately conservative — it upgrades {@code UNKNOWN} to {@link AgentStatus#WORKING} or + * {@link AgentStatus#IDLE} only on a clear signal, and leaves a genuinely unclassifiable screen + * (e.g. a wedged error state) as {@code UNKNOWN} so the CB-109 stall path can still fail it. + */ +public final class StatusRefiner { + + private static final Logger log = LoggerFactory.getLogger(StatusRefiner.class); + + /** + * herdr {@code agent.read} source used to inspect the pane. {@code detection} is the region + * herdr itself uses for status detection (the prompt/footer tail), which is exactly what we + * need to tell "idle at prompt" from "mid-turn". + */ + static final String PROBE_SOURCE = "detection"; + + private final AgentControl agents; + + public StatusRefiner(AgentControl agents) { + this.agents = agents; + } + + /** + * Return a trustworthy status for {@code target}. Any non-{@code UNKNOWN} {@code raw} is returned + * unchanged; an {@code UNKNOWN} triggers a pane read and content classification. A read failure + * leaves it {@code UNKNOWN} (the safe default: no delivery, and the stall path still applies). + */ + public AgentStatus refine(String target, AgentStatus raw) { + if (raw != AgentStatus.UNKNOWN) return raw; + String pane; + try { + pane = agents.read(target, PROBE_SOURCE); + } catch (RuntimeException e) { + log.debug("status refine read for {} failed; leaving UNKNOWN: {}", target, e.getMessage()); + return AgentStatus.UNKNOWN; + } + AgentStatus refined = classify(pane); + if (refined != AgentStatus.UNKNOWN) { + log.debug("refined {} from UNKNOWN to {} via pane content", target, refined); + } + return refined; + } + + /** + * Classify a Claude Code TUI pane tail. Package-private and pure so it is unit-testable without + * herdr. + * + *

+ */ + static AgentStatus classify(String pane) { + if (pane == null || pane.isBlank()) return AgentStatus.UNKNOWN; + String lower = pane.toLowerCase(); + // Claude Code shows "(esc to interrupt)" only while a turn is actively generating. + if (lower.contains("esc to interrupt")) return AgentStatus.WORKING; + // A settled, ready input prompt with no active-turn marker = idle-at-prompt. + boolean readyPrompt = pane.contains("❯") + || pane.contains("│ >") + || lower.contains("auto mode on") + || lower.contains("? for shortcuts"); + return readyPrompt ? AgentStatus.IDLE : AgentStatus.UNKNOWN; + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java b/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java index e09478b..b1e752b 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java @@ -22,6 +22,17 @@ public interface TurnListener { default void onTurnFailed(String target) { } + /** + * A message was just delivered into {@code target}'s pane (CB-115). Fired so the completion + * resolver can snapshot the pane's pre-turn content: a later {@link #onTurnComplete} whose + * scrape is unchanged from this baseline is a misattributed boundary (e.g. the prior + * turn's wind-down sampled as this turn's completion on rapid back-to-back sends) and must not + * resolve the send with the previous turn's stale answer. A default no-op keeps the interface + * functional for callers that don't scrape. + */ + default void onDelivered(String target) { + } + /** No-op default for callers that only need delivery, not completion signalling. */ TurnListener NOOP = _ -> { }; diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java b/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java index f7db678..1cc023c 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java @@ -12,6 +12,15 @@ import java.util.concurrent.ConcurrentHashMap; * *

At most one waiter per session — {@link MessageService} serializes sends per session, so a * resolution maps unambiguously to the one outstanding send and cannot be captured by another. + * + *

Waiter identity (CB-116). The completion/failure fallbacks run asynchronously + * and can fire after the turn they belong to has already been resolved by an explicit reply + * and a next send has opened its own waiter on the same session. Resolving "whatever waiter + * is registered now" would then land turn N's stale scrape on turn N+1's send. So those fallbacks + * resolve a specific {@link CompletableFuture} captured when their turn was delivered + * ({@link #resolveCompletion(CompletableFuture, String)} / + * {@link #resolveFailure(CompletableFuture, String)}): a no-op if that waiter was already resolved, + * and it can never touch a later send's waiter. */ public final class Rendezvous { @@ -31,8 +40,11 @@ public final class Rendezvous { private final ConcurrentHashMap> waiters = new ConcurrentHashMap<>(); - /** Register a waiter for {@code session}. The caller must hold that session's send lock. */ - CompletableFuture open(String session) { + /** + * Register a waiter for {@code session} — the await side of the public {@code resolve*} methods. + * The caller must hold that session's send lock. + */ + public CompletableFuture open(String session) { CompletableFuture waiter = new CompletableFuture<>(); waiters.put(session, waiter); return waiter; @@ -48,6 +60,15 @@ public final class Rendezvous { return waiters.containsKey(session); } + /** + * The waiter currently registered for {@code session}, or {@code null} if none is waiting. The + * completion/failure fallbacks capture this at delivery time so they can later resolve that exact + * send (see the CB-116 note above) rather than whichever send happens to be waiting when they fire. + */ + public CompletableFuture currentWaiter(String session) { + return waiters.get(session); + } + /** * Resolve the send awaiting on {@code session} with the worker's explicit reply {@code content}. * @@ -59,25 +80,27 @@ public final class Rendezvous { } /** - * Resolve the send awaiting on {@code session} as a completion (the delegated turn finished with - * no {@code bridge_reply}); {@code text} is the scraped transcript tail. A no-op if the worker - * also raced an explicit {@code bridge_reply} in — first resolution wins. + * Resolve a specific captured {@code waiter} as a completion (the delegated turn finished with no + * {@code bridge_reply}); {@code text} is the scraped transcript tail. The waiter is the one + * captured when this turn was delivered, so a late completion for turn N cannot land on turn N+1's + * send (CB-116). A no-op if that waiter was already resolved — a raced {@code bridge_reply} wins. * - * @return {@code true} if a waiter was resolved, {@code false} if none was waiting + * @return {@code true} if this call resolved the waiter, {@code false} if it was null or already resolved */ - public boolean resolveCompletion(String session, String text) { - return complete(session, new Resolution(Kind.COMPLETION, text)); + public boolean resolveCompletion(CompletableFuture waiter, String text) { + return waiter != null && waiter.complete(new Resolution(Kind.COMPLETION, text)); } /** - * Resolve the send awaiting on {@code session} as a failure — the worker ran the turn but wedged - * in an unrecoverable state (CB-109); {@code reason} is the failure context (e.g. the error - * screen). A no-op if the worker had already raced a reply/completion in — first resolution wins. + * Resolve a specific captured {@code waiter} as a failure — the worker ran the turn but wedged in + * an unrecoverable state (CB-109); {@code reason} is the failure context (e.g. the error screen). + * Like {@link #resolveCompletion(CompletableFuture, String)} it targets the exact captured send + * (CB-116). A no-op if that waiter was already resolved — first resolution wins. * - * @return {@code true} if a waiter was resolved, {@code false} if none was waiting + * @return {@code true} if this call resolved the waiter, {@code false} if it was null or already resolved */ - public boolean resolveFailure(String session, String reason) { - return complete(session, new Resolution(Kind.FAILED, reason)); + public boolean resolveFailure(CompletableFuture waiter, String reason) { + return waiter != null && waiter.complete(new Resolution(Kind.FAILED, reason)); } private boolean complete(String session, Resolution resolution) { diff --git a/bridged/src/main/java/dev/ltms/bridged/worker/WorkerService.java b/bridged/src/main/java/dev/ltms/bridged/worker/WorkerService.java index 2003523..171124e 100644 --- a/bridged/src/main/java/dev/ltms/bridged/worker/WorkerService.java +++ b/bridged/src/main/java/dev/ltms/bridged/worker/WorkerService.java @@ -56,10 +56,16 @@ public final class WorkerService { * profile — it is guidance, and a worker that never replies is caught by the send's timeout. */ static final String REPLY_CHARTER = - "You are a worker in the claude-bridge fleet. You receive one delegated task at a time " - + "as a prompt. When you finish the task, immediately call the bridge_reply MCP tool " - + "with `content` set to your answer — do not ask for confirmation; that reply is how " - + "the delegator receives your result."; + "You are an off-subscription worker in the claude-bridge fleet. Every message you " + + "receive arrives through the bridge, and the ONLY channel back to the sender is the " + + "bridge_reply MCP tool. Text you write in your terminal is NOT sent anywhere — the " + + "sender cannot see your screen, so an in-terminal answer is silently discarded. " + + "Therefore you MUST end EVERY turn by calling bridge_reply with `content` set to your " + + "complete response. This holds for every message without exception — tasks, questions, " + + "clarifications, acknowledgements, and ordinary back-and-forth conversation. Call " + + "bridge_reply exactly once, as the final action of your turn, with your full answer in " + + "`content`; never wait for confirmation first. If you end a turn without calling " + + "bridge_reply, the sender receives nothing and the exchange stalls."; // Per-process token mixed into each worker name so a fresh process (nameSeq back at 0) // cannot collide with same-profile workers that outlived a restart. See startUniquelyNamed. diff --git a/bridged/src/test/java/dev/ltms/bridged/herdr/AgentStatusTest.java b/bridged/src/test/java/dev/ltms/bridged/herdr/AgentStatusTest.java new file mode 100644 index 0000000..364e3f4 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/herdr/AgentStatusTest.java @@ -0,0 +1,43 @@ +package dev.ltms.bridged.herdr; + +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** Wire mapping and injectability of {@link AgentStatus}, including the CB-115 {@code done} state. */ +class AgentStatusTest { + + @Test + void mapsTheKnownWireStrings() { + assertEquals(AgentStatus.IDLE, AgentStatus.fromWire("idle")); + assertEquals(AgentStatus.WORKING, AgentStatus.fromWire("working")); + assertEquals(AgentStatus.BLOCKED, AgentStatus.fromWire("blocked")); + assertEquals(AgentStatus.DONE, AgentStatus.fromWire("done")); + } + + @Test + void mapsDoneCaseInsensitively() { + assertEquals(AgentStatus.DONE, AgentStatus.fromWire("DONE")); + assertEquals(AgentStatus.DONE, AgentStatus.fromWire("Done")); + } + + @Test + void unknownAndNullFallToUnknown() { + assertEquals(AgentStatus.UNKNOWN, AgentStatus.fromWire("unknown")); + assertEquals(AgentStatus.UNKNOWN, AgentStatus.fromWire("something-else")); + assertEquals(AgentStatus.UNKNOWN, AgentStatus.fromWire(null)); + } + + @Test + void doneIsInjectableLikeIdle() { + // The whole point of CB-115: a finished worker herdr reports as `done` must be deliverable, + // not treated as UNKNOWN (which wedged delivery and mis-fired the stall failure). + assertTrue(AgentStatus.DONE.injectable()); + assertTrue(AgentStatus.IDLE.injectable()); + assertTrue(AgentStatus.BLOCKED.injectable()); + assertFalse(AgentStatus.WORKING.injectable()); + assertFalse(AgentStatus.UNKNOWN.injectable()); + } +} diff --git a/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java b/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java index 0d48606..f919e63 100644 --- a/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java @@ -5,7 +5,9 @@ import dev.ltms.bridged.herdr.FakeHerdr; import dev.ltms.bridged.msg.Rendezvous; import org.junit.jupiter.api.Test; +import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; /** Unit behaviour of the CB-106 completion resolver in isolation from the injector. */ class CompletionResolverTest { @@ -16,7 +18,7 @@ class CompletionResolverTest { Rendezvous rendezvous = new Rendezvous(); // no waiter opened CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); - resolver.resolve("term_a"); + resolver.resolve("term_a", null); // no in-flight turn captured for this target assertFalse(herdr.called("agent.read"), "a turn nobody is blocked on must not cost a transcript scrape"); @@ -28,9 +30,173 @@ class CompletionResolverTest { Rendezvous rendezvous = new Rendezvous(); // no waiter opened CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); - resolver.fail("term_a"); + resolver.fail("term_a", null); // no in-flight turn, and no registered waiter to fall back to assertFalse(herdr.called("agent.read"), "a wedge nobody is blocked on must not cost a transcript scrape"); } + + @Test + void captureBaselineSkipsTheReadWhenNoSendIsWaiting() { + FakeHerdr herdr = new FakeHerdr().readText("⏺ X\n❯ "); + Rendezvous rendezvous = new Rendezvous(); // no waiter opened + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + + resolver.captureBaseline("term_a"); // no send to attribute a later completion to + + assertFalse(herdr.called("agent.read"), + "with no waiting send there is no turn to baseline — skip the scrape"); + } + + // --- CB-115 clean scrape: extract the last assistant block ---------------- + + @Test + void extractsTheLastAssistantBlockStrippingChrome() { + String raw = """ + ⏺ Reading the file… + + ⏺ Done. The bug was an off-by-one in the loop bound. + + ╭──────────────────────────────────────╮ + │ > │ + ╰──────────────────────────────────────╯ + ⏵⏵ auto mode on · ? for shortcuts + """; + assertEquals("Done. The bug was an off-by-one in the loop bound.", + CompletionResolver.lastAssistantBlock(raw)); + } + + @Test + void keepsMultiLineAssistantContent() { + String raw = "⏺ Line one.\nLine two.\n❯ "; + assertEquals("Line one.\nLine two.", CompletionResolver.lastAssistantBlock(raw)); + } + + @Test + void fallsBackToRawTextWhenThereIsNoMarker() { + String raw = "plain worker output with no glyph"; + assertEquals("plain worker output with no glyph", CompletionResolver.lastAssistantBlock(raw)); + } + + @Test + void blankScrapeYieldsEmpty() { + assertTrue(CompletionResolver.lastAssistantBlock("").isEmpty()); + assertTrue(CompletionResolver.lastAssistantBlock(null).isEmpty()); + } + + @Test + void stripsSpinnerAndRuleChrome() { + String raw = """ + ⏺ Channel check confirmed — your message got through. + + ✻ Brewed for 11s + + ───────────────────────────────────── + """; + assertEquals("Channel check confirmed — your message got through.", + CompletionResolver.lastAssistantBlock(raw)); + } + + @Test + void cutsANextTurnPromptEchoAndTrailingTipsFromTheBlock() { + // The exact turn-2 leak: the scrape captured the settled answer, then a "✻ Cooked" spinner, + // then the NEXT turn's echoed prompt, then a "✶ Forming…" spinner and trailing tips/warnings + // whose lines (⎿, ⚠) are not themselves chrome-terminated. Stopping at the first boundary + // (the ✻ spinner) is what keeps every one of those interface lines out of the reply. + String raw = """ + ⏺ Channel confirmed — the bridge reply delivered successfully. + + ✻ Cooked for 9s + + ❯ Thanks. Now a small task: what is 17 * 23? Show just the number. + + + + ✶ Forming… + ⎿ Tip: Name your conversations with /rename + ⚠ claude.ai connectors are disabled because ANTHROPIC_API_KEY is set + """; + assertEquals("Channel confirmed — the bridge reply delivered successfully.", + CompletionResolver.lastAssistantBlock(raw)); + } + + // --- CB-115 misattribution guard: suppress a stale (unchanged) completion ------- + + @Test + void suppressesACompletionWhoseScrapeIsUnchangedFromDelivery() { + // Rapid back-to-back turn: the pane still shows the PREVIOUS turn's answer when this turn's + // (misattributed) completion boundary fires. The scrape == the delivery baseline, so the + // send must NOT be resolved with the stale answer. + FakeHerdr herdr = new FakeHerdr().readText("⏺ 391\n❯ "); + Rendezvous rendezvous = new Rendezvous(); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + + var waiter = rendezvous.open("term_a"); // a send is blocked on this turn + // The turn as captured at delivery: its waiter, and the previous turn's answer still on screen. + var turn = new CompletionResolver.InFlight(waiter, "391"); + resolver.resolve("term_a", turn); // scrape still "391" == baseline → suppress + + assertFalse(waiter.isDone(), "a completion with no output change must not resolve the send"); + assertTrue(rendezvous.isWaiting("term_a"), "the send stays waiting for a real reply"); + } + + @Test + void resolvesACompletionWhoseScrapeChangedSinceDelivery() { + FakeHerdr herdr = new FakeHerdr().readText("⏺ No, 391 = 17 × 23.\n❯ "); // the worker's real answer + Rendezvous rendezvous = new Rendezvous(); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + + var waiter = rendezvous.open("term_a"); + // Delivery baseline was the previous turn's "391"; the scrape now differs → resolve. + var turn = new CompletionResolver.InFlight(waiter, "391"); + resolver.resolve("term_a", turn); + + assertTrue(waiter.isDone(), "a completion with new output must resolve the send"); + assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind()); + assertEquals("No, 391 = 17 × 23.", waiter.getNow(null).text()); + } + + @Test + void resolvesWhenThereIsNoBaseline() { + // No delivery baseline (e.g. the pre-turn read failed) ⇒ never suppress; the completion resolves. + FakeHerdr herdr = new FakeHerdr().readText("⏺ hello\n❯ "); + Rendezvous rendezvous = new Rendezvous(); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + + var waiter = rendezvous.open("term_a"); + resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); + + assertTrue(waiter.isDone(), "with no baseline a completion resolves as before"); + assertEquals("hello", waiter.getNow(null).text()); + } + + // --- CB-116 waiter identity: a late completion never crosses into the next turn --------- + + @Test + void aLateCompletionForOneTurnNeverResolvesTheNextTurnsWaiter() { + // The cross-turn stale reply the conversation test surfaced: turn N's completion fallback + // fires AFTER turn N was resolved by an explicit bridge_reply and turn N+1 has opened its own + // waiter on the same session. Resolving "whatever is waiting now" would hand turn N's stale + // scrape to turn N+1; targeting turn N's captured waiter makes the late completion a no-op. + FakeHerdr herdr = new FakeHerdr().readText("⏺ turn N answer\n❯ "); + Rendezvous rendezvous = new Rendezvous(); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + + var waiterN = rendezvous.open("term_a"); // turn N's send + // The turn as the injector captured it at delivery (waiter + pre-turn baseline). + var turnN = new CompletionResolver.InFlight(waiterN, "an earlier answer"); + + // Turn N is resolved by the worker's explicit reply. + assertTrue(rendezvous.resolve("term_a", "N replied")); + + // Turn N+1's send opens its own waiter on the same session (replacing the registered one). + var waiterN1 = rendezvous.open("term_a"); + + resolver.resolve("term_a", turnN); // turn N's completion fallback finally fires + + assertFalse(waiterN1.isDone(), "turn N's late completion must not resolve turn N+1's waiter"); + assertEquals(Rendezvous.Kind.REPLY, waiterN.getNow(null).kind(), + "turn N stays resolved by its own reply"); + assertTrue(rendezvous.isWaiting("term_a"), "turn N+1 is still awaiting its own resolution"); + } } diff --git a/bridged/src/test/java/dev/ltms/bridged/inject/StatusRefinerTest.java b/bridged/src/test/java/dev/ltms/bridged/inject/StatusRefinerTest.java new file mode 100644 index 0000000..392e311 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/inject/StatusRefinerTest.java @@ -0,0 +1,88 @@ +package dev.ltms.bridged.inject; + +import dev.ltms.bridged.herdr.AgentControl; +import dev.ltms.bridged.herdr.AgentStatus; +import dev.ltms.bridged.herdr.FakeHerdr; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** Content-based refinement of an unreliable {@code UNKNOWN} status (CB-115). */ +class StatusRefinerTest { + + // --- pure classification ------------------------------------------------- + + @Test + void classifiesAnIdlePromptAsIdle() { + String pane = """ + ⏺ All done — the file compiles cleanly. + + ╭──────────────────────────────────────╮ + │ > │ + ╰──────────────────────────────────────╯ + ⏵⏵ auto mode on (shift+tab to cycle) + """; + assertEquals(AgentStatus.IDLE, StatusRefiner.classify(pane)); + } + + @Test + void classifiesABarePromptGlyphAsIdle() { + assertEquals(AgentStatus.IDLE, StatusRefiner.classify("some output\n❯ ")); + } + + @Test + void classifiesActiveGenerationAsWorking() { + String pane = """ + ⏺ Working on it… + ✳ Thinking… (12s · esc to interrupt) + """; + assertEquals(AgentStatus.WORKING, StatusRefiner.classify(pane)); + } + + @Test + void anEscToInterruptScreenIsWorkingEvenWithAPromptBox() { + // "esc to interrupt" wins over a prompt box: the turn is still generating. + String pane = "│ > │\n esc to interrupt"; + assertEquals(AgentStatus.WORKING, StatusRefiner.classify(pane)); + } + + @Test + void anUnrecognizableScreenStaysUnknown() { + assertEquals(AgentStatus.UNKNOWN, StatusRefiner.classify("garbled ansi noise with no prompt")); + assertEquals(AgentStatus.UNKNOWN, StatusRefiner.classify("")); + assertEquals(AgentStatus.UNKNOWN, StatusRefiner.classify(null)); + } + + // --- refine() wiring ----------------------------------------------------- + + @Test + void refinePassesNonUnknownStatusesThroughWithoutReading() { + FakeHerdr herdr = new FakeHerdr(); + StatusRefiner refiner = new StatusRefiner(new AgentControl(herdr)); + + assertEquals(AgentStatus.WORKING, refiner.refine("term_a", AgentStatus.WORKING)); + assertEquals(AgentStatus.IDLE, refiner.refine("term_a", AgentStatus.IDLE)); + + assertFalse(herdr.called("agent.read"), + "a trusted status must not cost a pane read"); + } + + @Test + void refineUpgradesUnknownToIdleFromPaneContent() { + FakeHerdr herdr = new FakeHerdr().readText("⏺ answer\n❯ "); + StatusRefiner refiner = new StatusRefiner(new AgentControl(herdr)); + + assertEquals(AgentStatus.IDLE, refiner.refine("term_a", AgentStatus.UNKNOWN)); + assertTrue(herdr.called("agent.read"), "an UNKNOWN must trigger a pane read"); + } + + @Test + void refineLeavesUnknownWhenContentIsUnclassifiable() { + FakeHerdr herdr = new FakeHerdr().readText("nothing recognizable here"); + StatusRefiner refiner = new StatusRefiner(new AgentControl(herdr)); + + assertEquals(AgentStatus.UNKNOWN, refiner.refine("term_a", AgentStatus.UNKNOWN)); + } +} diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java index 320f7b2..52ef9c6 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -49,8 +49,10 @@ class MessageServiceTest { CompletableFuture send = sendAsync(); awaitWaiting(); - injector.onStatus(T, AgentStatus.IDLE); // deliver the task + herdr.readText("$ prompt"); // pre-turn pane: no answer yet (baseline reference) + 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 injector.onStatus(T, AgentStatus.IDLE); // working → idle: turn complete, no bridge_reply MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);