From ab0cc71aa40e71dad3ad25199694d921bfc1c069 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 10 Sep 2026 11:24:41 +0700 Subject: [PATCH] fleetd #418: barrier the throw-path push-loop test on state decide() reads MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit anAskThatLeavesByThrowingStillClosesItsQuestion barriered on Phase.ASKING, which markAsyncQuestion sets in ask()'s FIRST step. The assertion right after it depends on ask()'s THIRD step (pushLoop.onQuestionOpened), which is what actually populates ReplyPushLoop's pendingQuestions map. Under load the asker thread can be descheduled between those two steps, so the barrier released before decide() had anything to see, and it correctly returned STOP instead of the expected INJECT. Add ReplyPushLoop#pendingQuestionTurnIdsForTest, a package-private test seam (modeled on MessageService#isCompletionStampedForTest) exposing the private pendingQuestionTurnIdsFor. The test now waits for its own turnId to appear there before asserting on decide() — not for decide() itself to return INJECT, which would make the barrier assert nothing. Checked every other awaitTicketPhaseOn(..., Phase.ASKING) in the file (two, in the CB-582 nudge tests): both are followed by a real awaitNudge() that waits for an actual agent.prompt push-loop call before any assertion depends on push-loop state, so they are not exposed to this race. No production behaviour changed. --- .../dev/ltms/fleet/msg/ReplyPushLoop.java | 20 ++++++++++++++ .../ltms/fleet/msg/MessageServiceTest.java | 26 ++++++++++++++++++- 2 files changed, 45 insertions(+), 1 deletion(-) 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 -- 2.52.0