diff --git a/bridged/docs/CB-307-Push-Loop.md b/bridged/docs/CB-307-Push-Loop.md new file mode 100644 index 0000000..faa4f2a --- /dev/null +++ b/bridged/docs/CB-307-Push-Loop.md @@ -0,0 +1,137 @@ +# CB-307 — Active push-to-primary + reminder loop (the reliability layer) + +**Status:** design (2026-07-19). Builds directly on the shipped durable landing zone +(`AmqpReplyInbox`, main `2bc5f3a`, dogfooded live). gitea #5. + +## Why this exists + +Stage 2 gave a worker→primary reply a **durable place to wait** when no `bridge_send` is +open: it lands in `agent..inbox` on the broker and survives a daemon bounce. But +delivery is still **pull** — the primary only sees the reply if it happens to call +`bridge_poll(target)` / `GET /sessions/{id}/replies`. A reply can sit indefinitely while +the primary works on something else. + +This layer makes delivery **active**: the bridge *pushes* a nudge to the primary the moment +a reply lands, and keeps reminding (bounded) until the primary drains it. At-least-once, +dedup by `msgId`, and — critically — it never loses the reply even if every push fails, +because the durable inbox is the backstop. + +## The hard constraint it works around + +The bridge is an MCP **server**; the primary is an MCP **client**. A server cannot call +into a client. So "push to the primary" cannot be an MCP response — it needs a *sideband* +channel. The chosen channel: **inject a synthetic user-turn into the primary's own herdr +terminal pane** — the same mechanism the bridge already uses to deliver tasks to workers, +pointed at the primary's pane instead. + +```mermaid +flowchart LR + W["worker"] -->|"bridge_reply (no open send)"| MS["MessageService.reply"] + MS -->|"inbox.publish"| INBOX[("agent.<target>.inbox
(durable, LavinMQ)")] + MS -->|"notify"| LOOP["ReplyPushLoop"] + LOOP -->|"status-gated inject"| PANE["primary's herdr pane"] + PANE -->|"primary drains"| DRAIN["bridge_poll(target)
= peek + ack"] + DRAIN -->|"inbox now empty"| LOOP + LOOP -.->|"still non-empty →
re-inject on backoff"| PANE + classDef store fill:#2c5282,stroke:#1a365d,color:#ffffff; + class INBOX store +``` + +*Figure 1 — a reply lands in the durable inbox; the push loop nudges the primary's pane; +the primary's drain acks it; a still-full inbox triggers a bounded re-nudge.* + +## Three increments + +### Increment 1 — learn & store the primary's terminal_id + +**Finding (seam map):** `ConnectionIdentity.resolve(remoteAddr, remotePort)` already returns +the caller's herdr `terminal_id` for *every* MCP call, via `PaneLocator.terminalForPid` +(walks `pane.list`, matches the caller PID to a pane's process tree). It is non-null whenever +the caller runs in a herdr pane on this host. Today it's discarded for the primary +(`presence.markPresent` is a no-op on it). + +**Plan:** a single-slot `PrimaryRegistry` (thread-safe) holding the primary's `terminal_id`. +Populate it from the **orchestration-side** MCP tools — `bridge_send`, `bridge_spawn`, +`bridge_poll`, `bridge_list`, `bridge_status`, `bridge_profiles` — capturing +`callerTerminal(exchange)` when it is (a) non-null and (b) **not** a registered worker +session in `SessionManager`. That caller is, by construction, the primary. Worker-side tools +(`bridge_reply`, `bridge_ask`) never set it. + +- **Config override / pin:** a `primary: { terminal: "" }` block in `BridgedConfig` + (nested record, same shape as `Broker`). Lets an operator pin it, or supply it when + derivation can't (see degrade case). +- **Degrade:** if the primary is off-host or in a non-herdr terminal, `terminalForPid` + returns null and no override is set → **the registry stays empty → the push loop is a + no-op and we fall back to pull** (today's behaviour). The reply is never lost; it's just + not actively pushed. This is a safe, explicit degradation, not a failure. + +### Increment 2 — the push loop + +A `ReplyPushLoop` component, notified at the single no-waiter call site +(`MessageService.reply` → the `inbox.publish` branch, `MessageService.java:192`). + +- **Inject a nudge, not the payload.** The injected turn tells the primary *to drain* + (e.g. "Worker `` returned a reply — run `bridge_poll(target=)` to collect + it"), it does **not** carry the reply text. Rationale: replies can be large/multiline and + terminal injection would mangle them; the drain response is the clean transport. Keeps the + push idempotent — re-nudging is harmless. +- **Ack = drain.** The primary draining (`drainReplies` = peek + ack) is the acknowledgement. + The loop's **stop condition is `inbox.peek(target).isEmpty()`** — the reply is gone from the + inbox because it was acked. No new `bridge_ack` tool needed for v1 (see Increment 3). +- **Status-gated injection.** Never inject mid-turn. Reuse the `Injector`/`StatusPoller` + discipline: deliver only when the primary's terminal samples **injectable** (IDLE/BLOCKED), + one message per turn. **Readiness caveat (below).** +- **Bounded reminder / backoff.** While `peek(target)` stays non-empty, re-inject on a + backoff schedule up to a cap (N reminders or a max duration; config + `primary.push_reminders` / `primary.push_backoff_ms`). After the cap, **stop reminding** — + the reply remains in the durable inbox and the next natural poll (or a later worker reply's + nudge) still surfaces it. Bounded so the bridge never spams the primary. + +### Increment 3 — optional per-`msgId` `bridge_ack` tool (deferred) + +Drain-as-ack is coarse: it clears *all* pending replies for a target at once. If finer +control is ever needed (ack one reply, leave others held), add a `bridge_ack(msgId)` tool +mapping to `inbox.ack(target, msgId)` — the port already supports per-`msgId` ack. Not built +in v1; the stop-on-empty loop is sufficient. + +## The two subtleties (decided here) + +1. **Which caller is "the primary"?** Connection-derived, not self-reported: the caller whose + resolved terminal is non-null **and not a registered worker session**, seen on an + orchestration-side tool. This never mislabels a worker (workers are in `SessionManager`) + and needs no new env var or argument (identity stays connection-derived, per the existing + `BridgeMcp` invariant). + +2. **Readiness-gate mismatch.** The existing `Injector` gates delivery on + `ready.test(target)` = `WorkerPresence` (the *worker's* MCP connected). The primary is not + in `WorkerPresence`. So the primary push path uses a **different liveness signal**: the + primary is provably MCP-connected at the moment it calls us (that's how we learned its + terminal), and we still status-gate on its terminal sampling **injectable** via the poller. + Concretely, the primary push path either (a) uses a dedicated `Injector` instance whose + `ready` predicate is "primary terminal is known" (always true once registered), or + (b) a lighter direct `AgentControl.send` guarded by a `StatusPoller` injectable sample. + Decision: **(a)** — reuse the `Injector` queue/one-per-turn/backoff machinery with a + primary-appropriate `ready` predicate, rather than reimplement gating. + +## Boundary note + +This is the first time the bridge **writes into the primary's pane** — a new direction of +control. It stays within the communication-bus identity: the injection is a **nudge** (a +synthetic "go drain your replies" turn), **status-gated** so it never interrupts a turn, +**bounded** so it never spams, carries **no env** and **never crosses the subscription +boundary**. The bridge is signalling the primary that it has mail — not driving its work. + +## Test plan + +- **Unit (hermetic):** `PrimaryRegistry` set/clear/override; the "caller is primary iff + non-null terminal AND not a registered session" predicate; the loop's stop-on-empty and + bounded-reminder logic with an injected clock + a fake injector (no real herdr). +- **Live dogfood (primary-side):** with the daemon on the broker jar + a real worker, + delegate a task, let the worker reply after the `bridge_send` window closes, and observe the + bridge inject a drain nudge into *this* primary pane; confirm draining stops the reminders; + confirm an unreachable primary (registry empty) degrades to pull with no loss. + +## Out of scope + +Multi-host (CB-308) — the push loop is local-only; a remote primary is reached by its own +local gateway, not cross-host injection. Federation reuses this loop per-gateway.