Files
fleetd/bridged/docs/CB-307-Push-Loop.md
Dai Ha 131e7b1ccd CB-307: lock push-loop injection to dedicated status-gated loop (mechanism b)
Chose a small dedicated scheduled loop over AgentControl.send guarded by an
injectable status check, instead of reusing the worker Injector (which couples
to WorkerPresence/StatusPoller). Isolated + unit-testable via injected clock.
Records live ground truth: this primary resolves to term_656c8cc03e1f0b1 (w2:pY)
— confirms the primary runs in a herdr pane so the push path is exercisable.
2026-07-19 09:44:21 +02:00

143 lines
8.6 KiB
Markdown

# 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.<target>.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.&lt;target&gt;.inbox<br/>(durable, LavinMQ)")]
MS -->|"notify"| LOOP["ReplyPushLoop"]
LOOP -->|"status-gated inject"| PANE["primary's herdr pane"]
PANE -->|"primary drains"| DRAIN["bridge_poll(target)<br/>= peek + ack"]
DRAIN -->|"inbox now empty"| LOOP
LOOP -.->|"still non-empty →<br/>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: "<id>" }` 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 `<target>` returned a reply — run `bridge_poll(target=<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 (mechanism (b), chosen).** A dedicated lightweight scheduled loop,
**not** the worker `Injector`. It injects via `AgentControl.send(primaryTerminal, nudge)`
(the same herdr `agent.send` = `pane send-text` + submit that delivers to workers) only when
`AgentControl.status(primaryTerminal).injectable()` (IDLE/BLOCKED) — never mid-turn. This keeps
the primary path fully isolated from `WorkerPresence`/`StatusPoller` (which are worker-scoped),
and makes it unit-testable with a fake `AgentControl` + an injected clock (per the CB-306
`LongSupplier` clock + `Runnable` sleeper seam). Rejected (a) reuse-the-Injector: it would force
the primary terminal into the worker poller set and couple to worker-presence semantics — more
integration surface, harder to test, no real gain for a bounded reminder.
- **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 → dedicated loop.** The existing `Injector` gates delivery on
`ready.test(target)` = `WorkerPresence` (the *worker's* MCP connected). The primary is not
in `WorkerPresence`, so reusing `Injector` would mean forcing the primary terminal into the
worker `StatusPoller` set and swapping the `ready` predicate — extra integration surface with
worker-scoped machinery. Decision: **mechanism (b)** — a small dedicated scheduled loop that
calls `AgentControl.status(primaryTerminal).injectable()` then `AgentControl.send(...)`, with
an injected clock. Isolated from worker presence, trivially unit-testable, sufficient for a
bounded reminder. (Verified live: this primary resolves to `term_656c8cc03e1f0b1`, pane
`w2:pY` — the primary genuinely runs in a herdr pane on this host, so the path is exercisable.)
## 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.