CB-105: connection-based MCP caller identity (peer PID → herdr pane)
Resolve who is calling an MCP tool from the connection, not a spoofable argument (per docs/MCP-Contract.md). ConnectionIdentity ties the loopback peer PID (LsofPeerPidLookup) to a herdr pane (PaneLocator via pane.list/pane.process_info) → the caller's terminal_id; a caller owning no pane is the primary. bridge_reply now takes only content and resolves the worker from the connection (no sessionId). Live-verified: real PID→pane (contract test), and a non-pane MCP caller correctly gets a workers-only error. Single-host; the token path stays the split-host fallback. 58 tests green, IDE-clean.
This commit is contained in:
@@ -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());
|
||||
|
||||
@@ -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 <em>which worker</em>
|
||||
* is calling without the worker sending anything spoofable.
|
||||
*
|
||||
* <p>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;
|
||||
}
|
||||
}
|
||||
@@ -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<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"));
|
||||
})
|
||||
// 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() {
|
||||
|
||||
@@ -0,0 +1,40 @@
|
||||
package dev.ltms.bridged.mcp;
|
||||
|
||||
import dev.ltms.bridged.herdr.PaneLocator;
|
||||
|
||||
/**
|
||||
* Resolves <em>who is calling</em> 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}.
|
||||
*
|
||||
* <p>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);
|
||||
}
|
||||
}
|
||||
@@ -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<pid>' 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;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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);
|
||||
}
|
||||
@@ -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<Call> 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",
|
||||
|
||||
@@ -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}.
|
||||
*
|
||||
* <p>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());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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));
|
||||
}
|
||||
}
|
||||
@@ -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));
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user