From dd906526c0010fe3498b56c647d2c985502f8cd9 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 13 Aug 2026 19:19:14 +0200 Subject: [PATCH] CB-548: guard primary singleton to PRIMARY callers; open waiter before enqueue MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Fix two handler-level bugs found in PR #21: - Only PRIMARY callers may update PrimaryRegistry.record (the legacy singleton 'primary' fallback for no-delegation inbox nudges). An architect SEND previously recorded its terminal as the fallback; the per-target delegation map does not cure the singleton. New BridgeMcp.recordPrimarySingleton uses the resolved role (caller.isPrimary()) — named leads (PRIMARY) still record, architects never do. - MessageService.send now opens the rendezvous waiter BEFORE queueing delivery, fixing both the enqueue-before-open fast-reply race (a fast reply no longer orphans into the inbox) and callback-failure ordering: a throwing onAccepted (public callback) fails the send cleanly with no stale waiter and no queued, orphanable message. Tests: architect SEND vs lead SEND primary-singleton regression; throwing onAccepted leaves no stale waiter or queued orphan. --- .../java/dev/ltms/bridged/mcp/BridgeMcp.java | 28 +++++++++- .../dev/ltms/bridged/msg/MessageService.java | 56 +++++++++++-------- .../dev/ltms/bridged/mcp/BridgeMcpTest.java | 33 +++++++++++ .../ltms/bridged/msg/MessageServiceTest.java | 23 ++++++++ 4 files changed, 115 insertions(+), 25 deletions(-) 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 fb45f93..64785e8 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java @@ -127,7 +127,11 @@ public final class BridgeMcp { str(req.arguments(), "sessionId")); if (denied != null) return denied; String caller = callerTerminal(exchange); - if (caller != null) primaryRegistry.record(caller); + // CB-548: only a PRIMARY caller may claim the legacy singleton "primary" fallback. + // An architect delegates as its own pane but must never become the fallback that + // no-delegation inbox nudges target as if it were the primary (the per-target + // delegation map does not cure the singleton). + recordPrimarySingleton(primaryRegistry, caller, principal(exchange)); Map a = req.arguments(); String target = str(a, "sessionId"); String content = str(a, "content"); @@ -188,7 +192,9 @@ public final class BridgeMcp { McpSchema.CallToolResult denied = deny(exchange, Authz.Action.SPAWN, null); if (denied != null) return denied; String caller = callerTerminal(exchange); - if (caller != null) primaryRegistry.record(caller); + // SPAWN is already auth-gated to PRIMARY (architects can never call it), but + // enforce the same invariant here: only a PRIMARY may claim the legacy singleton. + recordPrimarySingleton(primaryRegistry, caller, principal(exchange)); 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). @@ -307,6 +313,24 @@ public final class BridgeMcp { return error(reason + ": " + caller.describe() + " may not " + action); } + /** + * Update the legacy singleton "primary" fallback used for no-delegation inbox nudges (CB-548). + * + *

