CB-108: fleet-management MCP tools (bridge_spawn / bridge_list / bridge_stop)

The primary could delegate to a worker but not create or reap one over MCP —
spawning was a raw REST POST /workers. BridgeMcp now adapts WorkerService so a
worker's whole lifecycle runs through MCP: bridge_spawn returns the new worker's
sessionId (for bridge_send) and paneId (for bridge_stop); bridge_list projects the
tracked workers; bridge_stop tears one down. The subscription boundary stays
enforced inside WorkerService (bridge_spawn surfaces a guard breach as a tool error
without touching herdr). Tools are thin static adapters, unit-tested by parity.
This commit is contained in:
Dai Ha
2026-07-15 14:22:50 +02:00
parent 8ed2370fbe
commit a97c287aee
3 changed files with 146 additions and 2 deletions
@@ -69,7 +69,7 @@ public final class Bridged {
// 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.
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), new LsofPeerPidLookup());
BridgeMcp mcp = new BridgeMcp(messages, rendezvous, identity);
BridgeMcp mcp = new BridgeMcp(messages, rendezvous, workers, identity);
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());
@@ -1,8 +1,11 @@
package dev.ltms.bridged.mcp;
import dev.ltms.bridged.guard.GuardException;
import dev.ltms.bridged.herdr.Agent;
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.modelcontextprotocol.common.McpTransportContext;
import io.modelcontextprotocol.json.McpJsonMapper;
import io.modelcontextprotocol.json.jackson3.JacksonMcpJsonMapperSupplier;
@@ -11,8 +14,10 @@ import io.modelcontextprotocol.server.McpSyncServer;
import io.modelcontextprotocol.server.McpSyncServerExchange;
import io.modelcontextprotocol.server.transport.HttpServletStreamableServerTransportProvider;
import io.modelcontextprotocol.spec.McpSchema;
import com.fasterxml.jackson.databind.ObjectMapper;
import jakarta.servlet.http.HttpServlet;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
@@ -22,6 +27,10 @@ import java.util.Map;
* 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>Beyond delegation the primary also manages the fleet here (CB-108): {@code bridge_spawn} /
* {@code bridge_list} / {@code bridge_stop} adapt {@link WorkerService} so a worker's whole lifecycle
* is driven through MCP, with the subscription boundary still enforced inside {@code WorkerService}.
*
* <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.
@@ -30,6 +39,7 @@ public final class BridgeMcp {
private static final long DEFAULT_TIMEOUT_MS = 25_000;
private static final long MAX_TIMEOUT_MS = 120_000;
private static final ObjectMapper MAPPER = new ObjectMapper(); // worker-view JSON projections
/** Transport-context key under which the extractor stashes the resolved caller identity. */
static final String CALLER_TERMINAL = "callerTerminal";
@@ -37,7 +47,8 @@ public final class BridgeMcp {
private final HttpServletStreamableServerTransportProvider transport;
private final McpSyncServer server;
public BridgeMcp(MessageService messages, Rendezvous rendezvous, ConnectionIdentity identity) {
public BridgeMcp(MessageService messages, Rendezvous rendezvous, WorkerService workers,
ConnectionIdentity identity) {
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
this.transport = HttpServletStreamableServerTransportProvider.builder()
.jsonMapper(json)
@@ -64,6 +75,10 @@ public final class BridgeMcp {
status(messages, str(req.arguments(), "sessionId")))
.toolCall(pollTool(), (_, req) ->
poll(messages, str(req.arguments(), "ticket")))
// Fleet management (CB-108): spawn/list/stop over WorkerService.
.toolCall(spawnTool(), (_, _) -> spawn(workers))
.toolCall(listTool(), (_, _) -> listWorkers(workers))
.toolCall(stopTool(), (_, req) -> stop(workers, str(req.arguments(), "paneId")))
.build();
}
@@ -174,6 +189,59 @@ public final class BridgeMcp {
}
}
// --- fleet management logic (CB-108) -------------------------------------------------------
/** {@code bridge_spawn}: launch a guard-checked worker and return its session id + pane id. */
static McpSchema.CallToolResult spawn(WorkerService workers) {
try {
return text(json(workerView(workers.spawn())));
} catch (GuardException e) {
return error("subscription boundary: " + e.getMessage());
} catch (HerdrException e) {
return error("herdr error spawning worker: " + e.getMessage());
}
}
/** {@code bridge_list}: every worker herdr tracks (session id, pane, status). */
static McpSchema.CallToolResult listWorkers(WorkerService workers) {
try {
List<Map<String, Object>> out = workers.list().stream().map(BridgeMcp::workerView).toList();
return text(json(Map.of("workers", out)));
} catch (HerdrException e) {
return error("herdr error listing workers: " + e.getMessage());
}
}
/** {@code bridge_stop}: tear a worker down by its pane id. */
static McpSchema.CallToolResult stop(WorkerService workers, String paneId) {
if (isBlank(paneId)) {
return error("paneId is required");
}
try {
workers.stop(paneId);
return text("stopped " + paneId);
} catch (HerdrException e) {
return error("herdr error stopping " + paneId + ": " + e.getMessage());
}
}
/** Projection a delegator can act on: sessionId (for bridge_send) + paneId (for bridge_stop). */
private static Map<String, Object> workerView(Agent a) {
Map<String, Object> m = new LinkedHashMap<>();
m.put("sessionId", a.terminalId()); // the id bridge_send / bridge_status take
m.put("paneId", a.paneId());
m.put("status", a.status().name().toLowerCase());
return m;
}
private static String json(Object o) {
try {
return MAPPER.writeValueAsString(o);
} catch (Exception e) {
return String.valueOf(o);
}
}
// --- tool schemas --------------------------------------------------------------------------
private static McpSchema.Tool sendTool() {
@@ -199,6 +267,27 @@ public final class BridgeMcp {
List.of("ticket")));
}
private static McpSchema.Tool spawnTool() {
return tool("bridge_spawn",
"Spawn a new off-subscription worker session for the configured profile. Returns its "
+ "sessionId (use with bridge_send) and paneId (use with bridge_stop).",
objectSchema(Map.of(), List.of()));
}
private static McpSchema.Tool listTool() {
return tool("bridge_list",
"List the worker sessions the bridge tracks — each with its sessionId, paneId, and status.",
objectSchema(Map.of(), List.of()));
}
private static McpSchema.Tool stopTool() {
return tool("bridge_stop",
"Tear down a worker session by its paneId (from bridge_spawn or bridge_list).",
objectSchema(Map.of(
"paneId", stringProp("The worker's paneId to stop")),
List.of("paneId")));
}
private static McpSchema.Tool replyTool() {
// No session/target arg — the worker's identity is resolved from the connection.
return tool("bridge_reply",
@@ -1,13 +1,18 @@
package dev.ltms.bridged.mcp;
import dev.ltms.bridged.config.BridgedConfig;
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.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.worker.WorkerService;
import io.modelcontextprotocol.spec.McpSchema;
import org.junit.jupiter.api.Test;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
@@ -29,6 +34,14 @@ class BridgeMcpTest {
return ((McpSchema.TextContent) r.content().getFirst()).text();
}
private static WorkerService workerService(FakeHerdr h, String baseUrl, Set<String> allow) {
BridgedConfig.Worker cfg = new BridgedConfig.Worker(
"ltms-local", baseUrl, "coder", null, "BRIDGED_WORKER_TOKEN", null,
"tab", "bridged-workers", "worker: {profile} #{n}", null);
return new WorkerService(new AgentControl(h), new WorkspaceControl(h),
new SubscriptionGuard(allow), cfg, _ -> "tok");
}
@Test
void sendThenReplyRoundTrips() throws Exception {
// bridge_send blocks; bridge_reply resolves it with the worker's structured answer.
@@ -106,6 +119,48 @@ class BridgeMcpTest {
assertTrue(textOf(res).contains("no send is awaiting"));
}
@Test
void spawnReturnsTheNewWorkersSessionAndPane() {
FakeHerdr h = new FakeHerdr();
McpSchema.CallToolResult res = BridgeMcp.spawn(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
assertNotEquals(Boolean.TRUE, res.isError());
String out = textOf(res);
assertTrue(out.contains("\"sessionId\":\"term_new\""), out);
assertTrue(out.contains("\"paneId\":\"w9:pW\""), out);
}
@Test
void spawnRejectsAnOffAllowlistProfileWithoutTouchingHerdr() {
FakeHerdr h = new FakeHerdr();
McpSchema.CallToolResult res = BridgeMcp.spawn(workerService(h, "https://api.anthropic.com", Set.of("gx00.gw")));
assertTrue(res.isError());
assertTrue(textOf(res).contains("subscription boundary"));
assertFalse(h.called("agent.start"), "the guard must block before any spawn");
}
@Test
void listReportsTrackedWorkers() {
FakeHerdr h = new FakeHerdr();
McpSchema.CallToolResult res = BridgeMcp.listWorkers(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
assertNotEquals(Boolean.TRUE, res.isError());
assertTrue(textOf(res).contains("\"sessionId\":\"term_a\""), textOf(res));
}
@Test
void stopTearsDownAWorkerByPane() {
FakeHerdr h = new FakeHerdr();
McpSchema.CallToolResult res = BridgeMcp.stop(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), "w9:pW");
assertNotEquals(Boolean.TRUE, res.isError());
assertEquals("stopped w9:pW", textOf(res));
assertTrue(h.called("pane.close"));
}
@Test
void stopRequiresAPaneId() {
FakeHerdr h = new FakeHerdr();
assertTrue(BridgeMcp.stop(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), " ").isError());
}
@Test
void statusReportsLiveAgentStatus() {
FakeHerdr blocked = new FakeHerdr().agentStatus("blocked");