#324: read task.turnId once in finishAsyncTask; analysis of the wider unlocked-ask asymmetry #327

Closed
agent wants to merge 0 commits from worker/fix-324-3e9bbf-14 into main
Member

Fix (#324)

finishAsyncTask(Task, Reply) read the volatile Task.turnId field twice — once for the
null check, once as the ConcurrentHashMap.remove key. answer() calls this while holding
sessionLocks for the target; ask()'s own timeout path mutates the same field with no
lock at all
, via clearAsyncQuestion(turnId, true). volatile makes each read individually
fresh, but not the pair atomic — so the field can go null between the two reads and
asyncTasksByTurn.remove(null, task) throws NullPointerException on the lead's own
answer() call, even though task.future.complete(result) on the line above already ran (the
answer was in fact delivered).

Fix: capture task.turnId into a local once, use that local for both the check and the
removal.

private void finishAsyncTask(Task task, Reply result) {
    task.future.complete(result);
    String turnId = task.turnId;
    if (turnId != null) {
        asyncTasksByTurn.remove(turnId, task);
    }
}

Test — mutation-proof, interleaving forced deterministically

A real race between two independent threads (answer()'s lock-holding thread and ask()'s
unlocked timeout thread) can't be relied on to land in a specific few-instruction window on
every run. Calling the two methods in sequence proves nothing, since the bug is specifically
about what happens between two reads inside one method body.

So the test drives the real sequence (worker asks → primary answers → worker's real reply
arrives) and installs a package-private test hook (finishAsyncTaskRaceHook) that fires at
exactly the point between finishAsyncTask's former two reads. The hook runs
forgetTurnForTest(turnId), which is a thin wrapper that calls the identical production
cleanup ask()'s own timeout catch block runs — clearAsyncQuestion(turnId, true) — so the
test doesn't hand-roll an approximation of the race; it forces the real mutation at the real
moment.

Both the hook field and its setter/forwarder are null/no-ops in production and only
reachable package-privately from the test.

What the test proves: given that exact interleaving, answer() must not throw, and the
async ticket must still resolve to the worker's real reply.
What it does not prove: that the interleaving is reachable in production on its own
schedule — that is established by reading the code (see the analysis below), not by this
test, since forcing an interleaving via a hook is not the same as two independent threads
racing without help.

Mutation proof

Reverted the fix (restored the double read, keeping the hook call in the same relative
position so it still lands between the check and the use), ran only the new test:

[ERROR] Tests run: 1, Failures: 0, Errors: 1, Skipped: 0, Time elapsed: 0.240 s <<< FAILURE! -- in dev.ltms.fleet.msg.MessageServiceTest
[ERROR] dev.ltms.fleet.msg.MessageServiceTest.finishAsyncTaskSurvivesTurnIdGoingNullBetweenItsTwoReads -- Time elapsed: 0.066 s <<< ERROR!
java.util.concurrent.ExecutionException: java.lang.NullPointerException
	at java.base/java.util.concurrent.CompletableFuture.wrapInExecutionException(CompletableFuture.java:345)
	...
	at dev.ltms.fleet.msg.MessageServiceTest.finishAsyncTaskSurvivesTurnIdGoingNullBetweenItsTwoReads(MessageServiceTest.java:984)
Caused by: java.lang.NullPointerException
	at java.base/java.util.concurrent.ConcurrentHashMap.remove(ConcurrentHashMap.java:1566)
	at dev.ltms.fleet.msg.MessageService.finishAsyncTask(MessageService.java:1257)
	at dev.ltms.fleet.msg.MessageService.finishAsyncTask(MessageService.java:1293)
	at dev.ltms.fleet.msg.MessageService.answer(MessageService.java:1011)
	at dev.ltms.fleet.msg.MessageServiceTest.lambda$finishAsyncTaskSurvivesTurnIdGoingNullBetweenItsTwoReads$2(MessageServiceTest.java:978)

Then restored the fix and re-ran the full build — green (see below).

Analysis — the four named compound actions, plus one I found

The candidate fix (local read) stops the crash but leaves the asymmetry in place. Here is the
per-instance analysis the ticket asked for, with a named path in for every one I call a
defect, and a reason for every one I call safe.

1. finishAsyncTask(Task, Reply) (:1234 pre-fix) — DEFECT, confirmed. Path in: exactly as
in the issue. Fixed above.

2. markAskTimedOut (:1208-1213) — SAFE. It is a get-then-check-then-set
(task.askTimedOut = true), guarded by turnId.equals(task.turnId). Its only writer-side
caller is ask()'s own thread (including a coalesced duplicate caller on the same ticket,
per the method's own doc). Two concurrent callers racing on it produce, at worst, a
redundant identical write — the guard is symmetric and the write is idempotent, so no
interleaving of concurrent calls to this method produces wrong state. Nothing under
sessionLocks (send/answer) ever reads or writes askTimedOut, so there is no
locked-vs-unlocked asymmetry to exploit here the way there is for turnId. No reachable path
in.

3. clearAsyncQuestion (:1216-1230) — SAFE IN ITSELF, but it is the root cause of risk in
its readers.
Its own two writes (task.question = null, and — if forgetTurn —
asyncTasksByTurn.remove(turnId, task) + task.turnId = null) all use the parameter
turnId, never a fresh read of task.turnId, so its own critical section is internally
torn-read-free regardless of interleaving with concurrent callers (the guard again makes
concurrent duplicate calls harmless). But because it runs unlocked (called only from
ask()) and nulls task.turnId / removes the asyncTasksByTurn entry, it is exactly what
makes items 1 and 4 unsafe.

4. answer()'s get-then-clear-then-finish across three calls (:985, :994, :1011) — AT
RISK, and only partly addressed by the local-read fix.
Beyond the confirmed NPE (which
needs the removal to land inside finishAsyncTask's own two-read window), the same race
against ask()'s unlocked clearAsyncQuestion(turnId, true) has a second, quieter outcome
that the local-read fix does not touch:

  • finishAsyncTask(String turnId, Reply) (:1242, the overload answer() calls at :1011) does
    its own asyncTasksByTurn.get(turnId) lookup, independent of the earlier lookup at
    answer():985. If ask()'s unlocked cleanup removes the asyncTasksByTurn entry
    before this lookup runs (not during the two reads inside finishAsyncTask(Task, Reply), but strictly before it is even called), the lookup returns null, the two-arg
    overload silently no-ops, and task.future (the thing fleet_poll{ticket} watches) is
    never completed by this path.

  • Because ask()'s cleanup on its timeout path is just two lightweight method calls, it
    will typically finish well before the worker does any further work and calls its real
    fleet_reply — so this "entry already gone" outcome is, if anything, more likely than
    the narrow crash window, for the same underlying race.

  • The lead's own answer() call still returns correctly in this case — the worker's real
    reply resolves answer()'s own forward waiter (opened at :980) directly via
    rendezvous.resolve()'s fast path in reply(), which never reaches the
    askAnsweredAsyncTasks() recovery designed for #137/#307 (that recovery only fires when
    reply()'s fast path finds no live waiter — here there is one: answer()'s). So the
    existing askTimedOut-based recovery mechanism does not cover this case.

  • The async ticket is left open (question == null, future not done) until whatever
    eventually calls abandon() on the session (an explicit fleet_stop, or the idle
    reaper) — which then resolves it as a misleading Outcome.WORKER_FAILED ("session
    released before it replied"), long after the worker actually replied successfully.

    Path in: the primary answers a fleet_ask right as its own ~55s window is
    independently lapsing on the worker's side (the same race the issue names), and the
    worker finishes its resumed turn and calls its real fleet_reply quickly enough after
    being unblocked that answer()'s finishAsyncTask(turnId, result) call still happens
    before session teardown — which, given how little work ask()'s cleanup does, is the
    ordinary case for this race, not a corner of it.

5. Bonus finding (not one of the four named, same shape, found while reading the same
lines) — reply()'s consumption of askAnsweredAsyncTasks(), :465-471:

if (orphan.turnId != null) {
    asyncTasksByTurn.remove(orphan.turnId, orphan);
}

reads the same volatile orphan.turnId twice — once for the null check, once as the removal
key — with no lock, structurally identical to the original bug. Path in: a worker whose
fleet_ask is still open (turnId stamped, not yet answered) sends a stray/duplicate
fleet_reply that reaches reply()'s slow path (reachable right after resolveQuestion
closes the forward waiter and before answer() opens a fresh one, so no live waiter exists
for the fast path to catch it) at the same moment that same fleet_ask's own timeout
independently fires and ask()'s catch block calls clearAsyncQuestion(turnId, true),
nulling orphan.turnId in between reply()'s two reads. I have not fixed this — it's
outside this ticket's stated scope (finishAsyncTask only) and I was told to report,
not build. Flagging it here since it's the same pattern with a plausible path in.

Is a local read enough, or is the unlocked ask() side wrong?

The unlocked side is wrong, in a way the local read doesn't fully close. The local read
closes the loud failure (item 1's crash). It does not close the quiet one (item 4's
silently-uncompleted ticket) or the twin instance of the loud one (item 5). Both of those
exist because a second, independent lookup by turnId (not just a second read of the
already-fetched field) can race ask()'s unlocked forgetting — and "read once" only helps
a caller that already holds the Task reference; it can't help a caller that still has to
look the task up by turnId in the first place.

If I were building a structural fix (not doing so here, per the ticket): make "find the
task by turnId AND detach it" a single atomic operation against asyncTasksByTurn, and stop
treating Task.turnId as a field a second reader can safely re-read for a decision. Concretely,
something like asyncTasksByTurn.remove(turnId) (or computeIfPresent) being the only
place a task is ever matched-and-detached by turnId, with Task.turnId becoming
write-once/informational only. Cost: touches roughly 5-6 call sites in this file
(markAsyncQuestion, markAskTimedOut, clearAsyncQuestion, both finishAsyncTask
overloads, reply()'s recovery path, answer()), needs re-verification against every
existing #282 (chained-ask) and #307 (askTimedOut forgetting) test, and is exactly the kind
of "different structure" the ticket says not to build in this pass. I have not built it — I'm
reporting it and leaving ask()'s locking question (invariant 1: don't make it take
sessionLocks without proving no deadlock) untouched, since it's the same "do not build,
report" instruction.

Build

[INFO] Tests run: 74, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 5.442 s -- in dev.ltms.fleet.msg.MessageServiceTest
...
[INFO] Tests run: 1337, Failures: 0, Errors: 0, Skipped: 0
[INFO] BUILD SUCCESS

Run in the worker's own worktree with cd fleetd && mvn clean install, unpiped, full output read.

Files changed

  • fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java — the single-local-read fix in
    finishAsyncTask(Task, Reply), plus the package-private test seam
    (finishAsyncTaskRaceHook, setFinishAsyncTaskRaceHookForTest, forgetTurnForTest), all
    inert in production.
  • fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java — one new test,
    finishAsyncTaskSurvivesTurnIdGoingNullBetweenItsTwoReads.

Caveat for review

Items 4 (the silent-drop variant on answer()'s own path) and 5 (the twin bug in reply())
are not fixed in this PR — they're reported per the ticket's instruction not to build a
structural fix here. The lead should decide whether either warrants its own follow-up ticket.

## Fix (#324) `finishAsyncTask(Task, Reply)` read the volatile `Task.turnId` field twice — once for the null check, once as the `ConcurrentHashMap.remove` key. `answer()` calls this while holding `sessionLocks` for the target; `ask()`'s own timeout path mutates the same field with **no lock at all**, via `clearAsyncQuestion(turnId, true)`. `volatile` makes each read individually fresh, but not the pair atomic — so the field can go `null` between the two reads and `asyncTasksByTurn.remove(null, task)` throws `NullPointerException` on the lead's own `answer()` call, even though `task.future.complete(result)` on the line above already ran (the answer was in fact delivered). **Fix:** capture `task.turnId` into a local once, use that local for both the check and the removal. ```java private void finishAsyncTask(Task task, Reply result) { task.future.complete(result); String turnId = task.turnId; if (turnId != null) { asyncTasksByTurn.remove(turnId, task); } } ``` ## Test — mutation-proof, interleaving forced deterministically A real race between two independent threads (`answer()`'s lock-holding thread and `ask()`'s unlocked timeout thread) can't be relied on to land in a specific few-instruction window on every run. Calling the two methods in sequence proves nothing, since the bug is specifically about what happens *between* two reads inside one method body. So the test drives the real sequence (worker asks → primary answers → worker's real reply arrives) and installs a package-private test hook (`finishAsyncTaskRaceHook`) that fires at exactly the point between `finishAsyncTask`'s former two reads. The hook runs `forgetTurnForTest(turnId)`, which is a thin wrapper that calls the **identical** production cleanup `ask()`'s own timeout catch block runs — `clearAsyncQuestion(turnId, true)` — so the test doesn't hand-roll an approximation of the race; it forces the real mutation at the real moment. Both the hook field and its setter/forwarder are `null`/no-ops in production and only reachable package-privately from the test. **What the test proves:** given that exact interleaving, `answer()` must not throw, and the async ticket must still resolve to the worker's real reply. **What it does not prove:** that the interleaving is reachable in production on its own schedule — that is established by reading the code (see the analysis below), not by this test, since forcing an interleaving via a hook is not the same as two independent threads racing without help. ### Mutation proof Reverted the fix (restored the double read, keeping the hook call in the same relative position so it still lands between the check and the use), ran only the new test: ``` [ERROR] Tests run: 1, Failures: 0, Errors: 1, Skipped: 0, Time elapsed: 0.240 s <<< FAILURE! -- in dev.ltms.fleet.msg.MessageServiceTest [ERROR] dev.ltms.fleet.msg.MessageServiceTest.finishAsyncTaskSurvivesTurnIdGoingNullBetweenItsTwoReads -- Time elapsed: 0.066 s <<< ERROR! java.util.concurrent.ExecutionException: java.lang.NullPointerException at java.base/java.util.concurrent.CompletableFuture.wrapInExecutionException(CompletableFuture.java:345) ... at dev.ltms.fleet.msg.MessageServiceTest.finishAsyncTaskSurvivesTurnIdGoingNullBetweenItsTwoReads(MessageServiceTest.java:984) Caused by: java.lang.NullPointerException at java.base/java.util.concurrent.ConcurrentHashMap.remove(ConcurrentHashMap.java:1566) at dev.ltms.fleet.msg.MessageService.finishAsyncTask(MessageService.java:1257) at dev.ltms.fleet.msg.MessageService.finishAsyncTask(MessageService.java:1293) at dev.ltms.fleet.msg.MessageService.answer(MessageService.java:1011) at dev.ltms.fleet.msg.MessageServiceTest.lambda$finishAsyncTaskSurvivesTurnIdGoingNullBetweenItsTwoReads$2(MessageServiceTest.java:978) ``` Then restored the fix and re-ran the full build — green (see below). ## Analysis — the four named compound actions, plus one I found The candidate fix (local read) stops the crash but leaves the asymmetry in place. Here is the per-instance analysis the ticket asked for, with a named path in for every one I call a defect, and a reason for every one I call safe. **1. `finishAsyncTask(Task, Reply)` (:1234 pre-fix) — DEFECT, confirmed.** Path in: exactly as in the issue. Fixed above. **2. `markAskTimedOut` (:1208-1213) — SAFE.** It is a get-then-check-then-set (`task.askTimedOut = true`), guarded by `turnId.equals(task.turnId)`. Its only writer-side caller is `ask()`'s own thread (including a coalesced duplicate caller on the same ticket, per the method's own doc). Two concurrent callers racing on it produce, at worst, a redundant identical write — the guard is symmetric and the write is idempotent, so no interleaving of concurrent calls to this method produces wrong state. Nothing under `sessionLocks` (`send`/`answer`) ever reads or writes `askTimedOut`, so there is no locked-vs-unlocked asymmetry to exploit here the way there is for `turnId`. No reachable path in. **3. `clearAsyncQuestion` (:1216-1230) — SAFE IN ITSELF, but it is the root cause of risk in its readers.** Its own two writes (`task.question = null`, and — if `forgetTurn` — `asyncTasksByTurn.remove(turnId, task)` + `task.turnId = null`) all use the **parameter** `turnId`, never a fresh read of `task.turnId`, so its own critical section is internally torn-read-free regardless of interleaving with concurrent callers (the guard again makes concurrent duplicate calls harmless). But because it runs **unlocked** (called only from `ask()`) and nulls `task.turnId` / removes the `asyncTasksByTurn` entry, it is exactly what makes items 1 and 4 unsafe. **4. `answer()`'s get-then-clear-then-finish across three calls (:985, :994, :1011) — AT RISK, and only *partly* addressed by the local-read fix.** Beyond the confirmed NPE (which needs the removal to land inside `finishAsyncTask`'s own two-read window), the **same race** against `ask()`'s unlocked `clearAsyncQuestion(turnId, true)` has a second, quieter outcome that the local-read fix does **not** touch: - `finishAsyncTask(String turnId, Reply)` (:1242, the overload `answer()` calls at :1011) does its **own** `asyncTasksByTurn.get(turnId)` lookup, independent of the earlier lookup at `answer()`:985. If `ask()`'s unlocked cleanup removes the `asyncTasksByTurn` entry **before** this lookup runs (not *during* the two reads inside `finishAsyncTask(Task, Reply)`, but strictly *before* it is even called), the lookup returns `null`, the two-arg overload silently no-ops, and `task.future` (the thing `fleet_poll{ticket}` watches) is **never completed** by this path. - Because `ask()`'s cleanup on its timeout path is just two lightweight method calls, it will typically finish well before the worker does any further work and calls its real `fleet_reply` — so this "entry already gone" outcome is, if anything, **more likely** than the narrow crash window, for the same underlying race. - The lead's own `answer()` call still returns correctly in this case — the worker's real reply resolves `answer()`'s own forward waiter (opened at :980) directly via `rendezvous.resolve()`'s fast path in `reply()`, which never reaches the `askAnsweredAsyncTasks()` recovery designed for #137/#307 (that recovery only fires when `reply()`'s fast path finds **no** live waiter — here there is one: `answer()`'s). So the existing `askTimedOut`-based recovery mechanism does not cover this case. - The async ticket is left open (`question == null`, `future` not done) until whatever eventually calls `abandon()` on the session (an explicit `fleet_stop`, or the idle reaper) — which then resolves it as a misleading `Outcome.WORKER_FAILED` ("session released before it replied"), long after the worker actually replied successfully. **Path in:** the primary answers a `fleet_ask` right as its own ~55s window is independently lapsing on the worker's side (the same race the issue names), and the worker finishes its resumed turn and calls its real `fleet_reply` quickly enough after being unblocked that `answer()`'s `finishAsyncTask(turnId, result)` call still happens before session teardown — which, given how little work `ask()`'s cleanup does, is the ordinary case for this race, not a corner of it. **5. Bonus finding (not one of the four named, same shape, found while reading the same lines) — `reply()`'s consumption of `askAnsweredAsyncTasks()`, :465-471:** ```java if (orphan.turnId != null) { asyncTasksByTurn.remove(orphan.turnId, orphan); } ``` reads the same volatile `orphan.turnId` twice — once for the null check, once as the removal key — with no lock, structurally identical to the original bug. **Path in:** a worker whose `fleet_ask` is still open (turnId stamped, not yet answered) sends a stray/duplicate `fleet_reply` that reaches `reply()`'s slow path (reachable right after `resolveQuestion` closes the forward waiter and before `answer()` opens a fresh one, so no live waiter exists for the fast path to catch it) at the same moment that same `fleet_ask`'s own timeout independently fires and `ask()`'s catch block calls `clearAsyncQuestion(turnId, true)`, nulling `orphan.turnId` in between `reply()`'s two reads. I have **not** fixed this — it's outside this ticket's stated scope (`finishAsyncTask` only) and I was told to report, not build. Flagging it here since it's the same pattern with a plausible path in. ## Is a local read enough, or is the unlocked `ask()` side wrong? The unlocked side is wrong, in a way the local read doesn't fully close. The local read closes the loud failure (item 1's crash). It does not close the quiet one (item 4's silently-uncompleted ticket) or the twin instance of the loud one (item 5). Both of those exist because a *second, independent* lookup by `turnId` (not just a second read of the already-fetched field) can race `ask()`'s unlocked forgetting — and "read once" only helps a caller that already holds the `Task` reference; it can't help a caller that still has to look the task up by `turnId` in the first place. **If I were building a structural fix (not doing so here, per the ticket):** make "find the task by turnId AND detach it" a single atomic operation against `asyncTasksByTurn`, and stop treating `Task.turnId` as a field a second reader can safely re-read for a decision. Concretely, something like `asyncTasksByTurn.remove(turnId)` (or `computeIfPresent`) being the *only* place a task is ever matched-and-detached by turnId, with `Task.turnId` becoming write-once/informational only. Cost: touches roughly 5-6 call sites in this file (`markAsyncQuestion`, `markAskTimedOut`, `clearAsyncQuestion`, both `finishAsyncTask` overloads, `reply()`'s recovery path, `answer()`), needs re-verification against every existing #282 (chained-ask) and #307 (`askTimedOut` forgetting) test, and is exactly the kind of "different structure" the ticket says not to build in this pass. I have not built it — I'm reporting it and leaving `ask()`'s locking question (invariant 1: don't make it take `sessionLocks` without proving no deadlock) untouched, since it's the same "do not build, report" instruction. ## Build ``` [INFO] Tests run: 74, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 5.442 s -- in dev.ltms.fleet.msg.MessageServiceTest ... [INFO] Tests run: 1337, Failures: 0, Errors: 0, Skipped: 0 [INFO] BUILD SUCCESS ``` Run in the worker's own worktree with `cd fleetd && mvn clean install`, unpiped, full output read. ## Files changed - `fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java` — the single-local-read fix in `finishAsyncTask(Task, Reply)`, plus the package-private test seam (`finishAsyncTaskRaceHook`, `setFinishAsyncTaskRaceHookForTest`, `forgetTurnForTest`), all inert in production. - `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java` — one new test, `finishAsyncTaskSurvivesTurnIdGoingNullBetweenItsTwoReads`. ## Caveat for review Items 4 (the silent-drop variant on `answer()`'s own path) and 5 (the twin bug in `reply()`) are **not fixed** in this PR — they're reported per the ticket's instruction not to build a structural fix here. The lead should decide whether either warrants its own follow-up ticket.
agent added 1 commit 2026-09-04 09:40:03 +02:00
#324: read task.turnId once in finishAsyncTask to stop an NPE from ask()'s unlocked forgetting
CI / contract (pull_request) Successful in 1m28s
CI / build (pull_request) Successful in 2m1s
e5cb51a90e
answer() holds sessionLocks while finishAsyncTask reads the volatile Task.turnId twice — once to
check it is non-null, once as the ConcurrentHashMap.remove key. ask()'s own timeout path mutates
the same field with no lock, via clearAsyncQuestion(turnId, true). volatile makes each read fresh
but not the pair atomic, so the field can go null between the two reads and remove(null, task)
throws NullPointerException on the lead's own answer() call, even though the reply already
completed on the line above.

Capture task.turnId into a local once and use that for both the check and the removal.

Added a package-private test seam (finishAsyncTaskRaceHook + forgetTurnForTest) so a test can force
the exact interleaving deterministically, by running the identical clearAsyncQuestion(turnId, true)
cleanup ask() uses, at the point between finishAsyncTask's former two reads. Both are inert (null)
in production.
ltms closed this pull request 2026-09-04 09:53:02 +02:00
Some checks are pending
CI / contract (pull_request) Successful in 1m28s
CI / build (pull_request) Successful in 2m1s

Pull request closed

Sign in to join this conversation.