From 5b26caca0c705252a4e7228cc46035c2e39f743a Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Wed, 15 Jul 2026 14:14:04 +0200 Subject: [PATCH] =?UTF-8?q?CB-106:=20completion=20fallback=20=E2=80=94=20r?= =?UTF-8?q?esolve=20a=20send=20when=20the=20worker's=20turn=20ends=20witho?= =?UTF-8?q?ut=20bridge=5Freply?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The blocking send previously resolved only on an explicit bridge_reply; a real delegated task (edit files, run a build) finishes and goes idle without ever calling it, so the send always timed out. The injector now reports a confirmed working -> idle turn boundary via a TurnListener; CompletionResolver scrapes the worker's transcript tail and resolves the awaiting send (Rendezvous.resolveCompletion, Kind.COMPLETION -> Outcome.COMPLETED_UNREPLIED), surfaced as replySource=transcript at the REST/MCP edges. Completion is synthesized only from a confirmed turn (a sampled 'working'), never from the pickup-grace path, so it can't race the explicit reply or fire on a turn that never ran. --- .../main/java/dev/ltms/bridged/Bridged.java | 7 +- .../bridged/inject/CompletionResolver.java | 78 ++++++++++++++++ .../dev/ltms/bridged/inject/Injector.java | 91 ++++++++++++++----- .../dev/ltms/bridged/inject/TurnListener.java | 19 ++++ .../java/dev/ltms/bridged/mcp/BridgeMcp.java | 8 +- .../dev/ltms/bridged/msg/MessageService.java | 26 ++++-- .../java/dev/ltms/bridged/msg/Rendezvous.java | 57 +++++++++--- .../dev/ltms/bridged/rest/BridgedApp.java | 8 +- .../dev/ltms/bridged/herdr/FakeHerdr.java | 9 ++ .../inject/CompletionResolverTest.java | 24 +++++ .../dev/ltms/bridged/inject/InjectorTest.java | 45 +++++++-- .../ltms/bridged/msg/MessageServiceTest.java | 74 +++++++++++++++ 12 files changed, 395 insertions(+), 51 deletions(-) create mode 100644 bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java create mode 100644 bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java create mode 100644 bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java create mode 100644 bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 498d09f..665b9cf 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -6,6 +6,7 @@ import dev.ltms.bridged.herdr.AgentControl; import dev.ltms.bridged.herdr.PaneLocator; import dev.ltms.bridged.herdr.UnixSocketHerdrClient; import dev.ltms.bridged.herdr.WorkspaceControl; +import dev.ltms.bridged.inject.CompletionResolver; import dev.ltms.bridged.inject.Injector; import dev.ltms.bridged.inject.StatusPoller; import dev.ltms.bridged.mcp.BridgeMcp; @@ -54,12 +55,14 @@ public final class Bridged { // Status-gated injector (CB-103): the single writer into workers, fed by a poller. // The blocking message endpoint (CB-104) is the producer; the poller is inert until then. - Injector injector = new Injector(agents); + // CB-106: a confirmed turn completion resolves a blocked send whose worker never replied. + Rendezvous rendezvous = new Rendezvous(); + CompletionResolver completion = new CompletionResolver(agents, rendezvous); + Injector injector = new Injector(agents, completion); StatusPoller poller = new StatusPoller(agents, injector, INJECT_POLL_MILLIS); poller.start(); Runtime.getRuntime().addShutdownHook(new Thread(poller::stop)); - Rendezvous rendezvous = new Rendezvous(); MessageService messages = new MessageService(agents, injector, rendezvous); // MCP server face (CB-105): bridge_send/bridge_reply/bridge_status, mounted at /mcp. diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java new file mode 100644 index 0000000..2ae1d50 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java @@ -0,0 +1,78 @@ +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; + +/** + * 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. + * + *

