202 lines
12 KiB
Markdown
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`).
|
|
|