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
Member

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 it

  • New MailboxState(coordId, exists, pending, consumers) record with an absent(coordId) factory.
  • LeadMailbox.inspect opens a throwaway probe channel per call (never the long-lived channel/publishChannel), does a passive queue declare, and returns absent on any IOException instead 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 — if inspect reused publish's channel, one miss would silently break every future publish on that instance.
  • Candidate mechanism from the ticket (a disposable channel for the passive declare): accepted as designed, with the invariant proved by a real-broker test (below), not just asserted in a comment.

(b) FleetConfig.Coordinator.peers: List<String>

  • Defaults to empty when absent or null; blank entries are dropped, matching the existing compact-constructor style. The live shape (coordinator: {uriEnv, selfId}, no peers key) still parses unchanged — covered by coordinatorBlockWithNoPeersKeyStillParses.
  • Did not add a second constructor overload to Coordinator — MemberCredentials needed a @JsonCreator to 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_list reports coordination state

  • New FleetMcp.CoordinationSource(LeadChannel, List<String> peers) record with a none() factory, following the existing OutageSource/QuarantineSource/LeadSeatSource "Source record" idiom rather than adding a 6th positional parameter to listFleet's overload chain (the former String selfCoordId last-parameter slot is now CoordinationSource coordination, threaded through the shorter overloads as CoordinationSource.none()).
  • Reports: own coord-id, own mailbox state, held messages (msgId/from/preview only, capped at 80 chars — never the full body), and one row per configured peer with reachability/pending/consumers.
  • Design choice — never slow or fail 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 to MailboxState.absent() for that peer row rather than propagating. fleet_list itself has no new way to throw or block past that ceiling.

(d) Honest fleet_send{coordId} wording

  • Success text now says "published ... and durably confirmed by the broker" instead of "delivered" — a publisher confirm proves durably queued, not read.
  • When the target mailbox exists with zero consumers, the (still-successful) result appends a warning that nobody is currently reading it; a missing/unreachable mailbox is unaffected (still probed with the same bounded timeout/degrade behavior as (c)).

Tests

Hermetic (no broker), all in the default mvn clean install run:

  • LeadMailboxTest's hermetic siblings via FakeLeadChannel (withMailbox support added) exercise the fleet_list/fleet_send paths 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 with mvn test -Pcontract) in LeadMailboxTest:

  • inspectReportsAnOwnedMailboxAsExistingWithItsOwnConsumer
  • inspectReportsAMissingMailboxAsAbsentRatherThanThrowing
  • inspectReportsPendingMessagesAndZeroConsumersWhenNobodyIsReadingAnymore
  • inspectingAMissingMailboxNeverBreaksPublishOnTheSameInstance — the point of the ticket: misses inspect() on a queue that has never existed (the exact 404-closes-the-channel case), then calls publish() on the same LeadMailbox instance and proves the message still lands, then proves a second inspect() (this time of an existing queue) still works too.

Build results (this worktree, unpiped)

Default build:

mvn clean install
...
Tests run: 1386, Failures: 0, Errors: 0, Skipped: 0
BUILD SUCCESS

Contract tests (Docker confirmed available via docker info; run separately, not part of the above):

mvn -Pcontract test -Dtest=LeadMailboxTest
...
Tests run: 9, Failures: 0, Errors: 0, Skipped: 0

