A reply delivered while release() is running is stranded unacked forever: #298 closed "already held", not "arrives during release" #318

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

Found by a delegated hunter. I read both methods and confirmed the interleaving.

Fourth instance today of "a one-way gate is not a gate", and the second where the earlier fix closed exactly the direction its own incident came from.

The two halves do not share a lock

AmqpReplyInbox.release — everything under channelLock:

synchronized (channelLock) {
    String tag = consumerTags.remove(target);
    if (tag != null) { channel.basicCancel(tag); }
    var perTarget = held.remove(target);
    if (perTarget != null) {
        synchronized (perTarget) {
            for (Held h : perTarget.values()) {
                channel.basicNack(h.deliveryTag(), false, true); // requeue, don't drop
            }
        }
    }
}

deliverCallback — the insertion into held takes no channelLock:

var perTarget = held.computeIfAbsent(target, _ -> new LinkedHashMap<>());
boolean duplicate;
synchronized (perTarget) {
    if (perTarget.containsKey(msgId)) { duplicate = true; }
    else { perTarget.put(msgId, new Held(tag, new InboxMessage(msgId, target, content))); duplicate = false; }
}

It only takes channelLock in the duplicate branch, for the ack.

The path in

  1. The idle reaper (or fleet_stop) releases a worker. MessageService.abandon runs first and finds nothing to recover — the message has not arrived yet. Then replyInbox.release(target).
  2. Concurrently the worker calls fleet_reply. There is no live rendezvous waiter, so it falls to inbox.publish(...), which uses the separate publish channel and succeeds.
  3. The broker dispatches to the target's consumer before basicCancel takes effect. The delivery lands in the client's consumer work pool — a separate thread. basicCancel stops new dispatches; it does not flush ones already handed to the pool. That is precisely why release bothers to requeue at all, and the method's own javadoc says so.
  4. release has already run held.remove(target) and is iterating the old map. The work-pool thread now runs computeIfAbsent, finds nothing, and creates a new map under the same key.
  5. release finishes iterating the old map and returns. The entry in the new map is never nacked.

Result: the message stays delivered-but-unacked on a channel that lives until the whole inbox closes. It is never requeued, so it is never redelivered. peek(target) would still see it, but nothing peeks a released target again — the session is gone from the roster and its terminal id is never reused. The reply is invisible until the daemon restarts.

Direction of harm: silent loss. A worker's real report disappears with no log line naming it. This is the same end state as fleetd #307, reached by a different route.

Why the existing fix does not cover it

