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
This commit is contained in:
Dai Ha
2026-07-19 09:59:02 +02:00
parent 131e7b1ccd
commit a1aecbf4fc
6 changed files with 243 additions and 5 deletions
@@ -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
@@ -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);
}
}
@@ -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<String, Object> 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<String, Object> 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).
@@ -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}.
*
* <p>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.
*
* <p>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<String> 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:
* <ul>
* <li>the registry is pinned (config override),
* <li>{@code terminalId} is {@code null} or blank (non-herdr caller).
* </ul>
*/
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<String> 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;
}
}
@@ -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");
}
}
@@ -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());
}
}