CB-307: reliable worker→primary delivery — push worker messages to main with ACK + reminder (at-least-once) #5

Closed
opened 2026-07-18 15:01:21 +02:00 by ltms · 6 comments
Owner

Problem — the delivery asymmetry

Delivery in the bridge is not symmetric:

  • Bridge → worker is a push: MessageService calls Injector.enqueue(target, content), and
    the injector delivers into the worker's herdr pane when the worker is idle/blocked. The bridge
    actively drives the message into the peer.
  • Bridge → primary is a pull: a worker's reply/question only lands if the primary already has
    an open blocking bridge_send for the injector/rendezvous to resolve
    (Rendezvous.resolveQuestion / reply resolution). If the primary dispatched async (block:false)
    or the send already returned, the message just sits on the ticket and the primary must poll
    (bridge_poll / GET /tasks/{ticket} → phase) to discover it. There is no injector path into
    main
    .

Root cause: MCP is client-initiated. The bridge is an MCP server and the primary is an MCP
client; a server cannot call into a client. So the bridge can push to workers (it drives their
herdr panes) but cannot push to the primary — it can only answer when the primary calls.

Consequence: if the primary isn't actively polling and has no open send, a completed worker turn can
sit undelivered — perceived as a "communication break" even though transport is healthy.

Proposal — push + ACK + remind (at-least-once to the primary)

