Files
fleetd/docs/CB-307-Reliable-Delivery.md

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 → 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.

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).