Wired as the {@link Injector}'s {@link TurnListener}; {@link #onTurnComplete} hands off to a + * virtual thread so the scrape's herdr round-trip never stalls the status poller. + */ +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; + + public CompletionResolver(AgentControl agents, Rendezvous rendezvous) { + this.agents = agents; + this.rendezvous = rendezvous; + } + + @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)); + } + + /** 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 + } + String tail; + try { + tail = clip(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 = ""; + } + if (rendezvous.resolveCompletion(target, tail)) { + log.debug("resolved send to {} via turn-completion fallback ({} chars scraped)", + target, tail.length()); + } + } + + 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); + } +} 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 5cdb47d..d0af9bc 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java @@ -32,6 +32,14 @@ import java.util.stream.Collectors; * a transient {@code unknown} — counts as a real pickup, so a detection glitch can't prematurely * release the latch. Perfectly reliable turn boundaries require a herdr {@code events.subscribe} * stream; that is the intended upgrade and would replace only the sampling, not this queue. + * + *

Turn completion (CB-106). Beyond delivery, the injector reports when a + * delegated turn finishes: after a delivery is picked up (a real {@code working} sample), + * the next injectable sample is a confirmed {@code working → idle} boundary and fires + * {@link TurnListener#onTurnComplete}. Completion is only ever synthesized from a confirmed + * turn — the pickup-grace path (a turn too fast to sample) unwedges the queue but does not fire + * completion, since without a sampled {@code working} there is no trustworthy "the worker just + * finished the task" signal to act on. */ public final class Injector { @@ -46,10 +54,18 @@ public final class Injector { private static final int PICKUP_GRACE_POLLS = 8; private final AgentControl agents; + private final TurnListener turnListener; private final ConcurrentHashMap targets = new ConcurrentHashMap<>(); + /** Delivery only; completion signalling is a no-op. */ public Injector(AgentControl agents) { + this(agents, TurnListener.NOOP); + } + + /** Delivery plus turn-completion signalling to {@code turnListener} (CB-106). */ + public Injector(AgentControl agents, TurnListener turnListener) { this.agents = agents; + this.turnListener = turnListener; } /** A pending message and the future that completes when it has been delivered. */ @@ -59,8 +75,10 @@ public final class Injector { /** Per-worker delivery state, guarded by its own monitor (single writer per worker). */ private static final class Target { final Deque queue = new ArrayDeque<>(); - boolean awaitingPickup; // sent a message, waiting for the worker to pick it up + boolean awaitingPickup; // sent a message, waiting for the worker to pick it up int injectableSincePickup; // consecutive injectable samples while awaitingPickup + boolean awaitingCompletion; // a delivered message's turn is not yet known-complete + boolean turnObserved; // saw a real `working` sample since that delivery (turn ran) synchronized void add(Pending p) { queue.add(p); @@ -99,33 +117,53 @@ public final class Injector { Pending sent = null; RuntimeException sendError = null; + boolean turnCompleted = false; synchronized (t) { if (status == AgentStatus.WORKING) { - // Definitive pickup: the worker is busy on our last message. + // Definitive pickup: the worker is busy on our last message, and (if a delivery is + // outstanding) a real turn is now confirmed to be running. t.awaitingPickup = false; t.injectableSincePickup = 0; + if (t.awaitingCompletion) t.turnObserved = true; } else if (status.injectable()) { // IDLE or BLOCKED if (t.awaitingPickup && ++t.injectableSincePickup >= PICKUP_GRACE_POLLS) { // Pickup edge was never sampled (turn faster than the poll, or status lag). - // The worker has plainly moved on — release the latch rather than wedge. + // Release the latch rather than wedge — and give up on synthesizing a completion + // for this message, since without a confirmed `working` we cannot trust that a + // task-processing turn actually ran. t.awaitingPickup = false; t.injectableSincePickup = 0; + t.awaitingCompletion = false; + t.turnObserved = false; } if (!t.awaitingPickup) { - Pending p = t.queue.peek(); - if (p != null) { - try { - agents.send(target, p.text()); - t.queue.poll(); - t.awaitingPickup = true; - t.injectableSincePickup = 0; - sent = p; - } catch (RuntimeException e) { - // Delivery failed at herdr; drop the poisoned message and surface it - // rather than blocking the queue behind it. - t.queue.poll(); - sent = p; - sendError = e; + // A confirmed turn (a `working` sample was seen) that has now returned to idle is + // a trustworthy `working → idle` completion boundary. + if (t.awaitingCompletion && t.turnObserved) { + t.awaitingCompletion = false; + t.turnObserved = false; + turnCompleted = true; + } + // Deliver the next queued message only once the prior turn is fully settled, so a + // completion is never confused with the pickup of the following message. + if (!t.awaitingCompletion) { + Pending p = t.queue.peek(); + if (p != null) { + try { + agents.send(target, p.text()); + t.queue.poll(); + t.awaitingPickup = true; + t.awaitingCompletion = true; + t.turnObserved = false; + t.injectableSincePickup = 0; + sent = p; + } catch (RuntimeException e) { + // Delivery failed at herdr; drop the poisoned message and surface it + // rather than blocking the queue behind it. + t.queue.poll(); + sent = p; + sendError = e; + } } } } @@ -133,13 +171,18 @@ public final class Injector { // UNKNOWN (and any other non-injectable, non-working): do nothing — neither a safe // window nor a reliable pickup signal, so we must not deliver or release the latch. - // Reclaim the entry once the worker is fully quiescent, so the map cannot grow without - // bound across many short-lived workers. - if (t.queue.isEmpty() && !t.awaitingPickup) { + // Reclaim the entry once the worker is fully quiescent (nothing queued, no pickup or + // completion awaited), so the map cannot grow without bound across short-lived workers. + if (t.queue.isEmpty() && !t.awaitingPickup && !t.awaitingCompletion) { targets.remove(target, t); } } + // Fire listeners after releasing the monitor so a continuation (which may call herdr) never + // runs on the poller thread while it holds the target lock. + if (turnCompleted) { + turnListener.onTurnComplete(target); + } if (sent != null) { if (sendError != null) { log.warn("inject to {} failed, dropped message: {}", target, sendError.getMessage()); @@ -150,12 +193,16 @@ public final class Injector { } } - /** Targets the poller must keep sampling: those with a queued message or an awaited pickup. */ + /** + * Targets the poller must keep sampling: those with a queued message, an awaited pickup, or an + * awaited turn completion (so the {@code working → idle} boundary is observed). + */ public Set activeTargets() { return targets.entrySet().stream() .filter(e -> { synchronized (e.getValue()) { - return !e.getValue().queue.isEmpty() || e.getValue().awaitingPickup; + Target t = e.getValue(); + return !t.queue.isEmpty() || t.awaitingPickup || t.awaitingCompletion; } }) .map(java.util.Map.Entry::getKey) diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java b/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java new file mode 100644 index 0000000..d3a0f04 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java @@ -0,0 +1,19 @@ +package dev.ltms.bridged.inject; + +/** + * Notified when a worker's delegated turn is observed to complete — a confirmed + * {@code WORKING → IDLE} transition after a delivery. This is the CB-106 completion signal the + * {@code CompletionResolver} uses to resolve a blocked send whose worker never called + * {@code bridge_reply}. Kept as a seam so the {@link Injector} needs no dependency on the message + * layer and stays unit-testable with a capturing fake. + */ +@FunctionalInterface +public interface TurnListener { + + /** A worker's delegated turn finished (worker returned to idle after visibly working). */ + void onTurnComplete(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/mcp/BridgeMcp.java b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java index 17398b7..51b5406 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java @@ -93,9 +93,15 @@ public final class BridgeMcp { long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs); try { MessageService.Reply r = messages.send(sessionId, content, timeout); - if (r.completed()) { + if (r.outcome() == MessageService.Outcome.REPLIED) { return text(r.text()); } + if (r.outcome() == MessageService.Outcome.COMPLETED_UNREPLIED) { + // The worker's turn finished but it never called bridge_reply — hand back the + // scraped transcript tail, flagged so the primary knows it isn't a structured reply. + return text("[worker finished without a structured bridge_reply — transcript tail follows]\n" + + r.text()); + } return text("[no reply within " + timeout + "ms — worker " + r.outcome().name().toLowerCase().replace("timed_out_", "") + "; retry or poll status]"); } catch (HerdrException e) { diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java index a45dc7e..fe1fb0e 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -34,8 +34,13 @@ public final class MessageService { /** Outcome of a blocking send. */ public enum Outcome { - /** The worker replied; {@code text} holds the structured answer. */ + /** The worker called {@code bridge_reply}; {@code text} holds the structured answer. */ REPLIED, + /** + * The worker's delegated turn finished without a {@code bridge_reply} (CB-106 fallback); + * {@code text} is the scraped transcript tail rather than a structured answer. + */ + COMPLETED_UNREPLIED, /** Timed out after the message was delivered — the worker is still working. */ TIMED_OUT_WORKING, /** Timed out before delivery — the message is still queued for the worker. */ @@ -44,10 +49,16 @@ public final class MessageService { BUSY } - /** @param text the worker's reply when {@link #outcome} is {@link Outcome#REPLIED}, else {@code null} */ + /** + * @param outcome how the send ended + * @param text the worker's answer when {@link #completed()} (a structured {@code bridge_reply} + * for {@link Outcome#REPLIED}, a scraped transcript tail for + * {@link Outcome#COMPLETED_UNREPLIED}), else {@code null} + */ public record Reply(Outcome outcome, String text) { + /** Whether the worker's turn actually finished with an answer (replied or scraped). */ public boolean completed() { - return outcome == Outcome.REPLIED; + return outcome == Outcome.REPLIED || outcome == Outcome.COMPLETED_UNREPLIED; } } @@ -80,10 +91,13 @@ public final class MessageService { } try { CompletableFuture delivered = injector.enqueue(target, content); - CompletableFuture reply = rendezvous.open(target); + CompletableFuture reply = rendezvous.open(target); try { - String text = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS); - return new Reply(Outcome.REPLIED, text); + Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS); + Outcome outcome = r.kind() == Rendezvous.Kind.REPLY + ? Outcome.REPLIED + : Outcome.COMPLETED_UNREPLIED; + return new Reply(outcome, r.text()); } catch (TimeoutException e) { boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally(); log.debug("send to {} timed out (delivered={})", target, wasDelivered); 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 b42a84c..5346ca8 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java @@ -4,38 +4,71 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; /** - * The reply rendezvous: where a blocking {@code bridge_send} awaits the worker's structured - * {@code bridge_reply} for the same session. The sending (primary) request thread {@link #open}s - * a waiter; the worker's reply (arriving on a different thread via {@code POST - * /sessions/{id}/reply}) {@link #resolve}s it. + * The reply rendezvous: where a blocking {@code bridge_send} awaits how the worker's delegated turn + * ends. The sending (primary) request thread {@link #open}s a waiter; it is resolved either by the + * worker's explicit {@code bridge_reply} ({@link #resolve}, arriving on a different thread via + * {@code POST /sessions/{id}/reply}) or — the CB-106 fallback — by the injector observing the + * worker's delegated turn return to idle without a reply ({@link #resolveCompletion}). * *

At most one waiter per session — {@link MessageService} serializes sends per session, so a - * reply maps unambiguously to the one outstanding send and cannot be captured by another. + * resolution maps unambiguously to the one outstanding send and cannot be captured by another. */ public final class Rendezvous { - private final ConcurrentHashMap> waiters = new ConcurrentHashMap<>(); + /** How a delegated turn ended. */ + public enum Kind { + /** The worker called {@code bridge_reply} with a structured answer. */ + REPLY, + /** The worker's turn finished without a {@code bridge_reply}; {@code text} is a scrape. */ + COMPLETION + } + + /** The resolved outcome of a send: its {@link Kind} and the associated text. */ + public record Resolution(Kind kind, String text) { + } + + 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) { - CompletableFuture waiter = new CompletableFuture<>(); + CompletableFuture open(String session) { + CompletableFuture waiter = new CompletableFuture<>(); waiters.put(session, waiter); return waiter; } /** Remove {@code waiter} for {@code session} (only if it is still the registered one). */ - void close(String session, CompletableFuture waiter) { + void close(String session, CompletableFuture waiter) { waiters.remove(session, waiter); } + /** Whether a send is currently awaiting a resolution for {@code session}. */ + public boolean isWaiting(String session) { + return waiters.containsKey(session); + } + /** - * Resolve the send awaiting on {@code session} with the worker's reply {@code content}. + * Resolve the send awaiting on {@code session} with the worker's explicit reply {@code content}. * * @return {@code true} if a waiter was resolved; {@code false} if none was waiting (a late or * spurious reply — e.g. the send already timed out) */ public boolean resolve(String session, String content) { - CompletableFuture waiter = waiters.get(session); - return waiter != null && waiter.complete(content); + return complete(session, new Resolution(Kind.REPLY, content)); + } + + /** + * 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. + * + * @return {@code true} if a waiter was resolved, {@code false} if none was waiting + */ + public boolean resolveCompletion(String session, String text) { + return complete(session, new Resolution(Kind.COMPLETION, text)); + } + + private boolean complete(String session, Resolution resolution) { + CompletableFuture waiter = waiters.get(session); + return waiter != null && waiter.complete(resolution); } } diff --git a/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java b/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java index 5360beb..c81c339 100644 --- a/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java +++ b/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java @@ -151,7 +151,11 @@ public final class BridgedApp { try { MessageService.Reply reply = messages.send(id, content, timeout); if (reply.completed()) { - ctx.status(200).json(Map.of("sessionId", id, "reply", reply.text())); + // replySource distinguishes a structured bridge_reply from the CB-106 completion + // fallback (a scrape of the worker's transcript when it finished without replying). + String source = reply.outcome() == MessageService.Outcome.REPLIED ? "reply" : "transcript"; + ctx.status(200).json(Map.of( + "sessionId", id, "reply", reply.text(), "replySource", source)); } else { ctx.status(202).json(Map.of( "sessionId", id, @@ -159,7 +163,7 @@ public final class BridgedApp { case TIMED_OUT_WORKING -> "working"; case TIMED_OUT_QUEUED -> "queued"; case BUSY -> "busy"; - case REPLIED -> "done"; // unreachable + case REPLIED, COMPLETED_UNREPLIED -> "done"; // unreachable (completed() branch) }, "detail", "no reply within " + timeout + "ms; poll status or retry")); } diff --git a/bridged/src/test/java/dev/ltms/bridged/herdr/FakeHerdr.java b/bridged/src/test/java/dev/ltms/bridged/herdr/FakeHerdr.java index 103468c..f32bc39 100644 --- a/bridged/src/test/java/dev/ltms/bridged/herdr/FakeHerdr.java +++ b/bridged/src/test/java/dev/ltms/bridged/herdr/FakeHerdr.java @@ -28,6 +28,7 @@ public final class FakeHerdr implements HerdrClient { private String paneCloseErrorCode = null; private String agentSendErrorCode = null; private volatile String agentStatus = "idle"; // steady-state agent.get status + private volatile String readText = "worker transcript tail"; // canned agent.read output public FakeHerdr healthy(boolean h) { this.healthy = h; @@ -58,6 +59,12 @@ public final class FakeHerdr implements HerdrClient { return this; } + /** The text {@code agent.read} returns (the CB-106 completion scrape). */ + public FakeHerdr readText(String text) { + this.readText = text; + return this; + } + /** Make {@code agent.send} fail with this herdr error code. */ public FakeHerdr agentSendFailsWith(String code) { this.agentSendErrorCode = code; @@ -111,6 +118,8 @@ public final class FakeHerdr implements HerdrClient { {"type":"agent_info","agent":{"terminal_id":"term_a","agent":"claude", "agent_status":"%s","workspace_id":"w2","tab_id":"w2:t7","pane_id":"w2:p7"}}""") .formatted(agentStatus)); + case "agent.read" -> mapper.readTree(mapper.writeValueAsString( + java.util.Map.of("type", "agent_read", "read", java.util.Map.of("text", readText)))); case "agent.start" -> { long starts = calls.stream().filter(c -> c.method().equals("agent.start")).count(); if (starts <= agentNameTakenFor) { diff --git a/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java b/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java new file mode 100644 index 0000000..4fec392 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java @@ -0,0 +1,24 @@ +package dev.ltms.bridged.inject; + +import dev.ltms.bridged.herdr.AgentControl; +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.assertFalse; + +/** Unit behaviour of the CB-106 completion resolver in isolation from the injector. */ +class CompletionResolverTest { + + @Test + void skipsTheScrapeWhenNoSendIsWaiting() { + FakeHerdr herdr = new FakeHerdr(); + Rendezvous rendezvous = new Rendezvous(); // no waiter opened + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + + resolver.resolve("term_a"); + + assertFalse(herdr.called("agent.read"), + "a turn nobody is blocked on must not cost a transcript scrape"); + } +} diff --git a/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java b/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java index 6653cf8..5f46da6 100644 --- a/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java @@ -6,8 +6,10 @@ import dev.ltms.bridged.herdr.FakeHerdr; import dev.ltms.bridged.herdr.HerdrException; import org.junit.jupiter.api.Test; +import java.util.ArrayList; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; @@ -126,17 +128,48 @@ class InjectorTest { } @Test - void activeWhileQueuedOrInFlightThenQuietAfterPickup() { + void activeWhileQueuedOrInFlightThenQuietAfterTurnCompletes() { assertTrue(injector.activeTargets().isEmpty()); injector.enqueue(T, "x"); - assertEquals(java.util.Set.of(T), injector.activeTargets(), "active while a message is queued"); + assertEquals(Set.of(T), injector.activeTargets(), "active while a message is queued"); - injector.onStatus(T, AgentStatus.IDLE); // delivers; still in-flight (awaiting pickup) - assertEquals(java.util.Set.of(T), injector.activeTargets(), + injector.onStatus(T, AgentStatus.IDLE); // delivers; awaiting pickup + assertEquals(Set.of(T), injector.activeTargets(), "stays active so the poller can observe the worker pick the message up"); - injector.onStatus(T, AgentStatus.WORKING); // pickup observed → in-flight cleared - assertTrue(injector.activeTargets().isEmpty(), "quiet once queue is empty and pickup is seen"); + injector.onStatus(T, AgentStatus.WORKING); // pickup observed; now awaiting turn completion + assertEquals(Set.of(T), injector.activeTargets(), + "stays active after pickup so the working→idle completion boundary is observed"); + + injector.onStatus(T, AgentStatus.IDLE); // working → idle: turn complete + assertTrue(injector.activeTargets().isEmpty(), "quiet once the delegated turn has completed"); + } + + @Test + void firesTurnCompleteOnAConfirmedWorkingThenIdle() { + List completed = new ArrayList<>(); + Injector inj = new Injector(new AgentControl(herdr), completed::add); + inj.enqueue(T, "task"); + + inj.onStatus(T, AgentStatus.IDLE); // deliver + inj.onStatus(T, AgentStatus.WORKING); // pickup + turn running + assertEquals(List.of(), completed, "no completion until the turn returns to idle"); + + inj.onStatus(T, AgentStatus.IDLE); // working → idle: turn complete + assertEquals(List.of(T), completed, "a confirmed working→idle fires exactly one completion"); + } + + @Test + void doesNotSynthesizeCompletionFromAnUnconfirmedTurn() { + List completed = new ArrayList<>(); + Injector inj = new Injector(new AgentControl(herdr), completed::add); + inj.enqueue(T, "task"); + + // Deliver, then only ever idle — a `working` sample is never seen. The pickup grace unwedges + // the queue but must NOT invent a completion: without a sampled turn there is no trustworthy + // "the worker finished the task" signal, so the send should fall through to its timeout. + for (int i = 0; i < 15; i++) inj.onStatus(T, AgentStatus.IDLE); + assertEquals(List.of(), completed, "no completion is synthesized from an unconfirmed turn"); } @Test diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java new file mode 100644 index 0000000..948acdc --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -0,0 +1,74 @@ +package dev.ltms.bridged.msg; + +import dev.ltms.bridged.herdr.AgentControl; +import dev.ltms.bridged.herdr.AgentStatus; +import dev.ltms.bridged.herdr.FakeHerdr; +import dev.ltms.bridged.inject.CompletionResolver; +import dev.ltms.bridged.inject.Injector; +import org.junit.jupiter.api.Test; + +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * The message layer's resolution paths (CB-104 reply + CB-106 completion fallback). The turn is + * driven deterministically by feeding {@code onStatus} rather than running a real poller. + */ +class MessageServiceTest { + + private static final String T = "term_a"; + + private final FakeHerdr herdr = new FakeHerdr().readText("BUILD GREEN: 391 files"); + private final AgentControl agents = new AgentControl(herdr); + private final Rendezvous rendezvous = new Rendezvous(); + private final CompletionResolver completion = new CompletionResolver(agents, rendezvous); + private final Injector injector = new Injector(agents, completion); + private final MessageService messages = new MessageService(agents, injector, rendezvous); + + /** Run {@code send} on a background thread; the current thread drives the worker's turn. */ + private CompletableFuture sendAsync() { + return CompletableFuture.supplyAsync(() -> messages.send(T, "do the task", 5000)); + } + + private void awaitWaiting() throws InterruptedException { + long deadline = System.currentTimeMillis() + 2000; + while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) { + //noinspection BusyWait + Thread.sleep(5); + } + assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter"); + } + + @Test + void completionFallbackResolvesATurnThatNeverCalledBridgeReply() throws Exception { + CompletableFuture send = sendAsync(); + awaitWaiting(); + + injector.onStatus(T, AgentStatus.IDLE); // deliver the task + injector.onStatus(T, AgentStatus.WORKING); // worker picks it up and works + injector.onStatus(T, AgentStatus.IDLE); // working → idle: turn complete, no bridge_reply + + MessageService.Reply reply = send.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.COMPLETED_UNREPLIED, reply.outcome(), + "an unreplied but finished turn resolves via the completion fallback"); + assertEquals("BUILD GREEN: 391 files", reply.text(), "the scraped transcript tail is returned"); + assertTrue(reply.completed(), "a scraped completion still counts as completed"); + } + + @Test + void explicitBridgeReplyResolvesAsReplied() throws Exception { + CompletableFuture send = sendAsync(); + awaitWaiting(); + + injector.onStatus(T, AgentStatus.IDLE); // deliver + injector.onStatus(T, AgentStatus.WORKING); // worker working + assertTrue(rendezvous.resolve(T, "LGTM ship it"), "an explicit reply resolves the send"); + + MessageService.Reply reply = send.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.REPLIED, reply.outcome()); + assertEquals("LGTM ship it", reply.text()); + } +}