sendAsync's executor catch swallows every exception thrown after the future completes, and that is why a stranded async ticket is invisible #329

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

Follow-up to #324. Three findings, all in MessageService, all measured by me during that merge. They are one unit because the first one is what hides the second.

F2 — the silent sink. Fix this one first.

MessageService.java:1071-1082 (sendAsync)
asyncExecutor.submit(() -> {
    try {
        Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task);
        if (result.outcome() == Outcome.QUESTION) {
            // Keep the accepted owner until answer() finishes it.
        } else {
            finishAsyncTask(task, result);
        }
    } catch (Throwable t) {
        task.future.completeExceptionally(t);
    }
});

finishAsyncTask calls task.future.complete(result) on its first line. So by the time anything later in that method throws, the future is already completed. completeExceptionally on a completed future returns false and does nothing. Nothing logs. The exception is gone.

Measured, not reasoned. I dropped the turnId != null guard in finishAsyncTask and ran the suite:

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

Zero NullPointerException in the log. I did not accept that. I re-ran with the guard dropped and a temporary log.error in that 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. 74/74 tests green. Not one log line.

Direction of harm: this is not one bug. It is a blind spot over the whole async region. Any defect that throws after the future completes produces a passing test suite and a silent daemon. That is how the turnId != null guard came to be load-bearing and completely unpinned — task.turnId is only ever set in markAsyncQuestion (:1186), so it is null for every async ticket whose worker never called fleet_ask, and removing the guard breaks every one of those with no visible sign.

Goal: an exception thrown inside that executor task must reach a log, whether or not the future was already completed.

Invariant: do not stop completing the future first. finishAsyncTask completing before it cleans up is what makes #324's crash harmless to the caller, and moving the completion later would lose an answer that was already delivered.

Candidate mechanism, as a candidate only: if task.future.completeExceptionally(t) returns false, log at error with the ticket and the target. Decide it yourself and justify it. If you find a better shape — for example not letting the cleanup throw at all — say so.

F1 — an async ticket is silently stranded. I confirmed this by experiment.

This is the report's item 4 from #324, which that PR reported but did not fix. I did not leave it as an argument. I wrote a throwaway test using the forgetTurnForTest seam that #324 merged, firing ask()'s own timeout cleanup after answer() unblocked the worker but before the worker's real reply arrived:

PROBE-ITEM4 phase=PENDING reply=null

The worker replied. answer() returned REPLIED. The async ticket stayed PENDING with a null reply.

Why the local-read fix cannot help. answer() calls finishAsyncTask(turnId, result) at :1011, and that overload does its own asyncTasksByTurn.get(turnId) at :1290. Reading a field once helps a caller that already holds the Task. It does nothing for a caller that must still look the task up — that lookup itself races ask()'s unlocked forgetting.

Path in. ask()'s ticket.answer().get(timeoutMillis) can time out at the same instant answer() completes that future. I read this myself at :924-936: rendezvous.answerAsk returns true and ask() still runs markAskTimedOut + clearAsyncQuestion(turnId, true) under no lock, which removes the asyncTasksByTurn entry. answer() then waits for the worker's real reply, gets it, and its lookup finds nothing.

Direction of harm: silent, and worse than #324 was. The crash made noise. Here fleet_poll{ticket} reports PENDING forever, and teardown's abandon() eventually resolves it as WORKER_FAILED — "session released before it replied" — long after the worker actually replied. A lead reading that is told a lie about its own worker. The askAnsweredAsyncTasks() recovery built for #137/#307 does not catch it: that only fires when reply() finds no live waiter, and here answer()'s waiter is live.

Goal: a worker's real reply must complete its async ticket, even when ask()'s unlocked cleanup removed the turn mapping first.

