CB-307 Stage 1: reply-inbox port + in-memory adapter — hold stranded worker replies instead of dropping them

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.
This commit is contained in:
Dai Ha
2026-07-18 21:12:51 +02:00
parent da5a987df0
commit ba6b4a5da9
10 changed files with 480 additions and 67 deletions
@@ -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.
@@ -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<String, Object> 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<String, Object> 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() {
@@ -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.
*
* <p><strong>This is soft-state, NOT persistence.</strong> 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<String, LinkedHashMap<String, InboxMessage>> 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<InboxMessage> 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);
}
}
}
}
@@ -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<String, ReentrantLock> sessionLocks = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, Task> 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 <em>not</em> a failure — the reply is held for later drain.
*
* <p><strong>Do NOT use this for mid-turn questions.</strong> {@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<ReplyInbox.InboxMessage> 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.
@@ -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.
*
* <p><strong>This interface is the port.</strong> {@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<InboxMessage> 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);
}
@@ -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()));
}
/**
@@ -57,13 +57,16 @@ class BridgeMcpTest {
CompletableFuture<McpSchema.CallToolResult> 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)));
}
@@ -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<Exception> 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();
}
}
}
@@ -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<MessageService.Reply> 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<MessageService.Reply> 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
}
}
@@ -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<String> 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<String> reply = postJson(port, "/sessions/term_a/reply", "{\"content\":\"LGTM ship it\"}");
assertEquals(200, reply.statusCode());
HttpResponse<String> 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<String> 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<String> 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<String> 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<String> 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<String> res = req(port, "GET", "/sessions/term_a/replies");
assertEquals(200, res.statusCode());
assertEquals(0, mapper.readTree(res.body()).get("replies").size());
}
@Test