AmqpReplyInbox.release() drops an unacked reply that the broker still holds — a worker's report is lost silently #298

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

Found by a read-only hunt over the message layer. I verified every step of the path below myself, including the AMQP semantics the fix turns on.

This is a silent message-loss defect: a worker's real fleet_reply content becomes permanently invisible, with no error and no log line.

The bug

AmqpReplyInbox has one shared Channel — private final Channel channel (:93), created once at :160 and used by every target. release(target) cancels that target's consumer and drops its local record of held deliveries:

public void release(String target) {
    synchronized (channelLock) {
        String tag = consumerTags.remove(target);
        held.remove(target);                   // stale delivery tags must not survive release
        if (tag == null) return;
        try {
            channel.basicCancel(tag);
        } catch (IOException e) { ... }
    }
}

That comment is where the reasoning went wrong. The delivery tags are not stale. Cancelling a consumer does not requeue the messages already delivered to it — in AMQP those stay unacked, attached to the still-open channel, until the channel or connection closes. release closes neither. So the entry is dropped locally while the broker still considers it outstanding: never acked, never nacked, never requeued, and no longer reachable by fleet_poll or peek.

Why the two siblings that look identical are actually safe

Both do a bare held.clear(), and both are correct, because each has a precondition release lacks:

  • handleRecovery (:176-190) runs after a real connection drop. The broker has already requeued every unacked delivery on that connection, so clearing local state is exactly right.
  • close() (:433-451) calls channel.close() at :437, which requeues them.

release() copied the pattern without the precondition that makes it safe. Same shape as #274 / #283 / #293, one layer out: the right action on one branch, the same action on a sibling where it is wrong.

InMemoryReplyInbox's release is not a counterexample — for that soft-state adapter, dropping on release is the documented, correct behaviour. Do not "fix" it.

The path into the bad state — verified end to end

  1. MessageService.reply() (:414-455) takes a fleet_reply with no live waiter and falls through to inbox.publish(...). The message lands in AmqpReplyInbox.held, unacked.
  2. Before anyone calls fleet_poll to drain it, the lead issues a fresh send()/sendAsync() to the same target. send() (:797) clears the bookkeeping flag:
    // CB-640: this send now owns target's delivery, so any earlier stranded-reply or
    // still-queued fact no longer describes the live state — clear both rather than let
    // them outlive the send that supersedes them.
    strandedReplies.remove(target);
    
    It clears the flag. It never drains the actual inbox entry. That is the hinge.
  3. The new turn resolves normally through its live waiter. Nothing looks abnormal anywhere.
  4. The session is torn down — fleet_stop, or the idle reaper. Fleetd.java:597 calls messages.abandon(target, reason, true). abandon gates recovery on hasStrandedReply(target) (:628), which is now false, so recoverStrandedReply/drainReplies never runs for the old message.
  5. Fleetd.java:598 then calls replyInbox.release(target) — which drops the still-unacked entry.

What is left behind: one worker's real report, invisible to fleet_poll and peek, never redelivered. No bounded reaper recovers it. It is freed only when the whole AMQP connection actually drops — a broker restart, a network blip, or a daemon bounce. Calling own(target) again for a later worker reusing the same target name does not recover it: the stuck delivery stays tied to the old, cancelled consumer on the still-open channel.

Direction of harm: silent loss of the one artefact the whole bridge exists to carry. A lead sees a member that finished and reported nothing, and has no way to tell that from a member that genuinely never replied.

Test coverage

AmqpReplyInboxContractTest.releaseCancelsConsumerAndClearsHeld (:156-167) asserts only that peek(target) is empty after release. It never checks whether the broker still holds the message or has requeued it — so it passes on the buggy behaviour and would pass on the fix too. It is asserting the symptom, not the property.

unackedReplySurvivesRestartAndIsRedelivered in the same file does prove redelivery, but it closes the whole AmqpReplyInbox — i.e. the connection — so it exercises the close() path, not release(). That is why this gap survived.

