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
@@ -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");
}
}
@@ -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;
}
}