fleetd #575: widen answer()'s try so one finally covers its STALE_TURN exit
The waiter cleanup pair (asyncTasksByWaiter.remove + rendezvous.close) was duplicated: two sites sit in a finally, the third was hand-rolled inline before answer()'s early STALE_TURN return, structurally outside any finally. Both rendezvous.answerAsk and clearAsyncQuestion(turnId, false) are total (cannot throw), so the gap never leaked in practice. But the duplicate was untested: mutating it away left all 1765 tests green, while the two finally-protected sites are each killed by 8-36 tests. Same shape as #572. Fix: widen the try to wrap the Task registration and the STALE_TURN check, so the single finally covers every exit and the hand-rolled copy is gone. Added a race hook + regression test that deterministically reproduces the 'ask lapsed between the lookup and the unblock' case and proves the fix still returns STALE_TURN and cleans up exactly once.
This commit is contained in:
@@ -1135,21 +1135,33 @@ public final class MessageService {
|
||||
}
|
||||
try {
|
||||
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(workerSession);
|
||||
// #282: mirror send()'s registration (:802) so a SECOND fleet_ask inside this same
|
||||
// resumed turn can re-associate the async ticket with its new turnId via
|
||||
// markAsyncQuestion — without this, that second ask has no Task to attach to, and
|
||||
// markAsyncQuestion silently returns null.
|
||||
Task task = asyncTasksByTurn.get(turnId);
|
||||
if (task != null) {
|
||||
asyncTasksByWaiter.put(reply, task);
|
||||
}
|
||||
if (!rendezvous.answerAsk(turnId, content)) {
|
||||
asyncTasksByWaiter.remove(reply);
|
||||
rendezvous.close(workerSession, reply);
|
||||
return new Reply(Outcome.STALE_TURN, null); // lapsed between the lookup and the unblock
|
||||
}
|
||||
clearAsyncQuestion(turnId, false);
|
||||
// fleetd #575: this try used to open below, AFTER the Task lookup/registration and the
|
||||
// STALE_TURN early return that follows it — so that return was covered only by a
|
||||
// hand-rolled copy of the finally's own cleanup pair, not the finally itself. Widening the
|
||||
// try up to wrap the registration closes the gap structurally: every exit from here on,
|
||||
// STALE_TURN included, now runs through the one finally below exactly once, and the
|
||||
// duplicated pair is gone. This did not fix a live leak — see the ticket: neither
|
||||
// rendezvous.answerAsk nor clearAsyncQuestion(turnId, false) can throw, so nothing ever
|
||||
// actually left through the old gap uncovered — but #572 found this exact drift (one
|
||||
// finally asserted, an identical sibling not) on this same file, and the hand-rolled copy
|
||||
// was the wrong shape to keep regardless of whether it was ever exercised.
|
||||
try {
|
||||
// #282: mirror send()'s registration (:802) so a SECOND fleet_ask inside this same
|
||||
// resumed turn can re-associate the async ticket with its new turnId via
|
||||
// markAsyncQuestion — without this, that second ask has no Task to attach to, and
|
||||
// markAsyncQuestion silently returns null.
|
||||
Task task = asyncTasksByTurn.get(turnId);
|
||||
if (task != null) {
|
||||
asyncTasksByWaiter.put(reply, task);
|
||||
}
|
||||
if (answerAskLapseRaceHookForTest != null) {
|
||||
// Test-only (fleetd #575): see the field's own javadoc.
|
||||
answerAskLapseRaceHookForTest.run();
|
||||
}
|
||||
if (!rendezvous.answerAsk(turnId, content)) {
|
||||
return new Reply(Outcome.STALE_TURN, null); // lapsed between the lookup and the unblock
|
||||
}
|
||||
clearAsyncQuestion(turnId, false);
|
||||
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
|
||||
Reply result = new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
|
||||
// #282: this waiter can resolve with a FRESH question rather than a terminal reply —
|
||||
@@ -1649,6 +1661,30 @@ public final class MessageService {
|
||||
this.askTimeoutRaceHookForTest = hook;
|
||||
}
|
||||
|
||||
/**
|
||||
* Null in production; test seam for fleetd #575 — invoked from {@link #answer}, right after this
|
||||
* call's own {@code Task} registration and right before its {@code rendezvous.answerAsk(turnId,
|
||||
* content)} call. A test installs this to complete the SAME turnId's ask directly via {@link
|
||||
* Rendezvous#answerAsk} from inside that exact window, deterministically reproducing what a
|
||||
* second, concurrent {@code answer()} call racing to unblock the same ask can otherwise only
|
||||
* win by timing luck: this call's own {@code askSession(turnId)} lookup at the top already saw
|
||||
* the ask as open, but by the time it reaches {@code rendezvous.answerAsk} here, the other call
|
||||
* already completed it (or the worker's own {@code ask()} teardown already closed it) — so this
|
||||
* call must see {@code false} and return {@link Outcome#STALE_TURN}, exactly the "lapsed between
|
||||
* the lookup and the unblock" case named at that call site. Proves the #575 fix (widening this
|
||||
* method's try so a single finally covers this exit) does not change that outcome and still
|
||||
* cleans this call's own {@code reply} up exactly once.
|
||||
*/
|
||||
private volatile Runnable answerAskLapseRaceHookForTest;
|
||||
|
||||
/**
|
||||
* Test-only (fleetd #575): install {@link #answerAskLapseRaceHookForTest}. Package-private so the
|
||||
* test, in the same package, can reach it without widening any production API.
|
||||
*/
|
||||
void setAnswerAskLapseRaceHookForTest(Runnable hook) {
|
||||
this.answerAskLapseRaceHookForTest = hook;
|
||||
}
|
||||
|
||||
/** A new send must not open a waiter while an async ticket owns this worker's paused turn. */
|
||||
private boolean hasAsyncQuestion(String target) {
|
||||
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
|
||||
|
||||
@@ -506,6 +506,52 @@ class MessageServiceTest {
|
||||
"an answer to a turn that never existed (or already lapsed) is stale, not a hang");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #575. {@code answer()}'s STALE_TURN return for "the ask lapsed between the lookup and
|
||||
* the unblock" ({@code rendezvous.answerAsk(turnId, content)} returning {@code false} even though
|
||||
* this call's own {@code rendezvous.askSession(turnId)} check at the top saw the ask as open) used
|
||||
* to be covered only by a hand-rolled copy of the cleanup pair its own finally already runs, not
|
||||
* the finally itself — a structural gap, closed by widening the try up to cover the registration
|
||||
* above it. This test pins the exact race with {@link
|
||||
* MessageService#setAnswerAskLapseRaceHookForTest}, which fires right before this call's own
|
||||
* {@code rendezvous.answerAsk} call and completes the same turnId's ask directly — reproducing
|
||||
* what a second, concurrent {@code answer()} winning that race could otherwise only do by timing
|
||||
* luck. Proves the fix changed no behaviour on this path: it still returns {@code STALE_TURN},
|
||||
* and this call's own forward waiter is still closed exactly once (not left open, and not closed
|
||||
* twice — there is now only one cleanup site left to run).
|
||||
*/
|
||||
@Test
|
||||
void answerLosingTheRaceToAnAlreadyAnsweredAskStillReturnsStaleTurnAndCleansUpOnce() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
String turnId = asking.turnId();
|
||||
assertNotNull(turnId, "an ASKING view carries the turnId to answer on");
|
||||
|
||||
messages.setAnswerAskLapseRaceHookForTest(() -> rendezvous.answerAsk(turnId, "raced in first"));
|
||||
try {
|
||||
MessageService.Reply r = messages.answer(turnId, "too late", 500);
|
||||
assertEquals(MessageService.Outcome.STALE_TURN, r.outcome(),
|
||||
"an ask already answered by the race must be seen as lapsed, not double-delivered");
|
||||
} finally {
|
||||
messages.setAnswerAskLapseRaceHookForTest(null);
|
||||
}
|
||||
|
||||
// Cleanup ran exactly once: the forward waiter THIS call opened is closed, not leaked.
|
||||
assertNull(rendezvous.currentWaiter(T),
|
||||
"the forward waiter this answer() call opened must be closed after a STALE_TURN return");
|
||||
|
||||
// The worker's own ask() call, unblocked by the hook's direct answerAsk, still completes
|
||||
// normally — the race this test simulates does not strand it.
|
||||
MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.AskOutcome.ANSWERED, a.outcome());
|
||||
assertEquals("raced in first", a.answer());
|
||||
}
|
||||
|
||||
// --- timeout, answer, poll, and lock-contention edges ----------------------------------
|
||||
|
||||
@Test
|
||||
|
||||
Reference in New Issue
Block a user