Scope

Make release(target) return the target's held deliveries to the broker instead of dropping them, so a later own(target) — or any other consumer — can still get them. Requeue is the right disposition: the message was never processed.

Decide the mechanism yourself and justify it in your report. Nack-with-requeue over the held delivery tags before clearing is the obvious candidate; say why you chose what you chose, and what happens if that call itself fails.

Rules

  • Order matters. Whatever you do must happen while the tags are still valid, i.e. before the cancel and before held.remove. Say in your report why your ordering is correct.
  • Failure of the requeue must not leave release half-done. release is called during teardown (Fleetd.java:598), right after abandon. A throw there aborts a teardown that other cleanups depend on. Look at how #293 handled exactly this on the launcher side (HerdrPeerLauncher.stop()) and follow the house style: best-effort, log a WARN naming what leaked, do not mask the caller's real failure. Catch RuntimeException, not a narrower type — the reason is recorded in that method.
  • Do not change InMemoryReplyInbox. Dropping on release is correct there.
  • Do not change handleRecovery or close(). Both are correct.
  • Do not change MessageService. In particular do not touch the CB-640 flag-clearing at :797: making send() drain the old inbox entry is a different design decision with its own consequences, and it is mine to make, not this ticket's.

Prove it

Add a test that fails without the fix. The existing releaseCancelsConsumerAndClearsHeld asserts the wrong property — extend or replace it so it asserts the message is still recoverable after release, not merely absent from peek. Then keep a peek-is-empty assertion too, since that behaviour must not regress.

Mutation proof required: revert the fix, quote the real failure output, restore it.

Shape check

When done, look for the same shape — local state cleared as if the broker had requeued, when nothing made it requeue — in AmqpReplyInbox only. One line each, do not fix any of it.

Also checked, and clean

The same hunt looked hard at ReplyPushLoop, LeadCoordLoop, LeadHeartbeatLoop, LeadMailbox and CompletionResolver and found them all symmetric (deliver-then-ack, or fail-without-ack). It also raised a rendezvous double-open scenario in MessageService.answer; I checked that and it is a false positive — send() holds the session lock across its entire wait and its inner finally closes the waiter before the outer finally unlocks, so no window exists where a waiter is registered and the lock is free. answer() takes that same lock first and returns BUSY. Rendezvous.open's javadoc already states this. Recorded here so it is not re-investigated.

