#282: don't let a chained fleet_ask kill its own async ticket
CI / contract (pull_request) Successful in 1m6s
CI / build (pull_request) Successful in 1m38s

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: <ASKING> but was: <FAILED>", confirming it catches the
regression.
This commit is contained in:
Dai Ha
2026-09-04 10:52:54 +07:00
parent 66e5247b6d
commit 6ed70700a0
2 changed files with 82 additions and 1 deletions
@@ -927,7 +927,16 @@ public final class MessageService {
}
try {
CompletableFuture<Rendezvous.Resolution> 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<Rendezvous.Resolution> 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;
}
@@ -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.
*
* <p>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<MessageService.AskResult> 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<MessageService.Reply> 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<MessageService.AskResult> 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<MessageService.Reply> 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