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.
This commit is contained in:
Dai Ha
2026-08-13 18:55:54 +02:00
parent 86cf4c285a
commit ec3001796a
7 changed files with 233 additions and 21 deletions
@@ -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
@@ -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} <em>request</em> 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<MessageService.Reply> 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 <em>accepted</em> 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<MessageService.Reply> 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<MessageService.Reply> 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<MessageService.Reply> 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<MessageService.AskResult> 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<MessageService.Reply> 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
@@ -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<Rendezvous.Resolution> 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);