From 2bc5f3a057e71cf8a85b83241302ce8d64517428 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 19 Jul 2026 07:30:29 +0200 Subject: [PATCH] =?UTF-8?q?CB-307=20Stage=202:=20AmqpReplyInbox=20?= =?UTF-8?q?=E2=80=94=20durable,=20cross-restart=20reply=20delivery?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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..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). --- bridged/bridged.example.yaml | 11 + bridged/pom.xml | 46 ++++ .../main/java/dev/ltms/bridged/Bridged.java | 23 +- .../ltms/bridged/config/BridgedConfig.java | 26 +- .../dev/ltms/bridged/msg/AmqpReplyInbox.java | 233 ++++++++++++++++++ .../bridged/config/BridgedConfigTest.java | 35 +++ .../msg/AmqpReplyInboxContractTest.java | 113 +++++++++ 7 files changed, 484 insertions(+), 3 deletions(-) create mode 100644 bridged/src/main/java/dev/ltms/bridged/msg/AmqpReplyInbox.java create mode 100644 bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxContractTest.java diff --git a/bridged/bridged.example.yaml b/bridged/bridged.example.yaml index 87f6c6b..ab6644a 100644 --- a/bridged/bridged.example.yaml +++ b/bridged/bridged.example.yaml @@ -68,3 +68,14 @@ guard: # idleTtlSeconds: 300 # contextCap: 10 # 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..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 diff --git a/bridged/pom.xml b/bridged/pom.xml index e166175..d7923ee 100644 --- a/bridged/pom.xml +++ b/bridged/pom.xml @@ -24,6 +24,10 @@ 2.0.16 1.5.18 5.11.4 + 5.22.0 + 1.20.4 + 1.27.1 + 3.18.0 + + org.apache.commons + commons-compress + ${commons-compress.version} + + + + org.apache.commons + commons-lang3 + ${commons-lang3.version} + @@ -93,6 +113,16 @@ ${mcp.version} + + + com.rabbitmq + amqp-client + ${amqp.version} + + org.slf4j @@ -112,6 +142,22 @@ ${junit.version} test + + + + org.testcontainers + rabbitmq + ${testcontainers.version} + test + + + org.testcontainers + junit-jupiter + ${testcontainers.version} + test + diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index d9041a3..75b1374 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -15,9 +15,11 @@ import dev.ltms.bridged.mcp.BridgeMcp; import dev.ltms.bridged.mcp.ConnectionIdentity; import dev.ltms.bridged.mcp.LsofPeerPidLookup; import dev.ltms.bridged.mcp.LsofProcessCwdLookup; +import dev.ltms.bridged.msg.AmqpReplyInbox; import dev.ltms.bridged.msg.InMemoryReplyInbox; import dev.ltms.bridged.msg.MessageService; import dev.ltms.bridged.msg.Rendezvous; +import dev.ltms.bridged.msg.ReplyInbox; import dev.ltms.bridged.rest.BridgedApp; import dev.ltms.bridged.session.GitWorktrees; import dev.ltms.bridged.session.SessionManager; @@ -117,7 +119,18 @@ public final class Bridged { StatusPoller poller = new StatusPoller(agents, injector, INJECT_POLL_MILLIS); 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. // Caller identity is resolved from the connection (peer PID → herdr pane), not arguments. @@ -135,6 +148,14 @@ public final class Bridged { messages.close(); mcp.close(); 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(); })); 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 b7af213..5329997 100644 --- a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java +++ b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java @@ -31,6 +31,8 @@ import java.util.Set; * @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) * @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) public record BridgedConfig( @@ -43,7 +45,8 @@ public record BridgedConfig( String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs, - Integer spawnReadyPollMs) { + Integer spawnReadyPollMs, + Broker broker) { @JsonIgnoreProperties(ignoreUnknown = true) public record Bind(String host, int port) { @@ -166,6 +169,24 @@ public record BridgedConfig( 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 * {@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); Integer timeout = (spawnReadyTimeoutMs != null) ? spawnReadyTimeoutMs : 20000; 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); } } diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/AmqpReplyInbox.java b/bridged/src/main/java/dev/ltms/bridged/msg/AmqpReplyInbox.java new file mode 100644 index 0000000..316ff60 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/msg/AmqpReplyInbox.java @@ -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. + * + *

Mapping — consume-and-hold with deferred manual ack. Each target owns a durable + * queue {@code agent..inbox}. A manual-ack consumer pulls persistent messages off that queue + * into an in-memory held map (keyed by {@code msgId}) but does not 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. + * + *

Dedup. 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. + * + *

Visibility. 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. + * + *

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> held = new ConcurrentHashMap<>(); + /** Targets whose queue is declared and consumer is running. */ + private final Set 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 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()); + } + } +} diff --git a/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java b/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java index f798eab..031f31b 100644 --- a/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java @@ -90,4 +90,39 @@ class BridgedConfigTest { Files.writeString(f, "bind:\n port: 8080\nfutureFeature:\n enabled: true\n"); 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"); + } } diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxContractTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxContractTest.java new file mode 100644 index 0000000..b18f060 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxContractTest.java @@ -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}. + * + *

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 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 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 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 awaitPeek(AmqpReplyInbox inbox, String target) + throws InterruptedException { + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10); + List msgs = inbox.peek(target); + while (msgs.isEmpty() && System.nanoTime() < deadline) { + Thread.sleep(50); + msgs = inbox.peek(target); + } + return msgs; + } +}