From 2dff4b84a4b7a32eb89694d6a4fa792c6041afb9 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Fri, 4 Sep 2026 10:35:01 +0700 Subject: [PATCH] Audit: async ticket / rendezvous lifecycle in fleetd msg package --- AUDIT.md | 96 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 96 insertions(+) create mode 100644 AUDIT.md diff --git a/AUDIT.md b/AUDIT.md new file mode 100644 index 0000000..c258de0 --- /dev/null +++ b/AUDIT.md @@ -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.