From 4320c1ca5263fb40ce1ca8637a99cc60d5998ffe Mon Sep 17 00:00:00 2001 From: Kevin Nguyen Date: Mon, 10 Aug 2026 22:31:45 +0700 Subject: [PATCH] =?UTF-8?q?Chapter=2010:=20fold=20in=20the=20adversarial?= =?UTF-8?q?=20review=20=E2=80=94=20envelope,=20dual=20ack,=20queue=20migra?= =?UTF-8?q?tion?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Second-pass (Opus reviewer) findings, accepted and written through: - §2.1 the message envelope: AMQP-header fields + kinds; sig over exact body bytes + canonical header subset (no JSON canonicalization in the security path); verify to == routing key; reject-to-DLX on failure. CB-201's descope stays right single-host and wrong across hosts — stated. - §3 footnote corrected: the inbox keying INVERTS (sender-keyed as-built → recipient-keyed), so CB-308 is a migration, not a rename; .v2 queue names dodge the PRECONDITION_FAILED crash loop on in-place upgrade. - §5/§6 rewritten as the dual ack model: reply path at-least-once (ack after drain), forward path at-most-once (ack BEFORE inject) + INJECTED confirmation; basicQos so the backlog stays on the queue. - §7.4: expiry is gateway-enforced (expiresAt), broker TTL a backstop; NO_WAITER answered on arrival. - §10.1: caps made real by prefetch; expiry sweeper (head-of-line); DLQ consumer reports FAILED to the sender; DLQ itself capped. - §8 checklist: separate confirm/mandatory publish channel (return-before- confirm caveat), queue deletion on stop, host-level signed heartbeat with profile list, exclusive-lock takeover caveat. - §9: keying migration + multi-primary named as CB-500's delta. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_013ZGgxLQ2VpwZhEYoru8rkf --- 10-Cross-Host-Messaging.md | 163 +++++++++++++++++++++++++++---------- 1 file changed, 121 insertions(+), 42 deletions(-) diff --git a/10-Cross-Host-Messaging.md b/10-Cross-Host-Messaging.md index b1c233e..3c329e7 100644 --- a/10-Cross-Host-Messaging.md +++ b/10-Cross-Host-Messaging.md @@ -63,6 +63,33 @@ consumes the inboxes of its **co-located** entities and publishes to the inboxes carried in the roster, never in the routing key. `[as-built: keyed by paneId, single host]` → `[proposed: keyed by globalId]`. +### 2.1 The message envelope `[proposed]` + +Single-host, the envelope was deliberately descoped (CB-201): connection identity answers *who is +talking* on every call. A broker hop has no connection identity, so cross-host the envelope +returns — minimal, and carried in **AMQP headers**, not a JSON wrapper: + +| Field | Meaning | +|---|---| +| `v` | envelope schema version (§10.3) | +| `kind` | `SEND` · `REPLY` · `ASK` · `ANSWER` · `TOO_LATE` · `NO_WAITER` · `ABANDONED` · `FAILED` · `INJECTED` · `BROADCAST` · `PRESENCE` · `SPAWN` | +| `msgId` | unique per message — the dedup key | +| `traceId` | minted at the flow's first send, carried by every hop (§10.4) | +| `from` / `to` | sender / recipient global ids | +| `turnId` | conversation flows only (`ASK` / `ANSWER` / `TOO_LATE`) | +| `spawnId` | `SPAWN` only — doubles as the new worker's gid (CB-308 §7.8) | +| `ts` / `expiresAt` | publish time / gateway-enforced expiry — on **every** kind (also bounds replay) | +| `sig` | signature by the publishing gateway | + +The body stays opaque bytes (UTF-8 text today). **`sig` covers the exact body bytes plus the +canonical header subset above** — signing headers-plus-body keeps JSON canonicalization out of the +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 +`expiresAt` plus a bounded seen-set — the at-most-once forward path (§6) has no broker redelivery +for a replayed duplicate to hide behind. + ## 3. Broker topology — exchanges Five exchanges. Point-to-point messaging, presence, and control are separated so each has its own @@ -74,13 +101,20 @@ durability and fan-out semantics. | `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] | +| `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] | -*\* **[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.* +*\* **[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 +keyed by the **worker a primary-bound reply came from** (sender-keyed), drained per-worker. This +chapter keys the inbox by **recipient** (invariant 2). Same name shape, opposite meaning — so +CB-308 is a **migration, not a rename**. Migrated queues take a version-suffixed name +(`agent..inbox.v2`): AMQP refuses to redeclare an existing durable queue with new arguments +(`PRECONDITION_FAILED` — a crash loop on an in-place upgrade from v1.0.0), and the suffix keeps +old sender-keyed and new recipient-keyed queues apart while both exist. The per-worker drain +surface (`bridge_poll(target)`) survives by filtering on the envelope's `from` field (§2.1). This +is also where CB-201's story completes: descoping the envelope was right on one host (connection +identity routes everything) and wrong across hosts (a broker hop has none) — the envelope returns +as §2.1.* ```mermaid flowchart TB @@ -127,8 +161,8 @@ 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 | [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) | +| 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 | @@ -156,39 +190,57 @@ stays soft-state; only *messages* are durable. See [1. Architecture](1-Architect 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`. +5. **Two ack models — matched to what a duplicate costs.** *Reply/pull path* (into a main): + at-least-once — ack **after** drain; a redelivered reply is benign and deduped by `msgId`. + *Forward path* (a task brief into a worker): at-most-once — ack **before** the local inject; a + crash in the window loses the inject, which surfaces as a **visible** failure at the sender + (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 + 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). 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]` +## 6. Message lifecycle & ack semantics `[as-built for the reply path; forward path proposed]` + +Two lifecycles, split by what a duplicate would cost (invariant 5). The consumer runs with a +**`basicQos` prefetch bound**, so the backlog stays *on the queue* — where the caps and expiry of +§10.1 can act — instead of draining into gateway heap. ```mermaid sequenceDiagram autonumber participant SND as sender gateway - participant MX as bridge.msg - participant Q as agent.GID.inbox + participant Q as agent.GID.inbox.v2 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) + rect rgb(235, 244, 255) + Note over SND,E: REPLY / pull path — at-least-once (as-built shape) + SND->>Q: publish (confirms + mandatory, §8 step 6) + OWN->>Q: consume (prefetch-bounded, dedup by msgId) + OWN->>E: hold for MCP pull + E-->>OWN: drained + OWN->>Q: basicAck — a bounce BEFORE ack → redelivered, deduped + end + rect rgb(255, 245, 235) + Note over SND,E: FORWARD path (brief into a worker) — at-most-once + SND->>Q: publish (confirms + mandatory) + OWN->>Q: consume + OWN->>Q: basicAck FIRST — a redelivered brief can never double-inject + OWN->>E: status-gated herdr inject + OWN->>SND: publish INJECTED (correlated to msgId) + Note over SND: no INJECTED within bound → loud, fast failure (retry is the sender's call) + end ``` -*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."* +*Figure 3 — the reply path keeps CB-307's consume-and-hold with deferred ack (`AmqpReplyInbox` +today, `9-Implementation`); the forward path inverts the ack order so a crash costs a visible +loss, never a silent duplicate injection. `INJECTED` closes the blind window between ack and +inject.* ## 7. Use-case walkthroughs @@ -273,7 +325,13 @@ 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` as **short-lived** messages (per-message TTL) carrying the asking turn's `turn_id`: +`bridge.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 +publishes `NO_WAITER` straight back if there is none (the worker learns "nobody listening" in one +round trip, as it does synchronously today), and gateway B injects an `ANSWER` **only if that turn +is still waiting**: ```mermaid sequenceDiagram @@ -284,11 +342,11 @@ sequenceDiagram 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) + GB->>BR: publish ASK (key = gidMA, expiresAt, 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) + GA->>BR: publish ANSWER (key = gidW, expiresAt, turn_id) BR->>GB: route alt turn still waiting GB->>W: inject answer — same turn resumes @@ -319,33 +377,48 @@ On boot, a gateway declares and wires exactly its own slice: 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`, **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]` +3. For each **local** entity: declare `agent..inbox.v2` (durable, DLX = `bridge.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, + 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 → 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]` + 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 + check rests on them. `[proposed]` +6. Open a **separate publish channel** with **publisher confirms**, the **`mandatory`** flag, and + a return listener — an unroutable publish is a fast error, not a black hole, and confirms never + serialize the consume/ack channel. Ordering caveat: a *return* arrives **before** the confirm, + so "confirmed" ≠ "routed" — check the returned-set at confirm time. `mandatory` is false only + for `BROADCAST` (an empty group is legal silence). `[proposed]` +7. On session stop / reap: **delete the worker's inbox queue** — its `broadcast.*` bindings die + with it, so no broadcasts to the dead; `x-expires` collects queues orphaned by a crashed + gateway. `[proposed]` +8. Automatic connection + topology recovery re-declares queues and re-attaches consumers after a + broker blip (dedup by `msgId` prevents double-queue). `[as-built]` Exclusive-consumer caveat: + after a gateway crash the broker holds the stale lock until its heartbeat times out — keep the + broker heartbeat short (~10s) so takeover is quick. `[proposed]` ## 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 | +| 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` | | 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` | +| 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` | | 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 | +| 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 @@ -357,9 +430,15 @@ Cross-cutting rules settled in a design review of this chapter + CB-308. The sen 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. +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 + 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 + (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) — so the broker runs on a machine inside the trust circle, never a shared or rented one.