CB-205/CB-201: bridge_ask reverse rendezvous + lightweight question/turn kind

A worker can now pause its delegated turn to ask the primary a question and
resume the same turn with the answer — the reverse of bridge_send.

- Rendezvous: QUESTION kind carrying a turnId, plus a reverse-ask registry
  (openAsk/resolveQuestion/askSession/answerAsk/closeAsk).
- MessageService.ask(): surface a worker's question to the primary's open send,
  block for the answer. answer(): resolve the worker's ask by turnId, then block
  for its eventual bridge_reply (session derived from turnId, not an argument).
- BridgeMcp: bridge_ask tool (worker-only, identity from the connection);
  bridge_send routes a turnId to the answer path. Shared formatReply().
- BridgedApp REST parity: POST /sessions/{id}/ask, turnId on /message.
- Lightweight CB-201: the kind vocabulary is the QUESTION outcome + turn_id
  correlation, not a rigid from/to/corr envelope (the connection-identity
  mechanism already covers addressing more robustly).

Tests: MessageServiceTest ask/answer round-trip + NO_WAITER/timeout/stale;
new RendezvousTest for the reverse registry; BridgeMcp ask/answer parity.
mvn: 147 green.
This commit is contained in:
Dai Ha
2026-07-16 15:32:56 +02:00
parent 55ebd5b949
commit 358c6970b5
7 changed files with 578 additions and 55 deletions
@@ -40,6 +40,10 @@ public final class BridgeMcp {
private static final long DEFAULT_TIMEOUT_MS = 25_000;
private static final long MAX_TIMEOUT_MS = 120_000;
// bridge_ask blocks the WORKER's own MCP call, which its client caps near 60s — default under
// that so the bridge returns a clean timeout before the client severs the call (CB-205).
private static final long ASK_DEFAULT_TIMEOUT_MS = 55_000;
private static final long ASK_MAX_TIMEOUT_MS = 115_000;
private static final ObjectMapper MAPPER = new ObjectMapper(); // worker-view JSON projections
/** Transport-context key under which the extractor stashes the resolved caller identity. */
@@ -73,6 +77,12 @@ public final class BridgeMcp {
.capabilities(McpSchema.ServerCapabilities.builder().tools(true).build())
.toolCall(sendTool(), (_, req) -> {
Map<String, Object> a = req.arguments();
String turnId = str(a, "turnId");
if (turnId != null && !turnId.isBlank()) {
// Answering a worker's bridge_ask (CB-205): resolve its blocked question and
// block for the worker's reply as it resumes the same turn.
return answer(messages, turnId, str(a, "content"), timeoutMs(a));
}
// wait defaults to true (block for the reply); wait:false is fire-and-poll.
return Boolean.FALSE.equals(a.get("wait"))
? sendAsync(messages, str(a, "sessionId"), str(a, "content"))
@@ -81,6 +91,9 @@ public final class BridgeMcp {
// bridge_reply's identity is the CONNECTION, never an argument.
.toolCall(replyTool(), (exchange, req) ->
reply(rendezvous, callerTerminal(exchange), str(req.arguments(), "content")))
// bridge_ask (CB-205): a worker's mid-turn question — identity from the CONNECTION.
.toolCall(askTool(), (exchange, req) ->
ask(messages, callerTerminal(exchange), str(req.arguments(), "question"), timeoutMs(req.arguments())))
.toolCall(statusTool(), (_, req) ->
status(messages, str(req.arguments(), "sessionId")))
.toolCall(pollTool(), (_, req) ->
@@ -138,27 +151,70 @@ public final class BridgeMcp {
}
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
try {
MessageService.Reply r = messages.send(sessionId, content, timeout);
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());
}
if (r.outcome() == MessageService.Outcome.WORKER_FAILED) {
// The worker ran the turn then wedged (CB-109) — surface the error context.
return text("[worker failed — turn ended in an unrecoverable state]\n" + r.text());
}
return text("[no reply within " + timeout + "ms — worker "
+ r.outcome().name().toLowerCase().replace("timed_out_", "") + "; retry or poll status]");
return formatReply(messages.send(sessionId, content, timeout), timeout);
} catch (HerdrException e) {
return error("herdr error contacting session " + sessionId + ": " + e.getMessage());
}
}
/**
* {@code bridge_send} carrying a {@code turnId}: the primary's answer to a worker's
* {@code bridge_ask} (CB-205). Resolves the worker's blocked question and blocks for its reply as
* it resumes the same turn — surfaced to the primary identically to a normal send.
*/
static McpSchema.CallToolResult answer(MessageService messages, String turnId, String content, Long timeoutMs) {
if (isBlank(turnId) || isBlank(content)) {
return error("turnId and content are required to answer a worker's question");
}
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
return formatReply(messages.answer(turnId, content, timeout), timeout);
}
/**
* {@code bridge_ask} (CB-205): a worker pauses its delegated turn to ask the primary, blocking
* until the primary answers. The worker is identified by its connection ({@code callerTerminal}),
* never an argument — a {@code null} means the caller is not a known worker.
*/
static McpSchema.CallToolResult ask(MessageService messages, String callerTerminal, String question, Long timeoutMs) {
if (callerTerminal == null) {
return error("bridge_ask is for workers only — could not identify the calling worker "
+ "from the connection");
}
if (isBlank(question)) {
return error("question is required");
}
long timeout = Math.clamp(timeoutMs == null ? ASK_DEFAULT_TIMEOUT_MS : timeoutMs, 1, ASK_MAX_TIMEOUT_MS);
MessageService.AskResult r = messages.ask(callerTerminal, question, timeout);
return switch (r.outcome()) {
case ANSWERED -> text(r.answer());
case NO_WAITER -> error("no primary is awaiting this turn — bridge_ask only works while a "
+ "bridge_send delegation is open to answer it");
case TIMED_OUT -> text("[no answer within " + timeout + "ms — the primary did not respond; "
+ "proceed on your best judgement, then call bridge_reply to end the turn]");
};
}
/** Render a {@link MessageService.Reply} as a tool result — shared by {@link #send} and {@link #answer}. */
private static McpSchema.CallToolResult formatReply(MessageService.Reply r, long timeout) {
return switch (r.outcome()) {
case REPLIED -> text(r.text());
// 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.
case COMPLETED_UNREPLIED -> text(
"[worker finished without a structured bridge_reply — transcript tail follows]\n" + r.text());
// The worker ran the turn then wedged (CB-109) — surface the error context.
case WORKER_FAILED -> text("[worker failed — turn ended in an unrecoverable state]\n" + r.text());
// The worker paused mid-turn to ask (CB-205) — tell the primary how to answer in-turn.
case QUESTION -> text("[question] the worker paused to ask before it can finish:\n" + r.text()
+ "\n\nAnswer it by calling bridge_send again with turnId=\"" + r.turnId()
+ "\" and content set to your answer; the worker resumes the same turn.");
case STALE_TURN -> error("that question is no longer open — it timed out or was already "
+ "answered (turnId stale)");
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> text("[no reply within " + timeout + "ms — worker "
+ r.outcome().name().toLowerCase().replace("timed_out_", "") + "; retry or poll status]");
};
}
/**
* {@code bridge_send} with {@code wait:false}: delegate {@code content} and return a ticket
* immediately (fire-and-poll), so a long task isn't cut off by the caller's MCP call timeout.
@@ -298,14 +354,31 @@ public final class BridgeMcp {
return tool("bridge_send",
"Delegate a task to a worker session. By default blocks until the worker replies and "
+ "returns its reply (or a 'still working / queued' note on timeout). Pass wait:false "
+ "for a long task to return a ticket immediately, then poll it with bridge_poll.",
+ "for a long task to return a ticket immediately, then poll it with bridge_poll. To "
+ "answer a worker's bridge_ask, pass its turnId (with content) instead of sessionId.",
objectSchema(Map.of(
"sessionId", stringProp("The worker session id (herdr terminal_id) to delegate to"),
"content", stringProp("The task/message to send to the worker"),
"content", stringProp("The task/message to send to the worker (or your answer, with turnId)"),
"timeoutMs", Map.of("type", "integer", "description", "Max ms to wait for a reply (blocking mode)"),
"wait", Map.of("type", "boolean",
"description", "Block for the reply (default true); false returns a ticket to poll")),
List.of("sessionId", "content")));
"description", "Block for the reply (default true); false returns a ticket to poll"),
"turnId", stringProp("When answering a worker's bridge_ask, its question turnId — "
+ "routes your answer back into the same turn (omit for a normal delegation)")),
List.of("content")));
}
private static McpSchema.Tool askTool() {
// No target/session arg — the worker's identity is resolved from the connection.
return tool("bridge_ask",
"Pause your current delegated turn to ask the primary a question, blocking until it "
+ "answers — then resume the same turn with the answer. Use this when only the "
+ "primary has a decision or detail you need to continue. You do not address the "
+ "primary; identity is resolved from your connection.",
objectSchema(Map.of(
"question", stringProp("The question to put to the primary"),
"timeoutMs", Map.of("type", "integer",
"description", "Max ms to wait for the primary's answer")),
List.of("question")));
}
private static McpSchema.Tool pollTool() {
@@ -65,27 +65,60 @@ public final class MessageService {
* failure context (e.g. the error screen). Terminal, but not a successful completion.
*/
WORKER_FAILED,
/**
* The worker paused mid-turn to ask the primary a question (CB-205); {@code text} is the
* question and {@code turnId} correlates the answer. Not terminal — the primary answers with
* {@link #answer(String, String, long)} and the turn resumes.
*/
QUESTION,
/** 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. */
TIMED_OUT_QUEUED,
/** Another send to this session was in flight for the whole window. */
BUSY
BUSY,
/**
* An answer ({@link #answer(String, String, long)}) referenced a {@code turnId} that is no
* longer open — the worker's {@code bridge_ask} already timed out or was answered.
*/
STALE_TURN
}
/**
* @param outcome how the send ended
* @param outcome how the send ended (or paused)
* @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}
* {@link Outcome#COMPLETED_UNREPLIED}), or the question for {@link Outcome#QUESTION},
* else {@code null}
* @param turnId correlation id for a {@link Outcome#QUESTION} (answered via
* {@link #answer(String, String, long)}), else {@code null}
*/
public record Reply(Outcome outcome, String text) {
public record Reply(Outcome outcome, String text, String turnId) {
/** A reply with no correlation id (the common terminal outcomes). */
public Reply(Outcome outcome, String text) {
this(outcome, text, null);
}
/** Whether the worker's turn actually finished with an answer (replied or scraped). */
public boolean completed() {
return outcome == Outcome.REPLIED || outcome == Outcome.COMPLETED_UNREPLIED;
}
}
/** How a worker's {@code bridge_ask} (CB-205) resolved. */
public enum AskOutcome {
/** The primary answered; {@link AskResult#answer} carries it. */
ANSWERED,
/** No delegation was open to surface the question to — the worker has no one to ask. */
NO_WAITER,
/** The primary did not answer within the window. */
TIMED_OUT
}
/** The outcome of a worker's {@code bridge_ask}: how it resolved and (if answered) the answer. */
public record AskResult(AskOutcome outcome, String answer) {
}
/** Lifecycle phase of an async delegation ticket. */
public enum Phase {
/** Delegated and in flight — queued for the worker or being worked. */
@@ -147,12 +180,7 @@ public final class MessageService {
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target);
try {
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
Outcome outcome = switch (r.kind()) {
case REPLY -> Outcome.REPLIED;
case COMPLETION -> Outcome.COMPLETED_UNREPLIED;
case FAILED -> Outcome.WORKER_FAILED;
};
return new Reply(outcome, r.text());
return new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
} catch (TimeoutException e) {
boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally();
log.debug("send to {} timed out (delivered={})", target, wasDelivered);
@@ -171,6 +199,90 @@ public final class MessageService {
}
}
/**
* A worker's mid-turn question (CB-205 reverse rendezvous): surface {@code question} to the
* primary by resolving its open blocking {@code bridge_send}, then block this (worker) call until
* the primary answers via {@link #answer} or {@code timeoutMillis} elapses. Identity is the
* worker's own session — it does not address the primary.
*
* <p>Returns {@link AskOutcome#NO_WAITER} when no delegation is open to surface the question to
* (nothing to answer it), {@link AskOutcome#ANSWERED} with the primary's answer, or
* {@link AskOutcome#TIMED_OUT} if the primary stayed silent. The worker resumes its turn either
* way — an answered ask hands back the answer; an unanswered one leaves it to proceed alone.
*/
public AskResult ask(String workerSession, String question, long timeoutMillis) {
Rendezvous.AskTicket ticket = rendezvous.openAsk(workerSession);
// Register the reverse waiter first (above), then surface the question — so the answer, which
// can arrive the instant the primary reacts, always finds an open waiter to resolve.
if (!rendezvous.resolveQuestion(workerSession, question, ticket.turnId())) {
rendezvous.closeAsk(ticket.turnId());
return new AskResult(AskOutcome.NO_WAITER, null); // no primary is blocked on this worker
}
try {
String answer = ticket.answer().get(timeoutMillis, TimeUnit.MILLISECONDS);
return new AskResult(AskOutcome.ANSWERED, answer);
} catch (TimeoutException e) {
log.debug("bridge_ask from {} went unanswered in {}ms", workerSession, timeoutMillis);
return new AskResult(AskOutcome.TIMED_OUT, null);
} catch (ExecutionException e) {
Throwable cause = e.getCause();
throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new IllegalStateException("interrupted awaiting the primary's answer for " + workerSession, e);
} finally {
rendezvous.closeAsk(ticket.turnId());
}
}
/**
* The primary's answer to a worker's {@code bridge_ask} (CB-205): resolve the worker's blocked
* question identified by {@code turnId}, then — like a fresh {@link #send} — block for the worker's
* eventual {@code bridge_reply} as it finishes the resumed turn. The worker session is derived from
* {@code turnId}, never a caller argument.
*
* <p>Unlike {@link #send} this does not re-inject through the {@link Injector}: the worker is
* mid-turn (already picked up), so the answer flows back through its own open {@code bridge_ask}
* call, not a new status-gated delivery. The forward waiter is opened <em>before</em> the worker
* is unblocked so a reply that lands the instant it resumes is not lost.
*/
public Reply answer(String turnId, String content, long timeoutMillis) {
String workerSession = rendezvous.askSession(turnId);
if (workerSession == null) {
return new Reply(Outcome.STALE_TURN, null); // the ask lapsed (timed out or already answered)
}
long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L;
ReentrantLock lock = sessionLocks.computeIfAbsent(workerSession, _ -> new ReentrantLock());
if (!tryLock(lock, remainingMillis(deadlineNanos))) {
return new Reply(Outcome.BUSY, null);
}
try {
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(workerSession);
if (!rendezvous.answerAsk(turnId, content)) {
rendezvous.close(workerSession, reply);
return new Reply(Outcome.STALE_TURN, null); // lapsed between the lookup and the unblock
}
try {
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
return new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
} catch (TimeoutException e) {
// The worker resumed but hasn't replied yet — no completion fallback arms an answered
// turn (it never re-entered the injector), so a silent worker rides out the window.
return new Reply(Outcome.TIMED_OUT_WORKING, null);
} catch (ExecutionException e) {
Throwable cause = e.getCause();
throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new IllegalStateException("interrupted awaiting reply from " + workerSession, e);
} finally {
rendezvous.close(workerSession, reply);
}
} finally {
lock.unlock();
}
}
/**
* Fire-and-poll variant of {@link #send}: deliver {@code content} to {@code target} on a
* background virtual thread and return immediately with a ticket to {@link #poll}. This is how a
@@ -241,6 +353,16 @@ public final class MessageService {
asyncExecutor.shutdown();
}
/** Map a rendezvous {@link Rendezvous.Kind} onto its send {@link Outcome} (shared by send/answer). */
private static Outcome outcomeOf(Rendezvous.Kind kind) {
return switch (kind) {
case REPLY -> Outcome.REPLIED;
case COMPLETION -> Outcome.COMPLETED_UNREPLIED;
case FAILED -> Outcome.WORKER_FAILED;
case QUESTION -> Outcome.QUESTION;
};
}
private static boolean tryLock(ReentrantLock lock, long millis) {
try {
return lock.tryLock(Math.max(0, millis), TimeUnit.MILLISECONDS);
@@ -2,6 +2,7 @@ package dev.ltms.bridged.msg;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong;
/**
* The reply rendezvous: where a blocking {@code bridge_send} awaits how the worker's delegated turn
@@ -24,22 +25,47 @@ import java.util.concurrent.ConcurrentHashMap;
*/
public final class Rendezvous {
/** How a delegated turn ended. */
/** How a delegated turn ended (or paused). */
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 worker ran the turn then wedged (CB-109); {@code text} is the failure context. */
FAILED
FAILED,
/**
* The worker paused mid-turn to ask the primary a question (CB-205 reverse rendezvous);
* {@code text} is the question and {@code turnId} correlates the primary's answer back to
* the worker's blocked {@code bridge_ask}. Not terminal — the turn resumes after the answer.
*/
QUESTION
}
/** The resolved outcome of a send: its {@link Kind} and the associated text. */
public record Resolution(Kind kind, String text) {
/**
* The resolved outcome of a send: its {@link Kind}, the associated text, and — only for
* {@link Kind#QUESTION} — the {@code turnId} the primary answers with (else {@code null}).
*/
public record Resolution(Kind kind, String text, String turnId) {
/** A terminal resolution (reply / completion / failure) with no correlation id. */
public Resolution(Kind kind, String text) {
this(kind, text, null);
}
}
/** A worker's open mid-turn question: the worker session it belongs to and the answer future. */
private record AskWaiter(String session, CompletableFuture<String> answer) {
}
/** Handle to a just-opened reverse-rendezvous turn: the minted {@code turnId} and its answer future. */
public record AskTicket(String turnId, CompletableFuture<String> answer) {
}
private final ConcurrentHashMap<String, CompletableFuture<Resolution>> waiters = new ConcurrentHashMap<>();
/** Reverse rendezvous (CB-205): worker questions awaiting the primary's answer, keyed by {@code turnId}. */
private final ConcurrentHashMap<String, AskWaiter> asks = new ConcurrentHashMap<>();
private final AtomicLong askSeq = new AtomicLong();
/**
* Register a waiter for {@code session} — the await side of the public {@code resolve*} methods.
* The caller must hold that session's send lock.
@@ -79,6 +105,56 @@ public final class Rendezvous {
return complete(session, new Resolution(Kind.REPLY, content));
}
// --- reverse rendezvous (CB-205 bridge_ask) ------------------------------------------------
/**
* Open a reverse-rendezvous waiter for a worker's mid-turn question: mint a fresh {@code turnId}
* (scoped to the worker session), register the answer future under it, and hand both back. The
* caller then {@link #resolveQuestion surfaces the question} to the primary and blocks on the
* returned future until the primary {@link #answerAsk answers} — a no-op if it never does.
*/
public AskTicket openAsk(String session) {
CompletableFuture<String> answer = new CompletableFuture<>();
String turnId = session + "#" + askSeq.incrementAndGet();
asks.put(turnId, new AskWaiter(session, answer));
return new AskTicket(turnId, answer);
}
/**
* Surface a worker's mid-turn {@code question} by resolving the primary's open {@code bridge_send}
* with a {@link Kind#QUESTION} carrying {@code turnId}. Same session-keyed semantics as
* {@link #resolve}: the one outstanding send for {@code session} unblocks with the question.
*
* @return {@code true} if a send was awaiting (the question reached the primary); {@code false}
* if none was (no delegation is open to answer it)
*/
public boolean resolveQuestion(String session, String question, String turnId) {
return complete(session, new Resolution(Kind.QUESTION, question, turnId));
}
/** The worker session an outstanding ask {@code turnId} belongs to, or {@code null} if unknown/lapsed. */
public String askSession(String turnId) {
AskWaiter w = asks.get(turnId);
return w == null ? null : w.session();
}
/**
* Resolve a worker's blocked {@code bridge_ask} with the primary's {@code answer}, unblocking it
* to resume its turn.
*
* @return {@code true} if the ask was still open and got the answer; {@code false} if the
* {@code turnId} is unknown or the ask already lapsed (timed out / was answered)
*/
public boolean answerAsk(String turnId, String answer) {
AskWaiter w = asks.get(turnId);
return w != null && w.answer().complete(answer);
}
/** Drop a reverse-rendezvous turn once its {@code bridge_ask} has resolved (answered or lapsed). */
public void closeAsk(String turnId) {
asks.remove(turnId);
}
/**
* Resolve a specific captured {@code waiter} as a completion (the delegated turn finished with no
* {@code bridge_reply}); {@code text} is the scraped transcript tail. The waiter is the one
@@ -34,6 +34,9 @@ public final class BridgedApp {
/** Default blocking window for a message; kept under typical HTTP idle timeouts. */
private static final long DEFAULT_MESSAGE_TIMEOUT_MS = 25_000;
private static final long MAX_MESSAGE_TIMEOUT_MS = 120_000;
/** Blocking window for a worker's bridge_ask (CB-205); the worker's MCP client caps its own call. */
private static final long DEFAULT_ASK_TIMEOUT_MS = 55_000;
private static final long MAX_ASK_TIMEOUT_MS = 115_000;
private final HerdrClient herdr;
private final WorkerService workers;
@@ -69,8 +72,9 @@ public final class BridgedApp {
app.get("/profiles", this::profiles); // configured worker profiles
app.post("/workers", this::spawnWorker); // optional ?profile= or {"profile":…}
app.delete("/workers/{paneId}", this::stopWorker);
app.post("/sessions/{id}/message", this::sendMessage); // bridge_send (primary; blocking or wait:false)
app.post("/sessions/{id}/message", this::sendMessage); // bridge_send (primary; blocking, wait:false, or answer via turnId)
app.post("/sessions/{id}/reply", this::replyMessage); // bridge_reply (worker)
app.post("/sessions/{id}/ask", this::askMessage); // bridge_ask (worker → primary, CB-205)
app.get("/sessions/{id}/status", this::sessionStatus); // bridge_status
app.get("/tasks/{ticket}", this::taskStatus); // poll an async (wait:false) send
return app;
@@ -170,11 +174,13 @@ public final class BridgedApp {
private void sendMessage(Context ctx) {
String id = ctx.pathParam("id");
String content;
String turnId;
long timeout;
boolean wait;
try {
JsonNode body = mapper.readTree(ctx.body());
content = body.path("content").asText("");
turnId = body.path("turnId").asText(null);
timeout = body.path("timeoutMs").asLong(DEFAULT_MESSAGE_TIMEOUT_MS);
wait = body.path("wait").asBoolean(true); // default: block for the reply (CB-104)
} catch (Exception e) {
@@ -185,6 +191,13 @@ public final class BridgedApp {
ctx.status(400).json(Map.of("error", "bad_request", "detail", "content is required"));
return;
}
timeout = Math.clamp(timeout, 1, MAX_MESSAGE_TIMEOUT_MS);
// Answering a worker's bridge_ask (CB-205): always blocks, and derives the worker from turnId.
if (turnId != null && !turnId.isBlank()) {
writeReply(ctx, id, messages.answer(turnId, content, timeout), timeout);
return;
}
if (!wait) {
// Fire-and-poll (CB-107): return a ticket immediately; the caller polls GET /tasks/{ticket}.
@@ -192,32 +205,79 @@ public final class BridgedApp {
ctx.status(202).json(Map.of("sessionId", id, "ticket", ticket, "status", "accepted"));
return;
}
timeout = Math.clamp(timeout, 1, MAX_MESSAGE_TIMEOUT_MS);
try {
MessageService.Reply reply = messages.send(id, content, timeout);
if (reply.completed()) {
writeReply(ctx, id, messages.send(id, content, timeout), timeout);
} catch (HerdrException e) {
herdrError(ctx, e);
}
}
/**
* Render a {@link MessageService.Reply} onto the response — shared by a normal send and a
* bridge_ask answer. A structured/scraped completion is 200; a worker's mid-turn question a 202
* (with its {@code turnId}); a stale answer a 409; every other non-terminal outcome a typed 202.
*/
private void writeReply(Context ctx, String id, MessageService.Reply reply, long timeout) {
switch (reply.outcome()) {
case QUESTION -> ctx.status(202).json(Map.of(
"sessionId", id, "status", "question",
"question", reply.text(), "turnId", reply.turnId()));
case STALE_TURN -> ctx.status(409).json(Map.of(
"sessionId", id, "error", "stale_turn",
"detail", "that question is no longer open (timed out or already answered)"));
case REPLIED, COMPLETED_UNREPLIED -> {
// 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,
"status", switch (reply.outcome()) {
case TIMED_OUT_WORKING -> "working";
case TIMED_OUT_QUEUED -> "queued";
case BUSY -> "busy";
case WORKER_FAILED -> "failed";
case REPLIED, COMPLETED_UNREPLIED -> "done"; // unreachable (completed() branch)
},
"detail", reply.outcome() == MessageService.Outcome.WORKER_FAILED && reply.text() != null
? reply.text()
: "no reply within " + timeout + "ms; poll status or retry"));
ctx.status(200).json(Map.of("sessionId", id, "reply", reply.text(), "replySource", source));
}
} catch (HerdrException e) {
herdrError(ctx, e);
default -> ctx.status(202).json(Map.of(
"sessionId", id,
"status", switch (reply.outcome()) {
case TIMED_OUT_WORKING -> "working";
case TIMED_OUT_QUEUED -> "queued";
case BUSY -> "busy";
case WORKER_FAILED -> "failed";
default -> "done"; // unreachable (terminal outcomes handled above)
},
"detail", reply.outcome() == MessageService.Outcome.WORKER_FAILED && reply.text() != null
? reply.text()
: "no reply within " + timeout + "ms; poll status or retry"));
}
}
/**
* A worker's mid-turn question ({@code bridge_ask}, CB-205) — surfaces to the primary's open
* blocking send and blocks until it answers. 200 with the answer, 409 if no delegation is open,
* 202 if the primary stayed silent.
*/
private void askMessage(Context ctx) {
String id = ctx.pathParam("id");
String question;
long timeout;
try {
JsonNode body = mapper.readTree(ctx.body());
question = body.path("question").asText("");
timeout = body.path("timeoutMs").asLong(DEFAULT_ASK_TIMEOUT_MS);
} catch (Exception e) {
ctx.status(400).json(Map.of("error", "bad_request", "detail", "body must be JSON"));
return;
}
if (question.isBlank()) {
ctx.status(400).json(Map.of("error", "bad_request", "detail", "question is required"));
return;
}
timeout = Math.clamp(timeout, 1, MAX_ASK_TIMEOUT_MS);
MessageService.AskResult r = messages.ask(id, question, timeout);
switch (r.outcome()) {
case ANSWERED -> ctx.status(200).json(Map.of("sessionId", id, "answered", true, "answer", r.answer()));
case NO_WAITER -> ctx.status(409).json(Map.of(
"sessionId", id, "error", "no_pending_send",
"detail", "no primary is awaiting this turn to answer a question"));
case TIMED_OUT -> ctx.status(202).json(Map.of(
"sessionId", id, "status", "no_answer",
"detail", "the primary did not answer within " + timeout + "ms"));
}
}
@@ -120,6 +120,63 @@ class BridgeMcpTest {
assertTrue(textOf(res).contains("no send is awaiting"));
}
@Test
void askThenAnswerRoundTrips() throws Exception {
// The primary delegates and blocks; wait until its waiter is open before the worker asks.
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> BridgeMcp.send(messages, "term_a", "do X", 5000L));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting("term_a"), "the send must be waiting for the ask to surface to");
// The worker asks mid-turn; the call blocks for the primary's answer.
CompletableFuture<McpSchema.CallToolResult> ask = CompletableFuture.supplyAsync(
() -> BridgeMcp.ask(messages, "term_a", "which config?", 5000L));
// The primary's send unblocks with the question and a turnId to answer on.
McpSchema.CallToolResult q = send.get(6, TimeUnit.SECONDS);
assertNotEquals(Boolean.TRUE, q.isError());
String qt = textOf(q);
assertTrue(qt.contains("[question]"), qt);
String afterMarker = qt.substring(qt.indexOf("turnId=\"") + "turnId=\"".length());
String turnId = afterMarker.substring(0, afterMarker.indexOf('"'));
// The primary answers via bridge_send(turnId); this blocks again for the worker's reply.
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> BridgeMcp.answer(messages, turnId, "config.yaml", 5000L));
// The worker's ask returns the answer — it resumes the same turn.
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
// The resumed worker replies, resolving the answering send (retry past the reopen race).
McpSchema.CallToolResult reply = BridgeMcp.reply(rendezvous, "term_a", "done");
deadline = System.currentTimeMillis() + 3000;
while (Boolean.TRUE.equals(reply.isError()) && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
Thread.sleep(10);
reply = BridgeMcp.reply(rendezvous, "term_a", "done");
}
assertEquals("delivered", textOf(reply));
assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS)));
}
@Test
void askFromANonWorkerConnectionIsAnError() {
McpSchema.CallToolResult res = BridgeMcp.ask(messages, null, "which config?", 500L);
assertTrue(res.isError());
assertTrue(textOf(res).contains("workers only"), textOf(res));
}
@Test
void answerToAStaleTurnIsAnError() {
McpSchema.CallToolResult res = BridgeMcp.answer(messages, "term_a#999", "too late", 500L);
assertTrue(res.isError());
assertTrue(textOf(res).contains("no longer open"), textOf(res));
}
@Test
void spawnReturnsTheNewWorkersSessionAndPane() {
FakeHerdr h = new FakeHerdr();
@@ -13,6 +13,7 @@ import java.util.concurrent.TimeUnit;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
@@ -107,4 +108,69 @@ class MessageServiceTest {
"a delivered send whose worker vanishes fails instead of hanging to the timeout");
assertFalse(reply.completed());
}
// --- bridge_ask reverse rendezvous (CB-205) ------------------------------------------------
@Test
void askSurfacesAsAQuestionAndTheAnswerResumesTheSameTurn() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up, then pauses to ask
// The worker asks mid-turn on its own thread; the call blocks for the primary's answer.
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
// The primary's blocking send unblocks with the question and a turnId to answer on.
MessageService.Reply q = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
assertEquals("which config file?", q.text());
assertNotNull(q.turnId(), "a question carries a turnId to answer on");
// The primary answers via bridge_send(turnId); this blocks again for the worker's reply.
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
// The worker's ask returns the answer — it resumes the same turn.
MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.AskOutcome.ANSWERED, a.outcome());
assertEquals("config.yaml", a.answer());
// The resumed worker finishes with a structured reply, resolving the answering send.
awaitWaiting(); // the answering send has (re)opened its forward waiter
assertTrue(rendezvous.resolve(T, "done"), "the worker's final reply resolves the answering send");
MessageService.Reply done = answer.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.REPLIED, done.outcome());
assertEquals("done", done.text());
}
@Test
void askWithNoOpenDelegationReturnsNoWaiter() {
MessageService.AskResult r = messages.ask(T, "anyone listening?", 500);
assertEquals(MessageService.AskOutcome.NO_WAITER, r.outcome(),
"a question with no blocked send has no primary to answer it");
}
@Test
void askTimesOutWhenThePrimaryNeverAnswers() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
MessageService.AskResult r = messages.ask(T, "still there?", 200); // primary never answers
assertEquals(MessageService.AskOutcome.TIMED_OUT, r.outcome());
// The send itself already unblocked with the question the instant the ask surfaced.
MessageService.Reply q = send.get(2, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
}
@Test
void answeringAnUnknownTurnIsStale() {
MessageService.Reply r = messages.answer(T + "#999", "too late", 500);
assertEquals(MessageService.Outcome.STALE_TURN, r.outcome(),
"an answer to a turn that never existed (or already lapsed) is stale, not a hang");
}
}
@@ -0,0 +1,69 @@
package dev.ltms.bridged.msg;
import org.junit.jupiter.api.Test;
import java.util.concurrent.CompletableFuture;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* The reverse rendezvous (CB-205): the {@code bridge_ask} registry that lets a worker pause mid-turn
* to ask the primary. Unit-level — the message-layer round-trip is covered in {@link MessageServiceTest}.
*/
class RendezvousTest {
private static final String W = "term_a";
private final Rendezvous rendezvous = new Rendezvous();
@Test
void openAskMintsAUniqueTurnScopedToItsSession() {
Rendezvous.AskTicket t1 = rendezvous.openAsk(W);
Rendezvous.AskTicket t2 = rendezvous.openAsk(W);
assertNotEquals(t1.turnId(), t2.turnId(), "each ask gets its own turnId");
assertTrue(t1.turnId().startsWith(W + "#"), "the turnId is scoped to the worker session");
assertEquals(W, rendezvous.askSession(t1.turnId()));
assertEquals(W, rendezvous.askSession(t2.turnId()));
}
@Test
void answerAskCompletesTheWaitersFuture() {
Rendezvous.AskTicket t = rendezvous.openAsk(W);
assertTrue(rendezvous.answerAsk(t.turnId(), "config.yaml"), "answering an open ask succeeds");
assertEquals("config.yaml", t.answer().getNow(null), "the answer reaches the blocked worker");
}
@Test
void answerAskOnAnUnknownTurnIsFalse() {
assertFalse(rendezvous.answerAsk("no-such#1", "x"), "an answer to an unknown turn is a no-op");
}
@Test
void resolveQuestionResolvesAnOpenSendWithTheQuestionKindAndTurnId() {
CompletableFuture<Rendezvous.Resolution> send = rendezvous.open(W);
assertTrue(rendezvous.resolveQuestion(W, "which config?", W + "#7"),
"the question resolves the primary's open send");
Rendezvous.Resolution r = send.getNow(null);
assertEquals(Rendezvous.Kind.QUESTION, r.kind());
assertEquals("which config?", r.text());
assertEquals(W + "#7", r.turnId(), "the turnId rides along so the primary can answer");
}
@Test
void resolveQuestionWithNoOpenSendIsFalse() {
assertFalse(rendezvous.resolveQuestion(W, "anyone?", W + "#1"),
"no blocked send means no primary to surface the question to");
}
@Test
void closeAskRemovesTheTurn() {
Rendezvous.AskTicket t = rendezvous.openAsk(W);
rendezvous.closeAsk(t.turnId());
assertNull(rendezvous.askSession(t.turnId()), "a closed ask is forgotten");
assertFalse(rendezvous.answerAsk(t.turnId(), "late"), "a closed ask can no longer be answered");
}
}