A second fleet_ask in the same resumed turn kills its own async ticket: answer() completes it with QUESTION #282

Closed
opened 2026-09-04 05:38:40 +02:00 by ltms · 1 comment
Owner

Found by an audit of msg/. I verified the whole chain in the code myself — every claim below names the line I read.

The defect

answer() completes an async ticket's future with a QUESTION outcome when the worker asks a second fleet_ask inside the same resumed turn. The ticket is then terminally done, fleet_poll reports it FAILED, and the worker's real reply is stranded — while the worker is alive and still working.

The asymmetry

Two places finish an async task after a send resolves. One guards against QUESTION; its sibling does not.

sendAsync's worker lambda, MessageService.java:999-1005:

Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task);
if (result.outcome() == Outcome.QUESTION) {
    // Keep the accepted owner until answer() finishes it. markAsyncQuestion may run
    // just after resolveQuestion wakes this thread.
} else {
    finishAsyncTask(task, result);
}

answer(), MessageService.java:936-939 — no guard:

Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
Reply result = new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
finishAsyncTask(turnId, result);
return result;

The comment on the first one says the author knew a resolved send can carry QUESTION and must not finish the task. answer() was written as if its own resolution could never be a question — that is the assumption the bug breaks.

The sequence, and the line for each step

  1. Lead calls fleet_send{wait:false} → sendAsync creates task1/ticket1. Its worker thread enters send(), which registers the waiter: asyncTasksByWaiter.put(reply, task) — :802.
  2. Worker calls fleet_ask{"Q1"}. markAsyncQuestion (:1106-1113) finds task1 through asyncTasksByWaiter, sets task1.turnId = turnId1 and asyncTasksByTurn.put(turnId1, task1). sendAsync's lambda sees QUESTION and correctly skips finishing. ticket1 polls as asking. Correct so far.
  3. Lead answers: fleet_send{turnId1} → answer() opens a new waiter at :929 — and, unlike send(), never puts it in asyncTasksByWaiter. Then clearAsyncQuestion(turnId1, false) runs (:1117-1131): with forgetTurn = false it clears task1.question but keeps asyncTasksByTurn[turnId1] = task1 and keeps task1.turnId. answer() then blocks on its own reply.get().
  4. Worker, still in the same resumed turn, calls fleet_ask{"Q2"} before replying. markAsyncQuestion looks up asyncTasksByWaiter.get(waiter) — the waiter from step 3 was never registered, so it returns null and turnId2 is never tied to task1. But resolveQuestion (Rendezvous.java:176-178) only needs a live waiter for the session, not a Task, so it still completes answer()'s waiter with Resolution(QUESTION, "Q2", turnId2).
  5. answer() wakes, builds result = Reply(QUESTION, "Q2", turnId2) and calls finishAsyncTask(turnId1, result) unguarded. That overload (:1143-1148) looks up asyncTasksByTurn.get(turnId1) — still task1 from step 3 — and completes task1.future with a QUESTION reply.

What the operator sees

completed() is true only for REPLIED and COMPLETED_UNREPLIED (:123), so QUESTION falls to the failure branch: fleet_poll{ticket1} reports Phase.FAILED. finishAsyncTask(task, result) also removes the asyncTasksByTurn entry (:1138), so hasAsyncQuestion goes false and the ticket is detached from the live turn. The worker's eventual real fleet_reply has no waiter and no live task to land on — it is stranded in the inbox, unreachable from the ticket.

So: a live delegation is reported as failed, and its answer is lost. That is the worst direction for this to fail in.

The fix

Guard answer()'s finishAsyncTask the same way sendAsync's lambda already guards its own: skip when result.outcome() == Outcome.QUESTION.

That alone stops the ticket being killed. To also make the chained ask work rather than merely not break, answer() must register its new waiter in asyncTasksByWaiter (mirroring send() at :802) so markAsyncQuestion can re-associate task1 with turnId2. Do both, and say in your report what each one changes.

