diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java index 99f8f80..fb5b7f5 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -1846,7 +1846,8 @@ class MessageServiceTest { * but backed by {@link ManualScheduler} instead of a real timer (fleetd #608): its tick never * fires on its own — a test drives it explicitly via {@link ManualScheduler#runDueTasks()}. */ - private record ManualPushWiring(MessageService service, FakeHerdr leadHerdr, ManualScheduler scheduler) + private record ManualPushWiring(MessageService service, FakeHerdr leadHerdr, ManualScheduler scheduler, + ReplyPushLoop pushLoop) implements AutoCloseable { @Override public void close() { @@ -1861,14 +1862,20 @@ class MessageServiceTest { * {@code anAlreadyCollectedTicketProducesNoNudge}, which sets it to 1 to prove exactly that. */ private ManualPushWiring wireWithManualScheduler(int maxReminders, long backoffMs) { + return wireWithManualScheduler(maxReminders, backoffMs, System::nanoTime); + } + + /** As above, with an injectable clock for tests that exercise terminal-ticket pruning. */ + private ManualPushWiring wireWithManualScheduler(int maxReminders, long backoffMs, + java.util.function.LongSupplier nowNanos) { PrimaryRegistry registry = new PrimaryRegistry(null); registry.recordDelegation(T, LEAD); FakeHerdr leadHerdr = new FakeHerdr(); AgentControl leadAgents = new AgentControl(leadHerdr); ManualScheduler scheduler = new ManualScheduler(); ReplyPushLoop pushLoop = new ReplyPushLoop(registry, leadAgents, inbox, scheduler, maxReminders, backoffMs); - MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop, null, System::nanoTime); - return new ManualPushWiring(service, leadHerdr, scheduler); + MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop, null, nowNanos); + return new ManualPushWiring(service, leadHerdr, scheduler, pushLoop); } private void awaitNudge(FakeHerdr leadHerdr) throws InterruptedException { @@ -1959,24 +1966,35 @@ class MessageServiceTest { @Test void severalAsyncTicketsFinishingTogetherProduceOneCoalescedNudge() throws Exception { - try (var wiring = wireWithPushLoop(1, 300)) { // wide backoff: both tickets land before the tick fires - String first = wiring.service().sendAsync(T, "first task"); - awaitWaiting(); - injectDelivery(); - assertTrue(rendezvous.resolve(T, "first done")); - // Settle without polling: poll() itself marks a ticket collected (that's the point of - // anAlreadyCollectedTicketProducesNoNudge above) — using it here to detect completion - // would collect the ticket before the coalescing this test checks ever gets a chance. - Thread.sleep(100); + // A 1ms backoff is due immediately. ManualScheduler still cannot run it until this test + // explicitly calls runDueTasks(), so both terminal tickets join one scheduled tick. + try (var wiring = wireWithManualScheduler(1, 1)) { + String first; + String second; + java.util.concurrent.CountDownLatch firstTerminalReached = new java.util.concurrent.CountDownLatch(1); + wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(firstTerminalReached::countDown); + try { + first = wiring.service().sendAsync(T, "first task"); + awaitWaiting(); + injectDelivery(); + assertTrue(rendezvous.resolve(T, "first done")); + assertTrue(firstTerminalReached.await(5, TimeUnit.SECONDS), + "the first ticket never reached its terminal phase"); - String second = wiring.service().sendAsync(T, "second task"); - awaitWaiting(); - injectDelivery(); - assertTrue(rendezvous.resolve(T, "second done")); - Thread.sleep(100); + java.util.concurrent.CountDownLatch secondTerminalReached = new java.util.concurrent.CountDownLatch(1); + wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(secondTerminalReached::countDown); + second = wiring.service().sendAsync(T, "second task"); + awaitWaiting(); + injectDelivery(); + assertTrue(rendezvous.resolve(T, "second done")); + assertTrue(secondTerminalReached.await(5, TimeUnit.SECONDS), + "the second ticket never reached its terminal phase"); + } finally { + wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(null); + } - awaitNudge(wiring.leadHerdr()); - Thread.sleep(200); // settle — nothing more should arrive beyond the one coalesced nudge + assertEquals(1, wiring.scheduler().runDueTasks(), + "both terminal tickets must coalesce onto one scheduled tick"); long nudgeCount = wiring.leadHerdr().calls.stream() .filter(c -> c.method().equals("agent.prompt")).count(); assertEquals(1, nudgeCount, "two tickets finishing together must produce ONE nudge, not two"); @@ -2055,7 +2073,8 @@ class MessageServiceTest { @Test void answeringAQuestionStopsFurtherNudgesAboutIt() throws Exception { - try (var wiring = wireWithPushLoop(5, 50)) { + // A 1ms backoff is due immediately, but ManualScheduler only ticks when this test asks it to. + try (var wiring = wireWithManualScheduler(5, 1)) { String ticket = wiring.service().sendAsync(T, "task that asks"); awaitWaiting(); injectDelivery(); @@ -2063,8 +2082,9 @@ class MessageServiceTest { CompletableFuture ask = CompletableFuture.supplyAsync( () -> wiring.service().ask(T, "which config file?", 5000)); MessageService.TaskView asking = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.ASKING); + awaitQuestionPendingOn(wiring.pushLoop(), LEAD, asking.turnId()); - awaitNudge(wiring.leadHerdr()); + assertEquals(1, wiring.scheduler().runDueTasks(), "the open question must have one scheduled tick"); long callsBeforeAnswer = wiring.leadHerdr().calls.stream() .filter(c -> c.method().equals("agent.prompt")).count(); @@ -2072,12 +2092,22 @@ class MessageServiceTest { () -> wiring.service().answer(asking.turnId(), "config.yaml", 5000)); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); awaitWaiting(); + java.util.concurrent.CountDownLatch terminalReached = new java.util.concurrent.CountDownLatch(1); + wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(terminalReached::countDown); assertTrue(rendezvous.resolve(T, "done")); + try { + assertTrue(terminalReached.await(5, TimeUnit.SECONDS), + "the answered ticket never reached its terminal phase"); + } finally { + wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(null); + } answer.get(5, TimeUnit.SECONDS); - // Let several more ticks (and the ticket's own now-legitimate terminal nudge) fire — - // none of them may still name the question's turnId, which is closed. - Thread.sleep(300); + // Run the ticket's legitimate terminal nudge and every later scheduled tick through its + // reminder cap. None may still name the closed question. + for (int tick = 0; tick < 6; tick++) { + assertEquals(1, wiring.scheduler().runDueTasks(), "expected one scheduled reminder tick"); + } boolean anyNamesClosedQuestion = wiring.leadHerdr().calls.stream() .filter(c -> c.method().equals("agent.prompt")) .skip(callsBeforeAnswer) @@ -2160,35 +2190,45 @@ class MessageServiceTest { // decideTickets hits the cap and STOPs — activeLeads drops the lead, but (before the fix) // pendingTickets never drops the ticket. That is the exact "cap already STOPped" branch of // the bug report, reached deterministically rather than by timing it against a live tick. - try (var wiring = wireWithPushLoop(1, 50, clock::get)) { - String stale = wiring.service().sendAsync(T, "first task"); - awaitWaiting(); - injectDelivery(); - assertTrue(rendezvous.resolve(T, "stale result")); + // A 1ms backoff is due immediately, but ManualScheduler runs only the ticks below. + try (var wiring = wireWithManualScheduler(1, 1, clock::get)) { + String stale; + String fresh; + java.util.concurrent.CountDownLatch staleTerminalReached = new java.util.concurrent.CountDownLatch(1); + wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(staleTerminalReached::countDown); + try { + stale = wiring.service().sendAsync(T, "first task"); + awaitWaiting(); + injectDelivery(); + assertTrue(rendezvous.resolve(T, "stale result")); + assertTrue(staleTerminalReached.await(5, TimeUnit.SECONDS), + "the stale ticket never reached its terminal phase"); - // Let the reminder loop fire its one nudge and hit the cap (STOP removes it from - // activeLeads; pendingTickets is untouched either way — that asymmetry is the bug). - awaitNudge(wiring.leadHerdr()); - Thread.sleep(300); - assertTrue(wiring.leadHerdr().lastCall("agent.prompt").params().toString().contains(stale), - "sanity: the stale ticket's own reminder must have fired first"); + // Fire the stale ticket's one nudge, then its cap tick (STOP removes it from activeLeads; + // pendingTickets is untouched either way — that asymmetry is the bug). + assertEquals(1, wiring.scheduler().runDueTasks(), "the stale ticket must have one nudge tick"); + assertEquals(1, wiring.scheduler().runDueTasks(), "the stale ticket must have one cap tick"); + assertTrue(wiring.leadHerdr().lastCall("agent.prompt").params().toString().contains(stale), + "sanity: the stale ticket's own reminder must have fired first"); - // Cross the TTL — a real clock would need 10 minutes; the injected one does it instantly. - clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1)); + // Cross the TTL — a real clock would need 10 minutes; the injected one does it instantly. + clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1)); - // A second, unrelated ticket to the same target/lead reaches sendAsync, which prunes. - String fresh = wiring.service().sendAsync(T, "second task"); - awaitWaiting(); - injectDelivery(); - assertTrue(rendezvous.resolve(T, "fresh result")); - - // The fresh ticket restarts the (now-dormant) reminder loop with its own nudge. - long before = wiring.leadHerdr().calls.stream().filter(c -> c.method().equals("agent.prompt")).count(); - long deadline = System.currentTimeMillis() + 3000; - while (wiring.leadHerdr().calls.stream().filter(c -> c.method().equals("agent.prompt")).count() <= before - && System.currentTimeMillis() < deadline) { - Thread.sleep(10); + // A second, unrelated ticket to the same target/lead reaches sendAsync, which prunes. + java.util.concurrent.CountDownLatch freshTerminalReached = new java.util.concurrent.CountDownLatch(1); + wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(freshTerminalReached::countDown); + fresh = wiring.service().sendAsync(T, "second task"); + awaitWaiting(); + injectDelivery(); + assertTrue(rendezvous.resolve(T, "fresh result")); + assertTrue(freshTerminalReached.await(5, TimeUnit.SECONDS), + "the fresh ticket never reached its terminal phase"); + } finally { + wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(null); } + + // The fresh ticket restarts the now-dormant reminder loop with its own nudge. + assertEquals(1, wiring.scheduler().runDueTasks(), "the fresh ticket must have one nudge tick"); String latestNudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString(); assertTrue(latestNudge.contains(fresh), "the fresh ticket's nudge must still arrive: " + latestNudge); assertFalse(latestNudge.contains(stale),