Found by a read-only hunt over the message layer. **I verified every step of the path below myself**, including the AMQP semantics the fix turns on. This is a silent message-loss defect: a worker's real `fleet_reply` content becomes permanently invisible, with no error and no log line. ## The bug `AmqpReplyInbox` has **one** shared `Channel` — `private final Channel channel` (`:93`), created once at `:160` and used by every target. `release(target)` cancels that target's consumer and drops its local record of held deliveries: ```java public void release(String target) { synchronized (channelLock) { String tag = consumerTags.remove(target); held.remove(target); // stale delivery tags must not survive release if (tag == null) return; try { channel.basicCancel(tag); } catch (IOException e) { ... } } } ``` **That comment is where the reasoning went wrong.** The delivery tags are not stale. Cancelling a consumer does not requeue the messages already delivered to it — in AMQP those stay unacked, attached to the still-open channel, until the *channel or connection* closes. `release` closes neither. So the entry is dropped locally while the broker still considers it outstanding: never acked, never nacked, never requeued, and no longer reachable by `fleet_poll` or `peek`. ## Why the two siblings that look identical are actually safe Both do a bare `held.clear()`, and both are correct, because each has a precondition `release` lacks: - `handleRecovery` (`:176-190`) runs *after a real connection drop*. The broker has already requeued every unacked delivery on that connection, so clearing local state is exactly right. - `close()` (`:433-451`) calls `channel.close()` at `:437`, which requeues them. `release()` copied the pattern without the precondition that makes it safe. Same shape as #274 / #283 / #293, one layer out: the right action on one branch, the same action on a sibling where it is wrong. `InMemoryReplyInbox`'s release is **not** a counterexample — for that soft-state adapter, dropping on release is the documented, correct behaviour. Do not "fix" it. ## The path into the bad state — verified end to end 1. `MessageService.reply()` (`:414-455`) takes a `fleet_reply` with no live waiter and falls through to `inbox.publish(...)`. The message lands in `AmqpReplyInbox.held`, unacked. 2. Before anyone calls `fleet_poll` to drain it, the lead issues a fresh `send()`/`sendAsync()` to the same target. `send()` (`:797`) clears the bookkeeping flag: ```java // CB-640: this send now owns target's delivery, so any earlier stranded-reply or // still-queued fact no longer describes the live state — clear both rather than let // them outlive the send that supersedes them. strandedReplies.remove(target); ``` It clears the **flag**. It never drains the actual inbox entry. That is the hinge. 3. The new turn resolves normally through its live waiter. Nothing looks abnormal anywhere. 4. The session is torn down — `fleet_stop`, or the idle reaper. `Fleetd.java:597` calls `messages.abandon(target, reason, true)`. `abandon` gates recovery on `hasStrandedReply(target)` (`:628`), which is now **false**, so `recoverStrandedReply`/`drainReplies` never runs for the old message. 5. `Fleetd.java:598` then calls `replyInbox.release(target)` — which drops the still-unacked entry. **What is left behind:** one worker's real report, invisible to `fleet_poll` and `peek`, never redelivered. No bounded reaper recovers it. It is freed only when the whole AMQP connection actually drops — a broker restart, a network blip, or a daemon bounce. Calling `own(target)` again for a later worker reusing the same target name does **not** recover it: the stuck delivery stays tied to the old, cancelled consumer on the still-open channel. **Direction of harm:** silent loss of the one artefact the whole bridge exists to carry. A lead sees a member that finished and reported nothing, and has no way to tell that from a member that genuinely never replied. ## Test coverage `AmqpReplyInboxContractTest.releaseCancelsConsumerAndClearsHeld` (`:156-167`) asserts only that `peek(target)` is empty after release. It never checks whether the broker still holds the message or has requeued it — so it passes on the buggy behaviour and would pass on the fix too. It is asserting the symptom, not the property. `unackedReplySurvivesRestartAndIsRedelivered` in the same file *does* prove redelivery, but it closes the whole `AmqpReplyInbox` — i.e. the connection — so it exercises the `close()` path, not `release()`. That is why this gap survived. ## Scope Make `release(target)` return the target's held deliveries to the broker instead of dropping them, so a later `own(target)` — or any other consumer — can still get them. Requeue is the right disposition: the message was never processed. Decide the mechanism yourself and justify it in your report. Nack-with-requeue over the held delivery tags before clearing is the obvious candidate; say why you chose what you chose, and what happens if that call itself fails. ## Rules - **Order matters.** Whatever you do must happen while the tags are still valid, i.e. before the cancel and before `held.remove`. Say in your report why your ordering is correct. - **Failure of the requeue must not leave `release` half-done.** `release` is called during teardown (`Fleetd.java:598`), right after `abandon`. A throw there aborts a teardown that other cleanups depend on. Look at how #293 handled exactly this on the launcher side (`HerdrPeerLauncher.stop()`) and follow the house style: best-effort, log a WARN naming what leaked, do not mask the caller's real failure. Catch `RuntimeException`, not a narrower type — the reason is recorded in that method. - **Do not change `InMemoryReplyInbox`.** Dropping on release is correct there. - **Do not change `handleRecovery` or `close()`.** Both are correct. - **Do not change `MessageService`.** In particular do not touch the CB-640 flag-clearing at `:797`: making `send()` drain the old inbox entry is a different design decision with its own consequences, and it is mine to make, not this ticket's. ## Prove it Add a test that fails without the fix. The existing `releaseCancelsConsumerAndClearsHeld` asserts the wrong property — **extend or replace it so it asserts the message is still recoverable after `release`**, not merely absent from `peek`. Then keep a `peek`-is-empty assertion too, since that behaviour must not regress. Mutation proof required: revert the fix, quote the real failure output, restore it. ## Shape check When done, look for the same shape — **local state cleared as if the broker had requeued, when nothing made it requeue** — in `AmqpReplyInbox` only. One line each, **do not fix any of it**. ## Also checked, and clean The same hunt looked hard at `ReplyPushLoop`, `LeadCoordLoop`, `LeadHeartbeatLoop`, `LeadMailbox` and `CompletionResolver` and found them all symmetric (deliver-then-ack, or fail-without-ack). It also raised a `rendezvous` double-open scenario in `MessageService.answer`; **I checked that and it is a false positive** — `send()` holds the session lock across its entire wait and its inner `finally` closes the waiter *before* the outer `finally` unlocks, so no window exists where a waiter is registered and the lock is free. `answer()` takes that same lock first and returns `BUSY`. `Rendezvous.open`'s javadoc already states this. Recorded here so it is not re-investigated.
Author
Owner

