diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java b/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java index a920f4a..d30dc39 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java @@ -230,6 +230,26 @@ public final class ReplyPushLoop { .collect(Collectors.toUnmodifiableSet()); } + /** + * Test seam only (fleetd #418): carries no production behaviour, and nothing in this class + * calls it. Exposes {@link #pendingQuestionTurnIdsFor} — the exact state {@link #decide} reads + * to decide whether a question keeps a lead's schedule alive. + * + *

{@code MessageService.ask()} does three things in order before a question is fully open to + * this loop: it flips the ticket's {@code poll()} phase to {@code Phase.ASKING}, then resolves + * the reverse-rendezvous waiter, then calls {@link #onQuestionOpened}, which is what actually + * populates {@link #pendingQuestions}. A test that barriers on {@code Phase.ASKING} observes only + * the first of those three steps — under load the asker thread can be descheduled between steps + * one and three, so the barrier releases before this method's underlying map is populated, and + * {@link #decide} correctly reports nothing pending yet. A test that must order itself after the + * state {@link #decide} actually reads waits on this instead of on the phase. + * + * @return an unmodifiable snapshot; empty for a lead with no open questions + */ + Set pendingQuestionTurnIdsForTest(String lead) { + return pendingQuestionTurnIdsFor(lead); + } + private List pendingIncidentsFor(String lead) { return pendingIncidents.values().stream().filter(i -> lead.equals(i.key().lead())).toList(); } 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 580d7b7..72d8b4f 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -1802,7 +1802,11 @@ class MessageServiceTest { Thread asker = new Thread(() -> assertThrows(IllegalStateException.class, () -> service.ask(T, "which config file?", 30_000))); asker.start(); - awaitTicketPhaseOn(service, ticket, MessageService.Phase.ASKING); + // Phase.ASKING (markAsyncQuestion) is only the FIRST of ask()'s three steps; the assertion + // below depends on the THIRD (pushLoop.onQuestionOpened). Barrier on the push loop's own + // pending-question state instead of the phase — see ReplyPushLoop#pendingQuestionTurnIdsForTest. + MessageService.TaskView asking = awaitTicketPhaseOn(service, ticket, MessageService.Phase.ASKING); + awaitQuestionPendingOn(pushLoop, LEAD, asking.turnId()); assertEquals(ReplyPushLoop.Action.INJECT, pushLoop.decide(LEAD, 0, 0, 0), "the open question should be the one thing keeping this lead's schedule alive"); @@ -2094,6 +2098,26 @@ class MessageServiceTest { } } + /** + * fleetd #418: waits until {@code turnId} actually appears in {@code pushLoop}'s own open-question + * state for {@code lead} — {@link ReplyPushLoop#pendingQuestionTurnIdsForTest} — rather than until + * {@link MessageService#poll} reports {@link MessageService.Phase#ASKING}. {@code Phase.ASKING} is + * set by {@code markAsyncQuestion}, the FIRST of three steps {@code MessageService.ask()} performs; + * {@code pushLoop.onQuestionOpened} (the one that actually publishes to {@code pendingQuestions}) + * is the THIRD. Under load the asker thread can be descheduled between those two steps, so a + * barrier on the phase alone can release before the push loop has anything pending — a test that + * then asserts on {@link ReplyPushLoop#decide} is asserting on state that has not been published + * yet, not on the throwing path it is named for. + */ + private void awaitQuestionPendingOn(ReplyPushLoop pushLoop, String lead, String turnId) throws Exception { + long deadline = System.currentTimeMillis() + 3000; + while (!pushLoop.pendingQuestionTurnIdsForTest(lead).contains(turnId)) { + assertTrue(System.currentTimeMillis() < deadline, + "turnId " + turnId + " never appeared in the push loop's pending questions for " + lead); + Thread.sleep(5); + } + } + // --- CB-640: fleet health evidence accessors -------------------------------------------- @Test