diff --git a/10-Cross-Host-Messaging.md b/10-Cross-Host-Messaging.md index c382841..b1c233e 100644 --- a/10-Cross-Host-Messaging.md +++ b/10-Cross-Host-Messaging.md @@ -27,10 +27,14 @@ Every cross-host interaction is one of these. Each is a message flow between two | U5 | **Roster / presence** — every gateway announces its local agents so all build one union "who/where/status" view | each gateway → all | `PRESENCE` | | U6 | **Main-pair** — the two mains (Opus + cloud, CB-500 B) message each other | A ↔ A' | `SEND`/`REPLY` | | U7 | **Orchestrator ↔ mains** — the orchestrator tier (CB-500 C) drives/monitors its main sessions | O ↔ main | `SEND`/`REPLY` | +| U8 | **Broadcast / group announce** — a main publishes **once** to all workers (or a named group); the broker copies the message into every bound inbox | A → all | `BROADCAST` | U1–U3 are the [CB-307 rendezvous](9-Implementation) semantics, now spanning hosts. U4–U5 are the net-new federation control plane. U6–U7 reuse the *same* per-entity inbox as U1 — a main and an -orchestrator are just message-addressable entities with their own inbox. +orchestrator are just message-addressable entities with their own inbox. U8 is for **identical +announcements** ("everyone: stop", "everyone: report status") — tailored task briefs keep the +per-worker fan-out pattern ([Team](6-Team#parallel-fan-out-map--reduce)), which works unchanged +across hosts because each send routes to its recipient's inbox wherever it lives. ## 2. Entities & identity @@ -66,7 +70,7 @@ durability and fan-out semantics. | Exchange | Type | Durable | Purpose | Status | |---|---|---|---|---| -| `bridge.msg` | topic | yes | All entity→entity messages (U1–U3, U6, U7). Routing key = **recipient globalId**. | [proposed]* | +| `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] | @@ -123,7 +127,7 @@ 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` · `` | the agent's **local** gateway → **herdr inject** | durable | [as-built] naming; the queue survives idle and daemon bounce | +| Worker / sandboxed agent `gid` | `agent..inbox` | `bridge.msg` · `` **+ `broadcast.all` / `broadcast.` (U8)** | the agent's **local** gateway → **herdr inject** | durable | [as-built] naming; the queue 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) | | 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) | @@ -137,9 +141,12 @@ stays soft-state; only *messages* are durable. See [1. Architecture](1-Architect ## 5. Invariants -1. **Single consumer per inbox.** Only the co-located gateway consumes `agent..inbox` → - per-recipient FIFO ordering, and no double-injection. A queue with two consumers on two hosts - would double-deliver. +1. **Single consumer per inbox — broker-enforced.** Only the co-located gateway consumes + `agent..inbox` → per-recipient FIFO ordering, and no double-injection. Consumers are opened + with AMQP's **exclusive** flag, so a misconfigured second gateway fails loudly at connect time + 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 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*. @@ -152,6 +159,11 @@ stays soft-state; only *messages* are durable. See [1. Architecture](1-Architect 5. **At-least-once + idempotent.** Dedup by `msgId`; a manual-ack consumer acks **only after** the entity has the message (ack = drained/injected), so a gateway bounce redelivers rather than loses. Repeated failures dead-letter to `bridge.dlq`. +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). + This extends the single-host rule — *identity comes from the connection, never an argument* — + across the broker: cross-host, identity comes from the key. Decision record: + `docs/CB-308-Multi-Host-Federation.md` §7. ## 6. Message lifecycle & ack semantics `[as-built for the inbox path]` @@ -257,16 +269,61 @@ flowchart LR everyone's; both converge on the same eventually-consistent union view (CB-304 `rosterView`, federated). Transient queues + heartbeat TTL keep it soft-state.* +### 7.4 U2 — cross-host ask (live-only, with expiry) + +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` as **short-lived** messages (per-message TTL) carrying the asking turn's `turn_id`: + +```mermaid +sequenceDiagram + autonumber + participant W as worker (host B, parked mid-turn) + participant GB as gateway B + participant BR as bridge.msg + participant GA as gateway A + participant MA as main (host A) + W->>GB: bridge_ask(question) + GB->>BR: publish ASK (key = gidMA, TTL, turn_id) + BR->>GA: route + GA->>MA: surfaces on the open send / drain (local pull) + MA->>GA: answer on turn_id + GA->>BR: publish ANSWER (key = gidW, TTL, turn_id) + BR->>GB: route + alt turn still waiting + GB->>W: inject answer — same turn resumes + else turn gone (timeout / completed / recycled) + GB--xW: NOT injected + GB->>BR: publish TOO_LATE notice (key = gidMA) + end +``` + +*Figure 7 — U2 across hosts. The one failure case — an answer outliving its turn — is loud, not +weird: the stale answer is dropped (never injected into an unrelated turn) and the main is told.* + +### 7.5 U8 — broadcast + +The main publishes **once** with key `broadcast.all` (or `broadcast.`); the broker copies +the message into every inbox bound to that key, and each copy then follows the normal local +delivery rules (Figure 3) at its own gateway. A worker that joins later starts receiving +broadcasts the moment its inbox binds the key. Replies to a broadcast are just N ordinary U1 +replies. Use it for identical announcements; per-worker briefs stay on the N-send fan-out +([Team](6-Team#parallel-fan-out-map--reduce)). + ## 8. Setup checklist (per gateway) On boot, a gateway declares and wires exactly its own slice: -1. Connect to the broker URI (`broker.uri`, e.g. `amqp://…@lavinmq-host:5672/`). LavinMQ default; - RabbitMQ by URI swap. `[as-built]` +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` (idempotent). `[proposed]` -3. For each **local** entity: declare `agent..inbox` (durable, DLX = `bridge.dlx`), bind to - `bridge.msg` key ``, start a **manual-ack** consumer → inject/hold. `[as-built for the queue]` +3. For each **local** entity: declare `agent..inbox` (durable, DLX = `bridge.dlx`, **capped**: + `x-max-length` + per-message TTL, so a neglected inbox dead-letters instead of growing forever), + bind to `bridge.msg` keys `` **and** `broadcast.all` (+ its groups), start a **manual-ack, + exclusive** consumer → inject/hold. `[as-built for the queue; caps, broadcast bindings, and the + exclusive flag proposed]` 4. Declare `control.` (durable), bind to `bridge.control` key ``, consume → run the launcher locally. `[proposed]` 5. Declare `roster.` (exclusive, auto-delete), bind to `bridge.roster` key `roster.#`, @@ -285,11 +342,39 @@ On boot, a gateway declares and wires exactly its own slice: | Poison handling | none | `bridge.dlx` → `bridge.dlq` | | Remind/backoff | in-JVM `ReplyPushLoop` (CB-307 Stage 3) | optionally broker-driven via `bridge.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 | +| Consumer enforcement | convention | AMQP **exclusive** consumers — a second gateway fails loudly | **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. +## 10. Hardening rules (design review 2026-08-10) `[proposed]` + +Cross-cutting rules settled in a design review of this chapter + CB-308. The sender-side +decisions (message signing, profile ownership, repo provisioning, ask expiry, spawn dedup) are +recorded in `docs/CB-308-Multi-Host-Federation.md` §7; the broker-level rules live here: + +1. **Inbox caps.** Every `agent..inbox` carries `x-max-length` and a per-message TTL; + overflow and expired messages **dead-letter** to `bridge.dlq` (never vanish) and a metric + fires. Bounded growth even when a main never drains. +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) — + so the broker runs on a machine inside the trust circle, never a shared or rented one. +3. **Schema version.** Every message carries a schema version. Within a major version, unknown + fields are ignored — rolling upgrades across mixed-version gateways just work. A newer-major + message is refused **to the DLQ with a loud log line**, never half-parsed. +4. **Trace id.** The first send of a flow mints a trace id; every subsequent hop (send, ask, + answer, reply, spawn) carries it, and every log/audit line prints it. Debugging a three-host + flow = search one id on each host. +5. **Broker outage.** Same-host traffic never touches the broker (the routing fork) and keeps + working with the broker down — a written promise. A **remote** send while the broker is + unreachable **fails fast** with a clear error; the gateway never buffers on the broker's + behalf (it stays soft-state, so a crash cannot lose messages it claimed to deliver). + Auto-reconnect re-declares the topology; remote hosts read as unknown in the roster meanwhile. + ## Related pages - [1. Architecture](1-Architecture) — the two modes, the invariants, the MCP asymmetry.