diff --git a/1-Architecture.md b/1-Architecture.md index d228b8c..fdc54cc 100644 --- a/1-Architecture.md +++ b/1-Architecture.md @@ -148,9 +148,12 @@ may outrun a sane request timeout, use Mode 2.* > blocking `bridge_send` has already timed out (or was never opened), its `bridge_reply` finds no > live waiter to resolve. `bridged` now **holds that reply in a per-worker inbox** rather than > discarding it, and the primary collects it later keyed by target (`bridge_poll(target)` / -> `GET /sessions/{id}/replies`). Stage 1 is soft-state (in-memory, lost on a daemon restart); -> cross-restart durability arrives with the CB-307 Stage-2 broker adapter. See -> [Roadmap](8-Roadmap#delivery-reliability--multi-host-cb-306--cb-308). +> `GET /sessions/{id}/replies`). With a broker configured the reply is **durable** across a daemon +> restart (Stage 2, `AmqpReplyInbox` on LavinMQ); without one it is soft-state in-memory. And +> delivery is **active, not just pull** (Stage 3): the moment a reply lands with no open send, +> `bridged` injects a bounded, status-gated *drain nudge* into a same-host primary's own herdr pane +> — so the primary need not be polling to notice. An off-host / non-herdr primary keeps the pull +> path. See [Roadmap](8-Roadmap#delivery-reliability--multi-host-cb-306--cb-308). ### Mode 2 — asynchronous delivery (bridged-mediated) @@ -186,9 +189,10 @@ for the idle event, and injects.* > **Where a hook is used.** A `Stop`-hook appears in exactly the spots where neither an MCP > tool call nor a pane injection can serve — and in **every** case it targets `bridged`, never > a broker: -> - a **split-host primary that is not a herdr pane** (e.g. Opus on your Mac) — the only -> session `bridged` cannot inject into — runs a `Stop`-hook that **long-polls `bridged`** for -> queued messages; and +> - a **split-host primary that is not a herdr pane** (e.g. Opus on your Mac) — the one session +> `bridged` cannot inject into (a *same-host* primary **is** a herdr pane and gets the CB-307 +> Stage-3 drain nudge) — runs a `Stop`-hook that **long-polls `bridged`** for queued messages; +> and > - a **non-MCP (herdr-only) worker** may run a `Stop`-hook that **POSTs its reply to > `bridged`** at turn end, a structured alternative to scraping the pane (see > [Message Server](2-Message-Server) → *Reply envelope*). diff --git a/8-Roadmap.md b/8-Roadmap.md index f2a243f..6633126 100644 --- a/8-Roadmap.md +++ b/8-Roadmap.md @@ -143,14 +143,15 @@ primary. That asymmetry is the root of all three tickets. ```mermaid flowchart LR - cb306["CB-306 ✅
spawn-readiness gate"] --> cb307["CB-307
reliable worker→primary delivery"] + cb306["CB-306 ✅
spawn-readiness gate"] --> cb307["CB-307 ✅
reliable worker→primary delivery"] cb307 --> s1["Stage 1 ✅
ReplyInbox port + in-memory adapter"] - cb307 --> s2["Stage 2 ⏳
AMQP / LavinMQ adapter"] + cb307 --> s2["Stage 2 ✅
AMQP / LavinMQ durable adapter"] + cb307 --> s3["Stage 3 ✅
active push-to-primary + reminder"] s2 --> cb308["CB-308 ⏳
multi-host federation"] classDef done fill:#2f855a,stroke:#22543d,color:#ffffff; classDef todo fill:#b7791f,stroke:#7b341e,color:#ffffff; - class cb306,s1 done - class s2,cb308 todo + class cb306,s1,s2,s3 done + class cb308 todo ``` *Figure: CB-307's broker fabric is the foundation CB-308 stretches across hosts; CB-306 makes a @@ -170,12 +171,26 @@ spawn fail fast instead of stalling.* dependency. **Only terminal replies are queued** — `bridge_ask` (interactive) and the completion/failure fallbacks are deliberately *not* (would risk double-delivery). Live on the running daemon. - - **Stage 2 ⏳ (deferred, gitea #5 open).** `AmqpReplyInbox` behind the *same* port for - cross-restart durability — **LavinMQ** default (single Crystal binary; `com.rabbitmq:amqp-client` - works unchanged; native DLX + delayed-message exchange give the remind/backoff loop for free), - RabbitMQ interchangeable by URI. `broker:` config absent → in-memory, present → AMQP. Needs a - live broker to integration-test. **`bridged` stays soft-state — the broker owns message - durability, not the bus.** + - **Stage 2 ✅ (main `2bc5f3a`).** `AmqpReplyInbox` behind the *same* port for cross-restart + durability — **LavinMQ** default (single Crystal binary; `com.rabbitmq:amqp-client` works + unchanged; native DLX + delayed-message exchange), RabbitMQ interchangeable by URI. Mapping = + consume-and-hold with deferred manual ack: each target owns a durable queue `agent..inbox`; + a manual-ack consumer pulls persistent messages into an in-memory held map but does *not* ack until + the primary drains, so a `java -jar` bounce leaves them on the broker for redelivery. `broker:` + config absent → in-memory, present → AMQP. Contract test `@Tag("contract")` runs against a RabbitMQ + container. Dogfooded live (reply survived a daemon bounce, redelivered + acked exactly once). + **`bridged` stays soft-state — the broker owns message durability, not the bus.** + - **Stage 3 — active push-to-primary + reminder ✅ (main `d4c9704`, gitea #5 closed).** The durable + inbox is a *landing zone* but delivery was still **pull** (the primary had to poll). Stage 3 makes + it **active**: a `ReplyPushLoop` nudges the primary the moment a reply lands with no open send, and + reminds (bounded) until drained. Push channel resolved the crux design question — the same-host + primary **is** a herdr pane (its `terminal_id` resolves on every MCP call), so the loop injects a + *drain nudge* (not the payload) into the primary's own pane via `AgentControl.send`, **status-gated** + (only when `injectable()`, never mid-turn) and **bounded** (`primary.push_reminders`=5, + `push_backoff_ms`=15000). **Ack = drain**: stop when `inbox.peek(target).isEmpty()`. A single-slot + `PrimaryRegistry` learns the primary from orchestration-side tools; an off-host / non-herdr primary + leaves it empty → the loop is a no-op and delivery degrades to pull (reply never lost). Optional + per-`msgId` `bridge_ack` tool for finer control than drain-all. Dogfooded live end-to-end. - **CB-308 — multi-host federation ⏳ (design note, gitea #6; depends on CB-307 Stage 2).** A primary on host A delegating to workers on hosts B, C… with no host learning another's terminals. Built on diff --git a/9-Implementation.md b/9-Implementation.md index 016e59d..fc1e7bf 100644 --- a/9-Implementation.md +++ b/9-Implementation.md @@ -4,8 +4,8 @@ > classes, flows, and state machines in the source tree, as a companion to the design-level > [1. Architecture](1-Architecture) and [2. Message Server](2-Message-Server). Every enum, > constant, and route below was verified against source at main `3aa69a9`; the `msg`-layer -> **reply-inbox** (CB-307 Stage 1) and its drain surface landed at `ba6b4a5` and are folded in -> below. +> **reply-inbox** (CB-307 Stage 1 `ba6b4a5`), its **AMQP durable adapter** (Stage 2 `2bc5f3a`), and +> the **active push-to-primary loop** (Stage 3 `d4c9704`) are folded in below. `bridged` is a single-host Java 25 / Maven daemon: the sole gateway between an on-subscription **primary** (Opus) and off-subscription **workers**, speaking to the **herdr** PTY manager @@ -106,18 +106,21 @@ calls `bridge_reply`. This is where the turn state machine lives. | `TurnListener` | interface | Turn-lifecycle callbacks: `onDelivered`, `onTurnComplete`, `onTurnFailed`. | | `WorkerPresence` | class | Tracks which workers have connected the bridge MCP (`markPresent`, `isPresent`, `forget`). | -### `msg` — the service core (4 classes) +### `msg` — the service core (6 classes) Owns the forward rendezvous (`bridge_send` → `bridge_reply`) and the reverse rendezvous (`bridge_ask` → answer), plus async fire-and-poll — and, since CB-307, the **reply inbox** that -holds a worker's terminal reply when *no* forward send is open, instead of dropping it. +holds a worker's terminal reply when *no* forward send is open (instead of dropping it) and the +**push loop** that actively nudges the primary to drain it. | Class | Kind | Role | |---|---|---| | `MessageService` | class | Orchestrates send/reply/ask/answer + async dispatch (`send`, `answer`, `ask`, `sendAsync`, `poll`), and routes `bridge_reply` through **`reply`** (resolve an open send, else publish to the inbox) + **`drainReplies`** (peek-then-ack a target's held replies). | | `Rendezvous` | class | Low-level registry of forward waiters + reverse-ask futures (`open`, `resolve`, `resolveQuestion`, `openAsk`, `answerAsk`, `closeAsk`, `resolveCompletion`, `resolveFailure`). **Untouched by CB-307** — it stays a pure synchronization primitive; a `false` from `resolve` (no live waiter) is what triggers the inbox publish, one layer up in `MessageService`. | -| `ReplyInbox` | interface | The port (CB-307): `publish(target, msgId, content)` (idempotent, dedup by `msgId`), `peek(target)` (non-destructive FIFO snapshot), `ack(target, msgId)`. Nested `InboxMessage(msgId, target, content)` record. The Stage-2 AMQP/LavinMQ adapter will implement the *same* port for cross-restart durability. | -| `InMemoryReplyInbox` | class | Stage-1 default adapter — per-target FIFO in a `ConcurrentHashMap>`, thread-safe, dedup by `msgId`. **Soft-state, not persistence** — undrained replies are lost on a `java -jar` bounce, consistent with "bridged stays soft-state; the broker owns durability." | +| `ReplyInbox` | interface | The port (CB-307): `publish(target, msgId, content)` (idempotent, dedup by `msgId`), `peek(target)` (non-destructive FIFO snapshot), `ack(target, msgId)`. Nested `InboxMessage(msgId, target, content)` record. Two adapters implement the *same* port — soft-state in-memory and durable AMQP. | +| `InMemoryReplyInbox` | class | Stage-1 default adapter — per-target FIFO in a `ConcurrentHashMap>`, thread-safe, dedup by `msgId`. **Soft-state, not persistence** — undrained replies are lost on a `java -jar` bounce, consistent with "bridged stays soft-state; the broker owns durability." Selected when `broker:` config is absent. | +| `AmqpReplyInbox` | class | Stage-2 durable adapter (CB-307 `2bc5f3a`) — one durable queue `agent..inbox` per target; **consume-and-hold with deferred manual ack** (a manual-ack consumer pulls persistent messages into an in-memory held map but doesn't ack until the primary drains, so a bounce leaves them on the broker for redelivery). Dedup keys the held map by `msgId`; automatic connection + topology recovery re-declares queues and clears stale delivery-tags. LavinMQ default, RabbitMQ by URI swap. Selected when `broker:` config is present. | +| `ReplyPushLoop` | class | Stage-3 active push (CB-307 `d4c9704`) — a dedicated status-gated scheduled loop (mechanism (b), **not** the worker `Injector`, so it stays decoupled from `WorkerPresence`). `onReplyQueued(target)` (called on the no-waiter branch of `reply`) starts a bounded reminder loop: at each tick `decide(target, count)` returns `INJECT` / `WAIT_BUSY` / `STOP`, injecting a *drain nudge* into the primary's own pane via `AgentControl.send` only when `status(primary).injectable()`. Stops when `peek(target)` is empty (ack = drain) or the reminder cap is reached; idempotent per target. | **Enums (verbatim).** `Rendezvous.Kind`: `REPLY`, `COMPLETION`, `FAILED`, `QUESTION`. `MessageService.Outcome`: `REPLIED`, `COMPLETED_UNREPLIED`, `WORKER_FAILED`, `QUESTION`, @@ -137,7 +140,7 @@ is `worker.ClaudeCodeLauncher`; future adapters (e.g. Codex) implement the same | `SpawnRequest` | record | Spawn parameters: `profileName`, `requestedCwd`, `callerCwd`. A null/blank profile means "use the default"; a null/blank cwd means "inherit from config or caller". | | `Capability` | enum | Declared launcher capabilities: `MID_TURN_ASK`, `SELF_PR`, `WORKTREE`, `ORPHAN_REAP`. Stage A only *declares* them (advisory); verb-layer enforcement — a verb against a peer lacking a capability returning a clean "unsupported" rather than crashing — is planned, not yet wired. | -### `mcp` — the MCP north face (6 classes) +### `mcp` — the MCP north face (7 classes) Exposes `bridged` as a Streamable-HTTP MCP endpoint and resolves caller identity from the **connection**, not from tool arguments (unspoofable). @@ -146,6 +149,7 @@ Exposes `bridged` as a Streamable-HTTP MCP endpoint and resolves caller identity |---|---|---| | `BridgeMcp` | class | Builds the MCP server, registers `bridge_*` tools, holds thin tool adapters (`servlet`, `send`, `answer`, `ask`, `reply`, `spawn`). | | `ConnectionIdentity` | class | Resolves *who is calling*: peer PID → herdr pane → `terminal_id` (`resolve`, `callerTerminal`, `cwdForPid`). | +| `PrimaryRegistry` | class | Single-slot thread-safe holder of the primary's `terminal_id` (CB-307 Stage 3). `record(id)` learns it from orchestration-side tools (`bridge_send`/`bridge_spawn`) when the caller resolves to a non-null terminal that is *not* a registered worker; `primaryTerminal()` / `isKnown()` feed the push loop. Optional constructor-pin (`primary.terminal` config) for an operator override; stays empty (→ push degrades to pull) when the primary is off-host or non-herdr. | | `PeerPidLookup` | interface | Abstracts OS peer-PID lookup for a loopback source port. | | `LsofPeerPidLookup` | class | `lsof` impl; excludes bridged's own PID. | | `ProcessCwdLookup` | interface | Abstracts PID → cwd lookup. | @@ -155,7 +159,8 @@ Exposes `bridged` as a Streamable-HTTP MCP endpoint and resolves caller identity (worker → structured answer; if no send is open it is now **held in the reply inbox**, not errored) · `bridge_ask` (worker pauses to ask primary) · `bridge_status` · `bridge_poll` (async ticket; with an optional `target`, **drains that worker's held replies** from the inbox) · -`bridge_spawn` · `bridge_list` · `bridge_stop` · `bridge_profiles`. +`bridge_ack` (ack one held reply by `msgId` — finer than drain-all) · `bridge_spawn` · +`bridge_list` · `bridge_stop` · `bridge_profiles`. **Invariant:** `bridge_reply`/`bridge_ask` accept *no* identity argument; it comes only from the transport context. **CB-307 scope:** only a *terminal* `bridge_reply` with no open send is queued — `bridge_ask` (interactive; the worker blocks and can't consume a late answer) and the injector's @@ -250,8 +255,20 @@ its send has already timed out (the ~60s client window) or was never opened, so finds no waiter. Before CB-307 that reply was silently discarded (the worker got an error / REST `409`). Now `MessageService.reply` publishes it to the `ReplyInbox` under a fresh `msgId`, and the primary collects it later keyed by target — `bridge_poll(target)` or `GET /sessions/{id}/replies` -(peek → deliver → ack, so an in-flight failure re-surfaces it). Soft-state only: a daemon bounce -clears any undrained replies; cross-restart durability is CB-307 Stage 2 (the AMQP adapter). +(peek → deliver → ack, so an in-flight failure re-surfaces it). With `broker:` configured the +`AmqpReplyInbox` gives this **cross-restart durability** (Stage 2) — the held reply survives a +`java -jar` bounce and is redelivered; without it the in-memory adapter is soft-state (undrained +replies clear on restart). + +Delivery is no longer purely pull. When the reply lands with no waiter, `reply` also calls +`ReplyPushLoop.onReplyQueued(target)`, which **actively nudges the primary to drain** (Stage 3): it +injects a *"run `bridge_poll(target=…)`"* turn into the primary's own herdr pane — the same +`AgentControl.send` primitive that delivers to workers, pointed at the primary — but only when the +primary is `injectable()` (never mid-turn), bounded to `primary.push_reminders` nudges on a +`push_backoff_ms` schedule, and stopping the instant a drain empties the inbox. The nudge carries +*no payload* (the drain response is the clean transport), so re-nudging is idempotent, and if every +push fails the durable inbox is still the backstop. An off-host / non-herdr primary leaves +`PrimaryRegistry` empty, so the loop is a no-op and delivery cleanly degrades to pull. ## Flow: reverse rendezvous (`bridge_ask` → answer)