Chapter 10: fold in the adversarial review — envelope, dual ack, queue migration

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) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ZGgxLQ2VpwZhEYoru8rkf
2026-08-10 22:31:45 +07:00
parent 36bb86588a
commit 4320c1ca52
+121 -42
@@ -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.<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. | [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.<gid>.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.<target>.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.<gid>.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.<gid>.inbox` | `bridge.msg` · `<gid>` **+ `broadcast.all` / `broadcast.<group>` (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.<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) |
| 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 |
@@ -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<br/>(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.<gid>.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 `<gid>` **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.<gid>.inbox.v2` (durable, DLX = `bridge.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,
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 →
run the launcher locally. `[proposed]`
5. Declare `roster.<thisHost>` (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.<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
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.<target>.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.<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` |
| 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) | 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.<gid>.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.