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 3056344..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,21 +127,31 @@ 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(); - // CB-532: remember WHICH lead is waiting on this worker, so its reply nudge goes - // back to that lead rather than to whichever one happened to send first. - primaryRegistry.recordDelegation(str(a, "sessionId"), caller); + String target = str(a, "sessionId"); + String content = str(a, "content"); String turnId = str(a, "turnId"); if (turnId != null && !turnId.isBlank()) { // Answering a worker's bridge_ask (CB-205): resolve its blocked question and - // block for the worker's reply as it resumes the same turn. - return answer(messages, turnId, str(a, "content"), timeoutMs(a)); + // block for the worker's reply as it resumes the same turn. This is the same + // delegation, so ownership is left untouched (CB-548) — never re-recorded. + return answer(messages, turnId, content, timeoutMs(a)); } + // CB-548: delegator ownership (which lead's reply nudge this worker routes to, + // CB-532) is recorded only once the send is ACCEPTED — MessageService has won the + // session lock and queued delivery — via the accepted-delivery callback, never at + // request time. A concurrent sender that times out BUSY therefore cannot steal a + // live turn's reply routing without ever owning the turn. + Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, caller); // wait defaults to true (block for the reply); wait:false is fire-and-poll. return Boolean.FALSE.equals(a.get("wait")) - ? sendAsync(messages, str(a, "sessionId"), str(a, "content")) - : send(messages, str(a, "sessionId"), str(a, "content"), timeoutMs(a)); + ? sendAsync(messages, target, content, onAccepted) + : send(messages, target, content, timeoutMs(a), onAccepted); }) // bridge_reply's identity is the CONNECTION, never an argument — so the authz check // is "is this caller a worker at all", and it can only ever reply as itself. @@ -182,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). @@ -301,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); @@ -343,12 +373,22 @@ public final class BridgeMcp { /** {@code bridge_send}: delegate {@code content} to a worker session and block for its reply. */ static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content, Long timeoutMs) { + return send(messages, sessionId, content, timeoutMs, null); + } + + /** + * As {@link #send(MessageService, String, String, Long)}, wiring an accepted-delivery hook + * (CB-548): {@code onAccepted} records delegator ownership the instant the send is accepted, so + * a BUSY interloper never claims a turn it did not win. {@code null} disables recording. + */ + static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content, + Long timeoutMs, Runnable onAccepted) { if (isBlank(sessionId) || isBlank(content)) { return error("sessionId and content are required"); } long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs); try { - return formatReply(messages.send(sessionId, content, timeout), timeout); + return formatReply(messages.send(sessionId, content, timeout, onAccepted), timeout); } catch (HerdrException e) { return error("herdr error contacting session " + sessionId + ": " + e.getMessage()); } @@ -417,10 +457,19 @@ public final class BridgeMcp { * immediately (fire-and-poll), so a long task isn't cut off by the caller's MCP call timeout. */ static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content) { + return sendAsync(messages, sessionId, content, null); + } + + /** + * As {@link #sendAsync(MessageService, String, String)}, wiring the accepted-delivery hook + * (CB-548) so an async flooding send records delegator ownership exactly once it is accepted. + */ + static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content, + Runnable onAccepted) { if (isBlank(sessionId) || isBlank(content)) { return error("sessionId and content are required"); } - String ticket = messages.sendAsync(sessionId, content); + String ticket = messages.sendAsync(sessionId, content, onAccepted); return text("accepted — task delegated. Poll bridge_poll with ticket=" + ticket); } diff --git a/bridged/src/main/java/dev/ltms/bridged/mcp/PrimaryRegistry.java b/bridged/src/main/java/dev/ltms/bridged/mcp/PrimaryRegistry.java index 64b156a..a7d4210 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/PrimaryRegistry.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/PrimaryRegistry.java @@ -67,12 +67,14 @@ public final class PrimaryRegistry { } /** - * Record that {@code leadTerminal} delegated to worker {@code target} (CB-532). + * Record that {@code leadTerminal} owns the accepted delegation of worker {@code target} (CB-532). * - *

