fleetd #361: close the lead-coordination visibility gap #364
Closed
agent
wants to merge 0 commits from
worker/361-coord-visibility-655144-1 into main
pull from: worker/361-coord-visibility-655144-1
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: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-298-ce1acb-8
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#364
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/361-coord-visibility-655144-1"
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?
Closes #361: lead-to-lead AMQP coordination had a send half with tools and a receive half without.
What changed
(a)
LeadChannel.inspect(coordId)— read your own/a peer's mailbox state without owning itMailboxState(coordId, exists, pending, consumers)record with anabsent(coordId)factory.LeadMailbox.inspectopens a throwaway probe channel per call (never the long-livedchannel/publishChannel), does a passive queue declare, and returnsabsenton anyIOExceptioninstead of throwing. This matters because in AMQP 0-9-1 a passive declare of a missing queue closes the channel it ran on with a 404 — ifinspectreusedpublish's channel, one miss would silently break every future publish on that instance.(b)
FleetConfig.Coordinator.peers: List<String>coordinator: {uriEnv, selfId}, nopeerskey) still parses unchanged — covered bycoordinatorBlockWithNoPeersKeyStillParses.Coordinator—MemberCredentialsneeded a@JsonCreatorto disambiguate multiple constructors for Jackson's record deserializer, and that risk wasn't worth taking here. Updated the ~9 existing call sites to pass the 5th argument instead.(c)
fleet_listreports coordination stateFleetMcp.CoordinationSource(LeadChannel, List<String> peers)record with anone()factory, following the existingOutageSource/QuarantineSource/LeadSeatSource"Source record" idiom rather than adding a 6th positional parameter tolistFleet's overload chain (the formerString selfCoordIdlast-parameter slot is nowCoordinationSource coordination, threaded through the shorter overloads asCoordinationSource.none()).msgId/from/preview only, capped at 80 chars — never the full body), and one row per configured peer with reachability/pending/consumers.fleet_list: every peer probe runs on a bounded virtual-thread pool with a 1.5s timeout (PEER_PROBE_TIMEOUT_MS); a timeout, an exception, or a down broker all degrade toMailboxState.absent()for that peer row rather than propagating.fleet_listitself has no new way to throw or block past that ceiling.(d) Honest
fleet_send{coordId}wordingTests
Hermetic (no broker), all in the default
mvn clean installrun:LeadMailboxTest's hermetic siblings viaFakeLeadChannel(withMailboxsupport added) exercise thefleet_list/fleet_sendpaths for mailbox-exists, zero-consumer, and peer-row shapes (FleetMcpTest,FleetMcpLeadCoordTest).FleetConfigTest:coordinatorBlockWithNoPeersKeyStillParses,coordinatorPeersParsesAndDropsBlankEntries,coordinatorPeersDefaultsToEmptyWhenConstructedWithNull.Contract (
@Tag("contract"), real broker via Testcontainers, excluded from the default build, run explicitly withmvn test -Pcontract) inLeadMailboxTest:inspectReportsAnOwnedMailboxAsExistingWithItsOwnConsumerinspectReportsAMissingMailboxAsAbsentRatherThanThrowinginspectReportsPendingMessagesAndZeroConsumersWhenNobodyIsReadingAnymoreinspectingAMissingMailboxNeverBreaksPublishOnTheSameInstance— the point of the ticket: missesinspect()on a queue that has never existed (the exact 404-closes-the-channel case), then callspublish()on the sameLeadMailboxinstance and proves the message still lands, then proves a secondinspect()(this time of an existing queue) still works too.Build results (this worktree, unpiped)
Default build:
Contract tests (Docker confirmed available via
docker info; run separately, not part of the above):Out of scope (spotted, not investigated further)
FleetMcp.reply()(backingfleet_reply) returns the literal text"delivered"after callingmessages.reply(...), which either resolves a waitingfleet_sendor just queues the reply in an inbox for later drain — the same "queued, not necessarily read" gap this ticket fixes forfleet_send{coordId}.FleetApp.replyMessage(the RESTPOST /sessions/{id}/replyhandler) returns{"delivered": true}unconditionally after the samemessages.reply(...)call, for the same reason.ReplyPushLoop/LeadHeartbeatLoopboth callcountNudge("delivered")after a pane nudge attempt — worth checking whether that only counts "attempted" rather than "the pane actually saw it", but not looked into.None of these three were touched — they're reported here only because the brief asked for other instances of the same shape.
Review round 1 — three findings, all fixed on this branch
Finding 1 (must-fix): the fix reproduced the defect the ticket is about
MailboxState.absent()was returned both for a genuinely-absent mailbox and for "the probe couldnot determine anything" (timeout, unreachable broker, any other declare failure) — the exact
"reported stronger than established" shape #361 exists to fix, one level down. Worse for the
daemon's own mailbox: a failed self-probe rendered as
pending: 0, consumers: 0, indistinguishablefrom a genuinely empty, genuinely unread mailbox.
Fix:
LeadChannel.MailboxStatenow carries aPresenceenum (EXISTS/ABSENT/UNKNOWN)with
exists()/known()accessors.LeadMailbox.inspectclassifies a measured AMQP 404 (anIOExceptionwrapping aShutdownSignalExceptionwhoseChannel.Closereply code is404—confirmed against a real broker, not assumed from the spec text; see
LeadMailboxTest.passiveDeclareOfAMissingQueueThrowsAnIOExceptionWrappingA404ShutdownSignal) asABSENT, and every other declare failure asUNKNOWN.fleet_list'smailbox/peer rows nowrender a
statusof"exists"/"absent"/"unknown", andpending/consumersare includedonly when
status == "exists"— an unresolved probe can no longer render as a measured zero.The
"reachable"boolean on peer rows is gone, replaced by the samestatusfield, for the samereason.
Finding 2: a timed-out probe was never cancelled, and it held a channel
FleetMcp.probe'sfuture.get(timeoutMs)timed out the caller but left the submittedinspect()task running forever on its own virtual thread — which had already opened an AMQPchannel. Against a hung (not down) broker, every
fleet_listcall would orphan one more channeluntil the connection's channel-max (2047 by default) was exhausted, which would break
publish()too — the exact invariant this ticket proves safe against the 404 route, defeated by a different
route.
Fix:
probe()now holds theFutureand callscancel(true)on timeout/failure, so theorphaned task is interrupted rather than abandoned, and returns
MailboxState.unknown()(neverabsent()) on any timeout or exception — consistent with finding 1. Also corrected the wording:PEER_PROBE_POOLis a dedicated virtual-thread-per-task executor, not a bounded pool — virtualthreads make spawning it cheap, not bounded, and the comment now says so.
Test:
FleetMcpTest.aTimedOutProbeInterruptsTheOrphanedTaskRatherThanAbandoningIt— ahermetic
LeadChannelfake that blocks until interrupted, driven through a new package-privateFleetMcp.probe(LeadChannel, String, long timeoutMs)seam with a 100ms timeout, proving theorphaned task is actually interrupted rather than left running. No hung broker needed, per the
review's own suggestion.
Finding 3 (latent):
inspectpromised "never throws" and could throwLeadMailbox.inspectcaught onlyIOException.connection.createChannel()on an already-closedconnection throws
AlreadyClosedException— confirmed against a real broker inLeadMailboxTest.createChannelOnAnAlreadyClosedConnectionThrowsAnUncheckedException— whichextends
ShutdownSignalException, an uncheckedRuntimeException, so it would have escaped thecatch (IOException)and broken the "never throws" contract. Not reachable today (every currentcaller goes through
probe'scatch (Exception)), but the next direct caller would have believedthe contract.
Fix: both catch sites in
inspect(opening the probe channel, and the passive declare itself)now also catch
RuntimeExceptionand reportunknown().Test:
LeadMailboxTest.inspectReportsUnknownRatherThanThrowingWhenTheConnectionIsAlreadyClosed— opens a mailbox, closes it (tearing down the connection), then calls
inspect()on that sameinstance and asserts no exception escapes and the result is
UNKNOWN, not the plausible-lookingbut wrong
ABSENT.Build (this worktree, unpiped, full output read)
Default:
(1389 vs. the earlier 1386 — three new hermetic tests: one tri-state
fleet_listrendering test,one probe-cancellation test, one
sendToLead-stays-quiet-on-unknown test.)Contract (
mvn -Pcontract test -Dtest=LeadMailboxTest, real broker via Testcontainers):(12 vs. the earlier 9 — the two exception-shape-pinning tests plus the closed-connection
inspect()test.)Not touched, per the reviewer's instruction: the three "same shape" observations from the first
report (
fleet_reply's"delivered",FleetApp.replyMessage's{"delivered": true},countNudge("delivered")) — filed separately by the reviewer.Review round 2 — one finding:
isMissingQueue's false branch was unpinnedThe reviewer ran a mutation on
isMissingQueue(made it returntrueunconditionally — restoringthe exact overstatement #361 exists to fix, every declare failure reading as a confirmed absence)
and it stayed green:
mvn clean install→ 1389/1389, 0 failures. No test drove a non-404 declarefailure through
inspect, so the discriminator betweenMailboxState.absent()andMailboxState.unknown()— the whole point of round 1's fix — was unproven.Fix: widened
LeadMailbox.isMissingQueuefromprivateto package-private and addedLeadMailboxIsMissingQueueTest— 5 hermetic tests (no broker needed), covering the true case andall three false shapes the method's own javadoc lists:
ACCESS_REFUSEDon a realChannel.Close),ShutdownSignalExceptionwhose reason is aConnection.Close, not aChannel.Close(evenwhen its reply code happens to also be 404),
IOExceptionwith no cause at all, and (one more, for completeness) with an unrelated cause.Verified against the reviewer's own mutation, locally, before pushing:
Secondary (not blocking), also addressed: added a one-line javadoc note on
inspectbeinghonest about which of its two
catch (RuntimeException e)sites is proven end-to-end against areal broker (
createChannel(), viainspectReportsUnknownRatherThanThrowingWhenTheConnectionIs AlreadyClosed) and which stays purely defensive with no test (the declare-site catch, for aconnection-drops-mid-call race no test drives on purpose).
Build (this worktree, unpiped, full output read)
Default:
Contract (
mvn -Pcontract test -Dtest=LeadMailboxTest, real broker via Testcontainers) — unchangedfrom round 1, re-run to confirm nothing regressed:
Not changed this round, confirmed correct and left alone per the reviewer's instruction: the
Presenceenum,exists()/known(),mailboxView'sEXISTS-gating,Future.cancel(true), thepool-comment correction, and the two round-1 contract tests pinning the real exception shapes.
Lead-to-lead AMQP coordination had a send half with tools and a receive half without. This closes three blind spots: - LeadChannel gains inspect(coordId) -> MailboxState(exists, pending, consumers), implemented in LeadMailbox with a throwaway probe channel (never the long-lived publish/consume channels) so a passive-declare 404 on a missing queue can never take down publish() on the same instance. - FleetConfig.Coordinator gains peers: List<String> (defaults to empty, blank entries dropped) so a daemon can declare which peer coord-ids it expects to reach. - fleet_list reports coordination state via a new CoordinationSource (own coord-id, own mailbox state, held messages as msgId/from/preview only, and one row per configured peer with reachability/pending/ consumers), following the existing OutageSource/QuarantineSource "Source record with none()" idiom instead of growing listFleet's overload chain by another positional parameter. Every peer probe is bounded by a 1.5s timeout on a virtual-thread pool and degrades to absent rather than ever slowing or failing fleet_list. - fleet_send{coordId}'s success text now says "durably confirmed by the broker" instead of "delivered", and warns (while still reporting success) when the target mailbox has zero consumers attached. Tests: hermetic unit tests for exists/absent/zero-consumer/old-config- no-peers-key/new fleet_list shape using FakeLeadChannel, plus a @Tag("contract") LeadMailboxTest.inspectingAMissingMailboxNeverBreaks PublishOnTheSameInstance proving the invariant against a real broker.Reviewed. Changes requested — three findings, sent back to the implementer.
First, verification of the claims, because a worker's build is not evidence on its own. I rebuilt this branch in a clean detached worktree:
Identical to the reported numbers. The report was honest, and
inspectingAMissingMailboxNeverBreaksPublishOnTheSameInstanceis the right test for the stated invariant — it proves the 404-closes-the-channel hazard cannot takepublishdown.1. The fix reproduces the defect this ticket is about
MailboxState.absent()is returned for three different situations: the mailbox genuinely does not exist (established), the probe timed out (nothing established), and the broker was unreachable or the declare failed for any other reason (nothing established). The javadoc says so — "or the look otherwise fails (broker unreachable, timed out)" — andFleetMcp.probe'scatch (Exception)collapses all of it beforepeerViewemits"reachable": false.This issue exists because
fleet_sendsaid "delivered" when it had only established "published". Reporting a definite absence for a fact never established is the same shape, one level down. An operator cannot tell "fleet01 is down" from "my broker is slow", and those need opposite responses.Worse for the daemon's own mailbox:
coordinatorViewemitsself.pending()/self.consumers()with no flag, so a failed self-probe renders aspending: 0, consumers: 0— read as "my coordination loop is dead" when the truth is "I could not check". An unmeasured value must never render as a measured zero.2. A timed-out probe is never cancelled, and it holds a channel
PEER_PROBE_POOL.submit(...).get(timeout)times out but does not cancel. The orphaned task has already calledconnection.createChannel(), so it holds an open AMQP channel until the underlying call returns. Against a hung broker — not a down one, which fails fast — everyfleet_listorphans another channel. Channel-max is finite (default 2047), and exhausting it would breakpublish: the same invariant the PR proves safe against the 404 route, defeated by a different route.Being precise about confidence: the mechanism is certain (there is no
cancel); the consequence is reasoned, not measured. I did not stand up a hung broker.Also, the pool is described as "bounded".
newThreadPerTaskExecutoris unbounded — virtual threads make it cheap, not bounded.3.
inspectpromises "never throws" and can throw (latent)catch (IOException)missesAlreadyClosedException, whichconnection.createChannel()raises on a dead connection. Checked against the jar on the classpath:Not reachable today — every caller goes through
probe'scatch (Exception)— which is why it is ranked third. But the contract says "never throws" and the next direct caller will believe it.What is good and should not be lost in the rework
Coordinator.peersis properly back-compatible:null → List.of(), blanks dropped, so an existing config parses unchanged.fleet_send's zero-consumer warning stays a success, not an error. Correct: the queue is durable and the peer will get it on return.Merged to
mainas6e9e464.Verified on the merge itself, not taken from the report:
I also re-ran my own mutation against the new tests on the merged tree, rather than trusting the worker's paste of it:
The discriminator's false branch is pinned now. Before round 2 the same mutation was green at 1389.
One note for the record: the
"reachable"boolean that round 1 added never reachedmain, so replacing it withstatuscosts no reader anything —grep -rn reachable docs/ wiki/ CLAUDE.mdfinds no reference to it.Pull request closed