CB-307 Stage 2: AmqpReplyInbox — durable, cross-restart reply delivery

Behind the existing ReplyInbox port, add an AMQP-backed adapter selected by a
`broker:` block in config (absent → the in-memory soft-state inbox; present →
AMQP). Mapping is consume-and-hold with deferred manual ack: each target owns a
durable queue `agent.<target>.inbox`; a manual-ack consumer pulls persistent
messages into an in-memory held map (dedup by msgId) but does not ack; peek
returns the snapshot; ack acks the broker delivery-tag and drops it. A crash
before caller-ack leaves messages unacked, so the broker redelivers on
reconnect — genuine durability with the port contract preserved. bridged still
owns no persistence; the broker does.

- msg/AmqpReplyInbox: the adapter (single synchronized channel; recovery
  listener clears held on reconnect so fresh delivery-tags repopulate).
- config/BridgedConfig: nullable Broker(uri) record; isConfigured() gates it.
- Bridged.main: select adapter; close the AMQP connection in the ordered
  shutdown hook (no-op for the in-memory inbox).
- deps: com.rabbitmq:amqp-client (main); testcontainers rabbitmq/junit-jupiter
  (test). Pinned commons-compress 1.27.1 + commons-lang3 3.18.0 to clear the
  test-scope CVEs those pull. Production default LavinMQ; RabbitMQ URI-swap.
- tests: BridgedConfigTest broker-selection cases (hermetic); AmqpReplyInbox
  contract test (@Tag("contract"), Testcontainers RabbitMQ) proving
  publish/peek/ack, msgId dedup, and cross-restart redelivery. Excluded from
  the default build so `mvn clean install` stays hermetic (210 green).
