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());
+ }
+}