#318: release() no longer strands a delivery that lands while it is running #321

Closed
agent wants to merge 0 commits from worker/fix-318-76ca36-9 into main
Member

Fixes #318.

The bug (confirmed by reading code, not observed live)

AmqpReplyInbox.release(target) did held.remove(target) under channelLock, then iterated the map it took and nacked every entry. deliverCallback inserts into held via computeIfAbsent with no channelLock. basicCancel stops new dispatches but does not flush a delivery already handed to the RabbitMQ client's consumer work-pool thread — so that thread can still call deliverCallback at any point during, or after, release()'s body. With the old code, computeIfAbsent would find the key gone (removed) and create a brand-new map under it — one release() has already returned past its remove and will never look at again. That entry sits delivered-but-unacked on the channel until the whole inbox closes: never requeued, never redelivered, and peek() is never called again for a released target. Commit be5ba22 (#298) is the dedicated fix for this method, but its contract test calls awaitPeek before release() runs, so it only proves the "already in held" case, never a concurrent arrival.

I checked the issue's write-up against the actual code and agree with it — I did not find anything wrong in it.

Mechanism chosen, and why

The issue offered three candidates. I picked the second one ("replace remove-then-iterate with a single atomic held.compute(target, ...) that both takes the map and leaves behind a tombstone computeIfAbsent will see") and I am not using the third ("re-check held once more at the end of release") — that one only narrows the window, it doesn't close it, and the issue is explicit that choosing it must be said plainly rather than called a fix. I did not choose it.

Why the tombstone approach actually closes the window rather than narrowing it: ConcurrentHashMap.compute() and computeIfAbsent() for the same key are mutually exclusive with each other in the JDK implementation — both go through the same per-bin locking path (including the reservation-node path used for a key's very first insert), so whichever of release()'s held.compute(target, (_, v) -> RELEASED) and a concurrent deliverCallback's held.computeIfAbsent(target, ...) runs first is fully visible to the other. There is no gap for a delivery to slip through, regardless of how much wall-clock time separates the swap from release() actually returning.

  • release(): swaps held[target] to a shared RELEASED sentinel via compute() instead of remove(), capturing whatever was previously held (via an AtomicReference written inside the remapping function) to nack exactly as before.
  • deliverCallback: after computeIfAbsent, checks perTarget == RELEASED (identity, never mutates the sentinel) and nacks-with-requeue immediately instead of holding the message — mirroring the existing duplicate-delivery ack, so this does not add a wait on channelLock across a broker round trip (invariant 3): basicNack, like the existing basicAck in the duplicate branch, does not await a broker reply.
  • peek() / ack(): treat the sentinel as "nothing held" (matches "already released" semantics; avoids ever calling a mutator on the shared sentinel object).
  • own(): clears a stale tombstone (held.remove(target, RELEASED)) when re-owning a target, so the fix can never permanently poison a target if its id is ever reused, and so held doesn't otherwise retain a tombstone forever for every target that has ever been released. The issue's own text says the id is never reused, so this line is defensive, not required by the issue — flagging it as a small addition beyond the strict ask, directly motivated by the tombstone this fix introduces.

Deadlock/serialization check for invariant 3: release()'s pattern is "hold channelLock, transiently take the CHM per-key lock nested inside for the compute() call, release it, keep going under channelLock." deliverCallback's pattern is "take the CHM per-key lock alone for computeIfAbsent, release it, then separately take channelLock." Neither thread ever holds one of these locks while blocked acquiring the other in the reverse order, so there's no lock-order cycle. It does mean a delivery for a released target can be made to wait behind an in-flight release()'s basicCancel/nack round trip for that one target — same cost the duplicate-ack path already pays, and it can't affect any other target's channelLock-free normal delivery path.

Invariant 2 (release() must not gain a new throw): held.compute() doesn't throw anything held.remove() didn't. Invariant 4 (duplicate path): untouched, still gated by the same perTarget.containsKey(msgId) check for the non-released case. Invariant 1 (never drop): every new nack path sets requeue=true; peek/ack never call a mutator on the sentinel.

Test — what it does and does not prove

AmqpReplyInboxReleaseRaceTest (new) forces the actual interleaving with a latch, not a sequence of calls: a fake Channel (java.lang.reflect.Proxy, no mocking library on the classpath — matches the existing AmqpReplyInboxRecoveryRaceTest pattern) blocks the first basicNack call. That call is only reachable from inside release()'s nack loop, which is only reached after held.compute(...) has already swapped in the tombstone — so observing the block is direct, ordering-guaranteed proof the tombstone is in place and release() has not yet returned (still holding channelLock). A second thread then fires the same DeliverCallback own() registered, for a brand-new message on the same target, while release() is still stuck there.

What it proves: a delivery whose computeIfAbsent call is ordered strictly after release()'s tombstone swap, while release() is still running, is nacked-with-requeue rather than silently parked forever.
What it does not prove: it does not drive a real broker — basicNack is a recorded call on the fake channel, not a verified requeue-and-redeliver. That half (a nacked-with-requeue message really does come back to a later owner) is already covered by AmqpReplyInboxContractTest.releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery against a real broker, which I ran (see below) and which still passes unmodified.

Mutation proof (reverted, ran, restored): I reverted AmqpReplyInbox.java to the pre-fix version (git show HEAD:... back over the working file) and ran the new test alone. Real failure output:

[ERROR] Tests run: 1, Failures: 1, Errors: 0, Skipped: 0, Time elapsed: 0.349 s <<< FAILURE! -- in dev.ltms.fleet.msg.AmqpReplyInboxReleaseRaceTest
org.opentest4j.AssertionFailedError: both the pre-held m0 and the concurrently-arriving m1 must be nacked, got: [[tag=1 requeue=1]] ==> expected: <2> but was: <1>
	at AmqpReplyInboxReleaseRaceTest.deliveryArrivingWhileReleaseIsStillRunningIsNackedNotStranded:129
[INFO] BUILD FAILURE

m0 (already held before release() ran) gets nacked; m1 (delivered while release() was still running, after held.remove executed) never does — exactly the reported defect. I then restored the fix (cp back the saved fixed file) and confirmed the full suite is green again before committing.

Build

cd fleetd && mvn clean install (unpiped, $? read directly — no tail/head/tee in the command that produced this number):

[INFO] Tests run: 1324, Failures: 0, Errors: 0, Skipped: 0
[INFO] BUILD SUCCESS

Contract suite against a real broker (Docker was available in this worktree, so I ran it rather than describing it as unrun): mvn test -Pcontract -Dtest=AmqpReplyInboxContractTest against a Testcontainers RabbitMQ:

[INFO] Tests run: 8, Failures: 0, Errors: 0, Skipped: 0
[INFO] BUILD SUCCESS

All 8 pass, including releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery (#298's own test) — the fix does not regress the case #298 already closed.

Shape check (msg/, one line each, not fixed)

Requested by the issue as a separate, non-blocking check — I'll add findings as a follow-up PR comment once a background check finishes; nothing in that scope was touched by this change.

Files changed

  • fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java — the fix (release, deliverCallback, own, peek, ack)
  • fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxReleaseRaceTest.java — new, forces the interleaving

Caveat for review

The own() tombstone-clearing line is slightly beyond the strict letter of the ticket (which only asked to fix release()/deliverCallback) — flagging it as a deliberate, small, closely-related addition rather than scope creep: without it, every released target's held map entry becomes a permanent (if cheap) tombstone for the life of the daemon, and if a target id were ever reused (the issue says it isn't) that target would be silently nacked forever. Happy to drop it if the reviewer would rather keep this PR to exactly the two methods named in the issue.

Fixes #318. ## The bug (confirmed by reading code, not observed live) `AmqpReplyInbox.release(target)` did `held.remove(target)` under `channelLock`, then iterated the map it took and nacked every entry. `deliverCallback` inserts into `held` via `computeIfAbsent` with **no** `channelLock`. `basicCancel` stops *new* dispatches but does not flush a delivery already handed to the RabbitMQ client's consumer work-pool thread — so that thread can still call `deliverCallback` at any point during, or after, `release()`'s body. With the old code, `computeIfAbsent` would find the key gone (removed) and create a **brand-new map** under it — one `release()` has already returned past its `remove` and will never look at again. That entry sits delivered-but-unacked on the channel until the whole inbox closes: never requeued, never redelivered, and `peek()` is never called again for a released target. Commit `be5ba22` (#298) is the dedicated fix for this method, but its contract test calls `awaitPeek` before `release()` runs, so it only proves the "already in `held`" case, never a concurrent arrival. I checked the issue's write-up against the actual code and agree with it — I did not find anything wrong in it. ## Mechanism chosen, and why The issue offered three candidates. I picked the second one ("replace remove-then-iterate with a single atomic `held.compute(target, ...)` that both takes the map and leaves behind a tombstone `computeIfAbsent` will see") and I am **not** using the third ("re-check `held` once more at the end of `release`") — that one only narrows the window, it doesn't close it, and the issue is explicit that choosing it must be said plainly rather than called a fix. I did not choose it. Why the tombstone approach actually **closes** the window rather than narrowing it: `ConcurrentHashMap.compute()` and `computeIfAbsent()` for the *same key* are mutually exclusive with each other in the JDK implementation — both go through the same per-bin locking path (including the reservation-node path used for a key's very first insert), so whichever of `release()`'s `held.compute(target, (_, v) -> RELEASED)` and a concurrent `deliverCallback`'s `held.computeIfAbsent(target, ...)` runs first is fully visible to the other. There is no gap for a delivery to slip through, regardless of how much wall-clock time separates the swap from `release()` actually returning. - `release()`: swaps `held[target]` to a shared `RELEASED` sentinel via `compute()` instead of `remove()`, capturing whatever was previously held (via an `AtomicReference` written inside the remapping function) to nack exactly as before. - `deliverCallback`: after `computeIfAbsent`, checks `perTarget == RELEASED` (identity, never mutates the sentinel) and nacks-with-requeue immediately instead of holding the message — mirroring the existing duplicate-delivery ack, so this does **not** add a wait on `channelLock` across a broker round trip (invariant 3): `basicNack`, like the existing `basicAck` in the duplicate branch, does not await a broker reply. - `peek()` / `ack()`: treat the sentinel as "nothing held" (matches "already released" semantics; avoids ever calling a mutator on the shared sentinel object). - `own()`: clears a stale tombstone (`held.remove(target, RELEASED)`) when re-owning a target, so the fix can never permanently poison a target if its id is ever reused, and so `held` doesn't otherwise retain a tombstone forever for every target that has ever been released. The issue's own text says the id is never reused, so this line is defensive, not required by the issue — flagging it as a small addition beyond the strict ask, directly motivated by the tombstone this fix introduces. Deadlock/serialization check for invariant 3: `release()`'s pattern is "hold `channelLock`, transiently take the CHM per-key lock nested inside for the `compute()` call, release it, keep going under `channelLock`." `deliverCallback`'s pattern is "take the CHM per-key lock alone for `computeIfAbsent`, release it, *then* separately take `channelLock`." Neither thread ever holds one of these locks while blocked acquiring the other in the reverse order, so there's no lock-order cycle. It does mean a delivery for a *released* target can be made to wait behind an in-flight `release()`'s `basicCancel`/nack round trip for **that one target** — same cost the duplicate-ack path already pays, and it can't affect any other target's channelLock-free normal delivery path. Invariant 2 (release() must not gain a new throw): `held.compute()` doesn't throw anything `held.remove()` didn't. Invariant 4 (duplicate path): untouched, still gated by the same `perTarget.containsKey(msgId)` check for the non-released case. Invariant 1 (never drop): every new nack path sets `requeue=true`; `peek`/`ack` never call a mutator on the sentinel. ## Test — what it does and does not prove `AmqpReplyInboxReleaseRaceTest` (new) forces the actual interleaving with a latch, not a sequence of calls: a fake `Channel` (`java.lang.reflect.Proxy`, no mocking library on the classpath — matches the existing `AmqpReplyInboxRecoveryRaceTest` pattern) blocks the *first* `basicNack` call. That call is only reachable from inside `release()`'s nack loop, which is only reached *after* `held.compute(...)` has already swapped in the tombstone — so observing the block is direct, ordering-guaranteed proof the tombstone is in place and `release()` has not yet returned (still holding `channelLock`). A second thread then fires the *same* `DeliverCallback` `own()` registered, for a brand-new message on the same target, while `release()` is still stuck there. **What it proves:** a delivery whose `computeIfAbsent` call is ordered strictly after `release()`'s tombstone swap, while `release()` is still running, is nacked-with-requeue rather than silently parked forever. **What it does not prove:** it does not drive a real broker — `basicNack` is a recorded call on the fake channel, not a verified requeue-and-redeliver. That half (a nacked-with-requeue message really does come back to a later owner) is already covered by `AmqpReplyInboxContractTest.releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery` against a real broker, which I ran (see below) and which still passes unmodified. **Mutation proof (reverted, ran, restored):** I reverted `AmqpReplyInbox.java` to the pre-fix version (`git show HEAD:...` back over the working file) and ran the new test alone. Real failure output: ``` [ERROR] Tests run: 1, Failures: 1, Errors: 0, Skipped: 0, Time elapsed: 0.349 s <<< FAILURE! -- in dev.ltms.fleet.msg.AmqpReplyInboxReleaseRaceTest org.opentest4j.AssertionFailedError: both the pre-held m0 and the concurrently-arriving m1 must be nacked, got: [[tag=1 requeue=1]] ==> expected: <2> but was: <1> at AmqpReplyInboxReleaseRaceTest.deliveryArrivingWhileReleaseIsStillRunningIsNackedNotStranded:129 [INFO] BUILD FAILURE ``` m0 (already held before `release()` ran) gets nacked; m1 (delivered while `release()` was still running, after `held.remove` executed) never does — exactly the reported defect. I then restored the fix (`cp` back the saved fixed file) and confirmed the full suite is green again before committing. ## Build `cd fleetd && mvn clean install` (unpiped, `$?` read directly — no `tail`/`head`/`tee` in the command that produced this number): ``` [INFO] Tests run: 1324, Failures: 0, Errors: 0, Skipped: 0 [INFO] BUILD SUCCESS ``` **Contract suite against a real broker** (Docker was available in this worktree, so I ran it rather than describing it as unrun): `mvn test -Pcontract -Dtest=AmqpReplyInboxContractTest` against a Testcontainers RabbitMQ: ``` [INFO] Tests run: 8, Failures: 0, Errors: 0, Skipped: 0 [INFO] BUILD SUCCESS ``` All 8 pass, including `releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery` (#298's own test) — the fix does not regress the case #298 already closed. ## Shape check (msg/, one line each, not fixed) Requested by the issue as a separate, non-blocking check — I'll add findings as a follow-up PR comment once a background check finishes; nothing in that scope was touched by this change. ## Files changed - `fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java` — the fix (`release`, `deliverCallback`, `own`, `peek`, `ack`) - `fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxReleaseRaceTest.java` — new, forces the interleaving ## Caveat for review The `own()` tombstone-clearing line is slightly beyond the strict letter of the ticket (which only asked to fix `release()`/`deliverCallback`) — flagging it as a deliberate, small, closely-related addition rather than scope creep: without it, every released target's `held` map entry becomes a permanent (if cheap) tombstone for the life of the daemon, and if a target id were ever reused (the issue says it isn't) that target would be silently nacked forever. Happy to drop it if the reviewer would rather keep this PR to exactly the two methods named in the issue.
agent added 1 commit 2026-09-04 09:11:55 +02:00
#318: release() no longer strands a delivery that lands while it is running
CI / build (pull_request) Failing after 1m59s
CI / contract (pull_request) Successful in 2m14s
2926cd1784
AmqpReplyInbox.release() used held.remove(target) then iterated the old
map. A delivery landing on the consumer work-pool thread after the
remove (basicCancel does not flush one already handed to that pool) hit
deliverCallback's computeIfAbsent, found the key gone, and created a
brand-new map release() never looks at again — delivered-but-unacked
forever, never requeued, never redelivered (#298 only closed the
"already in held when release runs" case).

Fix: release() swaps in a RELEASED tombstone via held.compute(...)
instead of held.remove(...). ConcurrentHashMap serializes compute/
computeIfAbsent calls for the same key against each other, so whichever
of release() and a concurrent deliverCallback runs first is fully
visible to the other — no gap. deliverCallback checks for the
tombstone and nacks-with-requeue instead of recreating a map; peek/ack
treat it as empty; own() clears a stale tombstone so a target is never
poisoned if its id is ever reused (the issue's own text says id reuse
doesn't happen, but the tombstone would otherwise sit in `held` forever
either way).

New test AmqpReplyInboxReleaseRaceTest forces the actual interleaving
with a latch (blocks release() inside its nack loop, which is only
reachable after the tombstone swap, then fires a concurrent delivery)
rather than a sequential call — a sequential test would not have caught
this, since #298's own contract test forces settlement before release()
runs. Mutation-tested: reverting the fix makes this test fail with
"expected: <2> but was: <1>" (m1 never nacked); restored after
confirming that failure.
Author
Member

Shape check (msg/, one line each, not fixed)

Ran a read-only sweep over src/main/java/dev/ltms/fleet/msg/ for the exact shape this ticket fixed: two halves of one invariant guarded by different locks, or a lock on one side and none on the other.

Same shape found — not fixed here:

  • MessageService.Task.turnId/.question (+ asyncTasksByTurn): ask() mutates them via markAsyncQuestion/clearAsyncQuestion/markAskTimedOut (MessageService.java:905,908,933,934) while holding no sessionLocks lock at all; answer() mutates the same task's fields via clearAsyncQuestion/finishAsyncTask (:994,1011) while holding sessionLocks.get(workerSession) — same key, two different synchronization regimes for what should be one consistent "this turn's question/turnId" state.
  • MessageService.asyncTasksByTurn via finishAsyncTask(Task,...): called from sendAsync's background virtual thread (:1079) with no lock (the session lock was already released when send() returned), vs. answer()'s call to the same method (:1011) with sessionLocks held — same asymmetry as above.
  • Lower severity, flagged only for completeness: MessageService.strandedReplies/queuedDeliveries — send() touches them under lock, reply()/abandon() touch them unlocked. Each individual touch is a single atomic map op on a self-healing, best-effort flag (documented CB-640 semantics), not a value that can be permanently stranded, so this one reads as lower risk than the two above.

Checked and found symmetric/safe: AmqpReplyInbox.held (this PR's fix — now consistent), LeadMailbox.held, InMemoryReplyInbox.store, Rendezvous.waiters/asks/openAsksBySession (lock-free by design via atomic ConcurrentHashMap ops throughout), ReplyPushLoop's pending* maps (same), LeadHeartbeatLoop.idleSinceNanos/.quietCount (single-threaded by construction).

Not fixing any of this — reporting per the ticket's ask.

## Shape check (msg/, one line each, not fixed) Ran a read-only sweep over `src/main/java/dev/ltms/fleet/msg/` for the exact shape this ticket fixed: two halves of one invariant guarded by different locks, or a lock on one side and none on the other. **Same shape found — not fixed here:** - `MessageService.Task.turnId`/`.question` (+ `asyncTasksByTurn`): `ask()` mutates them via `markAsyncQuestion`/`clearAsyncQuestion`/`markAskTimedOut` (`MessageService.java:905,908,933,934`) while holding **no** `sessionLocks` lock at all; `answer()` mutates the same task's fields via `clearAsyncQuestion`/`finishAsyncTask` (`:994,1011`) **while holding `sessionLocks.get(workerSession)`** — same key, two different synchronization regimes for what should be one consistent "this turn's question/turnId" state. - `MessageService.asyncTasksByTurn` via `finishAsyncTask(Task,...)`: called from `sendAsync`'s background virtual thread (`:1079`) with **no** lock (the session lock was already released when `send()` returned), vs. `answer()`'s call to the same method (`:1011`) **with** `sessionLocks` held — same asymmetry as above. - Lower severity, flagged only for completeness: `MessageService.strandedReplies`/`queuedDeliveries` — `send()` touches them under lock, `reply()`/`abandon()` touch them unlocked. Each individual touch is a single atomic map op on a self-healing, best-effort flag (documented CB-640 semantics), not a value that can be permanently stranded, so this one reads as lower risk than the two above. **Checked and found symmetric/safe:** `AmqpReplyInbox.held` (this PR's fix — now consistent), `LeadMailbox.held`, `InMemoryReplyInbox.store`, `Rendezvous.waiters`/`asks`/`openAsksBySession` (lock-free by design via atomic `ConcurrentHashMap` ops throughout), `ReplyPushLoop`'s pending* maps (same), `LeadHeartbeatLoop.idleSinceNanos`/`.quietCount` (single-threaded by construction). Not fixing any of this — reporting per the ticket's ask.
ltms closed this pull request 2026-09-04 09:24:18 +02:00
Some checks are pending
CI / build (pull_request) Failing after 1m59s
CI / contract (pull_request) Successful in 2m14s

Pull request closed

Sign in to join this conversation.