diff --git a/10-Cross-Host-Messaging.md b/10-Cross-Host-Messaging.md index 2bf0266..02d3b68 100644 --- a/10-Cross-Host-Messaging.md +++ b/10-Cross-Host-Messaging.md @@ -6,6 +6,18 @@ > It answers three questions: **what are the cross-host communication use cases**, **how the AMQP > broker is laid out**, and **which exchange + queue belongs to each entity**. Where a piece exists > today it is marked **[as-built]**; the rest is **[proposed]** (§9 draws the line). +> +> **One piece of cross-host traffic already ships, and it is not in the design below.** +> Lead-to-lead coordination between daemons on different hosts works today, over a **shared AMQP +> vhost** configured by a `coordinator:` block. A lead sends to a peer with `fleet_send{coordId}`, +> and its own coord-id is reported by `fleet_list`. The code is `msg/LeadMailbox` and +> `msg/LeadCoordLoop`, wired in `Fleetd.java:405` and `Fleetd.java:546-551`; the config record is +> `FleetConfig.java:74-78` and `99-101`. +> +> That path deliberately does **not** use the exchange topology in §3. It is a direct durable +> mailbox per lead, not a federated fabric, and it carries **coordination only** — leads divide the +> map and share findings, they never assign each other work. Read §3 onward as the design for +> *agent* traffic across hosts, which is still proposed. The one rule that shapes everything: **the broker moves _messages and presence_, never keystrokes.** Delivery into a peer is always a *local* herdr injection (an agent) or a *local* MCP pull (a @@ -87,7 +99,7 @@ canonical header subset above** — signing headers-plus-body keeps JSON canonic security path entirely. Verification at the receiving gateway, in order: signature valid → the claimed `from` lives on the signing host per **signed** presence (§8 step 5) → **`to` matches the routing key** (else a captured message could be replayed into a different inbox). Any check fails → -`basicReject(requeue=false)` → `bridge.dlx`, loud log, never a wedge. Replay is bounded by +`basicReject(requeue=false)` → `fleet.dlx`, loud log, never a wedge. Replay is bounded by `expiresAt` plus a bounded seen-set — the at-most-once forward path (§6) has no broker redelivery for a replayed duplicate to hide behind. @@ -98,11 +110,11 @@ durability and fan-out semantics. | Exchange | Type | Durable | Purpose | Status | |---|---|---|---|---| -| `bridge.msg` | topic | yes | All entity→entity messages (U1–U3, U6, U7) **and broadcast (U8)**. Routing key = **recipient globalId**, or `broadcast.all` / `broadcast.` for U8. | [proposed]* | -| `bridge.roster` | topic | yes | Presence heartbeats (U5). Routing key = `roster..`. | [proposed] | -| `bridge.control` | direct | yes | Cross-host spawn/stop/lifecycle (U4). Routing key = **target host id**. | [proposed] | -| `bridge.dlx` | fanout | yes | Dead-letter sink for poison messages off any inbox. | [proposed] | -| `bridge.delay` | x-delayed-message | yes | Optional broker-driven remind/backoff — re-publishes into `bridge.msg` after a delay. Native on LavinMQ; a **plugin** on RabbitMQ — stays optional so "RabbitMQ by URI swap" holds. | [proposed] | +| `fleet.msg` | topic | yes | All entity→entity messages (U1–U3, U6, U7) **and broadcast (U8)**. Routing key = **recipient globalId**, or `broadcast.all` / `broadcast.` for U8. | [proposed]* | +| `fleet.roster` | topic | yes | Presence heartbeats (U5). Routing key = `roster..`. | [proposed] | +| `fleet.control` | direct | yes | Cross-host spawn/stop/lifecycle (U4). Routing key = **target host id**. | [proposed] | +| `fleet.dlx` | fanout | yes | Dead-letter sink for poison messages off any inbox. | [proposed] | +| `fleet.delay` | x-delayed-message | yes | Optional broker-driven remind/backoff — re-publishes into `fleet.msg` after a delay. Native on LavinMQ; a **plugin** on RabbitMQ — stays optional so "RabbitMQ by URI swap" holds. | [proposed] | *\* **[as-built, with a semantic migration]** CB-307 publishes to the **default exchange** (`""`) with routing key = queue name — but note what that queue *means* today: `agent..inbox` is @@ -120,10 +132,10 @@ as §2.1.* ```mermaid flowchart TB subgraph exch["Exchanges"] - msg["bridge.msg
(topic)"] - ros["bridge.roster
(topic)"] - ctl["bridge.control
(direct)"] - dlx["bridge.dlx
(fanout)"] + msg["fleet.msg
(topic)"] + ros["fleet.roster
(topic)"] + ctl["fleet.control
(direct)"] + dlx["fleet.dlx
(fanout)"] end subgraph gwA["gateway A (host A)"] inA["agent.<gidA>.inbox
(durable)"] @@ -135,7 +147,7 @@ flowchart TB ctlB["control.hostB
(durable)"] rosB["roster.hostB
(exclusive, transient)"] end - dlq["bridge.dlq
(durable)"] + dlq["fleet.dlq
(durable)"] msg -->|"key = gidA"| inA msg -->|"key = gidB"| inB ctl -->|"key = hostA"| ctlA @@ -152,7 +164,7 @@ flowchart TB ``` *Figure 2 — exchanges (top) route to per-entity and per-gateway queues (green). Inbox queues -dead-letter poison messages to `bridge.dlx` → `bridge.dlq` (red). Each gateway consumes only the +dead-letter poison messages to `fleet.dlx` → `fleet.dlq` (red). Each gateway consumes only the queues in its own box.* ## 4. Queues per entity @@ -162,12 +174,12 @@ co-located with that entity.** Gateways additionally own a control queue and a r | Owner (entity/scope) | Queue | Bound to (exchange · key) | Sole consumer | Durability | Notes | |---|---|---|---|---|---| -| Worker / sandboxed agent `gid` | `agent..inbox` | `bridge.msg` · `` **+ `broadcast.all` / `broadcast.` (U8)** | the agent's **local** gateway → **herdr inject** | durable | [proposed] — recipient-keyed `.v2` **migration** of CB-307's sender-keyed queue (§3 footnote); survives idle and daemon bounce; broadcast = extra bindings on the *same* queue | -| Primary / main `gid` | `agent..inbox` | `bridge.msg` · `` | its **local** gateway → **MCP pull** (resolve open send, else hold + nudge) | durable | identical pattern — a main is just an entity with an inbox (U6/U7); per-worker drain filters on envelope `from` | -| Orchestrator `gid` | `agent..inbox` | `bridge.msg` · `` | its **local** gateway → **MCP pull** | durable | top MCP client; same fabric | -| Gateway (host) | `control.` | `bridge.control` · `` | that gateway | durable | cross-host spawn/stop (U4) | -| Gateway (host) | `roster.` | `bridge.roster` · `roster.#` | that gateway | **transient** (exclusive, auto-delete) | presence is soft-state — union view, rebuilt from heartbeats | -| Fleet (shared) | `bridge.dlq` | `bridge.dlx` | ops / redelivery tooling | durable | poison messages after N redeliveries | +| Worker / sandboxed agent `gid` | `agent..inbox` | `fleet.msg` · `` **+ `broadcast.all` / `broadcast.` (U8)** | the agent's **local** gateway → **herdr inject** | durable | [proposed] — recipient-keyed `.v2` **migration** of CB-307's sender-keyed queue (§3 footnote); survives idle and daemon bounce; broadcast = extra bindings on the *same* queue | +| Primary / main `gid` | `agent..inbox` | `fleet.msg` · `` | its **local** gateway → **MCP pull** (resolve open send, else hold + nudge) | durable | identical pattern — a main is just an entity with an inbox (U6/U7); per-worker drain filters on envelope `from` | +| Orchestrator `gid` | `agent..inbox` | `fleet.msg` · `` | its **local** gateway → **MCP pull** | durable | top MCP client; same fabric | +| Gateway (host) | `control.` | `fleet.control` · `` | that gateway | durable | cross-host spawn/stop (U4) | +| Gateway (host) | `roster.` | `fleet.roster` · `roster.#` | that gateway | **transient** (exclusive, auto-delete) | presence is soft-state — union view, rebuilt from heartbeats | +| Fleet (shared) | `fleet.dlq` | `fleet.dlx` | ops / redelivery tooling | durable | poison messages after N redeliveries | Why the roster queue is **transient** while inboxes are **durable**: the broker owns *message* durability, **not** who/where/status. Presence is rebuilt from heartbeats on reconnect; a missed @@ -182,7 +194,7 @@ stays soft-state; only *messages* are durable. See [1. Architecture](1-Architect instead of silently splitting the stream (two consumers on one queue *round-robin* — each sees half). Fan-out to many recipients is always *copies into many inboxes* (U8), never two consumers on one queue. -2. **Publish by recipient, not by host.** A sender publishes to `bridge.msg` with key = recipient +2. **Publish by recipient, not by host.** A sender publishes to `fleet.msg` with key = recipient globalId; the broker routes to whichever gateway holds that inbox. Senders are oblivious to the recipient's host — the roster resolves *existence*, the broker resolves *location*. 3. **The final hop is a pull for MCP clients.** For a primary/main/orchestrator the broker makes the @@ -198,7 +210,7 @@ stays soft-state; only *messages* are durable. See [1. Architecture](1-Architect (loud, recoverable) instead of a silent duplicate injection (quiet, dangerous). The worker's gateway publishes an **`INJECTED`** confirmation once the inject lands; no `INJECTED` within a bound = fast explicit failure, not a send-timeout-long blind window. Poison messages are - explicitly rejected (`basicReject(requeue=false)`) → `bridge.dlx` — nothing is ever silently + explicitly rejected (`basicReject(requeue=false)`) → `fleet.dlx` — nothing is ever silently requeued forever. 6. **Signed messages.** Every cross-host message is signed by its publishing gateway; the receiver verifies the signature **and** that the claimed sender lives on the signing host (roster check). @@ -252,7 +264,7 @@ sequenceDiagram autonumber participant MA as main (host A) participant GA as gateway A - participant BR as bridge.msg + participant BR as fleet.msg participant GB as gateway B participant W as worker (host B) MA->>GA: fleet_send(gidW, content) @@ -279,7 +291,7 @@ sequenceDiagram autonumber participant MA as main (host A) participant GA as gateway A - participant CX as bridge.control + participant CX as fleet.control participant GB as gateway B participant SL as SandboxLauncher (B) MA->>GA: fleet_spawn(profile, host = B) @@ -288,7 +300,7 @@ sequenceDiagram GB->>SL: spawn locally (into a local sandbox) Note over SL: CB-306 readiness gate — wait until injectable SL-->>GB: PeerHandle(new gid) - GB->>GB: announce on bridge.roster (roster.hostB.newGid) + GB->>GB: announce on fleet.roster (roster.hostB.newGid) Note over MA: gid appears in every gateway's union roster,
then U1 addresses it normally ``` @@ -309,7 +321,7 @@ flowchart LR hbB["heartbeat local agents
roster.hostB.*"] viewB["union rosterView"] end - ros["bridge.roster (topic)"] + ros["fleet.roster (topic)"] hbA -->|"publish"| ros hbB -->|"publish"| ros ros -->|"roster.# → roster.hostA"| viewA @@ -326,7 +338,7 @@ federated). Transient queues + heartbeat TTL keep it soft-state.* Asks are the one flow that is deliberately **not** durable — the single-host rule kept: only terminal replies are queued, a live conversation is not. Cross-host, `ASK`/`ANSWER` traverse -`bridge.msg` carrying the asking turn's `turn_id` and an **`expiresAt`** that the *gateways* +`fleet.msg` carrying the asking turn's `turn_id` and an **`expiresAt`** that the *gateways* enforce at delivery time — a broker per-message TTL cannot expire a message already consumed into the held map, so broker TTL is only a backstop for a queue nobody is consuming. Two arrival checks preserve the single-host semantics: gateway A checks for an open waiter **on arrival** and @@ -339,7 +351,7 @@ sequenceDiagram autonumber participant W as worker (host B, parked mid-turn) participant GB as gateway B - participant BR as bridge.msg + participant BR as fleet.msg participant GA as gateway A participant MA as main (host A) W->>GB: fleet_ask(question) @@ -376,17 +388,17 @@ On boot, a gateway declares and wires exactly its own slice: 1. Connect to the broker URI over TLS (`amqps://…`), with this gateway's **own broker login**. LavinMQ default; RabbitMQ by URI swap. `[as-built for plain amqp://; TLS + per-gateway login proposed]` -2. Declare the shared exchanges `bridge.msg`, `bridge.roster`, `bridge.control`, `bridge.dlx` +2. Declare the shared exchanges `fleet.msg`, `fleet.roster`, `fleet.control`, `fleet.dlx` (idempotent). `[proposed]` -3. For each **local** entity: declare `agent..inbox.v2` (durable, DLX = `bridge.dlx`, +3. For each **local** entity: declare `agent..inbox.v2` (durable, DLX = `fleet.dlx`, **capped**: `x-max-length` + message TTL, `x-expires` as the orphan-GC backstop), bind to - `bridge.msg` keys `` **and** `broadcast.all` (+ its groups), start a **manual-ack, + `fleet.msg` keys `` **and** `broadcast.all` (+ its groups), start a **manual-ack, exclusive** consumer with a **`basicQos` prefetch bound** → inject/hold. The `.v2` suffix is the migration seam: AMQP refuses to redeclare an existing durable queue with new arguments (`PRECONDITION_FAILED`). `[proposed; the unsuffixed, uncapped queue is as-built]` -4. Declare `control.` (durable), bind to `bridge.control` key ``, consume → +4. Declare `control.` (durable), bind to `fleet.control` key ``, consume → run the launcher locally. `[proposed]` -5. Declare `roster.` (exclusive, auto-delete), bind to `bridge.roster` key `roster.#`, +5. Declare `roster.` (exclusive, auto-delete), bind to `fleet.roster` key `roster.#`, consume → maintain the union view; start heartbeating local agents **and a host-level entry** (`roster.._gateway`, carrying this gateway's profile list — an agentless host must still be visible as a spawn target). All presence messages are **signed**: invariant 6's roster @@ -408,22 +420,36 @@ On boot, a gateway declares and wires exactly its own slice: | Piece | Today `[as-built]` | Cross-host `[proposed]` | |---|---|---| -| Per-entity inbox | `agent..inbox` durable, manual-ack consume-and-hold, dedup by `msgId`, auto-recovery — on the **default exchange**, keyed by the *sender* of a primary-bound reply | **recipient-keyed** `agent..inbox.v2` on `bridge.msg` — a semantic **migration** (§3 footnote); per-worker drain preserved via envelope `from` | +| Per-entity inbox | `agent..inbox` durable, manual-ack consume-and-hold, dedup by `msgId`, auto-recovery — on the **default exchange**, keyed by the *sender* of a primary-bound reply | **recipient-keyed** `agent..inbox.v2` on `fleet.msg` — a semantic **migration** (§3 footnote); per-worker drain preserved via envelope `from` | | Consumer model | one daemon consumes all inboxes | **one gateway per host**, sole consumer of its **local** entities' inboxes | -| Presence | in-process roster (CB-304) | `bridge.roster` + `roster.` transient queues → federated union | -| Spawn | local call in-process | `bridge.control` + `control.` cross-host control message | -| Poison handling | none | `bridge.dlx` → `bridge.dlq` | -| Remind/backoff | in-JVM `ReplyPushLoop` (CB-307 Stage 3), keyed by `target` | rekeyed to envelope `from` once the inbox flips; optionally broker-driven via `bridge.delay` | +| Presence | in-process roster (CB-304) | `fleet.roster` + `roster.` transient queues → federated union | +| Spawn | local call in-process | `fleet.control` + `control.` cross-host control message | +| Poison handling | none | `fleet.dlx` → `fleet.dlq` | +| Remind/backoff | in-JVM `ReplyPushLoop` (CB-307 Stage 3), keyed by `target` | rekeyed to envelope `from` once the inbox flips; optionally broker-driven via `fleet.delay` | | Identity in key | `paneId` | `globalId` (UUID); host in roster metadata | | Broadcast | none (N separate sends) | `broadcast.*` bindings on every inbox (U8) — publish once, the broker copies | | Sender authenticity | connection identity (loopback) | per-gateway message signing + roster host check | -| Inbox bounds | unbounded | `x-max-length` + TTL → `bridge.dlq`, with a metric | +| Inbox bounds | unbounded | `x-max-length` + TTL → `fleet.dlq`, with a metric | | Consumer enforcement | convention | AMQP **exclusive** consumers — a second gateway fails loudly | | Multi-primary (U6/U7) | single-slot `PrimaryRegistry`, one nudge target | **CB-500's delta, not CB-308's** — N message-addressable MCP clients per gateway, per-client nudge targets | -**Net:** the *inbox* half is real and already multi-host-ready on a shared broker; cross-host adds a -**presence** plane, a **control** plane, **dead-lettering**, and a **global id** — no change to the -delivery asymmetry or the consume-and-hold contract. +**Lead-to-lead is the exception, and it already shipped.** The table above is about *agent* traffic. +Cross-host **lead-to-lead** coordination is real today and took a different route: a durable mailbox +per lead on a shared vhost, addressed by `coordId`, with no exchange topology, no presence plane and +no global id. It solved a smaller problem, so it needed a smaller mechanism. + +| Piece | Lead-to-lead today `[as-built]` | +|---|---| +| Address | `coordId`, from the `coordinator:` block; a lead's own is in `fleet_list` | +| Transport | one shared AMQP vhost, separate from the per-fleet vhost | +| Send | `fleet_send{coordId}` — mutually exclusive with `sessionId` and `turnId` | +| Receive | `msg/LeadCoordLoop`, status-gated like any other delivery | +| Carries | coordination only — never a task, never a brief | + +**Net:** the *inbox* half is real and already multi-host-ready on a shared broker, and lead-to-lead +runs across hosts today. What cross-host *agent* traffic still needs is a **presence** plane, a +**control** plane, **dead-lettering**, and a **global id** — with no change to the delivery asymmetry +or the consume-and-hold contract. ## 10. Hardening rules (design review 2026-08-10) `[proposed]` @@ -434,11 +460,11 @@ recorded in `docs/CB-308-Multi-Host-Federation.md` §7; the broker-level rules l 1. **Inbox caps — made real by prefetch.** Every inbox carries `x-max-length` and a message TTL, **and the consumer runs with a `basicQos` prefetch bound** — without it, consume-and-hold drains the queue into gateway heap and the caps guard an empty queue (an as-built gap, - ticketed). Overflow and expiry **dead-letter** to `bridge.dlq` (never vanish) and a metric + ticketed). Overflow and expiry **dead-letter** to `fleet.dlq` (never vanish) and a metric fires. A gateway **sweeper** acks held messages past their `expiresAt`, so one abandoned sender's backlog cannot wedge a shared inbox behind a full prefetch window (head-of-line). A **DLQ consumer** reports every dead-lettered forward brief back to its sender as `FAILED` — - the caps must not become a new silent loss channel — and `bridge.dlq` itself is capped + the caps must not become a new silent loss channel — and `fleet.dlq` itself is capped (anything with broker write access could otherwise fill it). 2. **Private broker, TLS.** Gateways connect over `amqps://` with per-gateway logins. Written rule: the broker's **disk holds readable task text** (the audit log deliberately does not) —