#298: AmqpReplyInbox.release() requeues held deliveries instead of dropping them
CI / contract (pull_request) Successful in 1m12s
CI / build (pull_request) Successful in 2m3s

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.
This commit is contained in:
Dai Ha
2026-09-04 12:14:00 +07:00
parent ba51e0c6cc
commit be5ba22c75
2 changed files with 78 additions and 8 deletions
@@ -208,18 +208,72 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
}
}
/**
* Release ownership of {@code target}: cancel its consumer, then nack-with-requeue every
* delivery still held for it instead of just dropping the local record.
*
* <p><strong>Cancelling a consumer does not requeue its in-flight deliveries.</strong> In AMQP,
* a delivery that was pushed to a consumer stays unacked, attached to the still-open
* {@link #channel}, until that channel or the connection closes — {@code basicCancel} alone does
* neither. So before this method existed with a requeue step, it dropped {@link #held}'s entries
* for {@code target} while the broker still considered them outstanding: never acked, never
* nacked, never requeued, and no longer reachable by {@link #peek} — permanently invisible. This
* is unlike {@link #handleRecovery} and {@link #close()}, whose bare {@code held.clear()} is
* correct because each has already made the broker requeue (a real connection drop, or
* {@code channel.close()} respectively) before clearing local state.
*
* <p><strong>Order: cancel first, then nack.</strong> A delivery tag stays valid for
* {@code basicNack} on this channel regardless of whether its consumer is still attached — only
* a channel/connection close invalidates it — so cancelling {@code target}'s consumer first does
* not risk the tags. Doing it the other way round does: nacking a delivery with {@code requeue}
* while its consumer is still active hands the message straight back to that <em>same</em>
* consumer the instant a prefetch slot frees up (confirmed against a real broker — see
* {@code AmqpReplyInboxContractTest.releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery}),
* which races this method's own {@code held.remove(target)}: the redelivery can land after the
* clear and leave a stale entry behind, so {@link #peek} is no longer reliably empty right after
* {@link #release}. Cancelling first closes that consumer, so the requeued message goes back to
* the queue for whichever consumer picks it up next (a later {@link #own}), not this one.
*
* <p><strong>Failure of the requeue is best-effort, not fatal.</strong> {@link #release} runs
* during teardown ({@code Fleetd} calls it right after {@code MessageService.abandon}), and a
* throw here would abort cleanups the caller depends on — the same argument fleetd #293 settled
* for {@code HerdrPeerLauncher.stop()}'s tab-close step. So a failed {@code basicNack} is logged
* at WARN, naming the target and delivery tag that leaked, and release proceeds; the delivery
* stays unacked on the broker rather than being silently dropped, so it is still recoverable by a
* later connection drop even though this release did not manage to requeue it immediately. A
* failed {@code basicCancel} still throws, unchanged from before this fix — that failure means
* the consumer may still be attached, so best-effort requeue is not attempted underneath it.
*/
@Override
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;
if (tag != null) {
try {
channel.basicCancel(tag);
} catch (IOException e) {
throw new IllegalStateException("cannot cancel consumer for " + target, e);
}
}
try {
channel.basicCancel(tag);
} catch (IOException e) {
throw new IllegalStateException("cannot cancel consumer for " + target, e);
var perTarget = held.remove(target);
if (perTarget != null) {
synchronized (perTarget) {
for (Held h : perTarget.values()) {
try {
channel.basicNack(h.deliveryTag(), false, true); // requeue, don't drop
} catch (IOException | RuntimeException e) {
// Caught broadly (not just IOException) for the same reason #293 catches
// RuntimeException in HerdrPeerLauncher.stop(): best-effort teardown must
// not be guarded only against the expected failure and bare against any
// other. The message stays unacked on the broker either way — not lost,
// just not proactively requeued — until a connection drop frees it.
log.warn("release({}): could not requeue held delivery (msgId={}, tag={})"
+ " back to the broker — it stays unacked until a connection"
+ " drop frees it: {}",
target, h.message().msgId(), h.deliveryTag(), e.getMessage());
}
}
}
}
}
}
@@ -154,7 +154,7 @@ class AmqpReplyInboxContractTest {
}
@Test
void releaseCancelsConsumerAndClearsHeld() throws Exception {
void releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery() throws Exception {
String target = "worker-release-" + System.nanoTime();
try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
inbox.own(target);
@@ -164,6 +164,22 @@ class AmqpReplyInboxContractTest {
inbox.release(target);
assertTrue(inbox.peek(target).isEmpty(),
"release clears the local held snapshot");
// fleetd #298: release() must not just drop the local record — the broker delivery was
// never acked, so cancelling the consumer alone leaves it unacked-but-orphaned on the
// still-open channel unless release() nacks it back with requeue=true. Prove the message
// is genuinely recoverable, not merely absent from peek: re-own the same target and
// confirm the broker redelivers it to the fresh consumer.
inbox.own(target);
List<ReplyInbox.InboxMessage> recovered = awaitPeek(inbox, target);
assertEquals(1, recovered.size(),
"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");
assertEquals("m1", recovered.getFirst().msgId());
assertEquals("release me", recovered.getFirst().content());
inbox.ack(target, "m1");
}
}