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