From f32f90ffb1bca806e903c54c86ee9135136c939d Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Tue, 4 Aug 2026 17:44:21 +0200 Subject: [PATCH] CB-520: split ReplyInbox into explicit own/release and publish halves --- .../main/java/dev/ltms/bridged/Bridged.java | 9 +- .../dev/ltms/bridged/msg/AmqpReplyInbox.java | 93 ++++++++++--------- .../ltms/bridged/msg/InMemoryReplyInbox.java | 24 +++++ .../java/dev/ltms/bridged/msg/ReplyInbox.java | 22 ++++- .../ltms/bridged/session/SessionManager.java | 46 +++++++-- .../dev/ltms/bridged/mcp/BridgeMcpTest.java | 13 ++- .../msg/AmqpReplyInboxContractTest.java | 57 +++++++++++- .../bridged/msg/InMemoryReplyInboxTest.java | 47 +++++++++- .../ltms/bridged/msg/MessageServiceTest.java | 10 +- .../ltms/bridged/msg/ReplyPushLoopTest.java | 1 + .../dev/ltms/bridged/rest/BridgedAppTest.java | 8 +- 11 files changed, 270 insertions(+), 60 deletions(-) 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(); -- 2.52.0