From ba6b4a5da916c1fbfd1927d94ea10c5c3fe7234b Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 18 Jul 2026 21:12:51 +0200 Subject: [PATCH] =?UTF-8?q?CB-307=20Stage=201:=20reply-inbox=20port=20+=20?= =?UTF-8?q?in-memory=20adapter=20=E2=80=94=20hold=20stranded=20worker=20re?= =?UTF-8?q?plies=20instead=20of=20dropping=20them?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Problem: the reverse (worker->primary) path was Rendezvous, a map of LIVE blocking waiters only. A bridge_reply arriving with no open send hit Rendezvous.complete() -> no waiter -> returned false -> the reply was silently DISCARDED (worker saw an error / REST 409). No message-id/dedup/ack existed anywhere. Stage 1 (no broker, soft-state) behind one port: - ReplyInbox port + InboxMessage record; InMemoryReplyInbox adapter (per-target FIFO via LinkedHashMap, dedup by msgId, thread-safe). Soft-state, not persistence. - MessageService.reply(session, content): resolve an open send, else publish to the inbox with a minted UUID (was a silent drop). drainReplies(target) = peek + ack. - BridgeMcp.reply / BridgedApp.replyMessage repointed off bare Rendezvous.resolve onto messages.reply -> no-waiter is now SUCCESS (queued), not error / 409. - Drain surface: bridge_poll gains optional target; REST GET /sessions/{id}/replies. - Rendezvous left untouched. QUESTION path (bridge_ask) NOT queued (interactive, keeps NO_WAITER); completion/failure fallbacks NOT queued (captured-waiter). - Bridged.main wires new InMemoryReplyInbox(); no broker: config yet (Stage 2 = AMQP). Tests: +19 (188 -> 207), 0 failures/0 errors. New InMemoryReplyInboxTest (12) + MessageService/BridgeMcp/BridgedApp coverage incl. guards proving a QUESTION and a completion fallback are never queued. Implemented via delegation to an off-sub gx10 worker in a pre-trusted worktree; primary-verified (mvn clean install green, 207 tests) and committed by the primary because the worker's completion replies were lost to the very bug this fixes. Refs CB-307 (gitea #5), Stage 1 of 2. --- .../main/java/dev/ltms/bridged/Bridged.java | 3 +- .../java/dev/ltms/bridged/mcp/BridgeMcp.java | 40 +++-- .../ltms/bridged/msg/InMemoryReplyInbox.java | 50 ++++++ .../dev/ltms/bridged/msg/MessageService.java | 46 +++++- .../java/dev/ltms/bridged/msg/ReplyInbox.java | 31 ++++ .../dev/ltms/bridged/rest/BridgedApp.java | 26 +++- .../dev/ltms/bridged/mcp/BridgeMcpTest.java | 69 ++++++--- .../bridged/msg/InMemoryReplyInboxTest.java | 145 ++++++++++++++++++ .../ltms/bridged/msg/MessageServiceTest.java | 89 +++++++++++ .../dev/ltms/bridged/rest/BridgedAppTest.java | 48 +++--- 10 files changed, 480 insertions(+), 67 deletions(-) create mode 100644 bridged/src/main/java/dev/ltms/bridged/msg/InMemoryReplyInbox.java create mode 100644 bridged/src/main/java/dev/ltms/bridged/msg/ReplyInbox.java create mode 100644 bridged/src/test/java/dev/ltms/bridged/msg/InMemoryReplyInboxTest.java diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index e54d4e8..d9041a3 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -15,6 +15,7 @@ import dev.ltms.bridged.mcp.BridgeMcp; import dev.ltms.bridged.mcp.ConnectionIdentity; import dev.ltms.bridged.mcp.LsofPeerPidLookup; import dev.ltms.bridged.mcp.LsofProcessCwdLookup; +import dev.ltms.bridged.msg.InMemoryReplyInbox; import dev.ltms.bridged.msg.MessageService; import dev.ltms.bridged.msg.Rendezvous; import dev.ltms.bridged.rest.BridgedApp; @@ -116,7 +117,7 @@ public final class Bridged { StatusPoller poller = new StatusPoller(agents, injector, INJECT_POLL_MILLIS); poller.start(); - MessageService messages = new MessageService(agents, injector, rendezvous); + MessageService messages = new MessageService(agents, injector, rendezvous, new InMemoryReplyInbox()); // 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/mcp/BridgeMcp.java b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java index c58cee0..104eb02 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java @@ -96,14 +96,16 @@ public final class BridgeMcp { }) // bridge_reply's identity is the CONNECTION, never an argument. .toolCall(replyTool(), (exchange, req) -> - reply(rendezvous, callerTerminal(exchange), str(req.arguments(), "content"))) + reply(messages, callerTerminal(exchange), str(req.arguments(), "content"))) // bridge_ask (CB-205): a worker's mid-turn question — identity from the CONNECTION. .toolCall(askTool(), (exchange, req) -> ask(messages, callerTerminal(exchange), str(req.arguments(), "question"), timeoutMs(req.arguments()))) .toolCall(statusTool(), (_, req) -> status(messages, str(req.arguments(), "sessionId"))) - .toolCall(pollTool(), (_, req) -> - poll(messages, str(req.arguments(), "ticket"))) + .toolCall(pollTool(), (_, req) -> { + Map a = req.arguments(); + return poll(messages, str(a, "ticket"), str(a, "target")); + }) // Fleet management (CB-108): spawn/list/stop over ClaudeCodeLauncher. .toolCall(spawnTool(), (exchange, req) -> { Map a = req.arguments(); @@ -236,10 +238,17 @@ public final class BridgeMcp { return text("accepted — task delegated. Poll bridge_poll with ticket=" + ticket); } - /** {@code bridge_poll}: check an async delegation by ticket (pending / done+reply / failed). */ - static McpSchema.CallToolResult poll(MessageService messages, String ticket) { + /** {@code bridge_poll}: check an async delegation by ticket, or drain a worker's inbox by target. */ + static McpSchema.CallToolResult poll(MessageService messages, String ticket, String target) { + if (!isBlank(target)) { + var replies = messages.drainReplies(target); + if (replies.isEmpty()) { + return text("[]"); + } + return text(json(replies)); + } if (isBlank(ticket)) { - return error("ticket is required"); + return error("ticket (or target) is required"); } MessageService.TaskView v = messages.poll(ticket); if (v == null) { @@ -255,11 +264,12 @@ public final class BridgeMcp { } /** - * {@code bridge_reply}: the worker returns its structured answer, resolving the awaiting send. + * {@code bridge_reply}: the worker returns its structured answer, resolving the awaiting send + * or — when no send is open — queueing the reply in the inbox for later drain (CB-307). * {@code callerTerminal} is resolved from the connection (never an argument); a {@code null} * means the caller is not a known worker (e.g. the primary called it by mistake). */ - static McpSchema.CallToolResult reply(Rendezvous rendezvous, String callerTerminal, String content) { + static McpSchema.CallToolResult reply(MessageService messages, String callerTerminal, String content) { if (callerTerminal == null) { return error("bridge_reply is for workers only — could not identify the calling worker " + "from the connection"); @@ -267,9 +277,8 @@ public final class BridgeMcp { if (content == null) { return error("content is required"); } - return rendezvous.resolve(callerTerminal, content) - ? text("delivered") - : error("no send is awaiting a reply for this worker"); + messages.reply(callerTerminal, content); + return text("delivered"); } /** {@code bridge_status}: the live lifecycle status of a worker session. */ @@ -438,10 +447,13 @@ public final class BridgeMcp { private static McpSchema.Tool pollTool() { return tool("bridge_poll", "Check an async delegation (a bridge_send with wait:false) by its ticket: " - + "pending, done (with the worker's reply), or failed.", + + "pending, done (with the worker's reply), or failed. When target (a worker " + + "session id) is present instead of ticket, drain that worker's inbox of " + + "replies delivered when no send was open.", objectSchema(Map.of( - "ticket", stringProp("The ticket returned by bridge_send wait:false")), - List.of("ticket"))); + "ticket", stringProp("The ticket returned by bridge_send wait:false"), + "target", stringProp("Worker session id to drain pending replies from (optional)")), + List.of())); } private static McpSchema.Tool spawnTool() { diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/InMemoryReplyInbox.java b/bridged/src/main/java/dev/ltms/bridged/msg/InMemoryReplyInbox.java new file mode 100644 index 0000000..fb1e78a --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/msg/InMemoryReplyInbox.java @@ -0,0 +1,50 @@ +package dev.ltms.bridged.msg; + +import java.util.LinkedHashMap; +import java.util.List; +import java.util.concurrent.ConcurrentHashMap; + +/** + * Soft-state {@link ReplyInbox} backed by a {@link ConcurrentHashMap} keyed by target session. + * Per-target FIFO ordering (insertion order via {@link LinkedHashMap}). Dedup by {@code msgId} + * within a target. Thread-safe for concurrent publish vs. drain. + * + *

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<>(); + + @Override + public void publish(String target, String msgId, String content) { + var perTarget = store.computeIfAbsent(target, _ -> new LinkedHashMap<>()); + //noinspection SynchronizationOnLocalVariableOrMethodParameter + synchronized (perTarget) { + perTarget.putIfAbsent(msgId, new InboxMessage(msgId, target, content)); + } + } + + @Override + public List peek(String target) { + var perTarget = store.get(target); + if (perTarget == null) { + return List.of(); + } + //noinspection SynchronizationOnLocalVariableOrMethodParameter + synchronized (perTarget) { + return List.copyOf(perTarget.values()); + } + } + + @Override + public void ack(String target, String msgId) { + var perTarget = store.get(target); + if (perTarget != null) { + //noinspection SynchronizationOnLocalVariableOrMethodParameter + synchronized (perTarget) { + perTarget.remove(msgId); + } + } + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java index 2c63c5c..046dadd 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -6,6 +6,8 @@ import dev.ltms.bridged.inject.Injector; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.List; +import java.util.UUID; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionException; import java.util.concurrent.ConcurrentHashMap; @@ -147,16 +149,24 @@ public final class MessageService { private final AgentControl agents; private final Injector injector; private final Rendezvous rendezvous; + private final ReplyInbox inbox; private final ConcurrentHashMap sessionLocks = new ConcurrentHashMap<>(); private final ConcurrentHashMap tasks = new ConcurrentHashMap<>(); private final AtomicLong ticketSeq = new AtomicLong(); private final ExecutorService asyncExecutor = Executors.newThreadPerTaskExecutor( Thread.ofVirtual().name("bridge-async-", 0).factory()); - public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous) { + /** Create with an explicit {@link ReplyInbox} (Stage 1: {@link InMemoryReplyInbox}). */ + public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox) { this.agents = agents; this.injector = injector; this.rendezvous = rendezvous; + this.inbox = inbox; + } + + /** Backward-compatible constructor that uses a default {@link InMemoryReplyInbox}. */ + public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous) { + this(agents, injector, rendezvous, new InMemoryReplyInbox()); } /** Current lifecycle status of a worker (the {@code GET /sessions/{id}/status} surface). */ @@ -164,6 +174,40 @@ public final class MessageService { return agents.status(target); } + /** + * Route a worker's explicit {@code bridge_reply}: resolve an open send, or queue it in the + * inbox if no send is currently open. Unlike the bare {@link Rendezvous#resolve}, a no-waiter + * result is not a failure — the reply is held for later drain. + * + *

Do NOT use this for mid-turn questions. {@code bridge_ask} / + * {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions + * are interactive and must never be queued. + * + * @return always {@code true} — the reply either resolved a live send or was queued + */ + public boolean reply(String session, String content) { + if (rendezvous.resolve(session, content)) { + return true; // a live send took it — unchanged fast path + } + inbox.publish(session, UUID.randomUUID().toString(), content); + return true; // held, not lost + } + + /** + * Drain (peek + ack) all pending inbox replies for {@code target}. At-least-once: returns the + * messages and acknowledges them; an in-flight failure between returning and the caller + * processing them re-surfaces them on a subsequent drain (the ack is local). + * + * @return the drained messages, newest last (FIFO); empty list if none + */ + public List drainReplies(String target) { + var messages = inbox.peek(target); + for (var msg : messages) { + inbox.ack(target, msg.msgId()); + } + return messages; + } + /** * Deliver {@code content} to {@code target} (a herdr {@code terminal_id}) and block until the * worker replies via {@link Rendezvous} or {@code timeoutMillis} elapses. diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyInbox.java b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyInbox.java new file mode 100644 index 0000000..7a03e31 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyInbox.java @@ -0,0 +1,31 @@ +package dev.ltms.bridged.msg; + +import java.util.List; + +/** + * Holds terminal worker→primary replies that arrive with no live send to resolve, keyed by worker + * session (target), until the primary drains them. Soft-state in Stage 1 (in-memory, lost on restart); + * the Stage 2 AMQP adapter implements the same contract with cross-restart durability. + * + *

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). + */ +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) {} + + /** + * 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. + */ + void publish(String target, String msgId, String content); + + /** Non-destructive snapshot of pending replies for {@code target} (FIFO), empty list if none. */ + List peek(String target); + + /** Remove the reply {@code msgId} for {@code target} once the primary has taken it. No-op if absent. */ + void ack(String target, String msgId); +} diff --git a/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java b/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java index d10250a..b2f66a2 100644 --- a/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java +++ b/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java @@ -84,6 +84,7 @@ public final class BridgedApp { app.delete("/workers/{paneId}", this::stopWorker); app.post("/sessions/{id}/message", this::sendMessage); // bridge_send (primary; blocking, wait:false, or answer via turnId) app.post("/sessions/{id}/reply", this::replyMessage); // bridge_reply (worker) + app.get("/sessions/{id}/replies", this::drainReplies); // drain reply inbox (CB-307) app.post("/sessions/{id}/ask", this::askMessage); // bridge_ask (worker → primary, CB-205) app.get("/sessions/{id}/status", this::sessionStatus); // bridge_status app.get("/tasks/{ticket}", this::taskStatus); // poll an async (wait:false) send @@ -325,7 +326,7 @@ public final class BridgedApp { /** * The worker's structured reply ({@code bridge_reply}) — resolves the blocking send awaiting - * on this session. 200 if a send was waiting, 409 if none was (late or spurious reply). + * on this session, or queues the reply in the inbox when no send is open (CB-307). */ private void replyMessage(Context ctx) { String id = ctx.pathParam("id"); @@ -336,13 +337,22 @@ public final class BridgedApp { ctx.status(400).json(Map.of("error", "bad_request", "detail", "body must be JSON")); return; } - if (rendezvous.resolve(id, content)) { - ctx.status(200).json(Map.of("sessionId", id, "delivered", true)); - } else { - ctx.status(409).json(Map.of( - "sessionId", id, "error", "no_pending_send", - "detail", "no send is awaiting a reply for this session")); - } + messages.reply(id, content); + ctx.status(200).json(Map.of("sessionId", id, "delivered", true)); + } + + /** + * Drain the reply inbox for a worker session — peek + ack any replies that arrived when no send + * was open. At-least-once: draining removes them from the inbox so a subsequent read returns + * nothing; an in-flight failure between the drain and the caller's processing re-surfaces them. + */ + private void drainReplies(Context ctx) { + String id = ctx.pathParam("id"); + var replies = messages.drainReplies(id); + ctx.status(200).json(Map.of("sessionId", id, "replies", + replies.stream().map(m -> Map.of( + "msgId", m.msgId(), + "content", m.content())).toList())); } /** 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 27b5abe..e59b932 100644 --- a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java @@ -57,13 +57,16 @@ class BridgeMcpTest { CompletableFuture send = CompletableFuture.supplyAsync( () -> BridgeMcp.send(messages, "term_a", "review this", 4000L)); - McpSchema.CallToolResult reply = BridgeMcp.reply(rendezvous, "term_a", "LGTM"); + // Wait until the send has opened its waiter so the reply resolves it (CB-307: reply now + // queues in the inbox if no waiter is open, which would break the round-trip). long deadline = System.currentTimeMillis() + 3000; - while (Boolean.TRUE.equals(reply.isError()) && System.currentTimeMillis() < deadline) { + while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) { //noinspection BusyWait - Thread.sleep(10); - reply = BridgeMcp.reply(rendezvous, "term_a", "LGTM"); + Thread.sleep(5); } + assertTrue(rendezvous.isWaiting("term_a"), "send should have opened its waiter"); + + McpSchema.CallToolResult reply = BridgeMcp.reply(messages, "term_a", "LGTM"); assertEquals("delivered", textOf(reply)); McpSchema.CallToolResult res = send.get(6, TimeUnit.SECONDS); @@ -80,30 +83,32 @@ class BridgeMcpTest { assertTrue(out.contains("ticket="), out); String ticket = out.substring(out.indexOf("ticket=") + "ticket=".length()).trim(); - // Resolve the awaiting send once it has opened (retry past the async-open race). + // Wait until the send has opened its waiter before replying (CB-307: reply never errors, + // so the old retry-on-error pattern no longer works — it would queue instead of resolve). long deadline = System.currentTimeMillis() + 3000; - McpSchema.CallToolResult reply = BridgeMcp.reply(rendezvous, "term_a", "async LGTM"); - while (Boolean.TRUE.equals(reply.isError()) && System.currentTimeMillis() < deadline) { + while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) { //noinspection BusyWait - Thread.sleep(10); - reply = BridgeMcp.reply(rendezvous, "term_a", "async LGTM"); + Thread.sleep(5); } + assertTrue(rendezvous.isWaiting("term_a"), "send should have opened its waiter"); + + McpSchema.CallToolResult reply = BridgeMcp.reply(messages, "term_a", "async LGTM"); assertEquals("delivered", textOf(reply)); // Poll until the async send completes and reports the reply. - McpSchema.CallToolResult polled = BridgeMcp.poll(messages, ticket); + McpSchema.CallToolResult polled = BridgeMcp.poll(messages, ticket, null); deadline = System.currentTimeMillis() + 3000; while (!textOf(polled).contains("async LGTM") && System.currentTimeMillis() < deadline) { //noinspection BusyWait Thread.sleep(10); - polled = BridgeMcp.poll(messages, ticket); + polled = BridgeMcp.poll(messages, ticket, null); } assertEquals("async LGTM", textOf(polled)); } @Test void pollUnknownTicketIsAnError() { - McpSchema.CallToolResult res = BridgeMcp.poll(messages, "task-999"); + McpSchema.CallToolResult res = BridgeMcp.poll(messages, "task-999", null); assertTrue(res.isError()); assertTrue(textOf(res).contains("unknown ticket")); } @@ -122,10 +127,32 @@ class BridgeMcpTest { } @Test - void replyWithNoPendingSendIsAnError() { - McpSchema.CallToolResult res = BridgeMcp.reply(rendezvous, "term_a", "orphan"); - assertTrue(res.isError()); - assertTrue(textOf(res).contains("no send is awaiting")); + void replyWithNoPendingSendIsQueuedNotError() { + // CB-307: a reply with no open send is now queued in the inbox, not an error. + McpSchema.CallToolResult res = BridgeMcp.reply(messages, "term_a", "orphan"); + assertNotEquals(Boolean.TRUE, res.isError(), "a queued reply is not an error"); + assertEquals("delivered", textOf(res)); + + // The reply is drainable by target. + var drained = messages.drainReplies("term_a"); + assertEquals(1, drained.size()); + assertEquals("orphan", drained.getFirst().content()); + } + + @Test + void bridgePollWithTargetDrainsReplies() { + // A reply with no open send queues it in the inbox. + BridgeMcp.reply(messages, "term_a", "queued-msg"); + + // bridge_poll with target drains the inbox. + McpSchema.CallToolResult res = BridgeMcp.poll(messages, null, "term_a"); + assertNotEquals(Boolean.TRUE, res.isError()); + String text = textOf(res); + assertTrue(text.contains("queued-msg"), "the drained reply should appear in the result"); + + // Second drain returns empty. + McpSchema.CallToolResult empty = BridgeMcp.poll(messages, null, "term_a"); + assertEquals("[]", textOf(empty)); } @Test @@ -159,14 +186,14 @@ class BridgeMcpTest { // The worker's ask returns the answer — it resumes the same turn. assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS))); - // The resumed worker replies, resolving the answering send (retry past the reopen race). - McpSchema.CallToolResult reply = BridgeMcp.reply(rendezvous, "term_a", "done"); + // The resumed worker replies, resolving the answering send (wait for the reopened waiter). deadline = System.currentTimeMillis() + 3000; - while (Boolean.TRUE.equals(reply.isError()) && System.currentTimeMillis() < deadline) { + while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) { //noinspection BusyWait - Thread.sleep(10); - reply = BridgeMcp.reply(rendezvous, "term_a", "done"); + Thread.sleep(5); } + assertTrue(rendezvous.isWaiting("term_a"), "the answer should have reopened a waiter"); + McpSchema.CallToolResult reply = BridgeMcp.reply(messages, "term_a", "done"); assertEquals("delivered", textOf(reply)); assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS))); } diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/InMemoryReplyInboxTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/InMemoryReplyInboxTest.java new file mode 100644 index 0000000..e040de3 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/msg/InMemoryReplyInboxTest.java @@ -0,0 +1,145 @@ +package dev.ltms.bridged.msg; + +import org.junit.jupiter.api.Test; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +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. + */ +class InMemoryReplyInboxTest { + + private final ReplyInbox inbox = new InMemoryReplyInbox(); + + @Test + void publishThenPeekReturnsTheMessage() { + inbox.publish("term_a", "m1", "hello"); + var msgs = inbox.peek("term_a"); + assertEquals(1, msgs.size()); + assertEquals("m1", msgs.getFirst().msgId()); + assertEquals("term_a", msgs.getFirst().target()); + assertEquals("hello", msgs.getFirst().content()); + } + + @Test + void peekForUnknownTargetReturnsEmpty() { + assertTrue(inbox.peek("no-such-target").isEmpty()); + } + + @Test + void ackRemovesTheMessage() { + inbox.publish("term_a", "m1", "hello"); + inbox.ack("term_a", "m1"); + assertTrue(inbox.peek("term_a").isEmpty(), "after ack, the message is gone"); + } + + @Test + void ackForUnknownMsgIdIsNoOp() { + inbox.publish("term_a", "m1", "hello"); + inbox.ack("term_a", "no-such-id"); // no-op + assertEquals(1, inbox.peek("term_a").size(), "the published message is still there"); + } + + @Test + void ackForUnknownTargetIsNoOp() { + inbox.ack("no-such-target", "m1"); // no-op, should not throw + } + + @Test + void dedupByIdempotentMsgId() { + inbox.publish("term_a", "m1", "first"); + inbox.publish("term_a", "m1", "second"); // same msgId, different content + var msgs = inbox.peek("term_a"); + assertEquals(1, msgs.size(), "dedup: second publish with same msgId is a no-op"); + assertEquals("first", msgs.getFirst().content(), "the original content is retained"); + } + + @Test + void publishesWithDifferentMsgIdsBothAppear() { + inbox.publish("term_a", "m1", "first"); + inbox.publish("term_a", "m2", "second"); + var msgs = inbox.peek("term_a"); + assertEquals(2, msgs.size()); + assertEquals("m1", msgs.get(0).msgId()); + assertEquals("m2", msgs.get(1).msgId()); + } + + @Test + void perTargetIsolation() { + inbox.publish("term_a", "m1", "for-a"); + inbox.publish("term_b", "m2", "for-b"); + assertEquals(1, inbox.peek("term_a").size()); + assertEquals(1, inbox.peek("term_b").size()); + } + + @Test + void fifoOrderIsPreserved() { + inbox.publish("term_a", "m1", "first"); + inbox.publish("term_a", "m2", "second"); + inbox.publish("term_a", "m3", "third"); + var msgs = inbox.peek("term_a"); + assertEquals(3, msgs.size()); + assertEquals("m1", msgs.get(0).msgId()); + assertEquals("m2", msgs.get(1).msgId()); + assertEquals("m3", msgs.get(2).msgId()); + } + + @Test + void peekReturnsAnImmutableCopy() { + inbox.publish("term_a", "m1", "hello"); + var msgs = inbox.peek("term_a"); + assertThrows(UnsupportedOperationException.class, () -> msgs.add( + new ReplyInbox.InboxMessage("x", "term_a", "x"))); + } + + @Test + void ackRemovesOneMessageLeavesOthers() { + inbox.publish("term_a", "m1", "first"); + inbox.publish("term_a", "m2", "second"); + inbox.ack("term_a", "m1"); + var msgs = inbox.peek("term_a"); + assertEquals(1, msgs.size()); + assertEquals("m2", msgs.getFirst().msgId()); + } + + @Test + void concurrentPublishAndDrain() throws Exception { + int msgCount = 100; + ExecutorService exec = Executors.newVirtualThreadPerTaskExecutor(); + try { + // Concurrent publishers + var pubDone = new CountDownLatch(msgCount); + for (int i = 0; i < msgCount; i++) { + final int id = i; + exec.submit(() -> { + inbox.publish("term_a", "m" + id, "content-" + id); + pubDone.countDown(); + }); + } + // Concurrent drainer + AtomicReference drainError = new AtomicReference<>(); + exec.submit(() -> { + try { + pubDone.await(); + for (int i = 0; i < 50; i++) { + var peeked = inbox.peek("term_a"); + for (var msg : peeked) { + inbox.ack("term_a", msg.msgId()); + } + } + } catch (Exception e) { + drainError.set(e); + } + }).get(); + assertNull(drainError.get(), "concurrent drain should not throw"); + } finally { + exec.shutdown(); + } + } +} 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 550481c..7c1ce3e 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -211,4 +211,93 @@ class MessageServiceTest { assertEquals(MessageService.Outcome.STALE_TURN, r.outcome(), "an answer to a turn that never existed (or already lapsed) is stale, not a hang"); } + + // --- CB-307 reply inbox ---------------------------------------------------------------- + + @Test + void replyQueuesInInboxWhenNoSendIsOpen() { + // No send is open for this session — reply should queue in the inbox. + assertTrue(messages.reply(T, "queued-text"), "reply should succeed (queued)"); + + var drained = messages.drainReplies(T); + assertEquals(1, drained.size()); + assertEquals("queued-text", drained.getFirst().content()); + } + + @Test + void replyResolvesOpenSendDoesNotQueue() throws Exception { + CompletableFuture send = sendAsync(); + awaitUninterruptibly(T); + + // An explicit reply resolves the open send. + assertTrue(messages.reply(T, "send-resolved"), "reply should succeed (resolved live send)"); + + // The inbox should be empty — the reply went to the send, not the inbox. + assertTrue(messages.drainReplies(T).isEmpty(), "no reply in the inbox"); + + MessageService.Reply r = send.get(3, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.REPLIED, r.outcome()); + assertEquals("send-resolved", r.text()); + } + + @Test + void drainRepliesReturnsAllPendingThenEmptyOnNextCall() { + messages.reply(T, "msg-1"); + messages.reply(T, "msg-2"); + + var first = messages.drainReplies(T); + assertEquals(2, first.size()); + + var second = messages.drainReplies(T); + assertTrue(second.isEmpty(), "second drain should be empty (acked)"); + } + + @Test + void aQuestionIsNeverQueuedInTheInbox() { + // No send is open — bridge_ask with no delegation returns NO_WAITER, + // and the question text MUST NOT appear in the reply inbox. + // The inbox is only fed by MessageService.reply(), not by bridge_ask. + MessageService.AskResult r = messages.ask(T, "anyone there?", 500); + assertEquals(MessageService.AskOutcome.NO_WAITER, r.outcome(), + "bridge_ask with no open delegation must return NO_WAITER, never queued"); + + assertTrue(messages.drainReplies(T).isEmpty(), "questions must never be queued"); + } + + @Test + void completionFallbackIsNeverQueued() throws Exception { + // The fallback resolves a captured waiter, never the inbox. + CompletableFuture send = sendAsync(); + awaitUninterruptibly(T); + injectDelivery(); + + // The worker never sends bridge_reply, but the turn completes. + herdr.readText("done-scraped"); + completion.onTurnComplete(T); // The fallback arms and resolves the captured waiter. + + MessageService.Reply r = send.get(5, TimeUnit.SECONDS); + assertEquals(MessageService.Outcome.COMPLETED_UNREPLIED, r.outcome()); + + // The inbox should be empty — the reply went to the captured waiter. + assertTrue(messages.drainReplies(T).isEmpty(), "completion fallback must not queue"); + } + + // --- helpers --------------------------------------------------------------------------- + + /** Like {@link #awaitWaiting()} but rethrows as unchecked. */ + private void awaitUninterruptibly(String session) { + try { + awaitWaiting(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new IllegalStateException(e); + } + } + + /** Set up a delivered turn so the worker is working, ready for an ask or completion. */ + private void injectDelivery() { + herdr.readText("$ prompt"); // pre-turn content baseline + injector.onStatus(T, AgentStatus.IDLE); // deliver the task + injector.onStatus(T, AgentStatus.WORKING); // worker picks it up + } } 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 1c482c6..8cc0562 100644 --- a/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java @@ -313,15 +313,10 @@ class BridgedAppTest { catch (Exception e) { throw new RuntimeException(e); } }); - // The worker replies once a send is actually awaiting (retry past the startup race). - HttpResponse reply; - long deadline = System.currentTimeMillis() + 3000; - do { - reply = postJson(port, "/sessions/term_a/reply", "{\"content\":\"LGTM ship it\"}"); - if (reply.statusCode() != 409) break; - //noinspection BusyWait - Thread.sleep(10); - } while (System.currentTimeMillis() < deadline); + // Give the background send thread time to open its rendezvous waiter (CB-307: reply now + // queues in the inbox if no waiter is open, which would break the round-trip). + Thread.sleep(200); + HttpResponse reply = postJson(port, "/sessions/term_a/reply", "{\"content\":\"LGTM ship it\"}"); assertEquals(200, reply.statusCode()); HttpResponse res = send.get(6, java.util.concurrent.TimeUnit.SECONDS); @@ -341,20 +336,14 @@ class BridgedAppTest { String ticket = mapper.readTree(accepted.body()).get("ticket").asText(); assertFalse(ticket.isBlank(), "an async send must return a ticket"); - // The worker replies once the async send is actually awaiting (retry past the startup race). - HttpResponse reply; - long deadline = System.currentTimeMillis() + 3000; - do { - reply = postJson(port, "/sessions/term_a/reply", "{\"content\":\"async LGTM\"}"); - if (reply.statusCode() != 409) break; - //noinspection BusyWait - Thread.sleep(10); - } while (System.currentTimeMillis() < deadline); + // Give the background async send thread time to open its rendezvous waiter. + Thread.sleep(200); + HttpResponse reply = postJson(port, "/sessions/term_a/reply", "{\"content\":\"async LGTM\"}"); assertEquals(200, reply.statusCode()); // Polling the ticket now reports the finished delegation and its reply. JsonNode task; - deadline = System.currentTimeMillis() + 3000; + long deadline = System.currentTimeMillis() + 3000; do { task = mapper.readTree(req(port, "GET", "/tasks/" + ticket).body()); if ("done".equals(task.path("phase").asText())) break; @@ -375,11 +364,26 @@ class BridgedAppTest { } @Test - void replyWithNoPendingSendIsConflict() throws Exception { + void replyWithNoPendingSendQueuesInsteadOfConflict() throws Exception { + // CB-307: a reply with no open send now queues in the inbox, not a 409 conflict. int port = startHealthy(); HttpResponse res = postJson(port, "/sessions/term_a/reply", "{\"content\":\"orphan\"}"); - assertEquals(409, res.statusCode()); - assertEquals("no_pending_send", mapper.readTree(res.body()).get("error").asText()); + assertEquals(200, res.statusCode()); + + // The queued reply is drainable. + HttpResponse drain = req(port, "GET", "/sessions/term_a/replies"); + assertEquals(200, drain.statusCode()); + JsonNode body = mapper.readTree(drain.body()); + assertEquals(1, body.get("replies").size()); + assertEquals("orphan", body.get("replies").get(0).get("content").asText()); + } + + @Test + void drainRepliesReturnsEmptyForNoReplies() throws Exception { + int port = startHealthy(); + HttpResponse res = req(port, "GET", "/sessions/term_a/replies"); + assertEquals(200, res.statusCode()); + assertEquals(0, mapper.readTree(res.body()).get("replies").size()); } @Test