diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 1f10f10..fc1131f 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -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()); diff --git a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java index db19a7b..fb0548e 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java @@ -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}. * + *

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

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. @@ -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> 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 workerView(Agent a) { + Map 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", diff --git a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java index ea611e6..b216db1 100644 --- a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java @@ -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 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");