answer() mutates a Task under sessionLocks while ask()'s timeout path mutates the same Task under no lock — finishAsyncTask can throw NPE on a null turnId #324

Closed
opened 2026-09-04 09:26:17 +02:00 by ltms · 1 comment
Owner

Found as the shape-check tail of #318, then narrowed by me to one concrete reachable failure. I read the whole path myself.

Read the severity honestly

I have not observed this. The window is between two reads of one field, and it needs a fleet_ask to lapse in the same instant the lead's answer lands. It is narrow. I am filing it because the failure is a thrown exception on the lead's own call, the fix is two lines, and the asymmetry behind it is worth a decision rather than another patch.

The asymmetry

answer() takes the target's lock before it touches anything:

MessageService.java:975   ReentrantLock lock = sessionLocks.computeIfAbsent(workerSession, _ -> new ReentrantLock());

and then, holding it, calls asyncTasksByTurn.get(turnId) (:985), clearAsyncQuestion(turnId, false) (:994) and finishAsyncTask(turnId, result) (:1011).

ask() takes no lock at all. I read the whole method — there is no sessionLocks acquisition anywhere in it. On its timeout path it mutates the same Task:

MessageService.java:933   markAskTimedOut(ticket.turnId());
MessageService.java:934   clearAsyncQuestion(ticket.turnId(), true);

clearAsyncQuestion(_, true) removes the task from asyncTasksByTurn and sets task.turnId = null (:1227-1228).

So one side of this invariant is locked and the other is not. Task.turnId, Task.question and Task.completedNanos are volatile, so this is not a visibility bug. It is a compound-action bug.

The concrete failure

MessageService.java:1234-1239
private void finishAsyncTask(Task task, Reply result) {
    task.future.complete(result);
    if (task.turnId != null) {
        asyncTasksByTurn.remove(task.turnId, task);
    }
}

task.turnId is read twice — once for the null check, once for the removal. volatile makes each read fresh; it does not make the pair atomic.

Path in. The lead answers a fleet_ask at the moment its ~55 s window lapses.

  1. Lead's thread, in answer(), holding the lock: finishAsyncTask(turnId, result) (:1011) → asyncTasksByTurn.get(turnId) returns the task (:1243) → finishAsyncTask(task, result).
  2. It completes the future (:1235), then reads task.turnId and sees non-null.
  3. Worker's thread, in ask()'s catch (TimeoutException), holding no lock: clearAsyncQuestion(ticket.turnId(), true) sets task.turnId = null (:1228).
  4. Lead's thread reads task.turnId again for the removal. It is now null. ConcurrentHashMap.remove(Object, Object) throws NullPointerException on a null key, by contract.

Direction of harm: a thrown exception out of answer(), on the lead's own fleet_send{turnId, content} call. The lead is told its answer failed. It did not — task.future.complete(result) already ran on the line before, so the ticket is completed and the answer delivered. A lead that reacts by re-sending is acting on a false failure. The stale asyncTasksByTurn entry also leaks, and hasAsyncQuestion (:1251) scans that map to decide whether a target is BUSY, so the target can be reported busy after its task is done.

Confidence: the code path is confirmed by reading. I have not reproduced the interleaving.

What I want

Goal: answer() must not be able to throw because ask() changed a field underneath it, and the two sides of this invariant must stop being guarded differently.

Invariants:

  1. ask() must not start blocking on the target's session lock. It is called from the worker's own turn and it already blocks for up to ~55 s on the answer. Making it contend with send/answer for the same lock risks a deadlock or a stalled turn. If your fix takes that lock in ask(), prove it cannot deadlock — do not assume it.
  2. clearAsyncQuestion's forgetTurn=true forgetting stays. #307's fix depends on it, and the javadoc explains why: it is what stops the target being reported BUSY forever.
  3. markAskTimedOut must still run before that forgetting. Also #307. Its own javadoc says a lookup afterwards finds nothing.
  4. Do not change what answer() returns on success. This ticket is about not throwing on a race, not about new semantics.

Candidate mechanism, as a candidate only: read task.turnId once into a local in finishAsyncTask and use that local for both the check and the removal. Decide it yourself and justify it.

