dd906526c0
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.
613 lines
30 KiB
Java
613 lines
30 KiB
Java
package dev.ltms.bridged.msg;
|
|
|
|
import dev.ltms.bridged.herdr.AgentControl;
|
|
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;
|
|
|
|
import java.util.concurrent.CompletableFuture;
|
|
import java.util.concurrent.TimeUnit;
|
|
|
|
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;
|
|
|
|
/**
|
|
* The message layer's resolution paths (CB-104 reply + CB-106 completion fallback). The turn is
|
|
* driven deterministically by feeding {@code onStatus} rather than running a real poller.
|
|
*/
|
|
class MessageServiceTest {
|
|
|
|
private static final String T = "term_a";
|
|
|
|
private final FakeHerdr herdr = new FakeHerdr().readText("BUILD GREEN: 391 files");
|
|
private final AgentControl agents = new AgentControl(herdr);
|
|
private final Rendezvous rendezvous = new Rendezvous();
|
|
private final CompletionResolver completion = new CompletionResolver(agents, rendezvous);
|
|
private final Injector injector = new Injector(agents, completion);
|
|
private final InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
|
private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
|
|
|
|
@BeforeEach
|
|
void setUp() {
|
|
// CB-520: the inbox only peeks/acks targets it owns.
|
|
inbox.own(T);
|
|
}
|
|
|
|
/** Run {@code send} on a background thread; the current thread drives the worker's turn. */
|
|
private CompletableFuture<MessageService.Reply> sendAsync() {
|
|
return CompletableFuture.supplyAsync(() -> messages.send(T, "do the task", 5000));
|
|
}
|
|
|
|
private void awaitWaiting() throws InterruptedException {
|
|
long deadline = System.currentTimeMillis() + 2000;
|
|
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
|
//noinspection BusyWait
|
|
Thread.sleep(5);
|
|
}
|
|
assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter");
|
|
}
|
|
|
|
@Test
|
|
void completionFallbackResolvesATurnThatNeverCalledBridgeReply() throws Exception {
|
|
CompletableFuture<MessageService.Reply> send = sendAsync();
|
|
awaitWaiting();
|
|
|
|
herdr.readText("$ prompt"); // pre-turn pane: no answer yet (baseline reference)
|
|
injector.onStatus(T, AgentStatus.IDLE); // deliver the task (baselines the pre-turn content)
|
|
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up and works
|
|
herdr.readText("BUILD GREEN: 391 files"); // the worker's turn produced new output
|
|
injector.onStatus(T, AgentStatus.IDLE); // working → idle: turn complete, no bridge_reply
|
|
|
|
MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);
|
|
assertEquals(MessageService.Outcome.COMPLETED_UNREPLIED, reply.outcome(),
|
|
"an unreplied but finished turn resolves via the completion fallback");
|
|
assertEquals("BUILD GREEN: 391 files", reply.text(), "the scraped transcript tail is returned");
|
|
assertTrue(reply.completed(), "a scraped completion still counts as completed");
|
|
}
|
|
|
|
@Test
|
|
void explicitBridgeReplyResolvesAsReplied() throws Exception {
|
|
CompletableFuture<MessageService.Reply> send = sendAsync();
|
|
awaitWaiting();
|
|
|
|
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
|
injector.onStatus(T, AgentStatus.WORKING); // worker working
|
|
assertTrue(rendezvous.resolve(T, "LGTM ship it"), "an explicit reply resolves the send");
|
|
|
|
MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);
|
|
assertEquals(MessageService.Outcome.REPLIED, reply.outcome());
|
|
assertEquals("LGTM ship it", reply.text());
|
|
}
|
|
|
|
@Test
|
|
void aWedgedWorkerResolvesTheSendAsFailedWithTheErrorContext() throws Exception {
|
|
herdr.readText("API Error: Unable to connect to API (ENOTFOUND)");
|
|
CompletableFuture<MessageService.Reply> send = sendAsync();
|
|
awaitWaiting();
|
|
|
|
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
|
injector.onStatus(T, AgentStatus.WORKING); // worker starts the turn
|
|
for (int i = 0; i < 130; i++) injector.onStatus(T, AgentStatus.UNKNOWN); // then wedges (CB-109)
|
|
|
|
MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);
|
|
assertEquals(MessageService.Outcome.WORKER_FAILED, reply.outcome());
|
|
assertFalse(reply.completed(), "a wedge is terminal but not a successful completion");
|
|
assertTrue(reply.text().contains("ENOTFOUND"), "the error screen is carried as the failure reason");
|
|
}
|
|
|
|
@Test
|
|
void aWorkerThatVanishesMidTurnResolvesTheSendAsFailed() throws Exception {
|
|
CompletableFuture<MessageService.Reply> send = sendAsync();
|
|
awaitWaiting();
|
|
|
|
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
|
injector.onStatus(T, AgentStatus.WORKING); // worker starts the turn
|
|
// The worker's pane crashes — the poller sees a *_not_found and drops it (CB-110).
|
|
injector.drop(T, new HerdrException("worker gone", "pane_not_found", null));
|
|
|
|
MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);
|
|
assertEquals(MessageService.Outcome.WORKER_FAILED, reply.outcome(),
|
|
"a delivered send whose worker vanishes fails instead of hanging to the timeout");
|
|
assertFalse(reply.completed());
|
|
}
|
|
|
|
// --- bridge_ask reverse rendezvous (CB-205) ------------------------------------------------
|
|
|
|
@Test
|
|
void askSurfacesAsAQuestionAndTheAnswerResumesTheSameTurn() throws Exception {
|
|
CompletableFuture<MessageService.Reply> send = sendAsync();
|
|
awaitWaiting();
|
|
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
|
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up, then pauses to ask
|
|
|
|
// The worker asks mid-turn on its own thread; the call blocks for the primary's answer.
|
|
CompletableFuture<MessageService.AskResult> ask =
|
|
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
|
|
|
|
// The primary's blocking send unblocks with the question and a turnId to answer on.
|
|
MessageService.Reply q = send.get(5, TimeUnit.SECONDS);
|
|
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
|
|
assertEquals("which config file?", q.text());
|
|
assertNotNull(q.turnId(), "a question carries a turnId to answer on");
|
|
|
|
// The primary answers via bridge_send(turnId); this blocks again for the worker's reply.
|
|
CompletableFuture<MessageService.Reply> answer =
|
|
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
|
|
|
|
// The worker's ask returns the answer — it resumes the same turn.
|
|
MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS);
|
|
assertEquals(MessageService.AskOutcome.ANSWERED, a.outcome());
|
|
assertEquals("config.yaml", a.answer());
|
|
|
|
// The resumed worker finishes with a structured reply, resolving the answering send.
|
|
awaitWaiting(); // the answering send has (re)opened its forward waiter
|
|
assertTrue(rendezvous.resolve(T, "done"), "the worker's final reply resolves the answering send");
|
|
MessageService.Reply done = answer.get(5, TimeUnit.SECONDS);
|
|
assertEquals(MessageService.Outcome.REPLIED, done.outcome());
|
|
assertEquals("done", done.text());
|
|
}
|
|
|
|
@Test
|
|
void duplicateAsksFromTheSameSessionCoalesceToOneTurn() throws Exception {
|
|
CompletableFuture<MessageService.Reply> send = sendAsync();
|
|
awaitWaiting();
|
|
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
|
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up, then pauses to ask
|
|
|
|
// A transport retry: two concurrent bridge_ask calls from the same worker session.
|
|
CompletableFuture<MessageService.AskResult> ask1 =
|
|
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
|
|
CompletableFuture<MessageService.AskResult> ask2 =
|
|
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
|
|
|
|
// The primary's single blocked send surfaces exactly ONE question (one turnId).
|
|
MessageService.Reply q = send.get(5, TimeUnit.SECONDS);
|
|
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
|
|
assertEquals("which config file?", q.text());
|
|
assertNotNull(q.turnId(), "only one turnId should be minted");
|
|
|
|
// The primary answers that one turnId; both asks unblock with the same answer.
|
|
CompletableFuture<MessageService.Reply> answer =
|
|
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
|
|
|
|
MessageService.AskResult a1 = ask1.get(5, TimeUnit.SECONDS);
|
|
MessageService.AskResult a2 = ask2.get(5, TimeUnit.SECONDS);
|
|
assertEquals(MessageService.AskOutcome.ANSWERED, a1.outcome());
|
|
assertEquals("config.yaml", a1.answer());
|
|
assertEquals(MessageService.AskOutcome.ANSWERED, a2.outcome());
|
|
assertEquals("config.yaml", a2.answer());
|
|
|
|
// The resumed worker finishes with a structured reply, resolving the answering send.
|
|
awaitWaiting();
|
|
assertTrue(rendezvous.resolve(T, "done"), "the worker's final reply resolves the answering send");
|
|
MessageService.Reply done = answer.get(5, TimeUnit.SECONDS);
|
|
assertEquals(MessageService.Outcome.REPLIED, done.outcome());
|
|
assertEquals("done", done.text());
|
|
}
|
|
|
|
@Test
|
|
void askWithNoOpenDelegationReturnsNoWaiter() {
|
|
MessageService.AskResult r = messages.ask(T, "anyone listening?", 500);
|
|
assertEquals(MessageService.AskOutcome.NO_WAITER, r.outcome(),
|
|
"a question with no blocked send has no primary to answer it");
|
|
}
|
|
|
|
@Test
|
|
void askTimesOutWhenThePrimaryNeverAnswers() throws Exception {
|
|
CompletableFuture<MessageService.Reply> send = sendAsync();
|
|
awaitWaiting();
|
|
injector.onStatus(T, AgentStatus.IDLE);
|
|
injector.onStatus(T, AgentStatus.WORKING);
|
|
|
|
MessageService.AskResult r = messages.ask(T, "still there?", 200); // primary never answers
|
|
assertEquals(MessageService.AskOutcome.TIMED_OUT, r.outcome());
|
|
|
|
// The send itself already unblocked with the question the instant the ask surfaced.
|
|
MessageService.Reply q = send.get(2, TimeUnit.SECONDS);
|
|
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
|
|
}
|
|
|
|
@Test
|
|
void answeringAnUnknownTurnIsStale() {
|
|
MessageService.Reply r = messages.answer(T + "#999", "too late", 500);
|
|
assertEquals(MessageService.Outcome.STALE_TURN, r.outcome(),
|
|
"an answer to a turn that never existed (or already lapsed) is stale, not a hang");
|
|
}
|
|
|
|
// --- timeout, answer, poll, and lock-contention edges ----------------------------------
|
|
|
|
@Test
|
|
void sendTimesOutBeforeDeliveryIsQueuedNotWorking() {
|
|
// Nothing ever delivers the message and nothing resolves the send, so the reply future
|
|
// times out with delivery still incomplete — the message is still queued for the worker.
|
|
MessageService.Reply r = messages.send(T, "never delivered", 50);
|
|
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, r.outcome(),
|
|
"an undelivered send that times out is still queued, not working");
|
|
assertNull(r.text());
|
|
}
|
|
|
|
@Test
|
|
void sendTimesOutAfterDeliveryIsStillWorking() throws Exception {
|
|
CompletableFuture<MessageService.Reply> send =
|
|
CompletableFuture.supplyAsync(() -> messages.send(T, "do the task", 300));
|
|
awaitWaiting();
|
|
injector.onStatus(T, AgentStatus.IDLE); // deliver — the delivered future now completes
|
|
injector.onStatus(T, AgentStatus.WORKING); // worker starts but never replies
|
|
// No rendezvous.resolve(T, ...) — the reply future rides out its short timeout.
|
|
|
|
MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
|
|
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, r.outcome(),
|
|
"a delivered send whose worker never replies times out as still working");
|
|
}
|
|
|
|
@Test
|
|
void answerTimesOutWhenTheResumedWorkerNeverReplies() throws Exception {
|
|
CompletableFuture<MessageService.Reply> send = sendAsync();
|
|
awaitWaiting();
|
|
injector.onStatus(T, AgentStatus.IDLE);
|
|
injector.onStatus(T, AgentStatus.WORKING);
|
|
|
|
CompletableFuture<MessageService.AskResult> ask =
|
|
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
|
|
|
|
MessageService.Reply q = send.get(5, TimeUnit.SECONDS);
|
|
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
|
|
assertNotNull(q.turnId());
|
|
|
|
// The primary answers, unblocking the worker; but the worker never sends the follow-up
|
|
// bridge_reply, so the answering send rides out its short window as still-working.
|
|
MessageService.Reply answer = messages.answer(q.turnId(), "config.yaml", 200);
|
|
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answer.outcome(),
|
|
"an answered worker that never replies times out as still working");
|
|
|
|
MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS);
|
|
assertEquals(MessageService.AskOutcome.ANSWERED, a.outcome());
|
|
assertEquals("config.yaml", a.answer());
|
|
}
|
|
|
|
@Test
|
|
void pollReturnsNullForAnUnknownTicket() {
|
|
assertNull(messages.poll("task-999999"), "a ticket that was never minted is unknown");
|
|
}
|
|
|
|
@Test
|
|
void pollReportsACompletedTicket() throws Exception {
|
|
String ticket = messages.sendAsync(T, "long task");
|
|
awaitWaiting();
|
|
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
|
injector.onStatus(T, AgentStatus.WORKING); // worker works
|
|
assertTrue(rendezvous.resolve(T, "async result"), "a reply resolves the async send");
|
|
|
|
// Wait for the background send to finish and publish a DONE view.
|
|
MessageService.TaskView view = null;
|
|
long deadline = System.currentTimeMillis() + 2000;
|
|
while (view == null || view.phase() != MessageService.Phase.DONE) {
|
|
if (System.currentTimeMillis() >= deadline) break;
|
|
view = messages.poll(ticket);
|
|
//noinspection BusyWait
|
|
Thread.sleep(5);
|
|
}
|
|
assertNotNull(view, "a resolved async send must become DONE");
|
|
assertEquals(MessageService.Phase.DONE, view.phase());
|
|
assertEquals("async result", view.reply(), "the completed ticket reports the reply");
|
|
assertEquals("reply", view.replySource(), "a structured bridge_reply is sourced from 'reply'");
|
|
}
|
|
|
|
@Test
|
|
void concurrentSendToSameSessionWhileFirstHoldsItIsBusy() throws Exception {
|
|
CompletableFuture<MessageService.Reply> first =
|
|
CompletableFuture.supplyAsync(() -> messages.send(T, "first", 5000));
|
|
awaitWaiting(); // the first send now holds the session lock, blocked on its reply
|
|
|
|
// A second send to the SAME session cannot take the lock within its short window.
|
|
MessageService.Reply busy = messages.send(T, "second", 100);
|
|
assertEquals(MessageService.Outcome.BUSY, busy.outcome(),
|
|
"a second send while another holds the session is busy, not a hang");
|
|
assertNull(busy.text());
|
|
|
|
// Release the first send so it resolves cleanly and the test thread is not left pinned.
|
|
injector.onStatus(T, AgentStatus.IDLE); // deliver the first message
|
|
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up
|
|
assertTrue(rendezvous.resolve(T, "first done"), "the first send resolves with a reply");
|
|
MessageService.Reply firstReply = first.get(5, TimeUnit.SECONDS);
|
|
assertEquals(MessageService.Outcome.REPLIED, firstReply.outcome());
|
|
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
|
|
void replyQueuesInInboxWhenNoSendIsOpen() {
|
|
// No send is open for this session — reply should queue in the inbox.
|
|
assertTrue(messages.reply(T, "queued-text"), "reply should succeed (queued)");
|
|
|
|
var drained = messages.drainReplies(T);
|
|
assertEquals(1, drained.size());
|
|
assertEquals("queued-text", drained.getFirst().content());
|
|
}
|
|
|
|
@Test
|
|
void replyResolvesOpenSendDoesNotQueue() throws Exception {
|
|
CompletableFuture<MessageService.Reply> send = sendAsync();
|
|
awaitUninterruptibly(T);
|
|
|
|
// An explicit reply resolves the open send.
|
|
assertTrue(messages.reply(T, "send-resolved"), "reply should succeed (resolved live send)");
|
|
|
|
// The inbox should be empty — the reply went to the send, not the inbox.
|
|
assertTrue(messages.drainReplies(T).isEmpty(), "no reply in the inbox");
|
|
|
|
MessageService.Reply r = send.get(3, TimeUnit.SECONDS);
|
|
assertEquals(MessageService.Outcome.REPLIED, r.outcome());
|
|
assertEquals("send-resolved", r.text());
|
|
}
|
|
|
|
@Test
|
|
void drainRepliesReturnsAllPendingThenEmptyOnNextCall() {
|
|
messages.reply(T, "msg-1");
|
|
messages.reply(T, "msg-2");
|
|
|
|
var first = messages.drainReplies(T);
|
|
assertEquals(2, first.size());
|
|
|
|
var second = messages.drainReplies(T);
|
|
assertTrue(second.isEmpty(), "second drain should be empty (acked)");
|
|
}
|
|
|
|
@Test
|
|
void aQuestionIsNeverQueuedInTheInbox() {
|
|
// No send is open — bridge_ask with no delegation returns NO_WAITER,
|
|
// and the question text MUST NOT appear in the reply inbox.
|
|
// The inbox is only fed by MessageService.reply(), not by bridge_ask.
|
|
MessageService.AskResult r = messages.ask(T, "anyone there?", 500);
|
|
assertEquals(MessageService.AskOutcome.NO_WAITER, r.outcome(),
|
|
"bridge_ask with no open delegation must return NO_WAITER, never queued");
|
|
|
|
assertTrue(messages.drainReplies(T).isEmpty(), "questions must never be queued");
|
|
}
|
|
|
|
@Test
|
|
void completionFallbackIsNeverQueued() throws Exception {
|
|
// The fallback resolves a captured waiter, never the inbox.
|
|
CompletableFuture<MessageService.Reply> send = sendAsync();
|
|
awaitUninterruptibly(T);
|
|
injectDelivery();
|
|
|
|
// The worker never sends bridge_reply, but the turn completes.
|
|
herdr.readText("done-scraped");
|
|
completion.onTurnComplete(T); // The fallback arms and resolves the captured waiter.
|
|
|
|
MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
|
|
assertEquals(MessageService.Outcome.COMPLETED_UNREPLIED, r.outcome());
|
|
|
|
// The inbox should be empty — the reply went to the captured waiter.
|
|
assertTrue(messages.drainReplies(T).isEmpty(), "completion fallback must not queue");
|
|
}
|
|
|
|
// --- helpers ---------------------------------------------------------------------------
|
|
|
|
/** Like {@link #awaitWaiting()} but rethrows as unchecked. */
|
|
private void awaitUninterruptibly(String session) {
|
|
try {
|
|
awaitWaiting();
|
|
} catch (InterruptedException e) {
|
|
Thread.currentThread().interrupt();
|
|
throw new IllegalStateException(e);
|
|
}
|
|
}
|
|
|
|
/** Set up a delivered turn so the worker is working, ready for an ask or completion. */
|
|
private void injectDelivery() {
|
|
herdr.readText("$ prompt"); // pre-turn content baseline
|
|
injector.onStatus(T, AgentStatus.IDLE); // deliver the task
|
|
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up
|
|
}
|
|
|
|
// --- CB-516: a released session must not leave a send hanging ------------------------------
|
|
|
|
/**
|
|
* The bug this fixes: tearing a worker down left its rendezvous waiter open, so a blocking send
|
|
* kept blocking and an async one kept reporting PENDING until the 30-minute async timeout —
|
|
* even though the worker provably no longer existed.
|
|
*/
|
|
@Test
|
|
void abandonFailsASendThatIsStillWaitingOnAReleasedSession() throws Exception {
|
|
CompletableFuture<MessageService.Reply> send =
|
|
CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000));
|
|
awaitWaiting();
|
|
|
|
assertTrue(messages.abandon(T, "session released"), "a live waiter is abandoned");
|
|
|
|
MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
|
|
assertEquals(MessageService.Outcome.WORKER_FAILED, r.outcome(),
|
|
"an abandoned send fails rather than riding out its timeout");
|
|
assertEquals("session released", r.text(), "the caller is told why");
|
|
}
|
|
|
|
@Test
|
|
void abandonIsANoOpWhenNobodyIsWaiting() {
|
|
assertFalse(messages.abandon(T, "session released"),
|
|
"no open send ⇒ nothing to abandon");
|
|
}
|
|
|
|
@Test
|
|
void abandonDoesNotOverwriteAnAlreadyResolvedSend() throws Exception {
|
|
CompletableFuture<MessageService.Reply> send =
|
|
CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000));
|
|
awaitWaiting();
|
|
assertTrue(rendezvous.resolve(T, "the real answer"));
|
|
|
|
assertFalse(messages.abandon(T, "session released"),
|
|
"a send already answered by the worker must not be clobbered");
|
|
MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
|
|
assertEquals("the real answer", r.text());
|
|
}
|
|
|
|
/** The async path is the one that hung: poll must report FAILED, not PENDING forever. */
|
|
@Test
|
|
void anAbandonedAsyncTaskPollsAsFailedNotPending() throws Exception {
|
|
String ticket = messages.sendAsync(T, "long task");
|
|
awaitWaiting();
|
|
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase());
|
|
|
|
messages.abandon(T, "session released");
|
|
|
|
MessageService.TaskView view = null;
|
|
long deadline = System.currentTimeMillis() + 3000;
|
|
while (System.currentTimeMillis() < deadline) {
|
|
view = messages.poll(ticket);
|
|
if (view.phase() != MessageService.Phase.PENDING) break;
|
|
Thread.sleep(10);
|
|
}
|
|
assertNotNull(view);
|
|
assertEquals(MessageService.Phase.FAILED, view.phase(),
|
|
"a delegation whose worker is gone must not keep reporting PENDING");
|
|
assertTrue(view.detail() != null && view.detail().contains("released"),
|
|
"and the detail says why, rather than 'worker unknown'");
|
|
}
|
|
}
|