078bde2c02
Injector.drop, CompletionResolver.fail, SessionManager.onFailed/reapIdle/
acquire spawn failures, and MessageService.abandon used to fail a member
or a caller's request with no log, a bare DEBUG, or a log that named only
the symptom ("session marked failed", "failed send via turn-stall
fallback"). Each now logs at WARN and names the real cause and the
numbers involved. Observability only — no behaviour changed.
266 lines
14 KiB
Java
266 lines
14 KiB
Java
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.
|
||
*
|
||
* <p>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.
|
||
*
|
||
* <p>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.
|
||
*
|
||
* <p>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.
|
||
*
|
||
* <p><strong>Waiter-specific resolution (CB-116).</strong> On delivery we also capture the exact
|
||
* {@link Rendezvous} waiter this turn belongs to, and the completion/failure fallbacks resolve
|
||
* <em>that</em> 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
|
||
* <em>next</em> 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.
|
||
*
|
||
* <p>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 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.
|
||
*
|
||
* <p>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<Rendezvous.Resolution> waiter, String baseline) {
|
||
}
|
||
|
||
private final ConcurrentHashMap<String, InFlight> 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<Rendezvous.Resolution> 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<Rendezvous.Resolution> 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(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;
|
||
}
|
||
// 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, 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<Rendezvous.Resolution> 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.warn("failing send to {} via turn-stall fallback: {}", target, reason);
|
||
}
|
||
}
|
||
|
||
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 <em>first</em>
|
||
* hard interface boundary below it — the spinner/status line, input box, {@code ❯} prompt (which
|
||
* may echo the <em>next</em> 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.
|
||
*
|
||
* <p>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 <em>not</em>
|
||
* 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");
|
||
}
|
||
}
|