From b39f700765aa4a813e78327a4d6cb2aed0b8cf1a Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 13 Aug 2026 18:55:54 +0200 Subject: [PATCH] CB-548: record delegator ownership only when a send is accepted Record PrimaryRegistry delegator ownership via a MessageService accepted-delivery hook (won the session lock + queued delivery), never at bridge_send request time, so a concurrent sender that times out BUSY cannot steal a live turn's reply routing. Make Rendezvous.open atomic fail-if-present so a double open trips loudly instead of replacing the waiter another send is blocked on. Answering a bridge_ask keeps the same ownership (no rewrite). Adds ownership/rendezvous regression tests. --- .../java/dev/ltms/bridged/mcp/BridgeMcp.java | 43 +++++-- .../dev/ltms/bridged/mcp/PrimaryRegistry.java | 12 +- .../dev/ltms/bridged/msg/MessageService.java | 36 +++++- .../java/dev/ltms/bridged/msg/Rendezvous.java | 21 +++- .../inject/CompletionResolverTest.java | 6 +- .../ltms/bridged/msg/MessageServiceTest.java | 114 ++++++++++++++++++ .../dev/ltms/bridged/msg/RendezvousTest.java | 22 ++++ 7 files changed, 233 insertions(+), 21 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 3056344..fb45f93 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java @@ -129,19 +129,25 @@ public final class BridgeMcp { String caller = callerTerminal(exchange); if (caller != null) primaryRegistry.record(caller); 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. @@ -343,12 +349,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 +433,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..de47898 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,21 @@ 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, 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 + * 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()); @@ -318,6 +333,11 @@ public final class MessageService { } 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(); + } CompletableFuture reply = rendezvous.open(target); try { Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS); @@ -440,9 +460,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/msg/MessageServiceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java index 66d11f2..46c7acd 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; @@ -321,6 +322,119 @@ 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()); + } + + /** + * 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);