Merged as 38c248e. Real merge built green at 1307 tests, and the @Tag("contract") suite green at 8 tests against a live RabbitMQ.

No correction needed from me. This is the best-executed unit of the batch.

The worker deviated from the ticket, and it was right to

My ticket suggested nacking before the cancel, on the reasoning that the tags must still be valid. The worker tried that, watched it fail against a real broker, and reversed the order — then said so plainly and asked me to double-check it.

The reasoning is correct: nacking with requeue while the consumer is still attached hands the message straight back to that same consumer as soon as a prefetch slot frees, which races release()'s own held.remove. The redelivery can land after the clear and leave a stale entry, so peek is no longer reliably empty right after release. Cancelling first closes that consumer, and a delivery tag stays valid for basicNack on an open channel whether or not its consumer is attached — only a channel or connection close invalidates it. So cancel-then-nack costs nothing and removes the race.

My ticket's suggested order was wrong, and only a live broker could have shown that. Reasoning about AMQP from the source would not have caught it. Flagging the deviation for review instead of quietly following the spec — or quietly ignoring it — is exactly right.

To be precise about what I checked: I verified the ordering argument and the tag-validity claim, and the tests pass both ways round on the fixed code. I did not independently reproduce the race, because it is timing-dependent. The evidence for it is the worker's quoted live failure, not my own run.

My mutation — the subtle half-fix, not the missing fix

The worker proved the fix by reverting release() wholesale. That proves something happens. It does not distinguish the fix from a plausible near-miss, so I mutated for that instead: basicNack(tag, false, true) → basicNack(tag, false, false), i.e. discard instead of requeue.

That is the dangerous mutant. The message still leaves the unacked state, so any assertion of the form "it is no longer stuck" passes. Only an assertion that the reply is genuinely recoverable can tell the two apart:

AmqpReplyInboxContractTest.releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery:175
  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>

Caught, against the live broker. The rewritten test asserts the property rather than the symptom — which was the whole complaint against the test it replaced.

Two details I checked rather than took on trust

  • The locking follows the documented convention. held is declared at :96 as "target → (msgId → held delivery). Per-target map is guarded by synchronizing on itself." The new synchronized (perTarget) matches it, and the whole method still runs under channelLock.
  • The tag == null path did not change behaviour. The old code removed from held before its early return, so a target with no consumer still got its entry dropped. The new code keeps that and adds the requeue.

The failure semantics are right, and they are not symmetric on purpose

