Files
fleetd/docs/CB-307-Reliable-Delivery.md
T
Dai Ha da5a987df0 CB-307/CB-308 design notes: reliable-delivery Stage-1 delegation spec + multi-host federation proposal
- docs/CB-307-Reliable-Delivery.md: ReplyInbox port + in-memory adapter spec
  (Stage 1, no broker); publish at the Rendezvous no-waiter drop seam, drain by target.
- docs/CB-308-Multi-Host-Federation.md: per-host gateway + per-agent broker channels
  + federated roster proposal (gitea #6), wiki-ready with theme-safe mermaid.
2026-07-18 20:24:16 +02:00

9.5 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 → returns false (msg/Rendezvous.java:212-215).
  • The content string is never retained — it is dropped. The worker is told it failed: BridgeMcp.reply returns error("no send is awaiting a reply for this worker") (mcp/BridgeMcp.java:270-272); REST returns 409 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 msgId within a target (a LinkedHashMap<msgId, InboxMessage> per target, or a deque + a seen-set — your call; preserve insertion order).
  • peek returns an immutable copy; ack removes by msgId. Thread-safe (concurrent publish vs. drain).
  • This is soft-state, NOT persistence. Lost on a java -jar bounce — 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(...) onto messages.reply(...):
    • BridgeMcp.reply (mcp/BridgeMcp.java:262-273) — on success return a normal ack; remove the error("no send is awaiting a reply…") branch (that case is now a successful queue).
    • BridgedApp.replyMessage (rest/BridgedApp.java:330-346) — return 200 (queued) instead of 409 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): peek the inbox, ack each returned msgId, hand back the List<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_poll to accept an optional target (worker session) and, when present, return that worker's drained replies — alongside a matching REST route GET /sessions/{id}/replies. Do not change send/answer semantics (do not drain inside send — that conflates "deliver to worker" with "collect its mail"). Keep the existing ticket-based bridge_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)

  1. New ReplyInbox + InboxMessage + InMemoryReplyInbox in dev.ltms.bridged.msg.
  2. bridge_reply with no open send now succeeds and queues (no more error / 409); the reply is later retrievable and identical.
  3. 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).
  4. QUESTION path unchanged — bridge_ask with no open send still returns NO_WAITER (add/keep a test proving a question is never queued).
  5. Completion/failure fallbacks unchanged.
  6. Unit tests covering: InMemoryReplyInbox publish/peek/ack/dedup/FIFO/concurrency; MessageService.reply queues on no-waiter and resolves-not-queues when a send is open; drainReplies returns+acks; the QUESTION-not-queued guard.
  7. 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=$?" — never mvn … | 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_diagnostics results. 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.json is --skip-worktree in 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/rest wiring, tests, and if you add config wiring in Bridged.java). Nothing else.
  • Work only inside your assigned worktree on your feature branch. The primary fast-forwards main after re-gating — do not touch main.
  • 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.