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

8.6 KiB

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.

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.