Commit be5ba22 (#298, "release() requeues held deliveries instead of dropping them") is the dedicated fix for this method. Its contract test, AmqpReplyInboxContractTest.releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery, calls awaitPeek to force the delivery to settle before release runs. So it proves the "already in held" case and never exercises a delivery landing concurrently. The gate closed the direction the incident came from.

Honest limits

Neither the hunter nor I reproduced this. It is confirmed by reading the two methods and the RabbitMQ client's dispatch model, not observed. The window is one work-pool thread against one release, so it needs a reply landing in roughly the same moment as a teardown. That is not exotic — the idle reaper fires on a timer with no regard for what a worker is doing — but it is narrow, and I am not claiming it has cost us a report.

What I want

Goal: a delivery that arrives for a target at any point during or after its release must not end up sitting unacked in a map nobody will ever read. Either it is requeued so a later session can get it, or it is refused — but it must not be silently retained.

Invariants:

  1. Never drop a message to close this. Requeue or leave it unacked deliberately; do not basicAck a reply nobody has read. #298 exists because dropping was the old behaviour.
  2. release must not become able to throw where it currently does not. Its javadoc is explicit that a failed basicNack is logged and release proceeds, because release is called from teardown paths that must complete. A failed basicCancel still throws — leave that as it is.
  3. Do not hold channelLock across a blocking broker round trip from the delivery thread. basicCancel and the nack loop already run under it; adding a delivery-side wait on the same lock risks serialising every delivery behind a teardown, or worse. If your fix takes that lock in deliverCallback, say why it cannot deadlock or stall the consumer pool.
  4. The duplicate-delivery path must keep working. A redelivered message that is already in held is acked and dropped so the broker stops resending it. Do not break that.

Candidate mechanisms, all candidates only — pick one and justify it:

  • Have deliverCallback consult a "this target is released" marker and nack the delivery immediately instead of holding it. The marker has to be set and read consistently with held.remove, which is the whole difficulty.
  • Replace the remove-then-iterate with a single atomic held.compute(target, ...) that both takes the map and leaves behind a tombstone computeIfAbsent will see.
  • Re-check held for the target once more at the end of release and nack anything that appeared. Simplest, and it narrows the window rather than closing it — if you choose this, say so plainly rather than describing it as a fix.

Decide it yourself and justify it. I do not know which is right. If you find a fourth option, or find that one of these cannot work, that is a good outcome — say so.

Rules

  • A test that only calls the two methods in sequence proves nothing here. The bug is an interleaving. Force it: a test double or a latch that lets a delivery land after held.remove has run but before release returns. If you cannot build that seam without changing production code beyond the fix, say so and describe exactly what your test does and does not prove.
  • If the AMQP contract suite can run in your worktree (it needs a broker), run it and quote the real output. If it cannot, say so plainly — do not describe an unrun suite as passing. Some claims in this area only show up against a real broker: fleetd #298's first attempted fix was correct-looking and wrong for exactly that reason.
  • Mutation proof required: revert the fix, quote the real failure output, restore it.
  • Never run git stash — the stash is shared across every worktree here.
  • Never run git worktree remove or git worktree prune — other workers are live in these worktrees.
  • Stage files explicitly; never git add -A. Never merge.
  • Run cd fleetd && mvn clean install unpiped, and quote the real Tests run: and BUILD lines. Never pipe maven through tail/head, and never read $? after a pipe.
  • Put your full report in the PR body as well as in your fleet_reply.

Shape check

When done, look in msg/ only for the same shape: two halves of one invariant guarded by different locks, or by a lock on one side and none on the other. One line each, do not fix any of it.

Found by a delegated hunter. I read both methods and confirmed the interleaving. Fourth instance today of "a one-way gate is not a gate", and the second where the earlier fix closed exactly the direction its own incident came from. ## The two halves do not share a lock `AmqpReplyInbox.release` — everything under `channelLock`: ```java synchronized (channelLock) { String tag = consumerTags.remove(target); if (tag != null) { channel.basicCancel(tag); } var perTarget = held.remove(target); if (perTarget != null) { synchronized (perTarget) { for (Held h : perTarget.values()) { channel.basicNack(h.deliveryTag(), false, true); // requeue, don't drop } } } } ``` `deliverCallback` — the insertion into `held` takes **no** `channelLock`: ```java var perTarget = held.computeIfAbsent(target, _ -> new LinkedHashMap<>()); boolean duplicate; synchronized (perTarget) { if (perTarget.containsKey(msgId)) { duplicate = true; } else { perTarget.put(msgId, new Held(tag, new InboxMessage(msgId, target, content))); duplicate = false; } } ``` It only takes `channelLock` in the `duplicate` branch, for the ack. ## The path in 1. The idle reaper (or `fleet_stop`) releases a worker. `MessageService.abandon` runs first and finds nothing to recover — the message has not arrived yet. Then `replyInbox.release(target)`. 2. Concurrently the worker calls `fleet_reply`. There is no live rendezvous waiter, so it falls to `inbox.publish(...)`, which uses the separate publish channel and succeeds. 3. The broker dispatches to the target's consumer before `basicCancel` takes effect. The delivery lands in the client's consumer work pool — a separate thread. `basicCancel` stops *new* dispatches; it does not flush ones already handed to the pool. That is precisely why `release` bothers to requeue at all, and the method's own javadoc says so. 4. `release` has already run `held.remove(target)` and is iterating the **old** map. The work-pool thread now runs `computeIfAbsent`, finds nothing, and creates a **new** map under the same key. 5. `release` finishes iterating the old map and returns. The entry in the new map is never nacked. **Result:** the message stays delivered-but-unacked on a channel that lives until the whole inbox closes. It is never requeued, so it is never redelivered. `peek(target)` would still see it, but nothing peeks a released target again — the session is gone from the roster and its terminal id is never reused. The reply is invisible until the daemon restarts. **Direction of harm:** silent loss. A worker's real report disappears with no log line naming it. This is the same end state as fleetd #307, reached by a different route. ## Why the existing fix does not cover it Commit `be5ba22` (#298, *"release() requeues held deliveries instead of dropping them"*) is the dedicated fix for this method. Its contract test, `AmqpReplyInboxContractTest.releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery`, calls `awaitPeek` to force the delivery to settle **before** `release` runs. So it proves the "already in `held`" case and never exercises a delivery landing concurrently. The gate closed the direction the incident came from. ## Honest limits Neither the hunter nor I reproduced this. It is confirmed by reading the two methods and the RabbitMQ client's dispatch model, not observed. The window is one work-pool thread against one release, so it needs a reply landing in roughly the same moment as a teardown. That is not exotic — the idle reaper fires on a timer with no regard for what a worker is doing — but it is narrow, and I am not claiming it has cost us a report. ## What I want **Goal:** a delivery that arrives for a target at any point during or after its `release` must not end up sitting unacked in a map nobody will ever read. Either it is requeued so a later session can get it, or it is refused — but it must not be silently retained. **Invariants:** 1. **Never drop a message to close this.** Requeue or leave it unacked deliberately; do not `basicAck` a reply nobody has read. #298 exists because dropping was the old behaviour. 2. **`release` must not become able to throw where it currently does not.** Its javadoc is explicit that a failed `basicNack` is logged and release proceeds, because release is called from teardown paths that must complete. A failed `basicCancel` still throws — leave that as it is. 3. **Do not hold `channelLock` across a blocking broker round trip from the delivery thread.** `basicCancel` and the nack loop already run under it; adding a delivery-side wait on the same lock risks serialising every delivery behind a teardown, or worse. If your fix takes that lock in `deliverCallback`, say why it cannot deadlock or stall the consumer pool. 4. **The duplicate-delivery path must keep working.** A redelivered message that is already in `held` is acked and dropped so the broker stops resending it. Do not break that. **Candidate mechanisms, all candidates only — pick one and justify it:** - Have `deliverCallback` consult a "this target is released" marker and nack the delivery immediately instead of holding it. The marker has to be set and read consistently with `held.remove`, which is the whole difficulty. - Replace the `remove`-then-iterate with a single atomic `held.compute(target, ...)` that both takes the map and leaves behind a tombstone `computeIfAbsent` will see. - Re-check `held` for the target once more at the end of `release` and nack anything that appeared. Simplest, and it narrows the window rather than closing it — if you choose this, say so plainly rather than describing it as a fix. **Decide it yourself and justify it.** I do not know which is right. If you find a fourth option, or find that one of these cannot work, that is a good outcome — say so. ## Rules - **A test that only calls the two methods in sequence proves nothing here.** The bug is an interleaving. Force it: a test double or a latch that lets a delivery land after `held.remove` has run but before `release` returns. If you cannot build that seam without changing production code beyond the fix, say so and describe exactly what your test does and does not prove. - If the AMQP contract suite can run in your worktree (it needs a broker), run it and quote the real output. If it cannot, say so plainly — do not describe an unrun suite as passing. Some claims in this area only show up against a real broker: fleetd #298's first attempted fix was correct-looking and wrong for exactly that reason. - Mutation proof required: revert the fix, quote the real failure output, restore it. - Never run `git stash` — the stash is shared across every worktree here. - Never run `git worktree remove` or `git worktree prune` — other workers are live in these worktrees. - Stage files explicitly; never `git add -A`. Never merge. - Run `cd fleetd && mvn clean install` **unpiped**, and quote the real `Tests run:` and `BUILD` lines. Never pipe maven through `tail`/`head`, and never read `$?` after a pipe. - Put your full report in the PR body as well as in your `fleet_reply`. ## Shape check When done, look in `msg/` only for the same shape: **two halves of one invariant guarded by different locks, or by a lock on one side and none on the other.** One line each, do **not** fix any of it.
Author
Owner

Merged as fa1f496. Pushed to main.

Build on main after the merge: Tests run: 1336, Failures: 0, Errors: 0, Skipped: 0 — BUILD SUCCESS, unpiped.

The mechanism is the right one

The worker chose the tombstone in the same map over a separate "released" set, and the reason it gives is the one that matters: ConcurrentHashMap already serializes compute and computeIfAbsent for the same key against each other, so the atomicity comes free and there is no second structure to keep in sync. A separate set would have reintroduced the two-things-that-must-agree shape that keeps producing defects here.

It also explicitly refused the third candidate — re-checking held once more at the end — on the grounds that it narrows the window rather than closing it, and said it would not call that a fix. That is the correct reading of the ticket.

It ran the contract suite against a real broker. Docker was available, so mvn test -Pcontract -Dtest=AmqpReplyInboxContractTest ran with Testcontainers RabbitMQ: Tests run: 8, Failures: 0 including #298's own test, unmodified and still green. That is the check that matters most here, because #298 is exactly the ticket whose fix this one sits next to, and broker redelivery timing is not something source reading settles.

What I checked myself

Every read site of held. RELEASED is a single shared LinkedHashMap instance, so any path that synchronizes on it and mutates it would corrupt state for every released target at once. I listed all six accesses on the branch:

196  held.clear()                  handleRecovery — clears, never iterates values
225  held.remove(target, RELEASED) own — two-arg, value-based, correct
307  held.compute(target, ...)     release — the swap
381  held.get(target)              peek   — guarded by `|| perTarget == RELEASED`
392  held.get(target)              ack    — guarded
425  held.computeIfAbsent(...)     deliverCallback — guarded

close() does not touch held at all. So no path mutates the sentinel, and the recovery clear() dropping tombstones is harmless: release already cancelled the consumer, and an explicitly cancelled consumer is not restored by automatic topology recovery, so no redelivery arrives for a released target.

My own mutations, one per half of the fix. The worker reverted the whole file. I reverted each half separately, to prove neither is carrying the other.

L — keep the tombstone, drop deliverCallback's RELEASED check:

AmqpReplyInboxReleaseRaceTest.deliveryArrivingWhileReleaseIsStillRunningIsNackedNotStranded:129
  both the pre-held m0 and the concurrently-arriving m1 must be nacked, got: [[tag=1 requeue=1]]
  ==> expected: <2> but was: <1>

M — keep the check, revert release to a plain held.remove(target): same failure, same line.

So each half fails on its own. Neither is decoration.

The own() addition is right — keeping it

The worker flagged held.remove(target, RELEASED) in own() as a deviation beyond the ask and offered to drop it. Keep it. Without it the fix leaks one tombstone per released target for the life of the process, and it is placed correctly — still under channelLock, after basicConsume returns, so no delivery under the fresh consumer can race it. The two-argument remove means it can only ever remove a tombstone, never a real map.

The shape check found something worth its own ticket

The worker reports the same asymmetric-locking shape in MessageService:

  • Task.turnId / Task.question and asyncTasksByTurn: ask() mutates them through markAsyncQuestion / clearAsyncQuestion / markAskTimedOut with no sessionLocks lock, while answer() mutates the same task's fields through clearAsyncQuestion / finishAsyncTask holding sessionLocks.
  • finishAsyncTask again: sendAsync's background virtual thread calls it unlocked; answer() calls it locked.

I touched markAskTimedOut myself in #307 two days ago and did not notice this. It is confirmed by reading only, not observed. I am filing it separately rather than folding it in here — it needs its own reachability analysis, and this ticket is done.

The lower-severity strandedReplies / queuedDeliveries note is correctly dismissed: each touch is a single atomic map operation on a documented (CB-640) self-healing best-effort flag.

Closing.

Merged as `fa1f496`. Pushed to `main`. Build on `main` after the merge: `Tests run: 1336, Failures: 0, Errors: 0, Skipped: 0` — `BUILD SUCCESS`, unpiped. ## The mechanism is the right one The worker chose the tombstone in the same map over a separate "released" set, and the reason it gives is the one that matters: `ConcurrentHashMap` already serializes `compute` and `computeIfAbsent` for the same key against each other, so the atomicity comes free and there is no second structure to keep in sync. A separate set would have reintroduced the two-things-that-must-agree shape that keeps producing defects here. It also explicitly refused the third candidate — re-checking `held` once more at the end — on the grounds that it narrows the window rather than closing it, and said it would not call that a fix. That is the correct reading of the ticket. **It ran the contract suite against a real broker.** Docker was available, so `mvn test -Pcontract -Dtest=AmqpReplyInboxContractTest` ran with Testcontainers RabbitMQ: `Tests run: 8, Failures: 0` including #298's own test, unmodified and still green. That is the check that matters most here, because #298 is exactly the ticket whose fix this one sits next to, and broker redelivery timing is not something source reading settles. ## What I checked myself **Every read site of `held`.** `RELEASED` is a single shared `LinkedHashMap` instance, so any path that synchronizes on it and mutates it would corrupt state for every released target at once. I listed all six accesses on the branch: ``` 196 held.clear() handleRecovery — clears, never iterates values 225 held.remove(target, RELEASED) own — two-arg, value-based, correct 307 held.compute(target, ...) release — the swap 381 held.get(target) peek — guarded by `|| perTarget == RELEASED` 392 held.get(target) ack — guarded 425 held.computeIfAbsent(...) deliverCallback — guarded ``` `close()` does not touch `held` at all. So no path mutates the sentinel, and the recovery `clear()` dropping tombstones is harmless: `release` already cancelled the consumer, and an explicitly cancelled consumer is not restored by automatic topology recovery, so no redelivery arrives for a released target. **My own mutations, one per half of the fix.** The worker reverted the whole file. I reverted each half separately, to prove neither is carrying the other. *L — keep the tombstone, drop `deliverCallback`'s `RELEASED` check:* ``` AmqpReplyInboxReleaseRaceTest.deliveryArrivingWhileReleaseIsStillRunningIsNackedNotStranded:129 both the pre-held m0 and the concurrently-arriving m1 must be nacked, got: [[tag=1 requeue=1]] ==> expected: <2> but was: <1> ``` *M — keep the check, revert `release` to a plain `held.remove(target)`:* same failure, same line. So each half fails on its own. Neither is decoration. ## The `own()` addition is right — keeping it The worker flagged `held.remove(target, RELEASED)` in `own()` as a deviation beyond the ask and offered to drop it. Keep it. Without it the fix leaks one tombstone per released target for the life of the process, and it is placed correctly — still under `channelLock`, after `basicConsume` returns, so no delivery under the fresh consumer can race it. The two-argument `remove` means it can only ever remove a tombstone, never a real map. ## The shape check found something worth its own ticket The worker reports the same asymmetric-locking shape in `MessageService`: - `Task.turnId` / `Task.question` and `asyncTasksByTurn`: `ask()` mutates them through `markAsyncQuestion` / `clearAsyncQuestion` / `markAskTimedOut` with **no** `sessionLocks` lock, while `answer()` mutates the same task's fields through `clearAsyncQuestion` / `finishAsyncTask` **holding** `sessionLocks`. - `finishAsyncTask` again: `sendAsync`'s background virtual thread calls it unlocked; `answer()` calls it locked. I touched `markAskTimedOut` myself in #307 two days ago and did not notice this. It is confirmed by reading only, not observed. I am filing it separately rather than folding it in here — it needs its own reachability analysis, and this ticket is done. The lower-severity `strandedReplies` / `queuedDeliveries` note is correctly dismissed: each touch is a single atomic map operation on a documented (CB-640) self-healing best-effort flag. Closing.
ltms closed this issue 2026-09-04 09:24:05 +02:00
Sign in to join this conversation.
1 Participants
Notifications
Due Date
No due date set.
Dependencies

No dependencies set.

Reference: fleet/fleetd#318