diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java index 9926aee..dd8106b 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -536,6 +536,14 @@ public final class MessageService { // Only the fresh owner tears down the shared turn; a duplicate must leave it open. if (ticket.fresh()) { rendezvous.closeAsk(ticket.turnId()); + // CB-582: tear the push loop's copy down at the same point, not only on the three + // paths that call clearAsyncQuestion. The answer future can complete exceptionally + // (ExecutionException) or the thread be interrupted, and both leave this method by + // throwing — the question would stay pending forever, keep being named in nudges + // until its own cap, and never be removed from the map. Already-closed is a no-op. + if (pushLoop != null) { + pushLoop.questionClosed(ticket.turnId()); + } } } } 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 eb2563e..4790377 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -939,6 +939,45 @@ class MessageServiceTest { } } + @Test + void anAskThatLeavesByThrowingStillClosesItsQuestion() throws Exception { + // CB-582 follow-up. ask() calls clearAsyncQuestion on three paths — no-waiter, timed out, + // and (from answer()) answered — but it can also leave by *throwing*: an interrupt while + // blocked on the answer, or an ExecutionException from the answer future. Those paths run + // only the finally block, so before the fix the push loop kept the question pending for + // good: named in every nudge until its own cap, then never removed from the map at all. + PrimaryRegistry registry = new PrimaryRegistry(null); + registry.recordDelegation(T, LEAD); + var scheduler = java.util.concurrent.Executors.newSingleThreadScheduledExecutor(); + // A backoff far longer than the test: the schedule is started but no tick ever fires, so + // decide() is read directly and nothing here depends on timing. + ReplyPushLoop pushLoop = new ReplyPushLoop(registry, new AgentControl(new FakeHerdr()), inbox, + scheduler, 5, 60_000); + MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop); + try { + String ticket = service.sendAsync(T, "task that asks"); + awaitWaiting(); + injectDelivery(); + + Thread asker = new Thread(() -> assertThrows(IllegalStateException.class, + () -> service.ask(T, "which config file?", 30_000))); + asker.start(); + awaitTicketPhaseOn(service, ticket, MessageService.Phase.ASKING); + assertEquals(ReplyPushLoop.Action.INJECT, pushLoop.decide(LEAD, 0, 0, 0), + "the open question should be the one thing keeping this lead's schedule alive"); + + asker.interrupt(); + asker.join(5000); + assertFalse(asker.isAlive(), "the interrupted ask should have left ask() by throwing"); + + assertEquals(ReplyPushLoop.Action.STOP, pushLoop.decide(LEAD, 0, 0, 0), + "an ask that threw must still close its question, or the loop nudges about it for good"); + } finally { + service.close(); + scheduler.shutdownNow(); + } + } + @Test void aFleetWithNoPushLoopConfiguredBehavesExactlyAsToday() throws Exception { // `messages` (the shared field) uses the no-pushLoop constructor — poll() must not throw,