From bcf15a85c9912d4eb767aed35576ea0c5dc48d7f Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 19 Jul 2026 07:08:02 +0200 Subject: [PATCH] wiki: CB-307 reply-inbox (as-built) + CB-308 multi-host federation - 9-Implementation: msg package 2->4 classes (ReplyInbox port + InMemoryReplyInbox), reply now HELD not dropped when no send is open; drain via bridge_poll(target) + GET /sessions/{id}/replies; forward-rendezvous stranded-reply note. - 8-Roadmap: new 'Delivery reliability & multi-host' section (CB-306 shipped, CB-307 Stage 1 shipped / Stage 2 deferred, CB-308 designed) + staging diagram. - 1-Architecture: multi-host federation topology (CB-308) + diagram; late-reply durability note (CB-307) in Mode 1. --- 1-Architecture.md | 58 +++++++++++++++++++++++++++++++++++++++++++++ 8-Roadmap.md | 56 +++++++++++++++++++++++++++++++++++++++++++ 9-Implementation.md | 37 ++++++++++++++++++++++------- 3 files changed, 142 insertions(+), 9 deletions(-) diff --git a/1-Architecture.md b/1-Architecture.md index 6f11cbc..d228b8c 100644 --- a/1-Architecture.md +++ b/1-Architecture.md @@ -144,6 +144,14 @@ sequenceDiagram *observers* (a human, a dashboard) in parallel; the primary never has to hold it. For work that may outrun a sane request timeout, use Mode 2.* +> **Late replies are no longer lost (CB-307).** If the worker finishes *after* the primary's +> blocking `bridge_send` has already timed out (or was never opened), its `bridge_reply` finds no +> live waiter to resolve. `bridged` now **holds that reply in a per-worker inbox** rather than +> discarding it, and the primary collects it later keyed by target (`bridge_poll(target)` / +> `GET /sessions/{id}/replies`). Stage 1 is soft-state (in-memory, lost on a daemon restart); +> cross-restart durability arrives with the CB-307 Stage-2 broker adapter. See +> [Roadmap](8-Roadmap#delivery-reliability--multi-host-cb-306--cb-308). + ### Mode 2 — asynchronous delivery (bridged-mediated) For traffic with no caller waiting — a detached progress report, an out-of-band question, a @@ -225,6 +233,56 @@ detail is in [Message Server](2-Message-Server) → *Worker session lifecycle*.* `bridged`'s MCP/HTTP endpoint**. The primary isn't a herdr pane, so async wake-ups use the `Stop`-hook-polls-`bridged` path above. Deployment diagrams: [Message Server](2-Message-Server) → *Deployment model*. +- **Multi-host federation (planned — CB-308).** Not one `bridged` fronting remote workers, but + **many gateways, one per host**, meshed over a broker. Each host runs its own `bridged` owning + its local herdr + registry; agents get a **host-unique global id** and a per-agent broker inbox + `agent..inbox` whose *owning* gateway is the sole consumer; a **federated roster** + (soft-state presence on a `roster.*` topic) gives every gateway an eventually-consistent + who/where/status view. Sending is a one-line fork — **local → inject via herdr today; remote → + publish to the broker**, which routes to the owning gateway. Crucially **invariant #2 holds + per host**: a session still talks only to its *local* gateway, so the MCP pull-asymmetry (a + primary is pulled-from, never called-into) survives the network unchanged. This builds directly + on the CB-307 broker fabric; the gating concern is the cross-host **trust model** (the broker + link becomes the security boundary). Design: **`docs/CB-308-Multi-Host-Federation.md`**; + staging in [Roadmap → Delivery reliability & multi-host](8-Roadmap#delivery-reliability--multi-host-cb-306--cb-308). + +```mermaid +flowchart TB + subgraph hostA["Host A"] + gA["gateway = bridged A
local herdr + registry"] + P["primary (MCP client)"] + P --- gA + end + subgraph hostB["Host B"] + gB["gateway = bridged B
local herdr + registry"] + WB["worker panes"] + gB --- WB + end + subgraph BR["broker (LavinMQ / AMQP) — CB-307 fabric"] + IN["agent.<id>.inbox queues"] + RO["roster.* presence topic"] + end + gA -->|"remote send → publish"| IN + IN -->|"owning gateway consumes"| gB + gB -->|"reply → publish (durable, msg id)"| IN + IN -->|"held until primary pulls"| gA + gA -->|"announce local agents"| RO + gB -->|"announce local agents"| RO + RO -->|"union view"| gA + RO -->|"union view"| gB + + classDef core fill:#2f855a,stroke:#22543d,color:#ffffff; + classDef ext fill:#2b6cb0,stroke:#1a365d,color:#ffffff; + classDef warn fill:#b7791f,stroke:#7b341e,color:#ffffff; + class gA,gB core + class P,WB ext + class IN,RO warn +``` + +*Figure: multi-host is an addressing + routing concern, not a new kind of peer. Each gateway owns +its own herdr and consumes only its own agents' inboxes; the broker routes between them; the +federated roster is soft-state (bridged owns who/where/status, the broker owns message durability). +No host ever sees another host's terminals.* **Security:** `bridged` is an agent-control surface — `bridge_send` runs arbitrary prompts and key-passthrough sends raw keystrokes into a live agent. Bind it to `localhost` + SSH tunnel, or diff --git a/8-Roadmap.md b/8-Roadmap.md index 3dfd76b..f2a243f 100644 --- a/8-Roadmap.md +++ b/8-Roadmap.md @@ -133,6 +133,62 @@ Stage A has landed on main: - **Stage B / CB-402:** a second in-tree coding-agent adapter (e.g. Codex) to prove the SPI holds. - **Stage C:** dynamic external plugin loading (`ServiceLoader`/jar discovery). This is gated by a trust/capability model — a launcher runs at daemon privilege and can inject env/tokens into peers, so third-party plugins are not enabled without that model. +## Delivery reliability & multi-host (CB-306 – CB-308) + +A cross-cutting track that came out of a **"communication break" review** of the reverse +(worker → primary) path. The bridge *pushes* to workers (the status-gated injector) but only +*pulls* to the primary — a worker reply resolves only an **already-open** blocking `bridge_send`; +MCP is client-initiated (bridge = server, primary = client), so the server can't call into the +primary. That asymmetry is the root of all three tickets. + +```mermaid +flowchart LR + cb306["CB-306 ✅
spawn-readiness gate"] --> cb307["CB-307
reliable worker→primary delivery"] + cb307 --> s1["Stage 1 ✅
ReplyInbox port + in-memory adapter"] + cb307 --> s2["Stage 2 ⏳
AMQP / LavinMQ adapter"] + s2 --> cb308["CB-308 ⏳
multi-host federation"] + classDef done fill:#2f855a,stroke:#22543d,color:#ffffff; + classDef todo fill:#b7791f,stroke:#7b341e,color:#ffffff; + class cb306,s1 done + class s2,cb308 todo +``` + +*Figure: CB-307's broker fabric is the foundation CB-308 stretches across hosts; CB-306 makes a +spawn fail fast instead of stalling.* + +- **CB-306 — spawn-readiness gate ✅ (main `7dd6c46`).** `ClaudeCodeLauncher.spawn` blocks until the + new worker's Claude has connected the bridge MCP (`AgentStatus.injectable()`) or a + `spawn_ready_timeout_ms` (default 20 000) elapses → on timeout it self-reaps the pane and throws + `PeerUnreachableException` (MCP tool error / HTTP 502). Opt-out with `timeout == 0`. Turns the old + pre-REPL folder-trust stall (surfacing as a ~60s send timeout) into a fast, explicit failure. + +- **CB-307 — reliable worker → primary delivery.** Behind one `ReplyInbox` port so the adapter is + swappable (see [Implementation → `msg`](9-Implementation#msg--the-service-core-4-classes)). + - **Stage 1 ✅ (main `ba6b4a5`).** `ReplyInbox` port + `InMemoryReplyInbox` (soft-state, dedup by + `msgId`). A `bridge_reply` with no open send is now **held** instead of dropped; the primary + drains it by target via `bridge_poll(target)` / `GET /sessions/{id}/replies`. No broker, no new + dependency. **Only terminal replies are queued** — `bridge_ask` (interactive) and the + completion/failure fallbacks are deliberately *not* (would risk double-delivery). Live on the + running daemon. + - **Stage 2 ⏳ (deferred, gitea #5 open).** `AmqpReplyInbox` behind the *same* port for + cross-restart durability — **LavinMQ** default (single Crystal binary; `com.rabbitmq:amqp-client` + works unchanged; native DLX + delayed-message exchange give the remind/backoff loop for free), + RabbitMQ interchangeable by URI. `broker:` config absent → in-memory, present → AMQP. Needs a + live broker to integration-test. **`bridged` stays soft-state — the broker owns message + durability, not the bus.** + +- **CB-308 — multi-host federation ⏳ (design note, gitea #6; depends on CB-307 Stage 2).** A primary + on host A delegating to workers on hosts B, C… with no host learning another's terminals. Built on + CB-307's broker fabric: a **per-host gateway** (evolved `bridged` owning its local herdr + + registry), **per-agent broker channels** `agent..inbox` (owning gateway = sole consumer), + a **federated roster** (soft-state presence on a `roster.*` topic = [CB-304](9-Implementation)'s + `rosterView`, federated), and a one-line **routing fork** (`local ? inject : publish`). Five + net-new deltas: global agent id, federated directory, gateway routing, cross-host spawn, and the + cross-host **trust model** (the gating concern — the broker link becomes the security boundary). + The MCP pull-asymmetry survives the network: the final hop into the primary is still a pull from + its *local* gateway. Full design: **`docs/CB-308-Multi-Host-Federation.md`**; topology sketch in + [Architecture → Topologies](1-Architecture#topologies). + ## Stage 1 — detailed tickets **Goal:** Opus, from its own subscription session, mounts `bridged` and gets a code review diff --git a/9-Implementation.md b/9-Implementation.md index 87cb472..016e59d 100644 --- a/9-Implementation.md +++ b/9-Implementation.md @@ -3,7 +3,9 @@ > **Scope.** This is the *as-built* code map of the `bridged` module — the actual packages, > classes, flows, and state machines in the source tree, as a companion to the design-level > [1. Architecture](1-Architecture) and [2. Message Server](2-Message-Server). Every enum, -> constant, and route below was verified against source at main `3aa69a9`. +> constant, and route below was verified against source at main `3aa69a9`; the `msg`-layer +> **reply-inbox** (CB-307 Stage 1) and its drain surface landed at `ba6b4a5` and are folded in +> below. `bridged` is a single-host Java 25 / Maven daemon: the sole gateway between an on-subscription **primary** (Opus) and off-subscription **workers**, speaking to the **herdr** PTY manager @@ -27,7 +29,7 @@ flowchart TB rest["rest.BridgedApp
REST · Javalin"] mcp["mcp.BridgeMcp
MCP · /mcp servlet"] end - msg["msg.MessageService + Rendezvous
service core · per-session rendezvous"] + msg["msg.MessageService + Rendezvous + ReplyInbox
service core · rendezvous + held-reply inbox"] inject["inject.Injector + StatusPoller
single-writer, status-gated delivery"] guard["guard.SubscriptionGuard
boundary enforcement"] ccl["worker.ClaudeCodeLauncher
first Claude adapter · implements PeerLauncher"] @@ -104,15 +106,18 @@ calls `bridge_reply`. This is where the turn state machine lives. | `TurnListener` | interface | Turn-lifecycle callbacks: `onDelivered`, `onTurnComplete`, `onTurnFailed`. | | `WorkerPresence` | class | Tracks which workers have connected the bridge MCP (`markPresent`, `isPresent`, `forget`). | -### `msg` — the service core (2 classes) +### `msg` — the service core (4 classes) Owns the forward rendezvous (`bridge_send` → `bridge_reply`) and the reverse rendezvous -(`bridge_ask` → answer), plus async fire-and-poll. +(`bridge_ask` → answer), plus async fire-and-poll — and, since CB-307, the **reply inbox** that +holds a worker's terminal reply when *no* forward send is open, instead of dropping it. | Class | Kind | Role | |---|---|---| -| `MessageService` | class | Orchestrates send/reply/ask/answer + async dispatch (`send`, `answer`, `ask`, `sendAsync`, `poll`). | -| `Rendezvous` | class | Low-level registry of forward waiters + reverse-ask futures (`open`, `resolve`, `resolveQuestion`, `openAsk`, `answerAsk`, `closeAsk`, `resolveCompletion`, `resolveFailure`). | +| `MessageService` | class | Orchestrates send/reply/ask/answer + async dispatch (`send`, `answer`, `ask`, `sendAsync`, `poll`), and routes `bridge_reply` through **`reply`** (resolve an open send, else publish to the inbox) + **`drainReplies`** (peek-then-ack a target's held replies). | +| `Rendezvous` | class | Low-level registry of forward waiters + reverse-ask futures (`open`, `resolve`, `resolveQuestion`, `openAsk`, `answerAsk`, `closeAsk`, `resolveCompletion`, `resolveFailure`). **Untouched by CB-307** — it stays a pure synchronization primitive; a `false` from `resolve` (no live waiter) is what triggers the inbox publish, one layer up in `MessageService`. | +| `ReplyInbox` | interface | The port (CB-307): `publish(target, msgId, content)` (idempotent, dedup by `msgId`), `peek(target)` (non-destructive FIFO snapshot), `ack(target, msgId)`. Nested `InboxMessage(msgId, target, content)` record. The Stage-2 AMQP/LavinMQ adapter will implement the *same* port for cross-restart durability. | +| `InMemoryReplyInbox` | class | Stage-1 default adapter — per-target FIFO in a `ConcurrentHashMap>`, thread-safe, dedup by `msgId`. **Soft-state, not persistence** — undrained replies are lost on a `java -jar` bounce, consistent with "bridged stays soft-state; the broker owns durability." | **Enums (verbatim).** `Rendezvous.Kind`: `REPLY`, `COMPLETION`, `FAILED`, `QUESTION`. `MessageService.Outcome`: `REPLIED`, `COMPLETED_UNREPLIED`, `WORKER_FAILED`, `QUESTION`, @@ -147,10 +152,14 @@ Exposes `bridged` as a Streamable-HTTP MCP endpoint and resolves caller identity | `LsofProcessCwdLookup` | class | `lsof`-based cwd lookup for spawn inheritance. | **Tool surface:** `bridge_send` (delegate + block; with `turnId`, answers an ask) · `bridge_reply` -(worker → structured answer) · `bridge_ask` (worker pauses to ask primary) · `bridge_status` · -`bridge_poll` (async ticket) · `bridge_spawn` · `bridge_list` · `bridge_stop` · `bridge_profiles`. +(worker → structured answer; if no send is open it is now **held in the reply inbox**, not +errored) · `bridge_ask` (worker pauses to ask primary) · `bridge_status` · `bridge_poll` (async +ticket; with an optional `target`, **drains that worker's held replies** from the inbox) · +`bridge_spawn` · `bridge_list` · `bridge_stop` · `bridge_profiles`. **Invariant:** `bridge_reply`/`bridge_ask` accept *no* identity argument; it comes only from the -transport context. +transport context. **CB-307 scope:** only a *terminal* `bridge_reply` with no open send is queued — +`bridge_ask` (interactive; the worker blocks and can't consume a late answer) and the injector's +completion/failure fallbacks (they target a *captured* waiter, CB-116) are **never** queued. ### `rest` · `worker` · `guard` · `config` — the edge (6 classes) @@ -177,6 +186,7 @@ transport context. | `POST /sessions/{id}/reply` | `bridge_reply` (worker) | | `POST /sessions/{id}/ask` | `bridge_ask` (worker → primary) | | `GET /sessions/{id}/status` | Worker lifecycle + MCP readiness | +| `GET /sessions/{id}/replies` | Drain a worker's **held replies** from the inbox (CB-307): peek-then-ack, second call returns `[]` | | `GET /tasks/{ticket}` | Poll an async (`wait:false`) send | ## Bootstrap wiring @@ -234,6 +244,15 @@ sequenceDiagram *Figure 3 — forward send. If the timeout fires first, the outcome is `TIMED_OUT_WORKING` (delivered) or `TIMED_OUT_QUEUED` (never delivered).* +**The stranded-reply case (CB-307).** The figure is the happy path — a live waiter exists. The +failure that motivated CB-307 is the *reverse*: a worker finishes and calls `bridge_reply` **after** +its send has already timed out (the ~60s client window) or was never opened, so `Rendezvous.resolve` +finds no waiter. Before CB-307 that reply was silently discarded (the worker got an error / REST +`409`). Now `MessageService.reply` publishes it to the `ReplyInbox` under a fresh `msgId`, and the +primary collects it later keyed by target — `bridge_poll(target)` or `GET /sessions/{id}/replies` +(peek → deliver → ack, so an in-flight failure re-surfaces it). Soft-state only: a daemon bounce +clears any undrained replies; cross-restart durability is CB-307 Stage 2 (the AMQP adapter). + ## Flow: reverse rendezvous (`bridge_ask` → answer) A worker pauses its own turn to ask the primary a question; the question surfaces on the primary's