From 4fa6553db57c31ffb6fcbd6f402ff60bbcd385cf Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 16 Aug 2026 16:55:40 +0200 Subject: [PATCH 1/2] CB-527/CB-528: bound AMQP prefetch and confirm publishes before claiming durable CB-527: basicQos(prefetch) on the consume channel before basicConsume, configurable via broker.prefetch (default 32), so an undrained inbox backlog stays on the broker instead of growing the JVM heap without limit. CB-528: publish moves to its own confirm-mode channel with mandatory=true and a return listener, so an unroutable or unconfirmed reply now throws instead of vanishing silently. The confirm callback checks the per-message returned flag (set by the return listener, which the broker always fires before the matching confirm) so an acked-but-returned publish is still reported as a failure. The ack path stays on its own channel/lock and never waits on a publish confirm. --- bridged/bridged.example.yaml | 8 +- .../main/java/dev/ltms/bridged/Bridged.java | 5 +- .../ltms/bridged/config/BridgedConfig.java | 15 +- .../dev/ltms/bridged/msg/AmqpReplyInbox.java | 187 +++++++++++++++++- .../msg/AmqpReplyInboxContractTest.java | 83 ++++++++ 5 files changed, 282 insertions(+), 16 deletions(-) diff --git a/bridged/bridged.example.yaml b/bridged/bridged.example.yaml index b4d730c..9b5b9eb 100644 --- a/bridged/bridged.example.yaml +++ b/bridged/bridged.example.yaml @@ -434,10 +434,14 @@ guard: # on a durable per-target queue (agent..inbox) and survive a restart — the broker # redelivers anything the primary had not yet drained. Production default is LavinMQ; a stock # RabbitMQ speaks the same AMQP 0-9-1, so it is a URI-only swap. -# uri → AMQP connection URI. No trailing slash ⇒ the default vhost "/"; an empty path ("/") -# is vhost "" and will NOT connect. Encode a named vhost as .../%2Fmyvhost. +# uri → AMQP connection URI. No trailing slash ⇒ the default vhost "/"; an empty path ("/") +# is vhost "" and will NOT connect. Encode a named vhost as .../%2Fmyvhost. +# prefetch → CB-527: consumer basicQos, capping how many unacked messages the inbox holds +# in-heap per owned target (the rest sits on the broker's durable queue instead of +# growing the JVM heap). Default 32 when omitted. # broker: # uri: amqp://guest:guest@127.0.0.1:5672 +# prefetch: 32 # Active push-to-primary (CB-307 Stage 3). When a worker reply lands with no open bridge_send, # the ReplyPushLoop injects a *drain nudge* (never the payload) into the primary's own herdr diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 412a61b..a59695b 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -348,8 +348,9 @@ public final class Bridged { // connection, so keep the reference to close it in the ordered shutdown hook. final ReplyInbox replyInbox; if (cfg.broker() != null && cfg.broker().isConfigured()) { - replyInbox = AmqpReplyInbox.open(cfg.broker().uri()); - log.info("reply inbox: AMQP broker (durable) at {}", cfg.broker().uri()); + replyInbox = AmqpReplyInbox.open(cfg.broker().uri(), cfg.broker().prefetchOrDefault()); + log.info("reply inbox: AMQP broker (durable) at {} (prefetch={})", + cfg.broker().uri(), cfg.broker().prefetchOrDefault()); } else { replyInbox = new InMemoryReplyInbox(); log.info("reply inbox: in-memory (soft-state)"); diff --git a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java index 9bbb187..3140e68 100644 --- a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java +++ b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java @@ -5,6 +5,7 @@ import com.fasterxml.jackson.core.JsonParser; import com.fasterxml.jackson.core.JsonToken; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.dataformat.yaml.YAMLFactory; +import dev.ltms.bridged.msg.AmqpReplyInbox; import dev.ltms.bridged.peer.MemberRole; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -532,16 +533,24 @@ public record BridgedConfig( * stays soft-state. Production default is LavinMQ; a stock RabbitMQ speaks the same AMQP 0-9-1 * and is a URI-only swap. * - * @param uri AMQP connection URI, e.g. {@code amqp://guest:guest@127.0.0.1:5672/}. Blank/{@code null} - * ⇒ the broker block is treated as absent (in-memory adapter). + * @param uri AMQP connection URI, e.g. {@code amqp://guest:guest@127.0.0.1:5672/}. Blank/ + * {@code null} ⇒ the broker block is treated as absent (in-memory adapter). + * @param prefetch CB-527: the consumer's {@code basicQos} prefetch count, bounding how many + * unacked messages the AMQP inbox holds in-heap per owned target. {@code null}/ + * non-positive ⇒ {@link AmqpReplyInbox#DEFAULT_PREFETCH}. */ @JsonIgnoreProperties(ignoreUnknown = true) - public record Broker(String uri) { + public record Broker(String uri, Integer prefetch) { /** True when a usable broker URI is configured (an empty block does not enable AMQP). */ public boolean isConfigured() { return uri != null && !uri.isBlank(); } + + /** The prefetch to use, defaulting to {@link AmqpReplyInbox#DEFAULT_PREFETCH} when unset. */ + public int prefetchOrDefault() { + return (prefetch != null && prefetch > 0) ? prefetch : AmqpReplyInbox.DEFAULT_PREFETCH; + } } /** diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/AmqpReplyInbox.java b/bridged/src/main/java/dev/ltms/bridged/msg/AmqpReplyInbox.java index d116f82..e586372 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/AmqpReplyInbox.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/AmqpReplyInbox.java @@ -7,6 +7,7 @@ import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.DeliverCallback; import com.rabbitmq.client.Recoverable; import com.rabbitmq.client.RecoveryListener; +import com.rabbitmq.client.Return; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -14,7 +15,13 @@ import java.io.IOException; import java.nio.charset.StandardCharsets; import java.util.LinkedHashMap; import java.util.List; +import java.util.NavigableMap; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentSkipListMap; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; /** * AMQP-backed {@link ReplyInbox} (CB-307 Stage 2): genuine cross-restart durability behind the same @@ -29,6 +36,27 @@ import java.util.concurrent.ConcurrentHashMap; * leaves them on the broker — it redelivers on reconnect. That is the durability the in-memory * adapter cannot give, with the port contract preserved. * + *

