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:
@@ -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
|
||||||
|
|||||||
@@ -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;
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user