fleetd #418: barrier the throw-path push-loop test on state decide() reads #419

Merged
ltms merged 1 commits from worker/418-588283-3 into main 2026-09-10 06:29:23 +02:00
2 changed files with 45 additions and 1 deletions
@@ -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.
*
* <p>{@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<String> pendingQuestionTurnIdsForTest(String lead) {
return pendingQuestionTurnIdsFor(lead);
}
private List<PendingIncident> pendingIncidentsFor(String lead) {
return pendingIncidents.values().stream().filter(i -> lead.equals(i.key().lead())).toList();
}
@@ -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