Ch.10: cross-host lead-to-lead already shipped; rename the proposed exchanges

The scope note said only the single-host reply inbox was as-built. That is
wrong: lead-to-lead coordination across hosts works today over a shared AMQP
vhost (coordinator: block, fleet_send{coordId}, LeadMailbox + LeadCoordLoop).

It deliberately does not use the exchange topology in section 3 — it is a
durable mailbox per lead, carrying coordination only, never a task. Said so at
the top and added a short as-built table to section 9, so the chapter cannot be
read as 'none of this is real yet'.

Also renamed the proposed exchanges bridge.* -> fleet.* with the product.
Dai Ha
2026-08-31 10:12:16 +07:00
parent 798d3a8ef2
commit 48c8a1e328
+68 -42
@@ -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.<group>` for U8. | [proposed]* |
| `bridge.roster` | topic | yes | Presence heartbeats (U5). Routing key = `roster.<host>.<globalId>`. | [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.<group>` for U8. | [proposed]* |
| `fleet.roster` | topic | yes | Presence heartbeats (U5). Routing key = `roster.<host>.<globalId>`. | [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.<target>.inbox` is
@@ -120,10 +132,10 @@ as §2.1.*
```mermaid
flowchart TB
subgraph exch["Exchanges"]
msg["bridge.msg<br/>(topic)"]
ros["bridge.roster<br/>(topic)"]
ctl["bridge.control<br/>(direct)"]
dlx["bridge.dlx<br/>(fanout)"]
msg["fleet.msg<br/>(topic)"]
ros["fleet.roster<br/>(topic)"]
ctl["fleet.control<br/>(direct)"]
dlx["fleet.dlx<br/>(fanout)"]
end
subgraph gwA["gateway A (host A)"]
inA["agent.&lt;gidA&gt;.inbox<br/>(durable)"]
@@ -135,7 +147,7 @@ flowchart TB
ctlB["control.hostB<br/>(durable)"]
rosB["roster.hostB<br/>(exclusive, transient)"]
end
dlq["bridge.dlq<br/>(durable)"]
dlq["fleet.dlq<br/>(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.<gid>.inbox` | `bridge.msg` · `<gid>` **+ `broadcast.all` / `broadcast.<group>` (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.<gid>.inbox` | `bridge.msg` · `<gid>` | 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.<gid>.inbox` | `bridge.msg` · `<gid>` | its **local** gateway → **MCP pull** | durable | top MCP client; same fabric |
| Gateway (host) | `control.<host>` | `bridge.control` · `<host>` | that gateway | durable | cross-host spawn/stop (U4) |
| Gateway (host) | `roster.<host>` | `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.<gid>.inbox` | `fleet.msg` · `<gid>` **+ `broadcast.all` / `broadcast.<group>` (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.<gid>.inbox` | `fleet.msg` · `<gid>` | 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.<gid>.inbox` | `fleet.msg` · `<gid>` | its **local** gateway → **MCP pull** | durable | top MCP client; same fabric |
| Gateway (host) | `control.<host>` | `fleet.control` · `<host>` | that gateway | durable | cross-host spawn/stop (U4) |
| Gateway (host) | `roster.<host>` | `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,<br/>then U1 addresses it normally
```
@@ -309,7 +321,7 @@ flowchart LR
hbB["heartbeat local agents<br/>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.<gid>.inbox.v2` (durable, DLX = `bridge.dlx`,
3. For each **local** entity: declare `agent.<gid>.inbox.v2` (durable, DLX = `fleet.dlx`,
**capped**: `x-max-length` + message TTL, `x-expires` as the orphan-GC backstop), bind to
`bridge.msg` keys `<gid>` **and** `broadcast.all` (+ its groups), start a **manual-ack,
`fleet.msg` keys `<gid>` **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.<thisHost>` (durable), bind to `bridge.control` key `<thisHost>`, consume →
4. Declare `control.<thisHost>` (durable), bind to `fleet.control` key `<thisHost>`, consume →
run the launcher locally. `[proposed]`
5. Declare `roster.<thisHost>` (exclusive, auto-delete), bind to `bridge.roster` key `roster.#`,
5. Declare `roster.<thisHost>` (exclusive, auto-delete), bind to `fleet.roster` key `roster.#`,
consume → maintain the union view; start heartbeating local agents **and a host-level entry**
(`roster.<thisHost>._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.<target>.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.<gid>.inbox.v2` on `bridge.msg` — a semantic **migration** (§3 footnote); per-worker drain preserved via envelope `from` |
| Per-entity inbox | `agent.<target>.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.<gid>.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.<host>` transient queues → federated union |
| Spawn | local call in-process | `bridge.control` + `control.<host>` 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.<host>` transient queues → federated union |
| Spawn | local call in-process | `fleet.control` + `control.<host>` 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) —