Merge #282 (PR #289): a chained second fleet_ask no longer kills its own async ticket

This commit is contained in:
Dai Ha
2026-09-04 10:56:42 +07:00
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