diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 797d76a..498d09f 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -3,11 +3,14 @@ package dev.ltms.bridged; import dev.ltms.bridged.config.BridgedConfig; import dev.ltms.bridged.guard.SubscriptionGuard; import dev.ltms.bridged.herdr.AgentControl; +import dev.ltms.bridged.herdr.PaneLocator; 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.mcp.ConnectionIdentity; +import dev.ltms.bridged.mcp.LsofPeerPidLookup; import dev.ltms.bridged.msg.MessageService; import dev.ltms.bridged.msg.Rendezvous; import dev.ltms.bridged.rest.BridgedApp; @@ -60,7 +63,9 @@ public final class Bridged { MessageService messages = new MessageService(agents, injector, rendezvous); // MCP server face (CB-105): bridge_send/bridge_reply/bridge_status, mounted at /mcp. - BridgeMcp mcp = new BridgeMcp(messages, rendezvous); + // 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); 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/herdr/PaneLocator.java b/bridged/src/main/java/dev/ltms/bridged/herdr/PaneLocator.java new file mode 100644 index 0000000..46d7fb3 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/herdr/PaneLocator.java @@ -0,0 +1,59 @@ +package dev.ltms.bridged.herdr; + +import com.fasterxml.jackson.databind.JsonNode; + +import java.util.Map; + +/** + * Resolves which herdr pane a process belongs to — the herdr half of connection-based MCP + * identity (CB-105). Given the PID that opened an MCP connection, {@link #terminalForPid} finds + * the agent pane whose process tree contains it, so {@code bridged} can tell which worker + * is calling without the worker sending anything spoofable. + * + *

herdr owns the PID→pane truth: {@code pane.process_info} reports each pane's {@code shell_pid} + * and foreground process PIDs. This scans agent panes; a spawn-time {@code pid→terminal} cache is + * the obvious optimization once wired into {@code WorkerService}. + */ +public final class PaneLocator { + + private final HerdrClient herdr; + + public PaneLocator(HerdrClient herdr) { + this.herdr = herdr; + } + + /** + * The {@code terminal_id} of the agent pane whose process tree contains {@code pid}, or + * {@code null} if no agent pane owns it (e.g. the caller is the primary, or off-host). + */ + public String terminalForPid(long pid) { + if (pid <= 0) { + return null; + } + for (JsonNode pane : herdr.call("pane.list", Map.of()).path("panes")) { + String paneId = pane.path("pane_id").asText(null); + if (paneId != null && paneOwnsPid(paneId, pid)) { + return pane.path("terminal_id").asText(null); + } + } + return null; + } + + private boolean paneOwnsPid(String paneId, long pid) { + JsonNode info; + try { + info = herdr.call("pane.process_info", Map.of("pane_id", paneId)).path("process_info"); + } catch (HerdrException e) { + return false; // pane vanished mid-scan — just skip it + } + if (info.path("shell_pid").asLong(-1) == pid) { + return true; + } + for (JsonNode p : info.path("foreground_processes")) { + if (p.path("pid").asLong(-1) == pid) { + return true; + } + } + return false; + } +} 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 e630146..17398b7 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java @@ -3,10 +3,12 @@ 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.common.McpTransportContext; 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.McpSyncServerExchange; import io.modelcontextprotocol.server.transport.HttpServletStreamableServerTransportProvider; import io.modelcontextprotocol.spec.McpSchema; import jakarta.servlet.http.HttpServlet; @@ -29,14 +31,21 @@ public final class BridgeMcp { private static final long DEFAULT_TIMEOUT_MS = 25_000; private static final long MAX_TIMEOUT_MS = 120_000; + /** Transport-context key under which the extractor stashes the resolved caller identity. */ + static final String CALLER_TERMINAL = "callerTerminal"; + private final HttpServletStreamableServerTransportProvider transport; private final McpSyncServer server; - public BridgeMcp(MessageService messages, Rendezvous rendezvous) { + public BridgeMcp(MessageService messages, Rendezvous rendezvous, ConnectionIdentity identity) { McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get(); this.transport = HttpServletStreamableServerTransportProvider.builder() .jsonMapper(json) .mcpEndpoint("/mcp") + // Resolve the caller's worker identity from the connection (peer PID → herdr pane), + // so bridge_reply needs no spoofable session argument. + .contextExtractor(req -> McpTransportContext.create(Map.of( + CALLER_TERMINAL, orEmpty(identity.callerTerminal(req.getRemoteAddr(), req.getRemotePort()))))) .build(); this.server = McpServer.sync(transport) .serverInfo("bridge", "0.1.0") @@ -45,15 +54,25 @@ public final class BridgeMcp { 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")); - }) + // bridge_reply's identity is the CONNECTION, never an argument. + .toolCall(replyTool(), (exchange, req) -> + reply(rendezvous, callerTerminal(exchange), str(req.arguments(), "content"))) .toolCall(statusTool(), (_, req) -> status(messages, str(req.arguments(), "sessionId"))) .build(); } + /** The worker identity resolved from this call's connection, or {@code null} if the primary. */ + private static String callerTerminal(McpSyncServerExchange exchange) { + Object v = exchange.transportContext().get(CALLER_TERMINAL); + String s = v == null ? null : v.toString(); + return (s == null || s.isBlank()) ? null : s; + } + + private static String orEmpty(String s) { + return s == null ? "" : s; + } + /** The Streamable-HTTP servlet to mount at {@code /mcp} on the daemon's Jetty. */ public HttpServlet servlet() { return transport; @@ -84,14 +103,22 @@ public final class BridgeMcp { } } - /** {@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"); + /** + * {@code bridge_reply}: the worker returns its structured answer, resolving the awaiting send. + * {@code callerTerminal} is resolved from the connection (never an argument); a {@code null} + * means the caller is not a known worker (e.g. the primary called it by mistake). + */ + static McpSchema.CallToolResult reply(Rendezvous rendezvous, String callerTerminal, String content) { + if (callerTerminal == null) { + return error("bridge_reply is for workers only — could not identify the calling worker " + + "from the connection"); } - return rendezvous.resolve(sessionId, content) + if (content == null) { + return error("content is required"); + } + return rendezvous.resolve(callerTerminal, content) ? text("delivered") - : error("no send is awaiting a reply for session " + sessionId); + : error("no send is awaiting a reply for this worker"); } /** {@code bridge_status}: the live lifecycle status of a worker session. */ @@ -120,13 +147,13 @@ public final class BridgeMcp { } private static McpSchema.Tool replyTool() { + // No session/target arg — the worker's identity is resolved from the connection. 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"))); + List.of("content"))); } private static McpSchema.Tool statusTool() { diff --git a/bridged/src/main/java/dev/ltms/bridged/mcp/ConnectionIdentity.java b/bridged/src/main/java/dev/ltms/bridged/mcp/ConnectionIdentity.java new file mode 100644 index 0000000..fca4705 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/ConnectionIdentity.java @@ -0,0 +1,40 @@ +package dev.ltms.bridged.mcp; + +import dev.ltms.bridged.herdr.PaneLocator; + +/** + * Resolves who is calling an MCP tool from the connection alone — the anti-spoofing + * identity model of the MCP contract. It ties the connection's loopback peer PID (from the OS) + * to a herdr agent pane (from herdr), yielding the caller's worker {@code terminal_id}. A caller + * that maps to no worker pane — the primary, or an off-host client — resolves to {@code null}. + * + *

Both sources are authoritative and unforgeable: the OS reports the real connecting PID, and + * herdr owns the PID→pane mapping. A worker cannot claim to be another worker, nor the primary. + * Single-host only (the herd shares the {@code bridged} host); the token path is the split-host + * fallback. + */ +public final class ConnectionIdentity { + + private final PaneLocator panes; + private final PeerPidLookup pids; + + public ConnectionIdentity(PaneLocator panes, PeerPidLookup pids) { + this.panes = panes; + this.pids = pids; + } + + /** + * The calling worker's {@code terminal_id}, or {@code null} if the caller is not a known + * on-host worker (treat as the primary). + */ + public String callerTerminal(String remoteAddr, int remotePort) { + if (!isLoopback(remoteAddr)) { + return null; // only same-host callers can be workers + } + return panes.terminalForPid(pids.pidForLocalPort(remotePort)); + } + + private static boolean isLoopback(String addr) { + return "127.0.0.1".equals(addr) || "::1".equals(addr) || "0:0:0:0:0:0:0:1".equals(addr); + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/mcp/LsofPeerPidLookup.java b/bridged/src/main/java/dev/ltms/bridged/mcp/LsofPeerPidLookup.java new file mode 100644 index 0000000..fa47546 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/LsofPeerPidLookup.java @@ -0,0 +1,58 @@ +package dev.ltms.bridged.mcp; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.io.BufferedReader; +import java.io.InputStreamReader; +import java.nio.charset.StandardCharsets; +import java.util.concurrent.TimeUnit; + +/** + * {@link PeerPidLookup} via {@code lsof} (present on macOS and Linux). For a loopback TCP source + * {@code port}, both the client and this daemon appear on that port — so we exclude our own PID + * and take the other end, which is the calling process. + */ +public final class LsofPeerPidLookup implements PeerPidLookup { + + private static final Logger log = LoggerFactory.getLogger(LsofPeerPidLookup.class); + + private final long selfPid = ProcessHandle.current().pid(); + + @Override + public long pidForLocalPort(int port) { + try { + Process p = new ProcessBuilder("lsof", "-nP", "-FpP", "-iTCP:" + port) + .redirectErrorStream(true).start(); + long found = -1; + try (BufferedReader r = new BufferedReader( + new InputStreamReader(p.getInputStream(), StandardCharsets.UTF_8))) { + long current = -1; + String line; + // -Fp emits records: a 'p' line, then the ports/files under that pid. + while ((line = r.readLine()) != null) { + if (line.startsWith("p")) { + current = parse(line.substring(1)); + } else if (current > 0 && current != selfPid) { + found = current; // first process on this port that isn't us = the client + } + } + } + if (!p.waitFor(2, TimeUnit.SECONDS)) { + p.destroyForcibly(); + } + return found; + } catch (Exception e) { + log.debug("lsof peer-pid lookup for port {} failed: {}", port, e.getMessage()); + return -1; + } + } + + private static long parse(String s) { + try { + return Long.parseLong(s.trim()); + } catch (NumberFormatException e) { + return -1; + } + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/mcp/PeerPidLookup.java b/bridged/src/main/java/dev/ltms/bridged/mcp/PeerPidLookup.java new file mode 100644 index 0000000..105713a --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/PeerPidLookup.java @@ -0,0 +1,13 @@ +package dev.ltms.bridged.mcp; + +/** + * Resolves the OS PID that owns a loopback TCP source port — the OS half of connection-based MCP + * identity. Java exposes no peer PID for a TCP socket, so this shells out. Injectable so + * {@link ConnectionIdentity} is testable without a real connection. + */ +@FunctionalInterface +public interface PeerPidLookup { + + /** The PID whose socket has local (source) {@code port} on loopback, or {@code -1} if unknown. */ + long pidForLocalPort(int port); +} 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 045c9d2..103468c 100644 --- a/bridged/src/test/java/dev/ltms/bridged/herdr/FakeHerdr.java +++ b/bridged/src/test/java/dev/ltms/bridged/herdr/FakeHerdr.java @@ -16,6 +16,9 @@ public final class FakeHerdr implements HerdrClient { public record Call(String method, Object params) { } + /** The foreground PID of the one agent pane (term_a) in the canned {@code pane.process_info}. */ + public static final long WORKER_PID = 4242; + private final ObjectMapper mapper = new ObjectMapper(); public final List calls = new ArrayList<>(); private boolean healthy = true; @@ -142,6 +145,21 @@ public final class FakeHerdr implements HerdrClient { case "pane.get" -> mapper.readTree(""" {"type":"pane_info","pane":{"pane_id":"w9:pW","workspace_id":"w9", "tab_id":"w9:t2","agent_status":"idle"}}"""); + case "pane.list" -> mapper.readTree(""" + {"type":"pane_list","panes":[ + {"pane_id":"w2:p7","terminal_id":"term_a","workspace_id":"w2","tab_id":"w2:t7","agent":"claude"}, + {"pane_id":"w2:p9","terminal_id":"term_shell","workspace_id":"w2","tab_id":"w2:t8"}]}"""); + case "pane.process_info" -> { + Object paneId = params instanceof java.util.Map m ? m.get("pane_id") : null; + yield "w2:p7".equals(paneId) + ? mapper.readTree((""" + {"type":"pane_process_info","process_info":{"pane_id":"w2:p7","shell_pid":%d, + "foreground_processes":[{"pid":%d,"name":"node","argv0":"claude"}]}}""") + .formatted(WORKER_PID, WORKER_PID)) + : mapper.readTree(""" + {"type":"pane_process_info","process_info":{"pane_id":"w2:p9","shell_pid":9001, + "foreground_processes":[]}}"""); + } case "pane.close" -> { if (paneCloseErrorCode != null) { throw new HerdrException("herdr error [" + paneCloseErrorCode + "]: pane.close failed", diff --git a/bridged/src/test/java/dev/ltms/bridged/herdr/PaneLocatorContractTest.java b/bridged/src/test/java/dev/ltms/bridged/herdr/PaneLocatorContractTest.java new file mode 100644 index 0000000..1ef0451 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/herdr/PaneLocatorContractTest.java @@ -0,0 +1,43 @@ +package dev.ltms.bridged.herdr; + +import com.fasterxml.jackson.databind.JsonNode; +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.Test; + +import java.nio.file.Files; +import java.util.List; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.*; +import static org.junit.jupiter.api.Assumptions.assumeTrue; + +/** + * Contract test for the herdr half of connection-based identity against a REAL herdr: spawn a + * harmless probe, read its actual {@code shell_pid} from {@code pane.process_info}, and confirm + * {@link PaneLocator} resolves that PID back to the probe's own {@code terminal_id}. + * + *

Tagged {@code contract}; run with {@code mvn test -Pcontract}. + */ +@Tag("contract") +class PaneLocatorContractTest { + + @Test + void resolvesTheTerminalOwningARealProcessPid() throws Exception { + assumeTrue(Files.exists(UnixSocketHerdrClient.defaultSocketPath()), "no herdr socket — skipping"); + try (UnixSocketHerdrClient herdr = UnixSocketHerdrClient.connect()) { + AgentControl agents = new AgentControl(herdr); + Agent probe = agents.start("__pidprobe__", List.of("bash", "-c", "sleep 20"), Map.of()); + try { + JsonNode info = herdr.call("pane.process_info", Map.of("pane_id", probe.paneId())) + .path("process_info"); + long shellPid = info.path("shell_pid").asLong(-1); + assertTrue(shellPid > 0, "probe pane should report a shell pid"); + + assertEquals(probe.terminalId(), new PaneLocator(herdr).terminalForPid(shellPid), + "a real PID must resolve back to its own pane's terminal_id"); + } finally { + agents.close(probe.paneId()); + } + } + } +} diff --git a/bridged/src/test/java/dev/ltms/bridged/herdr/PaneLocatorTest.java b/bridged/src/test/java/dev/ltms/bridged/herdr/PaneLocatorTest.java new file mode 100644 index 0000000..54262d4 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/herdr/PaneLocatorTest.java @@ -0,0 +1,27 @@ +package dev.ltms.bridged.herdr; + +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.*; + +/** Unit tests for PID → pane resolution (the herdr half of connection-based MCP identity). */ +class PaneLocatorTest { + + private final PaneLocator loc = new PaneLocator(new FakeHerdr()); + + @Test + void resolvesTerminalForAForegroundPid() { + assertEquals("term_a", loc.terminalForPid(FakeHerdr.WORKER_PID)); + } + + @Test + void nullForAPidInNoPane() { + assertNull(loc.terminalForPid(999_999)); + } + + @Test + void nullForNonPositivePid() { + assertNull(loc.terminalForPid(0)); + assertNull(loc.terminalForPid(-1)); + } +} diff --git a/bridged/src/test/java/dev/ltms/bridged/mcp/ConnectionIdentityTest.java b/bridged/src/test/java/dev/ltms/bridged/mcp/ConnectionIdentityTest.java new file mode 100644 index 0000000..1048524 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/mcp/ConnectionIdentityTest.java @@ -0,0 +1,34 @@ +package dev.ltms.bridged.mcp; + +import dev.ltms.bridged.herdr.FakeHerdr; +import dev.ltms.bridged.herdr.PaneLocator; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.*; + +/** Connection → caller-identity resolution, with the OS peer-PID lookup faked. */ +class ConnectionIdentityTest { + + private final FakeHerdr herdr = new FakeHerdr(); + + private ConnectionIdentity with(PeerPidLookup pids) { + return new ConnectionIdentity(new PaneLocator(herdr), pids); + } + + @Test + void resolvesWorkerFromLoopbackPeerPid() { + assertEquals("term_a", with(_ -> FakeHerdr.WORKER_PID).callerTerminal("127.0.0.1", 55555)); + } + + @Test + void nullForOffHostCaller() { + // A non-loopback peer can't be an on-host worker → treat as primary/unknown. + assertNull(with(_ -> FakeHerdr.WORKER_PID).callerTerminal("10.0.0.9", 55555)); + } + + @Test + void nullWhenPidOwnsNoPane() { + // e.g. the primary — its PID maps to no worker pane. + assertNull(with(_ -> 999_999).callerTerminal("127.0.0.1", 55555)); + } +}