Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha 2dff4b84a4 Audit: async ticket / rendezvous lifecycle in fleetd msg package 2026-09-04 10:35:01 +07:00
3 changed files with 97 additions and 82 deletions
+96
View File
@@ -0,0 +1,96 @@
# Audit: async ticket / rendezvous lifecycle (`fleetd/src/main/java/dev/ltms/fleet/msg/`)
Scope: `Rendezvous.java` and `MessageService.java` — the lifecycle of an async ticket
(`Task`) and a rendezvous waiter: create, send, ask, answer, resolve, timeout, abandon, prune.
## Main finding
```
1. fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java:938
2. issue: answer() completes an async ticket's future with a QUESTION outcome when the worker
asks a second fleet_ask in the same resumed turn, permanently mislabeling a live delegation
as failed and losing its real reply to fleet_poll.
3. fix: guard the finishAsyncTask(turnId, result) call at line 938 the same way sendAsync's
lambda already guards its own call (lines 999-1005): skip it when result.outcome() ==
Outcome.QUESTION, and instead re-associate the task with the new turnId (as markAsyncQuestion
does on the first ask).
4. severity: high
```
### Call sequence that reaches it
1. Lead: `fleet_send{sessionId: W, content: "task", wait:false}` → `sendAsync` creates `task1`
/ `ticket1`. Its worker thread calls `send(W, content, ASYNC_TIMEOUT_MS, onAccepted, task1)`,
which does `asyncTasksByWaiter.put(reply, task1)` (line 802) before blocking on
`reply.get()`.
2. Worker `W` calls `fleet_ask{"Q1"}` → `ask(W, "Q1", t)`. `markAsyncQuestion` finds `task1`
via `asyncTasksByWaiter`, stamps `task1.turnId = turnId1`,
`asyncTasksByTurn[turnId1] = task1`. `resolveQuestion` wakes step 1's `send()`, which returns
`Outcome.QUESTION`; `sendAsync`'s lambda sees `QUESTION` and deliberately does **not** call
`finishAsyncTask` (lines 1000-1005) — `ticket1` correctly polls `Phase.ASKING`.
3. Lead polls, sees `ASKING`, answers: `fleet_send{turnId: turnId1, content: "A1"}` →
`answer(turnId1, "A1", t)`. This opens a **new** waiter via `rendezvous.open(workerSession)`
(line 929) but — unlike `send()` — never puts it into `asyncTasksByWaiter`.
`answerAsk(turnId1, "A1")` unblocks the worker's `ask()` call.
`clearAsyncQuestion(turnId1, false)` clears `task1.question` but keeps
`asyncTasksByTurn[turnId1] = task1` (deliberate, per its own javadoc). `answer()` then blocks
on its own `reply.get()` (line 936).
4. Worker `W`, still in the same resumed turn, calls `fleet_ask{"Q2"}` again before replying →
a second `ask(W, "Q2", t)`. `openAsk` mints `turnId2`. `markAsyncQuestion` looks up
`asyncTasksByWaiter.get(waiter)` for the waiter `answer()` opened in step 3 — **not found**
(never registered), so `task == null`; `task1.turnId` stays `turnId1`, no
`asyncTasksByTurn[turnId2]` entry is ever created. `resolveQuestion` still succeeds (it only
needs a live waiter, not a `Task`) and wakes `answer(turnId1,...)`'s blocked `reply.get()`
with `Resolution(QUESTION, "Q2", turnId2)`.
5. `answer(turnId1,...)` (line 936-939): `result = Reply(Outcome.QUESTION, "Q2", turnId2)`;
`finishAsyncTask(turnId1, result)` looks up `task1` by the **original** `turnId1` (still
stamped from step 3) and unconditionally does `task1.future.complete(result)` — completing
`ticket1`'s future with a **QUESTION** outcome, then removes `asyncTasksByTurn[turnId1]`.
`answer()` returns `Outcome.QUESTION` to the lead's own `fleet_send{turnId1,...}` call
(correct, and separately answerable via `turnId2`), but `ticket1` is now terminally done.
### What goes wrong
- `fleet_poll{ticket1}` now hits the `f.isDone()` branch in `poll()` permanently. `Outcome.QUESTION`
is not `REPLIED`/`COMPLETED_UNREPLIED` (`r.completed()` is false) and carries no
`WORKER_FAILED`/`BACKEND_EXHAUSTED` reason, so it falls through to
`Phase.FAILED`, `detail = "no reply — question"` — even though the worker is alive and only
waiting on `turnId2`.
- `asyncTasksByTurn` no longer has any entry for `task1`/`target`, so
`hasAsyncQuestion(target)` goes back to `false` immediately, and `hasOrphanedDelegation`
no longer excludes this target's real state correctly either.
- If the worker's eventual real `fleet_reply` (after `turnId2` is answered, or times out and it
finishes on its own) is not captured by a chained direct `answer(turnId2,...)` call,
`reply()`'s fast path (`rendezvous.resolve`) finds no live waiter, `askAnsweredAsyncTasks`
finds no candidate (`task1.future.isDone()` is already true, so it is excluded), and the reply
is silently dropped into the inbox as a **stranded reply** — unreachable from `ticket1` and
from `abandon()`'s stranded-reply recovery (no open `matching` task exists any more).
### Confidence
High. I traced this with no races or interleavings assumed beyond the documented, deterministic
CB-205 chained-ask protocol that `outcomeOf`/`answer()` already generically support (mapping
`Kind.QUESTION` through `answer()`'s own return value is clearly intentional — see the
`Reply.turnId()` javadoc). I did not run the suite, but grepped
`src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java` for a test exercising a *second*
`fleet_ask` inside one resumed (answered) turn and found none — `asyncQuestionBelongsToTheTaskThatOwnsItsForwardWaiter`
and `askSurfacesAsAQuestionAndTheAnswerResumesTheSameTurn` both cover only a single ask per
turn. A `git log -p` on this file also turned up the CB-588 comment (now at
`sendAsync`, describing the `whenComplete` hook) which explicitly frames
`answer()`'s `finishAsyncTask(turnId, result)` as firing "once a QUESTION is resolved" —
i.e. the author modeled that call as inherently terminal, which is exactly the assumption this
bug violates when the resumed turn asks again.
## Secondary (much shorter)
1. **Root-cause detail, same defect as above** — `answer()` (line 929) never puts its freshly
opened waiter into `asyncTasksByWaiter`, unlike `send()` (line 802). Even if the outcome
guard above is added, a chained second ask still can't be re-attached to `task1` via the
normal `markAsyncQuestion` path without also fixing this registration gap.
2. **Low / shape only** — `pruneTerminalTickets()` (line ~1092) is only invoked from inside
`sendAsync()`. A fleet whose sessions stop receiving new async sends (e.g. everything now
goes through blocking `send()`, or the target churns and workers are torn down) never prunes
its already-terminal `tasks` entries past `TICKET_TTL_NANOS`. Not reachable as a "stuck"
ticket (tickets still resolve correctly), only as unbounded `tasks`/`asyncTasksByWaiter`-adjacent
memory growth over a long-lived daemon with no further `sendAsync` traffic; did not verify
this is realistic in production traffic patterns, flagging as a shape only.
@@ -927,16 +927,7 @@ 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
}
@@ -944,13 +935,7 @@ 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());
// #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);
}
finishAsyncTask(turnId, result);
return result;
} catch (TimeoutException e) {
// The worker resumed but hasn't replied yet — no completion fallback arms an answered
@@ -963,7 +948,6 @@ 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 {
@@ -1122,16 +1106,9 @@ 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,64 +975,6 @@ 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