fleetd#335: abandon() and sendAsync must not swallow per-task cleanup throws #351

Closed
agent wants to merge 0 commits from worker/fd335-a71c35-1 into main
Member

fleetd #335 — MessageService: three more places a thrown exception has nowhere to go.

Per-site verdict (read the actual implementations, not just the call sites)

Site 1 — abandon()'s loop over matching (reachable). The loop's recovery/"put back"
branch calls inbox.publish(...). AmqpReplyInbox.publish performs a real broker round trip
under publisherConfirm and throws IllegalStateException on an unroutable, unconfirmed, or
interrupted publish (AmqpReplyInbox.java:359/368/372). An uncaught throw there aborted the
loop, so every task after the throwing one in matching iteration order was left PENDING
forever with no other teardown path — the exact strand #335 describes.

The primary branch (asyncTasksByTurn.remove(...) + rendezvous.closeAsk(turnId)) is
different: both are plain ConcurrentHashMap operations on a non-null key with no
user-overridable code (Rendezvous.closeAsk at Rendezvous.java:199 just does asks.get /
openAsksBySession.remove / asks.remove), so neither can actually throw today. I left that
half unwrapped in spirit — the fix wraps the whole per-task cleanup so it defends the one call
that really can throw, without pretending the map ops need a catch of their own.

Fixed by decoupling "this task gets its own outcome" from "this task's best-effort cleanup
succeeds": task.future.complete(outcome) (and the resulting asyncFailed bookkeeping) now runs
unconditionally, before the try/catch around cleanup, so a thrown cleanup only loses that one
task's bookkeeping — logged via log.error — and the loop still reaches every remaining task.

Reaching the throwing branch by real timing needs a race (a second completion of the same task's
future landing between abandon() choosing it as recoveryTask and the loop reaching it) that
this file's own #137 follow-up investigation already found unreachable through the public API
today, not merely hard to time
(see the comment block above
abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer in MessageServiceTest). So the reproducing
test uses a test-only hook (abandonCleanupHookForTest) to inject the throw at that exact point —
the same technique this file already uses for the equally-narrow #324/#329 races
(finishAsyncTaskRaceHook, afterFinishAsyncTaskCompleteHookForTest).

Site 2 — sendAsync's task.future.whenComplete(...) (reachable). The returned stage is
discarded, so an uncaught throw from pushLoop.onTicketTerminal used to vanish with no log line
and no metric — the exact thing that hook exists to prevent. This one reproduces for real, no hook
needed: Fleetd's shutdown hook runs messages.close() (stops the async executor from taking
new work, but does not cancel a send already in flight) before pushLoop.close() (shuts its
scheduler down immediately via shutdownNow()). A ticket that completes in that narrow window
makes onTicketTerminal's own scheduler.schedule(...) (ReplyPushLoop.java:710) throw a genuine
RejectedExecutionException. Fixed with try/catch(Throwable) + log.error inside the
whenComplete action.

Site 3 — the two finally { asyncTasksByWaiter.remove(reply); rendezvous.close(target, reply); }
blocks in send() and answer() (not a defect).
asyncTasksByWaiter is a
ConcurrentHashMap; Rendezvous.close (Rendezvous.java:107) is waiters.remove(session, waiter) — again a plain map op on a non-null key, no user code involved. Neither call can throw
with the real implementations. Recommend closing as not-a-defect.

Mutation proof

  • Reverting site 1's try/catch: abandon() now propagates the injected exception straight out —
    MessageServiceTest#aPerTaskCleanupFailureDoesNotStrandTheRemainingMatchingTasks errors
    (RuntimeException: PROBE-335-SITE1 at MessageService.abandon).
  • Reverting site 2's try/catch: the RejectedExecutionException still fires (proving the trigger
    is real) but is no longer logged —
    MessageServiceTest#aTicketTerminalPushFailureDoesNotVanishSilently fails its log assertion
    (expected: <true> but was: <false>).

Both mutations were applied, run, observed to fail with the exact test above, then reverted before
this PR.

Build

cd fleetd && mvn clean install (unpiped): Tests run: 1365, Failures: 0, Errors: 0, Skipped: 0,
BUILD SUCCESS.

Files changed

  • fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java
  • fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java

No caveats for review beyond the site-3 verdict above (recommend closing that part of the ticket
as not-a-defect rather than merging a catch that can never fire).

