#298: AmqpReplyInbox.release() requeues held deliveries instead of dropping them #301

Closed
agent wants to merge 0 commits from worker/fix-298-ce1acb-8 into main
Member

Fixes fleetd #298: AmqpReplyInbox.release() drops an unacked reply the broker still holds.

The bug

release(target) cancelled the target's AMQP consumer and dropped held's local record of its
deliveries, but cancelling a consumer does not requeue the broker's in-flight deliveries — they
stay unacked, attached to the still-open channel, until the channel or connection actually closes.
So a held-but-undrained fleet_reply became permanently unreachable: never acked, never nacked,
never requeued, and invisible to peek/fleet_poll.

handleRecovery and close() were already correct (each makes the broker requeue — a real
connection drop, or channel.close() — before clearing local state). release() copied their bare
held.clear()/held.remove() pattern without that precondition.

The fix

release(target) now:

  1. Cancels the consumer first (channel.basicCancel(tag)).
  2. Then nack-with-requeue (channel.basicNack(tag, false, true)) every delivery still held for
    target, before dropping the local held entry.

Why nack-with-requeue, and why after the cancel (not before): requeue is the right disposition —
the message was never processed. I initially tried nack-before-cancel (the ticket's suggested order,
"before the cancel and before held.remove"), but a real broker run (AmqpReplyInboxContractTest
against RabbitMQ via Testcontainers) showed a race: while the consumer is still active, a
nack-with-requeue is delivered straight back to that same consumer the instant a prefetch slot frees
up, racing this method's own held.remove(target) — the redelivery can land after the clear and
leave a stale entry, so peek was no longer reliably empty right after release (I saw this fail
live: expected: <true> but was: <false> on the "release clears the local held snapshot"
assertion). Cancelling first avoids this: a delivery tag stays valid for basicNack on an open
channel regardless of whether its consumer is attached, so cancelling first costs nothing and closes
the race.

