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 {
|
try {
|
||||||
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(workerSession);
|
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(workerSession);
|
||||||
// #282: mirror send()'s registration (:802) so a SECOND fleet_ask inside this same
|
// fleetd #575: this try used to open below, AFTER the Task lookup/registration and the
|
||||||
// resumed turn can re-associate the async ticket with its new turnId via
|
// STALE_TURN early return that follows it — so that return was covered only by a
|
||||||
// markAsyncQuestion — without this, that second ask has no Task to attach to, and
|
// hand-rolled copy of the finally's own cleanup pair, not the finally itself. Widening the
|
||||||
// markAsyncQuestion silently returns null.
|
// try up to wrap the registration closes the gap structurally: every exit from here on,
|
||||||
Task task = asyncTasksByTurn.get(turnId);
|
// STALE_TURN included, now runs through the one finally below exactly once, and the
|
||||||
if (task != null) {
|
// duplicated pair is gone. This did not fix a live leak — see the ticket: neither
|
||||||
asyncTasksByWaiter.put(reply, task);
|
// 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
|
||||||
if (!rendezvous.answerAsk(turnId, content)) {
|
// finally asserted, an identical sibling not) on this same file, and the hand-rolled copy
|
||||||
asyncTasksByWaiter.remove(reply);
|
// was the wrong shape to keep regardless of whether it was ever exercised.
|
||||||
rendezvous.close(workerSession, reply);
|
|
||||||
return new Reply(Outcome.STALE_TURN, null); // lapsed between the lookup and the unblock
|
|
||||||
}
|
|
||||||
clearAsyncQuestion(turnId, false);
|
|
||||||
try {
|
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);
|
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
|
||||||
Reply result = new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
|
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 —
|
// #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;
|
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. */
|
/** A new send must not open a waiter while an async ticket owns this worker's paused turn. */
|
||||||
private boolean hasAsyncQuestion(String target) {
|
private boolean hasAsyncQuestion(String target) {
|
||||||
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.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");
|
"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 ----------------------------------
|
// --- timeout, answer, poll, and lock-contention edges ----------------------------------
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
|
|||||||
Reference in New Issue
Block a user