Test

Drive the real public sequence — sendAsync → ask → answer → ask again — and assert the ticket is not FAILED and the second question is visible on the ticket. Do not build the state by reaching into Rendezvous or the task maps; a test that does that would have passed the whole time this bug existed.

There is currently no test with two fleet_ask calls inside one resumed turn — asyncQuestionBelongsToTheTaskThatOwnsItsForwardWaiter covers only one ask per turn. Check that yourself before you start.

Reachability caveat, stated honestly

I confirmed every link above by reading the code. I have not reproduced this against a live worker. Chained asks are a supported CB-205 mechanism (outcomeOf and answer() handle QUESTION resolutions generically), so the path is real, but whether a member in practice asks twice in one resumed turn before replying is not something I measured. If you find the second ask cannot actually happen, say so and close this — "not reachable, here is why" is a good result.

Found by an audit of `msg/`. **I verified the whole chain in the code myself** — every claim below names the line I read. ## The defect `answer()` completes an async ticket's `future` with a `QUESTION` outcome when the worker asks a **second** `fleet_ask` inside the same resumed turn. The ticket is then terminally done, `fleet_poll` reports it `FAILED`, and the worker's real reply is stranded — while the worker is alive and still working. ## The asymmetry Two places finish an async task after a send resolves. One guards against `QUESTION`; its sibling does not. `sendAsync`'s worker lambda, `MessageService.java:999-1005`: ```java Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task); if (result.outcome() == Outcome.QUESTION) { // Keep the accepted owner until answer() finishes it. markAsyncQuestion may run // just after resolveQuestion wakes this thread. } else { finishAsyncTask(task, result); } ``` `answer()`, `MessageService.java:936-939` — no guard: ```java Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS); Reply result = new Reply(outcomeOf(r.kind()), r.text(), r.turnId()); finishAsyncTask(turnId, result); return result; ``` The comment on the first one says the author knew a resolved send can carry `QUESTION` and must not finish the task. `answer()` was written as if its own resolution could never be a question — that is the assumption the bug breaks. ## The sequence, and the line for each step 1. Lead calls `fleet_send{wait:false}` → `sendAsync` creates `task1`/`ticket1`. Its worker thread enters `send()`, which registers the waiter: `asyncTasksByWaiter.put(reply, task)` — **`:802`**. 2. Worker calls `fleet_ask{"Q1"}`. `markAsyncQuestion` (**`:1106-1113`**) finds `task1` through `asyncTasksByWaiter`, sets `task1.turnId = turnId1` and `asyncTasksByTurn.put(turnId1, task1)`. `sendAsync`'s lambda sees `QUESTION` and correctly skips finishing. `ticket1` polls as asking. Correct so far. 3. Lead answers: `fleet_send{turnId1}` → `answer()` opens a **new** waiter at **`:929`** — and, unlike `send()`, **never puts it in `asyncTasksByWaiter`**. Then `clearAsyncQuestion(turnId1, false)` runs (**`:1117-1131`**): with `forgetTurn = false` it clears `task1.question` but **keeps** `asyncTasksByTurn[turnId1] = task1` and keeps `task1.turnId`. `answer()` then blocks on its own `reply.get()`. 4. Worker, still in the same resumed turn, calls `fleet_ask{"Q2"}` before replying. `markAsyncQuestion` looks up `asyncTasksByWaiter.get(waiter)` — the waiter from step 3 was never registered, so it returns `null` and `turnId2` is never tied to `task1`. But `resolveQuestion` (**`Rendezvous.java:176-178`**) only needs a live waiter for the session, not a `Task`, so it still completes `answer()`'s waiter with `Resolution(QUESTION, "Q2", turnId2)`. 5. `answer()` wakes, builds `result = Reply(QUESTION, "Q2", turnId2)` and calls `finishAsyncTask(turnId1, result)` unguarded. That overload (**`:1143-1148`**) looks up `asyncTasksByTurn.get(turnId1)` — still `task1` from step 3 — and completes `task1.future` with a `QUESTION` reply. ## What the operator sees `completed()` is true only for `REPLIED` and `COMPLETED_UNREPLIED` (**`:123`**), so `QUESTION` falls to the failure branch: `fleet_poll{ticket1}` reports `Phase.FAILED`. `finishAsyncTask(task, result)` also removes the `asyncTasksByTurn` entry (**`:1138`**), so `hasAsyncQuestion` goes false and the ticket is detached from the live turn. The worker's eventual real `fleet_reply` has no waiter and no live task to land on — it is stranded in the inbox, unreachable from the ticket. So: a live delegation is reported as failed, and its answer is lost. That is the worst direction for this to fail in. ## The fix Guard `answer()`'s `finishAsyncTask` the same way `sendAsync`'s lambda already guards its own: skip when `result.outcome() == Outcome.QUESTION`. That alone stops the ticket being killed. To also make the chained ask *work* rather than merely not break, `answer()` must register its new waiter in `asyncTasksByWaiter` (mirroring `send()` at `:802`) so `markAsyncQuestion` can re-associate `task1` with `turnId2`. Do both, and say in your report what each one changes. ## Test Drive the real public sequence — `sendAsync` → `ask` → `answer` → `ask` again — and assert the ticket is **not** `FAILED` and the second question is visible on the ticket. Do not build the state by reaching into `Rendezvous` or the task maps; a test that does that would have passed the whole time this bug existed. There is currently no test with two `fleet_ask` calls inside one resumed turn — `asyncQuestionBelongsToTheTaskThatOwnsItsForwardWaiter` covers only one ask per turn. Check that yourself before you start. ## Reachability caveat, stated honestly I confirmed every link above by reading the code. I have **not** reproduced this against a live worker. Chained asks are a supported CB-205 mechanism (`outcomeOf` and `answer()` handle `QUESTION` resolutions generically), so the path is real, but whether a member in practice asks twice in one resumed turn before replying is not something I measured. If you find the second ask cannot actually happen, say so and close this — "not reachable, here is why" is a good result.
Author
Owner

