Files
fleetd/docs/CB-307-Reliable-Delivery.md
T
Dai Ha ecc590f344 CB-632 unit 5: rename the daemon and classes in the docs prose
Part of #145 (CB-632). Documentation only, plus one internal literal.

Unit 1 renamed the package and classes, which left every doc describing
classes that no longer exist. This fixes the prose across README.md,
docs/ and bridged/docs/ -- 18 files.

Renamed: dev.ltms.bridged -> dev.ltms.fleet, the five class names, and
"bridged" where it names the daemon as a product rather than a path.

Also renamed two literals, because a doc that disagrees with the code is
worse than one that is out of date:

  - bridged-local-noauth -> fleetd-local-noauth. A placeholder apiKey
    OpenCodeLauncher sends when a profile resolves no token, to a local
    endpoint that does not check it. No test asserts the old string.
  - the vnd.ltms.bridged.* media type in the M4 design doc. It appears
    in no Java file, so nothing implements it yet.

Deliberately NOT renamed, because each is still literally true today and
changes only at the cutover:

  - paths: bridged/, bridged.yaml, bridged.example.yaml, bridged.jar,
    .bridged-worktrees, deploy/dev.ltms.bridged.plist,
    scripts/redeploy-bridged.sh, bridged-launchd-wrapper.sh
  - bridged_* metric names -- renaming these after the monitoring is
    wired would break dashboard continuity, so they move before it is
  - bridge_* MCP tool names, which answer alongside fleet_* on purpose
  - BRIDGED_* environment variables, read by a file outside this repo

Method note: perl, not sed. BSD sed has no \b and no lookaround, and a
word-boundary expression there fails silently. The prose replace uses
(?<![\w./-])bridged(?![\w./-]) so it cannot touch a path or an
identifier, then every remaining hit was read by hand.

Verified: mvn clean install green, 51 classes, 878 tests, 0 failures.
2026-08-23 06:46:34 +02:00

11 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 fleet_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: FleetMcp.reply returns error("no send is awaiting a reply for this worker") (mcp/FleetMcp.java:270-272); REST returns 409 no_pending_send (rest/FleetApp.java:339-345).

This is the observed "communication break": a worker that finishes just after its fleet_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.fleet.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.fleet.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 "fleetd 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 fleet_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(...):
    • FleetMcp.reply (mcp/FleetMcp.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).
    • FleetApp.replyMessage (rest/FleetApp.java:330-346) — return 200 (queued) instead of 409 no_pending_send.

DO NOT touch the QUESTION path. fleet_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 fleet_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 fleet_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 Fleetd.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.fleet.msg.
  2. fleet_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 — fleet_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 Fleetd.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 fleetd
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).