From 358c6970b5787f19fcbcac1762236e968f9dc9e2 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 16 Jul 2026 15:32:56 +0200 Subject: [PATCH] CB-205/CB-201: bridge_ask reverse rendezvous + lightweight question/turn kind MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- .../java/dev/ltms/bridged/mcp/BridgeMcp.java | 113 +++++++++++--- .../dev/ltms/bridged/msg/MessageService.java | 142 ++++++++++++++++-- .../java/dev/ltms/bridged/msg/Rendezvous.java | 84 ++++++++++- .../dev/ltms/bridged/rest/BridgedApp.java | 102 ++++++++++--- .../dev/ltms/bridged/mcp/BridgeMcpTest.java | 57 +++++++ .../ltms/bridged/msg/MessageServiceTest.java | 66 ++++++++ .../dev/ltms/bridged/msg/RendezvousTest.java | 69 +++++++++ 7 files changed, 578 insertions(+), 55 deletions(-) create mode 100644 bridged/src/test/java/dev/ltms/bridged/msg/RendezvousTest.java 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 27a3d7d..f9ccf2a 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java @@ -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 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() { 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 e1fa187..4b98a40 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -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 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. + * + *

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. + * + *

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 before 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 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); 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 1cc023c..7163f49 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java @@ -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 answer) { + } + + /** Handle to a just-opened reverse-rendezvous turn: the minted {@code turnId} and its answer future. */ + public record AskTicket(String turnId, CompletableFuture answer) { } private final ConcurrentHashMap> waiters = new ConcurrentHashMap<>(); + /** Reverse rendezvous (CB-205): worker questions awaiting the primary's answer, keyed by {@code turnId}. */ + private final ConcurrentHashMap 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 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 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 04718b1..f3c64d4 100644 --- a/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java +++ b/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java @@ -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")); } } diff --git a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java index 28b9ad8..81cd6e8 100644 --- a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java @@ -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 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 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 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(); diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java index 52ef9c6..bd74066 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -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 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 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 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 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"); + } } diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/RendezvousTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/RendezvousTest.java new file mode 100644 index 0000000..6848a75 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/msg/RendezvousTest.java @@ -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 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"); + } +}