Called at {@code bridge_send} time, where both halves are known: the target is the tool's - * argument and the lead is resolved from the connection. Last writer wins — if a second lead - * takes over a worker, replies follow the lead that most recently delegated to it, which is the - * one waiting. + *

Called from the {@code MessageService} accepted-delivery hook — only after a send has won + * the session's send lock and queued delivery — where both halves are known (CB-548). It is + * deliberately not called at {@code bridge_send} request time: a concurrent sender that + * times out {@code BUSY} must not steal a live delegation's reply routing without ever owning + * the turn. Last writer wins — if a second lead's later send is accepted, replies follow the + * lead that most recently delegated to it, which is the one waiting. */ public void recordDelegation(String target, String leadTerminal) { if (target == null || target.isBlank() || leadTerminal == null || leadTerminal.isBlank()) { 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 befa426..1ceecb2 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -310,6 +310,22 @@ public final class MessageService { * worker replies via {@link Rendezvous} or {@code timeoutMillis} elapses. */ public Reply send(String target, String content, long timeoutMillis) { + return send(target, content, timeoutMillis, null); + } + + /** + * As {@link #send(String, String, long)}, but with an accepted-delivery hook. + * + *

{@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. + */ + public Reply send(String target, String content, long timeoutMillis, Runnable onAccepted) { long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L; ReentrantLock lock = sessionLocks.computeIfAbsent(target, _ -> new ReentrantLock()); @@ -317,22 +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); + // 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); } @@ -440,9 +470,21 @@ public final class MessageService { * @return the ticket to poll for the eventual result */ public String sendAsync(String target, String content) { + return sendAsync(target, content, null); + } + + /** + * As {@link #sendAsync(String, String)}, with the accepted-delivery hook of + * {@link #send(String, String, long, Runnable)} — the running {@code send} invokes {@code onAccepted} + * the moment it becomes the accepted target turn, so async flooding records delegator ownership + * exactly as the blocking path does (CB-548). + * + * @return the ticket to poll for the eventual result + */ + public String sendAsync(String target, String content, Runnable onAccepted) { String ticket = "task-" + ticketSeq.incrementAndGet(); - CompletableFuture future = - CompletableFuture.supplyAsync(() -> send(target, content, ASYNC_TIMEOUT_MS), asyncExecutor); + CompletableFuture future = CompletableFuture.supplyAsync( + () -> send(target, content, ASYNC_TIMEOUT_MS, onAccepted), asyncExecutor); tasks.put(ticket, new Task(target, future, System.nanoTime())); pruneTerminalTickets(); log.debug("async send {} -> {}", ticket, target); diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java b/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java index f3895ec..5136537 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java @@ -74,15 +74,30 @@ public final class Rendezvous { /** * Register a waiter for {@code session} — the await side of the public {@code resolve*} methods. * The caller must hold that session's send lock. + * + *

Atomic fail-if-present (CB-548): if a waiter is already registered for {@code session}, an + * {@link IllegalStateException} is thrown rather than replacing the first — so any future + * invariant violation fails loudly instead of silently swapping the waiter another send is + * blocked on. {@code MessageService} serializes sends per session (the send lock), so in correct + * code a double open is impossible; this is a tripwire for the day that no longer holds. */ public CompletableFuture open(String session) { CompletableFuture waiter = new CompletableFuture<>(); - waiters.put(session, waiter); + CompletableFuture existing = waiters.putIfAbsent(session, waiter); + if (existing != null) { + throw new IllegalStateException( + "rendezvous double-open for session " + session + " — a waiter is already registered"); + } return waiter; } - /** Remove {@code waiter} for {@code session} (only if it is still the registered one). */ - void close(String session, CompletableFuture waiter) { + /** + * Remove {@code waiter} for {@code session}, only if it is still the registered one. The + * symmetric complement of {@link #open}: a terminal send deregisters its waiter so the next + * send on the session may {@link #open} a fresh one (CB-548 makes double-open an error, so a + * successful {@code open} after a finished turn requires this close to have happened first). + */ + public void close(String session, CompletableFuture waiter) { waiters.remove(session, waiter); } diff --git a/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java b/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java index 4b67084..8462072 100644 --- a/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java @@ -286,10 +286,12 @@ class CompletionResolverTest { // The turn as the injector captured it at delivery (waiter + pre-turn baseline). var turnN = new CompletionResolver.InFlight(waiterN, "an earlier answer"); - // Turn N is resolved by the worker's explicit reply. + // Turn N is resolved by the worker's explicit reply, and its send deregisters the waiter. assertTrue(rendezvous.resolve("term_a", "N replied")); + rendezvous.close("term_a", waiterN); // the sender's finally, before the next turn opens - // Turn N+1's send opens its own waiter on the same session (replacing the registered one). + // Turn N+1's send opens its own waiter on the same session (CB-548: open fails if the + // previous waiter is still registered, so a clean turn deregisters it first as above). var waiterN1 = rendezvous.open("term_a"); resolver.resolve("term_a", turnN); // turn N's completion fallback finally fires 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 66d11f2..5668890 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -5,6 +5,7 @@ import dev.ltms.bridged.herdr.AgentStatus; import dev.ltms.bridged.herdr.FakeHerdr; import dev.ltms.bridged.herdr.HerdrException; import dev.ltms.bridged.inject.CompletionResolver; +import dev.ltms.bridged.mcp.PrimaryRegistry; import dev.ltms.bridged.inject.Injector; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -16,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; /** @@ -321,6 +323,141 @@ class MessageServiceTest { assertEquals("first done", firstReply.text()); } + // --- CB-548: delegator ownership is recorded only on an ACCEPTED send ---------------------- + + private static final String LEAD_L = "term_lead_l"; + private static final String LEAD_A = "term_lead_a"; + + /** + * The bug CB-548 fixes: L holds worker W, then architect A attempts W and times out BUSY. With + * delegator ownership recorded at {@code bridge_send} request time, A's rejected call + * would overwrite L — and W's late no-waiter reply would be pushed to A, who never owned the + * turn. The accepted-delivery hook must not fire for a BUSY send, so L stays the delegator. + */ + @Test + void busySenderDoesNotBecomeTheDelegatingOwner() throws Exception { + PrimaryRegistry reg = new PrimaryRegistry(null); + // L accepts a delegation to W: the send wins the lock and queues delivery → L is recorded. + CompletableFuture first = CompletableFuture.supplyAsync( + () -> messages.send(T, "first", 5000, () -> reg.recordDelegation(T, LEAD_L))); + awaitWaiting(); + assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(), + "an accepted send owns the delegation"); + + // A attempts W while L holds it → BUSY (lock never taken) → its hook never fires. + MessageService.Reply busy = messages.send(T, "second", 100, () -> reg.recordDelegation(T, LEAD_A)); + assertEquals(MessageService.Outcome.BUSY, busy.outcome()); + assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(), + "a BUSY send must not steal the delegator ownership it never earned"); + + // L completes so the test thread is not left pinned. + injector.onStatus(T, AgentStatus.IDLE); + injector.onStatus(T, AgentStatus.WORKING); + assertTrue(rendezvous.resolve(T, "first done")); + MessageService.Reply firstReply = first.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.REPLIED, firstReply.outcome()); + } + + /** + * Once L's accepted delegation is fully done, a later accepted send from A may + * legitimately become the new delegator — ownership follows the turn, not the first caller. + */ + @Test + void anAcceptedSendAfterThePriorOwnerFinishesBecomesTheNewOwner() throws Exception { + PrimaryRegistry reg = new PrimaryRegistry(null); + CompletableFuture first = CompletableFuture.supplyAsync( + () -> messages.send(T, "first", 5000, () -> reg.recordDelegation(T, LEAD_L))); + awaitWaiting(); + injector.onStatus(T, AgentStatus.IDLE); + injector.onStatus(T, AgentStatus.WORKING); + assertTrue(rendezvous.resolve(T, "first done")); + MessageService.Reply firstReply = first.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.REPLIED, firstReply.outcome()); + assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(), "L owned the first turn"); + + // L finished; A's later accepted send takes the delegation over. + CompletableFuture second = CompletableFuture.supplyAsync( + () -> messages.send(T, "second", 5000, () -> reg.recordDelegation(T, LEAD_A))); + awaitWaiting(); + assertEquals(LEAD_A, reg.nudgeTargetFor(T).orElseThrow(), + "an accepted send after the owner finished becomes the new delegator"); + + injector.onStatus(T, AgentStatus.IDLE); + injector.onStatus(T, AgentStatus.WORKING); + assertTrue(rendezvous.resolve(T, "second done")); + MessageService.Reply secondReply = second.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.REPLIED, secondReply.outcome()); + } + + /** + * CB-548 requirement: answering an existing {@code bridge_ask} is the SAME delegation, so it must + * not rewrite ownership. L accepted the send (owned), the worker paused to ask, and L answers via + * turnId — ownership stays L throughout; the answer path never touches the registry. + */ + @Test + void answeringAnAskDoesNotRewriteDelegatorOwnership() throws Exception { + PrimaryRegistry reg = new PrimaryRegistry(null); + CompletableFuture send = CompletableFuture.supplyAsync( + () -> messages.send(T, "do X", 5000, () -> reg.recordDelegation(T, LEAD_L))); + awaitWaiting(); + assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(), "L owns the delegation"); + + injector.onStatus(T, AgentStatus.IDLE); // deliver + injector.onStatus(T, AgentStatus.WORKING); // worker picks it up, then pauses to ask + + CompletableFuture ask = + CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000)); + MessageService.Reply q = send.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.QUESTION, q.outcome()); + assertNotNull(q.turnId()); + + // L answers the ask on the same turn; the answer path must not touch ownership. + CompletableFuture answer = + CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000)); + assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); + awaitWaiting(); // the answering send reopened its forward waiter + assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(), + "answering an ask keeps L as the delegator — ownership is not rewritten"); + + assertTrue(rendezvous.resolve(T, "done")); + 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. + */ + @Test + void asyncSendRecordsOwnershipOnAcceptance() throws Exception { + PrimaryRegistry reg = new PrimaryRegistry(null); + messages.sendAsync(T, "async task", () -> reg.recordDelegation(T, LEAD_L)); + awaitWaiting(); // the background send won the lock, queued, and opened its waiter + assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(), + "the async path records delegator ownership on acceptance, like the blocking path"); + } + // --- CB-307 reply inbox ---------------------------------------------------------------- @Test diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/RendezvousTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/RendezvousTest.java index 927a169..14199ee 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/RendezvousTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/RendezvousTest.java @@ -8,6 +8,8 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; /** @@ -20,6 +22,26 @@ class RendezvousTest { private final Rendezvous rendezvous = new Rendezvous(); + /** + * CB-548: an {@code open} is atomic fail-if-present so a double open can never replace the first + * waiter another send is blocked on. {@code MessageService} serializes sends per session, so in + * correct code this cannot happen — the rejection is a loud tripwire for an invariant violation, + * and the first waiter must survive it and stay resolvable. + */ + @Test + void openRejectsADoubleOpenAndKeepsTheFirstWaiterRegisteredAndResolvable() { + CompletableFuture first = rendezvous.open(W); + + assertThrows(IllegalStateException.class, () -> rendezvous.open(W), + "a second open while one is registered is rejected loudly, not a silent replace"); + + assertSame(first, rendezvous.currentWaiter(W), "the first waiter remains the registered one"); + assertTrue(rendezvous.resolve(W, "first wins"), "the first waiter is still resolvable"); + Rendezvous.Resolution r = first.getNow(null); + assertEquals(Rendezvous.Kind.REPLY, r.kind()); + assertEquals("first wins", r.text(), "the resolution lands on the first waiter, not the rejected one"); + } + @Test void openAskMintsAUniqueTurnScopedToItsSessionAndCoalescesDuplicates() { Rendezvous.AskTicket t1 = rendezvous.openAsk(W);