Invariants — all four from #324 still hold, plus one:

  1. ask() must not start blocking on the target's session lock. Prove no deadlock if you take it.
  2. clearAsyncQuestion's forgetTurn=true forgetting stays (#307 depends on it).
  3. markAskTimedOut must still run before that forgetting (#307).
  4. Do not change what answer() returns on success.
  5. The #282 chained-ask behaviour must keep working. A worker that calls a second fleet_ask inside a resumed turn moves its task to a fresh turnId, and answer()'s QUESTION guard at :1010 deliberately leaves the ticket open. A fix that completes the ticket whenever the lookup fails would break that — the lookup also legitimately fails in the chained case. Say how you tell the two apart.

Candidate mechanism, as a candidate only: answer() already holds a Task reference from its own lookup at :985. Passing that down instead of looking it up again is the obvious move. It is not obviously right — invariant 5 is exactly why. Decide it yourself and justify it.

F3 — the same double read #324 just fixed, still present

MessageService.java:468-471 (reply)
if (orphan.turnId != null) {
    asyncTasksByTurn.remove(orphan.turnId, orphan);
}

Two reads of the same volatile field, one for the check and one as the removal key, under no lock. Structurally identical to what #324 fixed at :1250. Read-confirmed by me; I have not reproduced it.

Fix it the same way. If you conclude the path in is unreachable here, say why — and remember that "unreachable" must name the thing that stops it, not just the absence of a known trigger.

Order

F2 first. While it stands, neither F1's fix nor the turnId != null guard can be pinned by any test, because the failure they cause is invisible. Once F2 logs, say in your report whether that changes how you tested F1.

Rules

  • Each of the three needs a test that fails without its fix. For F1, the forgetTurnForTest seam merged in #324 is already there — use it rather than building another.
  • Mutation proof required for each: revert the fix, quote the real failure output, restore it.
  • Run the full suite, not one test. #324's worker proved its fix against a single test; I re-ran it against all 1340 and learned that nothing else covered it. Do the same and report the whole-suite result.
  • 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.

Find what else has this shape — report it, do not fix it

F2 is "a catch that cannot report, because the thing it reports through is already resolved". Look for other places in this file where a catch writes to something that may already be terminal, or where a cleanup step runs after a future is completed. One line each, in your report. Do not go and fix them.

Already established, do not re-derive

  • Task.turnId, Task.question and Task.completedNanos are volatile. None of this is a visibility bug; all three findings are compound actions.
  • markAskTimedOut (:1208) and clearAsyncQuestion (:1216) were each checked in #324 and are safe in themselves — clearAsyncQuestion uses its turnId parameter throughout, never a fresh field read. clearAsyncQuestion is still the cause of the risk in its readers.
  • The strandedReplies / queuedDeliveries pair was checked and dismissed in #324: single atomic map operations on documented best-effort flags.
Follow-up to #324. Three findings, all in `MessageService`, all measured by me during that merge. They are one unit because the first one is what hides the second. ## F2 — the silent sink. Fix this one first. ```java MessageService.java:1071-1082 (sendAsync) asyncExecutor.submit(() -> { try { Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task); if (result.outcome() == Outcome.QUESTION) { // Keep the accepted owner until answer() finishes it. } else { finishAsyncTask(task, result); } } catch (Throwable t) { task.future.completeExceptionally(t); } }); ``` `finishAsyncTask` calls `task.future.complete(result)` on its **first line**. So by the time anything later in that method throws, the future is already completed. `completeExceptionally` on a completed future returns `false` and does nothing. Nothing logs. The exception is gone. **Measured, not reasoned.** I dropped the `turnId != null` guard in `finishAsyncTask` and ran the suite: ``` Tests run: 1340, Failures: 0, Errors: 0, Skipped: 0 BUILD SUCCESS ``` Zero `NullPointerException` in the log. I did not accept that. I re-ran with the guard dropped **and** a temporary `log.error` in that 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. 74/74 tests green. Not one log line.** **Direction of harm:** this is not one bug. It is a blind spot over the whole async region. Any defect that throws after the future completes produces a passing test suite and a silent daemon. That is how the `turnId != null` guard came to be load-bearing and completely unpinned — `task.turnId` is only ever set in `markAsyncQuestion` (:1186), so it is `null` for every async ticket whose worker never called `fleet_ask`, and removing the guard breaks every one of those with no visible sign. **Goal:** an exception thrown inside that executor task must reach a log, whether or not the future was already completed. **Invariant:** do not stop completing the future first. `finishAsyncTask` completing before it cleans up is what makes #324's crash harmless to the caller, and moving the completion later would lose an answer that was already delivered. **Candidate mechanism, as a candidate only:** if `task.future.completeExceptionally(t)` returns `false`, log at error with the ticket and the target. **Decide it yourself and justify it.** If you find a better shape — for example not letting the cleanup throw at all — say so. ## F1 — an async ticket is silently stranded. I confirmed this by experiment. This is the report's item 4 from #324, which that PR reported but did not fix. I did not leave it as an argument. I wrote a throwaway test using the `forgetTurnForTest` seam that #324 merged, firing `ask()`'s own timeout cleanup after `answer()` unblocked the worker but **before** the worker's real reply arrived: ``` PROBE-ITEM4 phase=PENDING reply=null ``` The worker replied. `answer()` returned `REPLIED`. The async ticket stayed `PENDING` with a `null` reply. **Why the local-read fix cannot help.** `answer()` calls `finishAsyncTask(turnId, result)` at :1011, and that overload does its **own** `asyncTasksByTurn.get(turnId)` at :1290. Reading a field once helps a caller that already holds the `Task`. It does nothing for a caller that must still look the task up — that lookup itself races `ask()`'s unlocked forgetting. **Path in.** `ask()`'s `ticket.answer().get(timeoutMillis)` can time out at the same instant `answer()` completes that future. I read this myself at :924-936: `rendezvous.answerAsk` returns `true` **and** `ask()` still runs `markAskTimedOut` + `clearAsyncQuestion(turnId, true)` under no lock, which removes the `asyncTasksByTurn` entry. `answer()` then waits for the worker's real reply, gets it, and its lookup finds nothing. **Direction of harm:** silent, and worse than #324 was. The crash made noise. Here `fleet_poll{ticket}` reports `PENDING` forever, and teardown's `abandon()` eventually resolves it as `WORKER_FAILED` — "session released before it replied" — long after the worker actually replied. A lead reading that is told a lie about its own worker. The `askAnsweredAsyncTasks()` recovery built for #137/#307 does not catch it: that only fires when `reply()` finds no live waiter, and here `answer()`'s waiter is live. **Goal:** a worker's real reply must complete its async ticket, even when `ask()`'s unlocked cleanup removed the turn mapping first. **Invariants — all four from #324 still hold, plus one:** 1. `ask()` must not start blocking on the target's session lock. Prove no deadlock if you take it. 2. `clearAsyncQuestion`'s `forgetTurn=true` forgetting stays (#307 depends on it). 3. `markAskTimedOut` must still run before that forgetting (#307). 4. Do not change what `answer()` returns on success. 5. **The #282 chained-ask behaviour must keep working.** A worker that calls a second `fleet_ask` inside a resumed turn moves its task to a fresh `turnId`, and `answer()`'s `QUESTION` guard at :1010 deliberately leaves the ticket open. A fix that completes the ticket whenever the lookup fails would break that — the lookup also legitimately fails in the chained case. **Say how you tell the two apart.** **Candidate mechanism, as a candidate only:** `answer()` already holds a `Task` reference from its own lookup at :985. Passing that down instead of looking it up again is the obvious move. It is not obviously right — invariant 5 is exactly why. **Decide it yourself and justify it.** ## F3 — the same double read #324 just fixed, still present ```java MessageService.java:468-471 (reply) if (orphan.turnId != null) { asyncTasksByTurn.remove(orphan.turnId, orphan); } ``` Two reads of the same volatile field, one for the check and one as the removal key, under no lock. Structurally identical to what #324 fixed at :1250. Read-confirmed by me; I have not reproduced it. Fix it the same way. If you conclude the path in is unreachable here, say why — and remember that "unreachable" must name the thing that stops it, not just the absence of a known trigger. ## Order F2 first. While it stands, neither F1's fix nor the `turnId != null` guard can be pinned by any test, because the failure they cause is invisible. Once F2 logs, say in your report whether that changes how you tested F1. ## Rules - Each of the three needs a test that fails without its fix. For F1, the `forgetTurnForTest` seam merged in #324 is already there — use it rather than building another. - Mutation proof required for each: revert the fix, quote the real failure output, restore it. - **Run the full suite, not one test.** #324's worker proved its fix against a single test; I re-ran it against all 1340 and learned that nothing else covered it. Do the same and report the whole-suite result. - 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`. ## Find what else has this shape — report it, do not fix it F2 is "a catch that cannot report, because the thing it reports through is already resolved". Look for other places in this file where a `catch` writes to something that may already be terminal, or where a cleanup step runs after a future is completed. One line each, in your report. Do not go and fix them. ## Already established, do not re-derive - `Task.turnId`, `Task.question` and `Task.completedNanos` are `volatile`. None of this is a visibility bug; all three findings are compound actions. - `markAskTimedOut` (:1208) and `clearAsyncQuestion` (:1216) were each checked in #324 and are safe **in themselves** — `clearAsyncQuestion` uses its `turnId` parameter throughout, never a fresh field read. `clearAsyncQuestion` is still the *cause* of the risk in its readers. - The `strandedReplies` / `queuedDeliveries` pair was checked and dismissed in #324: single atomic map operations on documented best-effort flags.
Author
Owner

Merged to main as 0c86503 (--no-ff; the branch was behind main, caught with
git merge-base --is-ancestor). Follow-up commit 4aa1fae corrects a comment — see below.

Build after the merge, unpiped: Tests run: 1351, Failures: 0, Errors: 0, Skipped: 0,
BUILD SUCCESS.

My own mutation

I ran one mutation of my own, different from the worker's, against the full suite.

Mutation T — drop F1's task != null guard in answer():

Tests run: 1351, Failures: 0, Errors: 4
FleetMcpTest.askThenAnswerRoundTrips
MessageServiceTest.answeringAnAskDoesNotRewriteDelegatorOwnership
MessageServiceTest.askSurfacesAsAQuestionAndTheAnswerResumesTheSameTurn
MessageServiceTest.duplicateAsksFromTheSameSessionCoalesceToOneTurn

All four: NullPointerException: Cannot read field "future" because "task" is null. The guard is
pinned on the blocking-ask path by four tests, not one. Reverted.

The fix is narrower than its comment claimed

The comment that landed with F1 said a null task means "this was never an async ticket". I checked
that and it is wrong. A genuine async ticket also arrives with task == null, because ask() drops
the asyncTasksByTurn entry in its catch (clearAsyncQuestion(turnId, true)) while
rendezvous.closeAsk(turnId) runs later in its finally. Between the two, the ask is still
answerable and the entry is gone.

A throwaway probe firing only that first half printed:

PROBE-GAP ask=ANSWERED
PROBE-GAP answer=REPLIED phase=PENDING reply=null

The stranded ticket this issue set out to fix, one step earlier in the same race. The probe used
forgetTurnForTest, so it omits markAskTimedOut; that cannot change the outcome, because
askTimedOut is read only by askAnsweredAsyncTasks, which reply() never reaches while
answer()'s own waiter is live. I have not raced the two real threads.

So F1 narrows this window, it does not close it. Opened as #334, and the comment in the code now
says so instead of claiming the guard is complete.

Same shape, reported by the worker, not fixed here

Three more sites in the same file where a failure has nowhere to go. Filed separately as #335:

  • whenComplete swallowing at :1100-1102
  • abandon's uncaught cleanup at :712-724
  • finally-overrides-return at :887-889 / :1057-1059

Closing. Three fixes in, one gap out (#334), three candidates out (#335).

Merged to `main` as `0c86503` (`--no-ff`; the branch was behind main, caught with `git merge-base --is-ancestor`). Follow-up commit `4aa1fae` corrects a comment — see below. Build after the merge, unpiped: `Tests run: 1351, Failures: 0, Errors: 0, Skipped: 0`, `BUILD SUCCESS`. ## My own mutation I ran one mutation of my own, different from the worker's, against the **full** suite. **Mutation T — drop F1's `task != null` guard** in `answer()`: ``` Tests run: 1351, Failures: 0, Errors: 4 FleetMcpTest.askThenAnswerRoundTrips MessageServiceTest.answeringAnAskDoesNotRewriteDelegatorOwnership MessageServiceTest.askSurfacesAsAQuestionAndTheAnswerResumesTheSameTurn MessageServiceTest.duplicateAsksFromTheSameSessionCoalesceToOneTurn ``` All four: `NullPointerException: Cannot read field "future" because "task" is null`. The guard is pinned on the blocking-ask path by four tests, not one. Reverted. ## The fix is narrower than its comment claimed The comment that landed with F1 said a `null` task means "this was never an async ticket". I checked that and it is wrong. A genuine async ticket also arrives with `task == null`, because `ask()` drops the `asyncTasksByTurn` entry in its **catch** (`clearAsyncQuestion(turnId, true)`) while `rendezvous.closeAsk(turnId)` runs later in its **finally**. Between the two, the ask is still answerable and the entry is gone. A throwaway probe firing only that first half printed: ``` PROBE-GAP ask=ANSWERED PROBE-GAP answer=REPLIED phase=PENDING reply=null ``` The stranded ticket this issue set out to fix, one step earlier in the same race. The probe used `forgetTurnForTest`, so it omits `markAskTimedOut`; that cannot change the outcome, because `askTimedOut` is read only by `askAnsweredAsyncTasks`, which `reply()` never reaches while `answer()`'s own waiter is live. I have not raced the two real threads. So **F1 narrows this window, it does not close it.** Opened as #334, and the comment in the code now says so instead of claiming the guard is complete. ## Same shape, reported by the worker, not fixed here Three more sites in the same file where a failure has nowhere to go. Filed separately as #335: - `whenComplete` swallowing at :1100-1102 - `abandon`'s uncaught cleanup at :712-724 - `finally`-overrides-return at :887-889 / :1057-1059 Closing. Three fixes in, one gap out (#334), three candidates out (#335).
ltms closed this issue 2026-09-04 10:35:24 +02:00
Sign in to join this conversation.
1 Participants
Notifications
Due Date
No due date set.
Dependencies

No dependencies set.

Reference: fleet/fleetd#329