From be5ba22c75ca3eaef1f793713a1f1424b7245204 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Fri, 4 Sep 2026 12:14:00 +0700 Subject: [PATCH] #298: AmqpReplyInbox.release() requeues held deliveries instead of dropping them MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- .../dev/ltms/fleet/msg/AmqpReplyInbox.java | 68 +++++++++++++++++-- .../fleet/msg/AmqpReplyInboxContractTest.java | 18 ++++- 2 files changed, 78 insertions(+), 8 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java b/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java index dfef5c2..cdd1bd2 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java @@ -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. + * + *

Cancelling a consumer does not requeue its in-flight deliveries. 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. + * + *

Order: cancel first, then nack. 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 same + * 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. + * + *

Failure of the requeue is best-effort, not fatal. {@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()); + } + } + } } } } diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxContractTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxContractTest.java index 865d7d5..b310062 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxContractTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxContractTest.java @@ -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 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"); } }