diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java
index 54094fb..45f71f3 100644
--- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java
+++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java
@@ -209,11 +209,16 @@ public final class Bridged {
MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox,
pushLoop, metrics);
+ // CB-520: the reply inbox only consumes for agents this gateway owns. own on acquire,
+ // release on teardown. Do this before CB-516 so the inbox is owned before any reply can land.
+ sessions.onAcquire(replyInbox::own);
// CB-516: releasing a worker must fail whatever send was waiting on it. Without this a
// torn-down delegation kept reporting PENDING until the 30-minute async timeout, and never
// reached /metrics — the delegation was unresolvable and nothing said so.
- sessions.onRelease(terminal ->
- messages.abandon(terminal, "the worker session was released before it replied"));
+ sessions.onRelease(terminal -> {
+ messages.abandon(terminal, "the worker session was released before it replied");
+ replyInbox.release(terminal);
+ });
// 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.
diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/AmqpReplyInbox.java b/bridged/src/main/java/dev/ltms/bridged/msg/AmqpReplyInbox.java
index 316ff60..d116f82 100644
--- a/bridged/src/main/java/dev/ltms/bridged/msg/AmqpReplyInbox.java
+++ b/bridged/src/main/java/dev/ltms/bridged/msg/AmqpReplyInbox.java
@@ -14,25 +14,30 @@ 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.
+ * Mapping — consume-and-hold with deferred manual ack. Each target has a durable
+ * queue {@code agent..inbox}. The gateway that owns the target starts a manual-ack consumer
+ * ({@link #own}) that 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 owning gateway 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.
+ *
+ * Ownership is explicit. {@link #own} declares the queue and starts the consumer;
+ * {@link #release} cancels it. {@link #publish} sends to the queue but does not imply ownership
+ * and does not attach a consumer. This split is required by CB-308 federation, where one gateway may
+ * publish to an agent owned by another gateway; in that case the publisher must not compete for
+ * deliveries.
*
*
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.
+ * double-queues.
*
*
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
@@ -51,12 +56,12 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
private final Connection connection;
private final Channel channel;
- /** All channel operations (publish/declare/ack) serialize on this — a Channel is not thread-safe. */
+ /** All channel operations (publish/declare/ack/cancel) 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();
+ /** Targets whose queue is declared and consumer is running, mapped to their broker consumer tag. */
+ private final ConcurrentHashMap consumerTags = 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) {}
@@ -103,16 +108,41 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
}
@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
- }
+ public void own(String target) {
+ synchronized (channelLock) {
+ if (consumerTags.containsKey(target)) {
+ return; // already owning this target
+ }
+ String queue = queueName(target);
+ try {
+ channel.queueDeclare(queue, true, false, false, null); // durable, non-exclusive, keep on idle
+ String tag = channel.basicConsume(queue, false, deliverCallback(target), _ -> { });
+ consumerTags.put(target, tag);
+ log.debug("AMQP inbox owns queue {} for target {}", queue, target);
+ } catch (IOException e) {
+ throw new IllegalStateException("cannot own queue " + queue, e);
}
}
+ }
+
+ @Override
+ public void release(String target) {
+ synchronized (channelLock) {
+ String tag = consumerTags.remove(target);
+ held.remove(target); // stale delivery tags must not survive release
+ if (tag == null) {
+ return;
+ }
+ try {
+ channel.basicCancel(tag);
+ } catch (IOException e) {
+ throw new IllegalStateException("cannot cancel consumer for " + target, e);
+ }
+ }
+ }
+
+ @Override
+ public void publish(String target, String msgId, String content) {
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
.messageId(msgId)
.deliveryMode(2) // persistent — survives a broker restart
@@ -129,7 +159,6 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
@Override
public List peek(String target) {
- ensureConsuming(target);
var perTarget = held.get(target);
if (perTarget == null) {
return List.of();
@@ -166,26 +195,6 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
}
}
- /** 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();
diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/InMemoryReplyInbox.java b/bridged/src/main/java/dev/ltms/bridged/msg/InMemoryReplyInbox.java
index fb1e78a..c1b73ec 100644
--- a/bridged/src/main/java/dev/ltms/bridged/msg/InMemoryReplyInbox.java
+++ b/bridged/src/main/java/dev/ltms/bridged/msg/InMemoryReplyInbox.java
@@ -2,6 +2,7 @@ package dev.ltms.bridged.msg;
import java.util.LinkedHashMap;
import java.util.List;
+import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
/**
@@ -9,12 +10,29 @@ import java.util.concurrent.ConcurrentHashMap;
* Per-target FIFO ordering (insertion order via {@link LinkedHashMap}). Dedup by {@code msgId}
* within a target. Thread-safe for concurrent publish vs. drain.
*
+ * Ownership is explicit. {@link #own} marks a target as locally owned so that
+ * {@link #peek} and {@link #ack} operate on it; {@link #publish} works whether or not the target is
+ * owned. {@link #release} clears the local snapshot. This mirrors the AMQP adapter's contract so the
+ * non-broker path stays interchangeable.
+ *
*
This is soft-state, NOT persistence. Lost on a {@code java -jar} bounce — that
* is correct and consistent with "bridged stays soft-state." The Stage-2 AMQP adapter replaces this.
*/
public final class InMemoryReplyInbox implements ReplyInbox {
private final ConcurrentHashMap> store = new ConcurrentHashMap<>();
+ private final Set owned = ConcurrentHashMap.newKeySet();
+
+ @Override
+ public void own(String target) {
+ owned.add(target);
+ }
+
+ @Override
+ public void release(String target) {
+ owned.remove(target);
+ store.remove(target);
+ }
@Override
public void publish(String target, String msgId, String content) {
@@ -27,6 +45,9 @@ public final class InMemoryReplyInbox implements ReplyInbox {
@Override
public List peek(String target) {
+ if (!owned.contains(target)) {
+ return List.of();
+ }
var perTarget = store.get(target);
if (perTarget == null) {
return List.of();
@@ -39,6 +60,9 @@ public final class InMemoryReplyInbox implements ReplyInbox {
@Override
public void ack(String target, String msgId) {
+ if (!owned.contains(target)) {
+ return;
+ }
var perTarget = store.get(target);
if (perTarget != null) {
//noinspection SynchronizationOnLocalVariableOrMethodParameter
diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyInbox.java b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyInbox.java
index 7a03e31..98c9b19 100644
--- a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyInbox.java
+++ b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyInbox.java
@@ -10,16 +10,36 @@ import java.util.List;
* This interface is the port. {@link InMemoryReplyInbox} is the Stage-1 adapter;
* an AMQP-backed adapter (Stage 2) must implement the same contract (idempotent publish, FIFO peek,
* at-least-once ack).
+ *
+ *
Ownership is explicit. A gateway {@link #own owns} the inbox for each agent it
+ * spawned; only the owner consumes and drains it. {@link #publish} sends a reply to the target's
+ * inbox but does not imply ownership or start a consumer. This separation is required by
+ * CB-308 federation, where one gateway may publish to an agent owned by another gateway.
*/
public interface ReplyInbox {
/** A queued reply: an idempotency id, the worker session it came from, and the reply text. */
record InboxMessage(String msgId, String target, String content) {}
+ /**
+ * Start owning (consuming) the inbox for {@code target}. Idempotent: multiple calls for the same
+ * target are no-ops. The owner is the only gateway that may {@link #peek} and {@link #ack} replies
+ * for this target.
+ */
+ void own(String target);
+
+ /**
+ * Stop owning (consuming) the inbox for {@code target}. Idempotent. Any replies held locally but
+ * not yet acked are dropped from the local snapshot; the underlying durable queue keeps
+ * unacked messages for redelivery when the target is re-owned.
+ */
+ void release(String target);
+
/**
* Queue {@code content} from worker {@code target} under {@code msgId}. Idempotent: publishing an
* already-present {@code msgId} for {@code target} is a no-op (dedup), so an at-least-once Stage-2
- * redelivery cannot double-queue.
+ * redelivery cannot double-queue. Publishing does not imply ownership and must not start a
+ * consumer.
*/
void publish(String target, String msgId, String content);
diff --git a/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java b/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java
index d60c24b..9135de5 100644
--- a/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java
+++ b/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java
@@ -48,8 +48,10 @@ public final class SessionManager implements TurnListener {
private final LongSupplier nowNanos;
private final int contextCap;
+ /** CB-520: notified with a terminalId on every acquire; no-op until wired. */
+ private final List> acquireListeners = new java.util.concurrent.CopyOnWriteArrayList<>();
/** CB-516: notified with a terminalId on every release; no-op until wired. */
- private volatile Consumer releaseListener = _ -> { };
+ private final List> releaseListeners = new java.util.concurrent.CopyOnWriteArrayList<>();
/** Backward-compatible constructor: shared-tree sessions, production git seam. */
public SessionManager(PeerLauncher launcher) {
@@ -129,6 +131,7 @@ public final class SessionManager implements TurnListener {
registry.put(handle.id(), session);
log.debug("acquired session id={} terminal={} profile={} owner={}",
handle.id(), handle.terminalId(), session.profile(), session.ownerTerminal());
+ notifyAcquired(session.terminalId());
return session;
}
return acquireWithWorktree(profile, requestedCwd, callerCwd, ownerTerminal, wt);
@@ -151,18 +154,44 @@ public final class SessionManager implements TurnListener {
}
}
+ /**
+ * Register a callback invoked with a session's {@code terminalId} whenever it is acquired
+ * (CB-520). This is the hook that lets the reply inbox {@code own} a target's queue.
+ */
+ public void onAcquire(Consumer listener) {
+ if (listener != null) {
+ acquireListeners.add(listener);
+ }
+ }
+
/**
* Register a callback invoked with a session's {@code terminalId} whenever it is released
* (CB-516). Every teardown path funnels through {@link #release}, so one hook covers the REST
* and MCP stop tools, the idle-TTL reaper, {@code recycle}, and shutdown drain alike.
*
- * Set rather than injected because {@code MessageService} — the intended listener — is
+ *
Added rather than injected because {@code MessageService} — one intended listener — is
* constructed after this manager (it needs the injector and rendezvous, which need the session
* presence view this manager exposes). Wiring it at construction would require breaking that
* cycle for one callback.
*/
public void onRelease(Consumer listener) {
- this.releaseListener = (listener == null) ? _ -> { } : listener;
+ if (listener != null) {
+ releaseListeners.add(listener);
+ }
+ }
+
+ /** A listener failure must never prevent the acquisition it is reacting to. */
+ private void notifyAcquired(String terminalId) {
+ if (terminalId == null) {
+ return;
+ }
+ for (Consumer listener : acquireListeners) {
+ try {
+ listener.accept(terminalId);
+ } catch (RuntimeException e) {
+ log.warn("acquire listener failed for terminal {}: {}", terminalId, e.toString());
+ }
+ }
}
/** A listener failure must never prevent the teardown it is reacting to. */
@@ -170,10 +199,12 @@ public final class SessionManager implements TurnListener {
if (terminalId == null) {
return;
}
- try {
- releaseListener.accept(terminalId);
- } catch (RuntimeException e) {
- log.warn("release listener failed for terminal {}: {}", terminalId, e.toString());
+ for (Consumer listener : releaseListeners) {
+ try {
+ listener.accept(terminalId);
+ } catch (RuntimeException e) {
+ log.warn("release listener failed for terminal {}: {}", terminalId, e.toString());
+ }
}
}
@@ -221,6 +252,7 @@ public final class SessionManager implements TurnListener {
registry.put(handle.id(), session);
log.debug("acquired worktree session id={} terminal={} profile={} branch={} path={}",
handle.id(), handle.terminalId(), session.profile(), session.branch(), session.worktree());
+ notifyAcquired(session.terminalId());
return session;
}
diff --git a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java
index d06ebc3..1b4c054 100644
--- a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java
+++ b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java
@@ -15,6 +15,8 @@ import dev.ltms.bridged.session.WorkerSession;
import dev.ltms.bridged.session.WorktreeRequest;
import dev.ltms.bridged.worker.ClaudeCodeLauncher;
import io.modelcontextprotocol.spec.McpSchema;
+import dev.ltms.bridged.msg.InMemoryReplyInbox;
+import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import java.util.Map;
@@ -31,10 +33,19 @@ import static org.junit.jupiter.api.Assertions.*;
*/
class BridgeMcpTest {
+ private static final String T = "term_a";
+
private final FakeHerdr herdr = new FakeHerdr();
private final AgentControl agents = new AgentControl(herdr);
private final Rendezvous rendezvous = new Rendezvous();
- private final MessageService messages = new MessageService(agents, new Injector(agents), rendezvous);
+ private final InMemoryReplyInbox inbox = new InMemoryReplyInbox();
+ private final MessageService messages = new MessageService(agents, new Injector(agents), rendezvous, inbox);
+
+ @BeforeEach
+ void setUp() {
+ // CB-520: the inbox only peeks/acks targets it owns.
+ inbox.own(T);
+ }
private static String textOf(McpSchema.CallToolResult r) {
return ((McpSchema.TextContent) r.content().getFirst()).text();
diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxContractTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxContractTest.java
index b18f060..47e27f9 100644
--- a/bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxContractTest.java
+++ b/bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxContractTest.java
@@ -1,5 +1,9 @@
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 org.junit.jupiter.api.Tag;
import org.junit.jupiter.api.Test;
import org.testcontainers.containers.RabbitMQContainer;
@@ -20,8 +24,8 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
* 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.
+ * ack removal, msgId dedup, explicit ownership, 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
@@ -41,6 +45,7 @@ class AmqpReplyInboxContractTest {
void publishThenPeekThenAck() throws Exception {
String target = "worker-pub-" + System.nanoTime();
try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
+ inbox.own(target);
inbox.publish(target, "m1", "hello primary");
List got = awaitPeek(inbox, target);
@@ -58,6 +63,7 @@ class AmqpReplyInboxContractTest {
void duplicateMsgIdIsNotDoubleQueued() throws Exception {
String target = "worker-dedup-" + System.nanoTime();
try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
+ inbox.own(target);
inbox.publish(target, "dup", "first");
awaitPeek(inbox, target);
inbox.publish(target, "dup", "second"); // same msgId — must be a no-op
@@ -74,8 +80,9 @@ class AmqpReplyInboxContractTest {
void unackedReplySurvivesRestartAndIsRedelivered() throws Exception {
String target = "worker-durable-" + System.nanoTime();
- // First "process life": publish, see it held, but crash before acking.
+ // First "process life": own, publish, see it held, but crash before acking.
try (AmqpReplyInbox first = AmqpReplyInbox.open(uri())) {
+ first.own(target);
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
@@ -83,6 +90,7 @@ class AmqpReplyInboxContractTest {
// Second "process life": a fresh connection to the same broker must be redelivered the reply.
try (AmqpReplyInbox second = AmqpReplyInbox.open(uri())) {
+ second.own(target);
List got = awaitPeek(second, target);
assertEquals(1, got.size(), "an unacked persistent reply is redelivered after restart");
assertEquals("persist-1", got.getFirst().msgId());
@@ -93,11 +101,43 @@ class AmqpReplyInboxContractTest {
// Third life: once acked, it is gone for good — durability is not endless replay.
try (AmqpReplyInbox third = AmqpReplyInbox.open(uri())) {
+ third.own(target);
Thread.sleep(500);
assertTrue(third.peek(target).isEmpty(), "an acked reply does not come back on the next restart");
}
}
+ @Test
+ void publishDoesNotAttachAConsumer() throws Exception {
+ String target = "worker-pub-no-consumer-" + System.nanoTime();
+ try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri());
+ Connection inspect = newConnection()) {
+ inbox.own(target);
+ inbox.publish(target, "m1", "published");
+ awaitPeek(inbox, target); // ensure the owner's consumer received it
+
+ try (Channel ch = inspect.createChannel()) {
+ AMQP.Queue.DeclareOk ok = ch.queueDeclare(queueName(target), true, false, false, null);
+ assertEquals(1, ok.getConsumerCount(),
+ "publish must not attach a consumer; only the owner's consumer should exist");
+ }
+ }
+ }
+
+ @Test
+ void releaseCancelsConsumerAndClearsHeld() throws Exception {
+ String target = "worker-release-" + System.nanoTime();
+ try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
+ inbox.own(target);
+ inbox.publish(target, "m1", "release me");
+ assertEquals(1, awaitPeek(inbox, target).size(), "owned target holds the reply");
+
+ inbox.release(target);
+ assertTrue(inbox.peek(target).isEmpty(),
+ "release clears the local held snapshot");
+ }
+ }
+
/** 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)
@@ -110,4 +150,15 @@ class AmqpReplyInboxContractTest {
}
return msgs;
}
+
+ /** A separate broker connection for inspecting queue state without disturbing the inbox. */
+ private static Connection newConnection() throws Exception {
+ ConnectionFactory factory = new ConnectionFactory();
+ factory.setUri(uri());
+ return factory.newConnection("contract-inspector");
+ }
+
+ private static String queueName(String target) {
+ return "agent." + target + ".inbox";
+ }
}
diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/InMemoryReplyInboxTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/InMemoryReplyInboxTest.java
index e040de3..676337b 100644
--- a/bridged/src/test/java/dev/ltms/bridged/msg/InMemoryReplyInboxTest.java
+++ b/bridged/src/test/java/dev/ltms/bridged/msg/InMemoryReplyInboxTest.java
@@ -1,5 +1,6 @@
package dev.ltms.bridged.msg;
+import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import java.util.concurrent.CountDownLatch;
@@ -10,13 +11,18 @@ import java.util.concurrent.atomic.AtomicReference;
import static org.junit.jupiter.api.Assertions.*;
/**
- * Unit tests for {@link InMemoryReplyInbox}: publish, peek, ack, dedup, FIFO ordering, and thread
- * safety under concurrent publish vs. drain.
+ * Unit tests for {@link InMemoryReplyInbox}: publish, peek, ack, dedup, FIFO ordering, thread
+ * safety under concurrent publish vs. drain, and explicit ownership.
*/
class InMemoryReplyInboxTest {
private final ReplyInbox inbox = new InMemoryReplyInbox();
+ @BeforeEach
+ void setUp() {
+ inbox.own("term_a");
+ }
+
@Test
void publishThenPeekReturnsTheMessage() {
inbox.publish("term_a", "m1", "hello");
@@ -72,6 +78,7 @@ class InMemoryReplyInboxTest {
@Test
void perTargetIsolation() {
+ inbox.own("term_b");
inbox.publish("term_a", "m1", "for-a");
inbox.publish("term_b", "m2", "for-b");
assertEquals(1, inbox.peek("term_a").size());
@@ -142,4 +149,40 @@ class InMemoryReplyInboxTest {
exec.shutdown();
}
}
+
+ @Test
+ void publishWithoutOwnDoesNotClaimOwnership() {
+ String unowned = "term_unowned";
+ inbox.publish(unowned, "m1", "hello");
+ // Without an owner, peek returns nothing — publish did not imply consume.
+ assertTrue(inbox.peek(unowned).isEmpty(),
+ "publishing to an unowned target must not make it peekable");
+ // Owning afterwards makes the already-published message visible.
+ inbox.own(unowned);
+ var msgs = inbox.peek(unowned);
+ assertEquals(1, msgs.size());
+ assertEquals("m1", msgs.getFirst().msgId());
+ }
+
+ @Test
+ void peekAndAckAreNoOpsForUnownedTarget() {
+ assertTrue(inbox.peek("term_not_owned").isEmpty());
+ inbox.ack("term_not_owned", "m1"); // no-op, should not throw
+ }
+
+ @Test
+ void releaseStopsConsumingAndClearsHeld() {
+ inbox.publish("term_a", "m1", "hello");
+ assertEquals(1, inbox.peek("term_a").size());
+ inbox.release("term_a");
+ assertTrue(inbox.peek("term_a").isEmpty(),
+ "after release, the local snapshot is cleared");
+ }
+
+ @Test
+ void ownIsIdempotent() {
+ inbox.own("term_a");
+ inbox.publish("term_a", "m1", "hello");
+ assertEquals(1, inbox.peek("term_a").size());
+ }
}
diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java
index 92f1b0c..66d11f2 100644
--- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java
+++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java
@@ -6,6 +6,7 @@ import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.inject.CompletionResolver;
import dev.ltms.bridged.inject.Injector;
+import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import java.util.concurrent.CompletableFuture;
@@ -30,7 +31,14 @@ class MessageServiceTest {
private final Rendezvous rendezvous = new Rendezvous();
private final CompletionResolver completion = new CompletionResolver(agents, rendezvous);
private final Injector injector = new Injector(agents, completion);
- private final MessageService messages = new MessageService(agents, injector, rendezvous);
+ private final InMemoryReplyInbox inbox = new InMemoryReplyInbox();
+ private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
+
+ @BeforeEach
+ void setUp() {
+ // CB-520: the inbox only peeks/acks targets it owns.
+ inbox.own(T);
+ }
/** Run {@code send} on a background thread; the current thread drives the worker's turn. */
private CompletableFuture sendAsync() {
diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java
index dbbf859..920d54e 100644
--- a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java
+++ b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java
@@ -45,6 +45,7 @@ class ReplyPushLoopTest {
void setUp() {
registry = new PrimaryRegistry(PRIMARY);
inbox = new InMemoryReplyInbox();
+ inbox.own(WORKER); // CB-520: the inbox only peeks/acks targets it owns
scheduler = Executors.newSingleThreadScheduledExecutor();
}
diff --git a/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java b/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java
index d546ca4..960b96b 100644
--- a/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java
+++ b/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java
@@ -10,6 +10,7 @@ import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.inject.Injector;
import dev.ltms.bridged.inject.StatusPoller;
import dev.ltms.bridged.inject.WorkerPresence;
+import dev.ltms.bridged.msg.InMemoryReplyInbox;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.session.FakeWorktrees;
@@ -74,7 +75,12 @@ class BridgedAppTest {
poller = new StatusPoller(agents, injector, 5); // delivers when the fake reports idle
poller.start();
Rendezvous rendezvous = new Rendezvous();
- MessageService messages = new MessageService(agents, injector, rendezvous);
+ InMemoryReplyInbox inbox = new InMemoryReplyInbox();
+ sessions.onAcquire(inbox::own);
+ // The static REST tests address "term_a" without acquiring it through SessionManager, so own
+ // it directly so the inbox contract holds for those endpoints.
+ inbox.own("term_a");
+ MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
app = new BridgedApp(herdr, workers, sessions, messages, this.presence, null)
.build().start("127.0.0.1", 0);
return app.port();