Out of scope (spotted, not investigated further)

  • FleetMcp.reply() (backing fleet_reply) returns the literal text "delivered" after calling messages.reply(...), which either resolves a waiting fleet_send or just queues the reply in an inbox for later drain — the same "queued, not necessarily read" gap this ticket fixes for fleet_send{coordId}.
  • FleetApp.replyMessage (the REST POST /sessions/{id}/reply handler) returns {"delivered": true} unconditionally after the same messages.reply(...) call, for the same reason.
  • ReplyPushLoop/LeadHeartbeatLoop both call countNudge("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 could
not 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, indistinguishable
from a genuinely empty, genuinely unread mailbox.

Fix: LeadChannel.MailboxState now carries a Presence enum (EXISTS/ABSENT/UNKNOWN)
with exists()/known() accessors. LeadMailbox.inspect classifies a measured AMQP 404 (an
IOException wrapping a ShutdownSignalException whose Channel.Close reply code is 404 —
confirmed against a real broker, not assumed from the spec text; see
LeadMailboxTest.passiveDeclareOfAMissingQueueThrowsAnIOExceptionWrappingA404ShutdownSignal) as
ABSENT, and every other declare failure as UNKNOWN. fleet_list's mailbox/peer rows now
render a status of "exists" / "absent" / "unknown", and pending/consumers are included
only 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 same status field, for the same
reason.

Finding 2: a timed-out probe was never cancelled, and it held a channel

FleetMcp.probe's future.get(timeoutMs) timed out the caller but left the submitted
inspect() task running forever on its own virtual thread — which had already opened an AMQP
channel. Against a hung (not down) broker, every fleet_list call would orphan one more channel
until 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 the Future and calls cancel(true) on timeout/failure, so the
orphaned task is interrupted rather than abandoned, and returns MailboxState.unknown() (never
absent()) on any timeout or exception — consistent with finding 1. Also corrected the wording:
PEER_PROBE_POOL is a dedicated virtual-thread-per-task executor, not a bounded pool — virtual
threads make spawning it cheap, not bounded, and the comment now says so.

Test: FleetMcpTest.aTimedOutProbeInterruptsTheOrphanedTaskRatherThanAbandoningIt — a
hermetic LeadChannel fake that blocks until interrupted, driven through a new package-private
FleetMcp.probe(LeadChannel, String, long timeoutMs) seam with a 100ms timeout, proving the
orphaned task is actually interrupted rather than left running. No hung broker needed, per the
review's own suggestion.

Finding 3 (latent): inspect promised "never throws" and could throw

LeadMailbox.inspect caught only IOException. connection.createChannel() on an already-closed
connection throws AlreadyClosedException — confirmed against a real broker in
LeadMailboxTest.createChannelOnAnAlreadyClosedConnectionThrowsAnUncheckedException — which
extends ShutdownSignalException, an unchecked RuntimeException, so it would have escaped the
catch (IOException) and broken the "never throws" contract. Not reachable today (every current
caller goes through probe's catch (Exception)), but the next direct caller would have believed
the contract.

Fix: both catch sites in inspect (opening the probe channel, and the passive declare itself)
now also catch RuntimeException and report unknown().

Test: LeadMailboxTest.inspectReportsUnknownRatherThanThrowingWhenTheConnectionIsAlreadyClosed
— opens a mailbox, closes it (tearing down the connection), then calls inspect() on that same
instance and asserts no exception escapes and the result is UNKNOWN, not the plausible-looking
but wrong ABSENT.

Build (this worktree, unpiped, full output read)

Default:

mvn clean install
...
Tests run: 1389, Failures: 0, Errors: 0, Skipped: 0
BUILD SUCCESS

(1389 vs. the earlier 1386 — three new hermetic tests: one tri-state fleet_list rendering test,
one probe-cancellation test, one sendToLead-stays-quiet-on-unknown test.)

Contract (mvn -Pcontract test -Dtest=LeadMailboxTest, real broker via Testcontainers):

Tests run: 12, Failures: 0, Errors: 0, Skipped: 0

(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 unpinned

The reviewer ran a mutation on isMissingQueue (made it return true unconditionally — restoring
the 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 declare
failure through inspect, so the discriminator between MailboxState.absent() and
MailboxState.unknown() — the whole point of round 1's fix — was unproven.

Fix: widened LeadMailbox.isMissingQueue from private to package-private and added
LeadMailboxIsMissingQueueTest — 5 hermetic tests (no broker needed), covering the true case and
all three false shapes the method's own javadoc lists:

  • a different reply code (403 ACCESS_REFUSED on a real Channel.Close),
  • a ShutdownSignalException whose reason is a Connection.Close, not a Channel.Close (even
    when its reply code happens to also be 404),
  • an IOException with no cause at all, and (one more, for completeness) with an unrelated cause.

Verified against the reviewer's own mutation, locally, before pushing:

BASELINE  mvn test -Dtest=LeadMailboxIsMissingQueueTest
          Tests run: 5, Failures: 0, Errors: 0, Skipped: 0 — BUILD SUCCESS

MUTATED   isMissingQueue(e) { return true; }   // the reviewer's exact mutation
          mvn test -Dtest=LeadMailboxIsMissingQueueTest
          Tests run: 5, Failures: 4, Errors: 0, Skipped: 0 — BUILD FAILURE

          LeadMailboxIsMissingQueueTest.aDifferentReplyCodeIsNotAMissingQueue:45
            a non-404 reply code must never be read as a confirmed absence — the mailbox's
            real state is unknown ==> expected: <false> but was: <true>
          LeadMailboxIsMissingQueueTest.aShutdownSignalWhoseReasonIsNotAChannelCloseIsNotAMissingQueue:57
            a ShutdownSignalException whose reason is not a Channel.Close must never be read as a
            missing queue, even if its reply code happens to be 404 ==> expected: <false> but was: <true>
          LeadMailboxIsMissingQueueTest.anIOExceptionWithNoCauseAtAllIsNotAMissingQueue:66
            an IOException with no ShutdownSignalException cause must never be read as a confirmed
            absence ==> expected: <false> but was: <true>
          LeadMailboxIsMissingQueueTest.anIOExceptionWithAnUnrelatedCauseIsNotAMissingQueue:74
            a cause that isn't even a ShutdownSignalException must never be read as a confirmed
            absence ==> expected: <false> but was: <true>

REVERTED  back to the real fix; mvn clean install → 1394/1394 (1389 + these 5), BUILD SUCCESS

Secondary (not blocking), also addressed: added a one-line javadoc note on inspect being
honest about which of its two catch (RuntimeException e) sites is proven end-to-end against a
real broker (createChannel(), via inspectReportsUnknownRatherThanThrowingWhenTheConnectionIs AlreadyClosed) and which stays purely defensive with no test (the declare-site catch, for a
connection-drops-mid-call race no test drives on purpose).

Build (this worktree, unpiped, full output read)

Default:

mvn clean install
...
Tests run: 1394, Failures: 0, Errors: 0, Skipped: 0
BUILD SUCCESS

Contract (mvn -Pcontract test -Dtest=LeadMailboxTest, real broker via Testcontainers) — unchanged
from round 1, re-run to confirm nothing regressed:

Tests run: 12, Failures: 0, Errors: 0, Skipped: 0

Not changed this round, confirmed correct and left alone per the reviewer's instruction: the
Presence enum, exists()/known(), mailboxView's EXISTS-gating, Future.cancel(true), the
pool-comment correction, and the two round-1 contract tests pinning the real exception shapes.

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 it** - New `MailboxState(coordId, exists, pending, consumers)` record with an `absent(coordId)` factory. - `LeadMailbox.inspect` opens a **throwaway probe channel** per call (never the long-lived `channel`/`publishChannel`), does a passive queue declare, and returns `absent` on any `IOException` instead 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 — if `inspect` reused `publish`'s channel, one miss would silently break every future publish on that instance. - **Candidate mechanism from the ticket (a disposable channel for the passive declare): accepted as designed**, with the invariant proved by a real-broker test (below), not just asserted in a comment. **(b) `FleetConfig.Coordinator.peers: List<String>`** - Defaults to empty when absent or null; blank entries are dropped, matching the existing compact-constructor style. The live shape (`coordinator: {uriEnv, selfId}`, no `peers` key) still parses unchanged — covered by `coordinatorBlockWithNoPeersKeyStillParses`. - Did **not** add a second constructor overload to `Coordinator` — `MemberCredentials` needed a `@JsonCreator` to 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_list` reports coordination state** - New `FleetMcp.CoordinationSource(LeadChannel, List<String> peers)` record with a `none()` factory, following the existing `OutageSource`/`QuarantineSource`/`LeadSeatSource` "Source record" idiom rather than adding a 6th positional parameter to `listFleet`'s overload chain (the former `String selfCoordId` last-parameter slot is now `CoordinationSource coordination`, threaded through the shorter overloads as `CoordinationSource.none()`). - Reports: own coord-id, own mailbox state, held messages (`msgId`/`from`/**preview only, capped at 80 chars** — never the full body), and one row per configured peer with reachability/pending/consumers. - **Design choice — never slow or fail `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 to `MailboxState.absent()` for that peer row rather than propagating. `fleet_list` itself has no new way to throw or block past that ceiling. **(d) Honest `fleet_send{coordId}` wording** - Success text now says "published ... and durably confirmed by the broker" instead of "delivered" — a publisher confirm proves durably queued, not read. - When the target mailbox exists with zero consumers, the (still-successful) result appends a warning that nobody is currently reading it; a missing/unreachable mailbox is unaffected (still probed with the same bounded timeout/degrade behavior as (c)). ## Tests Hermetic (no broker), all in the default `mvn clean install` run: - `LeadMailboxTest`'s hermetic siblings via `FakeLeadChannel` (`withMailbox` support added) exercise the `fleet_list`/`fleet_send` paths 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 with `mvn test -Pcontract`) in `LeadMailboxTest`: - `inspectReportsAnOwnedMailboxAsExistingWithItsOwnConsumer` - `inspectReportsAMissingMailboxAsAbsentRatherThanThrowing` - `inspectReportsPendingMessagesAndZeroConsumersWhenNobodyIsReadingAnymore` - **`inspectingAMissingMailboxNeverBreaksPublishOnTheSameInstance`** — the point of the ticket: misses `inspect()` on a queue that has never existed (the exact 404-closes-the-channel case), then calls `publish()` on the *same* `LeadMailbox` instance and proves the message still lands, then proves a second `inspect()` (this time of an existing queue) still works too. ### Build results (this worktree, unpiped) Default build: ``` mvn clean install ... Tests run: 1386, Failures: 0, Errors: 0, Skipped: 0 BUILD SUCCESS ``` Contract tests (Docker confirmed available via `docker info`; run separately, not part of the above): ``` mvn -Pcontract test -Dtest=LeadMailboxTest ... Tests run: 9, Failures: 0, Errors: 0, Skipped: 0 ``` ## Out of scope (spotted, not investigated further) - `FleetMcp.reply()` (backing `fleet_reply`) returns the literal text `"delivered"` after calling `messages.reply(...)`, which either resolves a waiting `fleet_send` or just queues the reply in an inbox for later drain — the same "queued, not necessarily read" gap this ticket fixes for `fleet_send{coordId}`. - `FleetApp.replyMessage` (the REST `POST /sessions/{id}/reply` handler) returns `{"delivered": true}` unconditionally after the same `messages.reply(...)` call, for the same reason. - `ReplyPushLoop`/`LeadHeartbeatLoop` both call `countNudge("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 could not 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`, indistinguishable from a genuinely empty, genuinely unread mailbox. **Fix:** `LeadChannel.MailboxState` now carries a `Presence` enum (`EXISTS`/`ABSENT`/`UNKNOWN`) with `exists()`/`known()` accessors. `LeadMailbox.inspect` classifies a **measured** AMQP 404 (an `IOException` wrapping a `ShutdownSignalException` whose `Channel.Close` reply code is `404` — confirmed against a real broker, not assumed from the spec text; see `LeadMailboxTest.passiveDeclareOfAMissingQueueThrowsAnIOExceptionWrappingA404ShutdownSignal`) as `ABSENT`, and every other declare failure as `UNKNOWN`. `fleet_list`'s `mailbox`/peer rows now render a `status` of `"exists"` / `"absent"` / `"unknown"`, and `pending`/`consumers` are included **only** 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 same `status` field, for the same reason. ### Finding 2: a timed-out probe was never cancelled, and it held a channel `FleetMcp.probe`'s `future.get(timeoutMs)` timed out the caller but left the submitted `inspect()` task running forever on its own virtual thread — which had already opened an AMQP channel. Against a *hung* (not down) broker, every `fleet_list` call would orphan one more channel until 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 the `Future` and calls `cancel(true)` on timeout/failure, so the orphaned task is interrupted rather than abandoned, and returns `MailboxState.unknown()` (never `absent()`) on any timeout or exception — consistent with finding 1. Also corrected the wording: `PEER_PROBE_POOL` is a dedicated virtual-thread-per-task executor, not a bounded pool — virtual threads make spawning it cheap, not bounded, and the comment now says so. **Test:** `FleetMcpTest.aTimedOutProbeInterruptsTheOrphanedTaskRatherThanAbandoningIt` — a hermetic `LeadChannel` fake that blocks until interrupted, driven through a new package-private `FleetMcp.probe(LeadChannel, String, long timeoutMs)` seam with a 100ms timeout, proving the orphaned task is actually interrupted rather than left running. No hung broker needed, per the review's own suggestion. ### Finding 3 (latent): `inspect` promised "never throws" and could throw `LeadMailbox.inspect` caught only `IOException`. `connection.createChannel()` on an already-closed connection throws `AlreadyClosedException` — confirmed against a real broker in `LeadMailboxTest.createChannelOnAnAlreadyClosedConnectionThrowsAnUncheckedException` — which extends `ShutdownSignalException`, an unchecked `RuntimeException`, so it would have escaped the `catch (IOException)` and broken the "never throws" contract. Not reachable today (every current caller goes through `probe`'s `catch (Exception)`), but the next direct caller would have believed the contract. **Fix:** both catch sites in `inspect` (opening the probe channel, and the passive declare itself) now also catch `RuntimeException` and report `unknown()`. **Test:** `LeadMailboxTest.inspectReportsUnknownRatherThanThrowingWhenTheConnectionIsAlreadyClosed` — opens a mailbox, closes it (tearing down the connection), then calls `inspect()` on that same instance and asserts no exception escapes and the result is `UNKNOWN`, not the plausible-looking but wrong `ABSENT`. ### Build (this worktree, unpiped, full output read) Default: ``` mvn clean install ... Tests run: 1389, Failures: 0, Errors: 0, Skipped: 0 BUILD SUCCESS ``` (1389 vs. the earlier 1386 — three new hermetic tests: one tri-state `fleet_list` rendering test, one probe-cancellation test, one `sendToLead`-stays-quiet-on-unknown test.) Contract (`mvn -Pcontract test -Dtest=LeadMailboxTest`, real broker via Testcontainers): ``` Tests run: 12, Failures: 0, Errors: 0, Skipped: 0 ``` (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 unpinned The reviewer ran a mutation on `isMissingQueue` (made it return `true` unconditionally — restoring the 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 declare failure through `inspect`, so the discriminator between `MailboxState.absent()` and `MailboxState.unknown()` — the whole point of round 1's fix — was unproven. **Fix:** widened `LeadMailbox.isMissingQueue` from `private` to package-private and added `LeadMailboxIsMissingQueueTest` — 5 hermetic tests (no broker needed), covering the true case and all three false shapes the method's own javadoc lists: - a different reply code (403 `ACCESS_REFUSED` on a real `Channel.Close`), - a `ShutdownSignalException` whose reason is a `Connection.Close`, not a `Channel.Close` (even when its reply code happens to also be 404), - an `IOException` with no cause at all, and (one more, for completeness) with an unrelated cause. **Verified against the reviewer's own mutation, locally, before pushing:** ``` BASELINE mvn test -Dtest=LeadMailboxIsMissingQueueTest Tests run: 5, Failures: 0, Errors: 0, Skipped: 0 — BUILD SUCCESS MUTATED isMissingQueue(e) { return true; } // the reviewer's exact mutation mvn test -Dtest=LeadMailboxIsMissingQueueTest Tests run: 5, Failures: 4, Errors: 0, Skipped: 0 — BUILD FAILURE LeadMailboxIsMissingQueueTest.aDifferentReplyCodeIsNotAMissingQueue:45 a non-404 reply code must never be read as a confirmed absence — the mailbox's real state is unknown ==> expected: <false> but was: <true> LeadMailboxIsMissingQueueTest.aShutdownSignalWhoseReasonIsNotAChannelCloseIsNotAMissingQueue:57 a ShutdownSignalException whose reason is not a Channel.Close must never be read as a missing queue, even if its reply code happens to be 404 ==> expected: <false> but was: <true> LeadMailboxIsMissingQueueTest.anIOExceptionWithNoCauseAtAllIsNotAMissingQueue:66 an IOException with no ShutdownSignalException cause must never be read as a confirmed absence ==> expected: <false> but was: <true> LeadMailboxIsMissingQueueTest.anIOExceptionWithAnUnrelatedCauseIsNotAMissingQueue:74 a cause that isn't even a ShutdownSignalException must never be read as a confirmed absence ==> expected: <false> but was: <true> REVERTED back to the real fix; mvn clean install → 1394/1394 (1389 + these 5), BUILD SUCCESS ``` **Secondary (not blocking), also addressed:** added a one-line javadoc note on `inspect` being honest about which of its two `catch (RuntimeException e)` sites is proven end-to-end against a real broker (`createChannel()`, via `inspectReportsUnknownRatherThanThrowingWhenTheConnectionIs AlreadyClosed`) and which stays purely defensive with no test (the declare-site catch, for a connection-drops-mid-call race no test drives on purpose). ### Build (this worktree, unpiped, full output read) Default: ``` mvn clean install ... Tests run: 1394, Failures: 0, Errors: 0, Skipped: 0 BUILD SUCCESS ``` Contract (`mvn -Pcontract test -Dtest=LeadMailboxTest`, real broker via Testcontainers) — unchanged from round 1, re-run to confirm nothing regressed: ``` Tests run: 12, Failures: 0, Errors: 0, Skipped: 0 ``` Not changed this round, confirmed correct and left alone per the reviewer's instruction: the `Presence` enum, `exists()`/`known()`, `mailboxView`'s `EXISTS`-gating, `Future.cancel(true)`, the pool-comment correction, and the two round-1 contract tests pinning the real exception shapes.
agent added 1 commit 2026-09-05 07:45:18 +02:00
fleetd #361: close the lead-coordination visibility gap
CI / contract (pull_request) Successful in 1m10s
CI / build (pull_request) Successful in 1m34s
0df34f3220
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.
Owner

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:

Tests run: 1386, Failures: 0, Errors: 0, Skipped: 0
BUILD SUCCESS

Identical to the reported numbers. The report was honest, and inspectingAMissingMailboxNeverBreaksPublishOnTheSameInstance is the right test for the stated invariant — it proves the 404-closes-the-channel hazard cannot take publish down.

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)" — and FleetMcp.probe's catch (Exception) collapses all of it before peerView emits "reachable": false.

This issue exists because fleet_send said "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: coordinatorView emits self.pending() / self.consumers() with no flag, so a failed self-probe renders as pending: 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 called connection.createChannel(), so it holds an open AMQP channel until the underlying call returns. Against a hung broker — not a down one, which fails fast — every fleet_list orphans another channel. Channel-max is finite (default 2047), and exhausting it would break publish: 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". newThreadPerTaskExecutor is unbounded — virtual threads make it cheap, not bounded.

3. inspect promises "never throws" and can throw (latent)

catch (IOException) misses AlreadyClosedException, which connection.createChannel() raises on a dead connection. Checked against the jar on the classpath:

$ javap -cp . com.rabbitmq.client.AlreadyClosedException     # amqp-client-5.32.0.jar
public class com.rabbitmq.client.AlreadyClosedException extends com.rabbitmq.client.ShutdownSignalException
public class com.rabbitmq.client.ShutdownSignalException extends java.lang.RuntimeException

Not reachable today — every caller goes through probe's catch (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

  • The throwaway-probe-channel mechanism is right, and proved live rather than reasoned.
  • Coordinator.peers is properly back-compatible: null → List.of(), blanks dropped, so an existing config parses unchanged.
  • Updating all ~9 call sites instead of adding a second constructor avoids the Jackson multi-constructor ambiguity — the right call, and the reasoning was stated.
  • 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.
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: ``` Tests run: 1386, Failures: 0, Errors: 0, Skipped: 0 BUILD SUCCESS ``` Identical to the reported numbers. The report was honest, and `inspectingAMissingMailboxNeverBreaksPublishOnTheSameInstance` is the right test for the stated invariant — it proves the 404-closes-the-channel hazard cannot take `publish` down. ## 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)"* — and `FleetMcp.probe`'s `catch (Exception)` collapses all of it before `peerView` emits `"reachable": false`. This issue exists because `fleet_send` said "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: `coordinatorView` emits `self.pending()` / `self.consumers()` with no flag, so a failed self-probe renders as `pending: 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 called `connection.createChannel()`, so it holds an open AMQP channel until the underlying call returns. Against a *hung* broker — not a down one, which fails fast — every `fleet_list` orphans another channel. Channel-max is finite (default 2047), and exhausting it would break `publish`: **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". `newThreadPerTaskExecutor` is unbounded — virtual threads make it cheap, not bounded. ## 3. `inspect` promises "never throws" and can throw (latent) `catch (IOException)` misses `AlreadyClosedException`, which `connection.createChannel()` raises on a dead connection. Checked against the jar on the classpath: ``` $ javap -cp . com.rabbitmq.client.AlreadyClosedException # amqp-client-5.32.0.jar public class com.rabbitmq.client.AlreadyClosedException extends com.rabbitmq.client.ShutdownSignalException public class com.rabbitmq.client.ShutdownSignalException extends java.lang.RuntimeException ``` Not reachable today — every caller goes through `probe`'s `catch (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 - The throwaway-probe-channel mechanism is right, and proved live rather than reasoned. - `Coordinator.peers` is properly back-compatible: `null → List.of()`, blanks dropped, so an existing config parses unchanged. - Updating all ~9 call sites instead of adding a second constructor avoids the Jackson multi-constructor ambiguity — the right call, and the reasoning was stated. - `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.
agent added 1 commit 2026-09-05 08:01:57 +02:00
fleetd #361 review: fix false-negative absent, uncancelled probes, and a throw contract gap
CI / contract (pull_request) Successful in 47s
CI / build (pull_request) Successful in 1m45s
c4d40fbc2b
Three findings from review of #364, fixed on the same branch:

1. LeadChannel.MailboxState.absent() was returned both for a genuinely
   absent mailbox AND for "the probe could not determine anything"
   (timeout, unreachable broker, other declare failure) -- exactly the
   overstatement #361 exists to fix, one level down. MailboxState now
   carries a Presence enum (EXISTS/ABSENT/UNKNOWN) with exists()/known()
   accessors; LeadMailbox.inspect classifies a real AMQP 404 (measured
   against a live broker, not assumed: an IOException wrapping a
   ShutdownSignalException whose Channel.Close reply code is 404) as
   ABSENT and everything else as UNKNOWN. fleet_list's mailbox/peer rows
   now render a "status" of exists/absent/unknown and only include
   pending/consumers when status is "exists", so an unresolved self- or
   peer-probe can never render as a measured zero.

2. FleetMcp.probe's get(timeoutMs) left a timed-out inspect() task
   running forever on its own virtual thread, holding the AMQP channel
   it had already opened -- against a hung (not down) broker this would
   orphan one channel per fleet_list call until the connection's
   channel-max was exhausted, breaking publish() too. probe() now holds
   the Future and calls cancel(true) on timeout/failure so the orphaned
   task is interrupted instead of abandoned, and now returns
   MailboxState.unknown() (never absent()) on timeout/exception.

3. LeadMailbox.inspect only caught IOException, but createChannel() on
   an already-closed connection throws AlreadyClosedException, an
   unchecked RuntimeException (measured against a live broker) -- so it
   could escape the "never throws" contract. Both places in inspect now
   also catch RuntimeException and report unknown().

Tests: MailboxState.exists()/absent()/unknown() call sites updated
across FleetMcpTest/FleetMcpLeadCoordTest; new hermetic tests cover the
tri-state fleet_list rendering (self-probe unknown, a peer that's
absent vs. one that's unknown) and probe cancellation (a LeadChannel
fake that blocks until interrupted, proving probe() doesn't just give
up on it); new @Tag("contract") LeadMailboxTest cases pin the real
exception shapes for both the 404 and the already-closed-connection
paths and prove inspect() reports unknown (never throws) when the
connection is already closed.
agent added 1 commit 2026-09-05 08:12:32 +02:00
fleetd #361 review round 2: pin isMissingQueue's false branch
CI / contract (pull_request) Successful in 1m5s
CI / build (pull_request) Successful in 1m34s
01492059d4
The reviewer's mutation (isMissingQueue always returns true) restored
the exact overstatement fleetd #361 exists to fix -- every declare
failure reading as a confirmed absence -- and still left mvn clean
install green (1389/1389), because no test drove a non-404 shape
through inspect(). The false branch was the whole discriminator
between MailboxState.absent() and MailboxState.unknown(), unpinned.

Widened LeadMailbox.isMissingQueue from private to package-private and
added LeadMailboxIsMissingQueueTest: five hermetic tests (no broker)
covering the true case and all three false shapes isMissingQueue's own
javadoc lists -- a different reply code, a ShutdownSignalException
whose reason isn't a Channel.Close, and an IOException with no such
cause at all (plus an IOException wrapping an unrelated exception
type). Re-ran the reviewer's exact mutation locally: 4 of 5 new tests
went red with the expected assertion messages; reverted, and mvn clean
install is green again at 1394/1394 (1389 + 5 new).

Also added a one-line javadoc note on LeadMailbox.inspect being honest
about which of its two RuntimeException catches is proven by a test
(the createChannel() one, end-to-end against a real broker) and which
stays purely defensive (the declare-site one, for a connection-drops-
mid-call race no test drives on purpose).
Owner

Merged to main as 6e9e464.

Verified on the merge itself, not taken from the report:

mvn clean install  ->  Tests run: 1394, Failures: 0, Errors: 0, Skipped: 0
                   ->  BUILD SUCCESS

I also re-ran my own mutation against the new tests on the merged tree, rather than trusting the worker's paste of it:

BASELINE   Tests run: 1394, Failures: 0            BUILD SUCCESS
MUTATED    isMissingQueue -> return true
           Tests run: 1394, Failures: 4            BUILD FAILURE

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 reached main, so replacing it with status costs no reader anything — grep -rn reachable docs/ wiki/ CLAUDE.md finds no reference to it.

Merged to `main` as `6e9e464`. Verified on the merge itself, not taken from the report: ``` mvn clean install -> Tests run: 1394, Failures: 0, Errors: 0, Skipped: 0 -> BUILD SUCCESS ``` I also re-ran my own mutation against the new tests on the merged tree, rather than trusting the worker's paste of it: ``` BASELINE Tests run: 1394, Failures: 0 BUILD SUCCESS MUTATED isMissingQueue -> return true Tests run: 1394, Failures: 4 BUILD FAILURE ``` 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 reached `main`, so replacing it with `status` costs no reader anything — `grep -rn reachable docs/ wiki/ CLAUDE.md` finds no reference to it.
ltms closed this pull request 2026-09-05 08:21:19 +02:00
Some checks are pending
CI / contract (pull_request) Successful in 1m5s
CI / build (pull_request) Successful in 1m34s

Pull request closed

Sign in to join this conversation.