Merged as bfabe13, with a follow-up commit e93b5f6 correcting one claim. Full build of the real merge: 1288 tests, BUILD SUCCESS.

One claim in the PR does not hold

The PR says each half of the fix is necessary on its own:

Alone, this lets markAsyncQuestion re-associate the task, but without guard #1 the ticket would still get killed by the unconditional finishAsyncTask before that association could ever be used.

I removed only the QUESTION guard, keeping the waiter registration and the stale-key removal, and ran MessageServiceTest. It passed. So that sentence is wrong.

The reason is in the ordering. ask() calls markAsyncQuestion at :860 and resolveQuestion at :861 — in that order. So the task has already moved to the new turnId, and the stale-key removal (the third change in this PR) has already dropped the old key, before answer()'s thread ever wakes. finishAsyncTask(turnId1, result) then looks up asyncTasksByTurn.get(turnId1), finds nothing, and does nothing. The guard never fires on the tested path.

What the new test proves is the pair, not either half. That is still a real proof of a real fix — the defect is gone — but the PR's account of why is not accurate, and I would rather correct it than let it stand as the explanation the next person reads.

What I did about it

I kept the guard and wrote the measurement into the code (e93b5f6), because both possible mistakes are live:

  • Someone reads it, sees it never fires, and deletes it as dead code — without re-checking the ordering it depends on.
  • Someone trusts it as the protection, and removes the waiter registration or the stale-key removal.

The comment now says it is defence in depth, names the :860/:861 ordering it rests on, and points at sendAsync's sibling guard, whose own comment (:1017) warns that "markAsyncQuestion may run just after resolveQuestion wakes this thread" — i.e. the author of that guard did not consider the ordering safe to rely on. Given that, keeping a cheap mirrored guard is right; claiming the test proves it is not.

What did verify cleanly

