CB-307/CB-308 design notes: reliable-delivery Stage-1 delegation spec + multi-host federation proposal
- docs/CB-307-Reliable-Delivery.md: ReplyInbox port + in-memory adapter spec (Stage 1, no broker); publish at the Rendezvous no-waiter drop seam, drain by target. - docs/CB-308-Multi-Host-Federation.md: per-host gateway + per-agent broker channels + federated roster proposal (gitea #6), wiki-ready with theme-safe mermaid.
This commit is contained in:
@@ -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<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.
|
||||
@@ -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<br/>(MCP client)"]
|
||||
daemon["bridged daemon<br/>127.0.0.1:8765"]
|
||||
reg["in-process registry<br/>keyed by PeerHandle.id() == paneId"]
|
||||
herdr["herdr<br/>(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.<globalId>.inbox`. The agent's **owning gateway is the only consumer** of its inbox.
|
||||
Senders publish to `agent.<id>.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:<br/>is globalId local?"}
|
||||
lookup -->|"yes"| local["inject via local herdr<br/>(today's Injector path)"]
|
||||
lookup -->|"no"| pub["publish agent.<id>.inbox<br/>(broker routes to owning gateway)"]
|
||||
pub --> consume["owning gateway consumes<br/>→ 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. `<host>/<paneId>` 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<br/>broker-based reliable delivery<br/>(single host first)"] --> cb308["CB-308<br/>multi-host federation<br/>(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:** `<host>/<paneId>` (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).
|
||||
Reference in New Issue
Block a user