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

202 lines
12 KiB
Markdown

# 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.
```java
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)`:
```java
/** 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 `REPLY`s 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:
```bash
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`).