CB-548: harden send ownership and rendezvous routing #23

Merged
ltms merged 2 commits from worker/cb-548-rendezvous-guard-rebased into main 2026-08-13 19:44:02 +02:00
8 changed files with 339 additions and 37 deletions
@@ -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<String, Object> 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<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).
@@ -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).
*
* <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);
@@ -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);
}
@@ -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).
*
* <p>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.
* <p>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 <em>not</em> 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()) {
@@ -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.
*
* <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.
*/
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<Void> 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<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);
}
@@ -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<Reply> future =
CompletableFuture.supplyAsync(() -> send(target, content, ASYNC_TIMEOUT_MS), asyncExecutor);
CompletableFuture<Reply> 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);
@@ -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.
*
* <p>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<Resolution> open(String session) {
CompletableFuture<Resolution> waiter = new CompletableFuture<>();
waiters.put(session, waiter);
CompletableFuture<Resolution> 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<Resolution> 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<Resolution> waiter) {
waiters.remove(session, waiter);
}
@@ -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
@@ -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");
}
}
@@ -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} <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());
}
/**
* 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
@@ -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);