#298: AmqpReplyInbox.release() requeues held deliveries instead of dropping them #301
Closed
agent
wants to merge 0 commits from
worker/fix-298-ce1acb-8 into main
pull from: worker/fix-298-ce1acb-8
merge into: fleet:main
fleet:main
fleet:worker/fleetd-612-unita-87807e-1
fleet:worker/612-b3-mcpwirings-da2b58-3
fleet:worker/612-b2-cb185-176d3a-2
fleet:worker/612-b1-completion-457459-1
fleet:worker/612-agaps-73a926-2
fleet:worker/608-sleeps-3a64ff-3
fleet:worker/621-b4520b-1
fleet:worker/618-b83894-2
fleet:worker/fleetd-615-e05481-5
fleet:worker/lead-autocompact-5f1ab2-3
fleet:worker/fleetd-613-f85deb-3
fleet:worker/fleetd-608-flaky-nudge-test-d0c2d1-3
fleet:worker/lead-context-gauge-ad404f-1
fleet:worker/gauge-wiring-9158c1-4
fleet:worker/redeploy-slowstart-ead0e5-5
fleet:worker/charter-bytes-13668c-6
fleet:worker/rollover-outcome-291483-2
fleet:worker/589-f64303-2
fleet:worker/593-1a8025-5
fleet:worker/589-fcd2aa-1
fleet:worker/568-9fdaa2-3
fleet:worker/571-attempted-outcome-5739f7-2
fleet:worker/581-completionresolver-cas-sites-0542b7-6
fleet:worker/562-loop-health-wiring-test-99611c-5
fleet:worker/562-surface-loop-health-7df5cc-4
fleet:worker/575-waiter-cleanup-sites-62ad80-1
fleet:worker/572-answer-lock-release-46a9ae-5
fleet:worker/567-probe-channel-leak-a38fc5-6
fleet:worker/551-record-before-send-7cbf56-1
fleet:worker/561-listener-fanout-survives-a-throw-61d538-2
fleet:worker/555-redeploy-main-flow-seam-65c2f5-2
fleet:worker/556-injector-owns-registration-e027a5-1
fleet:worker/552-post-restart-mktemp-abort-bc2672-4
fleet:worker/553-onstatus-completion-leak-0da881-2
fleet:worker/550-shasum-linux-196132-1
fleet:worker/538-loop-dies-on-error-4a5eeb-6
fleet:worker/426-health-coverage-ef1fd4-4
fleet:worker/504-failed-reported-clean-3cfd66-3
fleet:worker/537-capturedlog-close-e4c437-2
fleet:worker/459-broken-link-targets-cadc17-5
fleet:worker/535-appender-leak-fe74c1-1
fleet:worker/512-part2-shutdown-detection-434701-9
fleet:worker/529-logger-level-sweep-2a5533-8
fleet:worker/528-drain-gate-call-site-5de83d-7
fleet:charter/forge-mcp-vs-token
fleet:worker/521-swap-guard-unpinned-28e931-5
fleet:worker/519-probe-test-harness-d25ab8-4
fleet:worker/525-logger-level-leak-1b4eb0-6
fleet:worker/518-fleetmcp-resolver-wiring-8ef96c-1
fleet:worker/512-drain-complete-line-7edd71-3
fleet:worker/517-abort-branch-and-jar-id-41b641-2
fleet:worker/500-9e52c9-3
fleet:worker/509-4912f4-2
fleet:worker/511-9a4b23-1
fleet:worker/493-479f45-2
fleet:worker/505-03f8b2-1
fleet:worker/492-followup-detect-unclear
fleet:worker/501-a31fa0-7
fleet:worker/498-451d1c-5
fleet:worker/494-1015ce-2
fleet:worker/492-209647-1
fleet:worker/489-001902-2
fleet:worker/480-relative-handover-path-906323-1
fleet:worker/480-b-handover-skill-45bf1f-5
fleet:worker/474-followup-source-pin-f54a55-17
fleet:worker/474-charter-check-on-reload-f54a55-17
fleet:worker/466-quarantine-repeatcount-report
fleet:worker/393-opencode-skill-seeding-71854b-13
fleet:worker/469-canonical-tool-names-2a472a-16
fleet:worker/466-quarantine-escalation-5ae9c1-15
fleet:worker/446-hot-exhausted-pattern-0af580-6
fleet:worker/464-charter-tool-name-guard-a85635-12
fleet:worker/463-listfleet-default-fails-open-f1c76c-11
fleet:worker/458-invariant-5-by-purpose-862f9a-10
fleet:worker/439-coordinator-row-gate-bc032a-8
fleet:worker/449-herdr-protocol-576015-4
fleet:worker/450-abstract-spawn-599e1c-5
fleet:worker/437-ack-refuses-177d91-1
fleet:worker/444-placement-window-feb56a-2
fleet:worker/440-helddurable-derived-d462d7-13
fleet:worker/425-rework-placement-resolve-c58ba1-9
fleet:worker/421-lead-peek-held-msgs-cdbad2-10
fleet:worker/435-fixed-policy-cap-fe11de-12
fleet:worker/422-gate-state-observability-9e79d6-11
fleet:worker/431-memberregistry-live-readers-cdbad2-10
fleet:worker/424-architect-slot-hot-038b41-7
fleet:worker/422-model-gate-spawn-c29f48-6
fleet:worker/425-default-profile-live-f55534-8
fleet:worker/415-coverage-wording-2cbf9c-5
fleet:worker/416-3ad1da-1
fleet:worker/418-588283-3
fleet:worker/deterministic-stamp-race-409-3cb7b6-10
fleet:worker/armed-reads-live-config-404-ed931f-9
fleet:worker/reply-peer-refusal-391-5a34bd-7
fleet:worker/models-allowlist-aa9e9b-3
fleet:worker/ttl-stamp-race-399-f1122f-8
fleet:worker/scrub-receipt-400-316b3e-5
fleet:worker/exhaustion-detection-395-105105-6
fleet:worker/scrub-abort-394-316b3e-5
fleet:fix/scrub-uid-abort
fleet:worker/task-scrub-517574-2
fleet:worker/t386-clock-bd5b78-4
fleet:worker/t384-scrub-813790-5
fleet:worker/t381-cc-748314-2
fleet:worker/t373-336973-2
fleet:worker/t365-3920c5-3
fleet:worker/t358-6e989b-1
fleet:worker/t355-8b321c-1
fleet:worker/fleetd-369-hermetic-git-tests-e8b19a-3
fleet:worker/fleetd-368-stale-lead-binding-f5682e-2
fleet:worker/fleetd-360-deploy-units-0d3793-1
fleet:worker/359-dead-lead-tabs-f1253b-4
fleet:worker/362-worktree-skills-c03e51-3
fleet:worker/361-coord-visibility-655144-1
fleet:362-plugin-visibility-and-drift
fleet:worker/errscan-bed2ca-2
fleet:worker/amqp-log-identity-bed2ca-2
fleet:worker/withdefaults-guard-561704
fleet:worker/sleepguard-82076d-1
fleet:worker/fd334-9ee1b6-5
fleet:worker/fd348-f1ab27-4
fleet:worker/fd335-a71c35-1
fleet:worker/fd342-174a17-2
fleet:worker/fd345-490d0f-3
fleet:worker/fleetd-337-5ec7d4-21
fleet:worker/fleetd-341-af5a6b-24
fleet:worker/fleetd-339-5ca0a2-23
fleet:worker/fleetd-338-83a4a1-22
fleet:worker/fleetd-333-281f46-18
fleet:worker/fleetd-329-11bdbb-16
fleet:worker/fleetd-330-2770fb-17
fleet:worker/fix-326-50506e-15
fleet:worker/fix-324-3e9bbf-14
fleet:worker/fix-323-b8287d-13
fleet:worker/fix-316b-bd0860-11
fleet:worker/fix-318-76ca36-9
fleet:worker/fix-317-486aec-8
fleet:worker/fix-315-ce47c5-6
fleet:worker/fix-307-275890-6
fleet:worker/fix-308-b4f664-7
fleet:worker/fix-309-ec3939-8
fleet:worker/fix-310-7a3974-9
fleet:worker/fix-302-52ad0e-9
fleet:worker/fix-297-66bd11-7
fleet:worker/fix-296-104622-6
fleet:worker/fix-293-bare-closetab-eb22b5-3
fleet:worker/fix-280-gone-ask-lapse-bca98e-2
fleet:worker/fix-290-reapidle-guard-coverage-9b0dd1-1
fleet:worker/fix-285-trust-seed-8f3565-10
fleet:worker/fix-284-backend-error-seat-85912c-11
fleet:worker/fix-282-chained-ask-e6d0bb-8
fleet:worker/fix-283-teardown-leaks-f40dfa-9
fleet:worker/fix-281-pin-handler-actions-4921ac-7
fleet:worker/audit-rendezvous-lifecycle-d072ae-2
fleet:worker/audit-health-placement-1a2476-6
fleet:worker/audit-teardown-exits-e207a5-3
fleet:worker/audit-launcher-asymmetry-27e370-4
fleet:worker/audit-rest-authz-6ca53c-5
fleet:worker/investigate-275-abandon-asking-fdef52-8
fleet:worker/fix-274-worktree-leak-b0095d-7
fleet:worker/fix-273-exhausted-pattern-9665b5-6
fleet:worker/fleetd-267-model-check-bd8068-1
fleet:worker/fleetd-131-archunit-18b834-7
fleet:worker/fleetd-266-sshagent-rename-a014ff-6
fleet:worker/fleetd-184-uid-claim-8e1f31-4
fleet:worker/fleetd-184-warn-b381ee-10
fleet:worker/fleetd-184-docs-be1d12-9
fleet:worker/fleetd-257-9bf010-7
fleet:worker/fleetd-103-23a113-6
fleet:worker/fleetd-247-342356-5
fleet:worker/fleetd-116-04dea8-4
fleet:worker/fleetd-252-a830e0-3
fleet:worker/fleetd-111-7e8673-9
fleet:worker/fleetd-155c-f8ef4b-8
fleet:worker/fleetd-176-b928ca-3
fleet:worker/fleetd-249-7a7878-2
fleet:worker/cb248-composition-root-b-9acdf7-15
fleet:worker/cb148-envrc-default-fa6c82-12
fleet:worker/cb201-unit5-wiring-6c12e6-8
fleet:worker/cb241-fallback-echo-1175e9-11
fleet:worker/cb149-trust-dialog-2392a5-9
fleet:worker/cb134-148-overlay-visible-c9b986-10
fleet:worker/cb234-session-id-keyed-04e1fc-1
fleet:worker/cb201-unit3-nudge-abdf5c-6
fleet:worker/cb201-unit2-policy-c1102c-5
fleet:worker/cb201-unit4-outcome-a13bfa-7
fleet:worker/cb201-unit1-classifier-91b9b1-4
fleet:worker/cb201-227-refine-831980-3
fleet:worker/cb175-model-readback-0f085f-1
fleet:worker/cb222-charter-tmpdir-17f013-1
fleet:worker/cb226-architect-slot-race-cd3aa8-3
fleet:worker/cb224-worktree-root-group-024523-2
fleet:worker/cb-123-role-demotion-c600f7-2
fleet:worker/cb-219-opencode-roots-1f677e-1
fleet:worker/cb214-claude-session-id-b9eab4-4
fleet:worker/cb213-zdotdir-wrong-process-dd6de4-3
fleet:worker/cb211-exhaustion-classification-9546e0-2
fleet:worker/cb137-ambiguous-task-4df3d8-4
fleet:worker/cb209-agentsessionid-4dfdb6-2
fleet:worker/cb185-hostenvnames-2692b5-3
fleet:worker/cb206-opencode-sqlite-128718-2
fleet:worker/cb185-worktree-group-fc0c99-1
fleet:worker/cb-137-ask-ticket-e7760c-2
fleet:worker/cb-172-broker-uri-d36ae4-4
fleet:worker/cb-175-model-readback-76ead6-3
fleet:worker/cb-161-pane-ancestry-293510-1
fleet:worker/cb-164-rebase-885863-8
fleet:worker/cb-164-empty-scrape-false-success-1a80af-3
fleet:fix/cb-197-ticket-ttl-from-completion
fleet:worker/cb-189-remote-url-coverage-4692f3-1
fleet:worker/cb-185-blockers-027756-4
fleet:worker/cb-192-gap-log-11b631-2
fleet:worker/cb-633-fix-5f4396-3
fleet:worker/cb185-router-d6436d-3
fleet:worker/cb185-router-routing-gaps-9e9d33-3
fleet:worker/cb185-paneids-992586-2
fleet:worker/cb-633-allow-list-union-ed374b-1
fleet:worker/cb-157-credential-in-remote-url-496e44-2
fleet:worker/cb-641-health-herdr-evidence-8f1f54-6
fleet:worker/cb-640-health-msg-evidence-99c9cd-1
fleet:worker/cb-642-fleets-status-skill-bbbc40-5
fleet:cb-634-ide-mcp
fleet:worker/lead-comms-wiring-c014b9-7
fleet:worker/lead-mailbox-c19577-6
fleet:worker/autocompact-window-82bc2f-5
fleet:worker/cb-634-probe-18056f-4
fleet:worker/cb635-broker-urienv
fleet:worker/cb-632-config-retry-8e0efa-7
fleet:lead/cb-622e-claude-md
fleet:lead/cb-622-followup
fleet:worker/cb-622a-165dff-1
fleet:lead/cb-622d-opencode-mount
fleet:worker/cb-622b-717c67-2
fleet:worker/cb-622c-ab7759-3
fleet:worker/cb-617b2-20ca4b-3
fleet:worker/cb-617a-5c2f4a-1
fleet:worker/cb596-4e49ef-3
fleet:worker/cb586-10500c-1
fleet:worker/cb-606-b9343a-25
fleet:worker/cb604-1445f8-24
fleet:worker/cb582-477374-21
fleet:worker/cb584-8c2281-22
fleet:worker/cb600-e6b9a9-20
fleet:worker/cb602-ce257f-19
fleet:worker/cb601-b42837-18
fleet:worker/cb598-6c7ba7-17
fleet:worker/cb599-740fe4-16
fleet:worker/cb597-282224-15
fleet:worker/cb590fix-185e9a-10
fleet:worker/cb528-recovery-race
fleet:worker/cb594-96bead-8
fleet:worker/cb590-916766-2
fleet:worker/cb527-997d99-3
fleet:worker/cb592-env-leak-3cbf9c-1
fleet:worker/cb588-async-ticket-nudge-3218f7-5
fleet:worker/cb578b-9dcb13-6
fleet:worker/cb581-d24826-5
fleet:worker/m2-u5-ef8c42-15
fleet:worker/cb578a-516499-2
fleet:worker/cb576-01a04b-17
fleet:worker/cb579-lead-tab-acba06-20
fleet:worker/cb580-terminal-health-ed6058-21
fleet:worker/cb577-f36fdc-18
fleet:worker/cb573b-3db06f-16
fleet:worker/cb568c-f36fdc-18
fleet:worker/cb568-drop-cause-c3ac1c
fleet:worker/cb575-cancelled-notification-c3ac1c
fleet:worker/m4-sol-a2cbec-3
fleet:worker/cb574-async-ask-c3ac1c
fleet:worker/cb573-health-model-8ca857-14
fleet:worker/cb572-unknown-target-7f2e35-13
fleet:worker/u4-700706-9
fleet:worker/u3-b9fcb6-6
fleet:worker/u2-ef5b68-4
fleet:worker/u1-469dce-1-clean
fleet:worker/u1-469dce-1
fleet:worker/cb-564-health-events-70cf7e-2
fleet:worker/cb-565-recycle-drops-role-98e58f-3
fleet:worker/cb-563-missing-reply-df2866-1
fleet:worker/cb-562-readiness-gate-silent-6c23c9-3
fleet:worker/cb-560-architect-presence-da8155-1
fleet:worker/cb-561-architect-silent-off-a71cab-2
fleet:worker/cb-548-bind-architect-slot-fe1b8c-1
fleet:worker/parity-overlay-settings-5fb711-1
fleet:secrets-central-store
fleet:cb-559-hot-key-correction
fleet:cb-557-fleet-role-pools
fleet:worker/cb-553-maxload-explicit-spawn-305ee3-6
fleet:worker/cb-551-idle-lead-heartbeat-f1633c-1
fleet:worker/cb-544-drain-preserves-worktree-925fad-3
fleet:worker/cb-552-docs-sync-1cb9cf-4
fleet:worker/cb-548-rendezvous-guard-rebased
fleet:worker/cb-548-rendezvous-guard-116b53-10
fleet:worker/cb-548-authz-v2-586df6-8
fleet:worker/cb-548-authz-264363-5
fleet:salvage/cb-528b-codex-home
fleet:salvage/cb-528a-codex-launcher
fleet:CB-518-primary-flow
fleet:feature/peer-launcher-spi
fleet:cb-103-injector
No Reviewers
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#301
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 "worker/fix-298-ce1acb-8"
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?
Fixes fleetd #298:
AmqpReplyInbox.release()drops an unacked reply the broker still holds.The bug
release(target)cancelled the target's AMQP consumer and droppedheld's local record of itsdeliveries, but cancelling a consumer does not requeue the broker's in-flight deliveries — they
stay unacked, attached to the still-open
channel, until the channel or connection actually closes.So a held-but-undrained
fleet_replybecame permanently unreachable: never acked, never nacked,never requeued, and invisible to
peek/fleet_poll.handleRecoveryandclose()were already correct (each makes the broker requeue — a realconnection drop, or
channel.close()— before clearing local state).release()copied their bareheld.clear()/held.remove()pattern without that precondition.The fix
release(target)now:channel.basicCancel(tag)).channel.basicNack(tag, false, true)) every delivery still held fortarget, before dropping the localheldentry.Why nack-with-requeue, and why after the cancel (not before): requeue is the right disposition —
the message was never processed. I initially tried nack-before-cancel (the ticket's suggested order,
"before the cancel and before
held.remove"), but a real broker run (AmqpReplyInboxContractTestagainst RabbitMQ via Testcontainers) showed a race: while the consumer is still active, a
nack-with-requeue is delivered straight back to that same consumer the instant a prefetch slot frees
up, racing this method's own
held.remove(target)— the redelivery can land after the clear andleave a stale entry, so
peekwas no longer reliably empty right afterrelease(I saw this faillive:
expected: <true> but was: <false>on the "release clears the local held snapshot"assertion). Cancelling first avoids this: a delivery tag stays valid for
basicNackon an openchannel regardless of whether its consumer is attached, so cancelling first costs nothing and closes
the race.
Failure handling: a failed
basicNackis caught broadly (IOException | RuntimeException,matching the wide-catch house style
HerdrPeerLauncher.stop()established for #293) and logged atWARN naming the target/msgId/delivery-tag that leaked, without aborting
release()— it runs duringteardown (
Fleetdcalls it right afterMessageService.abandon), and a throw here would abortcleanups the caller depends on. The delivery stays unacked on the broker either way (not silently
dropped) — recoverable by a later connection drop even when the immediate requeue attempt fails. A
failed
basicCancelstill throws, unchanged from before this fix.Out of scope, untouched per the ticket:
InMemoryReplyInbox(drop-on-release is correct there),handleRecovery,close(), andMessageService(including the CB-640 flag-clearing at:797).Proof
Docker was available in this worktree (
docker infosucceeded), so I ran the real@Tag("contract")suite against a RabbitMQ container via Testcontainers (
mvn test -Pcontract) rather than relying onreasoning alone.
Extended
AmqpReplyInboxContractTest.releaseCancelsConsumerAndClearsHeld→releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery: keeps the original "peek is empty afterrelease" assertion, and adds a new one — re-
own()the same target afterrelease()and assert thedelivery is redelivered (msgId + content), proving it is genuinely recoverable, not merely absent
from
peek.Mutation proof (revert fix → real failure → restore):
AmqpReplyInbox.javato the pre-fixrelease()(bareheld.remove, no requeue) viacpfrom a saved copy of the pre-edit file (never usedgit stash— it's shared across worktreeson this machine).
mvn test -Pcontract -Dtest=AmqpReplyInboxContractTest#releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery:cpback), reran the same test:Tests run: 1, Failures: 0, Errors: 0, Skipped: 0,BUILD SUCCESS.AmqpReplyInboxContractTestclass with the fix in place:Tests run: 8, Failures: 0, Errors: 0, Skipped: 0,BUILD SUCCESS— all other contract behaviour (publish/peek/ack, dedup,cross-restart durability, prefetch, unroutable publish, no-consumer-on-publish) unaffected.
Full build
mvn clean install(hermetic, unpiped, contract tests excluded by default — no Docker required):Ran twice: once before touching the contract test/mutation work, once again after restoring the fix.
Both green with the identical
Tests run: 1303count.Shape check (per ticket — report only, not fixed)
Scanned
AmqpReplyInbox.javafor the same shape ("local state cleared as if the broker hadrequeued, when nothing made it requeue") outside
release():handleRecovery(held.clear()) — has the precondition (real connection drop already requeuedeverything) — not the same defect, already excluded by the ticket.
close()(channel.close(), no explicitheld.clear()at all) —channel.close()itselfrequeues, and the object is discarded afterward anyway — not the same defect, already excluded.
ack()— removes fromheldonly after a successfulbasicAck; onIOExceptionit restoresthe entry rather than dropping it — opposite direction of the bug, not the same shape.
deliverCallback(duplicate-delivery path) — acks the duplicate's tag before dropping it, doesn'tassume anything was requeued — not the same shape.
No second occurrence found; the one in
release()was the only instance of this shape in the file.Report
worker/fix-298-ce1acb-8/Users/dai.ha/LTMS/.bridged-worktrees/56872d-8(own worktree, confirmed viagit rev-parse --show-toplevelbefore editing)fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java,fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxContractTest.javamvn clean install→Tests run: 1303, Failures: 0, Errors: 0, Skipped: 0/BUILD SUCCESS.mvn test -Pcontract -Dtest=AmqpReplyInboxContractTest(Docker) →Tests run: 8, Failures: 0, Errors: 0, Skipped: 0/BUILD SUCCESS.cancel") because it failed live against a real broker (self-redelivery race — see above); I cancel
first, then nack. Please double-check that reasoning.
AmqpReplyInboxContractTestneeds Docker torun (
@Tag("contract"), excluded from the defaultmvn clean install) — it was available in thisworktree so I ran it directly rather than only reasoning about the fix.
Pull request closed