Failure handling: a failed basicNack is caught broadly (IOException | RuntimeException,
matching the wide-catch house style HerdrPeerLauncher.stop() established for #293) and logged at
WARN naming the target/msgId/delivery-tag that leaked, without aborting release() — it runs during
teardown (Fleetd calls it right after MessageService.abandon), and a throw here would abort
cleanups the caller depends on. The delivery stays unacked on the broker either way (not silently
dropped) — recoverable by a later connection drop even when the immediate requeue attempt fails. A
failed basicCancel still throws, unchanged from before this fix.

Out of scope, untouched per the ticket: InMemoryReplyInbox (drop-on-release is correct there),
handleRecovery, close(), and MessageService (including the CB-640 flag-clearing at :797).

Proof

Docker was available in this worktree (docker info succeeded), so I ran the real @Tag("contract")
suite against a RabbitMQ container via Testcontainers (mvn test -Pcontract) rather than relying on
reasoning alone.

Extended AmqpReplyInboxContractTest.releaseCancelsConsumerAndClearsHeld →
releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery: keeps the original "peek is empty after
release" assertion, and adds a new one — re-own() the same target after release() and assert the
delivery is redelivered (msgId + content), proving it is genuinely recoverable, not merely absent
from peek.

Mutation proof (revert fix → real failure → restore):

  • Reverted AmqpReplyInbox.java to the pre-fix release() (bare held.remove, no requeue) via
    cp from a saved copy of the pre-edit file (never used git stash — it's shared across worktrees
    on this machine).
  • Ran mvn test -Pcontract -Dtest=AmqpReplyInboxContractTest#releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery:
    org.opentest4j.AssertionFailedError: a reply held (but undrained) at release() time must still be
    recoverable — release() must requeue it, not silently drop it while the broker still considers it
    outstanding ==> expected: <1> but was: <0>
    Tests run: 1, Failures: 1, Errors: 0, Skipped: 0
    BUILD FAILURE
    
  • Restored the fix (cp back), reran the same test: Tests run: 1, Failures: 0, Errors: 0, Skipped: 0,
    BUILD SUCCESS.
  • Ran the full AmqpReplyInboxContractTest class with the fix in place: Tests run: 8, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS — all other contract behaviour (publish/peek/ack, dedup,
    cross-restart durability, prefetch, unroutable publish, no-consumer-on-publish) unaffected.

Full build

mvn clean install (hermetic, unpiped, contract tests excluded by default — no Docker required):

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

Ran twice: once before touching the contract test/mutation work, once again after restoring the fix.
Both green with the identical Tests run: 1303 count.

Shape check (per ticket — report only, not fixed)

Scanned AmqpReplyInbox.java for the same shape ("local state cleared as if the broker had
requeued, when nothing made it requeue") outside release():

  • handleRecovery (held.clear()) — has the precondition (real connection drop already requeued
    everything) — not the same defect, already excluded by the ticket.
  • close() (channel.close(), no explicit held.clear() at all) — channel.close() itself
    requeues, and the object is discarded afterward anyway — not the same defect, already excluded.
  • ack() — removes from held only after a successful basicAck; on IOException it restores
    the entry rather than dropping it — opposite direction of the bug, not the same shape.
  • deliverCallback (duplicate-delivery path) — acks the duplicate's tag before dropping it, doesn't
    assume anything was requeued — not the same shape.

No second occurrence found; the one in release() was the only instance of this shape in the file.

Report

  • PR: (this PR)
  • branch: worker/fix-298-ce1acb-8
  • root: /Users/dai.ha/LTMS/.bridged-worktrees/56872d-8 (own worktree, confirmed via git rev-parse --show-toplevel before editing)
  • files: fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java,
    fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxContractTest.java
  • build: mvn clean install → Tests run: 1303, Failures: 0, Errors: 0, Skipped: 0 / BUILD SUCCESS.
    mvn test -Pcontract -Dtest=AmqpReplyInboxContractTest (Docker) → Tests run: 8, Failures: 0, Errors: 0, Skipped: 0 / BUILD SUCCESS.
  • caveat for reviewer: I deviated from the ticket's literal suggested ordering ("nack before
    cancel") because it failed live against a real broker (self-redelivery race — see above); I cancel
    first, then nack. Please double-check that reasoning. AmqpReplyInboxContractTest needs Docker to
    run (@Tag("contract"), excluded from the default mvn clean install) — it was available in this
    worktree so I ran it directly rather than only reasoning about the fix.
Fixes fleetd #298: `AmqpReplyInbox.release()` drops an unacked reply the broker still holds. ## The bug `release(target)` cancelled the target's AMQP consumer and dropped `held`'s local record of its deliveries, but cancelling a consumer does **not** requeue the broker's in-flight deliveries — they stay unacked, attached to the still-open `channel`, until the channel or connection actually closes. So a held-but-undrained `fleet_reply` became permanently unreachable: never acked, never nacked, never requeued, and invisible to `peek`/`fleet_poll`. `handleRecovery` and `close()` were already correct (each makes the broker requeue — a real connection drop, or `channel.close()` — before clearing local state). `release()` copied their bare `held.clear()`/`held.remove()` pattern without that precondition. ## The fix `release(target)` now: 1. Cancels the consumer first (`channel.basicCancel(tag)`). 2. Then nack-with-requeue (`channel.basicNack(tag, false, true)`) every delivery still held for `target`, before dropping the local `held` entry. **Why nack-with-requeue, and why after the cancel (not before):** requeue is the right disposition — the message was never processed. I initially tried nack-before-cancel (the ticket's suggested order, "before the cancel and before `held.remove`"), but a real broker run (`AmqpReplyInboxContractTest` against RabbitMQ via Testcontainers) showed a race: while the consumer is still active, a nack-with-requeue is delivered straight back to that same consumer the instant a prefetch slot frees up, racing this method's own `held.remove(target)` — the redelivery can land after the clear and leave a stale entry, so `peek` was no longer reliably empty right after `release` (I saw this fail live: `expected: <true> but was: <false>` on the "release clears the local held snapshot" assertion). Cancelling first avoids this: a delivery tag stays valid for `basicNack` on an open channel regardless of whether its consumer is attached, so cancelling first costs nothing and closes the race. **Failure handling:** a failed `basicNack` is caught broadly (`IOException | RuntimeException`, matching the wide-catch house style `HerdrPeerLauncher.stop()` established for #293) and logged at WARN naming the target/msgId/delivery-tag that leaked, without aborting `release()` — it runs during teardown (`Fleetd` calls it right after `MessageService.abandon`), and a throw here would abort cleanups the caller depends on. The delivery stays unacked on the broker either way (not silently dropped) — recoverable by a later connection drop even when the immediate requeue attempt fails. A failed `basicCancel` still throws, unchanged from before this fix. Out of scope, untouched per the ticket: `InMemoryReplyInbox` (drop-on-release is correct there), `handleRecovery`, `close()`, and `MessageService` (including the CB-640 flag-clearing at `:797`). ## Proof Docker was available in this worktree (`docker info` succeeded), so I ran the real `@Tag("contract")` suite against a RabbitMQ container via Testcontainers (`mvn test -Pcontract`) rather than relying on reasoning alone. Extended `AmqpReplyInboxContractTest.releaseCancelsConsumerAndClearsHeld` → `releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery`: keeps the original "peek is empty after release" assertion, and adds a new one — re-`own()` the same target after `release()` and assert the delivery is redelivered (msgId + content), proving it is genuinely recoverable, not merely absent from `peek`. **Mutation proof** (revert fix → real failure → restore): - Reverted `AmqpReplyInbox.java` to the pre-fix `release()` (bare `held.remove`, no requeue) via `cp` from a saved copy of the pre-edit file (never used `git stash` — it's shared across worktrees on this machine). - Ran `mvn test -Pcontract -Dtest=AmqpReplyInboxContractTest#releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery`: ``` org.opentest4j.AssertionFailedError: a reply held (but undrained) at release() time must still be recoverable — release() must requeue it, not silently drop it while the broker still considers it outstanding ==> expected: <1> but was: <0> Tests run: 1, Failures: 1, Errors: 0, Skipped: 0 BUILD FAILURE ``` - Restored the fix (`cp` back), reran the same test: `Tests run: 1, Failures: 0, Errors: 0, Skipped: 0`, `BUILD SUCCESS`. - Ran the full `AmqpReplyInboxContractTest` class with the fix in place: `Tests run: 8, Failures: 0, Errors: 0, Skipped: 0`, `BUILD SUCCESS` — all other contract behaviour (publish/peek/ack, dedup, cross-restart durability, prefetch, unroutable publish, no-consumer-on-publish) unaffected. ## Full build `mvn clean install` (hermetic, unpiped, contract tests excluded by default — no Docker required): ``` Tests run: 1303, Failures: 0, Errors: 0, Skipped: 0 BUILD SUCCESS ``` Ran twice: once before touching the contract test/mutation work, once again after restoring the fix. Both green with the identical `Tests run: 1303` count. ## Shape check (per ticket — report only, not fixed) Scanned `AmqpReplyInbox.java` for the same shape ("local state cleared as if the broker had requeued, when nothing made it requeue") outside `release()`: - `handleRecovery` (`held.clear()`) — has the precondition (real connection drop already requeued everything) — not the same defect, already excluded by the ticket. - `close()` (`channel.close()`, no explicit `held.clear()` at all) — `channel.close()` itself requeues, and the object is discarded afterward anyway — not the same defect, already excluded. - `ack()` — removes from `held` only after a successful `basicAck`; on `IOException` it *restores* the entry rather than dropping it — opposite direction of the bug, not the same shape. - `deliverCallback` (duplicate-delivery path) — acks the duplicate's tag before dropping it, doesn't assume anything was requeued — not the same shape. No second occurrence found; the one in `release()` was the only instance of this shape in the file. ## Report - PR: (this PR) - branch: `worker/fix-298-ce1acb-8` - root: `/Users/dai.ha/LTMS/.bridged-worktrees/56872d-8` (own worktree, confirmed via `git rev-parse --show-toplevel` before editing) - files: `fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java`, `fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxContractTest.java` - build: `mvn clean install` → `Tests run: 1303, Failures: 0, Errors: 0, Skipped: 0` / `BUILD SUCCESS`. `mvn test -Pcontract -Dtest=AmqpReplyInboxContractTest` (Docker) → `Tests run: 8, Failures: 0, Errors: 0, Skipped: 0` / `BUILD SUCCESS`. - caveat for reviewer: I deviated from the ticket's literal suggested ordering ("nack before cancel") because it failed live against a real broker (self-redelivery race — see above); I cancel first, then nack. Please double-check that reasoning. `AmqpReplyInboxContractTest` needs Docker to run (`@Tag("contract")`, excluded from the default `mvn clean install`) — it was available in this worktree so I ran it directly rather than only reasoning about the fix.
agent added 1 commit 2026-09-04 07:14:46 +02:00
#298: AmqpReplyInbox.release() requeues held deliveries instead of dropping them
CI / contract (pull_request) Successful in 1m12s
CI / build (pull_request) Successful in 2m3s
be5ba22c75
release(target) used to cancel the target's consumer and clear held's local
record for it. Cancelling a consumer does not requeue the broker's in-flight
deliveries — they stay unacked on the still-open channel until a real
connection drop. So a held-but-undrained reply became permanently
unreachable: never acked, never nacked, never requeued, invisible to peek.

Fix: cancel the consumer first (so it can no longer receive redeliveries),
then nack-with-requeue every held delivery for that target before dropping
the local record. Nacking before the cancel was tried first but a real
broker demonstrated a race: the still-active consumer immediately received
the requeued message back, racing held.remove and leaving peek non-empty.
Cancelling first avoids that. A failed requeue is logged at WARN and does
not abort release(), matching the best-effort teardown style #293 settled
for HerdrPeerLauncher.stop().

Extends AmqpReplyInboxContractTest.releaseCancelsConsumer... to assert the
held delivery is recoverable via a later own(), not just absent from peek.
ltms closed this pull request 2026-09-04 07:23:58 +02:00
Some checks are pending
CI / contract (pull_request) Successful in 1m12s
CI / build (pull_request) Successful in 2m3s

Pull request closed

Sign in to join this conversation.