diff --git a/docs/CB-307-Reliable-Delivery.md b/docs/CB-307-Reliable-Delivery.md new file mode 100644 index 0000000..c4cbce6 --- /dev/null +++ b/docs/CB-307-Reliable-Delivery.md @@ -0,0 +1,164 @@ +# 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>` +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 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` keyed by target session; per-target FIFO ordering. +- Dedup by `msgId` within a target (a `LinkedHashMap` 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` (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. diff --git a/docs/CB-308-Multi-Host-Federation.md b/docs/CB-308-Multi-Host-Federation.md new file mode 100644 index 0000000..18e3eff --- /dev/null +++ b/docs/CB-308-Multi-Host-Federation.md @@ -0,0 +1,205 @@ +# CB-308 — Multi-Host Federation (Stage 5) + +**Status:** design note (proposal) +**Depends on:** CB-307 (broker-based reliable delivery) — CB-308 is the multi-host layer built *on* +CB-307's broker fabric. +**Relates to:** CB-401 (`PeerHandle` opaque id), CB-304 (`rosterView`), CB-306 (spawn-readiness), +CB-303 (lifecycle limits), CB-117 (orphan reap). + +## 1. Goal + +Let `claude-bridge` coordinate agents that live on **more than one host** — a primary on host A +delegating to workers on hosts B, C, … — without any host learning another host's terminals. The +bus stays a **provider-neutral communication fabric**; multi-host is an addressing + routing +concern, not a new kind of peer. + +The design rests on three pieces (the shape this ticket proposes): + +1. **Dedicated per-agent channels** — every agent has its own addressable inbox on the broker. +2. **A federated agent directory** — a global "who/where/status" lookup, assembled from per-host + presence, not a central database. +3. **A per-host gateway** — each host runs a `bridged` that owns its local herdr, registers/manages + its own sessions, and proxies messages to/from other hosts over the broker. + +## 2. What is single-host today (the assumptions to break) + +```mermaid +flowchart TB + subgraph host["Single host (today)"] + primary["primary
(MCP client)"] + daemon["bridged daemon
127.0.0.1:8765"] + reg["in-process registry
keyed by PeerHandle.id() == paneId"] + herdr["herdr
(local unix-socket PTY mux)"] + w1["worker pane wQ:p1"] + w2["worker pane wQ:p2"] + primary --> daemon + daemon --> reg + daemon --> herdr + herdr --> w1 + herdr --> w2 + end +``` + +*Figure 1 — everything is co-located and loopback.* + +Three concrete bake-ins assume one host: + +| Assumption | Where | Why it blocks multi-host | +|---|---|---| +| **herdr is local** | `herdr/` unix socket `~/.config/herdr/herdr.sock` | You cannot drive another host's PTYs → each host **must** own its herdr. This is why a per-host gateway is mandatory. | +| **registry is in-process, keyed by `paneId`** | `session/SessionManager` | `paneId` (e.g. `wQ:p2B`) is a herdr-local coordinate — meaningless off-host. Routing needs a host-unique id. | +| **loopback, no authn** | `rest/BridgedApp` binds `127.0.0.1:8765` | Fine on one host; the moment a second host can talk to a gateway, that link is a trust boundary. | + +## 3. Target architecture + +```mermaid +flowchart TB + subgraph hostA["HOST A"] + gA["gateway = bridged A"] + regA["local registry + herdr"] + primary["primary (MCP client)"] + gA --- regA + primary --- gA + end + subgraph hostB["HOST B"] + gB["gateway = bridged B"] + regB["local registry + herdr"] + wb["worker panes"] + gB --- regB + gB --- wb + end + subgraph broker["BROKER (LavinMQ / AMQP) — CB-307 fabric"] + inbox["agent.<id>.inbox queues"] + roster["roster.* presence topic"] + dlq["DLQ · delayed-retry (remind)"] + end + gA -->|"publish to agent.<id>.inbox"| inbox + gB -->|"publish to agent.<id>.inbox"| inbox + inbox -->|"owning gateway consumes"| gA + inbox -->|"owning gateway consumes"| gB + gA -->|"announce local agents"| roster + gB -->|"announce local agents"| roster + roster -->|"union view"| gA + roster -->|"union view"| gB +``` + +*Figure 2 — each gateway owns its local herdr + registry, consumes only its own agents' inboxes, +and announces its agents onto a shared presence topic. The broker routes; no host sees another +host's terminals.* + +### 3.1 Component mapping (the three pieces) + +- **Dedicated per-agent channels** = a per-agent AMQP routing key / queue, e.g. + `agent..inbox`. The agent's **owning gateway is the only consumer** of its inbox. + Senders publish to `agent..inbox` and never need to know the agent's host — the broker + routes to whichever gateway holds it. LavinMQ additionally gives durability, DLX, and a native + delayed-message exchange (the remind/backoff loop for free) — the same reasons CB-307 picked it. + +- **Federated agent directory** = a **soft-state, bridge-owned** roster, *not* a broker-stored + database. Per the persistence-boundary decision (bridged is soft-state; the broker owns *message* + durability, not *who/where/status*), each gateway announces its local agents `(globalId, host, + status, capabilities)` on a `roster.*` presence topic with periodic heartbeats. Every gateway + builds an eventually-consistent **union view** — literally CB-304's `rosterView`, federated. A + stale entry expires by missed heartbeat (reuses CB-303's idle/TTL thinking). + +- **Per-host gateway** = today's `bridged` daemon, evolved. It already registers/manages sessions + and controls its local herdr; multi-host adds exactly two responsibilities: (a) a broker client + that consumes its agents' inboxes and injects into local herdr, and (b) presence announce + + union-roster assembly. Evolution, not rewrite. + +### 3.2 Routing rule + +```mermaid +flowchart LR + send["bridge_send(globalId, msg)"] --> lookup{"directory:
is globalId local?"} + lookup -->|"yes"| local["inject via local herdr
(today's Injector path)"] + lookup -->|"no"| pub["publish agent.<id>.inbox
(broker routes to owning gateway)"] + pub --> consume["owning gateway consumes
→ injects into its local herdr"] +``` + +*Figure 3 — one fork: local agents keep today's in-process inject path; remote agents go over the +broker. A sender is oblivious to which branch it took.* + +## 4. What CB-307 already provides vs. what is net-new + +**CB-307 delivers the transport half** and is independently valuable on a single host: the AMQP +broker fabric, the `bridged → broker` client/adapter, at-least-once + idempotent (dedup-by-id) +delivery, DLQ, and delayed-retry (remind). That *is* the "proxy cross-host message" backbone; +extending the same broker from "worker→primary reliability" to "gateway↔gateway" is incremental. + +**Net-new for CB-308 (multi-host), five items:** + +1. **Global agent id** — decouple the routing key from `paneId`. CB-401's `PeerHandle` already + abstracts the routing id; make it host-unique (e.g. `/` or a UUID minted at spawn). + The registry and all verbs route on the global id. +2. **Federated directory** — presence announce + heartbeat + union roster over `roster.*` + (§3.1). +3. **Gateway routing** — the `local ? inject : publish` fork (§3.2), plus each gateway consuming + its own agents' inbox queues and injecting into local herdr. +4. **Cross-host spawn** — `spawn on host B` = publish a control request to B's control channel → + gateway B runs `ClaudeCodeLauncher.spawn` **locally** (CB-306's readiness gate becomes *more* + valuable here: the far side wants a positive "agent ready" before anyone sends) → announces the + new agent into the federated roster. +5. **Trust** — the broker connection is now the security boundary. A gateway injects env/tokens at + daemon privilege (the CB-401 Stage-C concern), so a **remote-triggered spawn/send** needs + authn/authz: who may act on which host, and which control channels a gateway will honour. + +## 5. The one thing the broker does NOT dissolve + +The MCP asymmetry survives the network. The primary is an MCP **client** to its **local** gateway; +it cannot be called into. A worker on B replying to a primary on A flows: + +```mermaid +sequenceDiagram + participant W as worker (host B) + participant GB as gateway B + participant BR as broker + participant GA as gateway A + participant P as primary (host A, MCP client) + W->>GB: bridge_reply + GB->>BR: publish primary-bound (durable, msg id) + BR->>GA: route to A's primary inbox + Note over GA: held durably until the primary pulls + P->>GA: blocking bridge_send resolves / bridge_poll + GA-->>P: reply (then ACK to broker) +``` + +*Figure 4 — the broker makes the middle hop lossless, ordered, and idempotent; the **final** hop +into the primary is still a **pull** (gateway A holds the message until the primary's blocking +`bridge_send` or `bridge_poll`). Cross-host neither improves nor worsens this — it just spans hosts. +This is precisely the gap CB-307 closes on one host and CB-308 stretches across hosts.* + +## 6. Staging & dependencies + +```mermaid +flowchart LR + cb307["CB-307
broker-based reliable delivery
(single host first)"] --> cb308["CB-308
multi-host federation
(this note)"] + cb308 --> a["global agent id"] + cb308 --> b["federated directory"] + cb308 --> c["gateway routing"] + cb308 --> d["cross-host spawn"] + cb308 --> e["cross-host trust model"] + classDef gate fill:#b7791f,stroke:#7b341e,color:#ffffff; + class e gate +``` + +*Figure 5 — CB-307 is the foundation; CB-308's five items build on it. The trust model (e) is the +gating concern before any host accepts remote control.* + +**Recommendation:** keep CB-307 scoped to single-host broker reliability (foundation, independently +useful), and build CB-308's items on top once the broker fabric exists. Choose CB-307's broker / +channel naming **multi-host-ready** now (per-agent routing keys, a `roster.*` topic namespace) so +CB-308 doesn't have to repaint the topology. + +## 7. Open questions + +- **Directory ground-truth:** pure soft-state presence (heartbeats) vs. also treating broker queue + existence as authoritative. Lean soft-state to preserve the persistence boundary; revisit if + split-brain roster views cause mis-routing. +- **Global id scheme:** `/` (human-legible, leaks host) vs. opaque UUID (clean, needs + the directory to resolve host). Probably UUID in the protocol, host as directory metadata. +- **Gateway discovery:** how gateways find the broker and each other (static config vs. discovery). +- **Trust model shape:** per-host shared secret vs. mTLS on the broker vs. a capability token per + control action — ties into CB-401 Stage-C. +- **Failure semantics:** a host/gateway dies mid-turn — how the federated roster reaps it (missed + heartbeat) and whether in-flight primary-bound messages survive (broker durability = yes).