CB-528: close the recovery race in AmqpReplyInbox #88
@@ -434,10 +434,14 @@ guard:
|
||||
# on a durable per-target queue (agent.<target>.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
|
||||
|
||||
@@ -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)");
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -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.
|
||||
*
|
||||
* <p><strong>Prefetch bounds the held backlog (CB-527).</strong> 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.
|
||||
*
|
||||
* <p><strong>Publishes require a confirmed, routable delivery (CB-528).</strong> {@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
|
||||
* <em>return</em> for an unroutable message before the <em>confirm</em> 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.
|
||||
*
|
||||
* <p><strong>Ownership is explicit.</strong> {@link #own} declares the queue and starts the consumer;
|
||||
* {@link #release} cancels it. {@link #publish} sends to the queue but does <em>not</em> 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,87 @@ 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<String, String> 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<Long, Pending> pendingBySeq = new ConcurrentSkipListMap<>();
|
||||
/**
|
||||
* 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<String, Pending> 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<Void> 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 +223,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 +236,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,17 +333,115 @@ 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<Long, Pending> 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.
|
||||
* 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
|
||||
* <em>new</em> 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.
|
||||
*
|
||||
* <p>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"));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private static String queueName(String target) {
|
||||
return QUEUE_PREFIX + target + QUEUE_SUFFIX;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
failPendingPublishesOnClose();
|
||||
try {
|
||||
channel.close();
|
||||
} 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) {
|
||||
|
||||
@@ -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.<target>.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<ReplyInbox.InboxMessage> 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<ReplyInbox.InboxMessage> awaitPeek(AmqpReplyInbox inbox, String target)
|
||||
|
||||
@@ -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).
|
||||
*
|
||||
* <p>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<Long> seqOrder = new CopyOnWriteArrayList<>();
|
||||
List<String> msgIdOrder = new CopyOnWriteArrayList<>();
|
||||
AtomicReference<ConfirmCallback> ackCallback = new AtomicReference<>();
|
||||
AtomicReference<ConfirmCallback> 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<Throwable> 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<Throwable> 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<Long> seqOrder = new CopyOnWriteArrayList<>();
|
||||
List<String> msgIdOrder = new CopyOnWriteArrayList<>();
|
||||
AtomicReference<ConfirmCallback> ackCallback = new AtomicReference<>();
|
||||
AtomicReference<ConfirmCallback> 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<Throwable> 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<Long> seqOrder, List<String> msgIdOrder,
|
||||
AtomicReference<ConfirmCallback> ackCallback,
|
||||
AtomicReference<ConfirmCallback> 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;
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user