CB-582: close the question when ask() leaves by throwing
Found reviewing CB-582 before the merge, not by the implementer. ask() clears its question on three paths — no-waiter, timed out, and (from answer()) answered. It can also leave by throwing: an interrupt while blocked on the answer, or an ExecutionException from the answer future. Those run only the finally block, which tore down the rendezvous turn but not the push loop's copy. The result was a question that stayed pending for good: named in every nudge until it hit its own cap, then left in pendingQuestions with no remover at all. Teardown now happens where the rendezvous teardown already happens, so the two cannot drift apart again. Closing a turnId that was never pending is a no-op, so the normal paths are unaffected. The new test fails on the pre-fix code with expected: <STOP> but was: <INJECT>.
This commit is contained in:
@@ -536,6 +536,14 @@ public final class MessageService {
|
||||
// Only the fresh owner tears down the shared turn; a duplicate must leave it open.
|
||||
if (ticket.fresh()) {
|
||||
rendezvous.closeAsk(ticket.turnId());
|
||||
// CB-582: tear the push loop's copy down at the same point, not only on the three
|
||||
// paths that call clearAsyncQuestion. The answer future can complete exceptionally
|
||||
// (ExecutionException) or the thread be interrupted, and both leave this method by
|
||||
// throwing — the question would stay pending forever, keep being named in nudges
|
||||
// until its own cap, and never be removed from the map. Already-closed is a no-op.
|
||||
if (pushLoop != null) {
|
||||
pushLoop.questionClosed(ticket.turnId());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -939,6 +939,45 @@ class MessageServiceTest {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void anAskThatLeavesByThrowingStillClosesItsQuestion() throws Exception {
|
||||
// CB-582 follow-up. ask() calls clearAsyncQuestion on three paths — no-waiter, timed out,
|
||||
// and (from answer()) answered — but it can also leave by *throwing*: an interrupt while
|
||||
// blocked on the answer, or an ExecutionException from the answer future. Those paths run
|
||||
// only the finally block, so before the fix the push loop kept the question pending for
|
||||
// good: named in every nudge until its own cap, then never removed from the map at all.
|
||||
PrimaryRegistry registry = new PrimaryRegistry(null);
|
||||
registry.recordDelegation(T, LEAD);
|
||||
var scheduler = java.util.concurrent.Executors.newSingleThreadScheduledExecutor();
|
||||
// A backoff far longer than the test: the schedule is started but no tick ever fires, so
|
||||
// decide() is read directly and nothing here depends on timing.
|
||||
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, new AgentControl(new FakeHerdr()), inbox,
|
||||
scheduler, 5, 60_000);
|
||||
MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop);
|
||||
try {
|
||||
String ticket = service.sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
Thread asker = new Thread(() -> assertThrows(IllegalStateException.class,
|
||||
() -> service.ask(T, "which config file?", 30_000)));
|
||||
asker.start();
|
||||
awaitTicketPhaseOn(service, ticket, MessageService.Phase.ASKING);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, pushLoop.decide(LEAD, 0, 0, 0),
|
||||
"the open question should be the one thing keeping this lead's schedule alive");
|
||||
|
||||
asker.interrupt();
|
||||
asker.join(5000);
|
||||
assertFalse(asker.isAlive(), "the interrupted ask should have left ask() by throwing");
|
||||
|
||||
assertEquals(ReplyPushLoop.Action.STOP, pushLoop.decide(LEAD, 0, 0, 0),
|
||||
"an ask that threw must still close its question, or the loop nudges about it for good");
|
||||
} finally {
|
||||
service.close();
|
||||
scheduler.shutdownNow();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void aFleetWithNoPushLoopConfiguredBehavesExactlyAsToday() throws Exception {
|
||||
// `messages` (the shared field) uses the no-pushLoop constructor — poll() must not throw,
|
||||
|
||||
Reference in New Issue
Block a user