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