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
No Branch/Tag Specified
main
worker/fleetd-612-unita-87807e-1
worker/612-b3-mcpwirings-da2b58-3
worker/612-b2-cb185-176d3a-2
worker/612-b1-completion-457459-1
worker/612-agaps-73a926-2
worker/608-sleeps-3a64ff-3
worker/621-b4520b-1
worker/618-b83894-2
worker/fleetd-615-e05481-5
worker/lead-autocompact-5f1ab2-3
worker/fleetd-613-f85deb-3
worker/fleetd-608-flaky-nudge-test-d0c2d1-3
worker/lead-context-gauge-ad404f-1
worker/gauge-wiring-9158c1-4
worker/redeploy-slowstart-ead0e5-5
worker/charter-bytes-13668c-6
worker/rollover-outcome-291483-2
worker/589-f64303-2
worker/593-1a8025-5
worker/589-fcd2aa-1
worker/568-9fdaa2-3
worker/571-attempted-outcome-5739f7-2
worker/581-completionresolver-cas-sites-0542b7-6
worker/562-loop-health-wiring-test-99611c-5
worker/562-surface-loop-health-7df5cc-4
worker/575-waiter-cleanup-sites-62ad80-1
worker/572-answer-lock-release-46a9ae-5
worker/567-probe-channel-leak-a38fc5-6
worker/551-record-before-send-7cbf56-1
worker/561-listener-fanout-survives-a-throw-61d538-2
worker/555-redeploy-main-flow-seam-65c2f5-2
worker/556-injector-owns-registration-e027a5-1
worker/552-post-restart-mktemp-abort-bc2672-4
worker/553-onstatus-completion-leak-0da881-2
worker/550-shasum-linux-196132-1
worker/538-loop-dies-on-error-4a5eeb-6
worker/426-health-coverage-ef1fd4-4
worker/504-failed-reported-clean-3cfd66-3
worker/537-capturedlog-close-e4c437-2
worker/459-broken-link-targets-cadc17-5
worker/535-appender-leak-fe74c1-1
worker/512-part2-shutdown-detection-434701-9
worker/529-logger-level-sweep-2a5533-8
worker/528-drain-gate-call-site-5de83d-7
charter/forge-mcp-vs-token
worker/521-swap-guard-unpinned-28e931-5
worker/519-probe-test-harness-d25ab8-4
worker/525-logger-level-leak-1b4eb0-6
worker/518-fleetmcp-resolver-wiring-8ef96c-1
worker/512-drain-complete-line-7edd71-3
worker/517-abort-branch-and-jar-id-41b641-2
worker/500-9e52c9-3
worker/509-4912f4-2
worker/511-9a4b23-1
worker/493-479f45-2
worker/505-03f8b2-1
worker/492-followup-detect-unclear
worker/501-a31fa0-7
worker/498-451d1c-5
worker/494-1015ce-2
worker/492-209647-1
worker/489-001902-2
worker/480-relative-handover-path-906323-1
worker/480-b-handover-skill-45bf1f-5
worker/474-followup-source-pin-f54a55-17
worker/474-charter-check-on-reload-f54a55-17
worker/466-quarantine-repeatcount-report
worker/393-opencode-skill-seeding-71854b-13
worker/469-canonical-tool-names-2a472a-16
worker/466-quarantine-escalation-5ae9c1-15
worker/446-hot-exhausted-pattern-0af580-6
worker/464-charter-tool-name-guard-a85635-12
worker/463-listfleet-default-fails-open-f1c76c-11
worker/458-invariant-5-by-purpose-862f9a-10
worker/439-coordinator-row-gate-bc032a-8
worker/449-herdr-protocol-576015-4
worker/450-abstract-spawn-599e1c-5
worker/437-ack-refuses-177d91-1
worker/444-placement-window-feb56a-2
worker/440-helddurable-derived-d462d7-13
worker/425-rework-placement-resolve-c58ba1-9
worker/421-lead-peek-held-msgs-cdbad2-10
worker/435-fixed-policy-cap-fe11de-12
worker/422-gate-state-observability-9e79d6-11
worker/431-memberregistry-live-readers-cdbad2-10
worker/424-architect-slot-hot-038b41-7
worker/422-model-gate-spawn-c29f48-6
worker/425-default-profile-live-f55534-8
worker/415-coverage-wording-2cbf9c-5
worker/416-3ad1da-1
worker/418-588283-3
worker/deterministic-stamp-race-409-3cb7b6-10
worker/armed-reads-live-config-404-ed931f-9
worker/reply-peer-refusal-391-5a34bd-7
worker/models-allowlist-aa9e9b-3
worker/ttl-stamp-race-399-f1122f-8
worker/scrub-receipt-400-316b3e-5
worker/exhaustion-detection-395-105105-6
worker/scrub-abort-394-316b3e-5
fix/scrub-uid-abort
worker/task-scrub-517574-2
worker/t386-clock-bd5b78-4
worker/t384-scrub-813790-5
worker/t381-cc-748314-2
worker/t373-336973-2
worker/t365-3920c5-3
worker/t358-6e989b-1
worker/t355-8b321c-1
worker/fleetd-369-hermetic-git-tests-e8b19a-3
worker/fleetd-368-stale-lead-binding-f5682e-2
worker/fleetd-360-deploy-units-0d3793-1
worker/359-dead-lead-tabs-f1253b-4
worker/362-worktree-skills-c03e51-3
worker/361-coord-visibility-655144-1
362-plugin-visibility-and-drift
worker/errscan-bed2ca-2
worker/amqp-log-identity-bed2ca-2
worker/withdefaults-guard-561704
worker/sleepguard-82076d-1
worker/fd334-9ee1b6-5
worker/fd348-f1ab27-4
worker/fd335-a71c35-1
worker/fd342-174a17-2
worker/fd345-490d0f-3
worker/fleetd-337-5ec7d4-21
worker/fleetd-341-af5a6b-24
worker/fleetd-339-5ca0a2-23
worker/fleetd-338-83a4a1-22
worker/fleetd-333-281f46-18
worker/fleetd-329-11bdbb-16
worker/fleetd-330-2770fb-17
worker/fix-326-50506e-15
worker/fix-324-3e9bbf-14
worker/fix-323-b8287d-13
worker/fix-316b-bd0860-11
worker/fix-318-76ca36-9
worker/fix-317-486aec-8
worker/fix-315-ce47c5-6
worker/fix-307-275890-6
worker/fix-308-b4f664-7
worker/fix-309-ec3939-8
worker/fix-310-7a3974-9
worker/fix-302-52ad0e-9
worker/fix-298-ce1acb-8
worker/fix-297-66bd11-7
worker/fix-296-104622-6
worker/fix-293-bare-closetab-eb22b5-3
worker/fix-280-gone-ask-lapse-bca98e-2
worker/fix-290-reapidle-guard-coverage-9b0dd1-1
worker/fix-285-trust-seed-8f3565-10
worker/fix-284-backend-error-seat-85912c-11
worker/fix-282-chained-ask-e6d0bb-8
worker/fix-283-teardown-leaks-f40dfa-9
worker/fix-281-pin-handler-actions-4921ac-7
worker/audit-rendezvous-lifecycle-d072ae-2
worker/audit-health-placement-1a2476-6
worker/audit-teardown-exits-e207a5-3
worker/audit-launcher-asymmetry-27e370-4
worker/audit-rest-authz-6ca53c-5
worker/investigate-275-abandon-asking-fdef52-8
worker/fix-274-worktree-leak-b0095d-7
worker/fix-273-exhausted-pattern-9665b5-6
worker/fleetd-267-model-check-bd8068-1
worker/fleetd-131-archunit-18b834-7
worker/fleetd-266-sshagent-rename-a014ff-6
worker/fleetd-184-uid-claim-8e1f31-4
worker/fleetd-184-warn-b381ee-10
worker/fleetd-184-docs-be1d12-9
worker/fleetd-257-9bf010-7
worker/fleetd-103-23a113-6
worker/fleetd-247-342356-5
worker/fleetd-116-04dea8-4
worker/fleetd-252-a830e0-3
worker/fleetd-111-7e8673-9
worker/fleetd-155c-f8ef4b-8
worker/fleetd-176-b928ca-3
worker/fleetd-249-7a7878-2
worker/cb248-composition-root-b-9acdf7-15
worker/cb148-envrc-default-fa6c82-12
worker/cb201-unit5-wiring-6c12e6-8
worker/cb241-fallback-echo-1175e9-11
worker/cb149-trust-dialog-2392a5-9
worker/cb134-148-overlay-visible-c9b986-10
worker/cb234-session-id-keyed-04e1fc-1
worker/cb201-unit3-nudge-abdf5c-6
worker/cb201-unit2-policy-c1102c-5
worker/cb201-unit4-outcome-a13bfa-7
worker/cb201-unit1-classifier-91b9b1-4
worker/cb201-227-refine-831980-3
worker/cb175-model-readback-0f085f-1
worker/cb222-charter-tmpdir-17f013-1
worker/cb226-architect-slot-race-cd3aa8-3
worker/cb224-worktree-root-group-024523-2
worker/cb-123-role-demotion-c600f7-2
worker/cb-219-opencode-roots-1f677e-1
worker/cb214-claude-session-id-b9eab4-4
worker/cb213-zdotdir-wrong-process-dd6de4-3
worker/cb211-exhaustion-classification-9546e0-2
worker/cb137-ambiguous-task-4df3d8-4
worker/cb209-agentsessionid-4dfdb6-2
worker/cb185-hostenvnames-2692b5-3
worker/cb206-opencode-sqlite-128718-2
worker/cb185-worktree-group-fc0c99-1
worker/cb-137-ask-ticket-e7760c-2
worker/cb-172-broker-uri-d36ae4-4
worker/cb-175-model-readback-76ead6-3
worker/cb-161-pane-ancestry-293510-1
worker/cb-164-rebase-885863-8
worker/cb-164-empty-scrape-false-success-1a80af-3
fix/cb-197-ticket-ttl-from-completion
worker/cb-189-remote-url-coverage-4692f3-1
worker/cb-185-blockers-027756-4
worker/cb-192-gap-log-11b631-2
worker/cb-633-fix-5f4396-3
worker/cb185-router-d6436d-3
worker/cb185-router-routing-gaps-9e9d33-3
worker/cb185-paneids-992586-2
worker/cb-633-allow-list-union-ed374b-1
worker/cb-157-credential-in-remote-url-496e44-2
worker/cb-641-health-herdr-evidence-8f1f54-6
worker/cb-640-health-msg-evidence-99c9cd-1
worker/cb-642-fleets-status-skill-bbbc40-5
cb-634-ide-mcp
worker/lead-comms-wiring-c014b9-7
worker/lead-mailbox-c19577-6
worker/autocompact-window-82bc2f-5
worker/cb-634-probe-18056f-4
worker/cb635-broker-urienv
worker/cb-632-config-retry-8e0efa-7
lead/cb-622e-claude-md
lead/cb-622-followup
worker/cb-622a-165dff-1
lead/cb-622d-opencode-mount
worker/cb-622b-717c67-2
worker/cb-622c-ab7759-3
worker/cb-617b2-20ca4b-3
worker/cb-617a-5c2f4a-1
worker/cb596-4e49ef-3
worker/cb586-10500c-1
worker/cb-606-b9343a-25
worker/cb604-1445f8-24
worker/cb582-477374-21
worker/cb584-8c2281-22
worker/cb600-e6b9a9-20
worker/cb602-ce257f-19
worker/cb601-b42837-18
worker/cb598-6c7ba7-17
worker/cb599-740fe4-16
worker/cb597-282224-15
worker/cb590fix-185e9a-10
worker/cb528-recovery-race
worker/cb594-96bead-8
worker/cb590-916766-2
worker/cb527-997d99-3
worker/cb592-env-leak-3cbf9c-1
worker/cb588-async-ticket-nudge-3218f7-5
worker/cb578b-9dcb13-6
worker/cb581-d24826-5
worker/m2-u5-ef8c42-15
worker/cb578a-516499-2
worker/cb576-01a04b-17
worker/cb579-lead-tab-acba06-20
worker/cb580-terminal-health-ed6058-21
worker/cb577-f36fdc-18
worker/cb573b-3db06f-16
worker/cb568c-f36fdc-18
worker/cb568-drop-cause-c3ac1c
worker/cb575-cancelled-notification-c3ac1c
worker/m4-sol-a2cbec-3
worker/cb574-async-ask-c3ac1c
worker/cb573-health-model-8ca857-14
worker/cb572-unknown-target-7f2e35-13
worker/u4-700706-9
worker/u3-b9fcb6-6
worker/u2-ef5b68-4
worker/u1-469dce-1-clean
worker/u1-469dce-1
worker/cb-564-health-events-70cf7e-2
worker/cb-565-recycle-drops-role-98e58f-3
worker/cb-563-missing-reply-df2866-1
worker/cb-562-readiness-gate-silent-6c23c9-3
worker/cb-560-architect-presence-da8155-1
worker/cb-561-architect-silent-off-a71cab-2
worker/cb-548-bind-architect-slot-fe1b8c-1
worker/parity-overlay-settings-5fb711-1
secrets-central-store
cb-559-hot-key-correction
cb-557-fleet-role-pools
worker/cb-553-maxload-explicit-spawn-305ee3-6
worker/cb-551-idle-lead-heartbeat-f1633c-1
worker/cb-544-drain-preserves-worktree-925fad-3
worker/cb-552-docs-sync-1cb9cf-4
worker/cb-548-rendezvous-guard-rebased
worker/cb-548-rendezvous-guard-116b53-10
worker/cb-548-authz-v2-586df6-8
worker/cb-548-authz-264363-5
salvage/cb-528b-codex-home
salvage/cb-528a-codex-launcher
CB-518-primary-flow
feature/peer-launcher-spi
cb-103-injector
v1.1.0
v1.0.0
Labels
Clear labels
blocked
needs-live-proof
ready-to-delegate
silent-default
Cannot start until something else lands. The body says what.
Merged and green, but never shown working on the running daemon. Not the same as done.
Scope, files and acceptance criteria are written. A worker can be briefed from the body alone.
A feature that compiles, passes tests, and ships turned off. Nine recurrences and counting.
No Label
Milestone
No items
No Milestone
Projects
Clear projects
No project
Notifications
Due Date
No due date set.
Dependencies
No dependencies set.
Reference: fleet/fleetd#5
Reference in New Issue
Block a user
Blocking a user prevents them from interacting with repositories, such as opening or commenting on pull requests or issues. Learn more about blocking a user.
Delete Branch "%!s()"
Deleting a branch is permanent. Although the deleted branch may continue to exist for a short time before it actually gets removed, it CANNOT be undone in most cases. Continue?
Problem — the delivery asymmetry
Delivery in the bridge is not symmetric:
MessageServicecallsInjector.enqueue(target, content), andthe injector delivers into the worker's herdr pane when the worker is idle/blocked. The bridge
actively drives the message into the peer.
an open blocking
bridge_sendfor 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 intomain.
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:
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.
on next MCP call carrying the message id). Unacked messages stay pending.
(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
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.)
fire-and-poll; poll becomes an optimization, not the only delivery path.
Relationship to CB-306
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.
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
bridgedstays soft-state. Onbridge_send/bridge_replyit publishes to a brokerqueue; it consumes + routes + acks. It holds no PEL, no queue, no dedup store of its own.
basic.ack/basic.nack+requeue),java -jarbounce causes today).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 → primaryis 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
bridgedas 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.
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
bridgedto a thin AMQP port and make the broker a deploy-time choice.com.rabbitmq:amqp-clientworks unchanged — only the connection URI differs. ~25 MB RAM for100M 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, nothand-rolled — and where stock RabbitMQ needs a plugin for delayed-message, LavinMQ ships it.
onto
bridge_ask, file-backed durability. Costs a new adapter (own protocol, not AMQP) and DLQis a build-it-yourself pattern. Revisit if/when we go true multi-host.
com.rabbitmq:amqp-client(unchanged)io.nats:jnatsredis.clients:jedisRuled 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:
bridgedstays a soft-state AMQP client/adapter; LavinMQ gives the light single-hostdefault without changing the code path RabbitMQ would use.
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 bareConcurrentHashMap<session, CompletableFuture<Resolution>>of live blocking waiters only. There is no queue and no store: when abridge_replyarrives and no send is open,Rendezvous.complete()(msg/Rendezvous.java:212-215) hitswaiters.get(session)==null→ returnsfalse→ the reply string is silently discarded (worker getserror("no send is awaiting a reply")/ REST 409no_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
ReplyInboxport:publish(target, msgId, content)/drain(target)/ack(target, msgId), dedup bymsgId.InMemoryReplyInboxdefault adapter (per-session in-process queue). This is soft-state, not persistence — lost on ajava -jarbounce, exactly consistent with "bridged stays soft-state." It does NOT make the bus a datastore owner.Rendezvous.complete(), replace the silentreturn false(no live waiter) withinbox.publish(session, msgId, content)— the reply is held, not dropped. Consume side plugs intoMessageService.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 bymsgIdon hand-back.CompletionResolverfallbacks route through the same inbox.Stage 2 —
AmqpReplyInboxadapter (LavinMQ default) behind the same port — deferred until a broker is reachableReplyInboxinterface, 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.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.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 → primarystays a pull (MCP asymmetry) in both stages — the inbox just makes it lossless + idempotent instead of best-effort.Stage 1 SHIPPED — merged to
main@ba6b4a5(2026-07-18)The reply-inbox port + in-memory adapter is in. A worker
bridge_replythat arrives with no open send is now held and drainable instead of silently discarded.What landed:
ReplyInboxport (publish/peek/ack, dedup bymsgId) +InboxMessagerecord +InMemoryReplyInboxadapter (per-target FIFO viaLinkedHashMap, 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.replyMessagerepointed off bareRendezvous.resolve→ no-waiter is now success (queued), noterror/409 no_pending_send.bridge_pollgains optionaltarget; RESTGET /sessions/{id}/replies.Rendezvousuntouched. QUESTION path (bridge_ask) is NOT queued (interactive → keepsNO_WAITER); completion/failure fallbacks are NOT queued (captured-waiter, double-delivery risk) — both guarded by tests.Gate:
mvn clean installBUILD SUCCESS, 207 tests, 0 failures / 0 errors (baseline 188 → +19, incl. a newInMemoryReplyInboxTest). 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.jsonandwiki/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 sameReplyInboxport +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 2 shipped —
AmqpReplyInbox(durable, cross-restart) @2bc5f3aonmainStage 1 (
ba6b4a5) added theReplyInboxport +InMemoryReplyInboxso a stranded worker reply is held instead of dropped, drained by the primary viabridge_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 bymsgId) but does not ack;peekreturns the snapshot;ackacks the broker delivery-tag and drops it. A crash/java -jarbounce 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 auriinbridged.yamlswaps 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.
mvn clean installhermetic, 210 green (no Docker; the AMQP test is@Tag("contract"), excluded).msgIddedup, and cross-restart redelivery (unacked reply survives closing the inbox, redelivered to a fresh connection, gone once acked).commons-compress 1.27.1+commons-lang3 3.18.0that 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.Delivered + live-dogfooded — closing
All three increments are on
main(worker delivery131e7b1..7c252b5, primary-gate cleanupd4c9704) 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_idis already resolved on every MCP call byConnectionIdentity/PaneLocator; we just capture it. Push = a dedicated status-gated loop (ReplyPushLoop, mechanism (b) — not the workerInjector, to stay decoupled fromWorkerPresence) that injects a drain nudge into the primary's own pane viaAgentControl.send.Increments (all shipped)
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. Optionalprimary.terminalconfig pin; degrades safely to pull if the primary is off-host (registry empty → loop is a no-op, reply never lost).MessageService.reply),onReplyQueuedstarts a bounded loop: inject a nudge only whenAgentControl.status(primary).injectable(), re-nudge on backoff up to a cap (primary.push_reminders=5,primary.push_backoff_ms=15000ms). Ack = drain: stop condition isinbox.peek(target).isEmpty().bridge_ack— per-msgIdack tool overReplyInbox.ack, for finer control than drain-all.Live trace (this session, primary
term_656c8cc03e1f0b1)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 installgreen (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.