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