org.slf4j
diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java
index c323297..797d76a 100644
--- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java
+++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java
@@ -7,6 +7,7 @@ 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.mcp.BridgeMcp;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.rest.BridgedApp;
@@ -57,7 +58,11 @@ public final class Bridged {
Rendezvous rendezvous = new Rendezvous();
MessageService messages = new MessageService(agents, injector, rendezvous);
- Javalin app = new BridgedApp(herdr, workers, messages, rendezvous).build();
+
+ // MCP server face (CB-105): bridge_send/bridge_reply/bridge_status, mounted at /mcp.
+ BridgeMcp mcp = new BridgeMcp(messages, rendezvous);
+ Runtime.getRuntime().addShutdownHook(new Thread(mcp::close));
+ Javalin app = new BridgedApp(herdr, workers, messages, rendezvous, mcp.servlet()).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/mcp/BridgeMcp.java b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java
new file mode 100644
index 0000000..e630146
--- /dev/null
+++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java
@@ -0,0 +1,181 @@
+package dev.ltms.bridged.mcp;
+
+import dev.ltms.bridged.herdr.HerdrException;
+import dev.ltms.bridged.msg.MessageService;
+import dev.ltms.bridged.msg.Rendezvous;
+import io.modelcontextprotocol.json.McpJsonMapper;
+import io.modelcontextprotocol.json.jackson3.JacksonMcpJsonMapperSupplier;
+import io.modelcontextprotocol.server.McpServer;
+import io.modelcontextprotocol.server.McpSyncServer;
+import io.modelcontextprotocol.server.transport.HttpServletStreamableServerTransportProvider;
+import io.modelcontextprotocol.spec.McpSchema;
+import jakarta.servlet.http.HttpServlet;
+
+import java.util.List;
+import java.util.Map;
+
+/**
+ * The MCP SERVER face (CB-105): a Streamable-HTTP MCP server whose tools are thin adapters
+ * over the same {@link MessageService}/{@link Rendezvous} the REST routes use — so the two are
+ * validated by parity, not by re-implementing behaviour. The primary Opus calls {@code bridge_send}
+ * / {@code bridge_status}; the worker calls {@code bridge_reply}.
+ *
+ * The tool logic lives in package-private static methods returning a
+ * {@link McpSchema.CallToolResult}, so it is unit-testable without standing up the HTTP transport;
+ * the SDK owns the wire protocol. Mount {@link #servlet()} at {@code /mcp} on the daemon's Jetty.
+ */
+public final class BridgeMcp {
+
+ private static final long DEFAULT_TIMEOUT_MS = 25_000;
+ private static final long MAX_TIMEOUT_MS = 120_000;
+
+ private final HttpServletStreamableServerTransportProvider transport;
+ private final McpSyncServer server;
+
+ public BridgeMcp(MessageService messages, Rendezvous rendezvous) {
+ McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
+ this.transport = HttpServletStreamableServerTransportProvider.builder()
+ .jsonMapper(json)
+ .mcpEndpoint("/mcp")
+ .build();
+ this.server = McpServer.sync(transport)
+ .serverInfo("bridge", "0.1.0")
+ .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));
+ })
+ .toolCall(replyTool(), (_, req) -> {
+ Map a = req.arguments();
+ return reply(rendezvous, str(a, "sessionId"), str(a, "content"));
+ })
+ .toolCall(statusTool(), (_, req) ->
+ status(messages, str(req.arguments(), "sessionId")))
+ .build();
+ }
+
+ /** The Streamable-HTTP servlet to mount at {@code /mcp} on the daemon's Jetty. */
+ public HttpServlet servlet() {
+ return transport;
+ }
+
+ /** Graceful shutdown of the MCP server. */
+ public void close() {
+ server.closeGracefully();
+ }
+
+ // --- tool logic (thin adapters over the services; unit-testable) ---------------------------
+
+ /** {@code bridge_send}: delegate {@code content} to a worker session and block for its reply. */
+ static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content, Long timeoutMs) {
+ if (isBlank(sessionId) || isBlank(content)) {
+ return error("sessionId and content are required");
+ }
+ long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
+ try {
+ MessageService.Reply r = messages.send(sessionId, content, timeout);
+ if (r.completed()) {
+ return text(r.text());
+ }
+ return text("[no reply within " + timeout + "ms — worker "
+ + r.outcome().name().toLowerCase().replace("timed_out_", "") + "; retry or poll status]");
+ } catch (HerdrException e) {
+ return error("herdr error contacting session " + sessionId + ": " + e.getMessage());
+ }
+ }
+
+ /** {@code bridge_reply}: the worker returns its structured answer, resolving the awaiting send. */
+ static McpSchema.CallToolResult reply(Rendezvous rendezvous, String sessionId, String content) {
+ if (isBlank(sessionId) || content == null) {
+ return error("sessionId and content are required");
+ }
+ return rendezvous.resolve(sessionId, content)
+ ? text("delivered")
+ : error("no send is awaiting a reply for session " + sessionId);
+ }
+
+ /** {@code bridge_status}: the live lifecycle status of a worker session. */
+ static McpSchema.CallToolResult status(MessageService messages, String sessionId) {
+ if (isBlank(sessionId)) {
+ return error("sessionId is required");
+ }
+ try {
+ return text(messages.status(sessionId).name().toLowerCase());
+ } catch (HerdrException e) {
+ return error("herdr error for session " + sessionId + ": " + e.getMessage());
+ }
+ }
+
+ // --- tool schemas --------------------------------------------------------------------------
+
+ 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.",
+ 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")),
+ List.of("sessionId", "content")));
+ }
+
+ private static McpSchema.Tool replyTool() {
+ return tool("bridge_reply",
+ "Return your structured answer for the task you were delegated, "
+ + "resolving the caller's blocked bridge_send.",
+ objectSchema(Map.of(
+ "sessionId", stringProp("Your own worker session id"),
+ "content", stringProp("Your reply/answer")),
+ List.of("sessionId", "content")));
+ }
+
+ private static McpSchema.Tool statusTool() {
+ return tool("bridge_status",
+ "Get the live lifecycle status (idle/working/blocked/unknown) of a worker session.",
+ objectSchema(Map.of(
+ "sessionId", stringProp("The worker session id to query")),
+ List.of("sessionId")));
+ }
+
+ // --- small helpers -------------------------------------------------------------------------
+
+ // The SDK 2.0.0 deprecates its own Tool builders without a stable replacement — isolate it here.
+ @SuppressWarnings("deprecation")
+ private static McpSchema.Tool tool(String name, String description, Map inputSchema) {
+ return McpSchema.Tool.builder(name).description(description).inputSchema(inputSchema).build();
+ }
+
+ private static Map objectSchema(Map properties, List required) {
+ return Map.of("type", "object", "properties", properties, "required", required);
+ }
+
+ private static Map stringProp(String description) {
+ return Map.of("type", "string", "description", description);
+ }
+
+ private static McpSchema.CallToolResult text(String s) {
+ return McpSchema.CallToolResult.builder().addTextContent(s == null ? "" : s).build();
+ }
+
+ private static McpSchema.CallToolResult error(String s) {
+ return McpSchema.CallToolResult.builder().addTextContent(s).isError(true).build();
+ }
+
+ private static String str(Map args, String key) {
+ Object v = args.get(key);
+ return v == null ? null : v.toString();
+ }
+
+ private static Long timeoutMs(Map args) {
+ Object v = args.get("timeoutMs");
+ return v instanceof Number n ? n.longValue() : null;
+ }
+
+ private static long clamp(long ms) {
+ return Math.clamp(ms, 1, MAX_TIMEOUT_MS);
+ }
+
+ private static boolean isBlank(String s) {
+ return s == null || s.isBlank();
+ }
+}
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 3624539..5360beb 100644
--- a/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java
+++ b/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java
@@ -11,6 +11,8 @@ import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.worker.WorkerService;
import io.javalin.Javalin;
import io.javalin.http.Context;
+import jakarta.servlet.http.HttpServlet;
+import org.eclipse.jetty.servlet.ServletHolder;
import java.util.ArrayList;
import java.util.LinkedHashMap;
@@ -36,18 +38,28 @@ public final class BridgedApp {
private final WorkerService workers;
private final MessageService messages;
private final Rendezvous rendezvous;
+ private final HttpServlet mcpServlet; // MCP Streamable-HTTP endpoint, mounted at /mcp (nullable)
private final ObjectMapper mapper = new ObjectMapper();
- public BridgedApp(HerdrClient herdr, WorkerService workers, MessageService messages, Rendezvous rendezvous) {
+ public BridgedApp(HerdrClient herdr, WorkerService workers, MessageService messages,
+ Rendezvous rendezvous, HttpServlet mcpServlet) {
this.herdr = herdr;
this.workers = workers;
this.messages = messages;
this.rendezvous = rendezvous;
+ this.mcpServlet = mcpServlet;
}
/** Wire routes onto a fresh, unstarted Javalin instance. Caller starts it. */
public Javalin build() {
- Javalin app = Javalin.create(cfg -> cfg.showJavalinBanner = false);
+ Javalin app = Javalin.create(cfg -> {
+ cfg.showJavalinBanner = false;
+ if (mcpServlet != null) {
+ // The MCP server shares the daemon's port; Jetty routes /mcp to its servlet.
+ cfg.jetty.modifyServletContextHandler(h ->
+ h.addServlet(new ServletHolder(mcpServlet), "/mcp"));
+ }
+ });
app.get("/healthz", this::healthz);
app.get("/sessions", this::sessions);
app.get("/agents", this::agents);
diff --git a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java
new file mode 100644
index 0000000..6eee470
--- /dev/null
+++ b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java
@@ -0,0 +1,81 @@
+package dev.ltms.bridged.mcp;
+
+import dev.ltms.bridged.herdr.AgentControl;
+import dev.ltms.bridged.herdr.FakeHerdr;
+import dev.ltms.bridged.inject.Injector;
+import dev.ltms.bridged.msg.MessageService;
+import dev.ltms.bridged.msg.Rendezvous;
+import io.modelcontextprotocol.spec.McpSchema;
+import org.junit.jupiter.api.Test;
+
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.TimeUnit;
+
+import static org.junit.jupiter.api.Assertions.*;
+
+/**
+ * Parity tests for the MCP tool adapters — they must produce the same outcomes as the REST routes,
+ * since both drive the same {@link MessageService}/{@link Rendezvous}. The MCP wire protocol itself
+ * is the SDK's concern; here we test the thin adapter logic directly.
+ */
+class BridgeMcpTest {
+
+ private final FakeHerdr herdr = new FakeHerdr();
+ private final AgentControl agents = new AgentControl(herdr);
+ private final Rendezvous rendezvous = new Rendezvous();
+ private final MessageService messages = new MessageService(agents, new Injector(agents), rendezvous);
+
+ private static String textOf(McpSchema.CallToolResult r) {
+ return ((McpSchema.TextContent) r.content().getFirst()).text();
+ }
+
+ @Test
+ void sendThenReplyRoundTrips() throws Exception {
+ // bridge_send blocks; bridge_reply resolves it with the worker's structured answer.
+ CompletableFuture send = CompletableFuture.supplyAsync(
+ () -> BridgeMcp.send(messages, "term_a", "review this", 4000L));
+
+ McpSchema.CallToolResult reply = BridgeMcp.reply(rendezvous, "term_a", "LGTM");
+ long deadline = System.currentTimeMillis() + 3000;
+ while (Boolean.TRUE.equals(reply.isError()) && System.currentTimeMillis() < deadline) {
+ //noinspection BusyWait
+ Thread.sleep(10);
+ reply = BridgeMcp.reply(rendezvous, "term_a", "LGTM");
+ }
+ assertEquals("delivered", textOf(reply));
+
+ McpSchema.CallToolResult res = send.get(6, TimeUnit.SECONDS);
+ assertNotEquals(Boolean.TRUE, res.isError());
+ assertEquals("LGTM", textOf(res));
+ }
+
+ @Test
+ void sendTimesOutWithAWorkingNote() {
+ McpSchema.CallToolResult res = BridgeMcp.send(messages, "term_a", "hi", 120L);
+ assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
+ assertTrue(textOf(res).contains("no reply"), "got: " + textOf(res));
+ }
+
+ @Test
+ void sendRejectsMissingArgs() {
+ assertTrue(BridgeMcp.send(messages, null, "hi", null).isError());
+ assertTrue(BridgeMcp.send(messages, "term_a", " ", null).isError());
+ }
+
+ @Test
+ void replyWithNoPendingSendIsAnError() {
+ McpSchema.CallToolResult res = BridgeMcp.reply(rendezvous, "term_a", "orphan");
+ assertTrue(res.isError());
+ assertTrue(textOf(res).contains("no send is awaiting"));
+ }
+
+ @Test
+ void statusReportsLiveAgentStatus() {
+ FakeHerdr blocked = new FakeHerdr().agentStatus("blocked");
+ AgentControl blockedAgents = new AgentControl(blocked);
+ McpSchema.CallToolResult res = BridgeMcp.status(
+ new MessageService(blockedAgents, new Injector(blockedAgents), rendezvous), "term_a");
+ assertNotEquals(Boolean.TRUE, res.isError());
+ assertEquals("blocked", textOf(res));
+ }
+}
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 b1d373f..01bbd27 100644
--- a/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java
+++ b/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java
@@ -62,7 +62,7 @@ class BridgedAppTest {
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);
+ app = new BridgedApp(herdr, workers, messages, rendezvous, null).build().start("127.0.0.1", 0);
return app.port();
}