From 6ed70700a00b56e183a102a2db0f7d9966cea3ea Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Fri, 4 Sep 2026 10:52:54 +0700 Subject: [PATCH] #282: don't let a chained fleet_ask kill its own async ticket MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit answer() opened a fresh forward waiter but, unlike send(), never registered it in asyncTasksByWaiter. So when a worker chained a second fleet_ask inside the same resumed turn (before calling fleet_reply), markAsyncQuestion had no Task to re-associate, and answer() then completed the async ticket's future with the second QUESTION as if it were a terminal reply — fleet_poll reported FAILED while the worker was still alive and mid-conversation. Fix: register answer()'s waiter in asyncTasksByWaiter (mirroring send()) so a chained ask can re-arm the ticket under its new turnId, and guard answer()'s finishAsyncTask call the same way sendAsync's own lambda already does (skip on Outcome.QUESTION). Also drop the stale asyncTasksByTurn entry left behind when markAsyncQuestion re-arms a task under a new turnId, a leak the fix makes reachable for the first time. Reachability confirmed by driving the exact sequence through the public API (sendAsync -> ask -> answer -> ask again) in a new test; reverting the production change makes it fail with "expected: but was: ", confirming it catches the regression. --- .../dev/ltms/fleet/msg/MessageService.java | 25 +++++++- .../ltms/fleet/msg/MessageServiceTest.java | 58 +++++++++++++++++++ 2 files changed, 82 insertions(+), 1 deletion(-) 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 5d92c2d..04e825b 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -927,7 +927,16 @@ 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 } @@ -935,7 +944,13 @@ public final class MessageService { try { Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS); Reply result = new Reply(outcomeOf(r.kind()), r.text(), r.turnId()); - finishAsyncTask(turnId, result); + // #282: this waiter can resolve with a FRESH question rather than a terminal reply — + // the worker chained a second fleet_ask before replying. Mirror sendAsync's own guard + // (:1000) and leave the ticket open (markAsyncQuestion above already re-armed it under + // the new turnId) instead of completing it here with a QUESTION "reply". + if (result.outcome() != Outcome.QUESTION) { + finishAsyncTask(turnId, result); + } return result; } catch (TimeoutException e) { // The worker resumed but hasn't replied yet — no completion fallback arms an answered @@ -948,6 +963,7 @@ public final class MessageService { Thread.currentThread().interrupt(); throw new IllegalStateException("interrupted awaiting reply from " + workerSession, e); } finally { + asyncTasksByWaiter.remove(reply); rendezvous.close(workerSession, reply); } } finally { @@ -1106,9 +1122,16 @@ public final class MessageService { private Task markAsyncQuestion(CompletableFuture waiter, String text, String turnId) { Task task = waiter == null ? null : asyncTasksByWaiter.get(waiter); if (task != null) { + String previousTurnId = task.turnId; task.question = new Reply(Outcome.QUESTION, text, turnId); task.turnId = turnId; asyncTasksByTurn.put(turnId, task); + // #282: a second fleet_ask in the same resumed turn re-arms an already-answered task + // (answer() re-registers it in asyncTasksByWaiter) under a FRESH turnId — drop the old + // key so asyncTasksByTurn does not keep growing by one stale entry per chained ask. + if (previousTurnId != null && !previousTurnId.equals(turnId)) { + asyncTasksByTurn.remove(previousTurnId, task); + } } return task; } 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 7ce8c22..6e3580c 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -975,6 +975,64 @@ class MessageServiceTest { assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome()); } + /** + * fleetd #282: a worker that chains a SECOND {@code fleet_ask} inside the same resumed turn — + * before it ever calls {@code fleet_reply} — used to kill its own async ticket. {@code answer()} + * opens a fresh forward waiter but (unlike {@code send()}) never registered it in + * {@code asyncTasksByWaiter}, so the second ask's {@code markAsyncQuestion} found no {@code Task} + * to re-associate. That waiter still resolved with the second {@code QUESTION} once the worker + * asked again, and {@code answer()} completed the ticket's future with that QUESTION "reply" + * unconditionally — so {@code fleet_poll} reported FAILED while the worker was still alive and + * the primary was mid-conversation with it. + * + *

Driven entirely through {@code MessageService}'s public API (sendAsync/ask/answer/poll) — + * never by reaching into {@link Rendezvous} or the task maps directly, so this test cannot pass + * for a reason unrelated to the real bug. + */ + @Test + void secondFleetAskInTheSameResumedTurnDoesNotKillTheAsyncTicket() throws Exception { + String ticket = messages.sendAsync(T, "task that asks twice"); + awaitWaiting(); + + // The worker's first fleet_ask. + CompletableFuture ask1 = + CompletableFuture.supplyAsync(() -> messages.ask(T, "Q1", 5000)); + MessageService.TaskView asking1 = awaitTicketPhase(ticket, MessageService.Phase.ASKING); + assertEquals("Q1", asking1.reply()); + + // The primary answers it — answer() resumes the turn and blocks for what comes next. + CompletableFuture answer1 = CompletableFuture.supplyAsync( + () -> messages.answer(asking1.turnId(), "a1", 5000)); + assertEquals("a1", ask1.get(5, TimeUnit.SECONDS).answer()); + + // Still in the SAME resumed turn — before replying — the worker asks again. + CompletableFuture ask2 = + CompletableFuture.supplyAsync(() -> messages.ask(T, "Q2", 5000)); + + // answer1's own call unblocks with the second QUESTION (documented QUESTION-chaining + // behaviour — see FleetMcp.answer's javadoc: "Answer it by calling fleet_send again with + // turnId=..."). The bug: this used to also kill the async ticket in the process. + MessageService.Reply firstAnswerResult = answer1.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.QUESTION, firstAnswerResult.outcome()); + String turnId2 = firstAnswerResult.turnId(); + + MessageService.TaskView asking2 = awaitTicketPhase(ticket, MessageService.Phase.ASKING); + assertEquals("Q2", asking2.reply(), + "the ticket must surface the SECOND question, not be dead/FAILED"); + assertEquals(turnId2, asking2.turnId()); + + // The primary answers the second question; the worker finally sends its real fleet_reply. + CompletableFuture answer2 = CompletableFuture.supplyAsync( + () -> messages.answer(turnId2, "a2", 5000)); + assertEquals("a2", ask2.get(5, TimeUnit.SECONDS).answer()); + awaitWaiting(); + assertTrue(rendezvous.resolve(T, "done")); + assertEquals(MessageService.Outcome.REPLIED, answer2.get(5, TimeUnit.SECONDS).outcome()); + + MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE); + assertEquals("done", done.reply()); + } + // --- CB-582: fleet_status pendingAsk() ------------------------------------------------------ @Test