Only {@link Role#PRIMARY} callers — the unnamed primary and named leads alike — may claim + * it. An architect delegates as its own pane but must never become the fallback: the per-target + * delegation map ({@code PrimaryRegistry#recordDelegation}) does not cure the singleton, so an + * architect left here would draw nudges that belong to a primary. The decision uses the resolved + * role, never name/kind sniffing. A null {@code caller} (legacy/no-auth path) records nothing. + * + *

Split out of the tool handlers so the guard is unit-testable without fabricating an SDK + * {@code McpSyncServerExchange} (same pattern as {@link #denyFor}/{@link #principalFrom}). + */ + static void recordPrimarySingleton(PrimaryRegistry registry, String callerTerminal, Principal caller) { + if (caller != null && caller.isPrimary()) { + registry.record(callerTerminal); + } + } + /** 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); diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java index de47898..1ceecb2 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -316,10 +316,11 @@ public final class MessageService { /** * As {@link #send(String, String, long)}, but with an accepted-delivery hook. * - *

{@code onAccepted} is invoked exactly once, the instant this send becomes the accepted - * target turn — it has won {@code target}'s send lock and queued delivery through the - * injector. It is not invoked when the send is {@link Outcome#BUSY} (lock never taken) - * nor before acceptance. A caller uses this to record that it now owns the delegation's + *

{@code onAccepted} is invoked exactly once, once this send has won {@code target}'s send + * lock and so become the accepted target turn — it runs before delivery is + * queued, so a throwing hook fails the send cleanly (the waiter it already opened is closed and + * nothing is left queued). It is not invoked when the send is {@link Outcome#BUSY} + * (lock never taken). A caller uses this to record that it now owns the delegation's * reply routing (CB-548: {@code PrimaryRegistry} delegator ownership) — recording only on * acceptance means a concurrent sender that times out {@code BUSY} can never steal ownership it * never earned. {@code null} disables the hook. @@ -332,27 +333,36 @@ public final class MessageService { return new Reply(Outcome.BUSY, null); // another send held the session the whole window } try { - CompletableFuture delivered = injector.enqueue(target, content); - // The send is now the accepted target turn — it holds the lock and has queued delivery. - // Record delegator ownership here, never at request time (CB-548). - if (onAccepted != null) { - onAccepted.run(); - } + // Open the waiter BEFORE queueing delivery (CB-548). A fast reply — the worker already + // injectable the instant we enqueue — otherwise arrives before the waiter is registered + // and orphans into the inbox while this send blocks to the timeout (the enqueue-before- + // open race). Opening first also means a throwing onAccepted (fired before enqueue) or an + // enqueue failure is safely closed by the finally below: nothing is left queued, and the + // failed send leaves no stale waiter behind. CompletableFuture reply = rendezvous.open(target); try { - Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS); - return recorded(new Reply(outcomeOf(r.kind()), r.text(), r.turnId())); - } catch (TimeoutException e) { - boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally(); - log.debug("send to {} timed out (delivered={})", target, wasDelivered); - return recorded(new Reply( - wasDelivered ? Outcome.TIMED_OUT_WORKING : Outcome.TIMED_OUT_QUEUED, null)); - } catch (ExecutionException e) { - Throwable cause = e.getCause(); - throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - throw new IllegalStateException("interrupted awaiting reply from " + target, e); + // The send has won the lock; the accepted-delivery hook records delegator ownership + // here (CB-548). It runs BEFORE enqueue so a throwing hook — onAccepted is now a + // public callback — fails the send without queuing a message that would orphan. + if (onAccepted != null) { + onAccepted.run(); + } + CompletableFuture delivered = injector.enqueue(target, content); + try { + Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS); + return recorded(new Reply(outcomeOf(r.kind()), r.text(), r.turnId())); + } catch (TimeoutException e) { + boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally(); + log.debug("send to {} timed out (delivered={})", target, wasDelivered); + return recorded(new Reply( + wasDelivered ? Outcome.TIMED_OUT_WORKING : Outcome.TIMED_OUT_QUEUED, null)); + } catch (ExecutionException e) { + Throwable cause = e.getCause(); + throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new IllegalStateException("interrupted awaiting reply from " + target, e); + } } finally { rendezvous.close(target, 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 5daccde..6be944c 100644 --- a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java @@ -484,4 +484,37 @@ class BridgeMcpTest { assertTrue(out.contains("\"architect\":\"lead-designer\""), out); assertTrue(out.contains("\"sessionId\":\"term_design\""), out); } + + /** + * CB-548: an architect SEND delegates as its own pane (recording the per-target delegation) but + * must NEVER become the legacy singleton "primary" fallback — the per-target map does not cure + * the singleton, so an architect left there would draw no-delegation inbox nudges meant for a + * primary. Only PRIMARY callers (the unnamed primary and named leads alike) may claim it, and + * the decision keys on the resolved role, not name/kind sniffing. + */ + @Test + void architectSendDoesNotClaimThePrimarySingletonButALeadSendStillCan() { + // Architect SEND: does not change the legacy primary fallback. + PrimaryRegistry reg = new PrimaryRegistry(null); + BridgeMcp.recordPrimarySingleton(reg, "term_design", Principal.architect("design", "term_design", 400)); + assertTrue(reg.primaryTerminal().isEmpty(), + "an architect must never become the legacy primary fallback"); + + // Lead SEND (a named PRIMARY) still claims it — preserved from CB-530/CB-532. + PrimaryRegistry leadReg = new PrimaryRegistry(null); + BridgeMcp.recordPrimarySingleton(leadReg, "term_lead_opus", Principal.leader("opus", "term_lead_opus", 100)); + assertEquals("term_lead_opus", leadReg.primaryTerminal().orElseThrow(), + "a named lead is a primary and may claim the fallback"); + + // Unnamed primary likewise. + PrimaryRegistry primaryReg = new PrimaryRegistry(null); + BridgeMcp.recordPrimarySingleton(primaryReg, "term_p", Principal.primary(50)); + assertEquals("term_p", primaryReg.primaryTerminal().orElseThrow(), + "an unnamed primary may claim the fallback"); + + // A null caller (legacy/no-auth path) records nothing. + PrimaryRegistry legacy = new PrimaryRegistry(null); + BridgeMcp.recordPrimarySingleton(legacy, "term_x", null); + assertTrue(legacy.primaryTerminal().isEmpty(), "no caller means nothing is recorded"); + } } diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java index 46c7acd..5668890 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -17,6 +17,7 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; /** @@ -422,6 +423,28 @@ class MessageServiceTest { assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome()); } + /** + * CB-548: {@code onAccepted} is a public callback, so a throwing one must not orphan the turn. + * The waiter is opened first, then the hook runs BEFORE delivery is queued — so a throw fails + * the send loudly, closes its waiter, and never enqueues a message the worker would pick up and + * reply into the void. + */ + @Test + void aThrowingAcceptedHookLeavesNoStaleWaiterOrQueuedOrphan() { + assertThrows(IllegalStateException.class, + () -> messages.send(T, "doomed", 500, + () -> { throw new IllegalStateException("ownership hook failed"); }), + "a throwing ownership hook fails the send loudly"); + + assertFalse(rendezvous.isWaiting(T), "the failed send must not leave a stale rendezvous waiter"); + // Give the injector a delivery window: with nothing enqueued, nothing may reach the worker. + injector.onStatus(T, AgentStatus.IDLE); + boolean doomedQueued = herdr.calls.stream() + .anyMatch(c -> c.method().equals("agent.prompt") + && String.valueOf(c.params()).contains("doomed")); + assertFalse(doomedQueued, "a throwing ownership hook must not leave a queued, orphanable message"); + } + /** * The async (fire-and-poll) path runs the same {@code send} on a background thread, so the * accepted-delivery hook must thread through it — ownership is recorded exactly as blocking sends.