Prefetch bounds the held backlog (CB-527). The consumer channel calls + * {@code basicQos} with a configurable prefetch count ({@link #DEFAULT_PREFETCH} unless the caller + * passes another value to {@link #open(String, int)}) before starting any consumer. Without a bound, + * the broker pushes its entire queue into {@link #held} the instant a target is {@link #own owned}, + * so an undrained primary grows the JVM heap without limit and any queue-level control + * ({@code x-max-length}, per-message TTL) never fires because the queue never actually holds a + * backlog. Prefetch keeps the backlog where it is visible — on the broker — until the owner drains it. + * + *

Publishes require a confirmed, routable delivery (CB-528). {@link #publish} + * runs on a channel separate from the consume/ack channel ({@link #channel}), so a slow or blocked + * publish confirm can never hold {@link #channelLock} and stall an ack — the ack path never waits on + * a publish confirm. That publish channel is in publisher-confirm mode and every publish sets the + * {@code mandatory} flag, so an unroutable publish (queue not declared, e.g. the owner never called + * {@link #own}) is returned by the broker instead of silently dropped. The broker sends the + * return for an unroutable message before the confirm that covers it — the ack/nack + * callback checks the returned-set at confirm time rather than assuming an ack means routed — so + * "confirmed" here means "durably queued", not merely "accepted by the broker". A returned or nacked + * (or un-confirmed within the timeout) publish surfaces as an {@link IllegalStateException} on the + * caller's thread; the caller — {@link MessageService#reply} — must not report success for a + * black-holed reply. + * *

Ownership is explicit. {@link #own} declares the queue and starts the consumer; * {@link #release} cancels it. {@link #publish} sends to the queue but does not imply ownership * and does not attach a consumer. This split is required by CB-308 federation, where one gateway may @@ -54,6 +82,12 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable { private static final String QUEUE_PREFIX = "agent."; private static final String QUEUE_SUFFIX = ".inbox"; + /** CB-527: the prefetch used when a caller does not pass an explicit value to {@link #open(String, int)}. */ + public static final int DEFAULT_PREFETCH = 32; + + /** How long {@link #publish} waits for its publisher confirm before failing the call (CB-528). */ + private static final long CONFIRM_TIMEOUT_MS = 10_000L; + private final Connection connection; private final Channel channel; /** All channel operations (publish/declare/ack/cancel) serialize on this — a Channel is not thread-safe. */ @@ -63,39 +97,81 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable { /** Targets whose queue is declared and consumer is running, mapped to their broker consumer tag. */ private final ConcurrentHashMap consumerTags = new ConcurrentHashMap<>(); + /** + * CB-528: a dedicated channel for {@link #publish}, kept separate from {@link #channel} (consume + * + ack) so a publish confirm round trip never blocks under {@link #channelLock} and stalls an ack. + */ + private final Channel publishChannel; + private final Object publishChannelLock = new Object(); + /** In-flight publishes awaiting their confirm, keyed by the publish channel's sequence number. */ + private final ConcurrentSkipListMap pendingBySeq = new ConcurrentSkipListMap<>(); + /** The same in-flight publishes, keyed by {@code msgId} — a broker {@code Return} carries no delivery tag. */ + private final ConcurrentHashMap pendingByMsgId = new ConcurrentHashMap<>(); + /** A message pulled off the broker but not yet acked: its delivery-tag plus the port payload. */ private record Held(long deliveryTag, InboxMessage message) {} - /** Connect to {@code uri} (e.g. {@code amqp://guest:guest@127.0.0.1:5672/}) and open the inbox. */ + /** A publish awaiting its confirm; {@link #returned} records whether the broker already returned it. */ + private static final class Pending { + final String msgId; + final CompletableFuture confirmed = new CompletableFuture<>(); + volatile boolean returned; + + Pending(String msgId) { + this.msgId = msgId; + } + } + + /** Connect to {@code uri} (e.g. {@code amqp://guest:guest@127.0.0.1:5672/}) with {@link #DEFAULT_PREFETCH}. */ public static AmqpReplyInbox open(String uri) { + return open(uri, DEFAULT_PREFETCH); + } + + /** As {@link #open(String)}, with an explicit consumer prefetch (CB-527: caps the held backlog per target). */ + public static AmqpReplyInbox open(String uri, int prefetch) { try { ConnectionFactory factory = new ConnectionFactory(); factory.setUri(uri); // Self-heal transient blips; topology recovery re-declares queues and re-attaches consumers. factory.setAutomaticRecoveryEnabled(true); factory.setTopologyRecoveryEnabled(true); - return new AmqpReplyInbox(factory.newConnection("bridged-reply-inbox")); + return new AmqpReplyInbox(factory.newConnection("bridged-reply-inbox"), prefetch); } catch (Exception e) { throw new IllegalStateException("cannot connect to AMQP broker at " + uri, e); } } - /** Wrap an already-open connection (injection seam for the contract test). */ + /** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (injection seam for the contract test). */ AmqpReplyInbox(Connection connection) { + this(connection, DEFAULT_PREFETCH); + } + + /** As above, with an explicit prefetch (injection seam for the contract test). */ + AmqpReplyInbox(Connection connection, int prefetch) { this.connection = connection; try { this.channel = connection.createChannel(); + // CB-527: bound the held backlog per owned target — must be set before any own()/basicConsume. + this.channel.basicQos(prefetch); + this.publishChannel = connection.createChannel(); + this.publishChannel.confirmSelect(); + this.publishChannel.addReturnListener(this::onReturn); + this.publishChannel.addConfirmListener(this::onAck, this::onNack); } catch (IOException e) { throw new IllegalStateException("cannot open AMQP channel", e); } // On automatic recovery the broker redelivers unacked messages with FRESH delivery-tags; the // tags we were holding are now stale. Drop the held snapshot so the re-attached consumer - // repopulates it with valid tags (dedup by msgId still prevents any double-queue). + // repopulates it with valid tags (dedup by msgId still prevents any double-queue). Any publish + // confirm still in flight when the connection dropped is equally stale — its sequence number + // meant nothing on the old channel and means nothing on the recovered one, so fail it now + // rather than let it silently ride out CONFIRM_TIMEOUT_MS. if (connection instanceof Recoverable recoverable) { recoverable.addRecoveryListener(new RecoveryListener() { @Override public void handleRecovery(Recoverable recoverable) { held.clear(); + failPendingPublishesOnRecovery(); log.info("AMQP connection recovered; cleared held replies for fresh redelivery"); } @@ -141,6 +217,12 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable { } } + /** + * Publish {@code content} and block until the broker's publisher confirm for it lands (CB-528). + * Throws {@link IllegalStateException} if the message is returned as unroutable, nacked, or not + * confirmed within {@link #CONFIRM_TIMEOUT_MS} — the caller must treat that as a failed publish, + * not a lost-and-forgotten one. + */ @Override public void publish(String target, String msgId, String content) { AMQP.BasicProperties props = new AMQP.BasicProperties.Builder() @@ -148,12 +230,35 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable { .deliveryMode(2) // persistent — survives a broker restart .contentType("text/plain") .build(); - try { - synchronized (channelLock) { - channel.basicPublish("", queueName(target), props, content.getBytes(StandardCharsets.UTF_8)); + Pending pending = new Pending(msgId); + long seq; + synchronized (publishChannelLock) { + seq = publishChannel.getNextPublishSeqNo(); + pendingBySeq.put(seq, pending); + pendingByMsgId.put(msgId, pending); + try { + publishChannel.basicPublish("", queueName(target), true, props, + content.getBytes(StandardCharsets.UTF_8)); + } catch (IOException e) { + pendingBySeq.remove(seq, pending); + pendingByMsgId.remove(msgId, pending); + throw new IllegalStateException("cannot publish reply to " + queueName(target), e); } - } catch (IOException e) { - throw new IllegalStateException("cannot publish reply to " + queueName(target), e); + } + try { + pending.confirmed.get(CONFIRM_TIMEOUT_MS, TimeUnit.MILLISECONDS); + } catch (ExecutionException e) { + Throwable cause = e.getCause(); + throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause); + } catch (TimeoutException e) { + throw new IllegalStateException("publish confirm for reply " + msgId + " to " + queueName(target) + + " timed out after " + CONFIRM_TIMEOUT_MS + "ms — broker may be unreachable or overloaded", e); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new IllegalStateException("interrupted awaiting publish confirm for " + msgId, e); + } finally { + pendingBySeq.remove(seq, pending); + pendingByMsgId.remove(msgId, pending); } } @@ -222,6 +327,65 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable { }; } + /** Broker return for an unroutable {@code mandatory} publish — arrives BEFORE its confirm (CB-528). */ + private void onReturn(Return r) { + String msgId = r.getProperties() == null ? null : r.getProperties().getMessageId(); + Pending pending = msgId == null ? null : pendingByMsgId.get(msgId); + if (pending != null) { + pending.returned = true; + } else { + log.warn("AMQP return for reply {} (routingKey={}, {} {}) with no matching in-flight publish" + + " — already resolved by a prior confirm", msgId, r.getRoutingKey(), r.getReplyCode(), + r.getReplyText()); + } + } + + private void onAck(long seq, boolean multiple) { + resolveConfirm(seq, multiple, true); + } + + private void onNack(long seq, boolean multiple) { + resolveConfirm(seq, multiple, false); + } + + /** + * Resolve every pending publish covered by this confirm (a single seq, or — {@code multiple} — + * every seq up to and including it). Checks {@link Pending#returned} at confirm time: since the + * broker's return for an unroutable message always precedes its confirm, an ack that arrives after + * a return means "confirmed but never routed", not "durably queued". + */ + private void resolveConfirm(long seq, boolean multiple, boolean ack) { + NavigableMap covered = multiple + ? pendingBySeq.headMap(seq, true) + : pendingBySeq.subMap(seq, true, seq, true); + for (var it = covered.entrySet().iterator(); it.hasNext(); ) { + Pending pending = it.next().getValue(); + it.remove(); + pendingByMsgId.remove(pending.msgId, pending); + if (ack && !pending.returned) { + pending.confirmed.complete(null); + } else if (ack) { + pending.confirmed.completeExceptionally(new IllegalStateException( + "reply " + pending.msgId + " was returned as unroutable (queue not declared/owned)")); + } else { + pending.confirmed.completeExceptionally(new IllegalStateException( + "broker nacked publish of reply " + pending.msgId)); + } + } + } + + /** Fail every publish still awaiting its confirm — their sequence numbers are stale after recovery. */ + private void failPendingPublishesOnRecovery() { + for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) { + Pending pending = it.next().getValue(); + it.remove(); + pendingByMsgId.remove(pending.msgId, pending); + pending.confirmed.completeExceptionally(new IllegalStateException( + "AMQP connection recovered mid-publish; confirm status of reply " + pending.msgId + + " is unknown")); + } + } + private static String queueName(String target) { return QUEUE_PREFIX + target + QUEUE_SUFFIX; } @@ -233,6 +397,11 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable { } catch (Exception e) { log.debug("AMQP channel close: {}", e.toString()); } + try { + publishChannel.close(); + } catch (Exception e) { + log.debug("AMQP publish channel close: {}", e.toString()); + } try { connection.close(); } catch (Exception e) { diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxContractTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxContractTest.java index a67fe48..1c28a54 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxContractTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxContractTest.java @@ -16,6 +16,7 @@ import java.util.List; import java.util.concurrent.TimeUnit; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; /** @@ -166,6 +167,88 @@ class AmqpReplyInboxContractTest { } } + @Test + void prefetchBoundsTheHeldBacklog() throws Exception { + String target = "worker-prefetch-" + System.nanoTime(); + int prefetch = 4; + int published = 10; + try (AmqpReplyInbox inbox = new AmqpReplyInbox(newConnection(), prefetch); + Connection inspect = newConnection()) { + inbox.own(target); + for (int i = 0; i < published; i++) { + inbox.publish(target, "m" + i, "payload " + i); + } + awaitHeldAtLeast(inbox, target, prefetch); + + long depth; + try (Channel ch = inspect.createChannel()) { + depth = ch.queueDeclarePassive(queueName(target)).getMessageCount(); + } + assertTrue(depth >= published - prefetch, + "broker should still hold at least " + (published - prefetch) + + " undelivered messages behind a prefetch of " + prefetch + ", saw " + depth); + + drainUntilEmpty(inbox, target, inspect); + } + } + + @Test + void unroutablePublishReportsFailureNotSilentSuccess() throws Exception { + String target = "worker-unroutable-" + System.nanoTime(); + try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) { + // Deliberately never own(target): the queue is never declared, so the default-exchange + // route to agent..inbox does not exist and the broker must return the publish. + IllegalStateException ex = assertThrows(IllegalStateException.class, + () -> inbox.publish(target, "m1", "nobody home")); + assertTrue(ex.getMessage() != null && ex.getMessage().toLowerCase().contains("unroutable"), + "expected an unroutable-publish failure, got: " + ex.getMessage()); + } + } + + @Test + void confirmedPublishDeliversNormally() throws Exception { + String target = "worker-confirm-" + System.nanoTime(); + try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) { + inbox.own(target); + inbox.publish(target, "m1", "confirmed delivery"); // must return normally: routed and confirmed + + List got = awaitPeek(inbox, target); + assertEquals(1, got.size()); + assertEquals("confirmed delivery", got.getFirst().content()); + inbox.ack(target, "m1"); + } + } + + /** Poll peek until at least {@code n} replies for {@code target} are held, or ~10s elapse. */ + @SuppressWarnings("BusyWait") + private static void awaitHeldAtLeast(AmqpReplyInbox inbox, String target, int n) throws InterruptedException { + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10); + while (inbox.peek(target).size() < n && System.nanoTime() < deadline) { + Thread.sleep(50); + } + } + + /** + * Repeatedly ack whatever is currently held (freeing prefetch slots for the next delivery) until + * both the local snapshot and the broker's own queue depth are empty, or ~10s elapse. + */ + @SuppressWarnings("BusyWait") + private static void drainUntilEmpty(AmqpReplyInbox inbox, String target, Connection inspect) throws Exception { + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10); + while (System.nanoTime() < deadline) { + for (ReplyInbox.InboxMessage msg : inbox.peek(target)) { + inbox.ack(target, msg.msgId()); + } + try (Channel ch = inspect.createChannel()) { + if (ch.queueDeclarePassive(queueName(target)).getMessageCount() == 0 && inbox.peek(target).isEmpty()) { + return; + } + } + Thread.sleep(50); + } + throw new AssertionError("did not drain " + target + " to empty within the deadline"); + } + /** Poll peek (broker delivery is async) until a reply for {@code target} appears or ~10s elapse. */ @SuppressWarnings("BusyWait") // deliberate poll for async broker delivery, bounded by the deadline private static List awaitPeek(AmqpReplyInbox inbox, String target) -- 2.52.0 From c553d795d8c9780d565952010b9851f0a72cd4b9 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 16 Aug 2026 17:24:43 +0200 Subject: [PATCH 2/2] CB-528: close the recovery race in AmqpReplyInbox MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit failPendingPublishesOnRecovery walked and cleared pendingBySeq/pendingByMsgId without holding publishChannelLock, so a publish() that registered while the sweep was still iterating could be failed even though it published on the already-recovered channel — a successful publish reported as failed, and since MessageService.reply mints a fresh msgId per retry, dedup can't catch the resulting duplicate. Guard the sweep with publishChannelLock: publish() only holds it for the seq/map-put/basicPublish, so the sweep can only ever wait for an in-flight basicPublish to return, never a broker round trip. Also make close() fail in-flight publishes immediately with a clear message instead of leaving them to idle out the 10s confirm timeout, and record the (currently unreachable) msgId-uniqueness assumption pendingByMsgId relies on. AmqpReplyInboxRecoveryRaceTest drives the sweep and a real publish() against each other directly (no live broker reconnect) using Proxy-backed fake AMQP channels and a large in-flight backlog to make the race window observable; confirmed it fails without the guard (reply "fresh" wrongly failed as "connection recovered mid-publish") and passes with it. --- .../dev/ltms/bridged/msg/AmqpReplyInbox.java | 60 +++- .../msg/AmqpReplyInboxRecoveryRaceTest.java | 264 ++++++++++++++++++ 2 files changed, 314 insertions(+), 10 deletions(-) create mode 100644 bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxRecoveryRaceTest.java diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/AmqpReplyInbox.java b/bridged/src/main/java/dev/ltms/bridged/msg/AmqpReplyInbox.java index e586372..5f5846d 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/AmqpReplyInbox.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/AmqpReplyInbox.java @@ -105,7 +105,13 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable { private final Object publishChannelLock = new Object(); /** In-flight publishes awaiting their confirm, keyed by the publish channel's sequence number. */ private final ConcurrentSkipListMap pendingBySeq = new ConcurrentSkipListMap<>(); - /** The same in-flight publishes, keyed by {@code msgId} — a broker {@code Return} carries no delivery tag. */ + /** + * The same in-flight publishes, keyed by {@code msgId} — a broker {@code Return} carries no delivery + * tag. Assumes {@code msgId} is unique per in-flight publish: a second {@link #publish} for a + * {@code msgId} still awaiting its confirm would overwrite this entry and misdirect + * {@link #onReturn}'s lookup. Not reachable today — {@code MessageService.reply} generates a fresh + * {@code UUID} per call — so no guard is added for it. + */ private final ConcurrentHashMap pendingByMsgId = new ConcurrentHashMap<>(); /** A message pulled off the broker but not yet acked: its delivery-tag plus the port payload. */ @@ -374,15 +380,48 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable { } } - /** Fail every publish still awaiting its confirm — their sequence numbers are stale after recovery. */ - private void failPendingPublishesOnRecovery() { - for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) { - Pending pending = it.next().getValue(); - it.remove(); - pendingByMsgId.remove(pending.msgId, pending); - pending.confirmed.completeExceptionally(new IllegalStateException( - "AMQP connection recovered mid-publish; confirm status of reply " + pending.msgId - + " is unknown")); + /** + * Fail every publish still awaiting its confirm — their sequence numbers are stale after recovery. + * Guarded by {@link #publishChannelLock}, the same lock {@link #publish} holds while it takes its + * sequence number and registers its {@link Pending}: without it, a {@link #publish} that starts + * after the connection has already recovered (so it publishes — and will be confirmed — on the + * new channel) can register between this sweep's iteration and its clear, and this sweep + * then fails a publish that actually succeeded. {@link #publish} only holds the lock for the + * seq/map-put/{@code basicPublish} — it awaits the confirm outside it — so this sweep can only ever + * wait for an in-flight {@code basicPublish} call to return, never for a broker round trip. No + * deadlock. + * + *

Package-private (rather than {@code private}) only so the unit test can drive it directly + * against a concurrent {@link #publish} without a live broker reconnect. + */ + void failPendingPublishesOnRecovery() { + synchronized (publishChannelLock) { + for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) { + Pending pending = it.next().getValue(); + it.remove(); + pendingByMsgId.remove(pending.msgId, pending); + pending.confirmed.completeExceptionally(new IllegalStateException( + "AMQP connection recovered mid-publish; confirm status of reply " + pending.msgId + + " is unknown")); + } + } + } + + /** + * Fail every publish still awaiting its confirm with a clear, immediate error instead of leaving it + * to time out after {@link #CONFIRM_TIMEOUT_MS} once the channels are closed underneath it. Guarded + * by {@link #publishChannelLock} for the same reason as {@link #failPendingPublishesOnRecovery}. + */ + private void failPendingPublishesOnClose() { + synchronized (publishChannelLock) { + for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) { + Pending pending = it.next().getValue(); + it.remove(); + pendingByMsgId.remove(pending.msgId, pending); + pending.confirmed.completeExceptionally(new IllegalStateException( + "AMQP reply inbox closed while publish of reply " + pending.msgId + + " was still awaiting its confirm")); + } } } @@ -392,6 +431,7 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable { @Override public void close() { + failPendingPublishesOnClose(); try { channel.close(); } catch (Exception e) { diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxRecoveryRaceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxRecoveryRaceTest.java new file mode 100644 index 0000000..4031ea6 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxRecoveryRaceTest.java @@ -0,0 +1,264 @@ +package dev.ltms.bridged.msg; + +import com.rabbitmq.client.AMQP; +import com.rabbitmq.client.Channel; +import com.rabbitmq.client.ConfirmCallback; +import com.rabbitmq.client.Connection; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; + +import java.lang.reflect.InvocationHandler; +import java.lang.reflect.Proxy; +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.AtomicLong; +import java.util.concurrent.atomic.AtomicReference; + +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * CB-528 follow-up: {@link AmqpReplyInbox#failPendingPublishesOnRecovery()} must not fail a publish + * that registers concurrently with (and was not yet visible when) the recovery sweep began — that + * would report a publish that actually succeeded as failed, and {@code MessageService.reply} retries + * with a fresh {@code msgId}, so the reply is delivered twice. A live broker reconnect cannot be + * forced reliably, so this drives {@link AmqpReplyInbox#failPendingPublishesOnRecovery} and a real + * {@link AmqpReplyInbox#publish} against each other directly, against fake AMQP channels built with + * {@link Proxy} (no mocking library is on the classpath). + * + *

Also covers Finding 2 (CB-528 follow-up): {@link AmqpReplyInbox#close()} must fail an in-flight + * publish promptly instead of leaving it to idle out the 10s confirm timeout. + */ +class AmqpReplyInboxRecoveryRaceTest { + + /** Large enough that the (unfixed) unsynchronized sweep's iteration is a real, observable window + * a concurrently-started publish can land in — not just a best case, single-entry sprint. */ + private static final int STALE_PUBLISHES = 100_000; + + @Test + @Timeout(30) + void recoverySweepDoesNotFailAPublishThatRegistersWhileItIsRunning() throws Exception { + AtomicLong seqCounter = new AtomicLong(); + List seqOrder = new CopyOnWriteArrayList<>(); + List msgIdOrder = new CopyOnWriteArrayList<>(); + AtomicReference ackCallback = new AtomicReference<>(); + AtomicReference nackCallback = new AtomicReference<>(); + + Channel publishChannel = fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback); + Channel consumeChannel = fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback); + Connection connection = fakeConnection(consumeChannel, publishChannel); + + AmqpReplyInbox inbox = new AmqpReplyInbox(connection, AmqpReplyInbox.DEFAULT_PREFETCH); + + // STALE_PUBLISHES in-flight publishes that never get confirmed — they sit in pendingBySeq / + // pendingByMsgId exactly like publishes whose confirm never arrived before a connection drop. + // Virtual threads make this many concurrent blocking publish() calls cheap. + CountDownLatch staleStarted = new CountDownLatch(STALE_PUBLISHES); + for (int i = 0; i < STALE_PUBLISHES; i++) { + String msgId = "stale-" + i; + Thread.ofVirtual().start(() -> { + staleStarted.countDown(); + try { + inbox.publish("worker-stale", msgId, "x"); + } catch (IllegalStateException expected) { + // resolved (failed by the sweep) — that is exactly what this thread is here for + } + }); + } + staleStarted.await(); + // Let the registrations (the synchronized put into pendingBySeq/pendingByMsgId) actually land + // for all of them before the sweep starts, so the sweep begins with a large, real backlog. + Thread.sleep(300); + + AtomicReference sweepError = new AtomicReference<>(); + Thread sweepThread = new Thread(() -> { + try { + inbox.failPendingPublishesOnRecovery(); + } catch (Throwable t) { + sweepError.set(t); + } + }, "recovery-sweep"); + sweepThread.start(); + // A short, deliberate head start: with STALE_PUBLISHES this large, the (unfixed) sweep's own + // iteration takes several milliseconds, so this guarantees the sweep has already begun — + // and, once guarded, is already holding publishChannelLock — before "fresh" attempts to + // register. Without this head start, "fresh" sometimes wins the race for the lock and + // registers before the sweep even starts, which is the accepted "already in flight when + // recovery fires" case (correctly failed either way) rather than the bug under test. + Thread.sleep(5); + + // This is the exact interleaving CB-528's follow-up describes: "the still-running recovery + // sweep" racing a publish that registers while it is mid-flight. + AtomicReference publishError = new AtomicReference<>(); + Thread freshThread = new Thread(() -> { + try { + inbox.publish("worker-fresh", "fresh", "hello"); + } catch (Throwable t) { + publishError.set(t); + } + }, "fresh-publish"); + freshThread.start(); + + sweepThread.join(20_000); + assertNull(sweepError.get(), "sweep threw: " + sweepError.get()); + + // Simulate the broker's real confirm for "fresh" now that the sweep is done, so a correct + // implementation's publish() returns normally instead of idling out CONFIRM_TIMEOUT_MS. Poll + // for the registration rather than checking once: freshThread may still be contending for + // publishChannelLock (behind the 20,000 stale threads' own lock acquisitions) even though the + // sweep itself has already finished. + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(9); + int idx = -1; + while (idx < 0 && System.nanoTime() < deadline) { + idx = msgIdOrder.indexOf("fresh"); + if (idx < 0) { + Thread.sleep(20); + } + } + if (idx >= 0 && ackCallback.get() != null) { + ackCallback.get().handle(seqOrder.get(idx), false); + } + freshThread.join(15_000); + + assertNull(publishError.get(), + "a publish that registered while the recovery sweep was running must not be failed by " + + "it, but got: " + publishError.get()); + } + + @Test + @Timeout(15) + void closeFailsInFlightPublishPromptlyInsteadOfWaitingOutTheConfirmTimeout() throws Exception { + AtomicLong seqCounter = new AtomicLong(); + List seqOrder = new CopyOnWriteArrayList<>(); + List msgIdOrder = new CopyOnWriteArrayList<>(); + AtomicReference ackCallback = new AtomicReference<>(); + AtomicReference nackCallback = new AtomicReference<>(); + + Channel publishChannel = fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback); + Channel consumeChannel = fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback); + Connection connection = fakeConnection(consumeChannel, publishChannel); + + AmqpReplyInbox inbox = new AmqpReplyInbox(connection, AmqpReplyInbox.DEFAULT_PREFETCH); + + CountDownLatch publishReturned = new CountDownLatch(1); + AtomicReference publishError = new AtomicReference<>(); + AtomicLong elapsedMillis = new AtomicLong(); + Thread publishThread = new Thread(() -> { + long start = System.nanoTime(); + try { + inbox.publish("worker-close", "never-confirmed", "x"); + } catch (Throwable t) { + publishError.set(t); + } finally { + elapsedMillis.set((System.nanoTime() - start) / 1_000_000); + publishReturned.countDown(); + } + }, "publish-during-close"); + publishThread.start(); + + Thread.sleep(200); // let publish() register before close() runs + inbox.close(); + + assertTrue(publishReturned.await(5, TimeUnit.SECONDS), "publish() did not return after close()"); + assertNotNull(publishError.get(), "a publish in flight when close() runs must fail, not hang"); + assertTrue(publishError.get().getMessage() != null + && publishError.get().getMessage().toLowerCase().contains("closed"), + "expected a clear closed-inbox message, got: " + publishError.get()); + assertTrue(elapsedMillis.get() < 5_000, + "close() should fail the in-flight publish promptly, not wait out the confirm timeout — took " + + elapsedMillis.get() + "ms"); + } + + /** A {@link Proxy}-backed {@link Channel}: only the calls {@link AmqpReplyInbox} actually makes + * are meaningfully implemented; everything else returns a harmless default. */ + private static Channel fakeChannel(AtomicLong seqCounter, List seqOrder, List msgIdOrder, + AtomicReference ackCallback, + AtomicReference nackCallback) { + InvocationHandler handler = (proxy, method, args) -> { + String name = method.getName(); + if (name.equals("getNextPublishSeqNo")) { + long value = seqCounter.incrementAndGet(); + seqOrder.add(value); + return value; + } + if (name.equals("basicPublish")) { + AMQP.BasicProperties props = (AMQP.BasicProperties) args[3]; + msgIdOrder.add(props.getMessageId()); + return null; + } + if (name.equals("addConfirmListener")) { + ackCallback.set((ConfirmCallback) args[0]); + nackCallback.set((ConfirmCallback) args[1]); + return null; + } + if (name.equals("equals")) { + return proxy == args[0]; + } + if (name.equals("hashCode")) { + return System.identityHashCode(proxy); + } + if (name.equals("toString")) { + return "FakeChannel"; + } + return defaultValue(method.getReturnType()); + }; + return (Channel) Proxy.newProxyInstance(AmqpReplyInboxRecoveryRaceTest.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(AmqpReplyInboxRecoveryRaceTest.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; + } +} -- 2.52.0