12 KiB
CB-307 Stage 1 — Reply-Inbox Port + In-Memory Adapter (delegation spec)
Ticket: gitea #5 (CB-307). Stage: 1 of 2 (see the issue's "Implementation staging" comment).
Scope of THIS delegation: the ReplyInbox port + the in-memory (soft-state) adapter, wired at the
exact drop seam so a worker's terminal reply is held instead of silently discarded when no primary
send is open. No broker, no new dependency, no infra — fully unit-testable and primary-gate-verifiable.
Stage 2 (the AMQP/LavinMQ adapter behind the same port) is explicitly out of scope here.
Read this whole spec before starting. The port and the publish seam are prescriptive (non-negotiable). Where a choice is genuinely open it says "DECISION" with the required default — follow the default and flag it in your completion message for the primary's review.
1. The bug this fixes (grounded in current code)
The reverse (worker→primary) path is Rendezvous — a ConcurrentHashMap<session, CompletableFuture<Resolution>>
of live blocking waiters only. No queue, no store. When a worker calls bridge_reply and no send
is currently open for that worker:
Rendezvous.resolve(session, content)→complete(...)→waiters.get(session) == null→ returnsfalse(msg/Rendezvous.java:212-215).- The
contentstring is never retained — it is dropped. The worker is told it failed:BridgeMcp.replyreturnserror("no send is awaiting a reply for this worker")(mcp/BridgeMcp.java:270-272); REST returns409 no_pending_send(rest/BridgedApp.java:339-345).
This is the observed "communication break": a worker that finishes just after its bridge_send timed
out (the ~60s sync window) replies into the void. There is no message-id, dedup, or ack anywhere in
the message path today.
2. What to build
2.1 The port — dev.ltms.bridged.msg.ReplyInbox
A thin interface owned by the msg layer. The in-memory adapter is Stage 1; the AMQP adapter (Stage 2)
implements the same interface, so keep it broker-agnostic.
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.
*/
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);
}
2.2 The default adapter — InMemoryReplyInbox
- Backed by a
ConcurrentHashMap<String, ...>keyed by target session; per-target FIFO ordering. - Dedup by
msgIdwithin a target (aLinkedHashMap<msgId, InboxMessage>per target, or a deque + a seen-set — your call; preserve insertion order). peekreturns an immutable copy;ackremoves bymsgId. Thread-safe (concurrent publish vs. drain).- This is soft-state, NOT persistence. Lost on a
java -jarbounce — that is correct and consistent with "bridged stays soft-state." Do not add any file/DB backing.
2.3 Publish seam — route reply through the service layer
Keep Rendezvous a pure synchronization primitive (do not give it an inbox field). Instead centralize
in MessageService, which already owns the Rendezvous and will own the ReplyInbox:
- Add
MessageService.reply(String session, String content):/** Route a worker's explicit bridge_reply: resolve an open send, or queue it in the inbox if none. */ 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); // was a silent drop return true; // held, not lost } - Repoint the two callers off the bare
rendezvous.resolve(...)ontomessages.reply(...):BridgeMcp.reply(mcp/BridgeMcp.java:262-273) — on success return a normal ack; remove theerror("no send is awaiting a reply…")branch (that case is now a successful queue).BridgedApp.replyMessage(rest/BridgedApp.java:330-346) — return200(queued) instead of409 no_pending_send.
DO NOT touch the QUESTION path. bridge_ask / rendezvous.resolveQuestion must keep today's
NO_WAITER behaviour — a mid-turn question is interactive (the worker blocks synchronously and cannot
consume a late answer), so it must never be queued. Only terminal REPLYs go to the inbox.
DO NOT queue the completion/failure fallbacks (resolveCompletion / resolveFailure,
Rendezvous.java:196-210). They target a captured waiter (CB-116); a false there means the turn was
already resolved or the scrape is a stale late duplicate — queueing it risks double-delivery. Leave them
exactly as they are. (Extending durability to completions is a deliberate Stage-2 consideration, not this.)
2.4 Drain seam — how the primary collects a stranded reply
The primary re-checks a worker it delegated to. Expose a drain keyed by worker session (target):
- Add
MessageService.drainReplies(String target):peekthe inbox,ackeach returnedmsgId, hand back theList<InboxMessage>(or just the contents). At-least-once: peek→deliver→ack (ack only after the caller has them, so an in-flight failure re-surfaces them). - DECISION (required default): expose via the existing poll verb, keyed by target. Extend
bridge_pollto accept an optionaltarget(worker session) and, when present, return that worker's drained replies — alongside a matching REST routeGET /sessions/{id}/replies. Do not changesend/answersemantics (do not drain insidesend— that conflates "deliver to worker" with "collect its mail"). Keep the existing ticket-basedbridge_poll(ticket)path working unchanged. If you see a cleaner surface, still ship this default and note the alternative for review.
3. Config
None for Stage 1. The in-memory adapter is the unconditional default — wire new InMemoryReplyInbox()
into MessageService in Bridged.main. Do not add a broker: config block (that arrives with the
Stage-2 AMQP adapter: absent → in-memory, present → AMQP).
4. Acceptance criteria (what the primary will verify)
- New
ReplyInbox+InboxMessage+InMemoryReplyInboxindev.ltms.bridged.msg. bridge_replywith no open send now succeeds and queues (no moreerror/409); the reply is later retrievable and identical.- The queued reply is drainable by the primary keyed by target; draining acks it (a second drain
returns nothing); dedup by
msgId(re-publishing the same id does not double-queue). - QUESTION path unchanged —
bridge_askwith no open send still returnsNO_WAITER(add/keep a test proving a question is never queued). - Completion/failure fallbacks unchanged.
- Unit tests covering:
InMemoryReplyInboxpublish/peek/ack/dedup/FIFO/concurrency;MessageService.replyqueues on no-waiter and resolves-not-queues when a send is open;drainRepliesreturns+acks; the QUESTION-not-queued guard. - All pre-existing tests still green (baseline is 188; your total must be ≥ 188 + your new tests).
5. Build & verification (worker side)
- Build with Maven from the worktree's
bridged/dir. Capture the exit code without a masking pipe (mvn clean install; echo "MVN_EXIT=$?"— nevermvn … | tail, which hides failures). - Read the real test totals from
target/surefire-reports/TEST-*.xml, not from stdout scroll. - You do not have IDE MCP access — do not claim
ide_diagnosticsresults. The primary runs the authoritative gate (IDE sync + diagnostics 0/0 +mvn clean install) before integrating. Your self-reported counts are inputs to that gate, not final facts.
6. Hard constraints (non-negotiable)
.mcp.jsonis--skip-worktreein your worktree — never edit,git add, or commit it.wiki/is a submodule — never run git in it; never touch it.- Commit only your feature changes (the new port/adapter, the
msg/mcp/restwiring, tests, and if you add config wiring inBridged.java). Nothing else. - Work only inside your assigned worktree on your feature branch. The primary fast-forwards
mainafter re-gating — do not touchmain. - Java 25 idioms are welcome (unnamed
_params, records). Keep the diff minimal and match surrounding style.
7. Definition of done (report back over the bridge)
Commit on your feature branch and reply with: the commit SHA, the surefire total (run/failures/errors), a
one-line note on the drain-surface decision (§2.4) you shipped, and confirmation that .mcp.json/wiki/
were untouched. The primary re-gates and integrates.
8. Running the contract tests (CB-521)
AmqpReplyInboxContractTest is the real-broker proof of the ReplyInbox port (eventual visibility, ack
removal, msgId dedup, cross-restart redelivery). It is @Tag("contract"), so the default
mvn test / mvn clean install skip it — that hermetic, Docker-free default is deliberate and
untouched. Run it explicitly when Docker (or a broker) is available:
cd bridged
mvn -Pcontract test -Dtest=AmqpReplyInboxContractTest # local: spins a RabbitMQ Testcontainers fixture
Two broker modes
| Mode | Trigger | Broker | Needs Docker? |
|---|---|---|---|
| Local | AMQP_URI unset |
Testcontainers starts rabbitmq:3.13-management |
Yes |
| CI / external | AMQP_URI set |
the broker at that URI (CI RabbitMQ service container) | No — binds straight to the URI, never touches Testcontainers |
In CI the broker is provided as a RabbitMQ service container and AMQP_URI points at it, so the
contract job runs the same assertions with no Docker on the runner and no skipped test
(see .gitea/workflows/ci.yml → contract). The build job stays hermetic and Docker-free — keep
that separation.
Docker-engine discovery (why the contract profile pins api.version)
Out of the box, Testcontainers 1.20.4's docker-java client defaults to Docker API 1.32 when no
version is requested. Modern engines reject that as too old — on this host's OrbStack (min API 1.40)
testcontainers fails with "Could not find a valid Docker environment … client version 1.32 is too
old" even though the docker CLI works (the CLI negotiates a newer API).
The contract Maven profile sets api.version=1.43 in surefire, which works on OrbStack and Docker
24+, and is overridable per host: mvn -Pcontract -Dapi.version=1.54 test …. It only applies under
-Pcontract, so the default build is unaffected. If your engine differs, set -Dapi.version to a
version ≥ your engine's minimum API (e.g. docker version shows API version).