CB-548: guard primary singleton to PRIMARY callers; open waiter before enqueue
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.
This commit is contained in:
@@ -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<String, Object> 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<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).
|
||||
@@ -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).
|
||||
*
|
||||
* <p>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.
|
||||
*
|
||||
* <p>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);
|
||||
|
||||
@@ -316,10 +316,11 @@ public final class MessageService {
|
||||
/**
|
||||
* As {@link #send(String, String, long)}, but with an accepted-delivery hook.
|
||||
*
|
||||
* <p>{@code onAccepted} is invoked exactly once, the instant this send becomes the <em>accepted
|
||||
* target turn</em> — it has won {@code target}'s send lock and queued delivery through the
|
||||
* injector. It is <em>not</em> invoked when the send is {@link Outcome#BUSY} (lock never taken)
|
||||
* nor before acceptance. A caller uses this to record that <em>it</em> now owns the delegation's
|
||||
* <p>{@code onAccepted} is invoked exactly once, once this send has won {@code target}'s send
|
||||
* lock and so become the <em>accepted target turn</em> — it runs <em>before</em> 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 <em>not</em> invoked when the send is {@link Outcome#BUSY}
|
||||
* (lock never taken). A caller uses this to record that <em>it</em> 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<Void> 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<Rendezvous.Resolution> 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<Void> 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);
|
||||
}
|
||||
|
||||
@@ -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");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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.
|
||||
|
||||
Reference in New Issue
Block a user