CB-105: MCP server (bridge_send/bridge_reply/bridge_status) over the REST core
Streamable-HTTP MCP server (io.modelcontextprotocol.sdk:mcp 2.0.0) mounted on the daemon's Jetty at /mcp, exposing three tools as thin adapters over MessageService/Rendezvous — the primary calls bridge_send/bridge_status, the worker calls bridge_reply. Tool logic in unit-testable static methods (parity tests); SDK owns the wire protocol. Resolves the Jackson 2/3 split by pinning jackson-annotations 3.0-rc5 (works for both our Jackson 2.19 and the SDK's Jackson 3). Live-verified: initialize handshake + tools/list return all three tools. Documents the residual Jackson-3 CVE (loopback, trusted clients). 52 tests green, IDE-clean.
This commit is contained in:
@@ -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);
|
||||
|
||||
@@ -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 <em>thin adapters</em>
|
||||
* 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}.
|
||||
*
|
||||
* <p>The tool <em>logic</em> 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<String, Object> a = req.arguments();
|
||||
return send(messages, str(a, "sessionId"), str(a, "content"), timeoutMs(a));
|
||||
})
|
||||
.toolCall(replyTool(), (_, req) -> {
|
||||
Map<String, Object> 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<String, Object> inputSchema) {
|
||||
return McpSchema.Tool.builder(name).description(description).inputSchema(inputSchema).build();
|
||||
}
|
||||
|
||||
private static Map<String, Object> objectSchema(Map<String, Object> properties, List<String> required) {
|
||||
return Map.of("type", "object", "properties", properties, "required", required);
|
||||
}
|
||||
|
||||
private static Map<String, Object> 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<String, Object> args, String key) {
|
||||
Object v = args.get(key);
|
||||
return v == null ? null : v.toString();
|
||||
}
|
||||
|
||||
private static Long timeoutMs(Map<String, Object> 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();
|
||||
}
|
||||
}
|
||||
@@ -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);
|
||||
|
||||
Reference in New Issue
Block a user