This commit is contained in:
Dai Ha
2026-07-19 07:30:29 +02:00
parent ba6b4a5da9
commit 2bc5f3a057
7 changed files with 484 additions and 3 deletions
+11
View File
@@ -68,3 +68,14 @@ guard:
# idleTtlSeconds: 300 # idleTtlSeconds: 300
# contextCap: 10 # contextCap: 10
# drainTimeoutSeconds: 5 # drainTimeoutSeconds: 5
# Durable reply delivery (CB-307 Stage 2). OMIT this block entirely to keep the default
# in-memory, soft-state reply inbox (late worker replies are held only until a daemon bounce).
# Set a broker uri to swap in the AMQP-backed inbox: worker replies with no open send are held
# 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.
# broker:
# uri: amqp://guest:guest@127.0.0.1:5672
+46
View File
@@ -24,6 +24,10 @@
<slf4j.version>2.0.16</slf4j.version> <slf4j.version>2.0.16</slf4j.version>
<logback.version>1.5.18</logback.version> <logback.version>1.5.18</logback.version>
<junit.version>5.11.4</junit.version> <junit.version>5.11.4</junit.version>
<amqp.version>5.22.0</amqp.version>
<testcontainers.version>1.20.4</testcontainers.version>
<commons-compress.version>1.27.1</commons-compress.version>
<commons-lang3.version>3.18.0</commons-lang3.version>
</properties> </properties>
<!-- <!--
@@ -62,6 +66,22 @@
<artifactId>jackson-annotations</artifactId> <artifactId>jackson-annotations</artifactId>
<version>3.0-rc5</version> <version>3.0-rc5</version>
</dependency> </dependency>
<!-- Testcontainers 1.20.4 pulls commons-compress 1.24.0 (test scope), which carries
CVE-2024-25710 (8.1) + CVE-2024-26308 — both fixed in 1.26.0. Pin the patched line.
Test-scope only (never shipped in the jar), but bumped per the CVE policy. -->
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-compress</artifactId>
<version>${commons-compress.version}</version>
</dependency>
<!-- Testcontainers 1.20.4 also pulls commons-lang3 3.16.0 (test scope): CVE-2025-48924
(uncontrolled recursion in ClassUtils), fixed in 3.18.0. Pin the patched line.
Test-scope only (never shipped in the jar), bumped per the CVE policy. -->
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-lang3</artifactId>
<version>${commons-lang3.version}</version>
</dependency>
</dependencies> </dependencies>
</dependencyManagement> </dependencyManagement>
@@ -93,6 +113,16 @@
<version>${mcp.version}</version> <version>${mcp.version}</version>
</dependency> </dependency>
<!-- Broker client (CB-307 Stage 2): AMQP 0-9-1. Default deploy targets LavinMQ; this same
client speaks to RabbitMQ unchanged (URI-only swap), so integration tests run against a
stock RabbitMQ container. Only wired when a broker: block is present in config; absent →
the in-memory ReplyInbox. -->
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>${amqp.version}</version>
</dependency>
<!-- Logging --> <!-- Logging -->
<dependency> <dependency>
<groupId>org.slf4j</groupId> <groupId>org.slf4j</groupId>
@@ -112,6 +142,22 @@
<version>${junit.version}</version> <version>${junit.version}</version>
<scope>test</scope> <scope>test</scope>
</dependency> </dependency>
<!-- Testcontainers RabbitMQ: spins a real broker for the @Tag("contract") AMQP integration
test only. Excluded from the default build (contract group), so `mvn clean install`
stays hermetic and green without Docker; run under -Pcontract with Docker present. -->
<dependency>
<groupId>org.testcontainers</groupId>
<artifactId>rabbitmq</artifactId>
<version>${testcontainers.version}</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.testcontainers</groupId>
<artifactId>junit-jupiter</artifactId>
<version>${testcontainers.version}</version>
<scope>test</scope>
</dependency>
</dependencies> </dependencies>
<build> <build>
@@ -15,9 +15,11 @@ import dev.ltms.bridged.mcp.BridgeMcp;
import dev.ltms.bridged.mcp.ConnectionIdentity; import dev.ltms.bridged.mcp.ConnectionIdentity;
import dev.ltms.bridged.mcp.LsofPeerPidLookup; import dev.ltms.bridged.mcp.LsofPeerPidLookup;
import dev.ltms.bridged.mcp.LsofProcessCwdLookup; import dev.ltms.bridged.mcp.LsofProcessCwdLookup;
import dev.ltms.bridged.msg.AmqpReplyInbox;
import dev.ltms.bridged.msg.InMemoryReplyInbox; import dev.ltms.bridged.msg.InMemoryReplyInbox;
import dev.ltms.bridged.msg.MessageService; import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous; import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.msg.ReplyInbox;
import dev.ltms.bridged.rest.BridgedApp; import dev.ltms.bridged.rest.BridgedApp;
import dev.ltms.bridged.session.GitWorktrees; import dev.ltms.bridged.session.GitWorktrees;
import dev.ltms.bridged.session.SessionManager; import dev.ltms.bridged.session.SessionManager;
@@ -117,7 +119,18 @@ public final class Bridged {
StatusPoller poller = new StatusPoller(agents, injector, INJECT_POLL_MILLIS); StatusPoller poller = new StatusPoller(agents, injector, INJECT_POLL_MILLIS);
poller.start(); poller.start();
MessageService messages = new MessageService(agents, injector, rendezvous, new InMemoryReplyInbox()); // CB-307: reply inbox. A broker: block (with a uri) selects the AMQP-backed durable adapter;
// absent, bridged stays soft-state on the in-memory inbox. The AMQP inbox owns a broker
// 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());
} else {
replyInbox = new InMemoryReplyInbox();
log.info("reply inbox: in-memory (soft-state)");
}
MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox);
// MCP server face (CB-105): bridge_send/bridge_reply/bridge_status, mounted at /mcp. // MCP server face (CB-105): bridge_send/bridge_reply/bridge_status, mounted at /mcp.
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments. // Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
@@ -135,6 +148,14 @@ public final class Bridged {
messages.close(); messages.close();
mcp.close(); mcp.close();
if (reaper != null) reaper.stop(); if (reaper != null) reaper.stop();
// Release the broker connection last among message resources (no-op for the in-memory inbox).
if (replyInbox instanceof AutoCloseable closeable) {
try {
closeable.close();
} catch (Exception e) {
log.debug("reply inbox close: {}", e.toString());
}
}
herdr.close(); herdr.close();
})); }));
@@ -31,6 +31,8 @@ import java.util.Set;
* @param spawnReadyTimeoutMs max ms to wait for a spawned worker to reach an injectable state * @param spawnReadyTimeoutMs max ms to wait for a spawned worker to reach an injectable state
* ({@code null} / 0 disables the poll gate — legacy non-blocking behaviour) * ({@code null} / 0 disables the poll gate — legacy non-blocking behaviour)
* @param spawnReadyPollMs poll interval while waiting for the worker to become injectable * @param spawnReadyPollMs poll interval while waiting for the worker to become injectable
* @param broker external AMQP broker for durable reply delivery ({@code null} → in-memory,
* soft-state {@code ReplyInbox}; present → the AMQP-backed adapter, CB-307 Stage 2)
*/ */
@JsonIgnoreProperties(ignoreUnknown = true) @JsonIgnoreProperties(ignoreUnknown = true)
public record BridgedConfig( public record BridgedConfig(
@@ -43,7 +45,8 @@ public record BridgedConfig(
String worktreeRoot, String worktreeRoot,
Lifecycle lifecycle, Lifecycle lifecycle,
Integer spawnReadyTimeoutMs, Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs) { Integer spawnReadyPollMs,
Broker broker) {
@JsonIgnoreProperties(ignoreUnknown = true) @JsonIgnoreProperties(ignoreUnknown = true)
public record Bind(String host, int port) { public record Bind(String host, int port) {
@@ -166,6 +169,24 @@ public record BridgedConfig(
public record Lifecycle(Integer idleTtlSeconds, Integer contextCap, Integer drainTimeoutSeconds) { public record Lifecycle(Integer idleTtlSeconds, Integer contextCap, Integer drainTimeoutSeconds) {
} }
/**
* External AMQP broker for durable, cross-restart reply delivery (CB-307 Stage 2). Its mere
* presence swaps the in-memory {@code ReplyInbox} for the AMQP-backed adapter; absent, bridged
* 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).
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Broker(String uri) {
/** True when a usable broker URI is configured (an empty block does not enable AMQP). */
public boolean isConfigured() {
return uri != null && !uri.isBlank();
}
}
/** /**
* Subscription boundary. Only these hosts may back a worker's * Subscription boundary. Only these hosts may back a worker's
* {@code ANTHROPIC_BASE_URL}; the primary must carry none. * {@code ANTHROPIC_BASE_URL}; the primary must carry none.
@@ -236,6 +257,7 @@ public record BridgedConfig(
Lifecycle l = lifecycle != null ? lifecycle : new Lifecycle(null, null, null); Lifecycle l = lifecycle != null ? lifecycle : new Lifecycle(null, null, null);
Integer timeout = (spawnReadyTimeoutMs != null) ? spawnReadyTimeoutMs : 20000; Integer timeout = (spawnReadyTimeoutMs != null) ? spawnReadyTimeoutMs : 20000;
Integer pollMs = (spawnReadyPollMs != null) ? spawnReadyPollMs : 300; Integer pollMs = (spawnReadyPollMs != null) ? spawnReadyPollMs : 300;
return new BridgedConfig(b, herdrSocket, worker, workers, defaultWorker, g, worktreeRoot, l, timeout, pollMs); // broker is left as-is: null (or an empty/blank uri) keeps the in-memory soft-state inbox.
return new BridgedConfig(b, herdrSocket, worker, workers, defaultWorker, g, worktreeRoot, l, timeout, pollMs, broker);
} }
} }
@@ -0,0 +1,233 @@
package dev.ltms.bridged.msg;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;
import com.rabbitmq.client.Recoverable;
import com.rabbitmq.client.RecoveryListener;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
/**
* AMQP-backed {@link ReplyInbox} (CB-307 Stage 2): genuine cross-restart durability behind the same
* port {@link InMemoryReplyInbox} implements as soft state.
*
* <p><strong>Mapping — consume-and-hold with deferred manual ack.</strong> Each target owns a durable
* queue {@code agent.<target>.inbox}. A manual-ack consumer pulls persistent messages off that queue
* into an in-memory <em>held</em> map (keyed by {@code msgId}) but does <em>not</em> ack them.
* {@link #peek} returns that snapshot; {@link #ack} acks the broker delivery-tag and drops the entry.
* Because messages stay unacked until the primary actually drains them, a crash (or a {@code java -jar}
* bounce) before caller-ack 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>Dedup.</strong> The consumer keys the held map by {@code msgId}; a redelivered duplicate
* (at-least-once, or a producer double-publish) is acked-and-dropped on arrival, so it never
* double-queues. {@link #publish} additionally short-circuits an already-held {@code msgId} — a
* fast path; the consumer-side check is the real guarantee.
*
* <p><strong>Visibility.</strong> Unlike the in-memory adapter, publish → broker → consumer is
* asynchronous, so a {@link #peek} immediately after {@link #publish} may not yet see the message
* (broker delivery latency). Callers that need the reply drained poll (as the primary already does);
* the contract test waits for visibility. This is inherent to broker-backed delivery, not a defect.
*
* <p>The default deploy targets LavinMQ; a stock RabbitMQ speaks the same AMQP 0-9-1 (URI-only swap),
* so the {@code @Tag("contract")} integration test runs against a RabbitMQ container.
*/
public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
private static final Logger log = LoggerFactory.getLogger(AmqpReplyInbox.class);
private static final String QUEUE_PREFIX = "agent.";
private static final String QUEUE_SUFFIX = ".inbox";
private final Connection connection;
private final Channel channel;
/** All channel operations (publish/declare/ack) 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. */
private final ConcurrentHashMap<String, LinkedHashMap<String, Held>> held = new ConcurrentHashMap<>();
/** Targets whose queue is declared and consumer is running. */
private final Set<String> consuming = ConcurrentHashMap.newKeySet();
/** 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. */
public static AmqpReplyInbox open(String uri) {
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"));
} 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). */
AmqpReplyInbox(Connection connection) {
this.connection = connection;
try {
this.channel = connection.createChannel();
} 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).
if (connection instanceof Recoverable recoverable) {
recoverable.addRecoveryListener(new RecoveryListener() {
@Override
public void handleRecovery(Recoverable recoverable) {
held.clear();
log.info("AMQP connection recovered; cleared held replies for fresh redelivery");
}
@Override
public void handleRecoveryStarted(Recoverable recoverable) {
// no-op: we act once recovery completes
}
});
}
}
@Override
public void publish(String target, String msgId, String content) {
ensureConsuming(target);
var perTarget = held.get(target);
if (perTarget != null) {
synchronized (perTarget) {
if (perTarget.containsKey(msgId)) {
return; // already held — producer-side fast dedup
}
}
}
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
.messageId(msgId)
.deliveryMode(2) // persistent — survives a broker restart
.contentType("text/plain")
.build();
try {
synchronized (channelLock) {
channel.basicPublish("", queueName(target), props, content.getBytes(StandardCharsets.UTF_8));
}
} catch (IOException e) {
throw new IllegalStateException("cannot publish reply to " + queueName(target), e);
}
}
@Override
public List<InboxMessage> peek(String target) {
ensureConsuming(target);
var perTarget = held.get(target);
if (perTarget == null) {
return List.of();
}
synchronized (perTarget) {
return perTarget.values().stream().map(Held::message).toList();
}
}
@Override
public void ack(String target, String msgId) {
var perTarget = held.get(target);
if (perTarget == null) {
return;
}
Held h;
synchronized (perTarget) {
h = perTarget.remove(msgId);
}
if (h == null) {
return; // never held (or already acked) — no-op
}
try {
synchronized (channelLock) {
channel.basicAck(h.deliveryTag(), false);
}
} catch (IOException e) {
// Ack didn't reach the broker: restore the entry so a later ack (or a redelivery after
// reconnect) can retry. Keeps the at-least-once contract — a reply is never silently lost.
synchronized (perTarget) {
perTarget.putIfAbsent(msgId, h);
}
throw new IllegalStateException("cannot ack reply " + msgId + " on " + queueName(target), e);
}
}
/** Declare the durable per-target queue and start its manual-ack consumer, once per target. */
private void ensureConsuming(String target) {
if (consuming.contains(target)) {
return;
}
synchronized (channelLock) {
if (!consuming.add(target)) {
return; // another thread just set it up
}
String queue = queueName(target);
try {
channel.queueDeclare(queue, true, false, false, null); // durable, non-exclusive, keep on idle
channel.basicConsume(queue, false, deliverCallback(target), _ -> { });
} catch (IOException e) {
consuming.remove(target);
throw new IllegalStateException("cannot consume queue " + queue, e);
}
}
}
private DeliverCallback deliverCallback(String target) {
return (_, delivery) -> {
String msgId = delivery.getProperties().getMessageId();
long tag = delivery.getEnvelope().getDeliveryTag();
if (msgId == null || msgId.isBlank()) {
msgId = Long.toHexString(tag); // synthesize an id so dedup still has a key
}
String content = new String(delivery.getBody(), StandardCharsets.UTF_8);
var perTarget = held.computeIfAbsent(target, _ -> new LinkedHashMap<>());
boolean duplicate;
synchronized (perTarget) {
if (perTarget.containsKey(msgId)) {
duplicate = true;
} else {
perTarget.put(msgId, new Held(tag, new InboxMessage(msgId, target, content)));
duplicate = false;
}
}
if (duplicate) {
// Redelivered duplicate: ack the new tag and drop it so the broker stops resending.
synchronized (channelLock) {
channel.basicAck(tag, false);
}
}
};
}
private static String queueName(String target) {
return QUEUE_PREFIX + target + QUEUE_SUFFIX;
}
@Override
public void close() {
try {
channel.close();
} catch (Exception e) {
log.debug("AMQP channel close: {}", e.toString());
}
try {
connection.close();
} catch (Exception e) {
log.debug("AMQP connection close: {}", e.toString());
}
}
}
@@ -90,4 +90,39 @@ class BridgedConfigTest {
Files.writeString(f, "bind:\n port: 8080\nfutureFeature:\n enabled: true\n"); Files.writeString(f, "bind:\n port: 8080\nfutureFeature:\n enabled: true\n");
assertDoesNotThrow(() -> BridgedConfig.load(f)); assertDoesNotThrow(() -> BridgedConfig.load(f));
} }
@Test
void absentBrokerBlockLeavesInboxSoftState(@TempDir Path dir) throws Exception {
Path f = dir.resolve("no-broker.yaml");
Files.writeString(f, "bind:\n port: 8080\n");
BridgedConfig cfg = BridgedConfig.load(f);
assertNull(cfg.broker(), "no broker: block → null → in-memory inbox is selected");
}
@Test
void brokerBlockWithUriEnablesAmqpAdapter(@TempDir Path dir) throws Exception {
Path f = dir.resolve("broker.yaml");
Files.writeString(f, """
bind:
port: 8080
broker:
uri: amqp://guest:guest@127.0.0.1:5672/
""");
BridgedConfig cfg = BridgedConfig.load(f);
assertNotNull(cfg.broker());
assertTrue(cfg.broker().isConfigured(), "a non-blank uri enables the AMQP adapter");
assertEquals("amqp://guest:guest@127.0.0.1:5672/", cfg.broker().uri());
}
@Test
void brokerBlockWithBlankUriStaysSoftState(@TempDir Path dir) throws Exception {
Path f = dir.resolve("broker-blank.yaml");
Files.writeString(f, "bind:\n port: 8080\nbroker:\n uri: \"\"\n");
BridgedConfig cfg = BridgedConfig.load(f);
assertNotNull(cfg.broker());
assertFalse(cfg.broker().isConfigured(), "an empty uri must not enable AMQP");
}
} }
@@ -0,0 +1,113 @@
package dev.ltms.bridged.msg;
import org.junit.jupiter.api.Tag;
import org.junit.jupiter.api.Test;
import org.testcontainers.containers.RabbitMQContainer;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;
import org.testcontainers.utility.DockerImageName;
import java.util.List;
import java.util.concurrent.TimeUnit;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* Contract test for {@link AmqpReplyInbox} against a REAL broker (a RabbitMQ container — the same
* AMQP 0-9-1 the production LavinMQ deploy speaks, URI-only swap). Tagged {@code contract} so it is
* excluded from {@code mvn test}/{@code mvn clean install} (which stay hermetic and need no Docker);
* run it with Docker present via {@code mvn test -Pcontract}.
*
* <p>It proves the port contract on genuine infrastructure: eventual visibility of a published reply,
* ack removal, msgId dedup, and — the reason Stage 2 exists — cross-restart durability: an unacked
* reply survives closing the inbox and is redelivered to a fresh connection.
*/
@Tag("contract")
@Testcontainers
class AmqpReplyInboxContractTest {
@Container
static final RabbitMQContainer BROKER =
new RabbitMQContainer(DockerImageName.parse("rabbitmq:3.13-management"));
private static String uri() {
// guest/guest against the mapped AMQP port. No trailing slash: an empty path is vhost "",
// which does not exist — omitting it selects the default vhost "/".
return "amqp://guest:guest@" + BROKER.getHost() + ":" + BROKER.getAmqpPort();
}
@Test
void publishThenPeekThenAck() throws Exception {
String target = "worker-pub-" + System.nanoTime();
try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
inbox.publish(target, "m1", "hello primary");
List<ReplyInbox.InboxMessage> got = awaitPeek(inbox, target);
assertEquals(1, got.size(), "the published reply should be held for drain");
assertEquals("m1", got.getFirst().msgId());
assertEquals(target, got.getFirst().target());
assertEquals("hello primary", got.getFirst().content());
inbox.ack(target, "m1");
assertTrue(inbox.peek(target).isEmpty(), "an acked reply is dropped");
}
}
@Test
void duplicateMsgIdIsNotDoubleQueued() throws Exception {
String target = "worker-dedup-" + System.nanoTime();
try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
inbox.publish(target, "dup", "first");
awaitPeek(inbox, target);
inbox.publish(target, "dup", "second"); // same msgId — must be a no-op
// Give any erroneous second delivery time to land, then assert still exactly one.
Thread.sleep(500);
List<ReplyInbox.InboxMessage> got = inbox.peek(target);
assertEquals(1, got.size(), "a repeated msgId must not double-queue");
assertEquals("first", got.getFirst().content(), "the first payload wins");
}
}
@Test
void unackedReplySurvivesRestartAndIsRedelivered() throws Exception {
String target = "worker-durable-" + System.nanoTime();
// First "process life": publish, see it held, but crash before acking.
try (AmqpReplyInbox first = AmqpReplyInbox.open(uri())) {
first.publish(target, "persist-1", "survive me");
assertEquals(1, awaitPeek(first, target).size());
// no ack — simulate a java -jar bounce with the reply still pending
}
// Second "process life": a fresh connection to the same broker must be redelivered the reply.
try (AmqpReplyInbox second = AmqpReplyInbox.open(uri())) {
List<ReplyInbox.InboxMessage> got = awaitPeek(second, target);
assertEquals(1, got.size(), "an unacked persistent reply is redelivered after restart");
assertEquals("persist-1", got.getFirst().msgId());
assertEquals("survive me", got.getFirst().content());
second.ack(target, "persist-1");
}
// Third life: once acked, it is gone for good — durability is not endless replay.
try (AmqpReplyInbox third = AmqpReplyInbox.open(uri())) {
Thread.sleep(500);
assertTrue(third.peek(target).isEmpty(), "an acked reply does not come back on the next restart");
}
}
/** 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)
throws InterruptedException {
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
List<ReplyInbox.InboxMessage> msgs = inbox.peek(target);
while (msgs.isEmpty() && System.nanoTime() < deadline) {
Thread.sleep(50);
msgs = inbox.peek(target);
}
return msgs;
}
}