Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ab0cc71aa4 |
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user