diff --git a/fleetd/fleetd.example.yaml b/fleetd/fleetd.example.yaml index 5788a06..9295804 100644 --- a/fleetd/fleetd.example.yaml +++ b/fleetd/fleetd.example.yaml @@ -821,10 +821,17 @@ guard: # across every daemon sharing this vhost. # prefetch → consumer basicQos, capping how many unacked messages the mailbox holds in-heap. # Default 32 when omitted. +# peers → fleetd #361: the coord-ids of the OTHER daemons on this vhost, declared by the +# operator (the daemon never guesses). fleet_list reports each one's live reachability +# (a passive queue check, never a presence protocol) alongside this daemon's own +# mailbox state. Omit, or leave empty, for a daemon with no known peers yet — an +# undeclared peer can still reach you and be reached by fleet_send, it just will not +# show up as a row in fleet_list. # coordinator: # uriEnv: LEAD_COORD_URI # selfId: mac-opus # prefetch: 32 +# peers: [fleet01-lead] # Active push-to-primary (CB-307 Stage 3). When a worker reply lands with no open fleet_send, # the ReplyPushLoop injects a *drain nudge* (never the payload) into the primary's own herdr diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index 6058535..b8b7283 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -652,7 +652,12 @@ public final class Fleetd { quarantineSource, leadMailbox, outageSource, - new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads))); + new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)), + // fleetd #361: the operator-declared peers this daemon's fleet_list should try to + // reach. Read from the SAME snapshot leadMailbox itself opened from (cfg.coordinator()), + // not the live config.get() — coordinator wiring is already boot-time-fixed (see + // leadMailbox above), so peers follows the same rule rather than half hot-reloading. + cfg.coordinator() == null ? List.of() : cfg.coordinator().peers()); // CB-637: the receive half. Only constructed when a lead mailbox actually opened — with no // coordinator (or an unreachable one) there is nothing to deliver, so no scheduler is diff --git a/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java b/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java index ef96eb1..452743c 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java +++ b/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java @@ -888,12 +888,18 @@ public record FleetConfig( * {@code null} ⇒ kept as {@code null} (no self id configured). * @param prefetch the consumer's {@code basicQos} prefetch count. {@code null}/non-positive ⇒ * {@link LeadMailbox#DEFAULT_PREFETCH}. + * @param peers fleetd #361: the coord-ids the operator declares as this daemon's peers — the + * daemon never guesses who else exists. {@code fleet_list} reports each one's + * live reachability. Blank entries are dropped; {@code null} ⇒ an empty list, so + * a config written before this field existed still parses unchanged. */ @JsonIgnoreProperties(ignoreUnknown = true) - public record Coordinator(String uri, String uriEnv, String selfId, Integer prefetch) { + public record Coordinator(String uri, String uriEnv, String selfId, Integer prefetch, List peers) { public Coordinator { selfId = (selfId == null || selfId.isBlank()) ? null : selfId; + peers = peers == null ? List.of() + : peers.stream().filter(p -> p != null && !p.isBlank()).toList(); } /** True when a {@code uriEnv} is configured by name, whether or not its variable resolves. */ diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index 9dda812..7b35930 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -44,6 +44,7 @@ import java.util.Objects; import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.TimeUnit; import java.util.function.BiFunction; import java.util.function.Function; import java.util.function.LongSupplier; @@ -101,6 +102,8 @@ public final class FleetMcp { private final LeadSeatSource leadSeats; /** CB-637: this daemon's lead-to-lead channel; {@code null} when no coordinator is configured. */ private final LeadChannel leadChannel; + /** fleetd #361: {@code coordinator.peers} — see {@link CoordinationSource}. Empty when unset. */ + private final List peers; /** Capacity facts used by {@code fleet_list}; production must supply the placement live count. */ public record CapacitySource(Function liveCount, Function maxLoad, @@ -164,6 +167,32 @@ public final class FleetMcp { public static LeadSeatSource none() { return new LeadSeatSource(_ -> 0); } } + /** + * fleetd #361: peer-visibility facts for {@code fleet_list}'s {@code coordinator} row — this + * daemon's own {@link LeadChannel} (for its self mailbox state and held messages) plus the + * coord-ids the operator has declared as peers ({@code coordinator.peers}). Bundled as its own + * Source, the same idiom as {@link OutageSource}/{@link QuarantineSource}/{@link LeadSeatSource}, + * so {@code listFleet}'s already-long overload chain gains exactly one new required parameter + * instead of a further bare positional argument. + * + *