A failed basicNack is caught broadly, logged at WARN naming target, msgId and tag, and does not abort release() — correct, since release runs mid-teardown right after MessageService.abandon, and #293 settled that argument. A failed basicCancel still throws, unchanged. That asymmetry is deliberate and the javadoc explains it: if the cancel failed, the consumer may still be attached, so attempting a best-effort requeue underneath it would reintroduce the very race the ordering exists to avoid.

Shape check came back clean, and its reasoning on ack() is worth keeping: it restores the entry when basicAck fails rather than dropping it — the opposite direction, so not this shape.

Merged as `38c248e`. Real merge built green at **1307 tests**, and the `@Tag("contract")` suite green at **8 tests** against a live RabbitMQ. No correction needed from me. This is the best-executed unit of the batch. ## The worker deviated from the ticket, and it was right to My ticket suggested nacking **before** the cancel, on the reasoning that the tags must still be valid. The worker tried that, watched it fail against a real broker, and reversed the order — then said so plainly and asked me to double-check it. The reasoning is correct: nacking with requeue while the consumer is still attached hands the message straight back to that **same** consumer as soon as a prefetch slot frees, which races `release()`'s own `held.remove`. The redelivery can land after the clear and leave a stale entry, so `peek` is no longer reliably empty right after `release`. Cancelling first closes that consumer, and a delivery tag stays valid for `basicNack` on an open channel whether or not its consumer is attached — only a channel or connection close invalidates it. So cancel-then-nack costs nothing and removes the race. **My ticket's suggested order was wrong, and only a live broker could have shown that.** Reasoning about AMQP from the source would not have caught it. Flagging the deviation for review instead of quietly following the spec — or quietly ignoring it — is exactly right. To be precise about what I checked: I verified the ordering *argument* and the tag-validity claim, and the tests pass both ways round on the fixed code. I did **not** independently reproduce the race, because it is timing-dependent. The evidence for it is the worker's quoted live failure, not my own run. ## My mutation — the subtle half-fix, not the missing fix The worker proved the fix by reverting `release()` wholesale. That proves something happens. It does not distinguish the fix from a plausible near-miss, so I mutated for that instead: `basicNack(tag, false, true)` → `basicNack(tag, false, false)`, i.e. **discard instead of requeue**. That is the dangerous mutant. The message still leaves the unacked state, so any assertion of the form "it is no longer stuck" passes. Only an assertion that the reply is genuinely *recoverable* can tell the two apart: ``` AmqpReplyInboxContractTest.releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery:175 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> ``` Caught, against the live broker. The rewritten test asserts the property rather than the symptom — which was the whole complaint against the test it replaced. ## Two details I checked rather than took on trust - **The locking follows the documented convention.** `held` is declared at `:96` as *"target → (msgId → held delivery). Per-target map is guarded by synchronizing on itself."* The new `synchronized (perTarget)` matches it, and the whole method still runs under `channelLock`. - **The `tag == null` path did not change behaviour.** The old code removed from `held` before its early return, so a target with no consumer still got its entry dropped. The new code keeps that and adds the requeue. ## The failure semantics are right, and they are not symmetric on purpose A failed `basicNack` is caught broadly, logged at WARN naming target, msgId and tag, and does not abort `release()` — correct, since `release` runs mid-teardown right after `MessageService.abandon`, and #293 settled that argument. A failed `basicCancel` still throws, unchanged. That asymmetry is deliberate and the javadoc explains it: if the cancel failed, the consumer may still be attached, so attempting a best-effort requeue underneath it would reintroduce the very race the ordering exists to avoid. Shape check came back clean, and its reasoning on `ack()` is worth keeping: it restores the entry when `basicAck` fails rather than dropping it — the opposite direction, so not this shape.
ltms closed this issue 2026-09-04 07:23:55 +02:00
Sign in to join this conversation.
1 Participants
Notifications
Due Date
No due date set.
Dependencies

No dependencies set.

Reference: fleet/fleetd#298