From a1aecbf4fc661898c1f4b6bd6499366b45f4cd90 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 19 Jul 2026 09:59:02 +0200 Subject: [PATCH] CB-307 Increment 1: PrimaryRegistry + config + wiring - PrimaryRegistry: thread-safe single-slot registry with pin support - Primary config record (last positional, like Broker) - Wire capture in BridgeMcp (bridge_send and bridge_spawn handlers) - Construct PrimaryRegistry in Bridged.main - Tests: PrimaryRegistryTest + BridgedConfigTest primary config cases --- .../main/java/dev/ltms/bridged/Bridged.java | 7 +- .../ltms/bridged/config/BridgedConfig.java | 21 +++- .../java/dev/ltms/bridged/mcp/BridgeMcp.java | 9 +- .../dev/ltms/bridged/mcp/PrimaryRegistry.java | 68 +++++++++++ .../bridged/config/BridgedConfigTest.java | 35 ++++++ .../ltms/bridged/mcp/PrimaryRegistryTest.java | 108 ++++++++++++++++++ 6 files changed, 243 insertions(+), 5 deletions(-) create mode 100644 bridged/src/main/java/dev/ltms/bridged/mcp/PrimaryRegistry.java create mode 100644 bridged/src/test/java/dev/ltms/bridged/mcp/PrimaryRegistryTest.java diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 75b1374..aad56e7 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -13,6 +13,7 @@ import dev.ltms.bridged.inject.TurnListener; import dev.ltms.bridged.inject.WorkerPresence; import dev.ltms.bridged.mcp.BridgeMcp; import dev.ltms.bridged.mcp.ConnectionIdentity; +import dev.ltms.bridged.mcp.PrimaryRegistry; import dev.ltms.bridged.mcp.LsofPeerPidLookup; import dev.ltms.bridged.mcp.LsofProcessCwdLookup; import dev.ltms.bridged.msg.AmqpReplyInbox; @@ -136,8 +137,12 @@ public final class Bridged { // Caller identity is resolved from the connection (peer PID → herdr pane), not arguments. ConnectionIdentity identity = new ConnectionIdentity( new PaneLocator(herdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup()); + // CB-307: learn the primary's terminal from orchestration tool calls (or pin from config). + PrimaryRegistry primaryRegistry = new PrimaryRegistry( + cfg.primary() != null ? cfg.primary().terminal() : null); // Cast to ClaudeCodeLauncher: BridgeMcp is not yet migrated to PeerLauncher (Stage A scope). - BridgeMcp mcp = new BridgeMcp(messages, rendezvous, (ClaudeCodeLauncher) workers, sessions, identity, presence); + BridgeMcp mcp = new BridgeMcp(messages, rendezvous, (ClaudeCodeLauncher) workers, sessions, identity, presence, + primaryRegistry); // CB-303 part 3: single ordered shutdown hook. Drain sessions first while herdr is still // open (so releases reach the daemon), then stop poller/message/mcp/reaper, and close herdr diff --git a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java index 5329997..714d803 100644 --- a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java +++ b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java @@ -33,6 +33,9 @@ import java.util.Set; * @param spawnReadyPollMs poll interval while waiting for the worker to become injectable * @param broker external AMQP broker for durable reply delivery ({@code null} → in-memory, * soft-state {@code ReplyInbox}; present → the AMQP-backed adapter, CB-307 Stage 2) + * @param primary optional pinned primary terminal config ({@code null} → derived from connection); + * a non-blank {@code terminal} seeds {@code PrimaryRegistry} and prevents + * connection-derived overrides, CB-307 */ @JsonIgnoreProperties(ignoreUnknown = true) public record BridgedConfig( @@ -46,7 +49,8 @@ public record BridgedConfig( Lifecycle lifecycle, Integer spawnReadyTimeoutMs, Integer spawnReadyPollMs, - Broker broker) { + Broker broker, + Primary primary) { @JsonIgnoreProperties(ignoreUnknown = true) public record Bind(String host, int port) { @@ -187,6 +191,18 @@ public record BridgedConfig( } } + /** + * Optional pinned primary terminal config (CB-307). When present with a non-blank + * {@code terminal}, the bridge uses this as the primary's herdr identity instead of + * deriving it from the MCP connection. Useful when the primary runs off-host or in a + * non-herdr terminal where connection-derived identity is unavailable. + * + * @param terminal the primary's herdr {@code terminal_id} ({@code null}/blank → derive) + */ + @JsonIgnoreProperties(ignoreUnknown = true) + public record Primary(String terminal) { + } + /** * Subscription boundary. Only these hosts may back a worker's * {@code ANTHROPIC_BASE_URL}; the primary must carry none. @@ -258,6 +274,7 @@ public record BridgedConfig( Integer timeout = (spawnReadyTimeoutMs != null) ? spawnReadyTimeoutMs : 20000; Integer pollMs = (spawnReadyPollMs != null) ? spawnReadyPollMs : 300; // broker is left as-is: null (or an empty/blank uri) keeps the in-memory soft-state inbox. - return new BridgedConfig(b, herdrSocket, worker, workers, defaultWorker, g, worktreeRoot, l, timeout, pollMs, broker); + // primary is left as-is: null defaults to connection-derived identity. + return new BridgedConfig(b, herdrSocket, worker, workers, defaultWorker, g, worktreeRoot, l, timeout, pollMs, broker, primary); } } 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 104eb02..057b548 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java @@ -61,7 +61,8 @@ public final class BridgeMcp { private final McpSyncServer server; public BridgeMcp(MessageService messages, Rendezvous rendezvous, ClaudeCodeLauncher workers, - SessionManager sessions, ConnectionIdentity identity, WorkerPresence presence) { + SessionManager sessions, ConnectionIdentity identity, WorkerPresence presence, + PrimaryRegistry primaryRegistry) { McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get(); this.transport = HttpServletStreamableServerTransportProvider.builder() .jsonMapper(json) @@ -81,7 +82,9 @@ public final class BridgeMcp { this.server = McpServer.sync(transport) .serverInfo("bridge", "0.1.0") .capabilities(McpSchema.ServerCapabilities.builder().tools(true).build()) - .toolCall(sendTool(), (_, req) -> { + .toolCall(sendTool(), (exchange, req) -> { + String caller = callerTerminal(exchange); + if (caller != null) primaryRegistry.record(caller); Map a = req.arguments(); String turnId = str(a, "turnId"); if (turnId != null && !turnId.isBlank()) { @@ -108,6 +111,8 @@ public final class BridgeMcp { }) // Fleet management (CB-108): spawn/list/stop over ClaudeCodeLauncher. .toolCall(spawnTool(), (exchange, req) -> { + String caller = callerTerminal(exchange); + if (caller != null) primaryRegistry.record(caller); Map a = req.arguments(); // CB-112: worker inherits the primary's cwd unless the call pins one. // CB-301: carry the caller's identity as the session owner (null for the primary). diff --git a/bridged/src/main/java/dev/ltms/bridged/mcp/PrimaryRegistry.java b/bridged/src/main/java/dev/ltms/bridged/mcp/PrimaryRegistry.java new file mode 100644 index 0000000..a04d86b --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/PrimaryRegistry.java @@ -0,0 +1,68 @@ +package dev.ltms.bridged.mcp; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.Optional; +import java.util.concurrent.atomic.AtomicReference; + +/** + * Single-slot, thread-safe registry for the primary's herdr {@code terminal_id}. + * + *

Populated from the caller terminal of orchestration-side MCP tools + * ({@code bridge_send}, {@code bridge_spawn}) — tools that only the primary calls. + * A pinned terminal (from config) seeds the registry at construction and makes + * subsequent {@link #record(String)} calls no-ops. + * + *

The push loop ({@code ReplyPushLoop}) uses {@link #isKnown()} to decide + * whether active nudging is possible; an empty registry means the primary is + * off-host or non-herdr and delivery falls back to pull. + */ +public final class PrimaryRegistry { + + private static final Logger log = LoggerFactory.getLogger(PrimaryRegistry.class); + + private final AtomicReference terminal = new AtomicReference<>(); + private final boolean pinned; + + /** + * @param pinnedTerminal an optional pinned terminal from config ({@code null}/blank = unpinned) + */ + public PrimaryRegistry(String pinnedTerminal) { + if (pinnedTerminal != null && !pinnedTerminal.isBlank()) { + this.terminal.set(pinnedTerminal); + this.pinned = true; + log.info("primary terminal pinned: {}", pinnedTerminal); + } else { + this.pinned = false; + } + } + + /** + * Record a terminal_id. No-op when: + *

    + *
  • the registry is pinned (config override), + *
  • {@code terminalId} is {@code null} or blank (non-herdr caller). + *
+ */ + public void record(String terminalId) { + if (pinned) return; + if (terminalId == null || terminalId.isBlank()) return; + String prev = terminal.getAndSet(terminalId); + if (prev == null) { + log.debug("primary terminal learned: {}", terminalId); + } else if (!prev.equals(terminalId)) { + log.debug("primary terminal changed: {} -> {}", prev, terminalId); + } + } + + /** The known primary terminal, or empty if not yet learned (and not pinned). */ + public Optional primaryTerminal() { + return Optional.ofNullable(terminal.get()); + } + + /** {@code true} once a terminal has been recorded (or was pinned at construction). */ + public boolean isKnown() { + return terminal.get() != null; + } +} diff --git a/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java b/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java index 031f31b..2ec3f90 100644 --- a/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java @@ -125,4 +125,39 @@ class BridgedConfigTest { assertNotNull(cfg.broker()); assertFalse(cfg.broker().isConfigured(), "an empty uri must not enable AMQP"); } + + @Test + void absentPrimaryBlockLeavesPrimaryNull(@TempDir Path dir) throws Exception { + Path f = dir.resolve("no-primary.yaml"); + Files.writeString(f, "bind:\n port: 8080\n"); + + BridgedConfig cfg = BridgedConfig.load(f); + assertNull(cfg.primary(), "no primary: block → null → connection-derived identity"); + } + + @Test + void primaryBlockWithTerminalPinsIdentity(@TempDir Path dir) throws Exception { + Path f = dir.resolve("primary-pinned.yaml"); + Files.writeString(f, """ + bind: + port: 8080 + primary: + terminal: term_fixed + """); + + BridgedConfig cfg = BridgedConfig.load(f); + assertNotNull(cfg.primary()); + assertEquals("term_fixed", cfg.primary().terminal()); + } + + @Test + void primaryBlockWithBlankTerminalDefaultsToDerived(@TempDir Path dir) throws Exception { + Path f = dir.resolve("primary-blank.yaml"); + Files.writeString(f, "bind:\n port: 8080\nprimary:\n terminal: \"\"\n"); + + BridgedConfig cfg = BridgedConfig.load(f); + assertNotNull(cfg.primary()); + assertTrue(cfg.primary().terminal() == null || cfg.primary().terminal().isBlank(), + "a blank terminal in yaml should be treated as absent — null or empty are equivalent"); + } } diff --git a/bridged/src/test/java/dev/ltms/bridged/mcp/PrimaryRegistryTest.java b/bridged/src/test/java/dev/ltms/bridged/mcp/PrimaryRegistryTest.java new file mode 100644 index 0000000..1a26485 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/mcp/PrimaryRegistryTest.java @@ -0,0 +1,108 @@ +package dev.ltms.bridged.mcp; + +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.*; + +/** + * Unit tests for {@link PrimaryRegistry}: pin vs record, isKnown transitions, + * null/blank guard. + */ +class PrimaryRegistryTest { + + @Test + void unpinnedInitiallyUnknown() { + var reg = new PrimaryRegistry(null); + assertFalse(reg.isKnown()); + assertTrue(reg.primaryTerminal().isEmpty()); + } + + @Test + void unpinnedAcceptsBlankAsAbsent() { + var reg = new PrimaryRegistry(""); + assertFalse(reg.isKnown()); + assertTrue(reg.primaryTerminal().isEmpty()); + } + + @Test + void pinnedFromConstruction() { + var reg = new PrimaryRegistry("term_fixed"); + assertTrue(reg.isKnown()); + assertEquals("term_fixed", reg.primaryTerminal().orElseThrow()); + } + + @Test + void recordWhenUnpinnedSetsTheTerminal() { + var reg = new PrimaryRegistry(null); + reg.record("term_abc"); + assertTrue(reg.isKnown()); + assertEquals("term_abc", reg.primaryTerminal().orElseThrow()); + } + + @Test + void recordWithNullDoesNothingWhenUnpinned() { + var reg = new PrimaryRegistry(null); + reg.record(null); + assertFalse(reg.isKnown()); + } + + @Test + void recordWithBlankDoesNothingWhenUnpinned() { + var reg = new PrimaryRegistry(null); + reg.record(" "); + assertFalse(reg.isKnown()); + } + + @Test + void recordOverwritesWhenUnpinned() { + var reg = new PrimaryRegistry(null); + reg.record("term_first"); + assertEquals("term_first", reg.primaryTerminal().orElseThrow()); + reg.record("term_second"); + assertEquals("term_second", reg.primaryTerminal().orElseThrow()); + } + + @Test + void recordIsIgnoredWhenPinned() { + var reg = new PrimaryRegistry("term_pinned"); + reg.record("term_other"); + assertEquals("term_pinned", reg.primaryTerminal().orElseThrow(), "pinned value must survive record"); + } + + @Test + void nullRecordIsIgnoredWhenPinned() { + var reg = new PrimaryRegistry("term_pinned"); + reg.record(null); + assertTrue(reg.isKnown()); + assertEquals("term_pinned", reg.primaryTerminal().orElseThrow()); + } + + @Test + void blankRecordIsIgnoredWhenPinned() { + var reg = new PrimaryRegistry("term_pinned"); + reg.record(" "); + assertTrue(reg.isKnown()); + assertEquals("term_pinned", reg.primaryTerminal().orElseThrow()); + } + + @Test + void isKnownFalseAfterConstructionWithNull() { + var reg = new PrimaryRegistry(null); + assertFalse(reg.isKnown()); + } + + @Test + void isKnownAfterRecord() { + var reg = new PrimaryRegistry(null); + reg.record("term_x"); + assertTrue(reg.isKnown()); + } + + @Test + void primaryTerminalRoundTrip() { + var reg = new PrimaryRegistry(null); + assertTrue(reg.primaryTerminal().isEmpty()); + reg.record("term_found"); + assertEquals("term_found", reg.primaryTerminal().get()); + } +}