diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index a59d7d3..c323297 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -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); diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java new file mode 100644 index 0000000..a45dc7e --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -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 structured reply 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}. + * + *

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

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 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 delivered = injector.enqueue(target, content); + CompletableFuture 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; + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java b/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java new file mode 100644 index 0000000..b42a84c --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java @@ -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. + * + *

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> waiters = new ConcurrentHashMap<>(); + + /** Register a waiter for {@code session}. The caller must hold that session's send lock. */ + CompletableFuture open(String session) { + CompletableFuture 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 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 waiter = waiters.get(session); + return waiter != null && waiter.complete(content); + } +} 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 4f08fda..3624539 100644 --- a/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java +++ b/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java @@ -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 view(Agent a) { Map m = new LinkedHashMap<>(); diff --git a/bridged/src/test/java/dev/ltms/bridged/herdr/FakeHerdr.java b/bridged/src/test/java/dev/ltms/bridged/herdr/FakeHerdr.java index b5f8e5e..045c9d2 100644 --- a/bridged/src/test/java/dev/ltms/bridged/herdr/FakeHerdr.java +++ b/bridged/src/test/java/dev/ltms/bridged/herdr/FakeHerdr.java @@ -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," 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 7517420..b1d373f 100644 --- a/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java @@ -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 postMessage(int port, String json) throws Exception { + return postJson(port, "/sessions/term_a/message", json); + } + + private HttpResponse 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 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 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 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 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 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 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 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 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();