From 4aed45de19f44190d298de5b8312ca7d7438acee Mon Sep 17 00:00:00 2001 From: Kevin Nguyen Date: Sat, 1 Aug 2026 23:24:54 +0700 Subject: [PATCH] CB-514: add MessageService coverage for timeout, answer, poll, and lock edges --- .../ltms/bridged/msg/MessageServiceTest.java | 101 ++++++++++++++++++ 1 file changed, 101 insertions(+) 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 7c1ce3e..2f4bb75 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -14,6 +14,7 @@ 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.assertTrue; /** @@ -212,6 +213,106 @@ class MessageServiceTest { "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 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 send = sendAsync(); + awaitWaiting(); + injector.onStatus(T, AgentStatus.IDLE); + injector.onStatus(T, AgentStatus.WORKING); + + CompletableFuture 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 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-307 reply inbox ---------------------------------------------------------------- @Test