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 cdd1bd2..ebdb068 100644
--- a/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java
+++ b/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java
@@ -22,6 +22,7 @@ import java.util.concurrent.ConcurrentSkipListMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
+import java.util.concurrent.atomic.AtomicReference;
/**
* AMQP-backed {@link ReplyInbox} (CB-307 Stage 2): genuine cross-restart durability behind the same
@@ -93,8 +94,23 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
private final Channel channel;
/** All channel operations (publish/declare/ack/cancel) serialize on this — a Channel is not thread-safe. */
private final Object channelLock = new Object();
- /** target → (msgId → held delivery). Per-target map is guarded by synchronizing on itself. */
+ /**
+ * target → (msgId → held delivery). Per-target map is guarded by synchronizing on itself.
+ *
+ *
CB-318 tombstone. The value {@link #RELEASED} is a reserved sentinel: it
+ * marks a target whose {@link #release} has already run, so {@link #deliverCallback} can tell a
+ * delivery landing after release() apart from a fresh target it has never seen. See both methods'
+ * javadoc for why a plain {@code held.remove(target)} is not enough.
+ */
private final ConcurrentHashMap> held = new ConcurrentHashMap<>();
+
+ /**
+ * CB-318 sentinel stored in {@link #held} for a target whose {@link #release} has already run.
+ * Never mutated — every read site compares it by reference ({@code ==}) before touching it as a
+ * map, because it is a single object shared across every released target and calling a mutator on
+ * it would corrupt state for all of them.
+ */
+ private static final LinkedHashMap RELEASED = new LinkedHashMap<>();
/** Targets whose queue is declared and consumer is running, mapped to their broker consumer tag. */
private final ConcurrentHashMap consumerTags = new ConcurrentHashMap<>();
@@ -201,6 +217,12 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
channel.queueDeclare(queue, true, false, false, null); // durable, non-exclusive, keep on idle
String tag = channel.basicConsume(queue, false, deliverCallback(target), _ -> { });
consumerTags.put(target, tag);
+ // CB-318: drop a stale RELEASED tombstone from a prior ownership of this same target
+ // string, so a delivery under this fresh consumer is held normally instead of being
+ // nacked forever by deliverCallback's RELEASED check. Safe to do here, still under
+ // channelLock: no delivery for the consumer tag just registered above can reach
+ // deliverCallback before this basicConsume call returns.
+ held.remove(target, RELEASED);
log.debug("AMQP inbox owns queue {} for target {}", queue, target);
} catch (IOException e) {
throw new IllegalStateException("cannot own queue " + queue, e);
@@ -243,6 +265,32 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
* 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.
+ *
+ * CB-318: {@code held.remove(target)} alone leaves a second window open. The
+ * bullet above already explains why cancelling first does not save a tag from going stale — but
+ * that only accounts for a delivery landing before this method starts touching {@link #held}.
+ * {@code basicCancel} stops new dispatches; it does not flush one already handed to the
+ * consumer work pool. So a delivery can still land on that pool's thread and reach
+ * {@link #deliverCallback} at any point during, or after, this method's body — and a plain
+ * {@code held.remove(target)} does nothing to stop it: {@code deliverCallback}'s
+ * {@code computeIfAbsent} finds the key gone and happily creates a brand-new map under it, which
+ * this method — already past its {@code remove} — never looks at again. That entry then sits
+ * delivered-but-unacked on {@link #channel} until the whole inbox closes: never requeued, never
+ * redelivered, and {@link #peek} is never called again for a target nothing owns any more.
+ *
+ *
The fix is {@link #held}{@code .compute(target, ...)} instead of {@code remove}: it takes
+ * whatever was held (to nack, same as before) and, in the same atomic step, leaves the
+ * {@link #RELEASED} tombstone behind instead of an absent key. {@code computeIfAbsent} and
+ * {@code compute} calls for the same key are mutually exclusive in {@link ConcurrentHashMap} —
+ * whichever of this call and a concurrent {@code deliverCallback} runs first is fully visible to
+ * the other, with no gap between them. So a delivery that loses the race sees a real map here and
+ * gets nacked by the loop below, same as always; a delivery that wins the race (runs first) is
+ * itself nacked by that same loop, once it settles into {@code held}. A delivery that arrives once
+ * this method has stored {@link #RELEASED} finds it via {@code computeIfAbsent} and refuses itself
+ * — see {@link #deliverCallback}. Either way nothing is silently retained forever, satisfying the
+ * ticket's invariant against dropping a message. This closes the window rather than merely
+ * narrowing it — correctness does not depend on how much time elapses between the swap and this
+ * method returning.
*/
@Override
public void release(String target) {
@@ -255,8 +303,13 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
throw new IllegalStateException("cannot cancel consumer for " + target, e);
}
}
- var perTarget = held.remove(target);
- if (perTarget != null) {
+ AtomicReference> previouslyHeld = new AtomicReference<>();
+ held.compute(target, (_, v) -> {
+ previouslyHeld.set(v);
+ return RELEASED;
+ });
+ var perTarget = previouslyHeld.get();
+ if (perTarget != null && perTarget != RELEASED) {
synchronized (perTarget) {
for (Held h : perTarget.values()) {
try {
@@ -326,7 +379,7 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
@Override
public List peek(String target) {
var perTarget = held.get(target);
- if (perTarget == null) {
+ if (perTarget == null || perTarget == RELEASED) {
return List.of();
}
synchronized (perTarget) {
@@ -337,7 +390,7 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
@Override
public void ack(String target, String msgId) {
var perTarget = held.get(target);
- if (perTarget == null) {
+ if (perTarget == null || perTarget == RELEASED) {
return;
}
Held h;
@@ -370,6 +423,20 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
}
String content = new String(delivery.getBody(), StandardCharsets.UTF_8);
var perTarget = held.computeIfAbsent(target, _ -> new LinkedHashMap<>());
+ if (perTarget == RELEASED) {
+ // CB-318: release() already ran for this target and left the RELEASED tombstone in
+ // held (see release()'s javadoc) — computeIfAbsent() is guaranteed to see it rather
+ // than recreate a fresh map, because ConcurrentHashMap serializes compute/
+ // computeIfAbsent calls for the same key against each other. Refuse the delivery
+ // instead of holding it somewhere release() will never look at again: requeue it, the
+ // same way release() nacks its own held entries, so a later owner (or a connection
+ // drop) can still recover it. This does not need channelLock across a broker round
+ // trip — basicNack, like the duplicate-ack case just below, does not wait for one.
+ synchronized (channelLock) {
+ channel.basicNack(tag, false, true);
+ }
+ return;
+ }
boolean duplicate;
synchronized (perTarget) {
if (perTarget.containsKey(msgId)) {
diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxReleaseRaceTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxReleaseRaceTest.java
new file mode 100644
index 0000000..298547b
--- /dev/null
+++ b/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxReleaseRaceTest.java
@@ -0,0 +1,264 @@
+package dev.ltms.fleet.msg;
+
+import com.rabbitmq.client.AMQP;
+import com.rabbitmq.client.Channel;
+import com.rabbitmq.client.Connection;
+import com.rabbitmq.client.DeliverCallback;
+import com.rabbitmq.client.Delivery;
+import com.rabbitmq.client.Envelope;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.Timeout;
+
+import java.lang.reflect.InvocationHandler;
+import java.lang.reflect.Proxy;
+import java.nio.charset.StandardCharsets;
+import java.time.Duration;
+import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * CB-318: a delivery landing on the consumer work-pool thread after
+ * {@link AmqpReplyInbox#release} has already swapped the target's {@code held} entry for its
+ * tombstone, but before {@code release()} itself returns, must be nacked-with-requeue —
+ * never silently retained in a fresh map {@code release()} has already stopped looking at.
+ *
+ * This forces the actual interleaving, not a sequence of calls. {@code release()}
+ * runs on its own thread and is made to block inside its nack loop, via a fake
+ * {@link Channel} whose {@code basicNack} blocks on its first invocation. That block is only
+ * reachable after {@code release()}'s {@code held.compute(...)} has already swapped in the
+ * {@code RELEASED} tombstone (the compute call happens strictly before the loop that calls
+ * {@code basicNack}), so observing it is direct, ordering-guaranteed proof that the tombstone is in
+ * place and {@code release()} has not yet returned — still holding {@code channelLock} — when a
+ * second thread fires {@code own()}'s captured {@link DeliverCallback} for a brand-new message on the
+ * same target. No mocking library is on the classpath, so the fake broker is a {@link Proxy}, the
+ * same pattern {@code AmqpReplyInboxRecoveryRaceTest} already uses.
+ *
+ *
What this does and does not prove. It proves that a delivery whose
+ * {@code computeIfAbsent} call is ordered strictly after {@code release()}'s tombstone swap — while
+ * {@code release()} is still running — is nacked-with-requeue rather than silently parked forever.
+ * It does not drive a real broker: {@code basicNack} here is a recorded call on a fake channel, not a
+ * verified requeue-and-redeliver. That half of the contract (a nacked-with-requeue delivery really
+ * does come back to a later owner) is already covered against a real broker by
+ * {@code AmqpReplyInboxContractTest.releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery}, which
+ * this test does not duplicate.
+ */
+class AmqpReplyInboxReleaseRaceTest {
+
+ @Test
+ @Timeout(15)
+ void deliveryArrivingWhileReleaseIsStillRunningIsNackedNotStranded() throws Exception {
+ String target = "worker-release-race";
+
+ List nacks = new CopyOnWriteArrayList<>(); // {deliveryTag, requeue(1/0)}
+ List acks = new CopyOnWriteArrayList<>();
+ AtomicReference deliverCallback = new AtomicReference<>();
+ CountDownLatch nackStarted = new CountDownLatch(1);
+ CountDownLatch releaseMayFinishNack = new CountDownLatch(1);
+ AtomicInteger nackCallCount = new AtomicInteger();
+
+ Channel consumeChannel = fakeConsumeChannel(deliverCallback, nacks, acks, nackCallCount,
+ nackStarted, releaseMayFinishNack);
+ Channel publishChannel = fakeInertChannel();
+ Connection connection = fakeConnection(consumeChannel, publishChannel);
+
+ AmqpReplyInbox inbox = new AmqpReplyInbox(connection, AmqpReplyInbox.DEFAULT_PREFETCH);
+ inbox.own(target);
+ assertTrue(deliverCallback.get() != null, "own() must have registered a DeliverCallback");
+
+ // Seed one already-held delivery (m0) so release()'s nack loop has something to iterate, and
+ // therefore somewhere to block, before it can return.
+ deliverCallback.get().handle("ctag", delivery(1L, "m0", "first"));
+
+ AtomicReference releaseError = new AtomicReference<>();
+ Thread releaseThread = new Thread(() -> {
+ try {
+ inbox.release(target);
+ } catch (Throwable t) {
+ releaseError.set(t);
+ }
+ }, "release-under-test");
+ releaseThread.start();
+
+ // This latch only fires from inside the fake channel's basicNack — i.e. from inside
+ // release()'s nack loop, which release()'s code only reaches AFTER held.compute(...) has
+ // already swapped in RELEASED. Waiting for it is direct proof the swap has happened and
+ // release() has not yet returned (it is stuck mid-loop, still holding channelLock).
+ assertTrue(nackStarted.await(10, TimeUnit.SECONDS),
+ "release() never reached its nack loop — it may not have started");
+
+ // The exact interleaving CB-318 describes: a delivery for a NEW message on the same target
+ // lands on the "consumer work-pool thread" (this second thread) while release() is still
+ // running. With the pre-fix code (a bare held.remove(target)) this created a brand-new map
+ // under computeIfAbsent that release() — already past its remove — never looks at again.
+ AtomicReference deliveryError = new AtomicReference<>();
+ Thread deliveryThread = new Thread(() -> {
+ try {
+ deliverCallback.get().handle("ctag", delivery(2L, "m1", "second"));
+ } catch (Throwable t) {
+ deliveryError.set(t);
+ }
+ }, "concurrent-delivery");
+ deliveryThread.start();
+
+ // Head start for the delivery thread to reach (and, on the fixed code, block on)
+ // channelLock — release() still holds it at this point, so a correct fix cannot have
+ // resolved m1's nack yet. Purely in-memory work (computeIfAbsent, a reference compare)
+ // separates deliveryThread.start() from that block point, so 300ms is a large margin, not a
+ // tight timing assumption.
+ Thread.sleep(300);
+ assertEquals(1, nackCallCount.get(),
+ "the concurrent delivery must not resolve its nack before release() gives up "
+ + "channelLock — if this is 2 already, the interleaving below is not being "
+ + "tested, only a sequential call");
+
+ releaseMayFinishNack.countDown(); // let release() finish nacking m0 and return
+ assertTrue(releaseThread.join(Duration.ofSeconds(10)), "release() did not finish");
+ assertTrue(deliveryThread.join(Duration.ofSeconds(10)), "the concurrent delivery did not finish");
+
+ assertNull(releaseError.get(), "release() threw: " + releaseError.get());
+ assertNull(deliveryError.get(), "the concurrent delivery threw: " + deliveryError.get());
+
+ assertEquals(2, nacks.size(),
+ "both the pre-held m0 and the concurrently-arriving m1 must be nacked, got: "
+ + nacks.stream().map(n -> "[tag=" + n[0] + " requeue=" + n[1] + "]").toList());
+ assertTrue(nacks.stream().allMatch(n -> n[1] == 1L),
+ "invariant 1 (never drop): every nack must set requeue=true");
+ assertTrue(nacks.stream().anyMatch(n -> n[0] == 1L), "m0's delivery tag must be nacked");
+ assertTrue(nacks.stream().anyMatch(n -> n[0] == 2L),
+ "m1 — delivered while release() was still running, after the tombstone swap — must be "
+ + "nacked, not silently retained in a map release() will never look at again");
+ assertTrue(acks.isEmpty(), "invariant 1 (never drop): a held reply must never be basicAck'd");
+ }
+
+ private static Delivery delivery(long tag, String msgId, String body) {
+ Envelope envelope = new Envelope(tag, false, "", "irrelevant");
+ AMQP.BasicProperties props = new AMQP.BasicProperties.Builder().messageId(msgId).build();
+ return new Delivery(envelope, props, body.getBytes(StandardCharsets.UTF_8));
+ }
+
+ /** A {@link Proxy}-backed consume {@link Channel}: blocks the FIRST {@code basicNack} call on
+ * {@code releaseMayFinishNack}, after signalling {@code nackStarted} — everything else records
+ * the call and returns a harmless default, matching the style already used by
+ * {@code AmqpReplyInboxRecoveryRaceTest}. */
+ private static Channel fakeConsumeChannel(AtomicReference deliverCallback,
+ List nacks, List acks,
+ AtomicInteger nackCallCount,
+ CountDownLatch nackStarted,
+ CountDownLatch releaseMayFinishNack) {
+ InvocationHandler handler = (proxy, method, args) -> {
+ String name = method.getName();
+ if (name.equals("basicConsume")) {
+ deliverCallback.set((DeliverCallback) args[2]);
+ return "ctag";
+ }
+ if (name.equals("basicNack")) {
+ long tag = (long) args[0];
+ boolean requeue = (boolean) args[2];
+ if (nackCallCount.incrementAndGet() == 1) {
+ nackStarted.countDown();
+ if (!releaseMayFinishNack.await(10, TimeUnit.SECONDS)) {
+ throw new IllegalStateException("test never released the nack latch");
+ }
+ }
+ nacks.add(new long[] {tag, requeue ? 1L : 0L});
+ return null;
+ }
+ if (name.equals("basicAck")) {
+ acks.add((long) args[0]);
+ return null;
+ }
+ if (name.equals("equals")) {
+ return proxy == args[0];
+ }
+ if (name.equals("hashCode")) {
+ return System.identityHashCode(proxy);
+ }
+ if (name.equals("toString")) {
+ return "FakeConsumeChannel";
+ }
+ return defaultValue(method.getReturnType());
+ };
+ return (Channel) Proxy.newProxyInstance(AmqpReplyInboxReleaseRaceTest.class.getClassLoader(),
+ new Class>[] {Channel.class}, handler);
+ }
+
+ /** A {@link Proxy}-backed {@link Channel} that answers every call with a harmless default — used
+ * as the publish channel, which this test never actually publishes on. */
+ private static Channel fakeInertChannel() {
+ InvocationHandler handler = (proxy, method, args) -> {
+ String name = method.getName();
+ if (name.equals("equals")) {
+ return proxy == args[0];
+ }
+ if (name.equals("hashCode")) {
+ return System.identityHashCode(proxy);
+ }
+ if (name.equals("toString")) {
+ return "FakeInertChannel";
+ }
+ return defaultValue(method.getReturnType());
+ };
+ return (Channel) Proxy.newProxyInstance(AmqpReplyInboxReleaseRaceTest.class.getClassLoader(),
+ new Class>[] {Channel.class}, handler);
+ }
+
+ /** A {@link Proxy}-backed {@link Connection} handing out {@code first} then {@code second} from
+ * successive {@code createChannel()} calls, matching {@link AmqpReplyInbox}'s constructor. */
+ private static Connection fakeConnection(Channel first, Channel second) {
+ AtomicInteger calls = new AtomicInteger();
+ InvocationHandler handler = (proxy, method, args) -> {
+ String name = method.getName();
+ if (name.equals("createChannel") && (args == null || args.length == 0)) {
+ return calls.getAndIncrement() == 0 ? first : second;
+ }
+ if (name.equals("equals")) {
+ return proxy == args[0];
+ }
+ if (name.equals("hashCode")) {
+ return System.identityHashCode(proxy);
+ }
+ if (name.equals("toString")) {
+ return "FakeConnection";
+ }
+ return defaultValue(method.getReturnType());
+ };
+ return (Connection) Proxy.newProxyInstance(AmqpReplyInboxReleaseRaceTest.class.getClassLoader(),
+ new Class>[] {Connection.class}, handler);
+ }
+
+ private static Object defaultValue(Class> type) {
+ if (!type.isPrimitive() || type == void.class) {
+ return null;
+ }
+ if (type == boolean.class) {
+ return Boolean.FALSE;
+ }
+ if (type == long.class) {
+ return 0L;
+ }
+ if (type == short.class) {
+ return (short) 0;
+ }
+ if (type == byte.class) {
+ return (byte) 0;
+ }
+ if (type == char.class) {
+ return (char) 0;
+ }
+ if (type == double.class) {
+ return 0.0d;
+ }
+ if (type == float.class) {
+ return 0.0f;
+ }
+ return 0;
+ }
+}