CB-104: blocking bridge_send with rendezvous reply (POST /sessions/{id}/message + /reply)

Per-session-serialized blocking send that enqueues via the CB-103 injector (poller delivers) and blocks on a rendezvous resolved by the worker's structured bridge_reply, or a typed 200/202 outcome. No status-polling completion, no terminal scrape. GET /sessions/{id}/status. Reworked from an initial poll+scrape draft after a high-effort review found the polling completion unreliable; all findings fixed.
This commit is contained in:
Dai Ha
2026-07-14 14:20:05 +02:00
parent bc04637694
commit 597ac2e562
6 changed files with 369 additions and 6 deletions
@@ -7,6 +7,8 @@ import dev.ltms.bridged.herdr.UnixSocketHerdrClient;
import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.inject.Injector;
import dev.ltms.bridged.inject.StatusPoller;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.rest.BridgedApp;
import dev.ltms.bridged.worker.WorkerService;
import io.javalin.Javalin;
@@ -53,7 +55,9 @@ public final class Bridged {
poller.start();
Runtime.getRuntime().addShutdownHook(new Thread(poller::stop));
Javalin app = new BridgedApp(herdr, workers).build();
Rendezvous rendezvous = new Rendezvous();
MessageService messages = new MessageService(agents, injector, rendezvous);
Javalin app = new BridgedApp(herdr, workers, messages, rendezvous).build();
app.start(cfg.bind().host(), cfg.bind().port());
log.info("bridged listening on {}:{}, herdr socket {}",
cfg.bind().host(), cfg.bind().port(), socket);
@@ -0,0 +1,117 @@
package dev.ltms.bridged.msg;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.inject.Injector;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.locks.ReentrantLock;
/**
* The blocking delegation feature (CB-104): deliver {@code content} into a worker and block until
* the worker returns a <em>structured reply</em> via {@code bridge_reply} (the {@link Rendezvous}),
* then hand that reply back. Delivery is the {@link Injector}'s job (the background poller sends it
* when the worker is injectable); this service never drives the injector or scrapes the terminal —
* completion is the worker's explicit reply, not a guess about {@code agent_status}.
*
* <p>Sends are serialized per session so exactly one reply can be outstanding per worker, which is
* 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.)
*/
public final class MessageService {
private static final Logger log = LoggerFactory.getLogger(MessageService.class);
/** Outcome of a blocking send. */
public enum Outcome {
/** The worker replied; {@code text} holds the structured answer. */
REPLIED,
/** 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
}
/** @param text the worker's reply when {@link #outcome} is {@link Outcome#REPLIED}, else {@code null} */
public record Reply(Outcome outcome, String text) {
public boolean completed() {
return outcome == Outcome.REPLIED;
}
}
private final AgentControl agents;
private final Injector injector;
private final Rendezvous rendezvous;
private final ConcurrentHashMap<String, ReentrantLock> sessionLocks = new ConcurrentHashMap<>();
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous) {
this.agents = agents;
this.injector = injector;
this.rendezvous = rendezvous;
}
/** Current lifecycle status of a worker (the {@code GET /sessions/{id}/status} surface). */
public AgentStatus status(String target) {
return agents.status(target);
}
/**
* Deliver {@code content} to {@code target} (a herdr {@code terminal_id}) and block until the
* worker replies via {@link Rendezvous} or {@code timeoutMillis} elapses.
*/
public Reply send(String target, String content, long timeoutMillis) {
long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L;
ReentrantLock lock = sessionLocks.computeIfAbsent(target, _ -> new ReentrantLock());
if (!tryLock(lock, remainingMillis(deadlineNanos))) {
return new Reply(Outcome.BUSY, null); // another send held the session the whole window
}
try {
CompletableFuture<Void> delivered = injector.enqueue(target, content);
CompletableFuture<String> reply = rendezvous.open(target);
try {
String text = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
return new Reply(Outcome.REPLIED, text);
} catch (TimeoutException e) {
boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally();
log.debug("send to {} timed out (delivered={})", target, wasDelivered);
return new Reply(wasDelivered ? Outcome.TIMED_OUT_WORKING : Outcome.TIMED_OUT_QUEUED, 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 " + target, e);
} finally {
rendezvous.close(target, reply);
}
} finally {
lock.unlock();
}
}
private static boolean tryLock(ReentrantLock lock, long millis) {
try {
return lock.tryLock(Math.max(0, millis), TimeUnit.MILLISECONDS);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new IllegalStateException("interrupted awaiting the session send lock", e);
}
}
private static long remainingMillis(long deadlineNanos) {
return (deadlineNanos - System.nanoTime()) / 1_000_000L;
}
}
@@ -0,0 +1,41 @@
package dev.ltms.bridged.msg;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
/**
* The reply rendezvous: where a blocking {@code bridge_send} awaits the worker's structured
* {@code bridge_reply} for the same session. The sending (primary) request thread {@link #open}s
* a waiter; the worker's reply (arriving on a different thread via {@code POST
* /sessions/{id}/reply}) {@link #resolve}s it.
*
* <p>At most one waiter per session — {@link MessageService} serializes sends per session, so a
* reply maps unambiguously to the one outstanding send and cannot be captured by another.
*/
public final class Rendezvous {
private final ConcurrentHashMap<String, CompletableFuture<String>> waiters = new ConcurrentHashMap<>();
/** Register a waiter for {@code session}. The caller must hold that session's send lock. */
CompletableFuture<String> open(String session) {
CompletableFuture<String> waiter = new CompletableFuture<>();
waiters.put(session, waiter);
return waiter;
}
/** Remove {@code waiter} for {@code session} (only if it is still the registered one). */
void close(String session, CompletableFuture<String> waiter) {
waiters.remove(session, waiter);
}
/**
* Resolve the send awaiting on {@code session} with the worker's reply {@code content}.
*
* @return {@code true} if a waiter was resolved; {@code false} if none was waiting (a late or
* spurious reply — e.g. the send already timed out)
*/
public boolean resolve(String session, String content) {
CompletableFuture<String> waiter = waiters.get(session);
return waiter != null && waiter.complete(content);
}
}
@@ -1,10 +1,13 @@
package dev.ltms.bridged.rest;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import dev.ltms.bridged.guard.GuardException;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.herdr.HerdrClient;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.worker.WorkerService;
import io.javalin.Javalin;
import io.javalin.http.Context;
@@ -25,12 +28,21 @@ import java.util.Map;
*/
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;
private final HerdrClient herdr;
private final WorkerService workers;
private final MessageService messages;
private final Rendezvous rendezvous;
private final ObjectMapper mapper = new ObjectMapper();
public BridgedApp(HerdrClient herdr, WorkerService workers) {
public BridgedApp(HerdrClient herdr, WorkerService workers, MessageService messages, Rendezvous rendezvous) {
this.herdr = herdr;
this.workers = workers;
this.messages = messages;
this.rendezvous = rendezvous;
}
/** Wire routes onto a fresh, unstarted Javalin instance. Caller starts it. */
@@ -41,6 +53,9 @@ 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}/reply", this::replyMessage); // bridge_reply (worker)
app.get("/sessions/{id}/status", this::sessionStatus); // bridge_status
return app;
}
@@ -97,6 +112,93 @@ public final class BridgedApp {
ctx.status(204);
}
/**
* The blocking delegation call (CB-104): inject {@code content} into the worker via the
* status-gated injector and block until the worker returns a structured {@code bridge_reply}.
* Times out with a typed 202 (working / queued / busy) rather than an error — the message may
* still land.
*/
private void sendMessage(Context ctx) {
String id = ctx.pathParam("id");
String content;
long timeout;
try {
JsonNode body = mapper.readTree(ctx.body());
content = body.path("content").asText("");
timeout = body.path("timeoutMs").asLong(DEFAULT_MESSAGE_TIMEOUT_MS);
} catch (Exception e) {
ctx.status(400).json(Map.of("error", "bad_request", "detail", "body must be JSON"));
return;
}
if (content.isBlank()) {
ctx.status(400).json(Map.of("error", "bad_request", "detail", "content is required"));
return;
}
timeout = Math.clamp(timeout, 1, MAX_MESSAGE_TIMEOUT_MS);
try {
MessageService.Reply reply = messages.send(id, content, timeout);
if (reply.completed()) {
ctx.status(200).json(Map.of("sessionId", id, "reply", reply.text()));
} 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 REPLIED -> "done"; // unreachable
},
"detail", "no reply within " + timeout + "ms; poll status or retry"));
}
} catch (HerdrException e) {
herdrError(ctx, e);
}
}
/**
* The worker's structured reply ({@code bridge_reply}) — resolves the blocking send awaiting
* on this session. 200 if a send was waiting, 409 if none was (late or spurious reply).
*/
private void replyMessage(Context ctx) {
String id = ctx.pathParam("id");
String content;
try {
content = mapper.readTree(ctx.body()).path("content").asText("");
} catch (Exception e) {
ctx.status(400).json(Map.of("error", "bad_request", "detail", "body must be JSON"));
return;
}
if (rendezvous.resolve(id, content)) {
ctx.status(200).json(Map.of("sessionId", id, "delivered", true));
} else {
ctx.status(409).json(Map.of(
"sessionId", id, "error", "no_pending_send",
"detail", "no send is awaiting a reply for this session"));
}
}
/** Live lifecycle status of a worker (MCP `bridge_status` wraps this in CB-105). */
private void sessionStatus(Context ctx) {
String id = ctx.pathParam("id");
try {
ctx.status(200).json(Map.of(
"sessionId", id,
"status", messages.status(id).name().toLowerCase()));
} catch (HerdrException e) {
herdrError(ctx, e);
}
}
/** 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")) {
ctx.status(404).json(Map.of("error", "session_not_found", "detail", e.getMessage()));
} else {
ctx.status(502).json(Map.of("error", "herdr_error", "detail", e.getMessage()));
}
}
/** Stable JSON projection of an agent (null-safe for the start-time shape). */
private static Map<String, Object> view(Agent a) {
Map<String, Object> m = new LinkedHashMap<>();
@@ -24,7 +24,7 @@ public final class FakeHerdr implements HerdrClient {
private int workerTabPaneCount = 1;
private String paneCloseErrorCode = null;
private String agentSendErrorCode = null;
private volatile String agentStatus = "idle"; // what agent.get reports
private volatile String agentStatus = "idle"; // steady-state agent.get status
public FakeHerdr healthy(boolean h) {
this.healthy = h;
@@ -61,6 +61,7 @@ public final class FakeHerdr implements HerdrClient {
return this;
}
/** Seed an additional workspace into {@code workspace.list} (e.g. a pre-existing worker space). */
public FakeHerdr withWorkspace(String id, String label) {
extraWorkspaces.add(("{\"workspace_id\":\"%s\",\"label\":\"%s\",\"focused\":false,"
@@ -7,6 +7,10 @@ import dev.ltms.bridged.guard.SubscriptionGuard;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.inject.Injector;
import dev.ltms.bridged.inject.StatusPoller;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.worker.WorkerService;
import io.javalin.Javalin;
import org.junit.jupiter.api.AfterEach;
@@ -33,9 +37,11 @@ class BridgedAppTest {
private final ObjectMapper mapper = new ObjectMapper();
private final HttpClient http = HttpClient.newHttpClient();
private Javalin app;
private StatusPoller poller;
@AfterEach
void stop() {
if (poller != null) poller.stop();
if (app != null) app.stop();
}
@@ -47,10 +53,16 @@ class BridgedAppTest {
BridgedConfig.Worker wcfg = new BridgedConfig.Worker(
"ltms-local", workerBaseUrl, "coder", null, "BRIDGED_WORKER_TOKEN", null,
placement, "bridged-workers", "worker: {profile} #{n}");
AgentControl agents = new AgentControl(herdr);
WorkerService workers = new WorkerService(
new AgentControl(herdr), new WorkspaceControl(herdr), new SubscriptionGuard(allow), wcfg,
agents, new WorkspaceControl(herdr), new SubscriptionGuard(allow), wcfg,
k -> "BRIDGED_WORKER_TOKEN".equals(k) ? "tok-abc" : null);
app = new BridgedApp(herdr, workers).build().start("127.0.0.1", 0);
Injector injector = new Injector(agents);
poller = new StatusPoller(agents, injector, 5); // delivers when the fake reports idle
poller.start();
Rendezvous rendezvous = new Rendezvous();
MessageService messages = new MessageService(agents, injector, rendezvous);
app = new BridgedApp(herdr, workers, messages, rendezvous).build().start("127.0.0.1", 0);
return app.port();
}
@@ -58,6 +70,18 @@ class BridgedAppTest {
return start(new FakeHerdr(), "http://gx00.gw:8000", Set.of("gx00.gw"));
}
private HttpResponse<String> postMessage(int port, String json) throws Exception {
return postJson(port, "/sessions/term_a/message", json);
}
private HttpResponse<String> postJson(int port, String path, String json) throws Exception {
HttpRequest r = HttpRequest.newBuilder(URI.create("http://127.0.0.1:" + port + path))
.header("Content-Type", "application/json")
.POST(HttpRequest.BodyPublishers.ofString(json)).build();
return http.send(r, HttpResponse.BodyHandlers.ofString());
}
private HttpResponse<String> req(int port, String method, String path) throws Exception {
HttpRequest.Builder b = HttpRequest.newBuilder(URI.create("http://127.0.0.1:" + port + path));
b = switch (method) {
@@ -85,7 +109,8 @@ class BridgedAppTest {
@Test
void healthzDegradedWhenHerdrDown() throws Exception {
int port = start(new FakeHerdr().healthy(false), "http://gx00.gw:8000", Set.of("gx00.gw"));
FakeHerdr down = new FakeHerdr().healthy(false);
int port = start(down, "http://gx00.gw:8000", Set.of("gx00.gw"));
HttpResponse<String> res = req(port, "GET", "/healthz");
assertEquals(503, res.statusCode());
assertEquals("degraded", mapper.readTree(res.body()).get("status").asText());
@@ -211,6 +236,79 @@ class BridgedAppTest {
assertEquals("w9:t2", params(herdr, "tab.close").get("tab_id"));
}
@Test
void messageReturnsTheWorkersStructuredReply() throws Exception {
// CB-104 (option C): the blocking send resolves on the worker's bridge_reply, not a scrape.
FakeHerdr herdr = new FakeHerdr().agentStatus("idle"); // poller delivers the injection
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
var send = java.util.concurrent.CompletableFuture.supplyAsync(() -> {
try { return postMessage(port, "{\"content\":\"review this\",\"timeoutMs\":4000}"); }
catch (Exception e) { throw new RuntimeException(e); }
});
// The worker replies once a 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\":\"LGTM ship it\"}");
if (reply.statusCode() != 409) break;
//noinspection BusyWait
Thread.sleep(10);
} while (System.currentTimeMillis() < deadline);
assertEquals(200, reply.statusCode());
HttpResponse<String> res = send.get(6, java.util.concurrent.TimeUnit.SECONDS);
assertEquals(200, res.statusCode());
assertEquals("LGTM ship it", mapper.readTree(res.body()).get("reply").asText());
// (injection via agent.send is covered deterministically by the timeout-working test)
}
@Test
void replyWithNoPendingSendIsConflict() throws Exception {
int port = startHealthy();
HttpResponse<String> res = postJson(port, "/sessions/term_a/reply", "{\"content\":\"orphan\"}");
assertEquals(409, res.statusCode());
assertEquals("no_pending_send", mapper.readTree(res.body()).get("error").asText());
}
@Test
void messageTimesOutQueuedWhenWorkerNeverInjectable() throws Exception {
FakeHerdr herdr = new FakeHerdr().agentStatus("working"); // never injectable → never delivered
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
HttpResponse<String> res = postMessage(port, "{\"content\":\"hi\",\"timeoutMs\":150}");
assertEquals(202, res.statusCode());
assertEquals("queued", mapper.readTree(res.body()).get("status").asText());
assertFalse(herdr.called("agent.send"), "no injection while the worker is mid-turn");
}
@Test
void messageTimesOutWorkingWhenDeliveredButNoReply() throws Exception {
FakeHerdr herdr = new FakeHerdr().agentStatus("idle"); // delivered, but nobody replies
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
HttpResponse<String> res = postMessage(port, "{\"content\":\"hi\",\"timeoutMs\":250}");
assertEquals(202, res.statusCode());
assertEquals("working", mapper.readTree(res.body()).get("status").asText());
assertTrue(herdr.called("agent.send"), "message was injected");
}
@Test
void messageRejectsBlankContent() throws Exception {
int port = startHealthy();
assertEquals(400, postMessage(port, "{}").statusCode());
}
@Test
void sessionStatusReportsLiveAgentStatus() throws Exception {
FakeHerdr herdr = new FakeHerdr().agentStatus("blocked");
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
HttpResponse<String> res = req(port, "GET", "/sessions/term_a/status");
assertEquals(200, res.statusCode());
assertEquals("blocked", mapper.readTree(res.body()).get("status").asText());
}
@Test
void stopWorkerInPanePlacementClosesOnlyThePane() throws Exception {
FakeHerdr herdr = new FakeHerdr();