From b091c51eee2282c5b270f7ca11d9dd35413e47a3 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 12 Sep 2026 18:56:02 +0700 Subject: [PATCH] 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. --- .../dev/ltms/fleet/msg/MessageService.java | 64 +++++++++++++++---- .../ltms/fleet/msg/MessageServiceTest.java | 46 +++++++++++++ 2 files changed, 96 insertions(+), 14 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java index 8773b5d..3345935 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -1135,21 +1135,33 @@ public final class MessageService { } try { CompletableFuture 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)); diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java index 4b83d04..97f2380 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -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 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