Make worker→primary delivery reliable rather than poll-dependent:

  1. Push: when a worker produces a message destined for the primary (reply, bridge_ask,
    completion), the bridge proactively delivers it to the primary — the same way it injects into a
    worker's pane. If the primary is itself a herdr-addressable session, deliver via the injector into
    the primary's pane; otherwise define the push channel explicitly.
  2. ACK: delivery is not considered done until the primary acknowledges it (an ack verb / ack
    on next MCP call carrying the message id). Unacked messages stay pending.
  3. Remind: a bounded reminder loop re-surfaces unacked messages to the primary on an interval
    (with backoff + a max attempts / TTL cap so it can't spin forever), until acked or expired. This
    is the "bridge reminds main about the worker's message" behavior.

Net effect: at-least-once delivery to the primary with dedup on the primary side by message id — the
primary no longer has to be polling at the right moment to avoid missing a worker turn.

Design questions to resolve

  • Push channel to the primary. If the on-subscription primary runs in a herdr pane the daemon can
    see, injector push is symmetric with worker delivery. If not, what is the channel? (This is the
    crux — the asymmetry exists precisely because the primary is an MCP client.)
  • Idempotency. Message ids + primary-side dedup so reminders don't double-apply.
  • Interaction with async dispatch (CB-107). The reminder loop is the reliability layer under
    fire-and-poll; poll becomes an optimization, not the only delivery path.
  • Backoff + TTL. Bounded reminders; expire + surface in roster rather than remind forever.
  • Ordering. Whether reminders must preserve per-worker ordering to the primary.

Relationship to CB-306

  • CB-306 = spawn-readiness (make reachability honest at spawn).
  • CB-307 = reliable delivery (make worker→primary messages survive a non-polling primary).

Together they close the two ends of the "communication break" class: a peer that isn't actually
reachable yet, and a message that has nowhere to be pushed.

Layer: msg (MessageService / Rendezvous) + inject + a primary-side ack contract.
Type: resilience / reliability. Stage: hardening. Originated from the delivery-asymmetry
review during CB-401 integration.

## Problem — the delivery asymmetry Delivery in the bridge is **not symmetric**: - **Bridge → worker** is a **push**: `MessageService` calls `Injector.enqueue(target, content)`, and the injector delivers into the worker's herdr pane when the worker is idle/blocked. The bridge actively drives the message into the peer. - **Bridge → primary** is a **pull**: a worker's reply/question only lands if the primary already has an **open blocking `bridge_send`** for the injector/rendezvous to resolve (`Rendezvous.resolveQuestion` / reply resolution). If the primary dispatched async (`block:false`) or the send already returned, the message just sits on the ticket and the **primary must poll** (`bridge_poll` / `GET /tasks/{ticket}` → `phase`) to discover it. There is **no injector path into main**. Root cause: MCP is client-initiated. The bridge is an MCP **server** and the primary is an MCP **client**; a server cannot call into a client. So the bridge can push to workers (it drives their herdr panes) but cannot push to the primary — it can only answer when the primary calls. Consequence: if the primary isn't actively polling and has no open send, a completed worker turn can sit undelivered — perceived as a "communication break" even though transport is healthy. ## Proposal — push + ACK + remind (at-least-once to the primary) Make worker→primary delivery **reliable** rather than poll-dependent: 1. **Push:** when a worker produces a message destined for the primary (reply, `bridge_ask`, completion), the bridge proactively delivers it to the primary — the same way it injects into a worker's pane. If the primary is itself a herdr-addressable session, deliver via the injector into the primary's pane; otherwise define the push channel explicitly. 2. **ACK:** delivery is not considered done until the primary **acknowledges** it (an ack verb / ack on next MCP call carrying the message id). Unacked messages stay pending. 3. **Remind:** a bounded reminder loop re-surfaces unacked messages to the primary on an interval (with backoff + a max attempts / TTL cap so it can't spin forever), until acked or expired. This is the "bridge reminds main about the worker's message" behavior. Net effect: at-least-once delivery to the primary with dedup on the primary side by message id — the primary no longer has to be polling at the right moment to avoid missing a worker turn. ## Design questions to resolve - **Push channel to the primary.** If the on-subscription primary runs in a herdr pane the daemon can see, injector push is symmetric with worker delivery. If not, what is the channel? (This is the crux — the asymmetry exists precisely because the primary is an MCP *client*.) - **Idempotency.** Message ids + primary-side dedup so reminders don't double-apply. - **Interaction with async dispatch (CB-107).** The reminder loop is the reliability layer *under* fire-and-poll; poll becomes an optimization, not the only delivery path. - **Backoff + TTL.** Bounded reminders; expire + surface in roster rather than remind forever. - **Ordering.** Whether reminders must preserve per-worker ordering to the primary. ## Relationship to CB-306 - **CB-306** = spawn-readiness (make *reachability* honest at spawn). - **CB-307** = reliable delivery (make *worker→primary* messages survive a non-polling primary). Together they close the two ends of the "communication break" class: a peer that isn't actually reachable yet, and a message that has nowhere to be pushed. **Layer:** `msg` (MessageService / Rendezvous) + `inject` + a primary-side ack contract. **Type:** resilience / reliability. **Stage:** hardening. Originated from the delivery-asymmetry review during CB-401 integration.
Author
Owner

Decision (2026-07-18): bridged does NOT own message persistence — delegate to an external broker

Design direction from the owner: the bridge must not carry the durable-store role. Persistence,
ack, redelivery, and the reminder/backoff loop are delegated to an external broker we already
operate (RabbitMQ / equivalent)
, not implemented inside bridged.

This overrides the earlier "embedded SQLite/Chronicle queue" option — an embedded store would make
the bus a datastore owner, which contradicts its identity (bridged owns transport / rendezvous /
lifecycle / routing; not durability).

Resulting shape — bridged as a stateless broker adapter

  • bridged stays soft-state. On bridge_send / bridge_reply it publishes to a broker
    queue; it consumes + routes + acks. It holds no PEL, no queue, no dedup store of its own.
  • The broker provides the CB-307 semantics natively:
    • at-least-once + manual ack (basic.ack / basic.nack+requeue),
    • redelivery + DLQ for poison messages,
    • the "remind" loop = dead-letter-TTL / delayed-message pattern (no hand-rolled reminder timer),
    • persistence across a daemon restart (the exact loss a java -jar bounce causes today).
  • Broker candidates: RabbitMQ (owner already runs it; native ack/DLQ/delayed-retry — best fit).
    NATS JetStream / Redis Streams are equivalents if the deployment changes; not Kafka for a
    single-host bus.

What the broker does NOT change

The last hop bridged → primary is still the MCP asymmetry — the primary is an MCP client;
neither the broker nor bridged can call into it. The broker guarantees the worker→primary message is
durable and redelivered-until-acked, but surfacing it to an idle primary still needs the
channel decision (inject into the primary's herdr pane, or the primary drains on next poll). Net
effect: that hop becomes lossless + idempotent, not best-effort. CB-306 (spawn-readiness) is
unaffected.

Scope change for this ticket: implement bridged as a broker client/adapter (publish /
consume / ack against RabbitMQ), NOT as a queue implementer. Keep the broker choice behind a thin
port so the transport stays swappable.

## Decision (2026-07-18): bridged does NOT own message persistence — delegate to an external broker Design direction from the owner: **the bridge must not carry the durable-store role.** Persistence, ack, redelivery, and the reminder/backoff loop are delegated to an **external broker we already operate (RabbitMQ / equivalent)**, not implemented inside `bridged`. This overrides the earlier "embedded SQLite/Chronicle queue" option — an embedded store would make the bus a datastore owner, which contradicts its identity (bridged owns transport / rendezvous / lifecycle / routing; not durability). ### Resulting shape — bridged as a stateless broker adapter - `bridged` stays **soft-state**. On `bridge_send` / `bridge_reply` it **publishes** to a broker queue; it **consumes + routes + acks**. It holds no PEL, no queue, no dedup store of its own. - The broker provides the CB-307 semantics natively: - at-least-once + **manual ack** (`basic.ack` / `basic.nack`+requeue), - **redelivery + DLQ** for poison messages, - the **"remind" loop** = dead-letter-TTL / delayed-message pattern (no hand-rolled reminder timer), - **persistence across a daemon restart** (the exact loss a `java -jar` bounce causes today). - Broker candidates: **RabbitMQ** (owner already runs it; native ack/DLQ/delayed-retry — best fit). NATS JetStream / Redis Streams are equivalents if the deployment changes; **not** Kafka for a single-host bus. ### What the broker does NOT change The last hop **`bridged → primary` is still the MCP asymmetry** — the primary is an MCP *client*; neither the broker nor bridged can call into it. The broker guarantees the worker→primary message is **durable and redelivered-until-acked**, but *surfacing* it to an idle primary still needs the channel decision (inject into the primary's herdr pane, or the primary drains on next poll). Net effect: that hop becomes **lossless + idempotent**, not best-effort. CB-306 (spawn-readiness) is unaffected. **Scope change for this ticket:** implement `bridged` as a **broker client/adapter** (publish / consume / ack against RabbitMQ), NOT as a queue implementer. Keep the broker choice behind a thin port so the transport stays swappable.
Author
Owner

Broker decision (2026-07-18): build to a thin AMQP port; default = LavinMQ (RabbitMQ too heavy)

RabbitMQ satisfies the semantics but drags an Erlang/BEAM runtime — too heavy for a single-host bus.
Delegated a research sweep; outcome:

Build bridged to a thin AMQP port and make the broker a deploy-time choice.

  • Default — LavinMQ (Apache 2.0, single Crystal binary). Speaks AMQP 0-9-1, so
    com.rabbitmq:amqp-client works unchanged — only the connection URI differs. ~25 MB RAM for
    100M enqueued msgs. Native dead-letter exchange + native delayed-message exchange
    (x-delayed-type): the CB-307 "remind unacked with backoff" loop is a broker feature here, not
    hand-rolled — and where stock RabbitMQ needs a plugin for delayed-message, LavinMQ ships it.
  • Interchangeable — RabbitMQ: same AMQP client, URI-only swap, for any org already running it.
  • Runner-up — NATS JetStream: even lighter (~20 MB Go binary), native request/reply that maps
    onto bridge_ask, file-backed durability. Costs a new adapter (own protocol, not AMQP) and DLQ
    is a build-it-yourself pattern. Revisit if/when we go true multi-host.
Broker Footprint ack/nack redeliver+DLQ delayed-retry durable Java client AMQP drop-in
LavinMQ 1 Crystal bin, ~25MB ✓ ✓ native DLX ✓ native ✓ disk com.rabbitmq:amqp-client (unchanged) ✓
NATS JetStream 1 Go bin, ~20MB ✓ pattern ✓ native backoff ✓ file io.nats:jnats ✗
Redis Streams redis + AOF ✓ pattern, no DLQ ✗ (approx) AOF-tuned redis.clients:jedis ✗
ActiveMQ Artemis 2nd JVM ✓ ✓ ✓ ✓ Qpid JMS AMQP 1.0 only

Ruled out: MQTT (no NACK/DLX/delay — pub/sub, not a work queue) · Kafka/Redpanda (offset-commit,
no per-message ack; Redpanda is BSL-licensed) · Redis Streams (no native delay; durability rides on
AOF fsync tuning, default ≤1s loss) · Memphis/Superstream (unmaintained as of 2026) · KubeMQ
community (deliberately refuses to start under k8s to force the paid tier).

Net: bridged stays a soft-state AMQP client/adapter; LavinMQ gives the light single-host
default without changing the code path RabbitMQ would use.

## Broker decision (2026-07-18): build to a thin AMQP port; default = **LavinMQ** (RabbitMQ too heavy) RabbitMQ satisfies the semantics but drags an Erlang/BEAM runtime — too heavy for a single-host bus. Delegated a research sweep; outcome: **Build `bridged` to a thin AMQP port and make the broker a deploy-time choice.** - **Default — LavinMQ** (Apache 2.0, single Crystal binary). Speaks AMQP 0-9-1, so `com.rabbitmq:amqp-client` works **unchanged** — only the connection URI differs. ~25 MB RAM for 100M enqueued msgs. Native **dead-letter exchange** + native **delayed-message exchange** (`x-delayed-type`): the CB-307 "remind unacked with backoff" loop is a broker feature here, not hand-rolled — and where stock RabbitMQ needs a *plugin* for delayed-message, LavinMQ ships it. - **Interchangeable — RabbitMQ**: same AMQP client, URI-only swap, for any org already running it. - **Runner-up — NATS JetStream**: even lighter (~20 MB Go binary), native request/reply that maps onto `bridge_ask`, file-backed durability. Costs a *new* adapter (own protocol, not AMQP) and DLQ is a build-it-yourself pattern. Revisit if/when we go true multi-host. | Broker | Footprint | ack/nack | redeliver+DLQ | delayed-retry | durable | Java client | AMQP drop-in | |---|---|---|---|---|---|---|---| | **LavinMQ** | 1 Crystal bin, ~25MB | ✓ | ✓ native DLX | ✓ native | ✓ disk | `com.rabbitmq:amqp-client` (unchanged) | ✓ | | NATS JetStream | 1 Go bin, ~20MB | ✓ | pattern | ✓ native backoff | ✓ file | `io.nats:jnats` | ✗ | | Redis Streams | redis + AOF | ✓ | pattern, no DLQ | ✗ (approx) | AOF-tuned | `redis.clients:jedis` | ✗ | | ActiveMQ Artemis | 2nd JVM | ✓ | ✓ | ✓ | ✓ | Qpid JMS | AMQP 1.0 only | **Ruled out:** MQTT (no NACK/DLX/delay — pub/sub, not a work queue) · Kafka/Redpanda (offset-commit, no per-message ack; Redpanda is BSL-licensed) · Redis Streams (no native delay; durability rides on AOF fsync tuning, default ≤1s loss) · Memphis/Superstream (unmaintained as of 2026) · KubeMQ community (deliberately refuses to start under k8s to force the paid tier). **Net:** `bridged` stays a soft-state AMQP **client/adapter**; LavinMQ gives the light single-host default without changing the code path RabbitMQ would use.
Author
Owner

Implementation staging (2026-07-18): two stages behind one port

Grounded a delivery-path map in the current code before delegating. The reverse (worker→primary) path is Rendezvous — a bare ConcurrentHashMap<session, CompletableFuture<Resolution>> of live blocking waiters only. There is no queue and no store: when a bridge_reply arrives and no send is open, Rendezvous.complete() (msg/Rendezvous.java:212-215) hits waiters.get(session)==null → returns false → the reply string is silently discarded (worker gets error("no send is awaiting a reply") / REST 409 no_pending_send). No message-id, dedup, or ack concept exists anywhere in the message path today.

To keep the broker-adapter identity but not block on standing up broker infra, split delivery behind one thin port and ship in two stages:

Stage 1 — reply-inbox port + in-memory (soft-state) adapter — delegate now, no infra

  • Define a ReplyInbox port: publish(target, msgId, content) / drain(target) / ack(target, msgId), dedup by msgId.
  • InMemoryReplyInbox default adapter (per-session in-process queue). This is soft-state, not persistence — lost on a java -jar bounce, exactly consistent with "bridged stays soft-state." It does NOT make the bus a datastore owner.
  • Wire it at the exact drop seam: in Rendezvous.complete(), replace the silent return false (no live waiter) with inbox.publish(session, msgId, content) — the reply is held, not dropped. Consume side plugs into MessageService.send/answer/poll (msg/MessageService.java:180-182,267,315): on open, drain a queued reply for the target so a newly-arriving send/poll catches up; ack by msgId on hand-back. CompletionResolver fallbacks route through the same inbox.
  • Fully gate-able without a broker (unit tests only). Fixes the observed "communication break" (stranded replies) on a single host today, and freezes the exact port the AMQP adapter later implements.

Stage 2 — AmqpReplyInbox adapter (LavinMQ default) behind the same port — deferred until a broker is reachable

  • Same ReplyInbox interface, backed by AMQP (com.rabbitmq:amqp-client, URI-only swap LavinMQ↔RabbitMQ per the broker decision above). Adds durability across restart, native DLX, native delayed-message = the remind/backoff loop.
  • Config broker: block absent → in-memory (Stage 1); present → AMQP. Needs a live LavinMQ to integration-test the primary hand-back, so it waits on infra.
  • Name channels multi-host-ready (per-agent routing keys, roster.* topic) so CB-308 (multi-host federation) builds on this without repainting the topology.

Net: Stage 1 stops the loss on one host with zero infra and is primary-gate-verifiable; Stage 2 swaps the adapter for cross-restart/cross-host durability once a broker is deployed. The last hop bridged → primary stays a pull (MCP asymmetry) in both stages — the inbox just makes it lossless + idempotent instead of best-effort.

## Implementation staging (2026-07-18): two stages behind one port Grounded a delivery-path map in the current code before delegating. The reverse (worker→primary) path is `Rendezvous` — a bare `ConcurrentHashMap<session, CompletableFuture<Resolution>>` of **live blocking waiters only**. There is **no queue and no store**: when a `bridge_reply` arrives and no send is open, `Rendezvous.complete()` (`msg/Rendezvous.java:212-215`) hits `waiters.get(session)==null` → returns `false` → **the reply string is silently discarded** (worker gets `error("no send is awaiting a reply")` / REST 409 `no_pending_send`). No message-id, dedup, or ack concept exists anywhere in the message path today. To keep the broker-adapter identity but not block on standing up broker infra, split delivery behind **one thin port** and ship in two stages: ### Stage 1 — reply-inbox port + in-memory (soft-state) adapter — **delegate now, no infra** - Define a `ReplyInbox` port: `publish(target, msgId, content)` / `drain(target)` / `ack(target, msgId)`, dedup by `msgId`. - **`InMemoryReplyInbox`** default adapter (per-session in-process queue). This is **soft-state, not persistence** — lost on a `java -jar` bounce, exactly consistent with "bridged stays soft-state." It does NOT make the bus a datastore owner. - Wire it at the exact drop seam: in `Rendezvous.complete()`, replace the silent `return false` (no live waiter) with `inbox.publish(session, msgId, content)` — the reply is **held, not dropped**. Consume side plugs into `MessageService.send`/`answer`/`poll` (`msg/MessageService.java:180-182,267,315`): on open, drain a queued reply for the target so a newly-arriving send/poll catches up; ack by `msgId` on hand-back. `CompletionResolver` fallbacks route through the same inbox. - **Fully gate-able without a broker** (unit tests only). Fixes the observed "communication break" (stranded replies) on a single host today, and freezes the exact port the AMQP adapter later implements. ### Stage 2 — `AmqpReplyInbox` adapter (LavinMQ default) behind the same port — **deferred until a broker is reachable** - Same `ReplyInbox` interface, backed by AMQP (`com.rabbitmq:amqp-client`, URI-only swap LavinMQ↔RabbitMQ per the broker decision above). Adds durability across restart, native DLX, native delayed-message = the remind/backoff loop. - Config `broker:` block absent → in-memory (Stage 1); present → AMQP. Needs a live LavinMQ to integration-test the primary hand-back, so it waits on infra. - Name channels **multi-host-ready** (per-agent routing keys, `roster.*` topic) so CB-308 (multi-host federation) builds on this without repainting the topology. Net: Stage 1 stops the loss on one host with zero infra and is primary-gate-verifiable; Stage 2 swaps the adapter for cross-restart/cross-host durability once a broker is deployed. The last hop `bridged → primary` stays a pull (MCP asymmetry) in both stages — the inbox just makes it lossless + idempotent instead of best-effort.
Author
Owner

Stage 1 SHIPPED — merged to main @ ba6b4a5 (2026-07-18)

The reply-inbox port + in-memory adapter is in. A worker bridge_reply that arrives with no open send is now held and drainable instead of silently discarded.

What landed:

  • ReplyInbox port (publish/peek/ack, dedup by msgId) + InboxMessage record + InMemoryReplyInbox adapter (per-target FIFO via LinkedHashMap, thread-safe, soft-state — not persistence).
  • MessageService.reply(session, content) resolves an open send or publishes to the inbox with a minted UUID (the old silent-drop path); drainReplies(target) = peek + ack.
  • BridgeMcp.reply / BridgedApp.replyMessage repointed off bare Rendezvous.resolve → no-waiter is now success (queued), not error / 409 no_pending_send.
  • Drain surface: bridge_poll gains optional target; REST GET /sessions/{id}/replies.
  • Rendezvous untouched. QUESTION path (bridge_ask) is NOT queued (interactive → keeps NO_WAITER); completion/failure fallbacks are NOT queued (captured-waiter, double-delivery risk) — both guarded by tests.

Gate: mvn clean install BUILD SUCCESS, 207 tests, 0 failures / 0 errors (baseline 188 → +19, incl. a new InMemoryReplyInboxTest). Implemented by an off-sub gx10 worker in a pre-trusted worktree; primary-verified and committed by the primary because the worker's own completion replies were lost to the very bug this fixes (a fitting confirmation of the problem). .mcp.json and wiki/ untouched.

Not live yet — the running daemon is still on the pre-CB-307 jar; a restart makes it active.

Remaining: Stage 2 (deferred)

AmqpReplyInbox (LavinMQ default) behind the same ReplyInbox port + broker: config (absent → in-memory, present → AMQP) for cross-restart / cross-host durability. Needs a live broker to integration-test → waits on infra. This ticket stays open for Stage 2.

## Stage 1 SHIPPED — merged to `main` @ `ba6b4a5` (2026-07-18) The reply-inbox port + in-memory adapter is in. A worker `bridge_reply` that arrives with no open send is now **held and drainable** instead of silently discarded. **What landed:** - `ReplyInbox` port (`publish`/`peek`/`ack`, dedup by `msgId`) + `InboxMessage` record + `InMemoryReplyInbox` adapter (per-target FIFO via `LinkedHashMap`, thread-safe, **soft-state — not persistence**). - `MessageService.reply(session, content)` resolves an open send or publishes to the inbox with a minted UUID (the old silent-drop path); `drainReplies(target)` = peek + ack. - `BridgeMcp.reply` / `BridgedApp.replyMessage` repointed off bare `Rendezvous.resolve` → no-waiter is now **success (queued)**, not `error` / `409 no_pending_send`. - Drain surface: `bridge_poll` gains optional `target`; REST `GET /sessions/{id}/replies`. - `Rendezvous` untouched. **QUESTION path (`bridge_ask`) is NOT queued** (interactive → keeps `NO_WAITER`); **completion/failure fallbacks are NOT queued** (captured-waiter, double-delivery risk) — both guarded by tests. **Gate:** `mvn clean install` BUILD SUCCESS, **207 tests, 0 failures / 0 errors** (baseline 188 → +19, incl. a new `InMemoryReplyInboxTest`). Implemented by an off-sub gx10 worker in a pre-trusted worktree; **primary-verified and committed by the primary** because the worker's own completion replies were lost to the very bug this fixes (a fitting confirmation of the problem). `.mcp.json` and `wiki/` untouched. **Not live yet** — the running daemon is still on the pre-CB-307 jar; a restart makes it active. ### Remaining: Stage 2 (deferred) `AmqpReplyInbox` (LavinMQ default) behind the **same** `ReplyInbox` port + `broker:` config (absent → in-memory, present → AMQP) for cross-restart / cross-host durability. Needs a live broker to integration-test → waits on infra. This ticket stays open for Stage 2.
Author
Owner

Stage 2 shipped — AmqpReplyInbox (durable, cross-restart) @ 2bc5f3a on main

Stage 1 (ba6b4a5) added the ReplyInbox port + InMemoryReplyInbox so a stranded worker reply is held instead of dropped, drained by the primary via bridge_poll(target) / GET /sessions/{id}/replies. Stage 2 puts genuine durability behind the same port.

Adapter — consume-and-hold + deferred manual ack. Each target owns a durable queue agent.<target>.inbox. A manual-ack consumer pulls persistent messages into an in-memory held map (dedup by msgId) but does not ack; peek returns the snapshot; ack acks the broker delivery-tag and drops it. A crash/java -jar bounce before the primary drains leaves messages unacked → the broker redelivers on reconnect. The port contract is preserved (idempotent publish, FIFO peek, at-least-once ack); bridged still owns no persistence — the broker does.

Selection. A broker: block with a uri in bridged.yaml swaps the in-memory inbox for the AMQP one; absent, bridged stays soft-state. Production default is LavinMQ; a stock RabbitMQ speaks the same AMQP 0-9-1 (URI-only swap).

Verified.

  • Default mvn clean install hermetic, 210 green (no Docker; the AMQP test is @Tag("contract"), excluded).
  • Contract test against a real RabbitMQ (Testcontainers), 3 green: publish/peek/ack, msgId dedup, and cross-restart redelivery (unacked reply survives closing the inbox, redelivered to a fresh connection, gone once acked).
  • IDE diagnostics 0/0 on all new/edited files; CVE re-check clean (pinned test-scope commons-compress 1.27.1 + commons-lang3 3.18.0 that testcontainers pulled).

Still open (why this doesn't close #5). This is the durable landing zone for replies; the primary still pulls (poll/REST). The active push-to-primary + ACK + bounded reminder loop from the proposal remains blocked by the same MCP asymmetry (server can't call into the client) — the broker fabric is now the natural substrate for it, and for the CB-308 multi-host federation. Leaving open for that push/reminder layer.

Files: msg/AmqpReplyInbox, config/BridgedConfig (Broker(uri) record), Bridged.main (adapter select + ordered-shutdown close), bridged.example.yaml (broker: doc), BridgedConfigTest + AmqpReplyInboxContractTest.

## Stage 2 shipped — `AmqpReplyInbox` (durable, cross-restart) @ `2bc5f3a` on `main` Stage 1 (`ba6b4a5`) added the `ReplyInbox` port + `InMemoryReplyInbox` so a stranded worker reply is **held** instead of dropped, drained by the primary via `bridge_poll(target)` / `GET /sessions/{id}/replies`. Stage 2 puts genuine durability behind the same port. **Adapter — consume-and-hold + deferred manual ack.** Each target owns a durable queue `agent.<target>.inbox`. A manual-ack consumer pulls persistent messages into an in-memory held map (dedup by `msgId`) but does **not** ack; `peek` returns the snapshot; `ack` acks the broker delivery-tag and drops it. A crash/`java -jar` bounce before the primary drains leaves messages unacked → the broker redelivers on reconnect. The port contract is preserved (idempotent publish, FIFO peek, at-least-once ack); **bridged still owns no persistence — the broker does.** **Selection.** A `broker:` block with a `uri` in `bridged.yaml` swaps the in-memory inbox for the AMQP one; absent, bridged stays soft-state. Production default is LavinMQ; a stock RabbitMQ speaks the same AMQP 0-9-1 (URI-only swap). **Verified.** - Default `mvn clean install` hermetic, **210 green** (no Docker; the AMQP test is `@Tag("contract")`, excluded). - Contract test against a real RabbitMQ (Testcontainers), **3 green**: publish/peek/ack, `msgId` dedup, and **cross-restart redelivery** (unacked reply survives closing the inbox, redelivered to a fresh connection, gone once acked). - IDE diagnostics 0/0 on all new/edited files; CVE re-check clean (pinned test-scope `commons-compress 1.27.1` + `commons-lang3 3.18.0` that testcontainers pulled). **Still open (why this doesn't close #5).** This is the *durable landing zone* for replies; the primary still **pulls** (poll/REST). The active **push-to-primary + ACK + bounded reminder loop** from the proposal remains blocked by the same MCP asymmetry (server can't call into the client) — the broker fabric is now the natural substrate for it, and for the [CB-308](../issues/6) multi-host federation. Leaving open for that push/reminder layer. Files: `msg/AmqpReplyInbox`, `config/BridgedConfig` (`Broker(uri)` record), `Bridged.main` (adapter select + ordered-shutdown close), `bridged.example.yaml` (`broker:` doc), `BridgedConfigTest` + `AmqpReplyInboxContractTest`.
Author
Owner

Delivered + live-dogfooded — closing

All three increments are on main (worker delivery 131e7b1..7c252b5, primary-gate cleanup d4c9704) and were dogfooded live on the running daemon (PID 93639, d4c9704, AMQP durable inbox on LavinMQ).

The push channel — design question resolved

The crux question ("if the on-subscription primary runs in a herdr pane the daemon can see, injector push is symmetric with worker delivery") is answered yes, live. The primary's terminal_id is already resolved on every MCP call by ConnectionIdentity/PaneLocator; we just capture it. Push = a dedicated status-gated loop (ReplyPushLoop, mechanism (b) — not the worker Injector, to stay decoupled from WorkerPresence) that injects a drain nudge into the primary's own pane via AgentControl.send.

Increments (all shipped)

  1. Learn the primary — single-slot PrimaryRegistry, populated from orchestration-side tools (bridge_send/bridge_spawn) when the caller resolves to a non-null terminal that is not a registered worker session. Optional primary.terminal config pin; degrades safely to pull if the primary is off-host (registry empty → loop is a no-op, reply never lost).
  2. Push + remind — on the no-waiter reply branch (MessageService.reply), onReplyQueued starts a bounded loop: inject a nudge only when AgentControl.status(primary).injectable(), re-nudge on backoff up to a cap (primary.push_reminders=5, primary.push_backoff_ms=15000ms). Ack = drain: stop condition is inbox.peek(target).isEmpty().
  3. bridge_ack — per-msgId ack tool over ReplyInbox.ack, for finer control than drain-all.

Live trace (this session, primary term_656c8cc03e1f0b1)

primary terminal learned: term_656c8cc03e1f0b1          ← Inc 1
send to <worker> timed out (delivered=true)             ← 2s blocking send closes waiter
push: starting reminder loop for <worker>               ← worker reply, NO waiter → inbox → Inc 2
push: primary ... is WORKING (not injectable), waiting   ← status-gated: never interrupts a turn
push: already active ... ignoring duplicate trigger      ← idempotent (5 eager replies, one loop)
push: nudge 1/5 sent to primary term_656c8cc03e1f0b1     ← INJECT fired the instant the primary went idle
                                                            (arrived as a real turn in the primary's pane)
drain → PUSHLOOP-DOGFOOD-OK                               ← ack = drain
push: inbox empty ... reminder loop ended                ← STOP; exactly 1 nudge, no spam

Every designed behavior confirmed against real infrastructure: auto-learn, no-waiter trigger, status-gating (WAIT_BUSY while the primary was mid-turn), idempotency, INJECT-on-idle, and ack=drain → bounded STOP. Gate: IDE diagnostics 0/0 on all changed files, mvn clean install green (242 tests, 0F/0E).

Boundary preserved: the injection is a nudge (not the payload), status-gated so it never interrupts a turn, bounded so it never spams, carries no env, never crosses the subscription boundary. First time the bridge writes into the primary's pane — it signals "you have mail", it does not drive the primary's work.

Closing as done. Multi-host (CB-308) remains out of scope — the loop is local-only; a remote primary is reached by its own local gateway.

## Delivered + live-dogfooded — closing All three increments are on `main` (worker delivery `131e7b1..7c252b5`, primary-gate cleanup `d4c9704`) and were **dogfooded live on the running daemon** (PID 93639, `d4c9704`, AMQP durable inbox on LavinMQ). ### The push channel — design question resolved The crux question ("if the on-subscription primary runs in a herdr pane the daemon can see, injector push is symmetric with worker delivery") is answered **yes, live**. The primary's `terminal_id` is already resolved on every MCP call by `ConnectionIdentity`/`PaneLocator`; we just capture it. Push = a **dedicated status-gated loop** (`ReplyPushLoop`, mechanism (b) — *not* the worker `Injector`, to stay decoupled from `WorkerPresence`) that injects a **drain nudge** into the primary's own pane via `AgentControl.send`. ### Increments (all shipped) 1. **Learn the primary** — single-slot `PrimaryRegistry`, populated from orchestration-side tools (`bridge_send`/`bridge_spawn`) when the caller resolves to a non-null terminal that is **not** a registered worker session. Optional `primary.terminal` config pin; degrades safely to pull if the primary is off-host (registry empty → loop is a no-op, reply never lost). 2. **Push + remind** — on the no-waiter reply branch (`MessageService.reply`), `onReplyQueued` starts a bounded loop: inject a nudge only when `AgentControl.status(primary).injectable()`, re-nudge on backoff up to a cap (`primary.push_reminders`=5, `primary.push_backoff_ms`=15000ms). **Ack = drain**: stop condition is `inbox.peek(target).isEmpty()`. 3. **`bridge_ack`** — per-`msgId` ack tool over `ReplyInbox.ack`, for finer control than drain-all. ### Live trace (this session, primary `term_656c8cc03e1f0b1`) ``` primary terminal learned: term_656c8cc03e1f0b1 ← Inc 1 send to <worker> timed out (delivered=true) ← 2s blocking send closes waiter push: starting reminder loop for <worker> ← worker reply, NO waiter → inbox → Inc 2 push: primary ... is WORKING (not injectable), waiting ← status-gated: never interrupts a turn push: already active ... ignoring duplicate trigger ← idempotent (5 eager replies, one loop) push: nudge 1/5 sent to primary term_656c8cc03e1f0b1 ← INJECT fired the instant the primary went idle (arrived as a real turn in the primary's pane) drain → PUSHLOOP-DOGFOOD-OK ← ack = drain push: inbox empty ... reminder loop ended ← STOP; exactly 1 nudge, no spam ``` Every designed behavior confirmed against real infrastructure: auto-learn, no-waiter trigger, status-gating (WAIT_BUSY while the primary was mid-turn), idempotency, INJECT-on-idle, and ack=drain → bounded STOP. Gate: IDE diagnostics 0/0 on all changed files, `mvn clean install` green (242 tests, 0F/0E). Boundary preserved: the injection is a **nudge** (not the payload), **status-gated** so it never interrupts a turn, **bounded** so it never spams, carries **no env**, never crosses the subscription boundary. First time the bridge writes into the primary's pane — it signals "you have mail", it does not drive the primary's work. Closing as done. Multi-host (CB-308) remains out of scope — the loop is local-only; a remote primary is reached by its own local gateway.
ltms closed this issue 2026-07-19 17:54:04 +02:00
Sign in to join this conversation.
1 Participants
Notifications
Due Date
No due date set.
Dependencies

No dependencies set.

Reference: fleet/fleetd#5