package dev.ltms.bridged.inject; import dev.ltms.bridged.herdr.AgentControl; 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 * task without ever calling {@code bridge_reply} — the common case for a real delegated coding task. * *

On a confirmed {@code working → idle} boundary it scrapes the worker's recent transcript and * resolves the awaiting send with that tail (a {@link Rendezvous.Kind#COMPLETION} resolution, so the * caller can tell a scrape from a structured reply). It scrapes only when a send is actually waiting * — a fleet worker's own turns, or a send that already timed out, cost no herdr traffic. An explicit * {@code bridge_reply} that raced in first wins; {@link Rendezvous#resolveCompletion} is then a no-op. * *

It also handles the CB-109 stall signal ({@link #onTurnFailed}): a worker that ran a turn then * 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. * *

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 { private static final Logger log = LoggerFactory.getLogger(CompletionResolver.class); /** * herdr {@code agent.read} source for the completion scrape. {@code recent} returns the tail of * the transcript (the worker's last output), which is what a delegator wants when the worker * didn't structure a reply. */ static final String SCRAPE_SOURCE = "recent"; /** Cap the scraped tail so a long transcript can't return an unbounded blob. */ static final int MAX_SCRAPE_CHARS = 4000; private static final String CLIPPED_PANE_TAIL_MARKER = "[Pane tail clipped: member did not call bridge_reply.]"; 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 { // Clip to the same cap resolve() applies to the tail (line ~134): the CB-115 misattribution // guard compares baseline.equals(tail), so both sides must be the same capped representation. // An unclipped baseline vs a clipped tail would never match for a >MAX_SCRAPE_CHARS block, // defeating the guard and letting a stale completion resolve the send. baseline = clip(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)); } /** The turn currently baselined for {@code target}, or {@code null} — a test hook for the captureBaseline path. */ InFlight inFlight(String target) { return inFlight.get(target); } @Override public void onTurnComplete(String 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)); } /** * Resolve the completed turn before adapter housekeeping can erase its rendered output. This is * intentionally synchronous and used only when a post-turn context reset is enabled; the normal * path remains off-loaded so polling is not blocked by a scrape. */ public void resolveBeforePostAction(String target) { resolve(target, inFlight.get(target)); } @Override public void onTurnFailed(String 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, 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; int originalLength = 0; boolean clipped = false; boolean scrapeFailed = false; try { String assistantBlock = lastAssistantBlock(agents.read(target, SCRAPE_SOURCE)); 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; } // 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 } String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail; if (rendezvous.resolveCompletion(waiter, completion)) { inFlight.remove(target, turn); if (clipped) { log.warn("completion scrape for {} clipped from {} chars to the {} char cap; " + "member did not call bridge_reply, so the pane tail is partial", target, originalLength, MAX_SCRAPE_CHARS); } 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, 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 { reason = clip(agents.read(target, SCRAPE_SOURCE)); } catch (RuntimeException e) { reason = ""; } if (reason.isBlank()) { // No screen to scrape — either the worker is stuck (CB-109) or gone (CB-110). reason = "worker did not reply; its turn ended in an unrecoverable state " + "(worker unreachable or stuck)"; } if (rendezvous.resolveFailure(waiter, reason)) { inFlight.remove(target, turn); log.debug("failed send to {} via turn-stall fallback", target); } } private static String clip(String s) { if (s == null) return ""; String trimmed = s.strip(); return trimmed.length() <= MAX_SCRAPE_CHARS ? 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"); } }