CB-106: completion fallback — resolve a send when the worker's turn ends without bridge_reply

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.
This commit is contained in:
Dai Ha
2026-07-15 14:14:04 +02:00
parent 6988bfe88f
commit 5b26caca0c
12 changed files with 395 additions and 51 deletions
@@ -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.
@@ -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.
*
* <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>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);
}
}
@@ -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.
*
* <p><strong>Turn completion (CB-106).</strong> Beyond delivery, the injector reports when a
* delegated turn <em>finishes</em>: 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 <em>confirmed</em>
* 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<String, Target> 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<Pending> 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<String> 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)
@@ -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 = _ -> {
};
}
@@ -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) {
@@ -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<Void> delivered = injector.enqueue(target, content);
CompletableFuture<String> reply = rendezvous.open(target);
CompletableFuture<Rendezvous.Resolution> 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);
@@ -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}).
*
* <p>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<String, CompletableFuture<String>> 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<String, CompletableFuture<Resolution>> waiters = new ConcurrentHashMap<>();
/** Register a waiter for {@code session}. The caller must hold that session's send lock. */
CompletableFuture<String> open(String session) {
CompletableFuture<String> waiter = new CompletableFuture<>();
CompletableFuture<Resolution> open(String session) {
CompletableFuture<Resolution> 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<String> waiter) {
void close(String session, CompletableFuture<Resolution> 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<String> 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<Resolution> waiter = waiters.get(session);
return waiter != null && waiter.complete(resolution);
}
}
@@ -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"));
}
@@ -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) {
@@ -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");
}
}
@@ -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<String> 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<String> 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
@@ -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<MessageService.Reply> 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<MessageService.Reply> 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<MessageService.Reply> 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());
}
}