fleetd #335 — MessageService: three more places a thrown exception has nowhere to go. ## Per-site verdict (read the actual implementations, not just the call sites) **Site 1 — `abandon()`'s loop over `matching` (reachable).** The loop's recovery/"put back" branch calls `inbox.publish(...)`. `AmqpReplyInbox.publish` performs a real broker round trip under `publisherConfirm` and throws `IllegalStateException` on an unroutable, unconfirmed, or interrupted publish (`AmqpReplyInbox.java:359/368/372`). An uncaught throw there aborted the loop, so every task after the throwing one in `matching` iteration order was left `PENDING` forever with no other teardown path — the exact strand `#335` describes. The *primary* branch (`asyncTasksByTurn.remove(...)` + `rendezvous.closeAsk(turnId)`) is different: both are plain `ConcurrentHashMap` operations on a non-null key with no user-overridable code (`Rendezvous.closeAsk` at `Rendezvous.java:199` just does `asks.get` / `openAsksBySession.remove` / `asks.remove`), so neither can actually throw today. I left that half unwrapped in spirit — the fix wraps the whole per-task cleanup so it defends the one call that really can throw, without pretending the map ops need a catch of their own. Fixed by decoupling "this task gets its own outcome" from "this task's best-effort cleanup succeeds": `task.future.complete(outcome)` (and the resulting `asyncFailed` bookkeeping) now runs unconditionally, *before* the try/catch around cleanup, so a thrown cleanup only loses that one task's bookkeeping — logged via `log.error` — and the loop still reaches every remaining task. Reaching the throwing branch by real timing needs a race (a second completion of the *same* task's future landing between `abandon()` choosing it as `recoveryTask` and the loop reaching it) that this file's own `#137` follow-up investigation already found **unreachable through the public API today, not merely hard to time** (see the comment block above `abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer` in `MessageServiceTest`). So the reproducing test uses a test-only hook (`abandonCleanupHookForTest`) to inject the throw at that exact point — the same technique this file already uses for the equally-narrow `#324`/`#329` races (`finishAsyncTaskRaceHook`, `afterFinishAsyncTaskCompleteHookForTest`). **Site 2 — `sendAsync`'s `task.future.whenComplete(...)` (reachable).** The returned stage is discarded, so an uncaught throw from `pushLoop.onTicketTerminal` used to vanish with no log line and no metric — the exact thing that hook exists to prevent. This one reproduces for real, no hook needed: `Fleetd`'s shutdown hook runs `messages.close()` (stops the async executor from taking *new* work, but does not cancel a send already in flight) before `pushLoop.close()` (shuts its scheduler down immediately via `shutdownNow()`). A ticket that completes in that narrow window makes `onTicketTerminal`'s own `scheduler.schedule(...)` (`ReplyPushLoop.java:710`) throw a genuine `RejectedExecutionException`. Fixed with `try/catch(Throwable)` + `log.error` inside the `whenComplete` action. **Site 3 — the two `finally { asyncTasksByWaiter.remove(reply); rendezvous.close(target, reply); }` blocks in `send()` and `answer()` (not a defect).** `asyncTasksByWaiter` is a `ConcurrentHashMap`; `Rendezvous.close` (`Rendezvous.java:107`) is `waiters.remove(session, waiter)` — again a plain map op on a non-null key, no user code involved. Neither call can throw with the real implementations. Recommend closing as not-a-defect. ## Mutation proof - Reverting site 1's try/catch: `abandon()` now propagates the injected exception straight out — `MessageServiceTest#aPerTaskCleanupFailureDoesNotStrandTheRemainingMatchingTasks` **errors** (`RuntimeException: PROBE-335-SITE1` at `MessageService.abandon`). - Reverting site 2's try/catch: the `RejectedExecutionException` still fires (proving the trigger is real) but is no longer logged — `MessageServiceTest#aTicketTerminalPushFailureDoesNotVanishSilently` **fails** its log assertion (`expected: <true> but was: <false>`). Both mutations were applied, run, observed to fail with the exact test above, then reverted before this PR. ## Build `cd fleetd && mvn clean install` (unpiped): `Tests run: 1365, Failures: 0, Errors: 0, Skipped: 0`, `BUILD SUCCESS`. ## Files changed - `fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java` - `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java` No caveats for review beyond the site-3 verdict above (recommend closing that part of the ticket as not-a-defect rather than merging a catch that can never fire).
agent added 1 commit 2026-09-04 11:38:50 +02:00
fleetd#335: abandon()'s per-task cleanup and sendAsync's terminal hook must not swallow throws
CI / contract (pull_request) Successful in 1m2s
CI / build (pull_request) Successful in 2m21s
887aca0183
Site 1 (abandon()'s matching loop, reachable): the recovery/put-back branch calls
inbox.publish, which AmqpReplyInbox implements as a real broker round trip that
throws IllegalStateException on an unroutable/unconfirmed/interrupted publish.
An uncaught throw there aborted the loop, stranding every task after it in
`matching` PENDING forever. Fixed by recording each task's own future.complete()
result before any cleanup runs, then wrapping the cleanup in try/catch so one
task's failure cannot stop its siblings from getting their outcome. Reaching the
throwing branch by real timing needs a race the file's own #137 follow-up already
found unreachable through the public API, so the reproducing test uses a
test-only hook (same technique as the existing fleetd #324/#329 hooks) to inject
the throw at that exact point.

Site 2 (sendAsync's task.future.whenComplete, reachable): the returned stage is
discarded, so an uncaught throw from pushLoop.onTicketTerminal vanished with no
log line. Reproduced for real: Fleetd's shutdown hook runs messages.close()
(stops the async executor from taking new work, but does not cancel a send
already in flight) before pushLoop.close() (shuts its scheduler down
immediately) — a ticket completing in that window makes onTicketTerminal's own
scheduler.schedule(...) throw a genuine RejectedExecutionException. Fixed with a
try/catch(Throwable) plus log.error inside the whenComplete action.

Site 3 (the two `finally { asyncTasksByWaiter.remove(reply); rendezvous.close(...);
}` blocks in send() and answer()): read Rendezvous.close/closeAsk and the
ConcurrentHashMap operations behind them — both are plain map ops on a non-null
key with no user-overridable code, so neither can throw. Left unchanged; not a
defect.

Mutation-proven: reverting either fix reproduces the failure it exists to catch
— removing site 1's try/catch aborts abandon() with the injected exception
(MessageServiceTest#aPerTaskCleanupFailureDoesNotStrandTheRemainingMatchingTasks
errors); removing site 2's try/catch leaves the RejectedExecutionException
unlogged (MessageServiceTest#aTicketTerminalPushFailureDoesNotVanishSilently
fails its log assertion). Full suite: mvn clean install, Tests run: 1365,
Failures: 0, Errors: 0, BUILD SUCCESS.
ltms closed this pull request 2026-09-04 11:49:28 +02:00
Some checks are pending
CI / contract (pull_request) Successful in 1m2s
CI / build (pull_request) Successful in 2m21s

Pull request closed

Sign in to join this conversation.