AmqpReplyInbox.release() drops an unacked reply that the broker still holds — a worker's report is lost silently #298
Closed
opened 2026-09-04 07:05:57 +02:00 by ltms
·
1 comment
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#298
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?
Found by a read-only hunt over the message layer. I verified every step of the path below myself, including the AMQP semantics the fix turns on.
This is a silent message-loss defect: a worker's real
fleet_replycontent becomes permanently invisible, with no error and no log line.The bug
AmqpReplyInboxhas one sharedChannel—private final Channel channel(:93), created once at:160and used by every target.release(target)cancels that target's consumer and drops its local record of held deliveries:That comment is where the reasoning went wrong. The delivery tags are not stale. Cancelling a consumer does not requeue the messages already delivered to it — in AMQP those stay unacked, attached to the still-open channel, until the channel or connection closes.
releasecloses neither. So the entry is dropped locally while the broker still considers it outstanding: never acked, never nacked, never requeued, and no longer reachable byfleet_pollorpeek.Why the two siblings that look identical are actually safe
Both do a bare
held.clear(), and both are correct, because each has a preconditionreleaselacks:handleRecovery(:176-190) runs after a real connection drop. The broker has already requeued every unacked delivery on that connection, so clearing local state is exactly right.close()(:433-451) callschannel.close()at:437, which requeues them.release()copied the pattern without the precondition that makes it safe. Same shape as #274 / #283 / #293, one layer out: the right action on one branch, the same action on a sibling where it is wrong.InMemoryReplyInbox's release is not a counterexample — for that soft-state adapter, dropping on release is the documented, correct behaviour. Do not "fix" it.The path into the bad state — verified end to end
MessageService.reply()(:414-455) takes afleet_replywith no live waiter and falls through toinbox.publish(...). The message lands inAmqpReplyInbox.held, unacked.fleet_pollto drain it, the lead issues a freshsend()/sendAsync()to the same target.send()(:797) clears the bookkeeping flag:fleet_stop, or the idle reaper.Fleetd.java:597callsmessages.abandon(target, reason, true).abandongates recovery onhasStrandedReply(target)(:628), which is now false, sorecoverStrandedReply/drainRepliesnever runs for the old message.Fleetd.java:598then callsreplyInbox.release(target)— which drops the still-unacked entry.What is left behind: one worker's real report, invisible to
fleet_pollandpeek, never redelivered. No bounded reaper recovers it. It is freed only when the whole AMQP connection actually drops — a broker restart, a network blip, or a daemon bounce. Callingown(target)again for a later worker reusing the same target name does not recover it: the stuck delivery stays tied to the old, cancelled consumer on the still-open channel.Direction of harm: silent loss of the one artefact the whole bridge exists to carry. A lead sees a member that finished and reported nothing, and has no way to tell that from a member that genuinely never replied.
Test coverage
AmqpReplyInboxContractTest.releaseCancelsConsumerAndClearsHeld(:156-167) asserts only thatpeek(target)is empty after release. It never checks whether the broker still holds the message or has requeued it — so it passes on the buggy behaviour and would pass on the fix too. It is asserting the symptom, not the property.unackedReplySurvivesRestartAndIsRedeliveredin the same file does prove redelivery, but it closes the wholeAmqpReplyInbox— i.e. the connection — so it exercises theclose()path, notrelease(). That is why this gap survived.Scope
Make
release(target)return the target's held deliveries to the broker instead of dropping them, so a laterown(target)— or any other consumer — can still get them. Requeue is the right disposition: the message was never processed.Decide the mechanism yourself and justify it in your report. Nack-with-requeue over the held delivery tags before clearing is the obvious candidate; say why you chose what you chose, and what happens if that call itself fails.
Rules
held.remove. Say in your report why your ordering is correct.releasehalf-done.releaseis called during teardown (Fleetd.java:598), right afterabandon. A throw there aborts a teardown that other cleanups depend on. Look at how #293 handled exactly this on the launcher side (HerdrPeerLauncher.stop()) and follow the house style: best-effort, log a WARN naming what leaked, do not mask the caller's real failure. CatchRuntimeException, not a narrower type — the reason is recorded in that method.InMemoryReplyInbox. Dropping on release is correct there.handleRecoveryorclose(). Both are correct.MessageService. In particular do not touch the CB-640 flag-clearing at:797: makingsend()drain the old inbox entry is a different design decision with its own consequences, and it is mine to make, not this ticket's.Prove it
Add a test that fails without the fix. The existing
releaseCancelsConsumerAndClearsHeldasserts the wrong property — extend or replace it so it asserts the message is still recoverable afterrelease, not merely absent frompeek. Then keep apeek-is-empty assertion too, since that behaviour must not regress.Mutation proof required: revert the fix, quote the real failure output, restore it.
Shape check
When done, look for the same shape — local state cleared as if the broker had requeued, when nothing made it requeue — in
AmqpReplyInboxonly. One line each, do not fix any of it.Also checked, and clean
The same hunt looked hard at
ReplyPushLoop,LeadCoordLoop,LeadHeartbeatLoop,LeadMailboxandCompletionResolverand found them all symmetric (deliver-then-ack, or fail-without-ack). It also raised arendezvousdouble-open scenario inMessageService.answer; I checked that and it is a false positive —send()holds the session lock across its entire wait and its innerfinallycloses the waiter before the outerfinallyunlocks, so no window exists where a waiter is registered and the lock is free.answer()takes that same lock first and returnsBUSY.Rendezvous.open's javadoc already states this. Recorded here so it is not re-investigated.Merged as
38c248e. Real merge built green at 1307 tests, and the@Tag("contract")suite green at 8 tests against a live RabbitMQ.No correction needed from me. This is the best-executed unit of the batch.
The worker deviated from the ticket, and it was right to
My ticket suggested nacking before the cancel, on the reasoning that the tags must still be valid. The worker tried that, watched it fail against a real broker, and reversed the order — then said so plainly and asked me to double-check it.
The reasoning is correct: nacking with requeue while the consumer is still attached hands the message straight back to that same consumer as soon as a prefetch slot frees, which races
release()'s ownheld.remove. The redelivery can land after the clear and leave a stale entry, sopeekis no longer reliably empty right afterrelease. Cancelling first closes that consumer, and a delivery tag stays valid forbasicNackon an open channel whether or not its consumer is attached — only a channel or connection close invalidates it. So cancel-then-nack costs nothing and removes the race.My ticket's suggested order was wrong, and only a live broker could have shown that. Reasoning about AMQP from the source would not have caught it. Flagging the deviation for review instead of quietly following the spec — or quietly ignoring it — is exactly right.
To be precise about what I checked: I verified the ordering argument and the tag-validity claim, and the tests pass both ways round on the fixed code. I did not independently reproduce the race, because it is timing-dependent. The evidence for it is the worker's quoted live failure, not my own run.
My mutation — the subtle half-fix, not the missing fix
The worker proved the fix by reverting
release()wholesale. That proves something happens. It does not distinguish the fix from a plausible near-miss, so I mutated for that instead:basicNack(tag, false, true)→basicNack(tag, false, false), i.e. discard instead of requeue.That is the dangerous mutant. The message still leaves the unacked state, so any assertion of the form "it is no longer stuck" passes. Only an assertion that the reply is genuinely recoverable can tell the two apart:
Caught, against the live broker. The rewritten test asserts the property rather than the symptom — which was the whole complaint against the test it replaced.
Two details I checked rather than took on trust
heldis declared at:96as "target → (msgId → held delivery). Per-target map is guarded by synchronizing on itself." The newsynchronized (perTarget)matches it, and the whole method still runs underchannelLock.tag == nullpath did not change behaviour. The old code removed fromheldbefore its early return, so a target with no consumer still got its entry dropped. The new code keeps that and adds the requeue.The failure semantics are right, and they are not symmetric on purpose
A failed
basicNackis caught broadly, logged at WARN naming target, msgId and tag, and does not abortrelease()— correct, sincereleaseruns mid-teardown right afterMessageService.abandon, and #293 settled that argument. A failedbasicCancelstill throws, unchanged. That asymmetry is deliberate and the javadoc explains it: if the cancel failed, the consumer may still be attached, so attempting a best-effort requeue underneath it would reintroduce the very race the ordering exists to avoid.Shape check came back clean, and its reasoning on
ack()is worth keeping: it restores the entry whenbasicAckfails rather than dropping it — the opposite direction, so not this shape.