That candidate fixes the crash and leaves the asymmetry. So answer this too, and it is the more important half:

  • Is a local read enough, or is the unlocked ask() side wrong? Look for other compound actions on a Task that one side performs locked and the other unlocked. markAskTimedOut (:1208-1213) does a get-then-check-then-set; clearAsyncQuestion (:1216-1230) does a get-then-check-then-two-writes; answer() does get-then-clear-then-finish across three calls. Say which of these are genuinely at risk and which are safe, and why. Name the path in for each one you call a defect — a compound action that no second thread can reach is not a defect.
  • If you conclude the right answer is a lock, or a different structure (for example making the task's turn a single immutable value that is swapped atomically rather than a field that is nulled), say so and say what it would cost. Do not build it in this ticket — report it and I will decide.

A tested, reported deviation is a good outcome.

Rules

  • Prove the crash with a test that fails without the fix. A latch or a test seam that nulls turnId between the two reads is the honest way; a test that merely calls the two methods in sequence proves nothing, because the bug is an interleaving. If you cannot build that seam without changing production code beyond the fix, say so and state exactly what your test does and does not prove.
  • Mutation proof required: revert the fix, quote the real failure output, restore it.
  • Never git stash — the stash is shared across every worktree here.
  • Never run git worktree remove or git worktree prune — other workers are live in those directories.
  • Stage files explicitly; never git add -A. Never merge.
  • Run cd fleetd && mvn clean install unpiped; quote the real Tests run: and BUILD lines. Never read $? after a pipe.
  • fleetd/fleetd.yaml is gitignored and absent from your worktree. Do not report on its contents.
  • Put your full report in the PR body as well as in your fleet_reply.

Already ruled out, do not re-report

The strandedReplies / queuedDeliveries pair was checked and dismissed: each touch is a single atomic map operation on a documented (CB-640) self-healing best-effort flag, not a compound action.

Found as the shape-check tail of #318, then narrowed by me to one concrete reachable failure. I read the whole path myself. ## Read the severity honestly I have not observed this. The window is between two reads of one field, and it needs a `fleet_ask` to lapse in the same instant the lead's answer lands. It is narrow. I am filing it because the failure is a thrown exception on the lead's own call, the fix is two lines, and the *asymmetry* behind it is worth a decision rather than another patch. ## The asymmetry `answer()` takes the target's lock before it touches anything: ```java MessageService.java:975 ReentrantLock lock = sessionLocks.computeIfAbsent(workerSession, _ -> new ReentrantLock()); ``` and then, holding it, calls `asyncTasksByTurn.get(turnId)` (:985), `clearAsyncQuestion(turnId, false)` (:994) and `finishAsyncTask(turnId, result)` (:1011). `ask()` takes **no** lock at all. I read the whole method — there is no `sessionLocks` acquisition anywhere in it. On its timeout path it mutates the same `Task`: ```java MessageService.java:933 markAskTimedOut(ticket.turnId()); MessageService.java:934 clearAsyncQuestion(ticket.turnId(), true); ``` `clearAsyncQuestion(_, true)` removes the task from `asyncTasksByTurn` and sets `task.turnId = null` (:1227-1228). So one side of this invariant is locked and the other is not. `Task.turnId`, `Task.question` and `Task.completedNanos` are `volatile`, so this is not a visibility bug. It is a compound-action bug. ## The concrete failure ```java MessageService.java:1234-1239 private void finishAsyncTask(Task task, Reply result) { task.future.complete(result); if (task.turnId != null) { asyncTasksByTurn.remove(task.turnId, task); } } ``` `task.turnId` is read **twice** — once for the null check, once for the removal. `volatile` makes each read fresh; it does not make the pair atomic. **Path in.** The lead answers a `fleet_ask` at the moment its ~55 s window lapses. 1. Lead's thread, in `answer()`, holding the lock: `finishAsyncTask(turnId, result)` (:1011) → `asyncTasksByTurn.get(turnId)` returns the task (:1243) → `finishAsyncTask(task, result)`. 2. It completes the future (:1235), then reads `task.turnId` and sees non-null. 3. Worker's thread, in `ask()`'s `catch (TimeoutException)`, holding **no** lock: `clearAsyncQuestion(ticket.turnId(), true)` sets `task.turnId = null` (:1228). 4. Lead's thread reads `task.turnId` again for the removal. It is now `null`. `ConcurrentHashMap.remove(Object, Object)` throws `NullPointerException` on a null key, by contract. **Direction of harm:** a thrown exception out of `answer()`, on the lead's own `fleet_send{turnId, content}` call. The lead is told its answer failed. It did not — `task.future.complete(result)` already ran on the line before, so the ticket is completed and the answer delivered. A lead that reacts by re-sending is acting on a false failure. The stale `asyncTasksByTurn` entry also leaks, and `hasAsyncQuestion` (:1251) scans that map to decide whether a target is BUSY, so the target can be reported busy after its task is done. **Confidence:** the code path is confirmed by reading. I have not reproduced the interleaving. ## What I want **Goal:** `answer()` must not be able to throw because `ask()` changed a field underneath it, and the two sides of this invariant must stop being guarded differently. **Invariants:** 1. **`ask()` must not start blocking on the target's session lock.** It is called from the worker's own turn and it already blocks for up to ~55 s on the answer. Making it contend with `send`/`answer` for the same lock risks a deadlock or a stalled turn. If your fix takes that lock in `ask()`, prove it cannot deadlock — do not assume it. 2. **`clearAsyncQuestion`'s `forgetTurn=true` forgetting stays.** #307's fix depends on it, and the javadoc explains why: it is what stops the target being reported BUSY forever. 3. **`markAskTimedOut` must still run before that forgetting.** Also #307. Its own javadoc says a lookup afterwards finds nothing. 4. **Do not change what `answer()` returns on success.** This ticket is about not throwing on a race, not about new semantics. **Candidate mechanism, as a candidate only:** read `task.turnId` once into a local in `finishAsyncTask` and use that local for both the check and the removal. **Decide it yourself and justify it.** That candidate fixes the crash and leaves the asymmetry. So answer this too, and it is the more important half: - **Is a local read enough, or is the unlocked `ask()` side wrong?** Look for other compound actions on a `Task` that one side performs locked and the other unlocked. `markAskTimedOut` (:1208-1213) does a get-then-check-then-set; `clearAsyncQuestion` (:1216-1230) does a get-then-check-then-two-writes; `answer()` does get-then-clear-then-finish across three calls. Say which of these are genuinely at risk and which are safe, and why. **Name the path in for each one you call a defect** — a compound action that no second thread can reach is not a defect. - If you conclude the right answer is a lock, or a different structure (for example making the task's turn a single immutable value that is swapped atomically rather than a field that is nulled), say so and say what it would cost. **Do not build it in this ticket** — report it and I will decide. A tested, reported deviation is a good outcome. ## Rules - Prove the crash with a test that fails without the fix. A latch or a test seam that nulls `turnId` between the two reads is the honest way; a test that merely calls the two methods in sequence proves nothing, because the bug is an interleaving. If you cannot build that seam without changing production code beyond the fix, say so and state exactly what your test does and does not prove. - Mutation proof required: revert the fix, quote the real failure output, restore it. - Never `git stash` — the stash is shared across every worktree here. - Never run `git worktree remove` or `git worktree prune` — other workers are live in those directories. - Stage files explicitly; never `git add -A`. Never merge. - Run `cd fleetd && mvn clean install` **unpiped**; quote the real `Tests run:` and `BUILD` lines. Never read `$?` after a pipe. - `fleetd/fleetd.yaml` is gitignored and absent from your worktree. Do not report on its contents. - Put your full report in the PR body as well as in your `fleet_reply`. ## Already ruled out, do not re-report The `strandedReplies` / `queuedDeliveries` pair was checked and dismissed: each touch is a single atomic map operation on a documented (CB-640) self-healing best-effort flag, not a compound action.
Author
Owner

Merged as 02e6aef (PR #327). Verified by me, not taken on the worker's word.

What I checked myself

The branch was behind main (git merge-base --is-ancestor said so), so I merged with --no-ff and rebuilt. Full build after the merge, unpiped:

Tests run: 1340, Failures: 0, Errors: 0, Skipped: 0
BUILD SUCCESS

Mutation P — restore the double read (full suite, not one test)

The worker ran its mutation against a single test. I ran the same idea against the whole suite, so I could see if anything else covers it:

Tests run: 1340, Failures: 0, Errors: 1
MessageServiceTest.finishAsyncTaskSurvivesTurnIdGoingNullBetweenItsTwoReads:984 » NullPointerException

Exactly one test pins it, and it is the new one. Nothing else covered this. The fix is real and pinned.

Mutation Q — drop the turnId != null guard

This one passed when it should have failed: 1340 tests green, zero failures, and zero NullPointerException anywhere in the log.

I did not accept that as "the guard is dead code". I re-ran with the guard dropped and a temporary log.error in sendAsync's executor catch:

PROBE-324Q swallowed by sendAsync executor
java.lang.NullPointerException: null
  at java.base/java.util.concurrent.ConcurrentHashMap.remove(ConcurrentHashMap.java:1566)
  at dev.ltms.fleet.msg.MessageService.finishAsyncTask(MessageService.java:1257)

19 real exceptions, on the ordinary path, and the suite stayed green (74/74). task.turnId is only ever set in markAsyncQuestion, so it is null for every async ticket whose worker never called fleet_ask. The guard is load-bearing on the plain path, and nothing pins it — because nothing can see it. That is finding F2 below.

Reachability — I re-derived it rather than trusting the report

The report says answerAsk can succeed while ask() still takes its timeout branch. I read ask() (:924-936) and confirm it: ticket.answer().get(timeoutMillis) can time out at the same instant answer() completes that future, so answerAsk returns true and ask() still runs markAskTimedOut + clearAsyncQuestion(turnId, true) under no lock. That is the window. The bug is reachable; it is not a paper defect.

Three findings that this fix does NOT close

F1 — the ticket is silently stranded (the report's item 4). I confirmed this by experiment, not by reading. I wrote a throwaway test using the merged forgetTurnForTest seam, firing ask()'s cleanup after answer() unblocked the worker but before the worker's real reply:

PROBE-ITEM4 phase=PENDING reply=null

The worker replied. answer() returned REPLIED. The async ticket stayed PENDING, with a null reply, until teardown would call abandon() and report the misleading WORKER_FAILED "session released before it replied". finishAsyncTask(String, Reply) does its own asyncTasksByTurn.get(turnId), and the local-read fix cannot help a caller that still has to look the task up. This is worse than #324 was: the crash made noise, this one is silent. The throwaway test is not in the merge.

F2 — sendAsync's executor catch is a silent sink. catch (Throwable t) { task.future.completeExceptionally(t); } runs after task.future.complete(result) has already succeeded, so completeExceptionally is a no-op and nothing logs. Any bug in that whole region is invisible to all 1340 tests. Measured above: 19 swallowed NPEs, no failure, no log line.

F3 — the same double read is still in reply() at :469-471 (if (orphan.turnId != null) { asyncTasksByTurn.remove(orphan.turnId, orphan); }). Structurally identical to what this ticket just fixed. Read-confirmed, not fixed.

F1, F2 and F3 go to their own ticket. F2 comes first: while it stands, F1 and the turnId != null guard cannot be pinned by any test, because the failure they cause is unobservable.

One correction to the report

None on the substance — the analysis held up on every point I re-checked. The report called item 4 "AT RISK"; the experiment above makes it confirmed, so I have raised it.

Note on the test seam

The fix adds a production volatile Runnable hook plus two package-private setters, null and unused outside tests. I accepted this: no public path reaches this interleaving without a real race, so a deterministic test needs a seam. It costs one volatile read per finishAsyncTask.

Closing.

Merged as `02e6aef` (PR #327). Verified by me, not taken on the worker's word. ## What I checked myself The branch was **behind main** (`git merge-base --is-ancestor` said so), so I merged with `--no-ff` and rebuilt. Full build after the merge, unpiped: ``` Tests run: 1340, Failures: 0, Errors: 0, Skipped: 0 BUILD SUCCESS ``` ### Mutation P — restore the double read (full suite, not one test) The worker ran its mutation against a single test. I ran the same idea against the whole suite, so I could see if anything else covers it: ``` Tests run: 1340, Failures: 0, Errors: 1 MessageServiceTest.finishAsyncTaskSurvivesTurnIdGoingNullBetweenItsTwoReads:984 » NullPointerException ``` Exactly one test pins it, and it is the new one. Nothing else covered this. The fix is real and pinned. ### Mutation Q — drop the `turnId != null` guard This one **passed when it should have failed**: 1340 tests green, zero failures, and zero `NullPointerException` anywhere in the log. I did not accept that as "the guard is dead code". I re-ran with the guard dropped **and** a temporary `log.error` in `sendAsync`'s executor `catch`: ``` PROBE-324Q swallowed by sendAsync executor java.lang.NullPointerException: null at java.base/java.util.concurrent.ConcurrentHashMap.remove(ConcurrentHashMap.java:1566) at dev.ltms.fleet.msg.MessageService.finishAsyncTask(MessageService.java:1257) ``` **19 real exceptions, on the ordinary path, and the suite stayed green (74/74).** `task.turnId` is only ever set in `markAsyncQuestion`, so it is `null` for every async ticket whose worker never called `fleet_ask`. The guard is load-bearing on the plain path, and nothing pins it — because nothing *can* see it. That is finding **F2** below. ## Reachability — I re-derived it rather than trusting the report The report says `answerAsk` can succeed while `ask()` still takes its timeout branch. I read `ask()` (:924-936) and confirm it: `ticket.answer().get(timeoutMillis)` can time out at the same instant `answer()` completes that future, so `answerAsk` returns `true` **and** `ask()` still runs `markAskTimedOut` + `clearAsyncQuestion(turnId, true)` under no lock. That is the window. The bug is reachable; it is not a paper defect. ## Three findings that this fix does NOT close **F1 — the ticket is silently stranded (the report's item 4). I confirmed this by experiment, not by reading.** I wrote a throwaway test using the merged `forgetTurnForTest` seam, firing `ask()`'s cleanup after `answer()` unblocked the worker but before the worker's real reply: ``` PROBE-ITEM4 phase=PENDING reply=null ``` The worker replied. `answer()` returned `REPLIED`. The async ticket stayed `PENDING`, with a `null` reply, until teardown would call `abandon()` and report the misleading `WORKER_FAILED` "session released before it replied". `finishAsyncTask(String, Reply)` does its **own** `asyncTasksByTurn.get(turnId)`, and the local-read fix cannot help a caller that still has to look the task up. **This is worse than #324 was:** the crash made noise, this one is silent. The throwaway test is not in the merge. **F2 — `sendAsync`'s executor `catch` is a silent sink.** `catch (Throwable t) { task.future.completeExceptionally(t); }` runs after `task.future.complete(result)` has already succeeded, so `completeExceptionally` is a no-op and nothing logs. Any bug in that whole region is invisible to all 1340 tests. Measured above: 19 swallowed NPEs, no failure, no log line. **F3 — the same double read is still in `reply()`** at :469-471 (`if (orphan.turnId != null) { asyncTasksByTurn.remove(orphan.turnId, orphan); }`). Structurally identical to what this ticket just fixed. Read-confirmed, not fixed. F1, F2 and F3 go to their own ticket. F2 comes first: while it stands, F1 and the `turnId != null` guard cannot be pinned by any test, because the failure they cause is unobservable. ## One correction to the report None on the substance — the analysis held up on every point I re-checked. The report called item 4 "AT RISK"; the experiment above makes it confirmed, so I have raised it. ## Note on the test seam The fix adds a production `volatile Runnable` hook plus two package-private setters, null and unused outside tests. I accepted this: no public path reaches this interleaving without a real race, so a deterministic test needs a seam. It costs one volatile read per `finishAsyncTask`. Closing.
ltms closed this issue 2026-09-04 09:52:56 +02:00
Sign in to join this conversation.
1 Participants
Notifications
Due Date
No due date set.
Dependencies

No dependencies set.

Reference: fleet/fleetd#324