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.
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -52,13 +52,18 @@ public final class BridgeMcp {
|
||||
.capabilities(McpSchema.ServerCapabilities.builder().tools(true).build())
|
||||
.toolCall(sendTool(), (_, req) -> {
|
||||
Map<String, Object> 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",
|
||||
|
||||
@@ -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).
|
||||
*
|
||||
* <p>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}).
|
||||
*
|
||||
* <p><strong>Async fire-and-poll (CB-107).</strong> 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 <em>ticket</em> 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<Reply> future, long createdNanos) {
|
||||
}
|
||||
|
||||
private final AgentControl agents;
|
||||
private final Injector injector;
|
||||
private final Rendezvous rendezvous;
|
||||
private final ConcurrentHashMap<String, ReentrantLock> sessionLocks = new ConcurrentHashMap<>();
|
||||
private final ConcurrentHashMap<String, Task> 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<Reply> 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<Reply> 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);
|
||||
|
||||
@@ -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<String, Object> 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")) {
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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<String> 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<String> 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<String> 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();
|
||||
|
||||
Reference in New Issue
Block a user