Every mailbox look this triggers goes through {@link LeadChannel#inspect}, which is + * specified to run on its own disposable channel — never the channel {@link LeadChannel#publish} + * or the consume loop depends on — so a peer that happens to be down, or the coordination broker + * itself being unreachable, can never take {@code fleet_send}/{@code LeadCoordLoop}'s own path + * down with it. See {@code FleetMcp.probe} for the additional timeout bound on top of that. + * + * @param leadChannel this daemon's own channel, or {@code null} when no coordinator is configured + * @param peers the coord-ids declared under {@code coordinator.peers}, or empty + */ + public record CoordinationSource(LeadChannel leadChannel, List peers) { + public CoordinationSource { + peers = peers == null ? List.of() : List.copyOf(peers); + } + + /** Inert source — no coordinator row is ever reported. */ + public static CoordinationSource none() { return new CoordinationSource(null, List.of()); } + } + /** * @param callers resolves each call's {@link Principal}; {@code null} disables authorization. * This surface needs its own enforcement: {@code /mcp} is a raw servlet on @@ -211,19 +240,34 @@ public final class FleetMcp { } /** - * As above, with fleetd #176 lead-seat facts (see {@link LeadSeatSource}). This is what - * {@code Fleetd.main} actually wires up. - * - * @param leadSeats required — pass {@link LeadSeatSource#none()} for a caller that does not want - * the feature, never a defaulting overload (the same rule {@code quarantine} and - * {@code outage} follow). + * As above, with fleetd #176 lead-seat facts (see {@link LeadSeatSource}). */ public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions, ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry, CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage, QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage, LeadSeatSource leadSeats) { + this(messages, workers, sessions, identity, presence, primaryRegistry, callers, metrics, capacity, + healthCoverage, quarantine, leadChannel, outage, leadSeats, List.of()); + } + + /** + * As above, with fleetd #361 {@code coordinator.peers} (see {@link CoordinationSource}). This is + * what {@code Fleetd.main} actually wires up. + * + * @param leadSeats required — pass {@link LeadSeatSource#none()} for a caller that does not want + * the feature, never a defaulting overload (the same rule {@code quarantine} and + * {@code outage} follow). + * @param peers the coord-ids declared under {@code coordinator.peers}; empty when unset or + * when {@code leadChannel} is {@code null}. + */ + public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions, + ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry, + CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage, + QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage, + LeadSeatSource leadSeats, List peers) { this.leadChannel = leadChannel; + this.peers = peers == null ? List.of() : List.copyOf(peers); this.capacity = capacity; this.quarantine = Objects.requireNonNull(quarantine, "quarantine"); this.outage = Objects.requireNonNull(outage, "outage"); @@ -361,7 +405,7 @@ public final class FleetMcp { return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage, leadSeats, callers == null ? Map.of() : callers.leads(), callerTerminal(exchange), - leadChannel == null ? null : leadChannel.selfCoordId()); + new CoordinationSource(leadChannel, peers)); }; BiFunction stopHandler = (exchange, req) -> { @@ -697,6 +741,15 @@ public final class FleetMcp { * turned into a tool error naming the coord-id. It is never allowed to escape as a crash: an * unreachable peer is an ordinary outcome of addressing a fleet you do not control. * + *

fleetd #361: the success text is honest about what "durably confirmed" does and + * does not mean. The broker's publisher confirm proves the message is durably queued — + * it says nothing about whether the peer's pane has, or ever will, receive it. After a + * successful publish this looks at the target mailbox's consumer count (via + * {@link LeadChannel#inspect}, bounded and never allowed to fail the call — see {@link #probe}) + * and appends a warning when it is zero: that is the observable form of "nobody is reading this + * right now". A zero-consumer publish is still reported as a SUCCESS, never an error — the + * message is safely queued and will be read once a daemon owning that coord-id connects. + * * @param leadChannel this daemon's channel, or {@code null} when no coordinator is configured */ static McpSchema.CallToolResult sendToLead(LeadChannel leadChannel, String coordId, String content, @@ -724,7 +777,16 @@ public final class FleetMcp { + ". Check that a daemon is running with coordinator.selfId=\"" + coordId + "\" and is connected to the same coordination broker."); } - return text("delivered to peer lead " + coordId + " (msgId " + msg.msgId() + ")"); + String result = "published to peer lead \"" + coordId + "\"'s mailbox and durably confirmed " + + "by the broker (msgId " + msg.msgId() + ")."; + LeadChannel.MailboxState state = probe(leadChannel, coordId); + if (state.exists() && state.consumers() == 0) { + result += " Warning: that mailbox currently has NO consumers attached — nobody is reading " + + "it right now. The message is safely queued and will be delivered once a daemon " + + "with coordinator.selfId=\"" + coordId + "\" is running and connected; until then " + + "it will not reach that lead's pane."; + } + return text(result); } /** @@ -1110,7 +1172,8 @@ public final class FleetMcp { static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages, CapacitySource capacity, HealthCoverageSource healthCoverage, QuarantineSource quarantine, Map leads, String selfTerm) { - return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, leads, selfTerm, null); + return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, leads, selfTerm, + CoordinationSource.none()); } /** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */ @@ -1119,33 +1182,33 @@ public final class FleetMcp { QuarantineSource quarantine, OutageSource outage, Map leads, String selfTerm) { return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage, - LeadSeatSource.none(), leads, selfTerm, null); + LeadSeatSource.none(), leads, selfTerm, CoordinationSource.none()); } /** - * As above, additionally reporting this daemon's own lead coordination id (CB-637) when one is - * configured and its channel opened. There is no peer-discovery surface yet — a lead addresses a - * peer by a coord-id it was told — so this row exists to answer the one question the operator - * cannot answer any other way: what is MY coord-id, the one a peer must use to reach me. It is - * omitted entirely when no coordinator is configured, so an ordinary fleet's output is unchanged. + * As above, additionally reporting this daemon's own lead coordination state (CB-637, fleetd + * #361) when a coordinator is configured and its channel opened — see {@link #coordinatorView} + * for the shape. Omitted entirely when no coordinator is configured, so an ordinary fleet's + * output is unchanged. * - * @param selfCoordId this daemon's coord-id, or {@code null} when lead coordination is off + * @param coordination this daemon's lead channel plus its declared peers, or + * {@link CoordinationSource#none()} when lead coordination is off */ static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages, CapacitySource capacity, HealthCoverageSource healthCoverage, QuarantineSource quarantine, Map leads, String selfTerm, - String selfCoordId) { + CoordinationSource coordination) { return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, OutageSource.none(), - LeadSeatSource.none(), leads, selfTerm, selfCoordId); + LeadSeatSource.none(), leads, selfTerm, coordination); } /** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */ static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages, CapacitySource capacity, HealthCoverageSource healthCoverage, QuarantineSource quarantine, OutageSource outage, - Map leads, String selfTerm, String selfCoordId) { + Map leads, String selfTerm, CoordinationSource coordination) { return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage, - LeadSeatSource.none(), leads, selfTerm, selfCoordId); + LeadSeatSource.none(), leads, selfTerm, coordination); } /** As above, plus fleetd #176 lead-seat facts (see {@link LeadSeatSource}). */ @@ -1153,7 +1216,7 @@ public final class FleetMcp { CapacitySource capacity, HealthCoverageSource healthCoverage, QuarantineSource quarantine, OutageSource outage, LeadSeatSource leadSeats, Map leads, String selfTerm, - String selfCoordId) { + CoordinationSource coordination) { try { Map live = workers.list().stream() .map(Agent.class::cast) @@ -1175,8 +1238,9 @@ public final class FleetMcp { Map result = new LinkedHashMap<>(); result.put("leads", leadRows); result.put("members", out); result.put("healthCoverage", healthCoverage.value().get()); - if (selfCoordId != null && !selfCoordId.isBlank()) { - result.put("coordinator", Map.of("selfId", selfCoordId, "configured", true)); + Map coordinatorRow = coordinatorView(coordination); + if (coordinatorRow != null) { + result.put("coordinator", coordinatorRow); } if (capacity.available()) result.put("capacity", profiles.stream() .map(profile -> capacityView(profile, capacity.liveCount(), capacity.maxLoad(), roster, messages, @@ -1187,6 +1251,140 @@ public final class FleetMcp { } } + /** + * fleetd #361: the {@code coordinator} row — this daemon's own coord-id and mailbox state, the + * messages currently held for it, and the live reachability of every operator-declared peer. + * {@code null} (the row is then omitted entirely) when lead coordination is off, so an ordinary + * fleet's {@code fleet_list} output is byte-identical to before this feature existed. + * + *

Every peer/self mailbox look goes through {@link #probe}, which bounds each + * {@link LeadChannel#inspect} call to {@link #PEER_PROBE_TIMEOUT_MS} and never lets it throw — + * a coordination broker that is down or slow degrades this row toward "unreachable"/"unknown" + * counts, it can never make {@code fleet_list} itself slow or fail. {@code held} comes from + * {@link LeadChannel#peek}, a pure in-memory read with no broker round trip, so it is never + * subject to that bound. + */ + private static Map coordinatorView(CoordinationSource coordination) { + LeadChannel channel = coordination.leadChannel(); + if (channel == null) { + return null; + } + String selfId = channel.selfCoordId(); + Map row = new LinkedHashMap<>(); + row.put("selfId", selfId); + row.put("configured", true); + row.put("mailbox", mailboxView(probe(channel, selfId))); + row.put("held", channel.peek().stream().map(FleetMcp::heldView).toList()); + row.put("peers", coordination.peers().stream().map(p -> peerView(channel, p)).toList()); + return row; + } + + /** + * fleetd #361: render a {@link LeadChannel.MailboxState} without ever presenting an unmeasured + * fact as a measured one. {@code status} is the tri-state itself — {@code "exists"}, + * {@code "absent"} (the broker positively confirmed no such queue), or {@code "unknown"} (the + * probe could not determine either way: down, unreachable, or timed out). {@code pending}/ + * {@code consumers} are included ONLY when {@code status == "exists"} — a reader must never see + * them default to {@code 0} for a mailbox this call never actually measured. This is the fix for + * the review finding that a collapsed {@code absent()} rendered a self-probe timeout as + * "pending: 0, consumers: 0", indistinguishable from an actually-empty, actually-unread mailbox. + */ + private static Map mailboxView(LeadChannel.MailboxState state) { + Map row = new LinkedHashMap<>(); + row.put("status", state.exists() ? "exists" : state.known() ? "absent" : "unknown"); + if (state.exists()) { + row.put("pending", state.pending()); + row.put("consumers", state.consumers()); + } + return row; + } + + /** One held-for-me message: enough to identify it and see roughly what it says, never the whole body. */ + private static Map heldView(LeadMessage m) { + Map row = new LinkedHashMap<>(); + row.put("msgId", m.msgId()); + row.put("from", m.from()); + row.put("preview", preview(m.content())); + return row; + } + + /** Cap a held message's content to a short preview — {@code fleet_list} must never dump a full body. */ + private static final int HELD_PREVIEW_MAX_CHARS = 80; + + private static String preview(String content) { + if (content == null) { + return ""; + } + return content.length() <= HELD_PREVIEW_MAX_CHARS + ? content + : content.substring(0, HELD_PREVIEW_MAX_CHARS) + "…"; + } + + /** + * One declared peer's row: its coord-id, then the same tri-state {@link #mailboxView} shape. + * Deliberately no boolean "reachable" field — that collapsed "confirmed gone" and "could not + * check" into the same {@code false}, which is exactly the review finding this row now avoids: + * an operator reading {@code status} can tell "fleet01 is down" (a {@code coordinator.selfId} + * nobody has ever run) apart from "my own broker is slow or unreachable right now". + */ + private static Map peerView(LeadChannel channel, String coordId) { + Map row = new LinkedHashMap<>(); + row.put("coordId", coordId); + row.putAll(mailboxView(probe(channel, coordId))); + return row; + } + + /** + * fleetd #361: how long {@code fleet_list} waits on any single {@link LeadChannel#inspect} call + * before giving up on it — see {@link #probe}. + */ + private static final long PEER_PROBE_TIMEOUT_MS = 1_500L; + + /** + * Dedicated pool for {@link LeadChannel#inspect} calls so a slow one blocks only its own virtual + * thread, never the MCP request thread calling {@code fleet_list}. Not a bounded pool — + * {@code newThreadPerTaskExecutor} starts a fresh virtual thread per call with no cap on how many + * run at once; virtual threads make that cheap, not bounded. What actually keeps a hung probe + * from accumulating forever is the {@link Future#cancel} in {@link #probe}, not a pool limit. + */ + private static final java.util.concurrent.ExecutorService PEER_PROBE_POOL = + java.util.concurrent.Executors.newThreadPerTaskExecutor(Thread.ofVirtual().name("fleet-peer-probe-", 0).factory()); + + /** + * fleetd #361: {@link LeadChannel#inspect}, bounded to {@link #PEER_PROBE_TIMEOUT_MS} and never + * allowed to throw or hang the caller — a coordination broker that is unreachable or slow + * degrades to {@link LeadChannel.MailboxState#unknown} (never {@code absent}: a timeout proves + * nothing about whether the mailbox exists) rather than making {@code fleet_list} slow or + * failing it. {@code inspect} itself is already specified to never throw, but this is the seam + * that also survives an implementation that does, or one that blocks indefinitely on a dead + * connection. + * + *

A timeout cancels the orphaned task rather than abandoning it. Before this, + * {@code get(timeout)} on a hung {@code inspect} left the submitted task running forever on its + * own virtual thread, holding the AMQP channel it had already opened — against a broker that + * hangs rather than fails fast, every {@code fleet_list} call would orphan one more channel until + * the connection's channel-max (2047 by default) was exhausted, which would break {@link + * LeadChannel#publish} too. {@link Future#cancel(boolean) cancel(true)} interrupts the orphaned + * task's thread; {@link LeadMailbox#inspect} has no interruptible wait of its own to catch that, + * but the underlying AMQP RPC continuation does block on one, so the interrupt reaches it and the + * task's {@code finally} still closes the probe channel it opened rather than leaking it forever. + */ + private static LeadChannel.MailboxState probe(LeadChannel channel, String coordId) { + return probe(channel, coordId, PEER_PROBE_TIMEOUT_MS); + } + + /** As {@link #probe(LeadChannel, String)}, with an explicit timeout — a seam for tests. */ + static LeadChannel.MailboxState probe(LeadChannel channel, String coordId, long timeoutMs) { + java.util.concurrent.Future future = + PEER_PROBE_POOL.submit(() -> channel.inspect(coordId)); + try { + return future.get(timeoutMs, TimeUnit.MILLISECONDS); + } catch (Exception e) { + future.cancel(true); // best-effort: don't leave a hung probe (and its channel) running forever + return LeadChannel.MailboxState.unknown(coordId); + } + } + /** * Capacity is advisory only. {@code reclaimable} says there is no bridge work, not that fleetd * may stop the member: the bridge has capacity facts but no work list, and choosing work needs @@ -1377,8 +1575,12 @@ public final class FleetMcp { + "routes your answer back into the same turn (omit for a normal delegation)"), "coordId", stringProp("A peer LEAD's coordination id — delivers content to that " + "lead's durable mailbox on the shared coordination broker, which works " - + "across hosts. Mutually exclusive with sessionId and turnId. Your own " - + "coordId is reported by fleet_list.")), + + "across hosts. Mutually exclusive with sessionId and turnId. Success means " + + "the message is durably queued and confirmed by the broker, with a warning " + + "if that mailbox has no consumers attached right now (queued, but nobody is " + + "reading it yet) — it does not mean the peer's pane has seen it. Your own " + + "coordId, this daemon's mailbox state, and every coordinator.peers entry's " + + "live reachability are reported by fleet_list's coordinator row.")), List.of("content"))); } @@ -1486,7 +1688,13 @@ public final class FleetMcp { + "for a quarantined profile's credential (see fleet_profiles), whatever its " + "maxLoad/live — with credentialId and quarantinedForSeconds naming the " + "quarantine, so 'free: 0, busy' can be told apart from 'free: 0, refusing " - + "for N seconds'.", + + "for N seconds'. When lead-to-lead coordination is configured, a 'coordinator' " + + "object reports this daemon's own coord-id ('selfId') and mailbox state " + + "('mailbox': pending/consumers), the messages currently held for it ('held': " + + "msgId/from/preview, never the full body), and one row per coordinator.peers " + + "coord-id ('peers': coordId/reachable, plus pending/consumers when reachable) — " + + "this is peer DISCOVERY for cross-host leads, distinct from the local 'leads' " + + "array above. It is omitted entirely when no coordinator is configured.", objectSchema(Map.of(), List.of())); } diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadChannel.java b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadChannel.java index 07664a9..db16aa3 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadChannel.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadChannel.java @@ -4,8 +4,8 @@ import java.util.List; /** * The lead-to-lead message channel this daemon speaks, as its callers need it — one lead's own - * mailbox: publish to a peer's coord-id, look at what has arrived for me, and ack what I have - * delivered. + * mailbox: publish to a peer's coord-id, look at what has arrived for me, ack what I have + * delivered, and (fleetd #361) inspect any coord-id's mailbox from the outside without owning it. * *

Extracted from {@link LeadMailbox} purely as a seam. {@code LeadMailbox} is the one production * implementation and owns a live AMQP connection, so a test that wanted to exercise the routing in @@ -37,4 +37,78 @@ public interface LeadChannel { /** This daemon's own lead coordination id — the mailbox it owns, and the {@code from} it sends as. */ String selfCoordId(); + + /** + * A non-destructive look at {@code coordId}'s mailbox — does it exist, how many messages are + * waiting on it, and how many consumers are attached — without owning, consuming, or otherwise + * changing it. {@code consumers == 0} on an existing mailbox is the observable form of "nobody + * is reading this right now": a publish to it will sit queued rather than reach a pane. + * + *

Never throws — this is a best-effort fact-finding call, not an operation a + * caller must handle failing. But it must never turn "I could not check" into a false negative: + * {@link MailboxState#absent(String)} means the broker positively confirmed there is no such + * queue, and {@link MailboxState#unknown(String)} — a distinct value — means the look could not + * be completed at all (broker unreachable, timed out, connection closed). A caller that + * collapses those two into one, as fleetd #361 initially did, cannot tell "that peer is down" + * from "I could not check", and a reader of {@code pending}/{@code consumers} cannot tell a + * measured zero from a zero standing in for "not measured". + * + *

Must never share fate with {@link #publish} or {@link #peek}/{@link #ack}. + * fleetd #361: in AMQP 0-9-1 a passive queue declare of a queue that does not exist closes the + * channel it was declared on with a 404. An implementation backed by a real broker connection + * must inspect on a channel it can afford to lose — never the channel {@link #publish} or the + * consume loop depends on — so that looking at a peer that happens to be down can never break + * this daemon's own send or receive path. + */ + MailboxState inspect(String coordId); + + /** + * The result of {@link #inspect}. {@code presence} tells apart three states a caller must not + * conflate: a confirmed-existing mailbox ({@link Presence#EXISTS}, the only case where + * {@code pending}/{@code consumers} are measured facts), a confirmed-absent one + * ({@link Presence#ABSENT} — the broker positively said "no such queue"), and one this call + * simply could not determine ({@link Presence#UNKNOWN} — broker unreachable, timed out, + * connection closed). {@code pending}/{@code consumers} are always {@code 0} and meaningless + * outside {@link Presence#EXISTS}; a renderer must gate on {@link #exists()} (or {@code + * presence} directly), never present them as measured otherwise. + * + * @param coordId the coord-id inspected + * @param presence whether the mailbox is confirmed to exist, confirmed absent, or unknown + * @param pending messages ready for delivery but not yet in a consumer's hands (0 unless EXISTS) + * @param consumers how many consumers are attached (0 unless EXISTS) + */ + record MailboxState(String coordId, Presence presence, int pending, int consumers) { + + /** Whether {@link #inspect} was able to reach a definite answer, of either kind. */ + public enum Presence { EXISTS, ABSENT, UNKNOWN } + + /** {@code true} only when the broker confirmed this exact queue is currently declared. */ + public boolean exists() { + return presence == Presence.EXISTS; + } + + /** + * {@code true} when {@link #inspect} reached a definite answer (exists or confirmed + * absent); {@code false} when it could not determine either way. A caller must never treat + * {@code !known()} the same as a confirmed absence — the mailbox may well exist. + */ + public boolean known() { + return presence != Presence.UNKNOWN; + } + + /** The broker confirmed this queue exists, with these measured counts. */ + public static MailboxState exists(String coordId, int pending, int consumers) { + return new MailboxState(coordId, Presence.EXISTS, pending, consumers); + } + + /** The broker positively confirmed there is no such queue (e.g. a 404 on passive declare). */ + public static MailboxState absent(String coordId) { + return new MailboxState(coordId, Presence.ABSENT, 0, 0); + } + + /** The look could not be completed — broker unreachable, timed out, or connection closed. */ + public static MailboxState unknown(String coordId) { + return new MailboxState(coordId, Presence.UNKNOWN, 0, 0); + } + } } diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java index 78ee0da..a2f3b64 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java @@ -10,6 +10,7 @@ import com.rabbitmq.client.DeliverCallback; import com.rabbitmq.client.Recoverable; import com.rabbitmq.client.RecoveryListener; import com.rabbitmq.client.Return; +import com.rabbitmq.client.ShutdownSignalException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -259,6 +260,89 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable { } } + /** + * fleetd #361: look at {@code coordId}'s mailbox on a fresh, immediately-closed throwaway + * channel — never {@link #channel} (consume/ack) or {@link #publishChannel} (publish). A + * passive queue declare of a queue that does not exist closes the channel it was declared on + * with a 404; using a disposable probe channel means that closure can never touch either + * long-lived channel this instance depends on for {@link #publish} or the consume loop. + * + *

Classifies failures rather than collapsing them, both measured against a real broker in + * {@code LeadMailboxTest} rather than assumed from the AMQP 0-9-1 spec text: + *

    + *
  • a genuine 404 — an {@link IOException} wrapping a {@link ShutdownSignalException} whose + * {@link AMQP.Channel.Close#getReplyCode()} is {@code 404} — reports + * {@link MailboxState#absent}; every other declare failure reports + * {@link MailboxState#unknown} instead of quietly becoming the same "absent" value; + *
  • {@code catch (RuntimeException e)} on both attempts matters as much as the checked + * catches: a connection that is already closed makes {@link Connection#createChannel()} + * throw {@link com.rabbitmq.client.AlreadyClosedException} (a {@link RuntimeException}, + * not an {@link IOException}) — an {@code inspect} that only caught {@code IOException} + * would let that escape, breaking the "never throws" contract this method promises. + *
+ * + *

Honesty about which catch is measured and which is defensive: the + * {@code createChannel()} catch above is exercised end-to-end against a real broker by + * {@code LeadMailboxTest.inspectReportsUnknownRatherThanThrowingWhenTheConnectionIsAlreadyClosed}. + * The second {@code catch (RuntimeException e)}, around the passive declare itself — for the + * narrower race where the connection drops between {@code createChannel()} succeeding + * and the declare landing — has no such test; reaching it needs a connection that dies at that + * exact instant, which is not a scenario this suite drives on purpose. It stays purely + * defensive: correct by the same reasoning as the first catch, but unproven the way the first + * one is proven. + */ + @Override + public MailboxState inspect(String coordId) { + String queue = queueName(coordId); + Channel probe; + try { + probe = connection.createChannel(); + } catch (IOException | RuntimeException e) { + log.debug("lead mailbox inspect: cannot open a probe channel for {}: {}", coordId, e.toString()); + return MailboxState.unknown(coordId); + } + try { + AMQP.Queue.DeclareOk declared = probe.queueDeclarePassive(queue); + return MailboxState.exists(coordId, declared.getMessageCount(), declared.getConsumerCount()); + } catch (IOException e) { + // The broker (or the client library) has already closed `probe` for us either way; only + // a confirmed 404 means "no such queue" — anything else (a different declare failure) is + // "could not determine", never silently reported as the same value as a genuine absence. + return isMissingQueue(e) ? MailboxState.absent(coordId) : MailboxState.unknown(coordId); + } catch (RuntimeException e) { + // E.g. the connection dropped between createChannel() and the declare landing. + log.debug("lead mailbox inspect: declare failed unexpectedly for {}: {}", coordId, e.toString()); + return MailboxState.unknown(coordId); + } finally { + try { + if (probe.isOpen()) { + probe.close(); + } + } catch (Exception e) { + log.debug("lead mailbox inspect: probe channel close for {}: {}", coordId, e.toString()); + } + } + } + + /** + * {@code true} only for the specific shape a missing-queue passive declare actually produces — + * measured against a real broker, not assumed from the spec text (see {@code + * LeadMailboxTest.passiveDeclareOfAMissingQueueThrowsAnIOExceptionWrappingA404ShutdownSignal}): + * an {@link IOException} whose cause is a {@link ShutdownSignalException} carrying an + * {@link AMQP.Channel.Close} reason with {@code replyCode == 404}. Any other shape (a different + * reply code, a {@code ShutdownSignalException} cause whose reason is not a + * {@code Channel.Close}, or no cause at all) is a declare failure of some other kind and must + * not be read as "confirmed absent" — pinned hermetically, with no broker needed, by + * {@code LeadMailboxIsMissingQueueTest} for exactly those three false shapes. Package-private + * (not {@code private}) so that test can call it directly. + */ + static boolean isMissingQueue(IOException e) { + if (!(e.getCause() instanceof ShutdownSignalException sse)) { + return false; + } + return sse.getReason() instanceof AMQP.Channel.Close close && close.getReplyCode() == AMQP.NOT_FOUND; + } + /** Convenience: {@link #peek} the current snapshot, then {@link #ack} every message in it. */ public List drain() { List snapshot = peek(); diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdLeadMailboxSelectionTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdLeadMailboxSelectionTest.java index 82cc333..f681bfa 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/FleetdLeadMailboxSelectionTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdLeadMailboxSelectionTest.java @@ -80,7 +80,7 @@ class FleetdLeadMailboxSelectionTest { @Test void opensTheMailboxWhenAUriAndSelfIdAreConfigured() { var opener = new RecordingOpener(); - var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null); + var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null, null); Fleetd.openLeadMailbox(coordinator, Map.of(), opener); @@ -94,7 +94,7 @@ class FleetdLeadMailboxSelectionTest { void honoursUriEnvOverALiteralUri() { var opener = new RecordingOpener(); var coordinator = new FleetConfig.Coordinator("amqp://stale:stale@old:5672/x", "COORD_URI", - "mac-opus", 8); + "mac-opus", 8, null); Fleetd.openLeadMailbox(coordinator, Map.of("COORD_URI", RESOLVED_URI), opener); @@ -105,7 +105,7 @@ class FleetdLeadMailboxSelectionTest { @Test void turnsOffWhenUriEnvDoesNotResolve() { var opener = new RecordingOpener(); - var coordinator = new FleetConfig.Coordinator(null, "COORD_URI", "mac-opus", null); + var coordinator = new FleetConfig.Coordinator(null, "COORD_URI", "mac-opus", null, null); assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener)); @@ -116,7 +116,7 @@ class FleetdLeadMailboxSelectionTest { void warnsAndStaysOffWhenSelfIdIsMissing() { var appender = captureFleetdLogs(); var opener = new RecordingOpener(); - var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, null, null); + var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, null, null, null); assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener)); @@ -131,7 +131,7 @@ class FleetdLeadMailboxSelectionTest { var appender = captureFleetdLogs(); var opener = new RecordingOpener(); opener.unreachable = true; - var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null); + var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null, null); assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener), "a down coordination broker turns the feature off; it must never take the daemon down"); diff --git a/fleetd/src/test/java/dev/ltms/fleet/config/ConfigRefTopLevelReportingCoverageTest.java b/fleetd/src/test/java/dev/ltms/fleet/config/ConfigRefTopLevelReportingCoverageTest.java index aea7bba..030258e 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/config/ConfigRefTopLevelReportingCoverageTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/config/ConfigRefTopLevelReportingCoverageTest.java @@ -104,7 +104,7 @@ class ConfigRefTopLevelReportingCoverageTest { v.put("configReload", new FleetConfig.ConfigReload(true, 10)); v.put("quarantineCooldownSeconds", 1800); v.put("memberCredentials", null); - v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-a", null, "self-a", 1)); + v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-a", null, "self-a", 1, null)); v.put("worktreeGroup", "group-a"); v.put("memberLoginShell", null); assertNamesMatchComponents(v); @@ -144,7 +144,7 @@ class ConfigRefTopLevelReportingCoverageTest { v.put("configReload", new FleetConfig.ConfigReload(false, 20)); v.put("quarantineCooldownSeconds", 3600); v.put("memberCredentials", null); - v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-b", null, "self-b", 2)); + v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-b", null, "self-b", 2, null)); v.put("worktreeGroup", "group-b"); v.put("memberLoginShell", null); assertNamesMatchComponents(v); diff --git a/fleetd/src/test/java/dev/ltms/fleet/config/FleetConfigTest.java b/fleetd/src/test/java/dev/ltms/fleet/config/FleetConfigTest.java index 4292139..df9331b 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/config/FleetConfigTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/config/FleetConfigTest.java @@ -1189,7 +1189,7 @@ class FleetConfigTest { @Test void coordinatorEffectiveUriHonorsUriEnv() { FleetConfig.Coordinator withEnv = new FleetConfig.Coordinator( - "amqp://stale-clear-text@127.0.0.1:5672/coord", "LEAD_COORD_URI", "fleet01-lead", null); + "amqp://stale-clear-text@127.0.0.1:5672/coord", "LEAD_COORD_URI", "fleet01-lead", null, null); assertEquals("amqp://from-env@127.0.0.1:5672/coord", withEnv.effectiveUri(Map.of("LEAD_COORD_URI", "amqp://from-env@127.0.0.1:5672/coord")), @@ -1200,12 +1200,61 @@ class FleetConfigTest { "a blank uriEnv variable must not fall back to the literal uri"); FleetConfig.Coordinator noEnv = new FleetConfig.Coordinator( - "amqp://guest:guest@127.0.0.1:5672/coord", null, null, null); + "amqp://guest:guest@127.0.0.1:5672/coord", null, null, null, null); assertEquals("amqp://guest:guest@127.0.0.1:5672/coord", noEnv.effectiveUri(Map.of()), "the literal uri is used when no uriEnv is configured"); assertEquals(LeadMailbox.DEFAULT_PREFETCH, noEnv.prefetchOrDefault()); } + /** + * fleetd #361: the live block ships as just {@code uriEnv} + {@code selfId} (see + * {@code Fleetd.example.yaml} / the operator's real {@code fleetd.yaml}, gitignored). That exact + * shape, with no {@code peers:} key at all, must keep parsing unchanged after this field is added. + */ + @Test + void coordinatorBlockWithNoPeersKeyStillParses(@TempDir Path dir) throws Exception { + Path f = dir.resolve("coordinator-no-peers.yaml"); + Files.writeString(f, """ + bind: + port: 8080 + coordinator: + uriEnv: COORD_AMQP_URI + selfId: mac + """); + + FleetConfig cfg = FleetConfig.load(f); + assertNotNull(cfg.coordinator()); + assertEquals("mac", cfg.coordinator().selfId()); + assertEquals(List.of(), cfg.coordinator().peers(), "no peers: key means no configured peers, never null"); + } + + @Test + void coordinatorPeersParsesAndDropsBlankEntries(@TempDir Path dir) throws Exception { + Path f = dir.resolve("coordinator-peers.yaml"); + Files.writeString(f, """ + bind: + port: 8080 + coordinator: + selfId: mac + uri: amqp://guest:guest@127.0.0.1:5672/coord + peers: + - fleet01 + - "" + - fleet02 + """); + + FleetConfig cfg = FleetConfig.load(f); + assertEquals(List.of("fleet01", "fleet02"), cfg.coordinator().peers(), + "a blank peer entry must be dropped, never kept as an empty coord-id"); + } + + @Test + void coordinatorPeersDefaultsToEmptyWhenConstructedWithNull() { + FleetConfig.Coordinator c = new FleetConfig.Coordinator( + "amqp://guest:guest@127.0.0.1:5672/coord", null, "mac", null, null); + assertEquals(List.of(), c.peers(), "a null peers list must default to empty, never NPE downstream"); + } + @Test void absentWorktreeGroupLeavesItNull(@TempDir Path dir) throws Exception { Path f = dir.resolve("no-worktree-group.yaml"); diff --git a/fleetd/src/test/java/dev/ltms/fleet/config/FleetConfigWithDefaultsPreservesEveryComponentTest.java b/fleetd/src/test/java/dev/ltms/fleet/config/FleetConfigWithDefaultsPreservesEveryComponentTest.java index 1f710a0..7bbfc3a 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/config/FleetConfigWithDefaultsPreservesEveryComponentTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/config/FleetConfigWithDefaultsPreservesEveryComponentTest.java @@ -92,7 +92,7 @@ class FleetConfigWithDefaultsPreservesEveryComponentTest { v.put("memberCredentials", new FleetConfig.MemberCredentials( FleetConfig.MemberCredentials.POLICY_DENY_BY_DEFAULT, List.of("git"), List.of("git", "ssh"), null)); - v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-guard", null, "self-guard", 3)); + v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-guard", null, "self-guard", 3, null)); v.put("worktreeGroup", "group-guard"); v.put("memberLoginShell", "/bin/zsh"); assertNamesMatchComponents(v); diff --git a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpLeadCoordTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpLeadCoordTest.java index 6a3b66c..2384c1c 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpLeadCoordTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpLeadCoordTest.java @@ -1,6 +1,7 @@ package dev.ltms.fleet.mcp; import dev.ltms.fleet.msg.FakeLeadChannel; +import dev.ltms.fleet.msg.LeadChannel; import dev.ltms.fleet.msg.LeadMessage; import io.modelcontextprotocol.spec.McpSchema; import org.junit.jupiter.api.Test; @@ -93,4 +94,56 @@ class FleetMcpLeadCoordTest { assertTrue(FleetMcp.sendToLead(channel, PEER, " ", null, null).isError()); assertEquals(0, channel.published().size()); } + + /** + * fleetd #361: "delivered" overstated what publish actually proves — the broker's confirm means + * durably queued, not read. A zero-consumer target is the observable form of "this will not + * reach a pane right now", so the (still-successful) result must say so. + */ + @Test + void warnsWhenThePeerMailboxHasNoConsumersButStillReportsSuccess() { + var channel = new FakeLeadChannel(SELF) + .withMailbox(PEER, LeadChannel.MailboxState.exists(PEER, 0, 0)); + + McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", null, null); + + assertFalse(res.isError(), "a zero-consumer mailbox is still a successful, durably-queued publish"); + String out = textOf(res); + assertTrue(out.contains("durably confirmed"), out); + assertTrue(out.toLowerCase().contains("no consumers"), () -> "must warn nobody is reading it: " + out); + } + + @Test + void staysQuietAboutConsumersWhenThePeerMailboxHasOne() { + var channel = new FakeLeadChannel(SELF) + .withMailbox(PEER, LeadChannel.MailboxState.exists(PEER, 0, 1)); + + McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", null, null); + + assertFalse(res.isError()); + String out = textOf(res); + assertTrue(out.contains("durably confirmed"), out); + assertFalse(out.toLowerCase().contains("no consumers"), () -> "a consumer IS attached: " + out); + } + + /** + * fleetd #361 review finding 1: an unmeasured fact must never render as a definite one. When + * the post-publish probe could not determine the mailbox's consumer count at all (broker slow, + * unreachable, or the probe timed out — {@link LeadChannel.MailboxState#unknown}), the result + * must stay just as quiet as the has-a-consumer case — never assert "no consumers" for a mailbox + * this call never actually measured. + */ + @Test + void staysQuietAboutConsumersWhenThePeerMailboxStateIsUnknown() { + var channel = new FakeLeadChannel(SELF) + .withMailbox(PEER, LeadChannel.MailboxState.unknown(PEER)); + + McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", null, null); + + assertFalse(res.isError(), "an unresolved post-publish probe must never turn a durably-confirmed publish into an error"); + String out = textOf(res); + assertTrue(out.contains("durably confirmed"), out); + assertFalse(out.toLowerCase().contains("no consumers"), + () -> "an unmeasured fact must never be reported as a definite zero-consumer mailbox: " + out); + } } diff --git a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java index 0e08820..554d767 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java @@ -10,6 +10,9 @@ import dev.ltms.fleet.herdr.FakeHerdr; import dev.ltms.fleet.herdr.PaneLocator; import dev.ltms.fleet.herdr.WorkspaceControl; import dev.ltms.fleet.inject.Injector; +import dev.ltms.fleet.msg.FakeLeadChannel; +import dev.ltms.fleet.msg.LeadChannel; +import dev.ltms.fleet.msg.LeadMessage; import dev.ltms.fleet.msg.MessageService; import dev.ltms.fleet.msg.Rendezvous; import dev.ltms.fleet.session.FakeWorktrees; @@ -34,7 +37,9 @@ import java.util.Map; import java.util.EnumSet; import java.util.Set; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.Function; import static org.junit.jupiter.api.Assertions.*; @@ -576,17 +581,48 @@ class FleetMcpTest { void listReportsThisDaemonsOwnCoordIdWhenLeadCoordinationIsOn() { FakeHerdr h = new FakeHerdr(); SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))); + FakeLeadChannel channel = new FakeLeadChannel("mac-opus") + .withMailbox("mac-opus", LeadChannel.MailboxState.exists("mac-opus", 0, 1)); McpSchema.CallToolResult res = FleetMcp.listFleet( workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null, FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"), - FleetMcp.QuarantineSource.none(), Map.of(), "", "mac-opus"); + FleetMcp.QuarantineSource.none(), Map.of(), "", + new FleetMcp.CoordinationSource(channel, List.of())); String out = textOf(res); - // There is no peer-discovery surface yet, so this row answers the one question an operator - // cannot answer any other way: which coord-id a peer must use to reach ME. + // fleetd #361: reports both which coord-id a peer must use to reach ME, and this daemon's + // own mailbox state (a self-diagnosis: is my own consumer actually attached?). assertTrue(out.contains("\"coordinator\""), out); assertTrue(out.contains("\"selfId\":\"mac-opus\""), out); + assertTrue(out.contains("\"mailbox\":{\"status\":\"exists\",\"pending\":0,\"consumers\":1}"), out); + assertTrue(out.contains("\"held\":[]"), out); + assertTrue(out.contains("\"peers\":[]"), out); + } + + /** + * fleetd #361 review finding 1: a self-probe that could not complete (broker unreachable, timed + * out) must never render the same as a measured "0 pending, 0 consumers" — that was exactly the + * bug: a reader could not tell "my mailbox is empty and idle" from "I could not check", and the + * second one is the far more alarming state. + */ + @Test + void listReportsAnUnresolvedSelfProbeAsUnknownNeverAsAMeasuredZero() { + FakeHerdr h = new FakeHerdr(); + SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))); + FakeLeadChannel channel = new FakeLeadChannel("mac-opus") + .withMailbox("mac-opus", LeadChannel.MailboxState.unknown("mac-opus")); + + McpSchema.CallToolResult res = FleetMcp.listFleet( + workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null, + FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"), + FleetMcp.QuarantineSource.none(), Map.of(), "", + new FleetMcp.CoordinationSource(channel, List.of())); + + String out = textOf(res); + assertTrue(out.contains("\"mailbox\":{\"status\":\"unknown\"}"), out); + assertFalse(out.contains("\"pending\""), "an unresolved probe must never carry a pending count at all: " + out); + assertFalse(out.contains("\"consumers\""), "an unresolved probe must never carry a consumers count at all: " + out); } @Test @@ -601,6 +637,110 @@ class FleetMcpTest { "an ordinary fleet's output must be unchanged by this feature"); } + @Test + void listReportsHeldMessagesWithATruncatedPreviewNeverTheFullBody() { + FakeHerdr h = new FakeHerdr(); + SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))); + String longContent = "x".repeat(200); + FakeLeadChannel channel = new FakeLeadChannel("mac-opus") + .hold(new LeadMessage("m1", "fleet01-lead", "mac-opus", longContent)); + + McpSchema.CallToolResult res = FleetMcp.listFleet( + workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null, + FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"), + FleetMcp.QuarantineSource.none(), Map.of(), "", + new FleetMcp.CoordinationSource(channel, List.of())); + + String out = textOf(res); + assertTrue(out.contains("\"msgId\":\"m1\""), out); + assertTrue(out.contains("\"from\":\"fleet01-lead\""), out); + assertFalse(out.contains(longContent), "fleet_list must never dump a held message's full body: " + out); + assertTrue(out.contains("x".repeat(80) + "…"), "expected an 80-char preview with an ellipsis: " + out); + } + + @Test + void listReportsEachDeclaredPeersLiveReachability() { + FakeHerdr h = new FakeHerdr(); + SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))); + FakeLeadChannel channel = new FakeLeadChannel("mac-opus") + .withMailbox("fleet01-lead", LeadChannel.MailboxState.exists("fleet01-lead", 2, 1)) + .withMailbox("fleet03-lead", LeadChannel.MailboxState.unknown("fleet03-lead")); + // "fleet02-lead" is declared as a peer but never configured on the fake — inspect() falls + // back to MailboxState.absent, exactly as a real down (never-run) peer would report. + // "fleet03-lead" IS configured, as unknown — a broker that could not be reached in time, + // which review finding 1 says must render distinctly from "fleet02-lead"'s confirmed absence. + + McpSchema.CallToolResult res = FleetMcp.listFleet( + workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null, + FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"), + FleetMcp.QuarantineSource.none(), Map.of(), "", + new FleetMcp.CoordinationSource(channel, List.of("fleet01-lead", "fleet02-lead", "fleet03-lead"))); + + String out = textOf(res); + assertTrue(out.contains("\"coordId\":\"fleet01-lead\",\"status\":\"exists\",\"pending\":2,\"consumers\":1"), out); + assertTrue(out.contains("\"coordId\":\"fleet02-lead\",\"status\":\"absent\""), out); + assertTrue(out.contains("\"coordId\":\"fleet03-lead\",\"status\":\"unknown\""), out); + assertFalse(out.contains("\"coordId\":\"fleet02-lead\",\"status\":\"absent\",\"pending\""), + "pending/consumers must be omitted, not faked as zero, for a confirmed-absent peer: " + out); + assertFalse(out.contains("\"coordId\":\"fleet03-lead\",\"status\":\"unknown\",\"pending\""), + "pending/consumers must be omitted, not faked as zero, for an unresolved peer probe: " + out); + } + + /** + * fleetd #361 review finding 2: {@code get(timeout)} alone times out the CALLER but leaves the + * submitted {@link LeadChannel#inspect} task running forever on its own virtual thread — against + * a hung (not down) broker every probe would orphan one more thread holding an AMQP channel + * until the connection's channel-max is exhausted, which would break {@code publish} too. This + * proves {@link FleetMcp#probe(LeadChannel, String, long)} does not merely give up on a slow + * task: it interrupts it, so the task does not go on running unbounded after the caller has + * already moved on. No hung broker needed — a {@link LeadChannel} fake that blocks until + * interrupted is enough to observe the same mechanism. + */ + @Test + void aTimedOutProbeInterruptsTheOrphanedTaskRatherThanAbandoningIt() throws Exception { + CountDownLatch started = new CountDownLatch(1); + AtomicBoolean wasInterrupted = new AtomicBoolean(false); + LeadChannel hangs = new LeadChannel() { + @Override + public void publish(String toCoordId, LeadMessage m) { } + + @Override + public List peek() { return List.of(); } + + @Override + public void ack(String msgId) { } + + @Override + public String selfCoordId() { return "mac-opus"; } + + @Override + public MailboxState inspect(String coordId) { + started.countDown(); + try { + Thread.sleep(60_000); + } catch (InterruptedException e) { + wasInterrupted.set(true); + Thread.currentThread().interrupt(); + } + return MailboxState.unknown(coordId); + } + }; + + LeadChannel.MailboxState result = FleetMcp.probe(hangs, "fleet01-lead", 100L); + + assertFalse(result.exists(), "a timed-out probe must never claim the mailbox exists"); + assertFalse(result.known(), "a timed-out probe proves nothing either way — it must report unknown"); + assertTrue(started.await(2, TimeUnit.SECONDS), "the probe task must actually have started"); + // The interrupt is delivered asynchronously to the orphaned task's own thread — poll briefly + // rather than assume it has already landed the instant probe() returns. + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(2); + while (!wasInterrupted.get() && System.nanoTime() < deadline) { + Thread.sleep(20); + } + assertTrue(wasInterrupted.get(), + "probe() must cancel the orphaned task (interrupt it) instead of leaving it to run forever"); + } + @Test void capacityUsesThePlacementLiveCount() { FakeHerdr h = new FakeHerdr(); @@ -911,7 +1051,8 @@ class FleetMcpTest { String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null, new FleetMcp.CapacitySource(profile -> 2, profile -> 3, () -> Set.of("sonnet"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"), - FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", null)); + FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", + FleetMcp.CoordinationSource.none())); assertTrue(out.contains("\"maxLoad\":3"), "maxLoad itself must be left untouched: " + out); assertTrue(out.contains("\"live\":2"), out); @@ -934,7 +1075,8 @@ class FleetMcpTest { String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null, new FleetMcp.CapacitySource(profile -> 0, profile -> 3, () -> Set.of("sonnet"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"), - FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", null)); + FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", + FleetMcp.CoordinationSource.none())); assertTrue(out.contains("\"live\":0"), out); assertTrue(out.contains("\"free\":3"), "the real gate never subtracts the lead's seat: " + out); @@ -988,7 +1130,7 @@ class FleetMcpTest { String out = textOf(FleetMcp.listFleet(composite, sm, null, new FleetMcp.CapacitySource(liveCount, p -> profiles.get(p).maxLoad(), profiles::keySet, () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"), FleetMcp.QuarantineSource.none(), - FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", null)); + FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", FleetMcp.CoordinationSource.none())); int reportedFree = extractInt(out, "free"); for (int i = 0; i < reportedFree; i++) { @@ -1017,7 +1159,7 @@ class FleetMcpTest { sessions, null, new FleetMcp.CapacitySource(profile -> 0, profile -> 2, () -> Set.of("terra"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"), FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(), - Map.of(), "", null)); + Map.of(), "", FleetMcp.CoordinationSource.none())); assertTrue(out.contains("\"free\":2"), out); assertFalse(out.contains("leadSeats"), "no lead shares this profile's credential: " + out); diff --git a/fleetd/src/test/java/dev/ltms/fleet/member/MemberEnvAllowListTest.java b/fleetd/src/test/java/dev/ltms/fleet/member/MemberEnvAllowListTest.java index d579419..ae8d64b 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/member/MemberEnvAllowListTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/member/MemberEnvAllowListTest.java @@ -146,7 +146,7 @@ class MemberEnvAllowListTest { return new FleetConfig(null, null, null, Map.of(), null, null, null, null, null, new FleetConfig.Broker(null, brokerUriEnv, null), null, null, null, null, null, null, null, null, null, - new FleetConfig.Coordinator(null, coordinatorUriEnv, null, null)).withDefaults(); + new FleetConfig.Coordinator(null, coordinatorUriEnv, null, null, null)).withDefaults(); } /** {@code LC_*} categories are infrastructure by prefix; everything else needs an exact match. */ diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/FakeLeadChannel.java b/fleetd/src/test/java/dev/ltms/fleet/msg/FakeLeadChannel.java index 38ee2e9..b603a79 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/FakeLeadChannel.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/FakeLeadChannel.java @@ -3,6 +3,8 @@ package dev.ltms.fleet.msg; import java.util.ArrayList; import java.util.Collections; import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; /** * Hermetic stand-in for {@link LeadChannel}: an in-memory mailbox that records what was published @@ -24,11 +26,19 @@ public final class FakeLeadChannel implements LeadChannel { private final List acked = Collections.synchronizedList(new ArrayList<>()); /** When set, every {@link #publish} throws it — the unroutable/nacked/timed-out peer. */ private volatile IllegalStateException publishFailure; + /** Canned {@link #inspect} results by coord-id — absent for any coord-id not configured here. */ + private final Map mailboxes = new ConcurrentHashMap<>(); public FakeLeadChannel(String selfCoordId) { this.selfCoordId = selfCoordId; } + /** Make {@link #inspect(String)} return {@code state} for {@code coordId} instead of "absent". */ + public FakeLeadChannel withMailbox(String coordId, MailboxState state) { + mailboxes.put(coordId, state); + return this; + } + /** Make every publish fail as an unreachable peer would. */ public FakeLeadChannel failPublishWith(String message) { this.publishFailure = new IllegalStateException(message); @@ -65,6 +75,11 @@ public final class FakeLeadChannel implements LeadChannel { return selfCoordId; } + @Override + public MailboxState inspect(String coordId) { + return mailboxes.getOrDefault(coordId, MailboxState.absent(coordId)); + } + public List published() { return List.copyOf(published); } diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/LeadMailboxIsMissingQueueTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/LeadMailboxIsMissingQueueTest.java new file mode 100644 index 0000000..aa2f50e --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/LeadMailboxIsMissingQueueTest.java @@ -0,0 +1,77 @@ +package dev.ltms.fleet.msg; + +import com.rabbitmq.client.ShutdownSignalException; +import com.rabbitmq.client.impl.AMQImpl; +import org.junit.jupiter.api.Test; + +import java.io.IOException; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * fleetd #361 review round 2: a mutation that made {@link LeadMailbox#isMissingQueue} return + * {@code true} unconditionally still left {@code mvn clean install} green — 1389 tests, 0 + * failures — because nothing exercised its false branch. That branch is the whole discriminator + * between {@link LeadChannel.MailboxState#absent} and {@link LeadChannel.MailboxState#unknown}; + * without a test pinning it, a future refactor that widens it back to "always true" (restoring the + * exact overstatement fleetd #361 exists to fix) would pass this suite. + * + *

Hermetic — no broker needed, per the review's own suggestion. {@code isMissingQueue} takes a + * plain {@link IOException}, so every input here is constructed directly rather than provoked from + * a live connection. The real 404 shape itself is still pinned against a real broker, in + * {@code LeadMailboxTest.passiveDeclareOfAMissingQueueThrowsAnIOExceptionWrappingA404ShutdownSignal} + * — this class covers the three false shapes {@link LeadMailbox#isMissingQueue}'s own javadoc + * lists, so both directions of the discriminator are proven somewhere. + */ +class LeadMailboxIsMissingQueueTest { + + @Test + void aConfirmedMissingQueueIsRecognized() { + ShutdownSignalException sse = new ShutdownSignalException(true, false, + new AMQImpl.Channel.Close(404, "NOT_FOUND - no queue 'lead.x.inbox' in vhost '/'", 50, 10), null); + IOException e = new IOException("channel error", sse); + + assertTrue(LeadMailbox.isMissingQueue(e), "a genuine 404 Channel.Close must be recognized as a missing queue"); + } + + @Test + void aDifferentReplyCodeIsNotAMissingQueue() { + // E.g. 403 ACCESS_REFUSED — the queue may well exist; this call was simply refused. + ShutdownSignalException sse = new ShutdownSignalException(true, false, + new AMQImpl.Channel.Close(403, "ACCESS_REFUSED", 50, 10), null); + IOException e = new IOException("channel error", sse); + + assertFalse(LeadMailbox.isMissingQueue(e), + "a non-404 reply code must never be read as a confirmed absence — the mailbox's real state is unknown"); + } + + @Test + void aShutdownSignalWhoseReasonIsNotAChannelCloseIsNotAMissingQueue() { + // A Connection.Close (a whole different broker-level shutdown) is still a ShutdownSignalException, + // but its reason is not a Channel.Close at all — must not be misread as "no such queue". + ShutdownSignalException sse = new ShutdownSignalException(true, false, + new AMQImpl.Connection.Close(404, "coincidentally 404, but this is a CONNECTION close", 10, 50), null); + IOException e = new IOException("connection error", sse); + + assertFalse(LeadMailbox.isMissingQueue(e), + "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"); + } + + @Test + void anIOExceptionWithNoCauseAtAllIsNotAMissingQueue() { + IOException e = new IOException("some other declare failure, no cause attached"); + + assertFalse(LeadMailbox.isMissingQueue(e), + "an IOException with no ShutdownSignalException cause must never be read as a confirmed absence"); + } + + @Test + void anIOExceptionWithAnUnrelatedCauseIsNotAMissingQueue() { + IOException e = new IOException("wrapped something else entirely", new RuntimeException("boom")); + + assertFalse(LeadMailbox.isMissingQueue(e), + "a cause that isn't even a ShutdownSignalException must never be read as a confirmed absence"); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/LeadMailboxTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/LeadMailboxTest.java index bc015ac..ccc5144 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/LeadMailboxTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/LeadMailboxTest.java @@ -1,5 +1,9 @@ package dev.ltms.fleet.msg; +import com.rabbitmq.client.AMQP; +import com.rabbitmq.client.Channel; +import com.rabbitmq.client.Connection; +import com.rabbitmq.client.ShutdownSignalException; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Tag; import org.junit.jupiter.api.Test; @@ -7,11 +11,14 @@ import org.testcontainers.containers.RabbitMQContainer; import org.testcontainers.junit.jupiter.Testcontainers; import org.testcontainers.utility.DockerImageName; +import java.io.IOException; import java.util.List; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -159,6 +166,145 @@ class LeadMailboxTest { } } + @Test + void inspectReportsAnOwnedMailboxAsExistingWithItsOwnConsumer() throws Exception { + // A LeadMailbox declares AND consumes its own queue the moment open() returns (see own()), + // so inspecting a coord-id this same process owns must always find exactly one consumer. + String self = coordId("lead-inspect-self"); + try (LeadMailbox mailbox = LeadMailbox.open(uri(), self)) { + LeadChannel.MailboxState state = mailbox.inspect(self); + assertEquals(self, state.coordId()); + assertTrue(state.exists(), "this daemon owns and has declared this exact queue"); + assertEquals(0, state.pending(), "nothing has been published to it yet"); + assertEquals(1, state.consumers(), "the mailbox's own constructor already attached a consumer"); + } + } + + @Test + void inspectReportsAMissingMailboxAsAbsentRatherThanThrowing() throws Exception { + String nobody = coordId("lead-inspect-nobody"); + try (LeadMailbox mailbox = LeadMailbox.open(uri(), coordId("lead-inspect-caller"))) { + LeadChannel.MailboxState state = mailbox.inspect(nobody); + assertEquals(LeadChannel.MailboxState.absent(nobody), state, + "a queue nobody has ever declared must report absent, never throw"); + assertTrue(state.known(), "a confirmed 404 IS a definite answer — this is not the unknown case"); + } + } + + /** + * fleetd #361 review finding 3: {@code inspect} is specified to never throw, but the original + * implementation caught only {@link IOException} — and {@link Connection#createChannel()} on an + * already-closed connection throws {@link com.rabbitmq.client.AlreadyClosedException}, an + * unchecked {@link RuntimeException} (pinned by {@code + * createChannelOnAnAlreadyClosedConnectionThrowsAnUncheckedException} above). This drives that + * exact scenario through the real {@link LeadMailbox#inspect} — not the raw client call — and + * checks both halves of finding 1 and finding 3 at once: no exception escapes, and the result is + * {@code UNKNOWN} rather than the wrong-but-plausible-looking {@code ABSENT}. + */ + @Test + void inspectReportsUnknownRatherThanThrowingWhenTheConnectionIsAlreadyClosed() throws Exception { + LeadMailbox mailbox = LeadMailbox.open(uri(), coordId("lead-inspect-dead-connection")); + mailbox.close(); // tears down the connection `inspect` will try to open a probe channel on + + LeadChannel.MailboxState state = mailbox.inspect(coordId("lead-inspect-irrelevant-target")); + + assertFalse(state.exists()); + assertFalse(state.known(), "a dead connection proves nothing about the target mailbox — it must be unknown, not absent"); + assertEquals(LeadChannel.MailboxState.Presence.UNKNOWN, state.presence()); + } + + @Test + void inspectReportsPendingMessagesAndZeroConsumersWhenNobodyIsReadingAnymore() throws Exception { + // Publish into a mailbox this test owns, then never consume from it, to prove `pending` and + // `consumers` really come off the broker rather than off this process's own in-memory state. + String to = coordId("lead-inspect-pending"); + String observerId = coordId("lead-inspect-observer"); + try (LeadMailbox owner = LeadMailbox.open(uri(), to); + LeadMailbox observer = LeadMailbox.open(uri(), observerId)) { + owner.publish(to, new LeadMessage("m1", "lead-from", to, "sitting in the queue")); + awaitPeek(owner); // make sure the broker has actually enqueued it before inspecting + } // `owner` closes here: its consumer disconnects, but the durable, unacked message stays queued. + + try (LeadMailbox observer = LeadMailbox.open(uri(), coordId("lead-inspect-observer-2"))) { + // The broker requeues `owner`'s unacked delivery asynchronously once its connection drops, + // so poll rather than assume the very first passive declare already sees the settled state. + LeadChannel.MailboxState state = awaitInspect(observer, to, s -> s.consumers() == 0); + assertTrue(state.exists()); + assertEquals(0, state.consumers(), "the only owner just closed — nobody is reading this anymore"); + assertEquals(1, state.pending(), "the unacked message must be requeued, never dropped"); + } + } + + /** + * fleetd #361's central invariant, proved rather than assumed: a passive queue declare of a + * missing queue closes ITS channel with a 404 in AMQP 0-9-1. {@link LeadMailbox#inspect} is + * specified to run on its own disposable channel for exactly this reason — this test is the one + * that actually exercises the failure mode and shows {@link LeadMailbox#publish} on the SAME + * instance is unaffected by it. + */ + @Test + void inspectingAMissingMailboxNeverBreaksPublishOnTheSameInstance() throws Exception { + String self = coordId("lead-invariant-self"); + try (LeadMailbox mailbox = LeadMailbox.open(uri(), self)) { + // Miss on a queue that has never existed — this is exactly the 404-closes-the-channel case. + LeadChannel.MailboxState missed = mailbox.inspect(coordId("lead-invariant-nobody-home")); + assertFalse(missed.exists()); + assertTrue(missed.known(), "a genuine 404 on a queue that never existed is a confirmed fact, not an unknown"); + + // publish() must still work on THIS SAME instance: if inspect() had reused `publishChannel` + // (or `channel`), the broker's 404 would have closed it out from underneath publish(). + LeadMessage sent = new LeadMessage("after-miss", "lead-from", self, "still alive"); + mailbox.publish(self, sent); + List got = awaitPeek(mailbox); + assertEquals(1, got.size(), "publish must still reach this mailbox's own queue after a missed inspect"); + assertEquals("after-miss", got.getFirst().msgId()); + + // And a second inspect() — of a mailbox that DOES exist this time — must also still work, + // proving the miss did not wedge inspect() itself either. + LeadChannel.MailboxState self2 = mailbox.inspect(self); + assertTrue(self2.exists()); + } + } + + /** + * Pins the exact exception shape {@link LeadMailbox#inspect} relies on to tell a genuine 404 + * (mailbox confirmed absent) apart from everything else (mailbox state unknown) — measured + * against a real broker rather than assumed from the AMQP 0-9-1 spec text. If this ever fails, + * the classification in {@code inspect} is reading the wrong shape and must be revisited. + */ + @Test + void passiveDeclareOfAMissingQueueThrowsAnIOExceptionWrappingA404ShutdownSignal() throws Exception { + try (Connection conn = LeadMailbox.connectionFactory(uri()).newConnection()) { + Channel probe = conn.createChannel(); + String missing = LeadMailbox.queueName(coordId("lead-404-shape")); + IOException thrown = assertThrows(IOException.class, () -> probe.queueDeclarePassive(missing)); + assertInstanceOf(ShutdownSignalException.class, thrown.getCause(), + () -> "expected the IOException to wrap a ShutdownSignalException, got: " + thrown); + ShutdownSignalException sse = (ShutdownSignalException) thrown.getCause(); + assertInstanceOf(AMQP.Channel.Close.class, sse.getReason(), + () -> "expected a Channel.Close reason: " + sse); + AMQP.Channel.Close close = (AMQP.Channel.Close) sse.getReason(); + assertEquals(404, close.getReplyCode(), () -> "expected AMQP NOT_FOUND (404): " + close); + assertFalse(probe.isOpen(), "the 404 must have closed the channel the declare ran on"); + } + } + + /** + * The other half of the same measurement: calling {@code createChannel()} on an + * already-closed connection — the shape {@link LeadMailbox#inspect} hits when the broker + * connection itself is gone — throws {@link com.rabbitmq.client.AlreadyClosedException}, an + * unchecked {@link RuntimeException}, not an {@link IOException}. An {@code inspect} that only + * caught {@code IOException} here would let this escape instead of reporting "unknown". + */ + @Test + void createChannelOnAnAlreadyClosedConnectionThrowsAnUncheckedException() throws Exception { + Connection conn = LeadMailbox.connectionFactory(uri()).newConnection(); + conn.close(); + RuntimeException thrown = assertThrows(RuntimeException.class, conn::createChannel); + assertInstanceOf(com.rabbitmq.client.AlreadyClosedException.class, thrown, + () -> "expected AlreadyClosedException, got: " + thrown); + } + /** Poll peek until at least one message is held, or ~10s elapse (broker delivery is async). */ @SuppressWarnings("BusyWait") private static List awaitPeek(LeadMailbox inbox) throws InterruptedException { @@ -170,4 +316,18 @@ class LeadMailboxTest { } return msgs; } + + /** Poll inspect(coordId) until it satisfies {@code done}, or ~10s elapse (broker state settles async). */ + @SuppressWarnings("BusyWait") + private static LeadChannel.MailboxState awaitInspect( + LeadMailbox observer, String coordId, java.util.function.Predicate done) + throws InterruptedException { + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10); + LeadChannel.MailboxState state = observer.inspect(coordId); + while (!done.test(state) && System.nanoTime() < deadline) { + Thread.sleep(50); + state = observer.inspect(coordId); + } + return state; + } }