From b646be108de1e100ab0c4d37763a1429eff49bcc Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Tue, 28 Jul 2026 16:33:44 +0200 Subject: [PATCH] =?UTF-8?q?wiki:=20add=20chapter=2010=20=E2=80=94=20Cross-?= =?UTF-8?q?Host=20Messaging=20&=20Broker=20Topology?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit New design chapter answering the cross-host broker questions: the 7 cross-host communication use cases (delegate+reply, mid-turn ask, stranded reply, cross-host spawn, presence, main-pair, orchestrator↔mains), the AMQP exchange/queue layout, and which exchange+queue belongs to each entity. Topology: 5 exchanges (bridge.msg topic, bridge.roster topic, bridge.control direct, bridge.dlx fanout, bridge.delay) → per-entity inbox agent..inbox (durable, consumed ONLY by the co-located gateway → herdr-inject for agents / MCP- pull for mains), per-gateway control. (durable) + roster. (transient), shared bridge.dlq. Invariants: single-consumer-per-inbox, publish-by-recipient, pull-final-hop for MCP clients, keystrokes-never-on-broker, soft-state roster (persistence boundary). Honest as-built vs proposed split: CB-307's agent..inbox consume-and-hold is real (default exchange, already cross- host-routable); net-new = globalId key + roster/control/DLX planes. 6 mmdc-validated theme-safe diagrams (entity model, topology, ack lifecycle, U1 delegate+reply, U4 cross-host spawn, U5 presence fanout). Sidebar updated. Anchored to msg.AmqpReplyInbox as-built; tracks CB-308 + CB-500 for proposed. --- 10-Cross-Host-Messaging.md | 302 +++++++++++++++++++++++++++++++++++++ _Sidebar.md | 1 + 2 files changed, 303 insertions(+) create mode 100644 10-Cross-Host-Messaging.md diff --git a/10-Cross-Host-Messaging.md b/10-Cross-Host-Messaging.md new file mode 100644 index 0000000..c382841 --- /dev/null +++ b/10-Cross-Host-Messaging.md @@ -0,0 +1,302 @@ +# 10. Cross-Host Messaging & Broker Topology + +> **Scope.** This chapter is the **broker design** for cross-host operation — the [CB-308 +> federation](8-Roadmap) and [CB-500 multi-tier](8-Roadmap) fabric — built on the **as-built** +> single-host reply inbox shipped by CB-307 (`msg.AmqpReplyInbox`, verified at `9-Implementation`). +> 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). + +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 +primary/main) done by the **gateway co-located with that peer**. The broker only carries the +*middle* hop between gateways. This is the same asymmetry CB-307 closes on one host, stretched across +hosts by CB-308 — see [1. Architecture](1-Architecture) and CB-500 §11. + +## 1. The cross-host communication use cases + +Every cross-host interaction is one of these. Each is a message flow between two **entities** +(a primary/main, an orchestrator, or a worker/sandboxed agent) that may live on different hosts. + +| # | Use case | Flow | Kind | +|---|---|---|---| +| U1 | **Delegate + reply** — a main on host A tasks a worker/sandboxed agent on host B, awaits its answer | A → B, then B → A | `SEND` → `REPLY` | +| U2 | **Mid-turn ask** — a worker on B pauses its turn to ask the main on A, resumes on the answer | B → A, then A → B (turn-scoped) | `ASK` → `ANSWER` | +| U3 | **Stranded / late reply** — B replies after A's blocking send timed out (or was never open); held durably until A pulls | B → A (held) | `REPLY` (deferred) | +| U4 | **Cross-host spawn** — A requests a new agent *on host B*; gateway B launches it locally and announces it | A → gateway B (control) | `SPAWN` | +| 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` | + +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. + +## 2. Entities & identity + +```mermaid +flowchart LR + subgraph e["Message-addressable entities (each owns ONE inbox)"] + orch["orchestrator
(MCP client)"] + main["primary / main
(MCP client)"] + wrk["worker / sandboxed agent
(herdr pane)"] + end + gw["gateway = bridged
(one per host)"] + orch -->|"delivered by local MCP pull"| gw + main -->|"delivered by local MCP pull"| gw + wrk -->|"delivered by local herdr inject"| gw + gw -->|"consumes its local entities' inboxes,
publishes to remote inboxes"| broker["BROKER (AMQP)"] + classDef bus fill:#2b6cb0,stroke:#1a365d,color:#ffffff; + class broker bus +``` + +*Figure 1 — the gateway is the only AMQP client. Entities never touch the broker: an agent is +**injected** into over local herdr, a main/orchestrator **pulls** over local MCP. The gateway +consumes the inboxes of its **co-located** entities and publishes to the inboxes of **remote** ones.* + +**Global id.** The routing key is the entity's **global id** — a UUID minted at spawn (CB-308), +*not* the herdr `paneId` (which is host-local and meaningless off-host). Host is directory metadata, +carried in the roster, never in the routing key. `[as-built: keyed by paneId, single host]` +→ `[proposed: keyed by globalId]`. + +## 3. Broker topology — exchanges + +Five exchanges. Point-to-point messaging, presence, and control are separated so each has its own +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.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. | [proposed] | + +*\* **[as-built]** CB-307 publishes to the **default exchange** (`""`) with routing key = queue name, +which *already* routes cross-host on a shared broker (any gateway that declared `agent..inbox` +receives it). `bridge.msg` is an **observability/routing upgrade**, not a correctness requirement: a +named topic exchange lets an audit/monitor consumer bind `#` to watch all traffic, and leaves room +for future kind-scoped keys. The queue names and routing below are unchanged from CB-307.* + +```mermaid +flowchart TB + subgraph exch["Exchanges"] + msg["bridge.msg
(topic)"] + ros["bridge.roster
(topic)"] + ctl["bridge.control
(direct)"] + dlx["bridge.dlx
(fanout)"] + end + subgraph gwA["gateway A (host A)"] + inA["agent.<gidA>.inbox
(durable)"] + ctlA["control.hostA
(durable)"] + rosA["roster.hostA
(exclusive, transient)"] + end + subgraph gwB["gateway B (host B)"] + inB["agent.<gidB>.inbox
(durable)"] + ctlB["control.hostB
(durable)"] + rosB["roster.hostB
(exclusive, transient)"] + end + dlq["bridge.dlq
(durable)"] + msg -->|"key = gidA"| inA + msg -->|"key = gidB"| inB + ctl -->|"key = hostA"| ctlA + ctl -->|"key = hostB"| ctlB + ros -->|"key roster.#"| rosA + ros -->|"key roster.#"| rosB + inA -.->|"poison"| dlx + inB -.->|"poison"| dlx + dlx --> dlq + classDef q fill:#2f855a,stroke:#22543d,color:#ffffff; + classDef dead fill:#9b2c2c,stroke:#742a2a,color:#ffffff; + class inA,inB,ctlA,ctlB,rosA,rosB q + class dlq,dlx dead +``` + +*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 +queues in its own box.* + +## 4. Queues per entity + +The core rule: **one inbox per message-addressable entity, consumed by exactly one gateway — the one +co-located with that entity.** Gateways additionally own a control queue and a roster queue. + +| 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 | +| 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) | +| 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 | + +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 +heartbeat expires an entry (CB-303 TTL thinking). This preserves the persistence boundary — bridged +stays soft-state; only *messages* are durable. See [1. Architecture](1-Architecture). + +## 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. +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*. +3. **The final hop is a pull for MCP clients.** For a primary/main/orchestrator the broker makes the + *middle* hop lossless + ordered + idempotent, but the *last* hop is still `bridge_poll` / drain — + the gateway **holds** the reply until the client pulls (consume-and-hold). The broker does not + dissolve the MCP asymmetry (§ CB-308 "the one thing the broker does NOT dissolve"). +4. **Keystrokes never traverse the broker.** Only messages + presence. Injection into an agent is a + *local* herdr write by its owning gateway (CB-500 §11). +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. Message lifecycle & ack semantics `[as-built for the inbox path]` + +```mermaid +sequenceDiagram + autonumber + participant SND as sender gateway + participant MX as bridge.msg + participant Q as agent.GID.inbox + participant OWN as owning gateway + participant E as entity (agent/main) + SND->>MX: publish(key = GID, msgId, content) + MX->>Q: route (durable, persistent) + OWN->>Q: manual-ack consumer pulls (dedup by msgId) + Note over OWN: held in memory, NOT yet acked + OWN->>E: deliver locally (herdr inject OR hold for MCP pull) + E-->>OWN: drained / injected (confirmed) + OWN->>Q: basicAck(msgId) + Note over OWN,Q: a bounce BEFORE ack → broker redelivers
(cross-restart durability, CB-307 Stage 2) +``` + +*Figure 3 — consume-and-hold with deferred manual ack. The message stays on the broker until the +entity has actually received it, so an in-flight failure re-surfaces it. Dedup by `msgId` makes +redelivery safe. This is exactly `AmqpReplyInbox` today (`9-Implementation`), generalized from +"primary-bound replies" to "any entity inbox."* + +## 7. Use-case walkthroughs + +### 7.1 U1 — cross-host delegate + reply + +```mermaid +sequenceDiagram + autonumber + participant MA as main (host A) + participant GA as gateway A + participant BR as bridge.msg + participant GB as gateway B + participant W as worker (host B) + MA->>GA: bridge_send(gidW, content) + GA->>GA: roster - gidW local? NO + GA->>BR: publish(key = gidW) + BR->>GB: route to agent.gidW.inbox + GB->>W: inject via B's LOCAL herdr + W-->>GB: bridge_reply(to gidMA) + GB->>BR: publish(key = gidMA, durable) + BR->>GA: route to agent.gidMA.inbox + Note over GA: held until MA pulls (MA is an MCP client) + MA->>GA: blocking send resolves / bridge_poll + GA-->>MA: reply +``` + +*Figure 4 — U1. Both injection points (into W on B, and the drain into MA on A) are local; only the +two middle hops cross the broker. U3 (stranded reply) is the same picture where step 8's "held" +outlives A's send window and MA collects it later by `bridge_poll(gidW)` / drain.* + +### 7.2 U4 — cross-host spawn (the control plane) + +```mermaid +sequenceDiagram + autonumber + participant MA as main (host A) + participant GA as gateway A + participant CX as bridge.control + participant GB as gateway B + participant SL as SandboxLauncher (B) + MA->>GA: bridge_spawn(profile, host = B) + GA->>CX: publish(key = hostB, SpawnRequest) + CX->>GB: route to control.hostB + 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) + Note over MA: gid appears in every gateway's union roster,
then U1 addresses it normally +``` + +*Figure 5 — U4. Spawn is a *control* message to host B's gateway, which runs the launcher **locally** +(CB-500 §11 — a sandbox is spawned by its own host's gateway). The readiness gate is more valuable +here: the far side wants a positive "ready" before anyone sends. The new agent enters the roster and +is then reachable by U1.* + +### 7.3 U5 — presence federation + +```mermaid +flowchart LR + subgraph gwA["gateway A"] + hbA["heartbeat local agents
roster.hostA.*"] + viewA["union rosterView"] + end + subgraph gwB["gateway B"] + hbB["heartbeat local agents
roster.hostB.*"] + viewB["union rosterView"] + end + ros["bridge.roster (topic)"] + hbA -->|"publish"| ros + hbB -->|"publish"| ros + ros -->|"roster.# → roster.hostA"| viewA + ros -->|"roster.# → roster.hostB"| viewB + classDef bus fill:#2b6cb0,stroke:#1a365d,color:#ffffff; + class ros bus +``` + +*Figure 6 — U5. Each gateway publishes heartbeats for its local agents and binds `roster.#` to see +everyone's; both converge on the same eventually-consistent union view (CB-304 `rosterView`, +federated). Transient queues + heartbeat TTL keep it soft-state.* + +## 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]` +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]` +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.#`, + consume → maintain the union view; start heartbeating local agents. `[proposed]` +6. Automatic connection + topology recovery re-declares queues and re-attaches consumers after a + broker blip (dedup by `msgId` prevents double-queue). `[as-built]` + +## 9. As-built (CB-307) vs proposed (CB-308 / CB-500) + +| 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** | key by **globalId**; front with named `bridge.msg` topic exchange for observability | +| 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) | optionally broker-driven via `bridge.delay` | +| Identity in key | `paneId` | `globalId` (UUID); host in roster metadata | + +**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. + +## Related pages + +- [1. Architecture](1-Architecture) — the two modes, the invariants, the MCP asymmetry. +- [8. Roadmap](8-Roadmap) — CB-307 (delivery), CB-308 (federation), CB-500 (multi-tier) staging. +- [9. Implementation](9-Implementation) — `msg.AmqpReplyInbox` / `ReplyPushLoop` as-built. + +--- +*Design chapter for the cross-host broker fabric. As-built pieces verified against `msg.AmqpReplyInbox` +at wiki-source parity; proposed pieces track CB-308 (`docs/CB-308-Multi-Host-Federation.md`) and +CB-500 (`docs/CB-500-Multi-Tier-Coordination.md`).* diff --git a/_Sidebar.md b/_Sidebar.md index f33807e..680ae93 100644 --- a/_Sidebar.md +++ b/_Sidebar.md @@ -13,6 +13,7 @@ 7. [Use Cases](7-Use-Cases) — the review scenario + mechanisms 8. [Roadmap](8-Roadmap) — stages, tech stack, tickets 9. [Implementation](9-Implementation) — as-built code map · classes · flows · state machines +10. [Cross-Host Messaging](10-Cross-Host-Messaging) — broker topology · exchanges · queues per entity --- 🟢 herdr-centric `bridged` · AgentAPI = fallback