From 8ed2370fbecd2eee0195341ba9095c30d855c00f Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Wed, 15 Jul 2026 14:19:14 +0200 Subject: [PATCH] CB-107: async fire-and-poll delegation (wait:false + ticket poll) A caller's MCP client caps a blocking bridge_send at ~60s, but a real delegated task runs for minutes. sendAsync runs the same blocking send on a background virtual thread and returns a ticket; poll(ticket) reports pending/done/failed. Async reuses the blocking path (and its per-target serialization), so it inherits reply + completion resolution for free. Surfaces: REST POST message wait:false -> 202 {ticket} + GET /tasks/{ticket}; MCP bridge_send wait flag + new bridge_poll. Terminal tickets are pruned after a TTL so the registry stays bounded. --- .../main/java/dev/ltms/bridged/Bridged.java | 1 + .../java/dev/ltms/bridged/mcp/BridgeMcp.java | 55 +++++++- .../dev/ltms/bridged/msg/MessageService.java | 119 +++++++++++++++++- .../dev/ltms/bridged/rest/BridgedApp.java | 32 ++++- .../dev/ltms/bridged/mcp/BridgeMcpTest.java | 37 ++++++ .../dev/ltms/bridged/rest/BridgedAppTest.java | 44 +++++++ 6 files changed, 280 insertions(+), 8 deletions(-) diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 665b9cf..1f10f10 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -64,6 +64,7 @@ public final class Bridged { Runtime.getRuntime().addShutdownHook(new Thread(poller::stop)); MessageService messages = new MessageService(agents, injector, rendezvous); + Runtime.getRuntime().addShutdownHook(new Thread(messages::close)); // MCP server face (CB-105): bridge_send/bridge_reply/bridge_status, mounted at /mcp. // Caller identity is resolved from the connection (peer PID → herdr pane), not arguments. 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 51b5406..db19a7b 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java @@ -52,13 +52,18 @@ public final class BridgeMcp { .capabilities(McpSchema.ServerCapabilities.builder().tools(true).build()) .toolCall(sendTool(), (_, req) -> { Map a = req.arguments(); - return send(messages, str(a, "sessionId"), 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")) + : send(messages, str(a, "sessionId"), str(a, "content"), timeoutMs(a)); }) // bridge_reply's identity is the CONNECTION, never an argument. .toolCall(replyTool(), (exchange, req) -> reply(rendezvous, callerTerminal(exchange), str(req.arguments(), "content"))) .toolCall(statusTool(), (_, req) -> status(messages, str(req.arguments(), "sessionId"))) + .toolCall(pollTool(), (_, req) -> + poll(messages, str(req.arguments(), "ticket"))) .build(); } @@ -109,6 +114,36 @@ public final class BridgeMcp { } } + /** + * {@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. + */ + static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content) { + if (isBlank(sessionId) || isBlank(content)) { + return error("sessionId and content are required"); + } + String ticket = messages.sendAsync(sessionId, content); + return text("accepted — task delegated. Poll bridge_poll with ticket=" + ticket); + } + + /** {@code bridge_poll}: check an async delegation by ticket (pending / done+reply / failed). */ + static McpSchema.CallToolResult poll(MessageService messages, String ticket) { + if (isBlank(ticket)) { + return error("ticket is required"); + } + MessageService.TaskView v = messages.poll(ticket); + if (v == null) { + return error("unknown ticket: " + ticket + " (never issued, or expired)"); + } + return switch (v.phase()) { + case DONE -> text(v.replySource() != null && v.replySource().equals("transcript") + ? "[done — worker finished without a structured bridge_reply; transcript tail follows]\n" + v.reply() + : v.reply()); + case PENDING -> text("[pending — " + v.detail() + "]"); + case FAILED -> text("[failed — " + v.detail() + "]"); + }; + } + /** * {@code bridge_reply}: the worker returns its structured answer, resolving the awaiting send. * {@code callerTerminal} is resolved from the connection (never an argument); a {@code null} @@ -143,15 +178,27 @@ public final class BridgeMcp { private static McpSchema.Tool sendTool() { return tool("bridge_send", - "Delegate a task to a worker session and block until it replies. " - + "Returns the worker's reply, or a 'still working / queued' note on timeout.", + "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.", 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"), - "timeoutMs", Map.of("type", "integer", "description", "Max ms to wait for a reply")), + "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"))); } + private static McpSchema.Tool pollTool() { + return tool("bridge_poll", + "Check an async delegation (a bridge_send with wait:false) by its ticket: " + + "pending, done (with the worker's reply), or failed.", + objectSchema(Map.of( + "ticket", stringProp("The ticket returned by bridge_send wait:false")), + List.of("ticket"))); + } + private static McpSchema.Tool replyTool() { // No session/target arg — the worker's identity is resolved from the connection. return tool("bridge_reply", 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 fe1fb0e..94ae3eb 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -7,10 +7,14 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.locks.ReentrantLock; /** @@ -24,14 +28,29 @@ import java.util.concurrent.locks.ReentrantLock; * what lets a reply map unambiguously to its send (no cross-talk between concurrent callers). * *

If the worker never replies within the timeout, the caller gets a typed "still working" / - * "queued" outcome — the message may still be mid-flight. (A herdr {@code events.subscribe} - * completion signal is the planned fallback resolver for workers that don't call {@code - * bridge_reply}; the rendezvous is the resolver seam it will plug into.) + * "queued" outcome — the message may still be mid-flight. A finished-but-unreplied turn is caught + * by the CB-106 completion fallback (see {@link Rendezvous#resolveCompletion}). + * + *

Async fire-and-poll (CB-107). A caller's MCP client caps a blocking call at + * ~60s, but a real delegated task runs for minutes. {@link #sendAsync} therefore runs the same + * blocking {@link #send} on a background virtual thread and hands back a ticket the caller + * polls with {@link #poll}. The blocking and async paths share one code path (and the same per-target + * serialization), so async inherits the reply + completion resolution behaviour for free. */ public final class MessageService { private static final Logger log = LoggerFactory.getLogger(MessageService.class); + /** + * The window a fire-and-poll send waits for resolution — generous, since no caller is blocked on + * it; a real delegated task resolves (reply or completion) well within this, and only a genuinely + * hung worker rides it out. + */ + private static final long ASYNC_TIMEOUT_MS = 30 * 60 * 1_000L; + + /** How long a finished (terminal) ticket is retained for polling before it is pruned. */ + private static final long TICKET_TTL_NANOS = 10 * 60 * 1_000_000_000L; + /** Outcome of a blocking send. */ public enum Outcome { /** The worker called {@code bridge_reply}; {@code text} holds the structured answer. */ @@ -62,10 +81,39 @@ public final class MessageService { } } + /** Lifecycle phase of an async delegation ticket. */ + public enum Phase { + /** Delegated and in flight — queued for the worker or being worked. */ + PENDING, + /** The worker's turn finished; {@link TaskView#reply} holds the answer. */ + DONE, + /** The delegation could not complete (timed out, worker gone, or busy). */ + FAILED + } + + /** + * A poll snapshot of an async delegation. + * + * @param reply the answer when {@link #phase} is {@link Phase#DONE}, else {@code null} + * @param replySource {@code "reply"} (structured {@code bridge_reply}) or {@code "transcript"} + * (completion scrape) when {@link Phase#DONE}, else {@code null} + * @param detail a human note (live worker status while pending, or the failure reason) + */ + public record TaskView(String ticket, Phase phase, String reply, String replySource, String detail) { + } + + /** An in-flight or finished async delegation, keyed by its ticket. */ + private record Task(String target, CompletableFuture future, long createdNanos) { + } + private final AgentControl agents; private final Injector injector; private final Rendezvous rendezvous; private final ConcurrentHashMap sessionLocks = new ConcurrentHashMap<>(); + private final ConcurrentHashMap tasks = new ConcurrentHashMap<>(); + private final AtomicLong ticketSeq = new AtomicLong(); + private final ExecutorService asyncExecutor = Executors.newThreadPerTaskExecutor( + Thread.ofVirtual().name("bridge-async-", 0).factory()); public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous) { this.agents = agents; @@ -116,6 +164,71 @@ public final class MessageService { } } + /** + * 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 + * long task is delegated without tripping the caller's MCP client call timeout. + * + * @return the ticket to poll for the eventual result + */ + public String sendAsync(String target, String content) { + String ticket = "task-" + ticketSeq.incrementAndGet(); + CompletableFuture future = + CompletableFuture.supplyAsync(() -> send(target, content, ASYNC_TIMEOUT_MS), asyncExecutor); + tasks.put(ticket, new Task(target, future, System.nanoTime())); + pruneTerminalTickets(); + log.debug("async send {} -> {}", ticket, target); + return ticket; + } + + /** + * Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket; + * otherwise a {@link Phase#PENDING} view (with the live worker status as detail), a + * {@link Phase#DONE} view carrying the reply, or a {@link Phase#FAILED} view with the reason. + */ + public TaskView poll(String ticket) { + Task task = tasks.get(ticket); + if (task == null) { + return null; + } + CompletableFuture f = task.future(); + if (!f.isDone()) { + return new TaskView(ticket, Phase.PENDING, null, null, "worker " + liveStatus(task.target())); + } + Reply r; + try { + r = f.getNow(null); + } catch (CompletionException | java.util.concurrent.CancellationException e) { + Throwable cause = (e instanceof CompletionException ce && ce.getCause() != null) ? ce.getCause() : e; + return new TaskView(ticket, Phase.FAILED, null, null, cause.getMessage()); + } + if (r.completed()) { + String source = r.outcome() == Outcome.REPLIED ? "reply" : "transcript"; + return new TaskView(ticket, Phase.DONE, r.text(), source, null); + } + return new TaskView(ticket, Phase.FAILED, null, null, "no reply — " + r.outcome().name().toLowerCase()); + } + + /** Best-effort live worker status for a pending poll; never throws (a lookup error is just noise). */ + private String liveStatus(String target) { + try { + return agents.status(target).name().toLowerCase(); + } catch (RuntimeException e) { + return "unknown"; + } + } + + /** Drop finished tickets older than the TTL so the registry cannot grow without bound. */ + private void pruneTerminalTickets() { + long cutoff = System.nanoTime() - TICKET_TTL_NANOS; + tasks.values().removeIf(t -> t.future().isDone() && t.createdNanos() < cutoff); + } + + /** Release the async executor. */ + public void close() { + asyncExecutor.shutdown(); + } + 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/rest/BridgedApp.java b/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java index c81c339..ae87208 100644 --- a/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java +++ b/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java @@ -65,9 +65,10 @@ public final class BridgedApp { app.get("/agents", this::agents); app.post("/workers", this::spawnWorker); app.delete("/workers/{paneId}", this::stopWorker); - app.post("/sessions/{id}/message", this::sendMessage); // bridge_send (primary, blocking) + app.post("/sessions/{id}/message", this::sendMessage); // bridge_send (primary; blocking or wait:false) app.post("/sessions/{id}/reply", this::replyMessage); // bridge_reply (worker) app.get("/sessions/{id}/status", this::sessionStatus); // bridge_status + app.get("/tasks/{ticket}", this::taskStatus); // poll an async (wait:false) send return app; } @@ -134,10 +135,12 @@ public final class BridgedApp { String id = ctx.pathParam("id"); String content; long timeout; + boolean wait; try { JsonNode body = mapper.readTree(ctx.body()); content = body.path("content").asText(""); 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) { ctx.status(400).json(Map.of("error", "bad_request", "detail", "body must be JSON")); return; @@ -146,6 +149,13 @@ public final class BridgedApp { ctx.status(400).json(Map.of("error", "bad_request", "detail", "content is required")); return; } + + if (!wait) { + // Fire-and-poll (CB-107): return a ticket immediately; the caller polls GET /tasks/{ticket}. + String ticket = messages.sendAsync(id, content); + ctx.status(202).json(Map.of("sessionId", id, "ticket", ticket, "status", "accepted")); + return; + } timeout = Math.clamp(timeout, 1, MAX_MESSAGE_TIMEOUT_MS); try { @@ -206,6 +216,26 @@ public final class BridgedApp { } } + /** Poll an async (wait:false) delegation by ticket. 404 for an unknown/expired ticket. */ + private void taskStatus(Context ctx) { + MessageService.TaskView v = messages.poll(ctx.pathParam("ticket")); + if (v == null) { + ctx.status(404).json(Map.of("error", "unknown_ticket", "detail", "no such task (or it has expired)")); + return; + } + Map body = new LinkedHashMap<>(); + body.put("ticket", v.ticket()); + body.put("phase", v.phase().name().toLowerCase()); + if (v.reply() != null) { + body.put("reply", v.reply()); + body.put("replySource", v.replySource()); + } + if (v.detail() != null) { + body.put("detail", v.detail()); + } + ctx.status(200).json(body); + } + /** Map a herdr failure: unknown target → 404, anything else → 502 (herdr is upstream). */ private static void herdrError(Context ctx, HerdrException e) { if (e.code() != null && e.code().endsWith("_not_found")) { 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 6eee470..ea611e6 100644 --- a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java @@ -49,6 +49,43 @@ class BridgeMcpTest { assertEquals("LGTM", textOf(res)); } + @Test + void asyncSendReturnsATicketThenPollReportsTheReply() throws Exception { + // wait:false parity — a ticket is issued, resolved by a reply, and surfaced by bridge_poll. + McpSchema.CallToolResult accepted = BridgeMcp.sendAsync(messages, "term_a", "do it"); + assertNotEquals(Boolean.TRUE, accepted.isError()); + String out = textOf(accepted); + assertTrue(out.contains("ticket="), out); + String ticket = out.substring(out.indexOf("ticket=") + "ticket=".length()).trim(); + + // Resolve the awaiting send once it has opened (retry past the async-open race). + long deadline = System.currentTimeMillis() + 3000; + McpSchema.CallToolResult reply = BridgeMcp.reply(rendezvous, "term_a", "async LGTM"); + while (Boolean.TRUE.equals(reply.isError()) && System.currentTimeMillis() < deadline) { + //noinspection BusyWait + Thread.sleep(10); + reply = BridgeMcp.reply(rendezvous, "term_a", "async LGTM"); + } + assertEquals("delivered", textOf(reply)); + + // Poll until the async send completes and reports the reply. + McpSchema.CallToolResult polled = BridgeMcp.poll(messages, ticket); + deadline = System.currentTimeMillis() + 3000; + while (!textOf(polled).contains("async LGTM") && System.currentTimeMillis() < deadline) { + //noinspection BusyWait + Thread.sleep(10); + polled = BridgeMcp.poll(messages, ticket); + } + assertEquals("async LGTM", textOf(polled)); + } + + @Test + void pollUnknownTicketIsAnError() { + McpSchema.CallToolResult res = BridgeMcp.poll(messages, "task-999"); + assertTrue(res.isError()); + assertTrue(textOf(res).contains("unknown ticket")); + } + @Test void sendTimesOutWithAWorkingNote() { McpSchema.CallToolResult res = BridgeMcp.send(messages, "term_a", "hi", 120L); diff --git a/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java b/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java index b3f1ad2..e204de8 100644 --- a/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java @@ -264,6 +264,50 @@ class BridgedAppTest { // (injection via agent.send is covered deterministically by the timeout-working test) } + @Test + void asyncSendReturnsATicketThenPollReportsTheReply() throws Exception { + // CB-107 fire-and-poll: wait:false returns a ticket immediately; the result is polled. + FakeHerdr herdr = new FakeHerdr().agentStatus("idle"); // poller delivers the injection + int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw")); + + HttpResponse accepted = postMessage(port, "{\"content\":\"do it\",\"wait\":false}"); + assertEquals(202, accepted.statusCode()); + String ticket = mapper.readTree(accepted.body()).get("ticket").asText(); + assertFalse(ticket.isBlank(), "an async send must return a ticket"); + + // The worker replies once the async send is actually awaiting (retry past the startup race). + HttpResponse reply; + long deadline = System.currentTimeMillis() + 3000; + do { + reply = postJson(port, "/sessions/term_a/reply", "{\"content\":\"async LGTM\"}"); + if (reply.statusCode() != 409) break; + //noinspection BusyWait + Thread.sleep(10); + } while (System.currentTimeMillis() < deadline); + assertEquals(200, reply.statusCode()); + + // Polling the ticket now reports the finished delegation and its reply. + JsonNode task; + deadline = System.currentTimeMillis() + 3000; + do { + task = mapper.readTree(req(port, "GET", "/tasks/" + ticket).body()); + if ("done".equals(task.path("phase").asText())) break; + //noinspection BusyWait + Thread.sleep(10); + } while (System.currentTimeMillis() < deadline); + assertEquals("done", task.get("phase").asText()); + assertEquals("async LGTM", task.get("reply").asText()); + assertEquals("reply", task.get("replySource").asText()); + } + + @Test + void pollUnknownTicketIs404() throws Exception { + int port = startHealthy(); + HttpResponse res = req(port, "GET", "/tasks/task-999"); + assertEquals(404, res.statusCode()); + assertEquals("unknown_ticket", mapper.readTree(res.body()).get("error").asText()); + } + @Test void replyWithNoPendingSendIsConflict() throws Exception { int port = startHealthy();