The fix itself is good and I checked the parts that matter:

  • The new asyncTasksByWaiter entry is removed in the finally and on the STALE_TURN early return, so it cannot leak.
  • markAsyncQuestion now drops the stale turnId key. That leak was reachable for the first time only because this fix made a task re-armable — the worker found it on its own, and it is the change that actually does the work here.
  • The worker checked the rest of msg/ for the same missing-guard shape and reported none: reply() and abandon() both build their own terminal Reply rather than forwarding a Rendezvous.Resolution, so neither can carry QUESTION through. I spot-checked that and agree.

The unexplained build failure

The worker reported that its very first mvn clean install exited 1 with output it could not read, and that every run afterwards was green. I did not reproduce it — my build of the real merge was green on the first attempt at 1288 tests. So it stays unexplained rather than diagnosed. Recording it here in case it shows up again; a flake nobody saw is not the same as no flake.

Merged as `bfabe13`, with a follow-up commit `e93b5f6` correcting one claim. Full build of the real merge: **1288 tests, BUILD SUCCESS**. ## One claim in the PR does not hold The PR says each half of the fix is necessary on its own: > Alone, this lets markAsyncQuestion re-associate the task, but without guard #1 the ticket would still get killed by the unconditional finishAsyncTask before that association could ever be used. I removed **only** the `QUESTION` guard, keeping the waiter registration and the stale-key removal, and ran `MessageServiceTest`. It **passed**. So that sentence is wrong. The reason is in the ordering. `ask()` calls `markAsyncQuestion` at `:860` and `resolveQuestion` at `:861` — in that order. So the task has already moved to the new `turnId`, and the stale-key removal (the third change in this PR) has already dropped the old key, before `answer()`'s thread ever wakes. `finishAsyncTask(turnId1, result)` then looks up `asyncTasksByTurn.get(turnId1)`, finds nothing, and does nothing. The guard never fires on the tested path. What the new test proves is the **pair**, not either half. That is still a real proof of a real fix — the defect is gone — but the PR's account of why is not accurate, and I would rather correct it than let it stand as the explanation the next person reads. ## What I did about it I kept the guard and wrote the measurement into the code (`e93b5f6`), because both possible mistakes are live: - Someone reads it, sees it never fires, and deletes it as dead code — without re-checking the ordering it depends on. - Someone trusts it as the protection, and removes the waiter registration or the stale-key removal. The comment now says it is defence in depth, names the `:860`/`:861` ordering it rests on, and points at `sendAsync`'s sibling guard, whose own comment (`:1017`) warns that "markAsyncQuestion may run just after resolveQuestion wakes this thread" — i.e. the author of that guard did not consider the ordering safe to rely on. Given that, keeping a cheap mirrored guard is right; claiming the test proves it is not. ## What did verify cleanly The fix itself is good and I checked the parts that matter: - The new `asyncTasksByWaiter` entry is removed in the `finally` **and** on the `STALE_TURN` early return, so it cannot leak. - `markAsyncQuestion` now drops the stale `turnId` key. That leak was reachable for the first time only because this fix made a task re-armable — the worker found it on its own, and it is the change that actually does the work here. - The worker checked the rest of `msg/` for the same missing-guard shape and reported none: `reply()` and `abandon()` both build their own terminal `Reply` rather than forwarding a `Rendezvous.Resolution`, so neither can carry `QUESTION` through. I spot-checked that and agree. ## The unexplained build failure The worker reported that its very first `mvn clean install` exited 1 with output it could not read, and that every run afterwards was green. **I did not reproduce it** — my build of the real merge was green on the first attempt at 1288 tests. So it stays unexplained rather than diagnosed. Recording it here in case it shows up again; a flake nobody saw is not the same as no flake.
ltms closed this issue 2026-09-04 06:00:04 +02:00
Sign in to join this conversation.
1 Participants
Notifications
Due Date
No due date set.
Dependencies

No dependencies set.

Reference: fleet/fleetd#282