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 0a73221..5c532a6 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -965,13 +965,41 @@ public final class MessageService { return new AskResult(AskOutcome.ANSWERED, answer); } catch (TimeoutException e) { log.debug("fleet_ask from {} went unanswered in {}ms", workerSession, timeoutMillis); - // fleetd #307: mark the task BEFORE clearAsyncQuestion(forgetTurn=true) below drops it out of - // asyncTasksByTurn and nulls its turnId — that forgetting is deliberate and stays (it is - // what keeps the target from staying BUSY forever), but it would otherwise also erase - // askAnsweredAsyncTasks' only signal that the worker's eventual real fleet_reply still - // belongs to this task, stranding it in the inbox with a false "never replied" verdict. - markAskTimedOut(ticket.turnId()); - clearAsyncQuestion(ticket.turnId(), true); + // Only the fresh owner tears down the shared turn (mirrors the finally block below). + // A duplicate's own timeoutMillis says nothing about whether the SHARED ask is actually + // done — it must leave the close/forget bookkeeping to the fresh owner, exactly as it + // already leaves closeAsk to it. + if (ticket.fresh()) { + // fleetd #334: close the ask turn BEFORE forgetting this task's turnId mapping below. + // Before this fix the order was reversed — the mapping was forgotten here first, and + // rendezvous.closeAsk only ran afterward, in the shared finally. A primary's answer() + // call racing this exact timeout could then find rendezvous.askSession(turnId) still + // non-null (the ask still "answerable") after the Task mapping was already gone: + // answer()'s own asyncTasksByTurn lookup returned null, its task != null guard skipped + // the completion, and the async ticket sat at PENDING forever even though answer() + // itself reported the worker's real reply. Closing here first removes that window: + // any answer() call that still observes askSession(turnId) != null is necessarily + // racing a point BEFORE the forgetting below runs (both happen on this one thread, in + // this order, with nothing that yields in between), so the Task mapping is still there + // for it to find; any call that observes askSession(turnId) == null now correctly + // bails out STALE_TURN (see answer()'s own top check) before ever reaching + // asyncTasksByTurn. rendezvous.closeAsk is idempotent — a no-op once the turn is + // already removed, see its own javadoc — so the shared finally below re-running it + // for this same fresh call is harmless. + rendezvous.closeAsk(ticket.turnId()); + if (askTimeoutRaceHookForTest != null) { + // Test-only (fleetd #334): see the field's own javadoc. + askTimeoutRaceHookForTest.run(); + } + // fleetd #307: mark the task BEFORE clearAsyncQuestion(forgetTurn=true) below drops it + // out of asyncTasksByTurn and nulls its turnId — that forgetting is deliberate and + // stays (it is what keeps the target from staying BUSY forever), but it would + // otherwise also erase askAnsweredAsyncTasks' only signal that the worker's eventual + // real fleet_reply still belongs to this task, stranding it in the inbox with a false + // "never replied" verdict. + markAskTimedOut(ticket.turnId()); + clearAsyncQuestion(ticket.turnId(), true); + } return new AskResult(AskOutcome.TIMED_OUT, null); } catch (ExecutionException e) { Throwable cause = e.getCause(); @@ -1070,19 +1098,28 @@ public final class MessageService { // the chained ask deliberately left open. // // A null task is NOT only "this was never an async ticket". That reading was in this - // comment when #329 merged and it is wrong. A genuine async ticket also lands here - // with task == null, because ask()'s timeout path runs clearAsyncQuestion(turnId, - // true) — which drops the asyncTasksByTurn entry — in its catch block, while - // rendezvous.closeAsk(turnId) runs later, in its finally. Between those two the ask - // is still answerable but the map entry is already gone, so the lookup at :991 - // returns null and this ticket is never completed. Measured on 2026-09-04: a probe - // firing only that first half before answer() runs printed - // "answer=REPLIED phase=PENDING reply=null" — the same stranded ticket #329 set out - // to fix, one step earlier in the same race. The probe used forgetTurnForTest, which - // omits ask()'s markAskTimedOut; that cannot change the outcome, because askTimedOut - // is read only by askAnsweredAsyncTasks, and reply() never reaches it while this - // method's own waiter is live. So #329 narrows this window rather than closing it. - // Open as fleetd #334 — do not read this guard as complete. + // comment when #329 merged and it is wrong; it is still not the whole story after + // #334. A genuine async ticket can still land here with task == null — a blocking + // (wait:true) send's fleet_ask never has a Task at all, so that case is expected and + // fine. What #334 fixed was a SECOND, unintended way to get here with task == null: + // ask()'s timeout path used to run clearAsyncQuestion(turnId, true) — which drops the + // asyncTasksByTurn entry — in its catch block, while rendezvous.closeAsk(turnId) ran + // later, in its finally. Between those two the ask was still answerable but the map + // entry was already gone, so the lookup at :1053 returned null and this ticket was + // never completed. Measured on 2026-09-04: a probe firing only that first half before + // answer() ran printed "answer=REPLIED phase=PENDING reply=null" — the same stranded + // ticket #329 set out to fix, one step earlier in the same race; the probe used + // forgetTurnForTest, which omits ask()'s markAskTimedOut, and that omission does not + // change the outcome, because askTimedOut is read only by askAnsweredAsyncTasks, and + // reply() never reaches it while this method's own waiter is live. #334's fix + // reorders ask()'s timeout catch to run closeAsk before the forgetting (see the + // fresh-owner block there), which removes this path entirely rather than narrowing + // it further: once closeAsk has run, rendezvous.askSession(turnId) is null and + // answer() returns STALE_TURN from its own top check, before it ever reaches this + // lookup — see aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket in + // MessageServiceTest, which pins the exact window with askTimeoutRaceHookForTest. + // So by the time this line runs, task == null means only the ordinary blocking-send + // case (or #329's own already-fixed race elsewhere) — not this one. if (result.outcome() != Outcome.QUESTION && task != null) { finishAsyncTask(task, result); } @@ -1479,6 +1516,29 @@ public final class MessageService { this.abandonCleanupHookForTest = hook; } + /** + * Null in production; test seam for fleetd #334 — invoked from {@link #ask}'s {@code + * TimeoutException} catch, only for the fresh owner, right after {@code rendezvous.closeAsk} + * has run and before {@link #markAskTimedOut} / {@link #clearAsyncQuestion} forget this task's + * turnId mapping. A test installs this to call {@link #answer} for the very same {@code turnId} + * synchronously from inside that exact window, deterministically reproducing the race a real + * concurrent {@code answer()} call could otherwise only win by timing luck: with the ask already + * closed, that call must see {@code rendezvous.askSession(turnId) == null} and return {@link + * Outcome#STALE_TURN} immediately, never reaching {@code asyncTasksByTurn} at all — proving the + * window fleetd #334 describes (mapping forgotten while the ask was still "answerable") is + * closed, rather than merely narrowed the way fleetd #329 narrowed the sibling race in {@link + * #finishAsyncTask}. + */ + private volatile Runnable askTimeoutRaceHookForTest; + + /** + * Test-only (fleetd #334): install {@link #askTimeoutRaceHookForTest}. Package-private so the + * test, in the same package, can reach it without widening any production API. + */ + void setAskTimeoutRaceHookForTest(Runnable hook) { + this.askTimeoutRaceHookForTest = 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 001870d..d965d10 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -377,6 +377,76 @@ class MessageServiceTest { assertEquals(MessageService.Outcome.QUESTION, q.outcome()); } + /** + * fleetd #334. {@code ask()}'s {@code TimeoutException} catch used to forget this task's + * {@code turnId} mapping ({@code clearAsyncQuestion(turnId, true)}) BEFORE closing the ask + * ({@code rendezvous.closeAsk}, in the shared {@code finally}). A primary's {@code answer()} + * call racing that exact window found the ask still "answerable" ({@code + * rendezvous.askSession(turnId)} still non-null) while the {@code Task} was already forgotten, + * so its {@code task != null} guard skipped the completion and the async ticket sat at + * {@code PENDING} forever even though {@code answer()} itself reported a result. The fix + * (closing the ask first) makes this window impossible: a racing {@code answer()} call either + * still finds the ask open (and the {@code Task} mapping guaranteed intact) or finds it already + * closed (and bails {@code STALE_TURN} before ever touching the {@code Task}). This test pins + * the exact window with {@code askTimeoutRaceHookForTest} and proves both invariants the ticket + * named: (1) a late/racing {@code answer()} sees the ask as already lapsed ({@code STALE_TURN}), + * never made answerable again, and (2) the async ticket still resolves {@code DONE} once the + * worker's real {@code fleet_reply} lands — it is never stranded {@code PENDING}. + */ + @Test + void aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket() throws Exception { + String ticket = messages.sendAsync(T, "long task"); + awaitWaiting(); + injector.onStatus(T, AgentStatus.IDLE); // deliver + injector.onStatus(T, AgentStatus.WORKING); // worker picks it up, then pauses to ask + + CompletableFuture ask = + CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 200)); + + // Wait for the question to actually surface (poll sees ASKING) before racing the timeout. + MessageService.TaskView asking = null; + long deadline = System.currentTimeMillis() + 2000; + while ((asking == null || asking.phase() != MessageService.Phase.ASKING) + && System.currentTimeMillis() < deadline) { + asking = messages.poll(ticket); + //noinspection BusyWait + Thread.sleep(5); + } + assertNotNull(asking, "the question must surface before the ask times out"); + String turnId = asking.turnId(); + assertNotNull(turnId, "an ASKING view carries the turnId to answer on"); + + CompletableFuture lateAnswer = new CompletableFuture<>(); + messages.setAskTimeoutRaceHookForTest(() -> + lateAnswer.complete(messages.answer(turnId, "too late", 500))); + try { + MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.AskOutcome.TIMED_OUT, a.outcome()); + + MessageService.Reply late = lateAnswer.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.STALE_TURN, late.outcome(), + "a late answer racing the timeout teardown must see the ask as already lapsed"); + + // The worker resumes on its own (per the ask() contract) and eventually sends its real + // fleet_reply; the async ticket must still resolve with it, not strand at PENDING. + assertTrue(messages.reply(T, "real result"), "the worker's real reply must still be accepted"); + } finally { + messages.setAskTimeoutRaceHookForTest(null); + } + + MessageService.TaskView done = null; + deadline = System.currentTimeMillis() + 2000; + while ((done == null || done.phase() == MessageService.Phase.PENDING) + && System.currentTimeMillis() < deadline) { + done = messages.poll(ticket); + //noinspection BusyWait + Thread.sleep(5); + } + assertNotNull(done); + assertEquals(MessageService.Phase.DONE, done.phase(), "the async ticket must not be stranded PENDING"); + assertEquals("real result", done.reply()); + } + @Test void answeringAnUnknownTurnIsStale() { MessageService.Reply r = messages.answer(T + "#999", "too late", 500);