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 8609276..7b35930 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -1273,16 +1273,32 @@ public final class FleetMcp { Map row = new LinkedHashMap<>(); row.put("selfId", selfId); row.put("configured", true); - LeadChannel.MailboxState self = probe(channel, selfId); - Map mailbox = new LinkedHashMap<>(); - mailbox.put("pending", self.pending()); - mailbox.put("consumers", self.consumers()); - row.put("mailbox", mailbox); + 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<>(); @@ -1304,16 +1320,17 @@ public final class FleetMcp { : content.substring(0, HELD_PREVIEW_MAX_CHARS) + "…"; } - /** One declared peer's row: its coord-id, whether its mailbox exists, and — only then — its counts. */ + /** + * 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) { - LeadChannel.MailboxState state = probe(channel, coordId); Map row = new LinkedHashMap<>(); row.put("coordId", coordId); - row.put("reachable", state.exists()); - if (state.exists()) { - row.put("pending", state.pending()); - row.put("consumers", state.consumers()); - } + row.putAll(mailboxView(probe(channel, coordId))); return row; } @@ -1325,7 +1342,10 @@ public final class FleetMcp { /** * 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}. + * 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()); @@ -1333,17 +1353,35 @@ public final class FleetMcp { /** * 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#absent} 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. + * 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 PEER_PROBE_POOL.submit(() -> channel.inspect(coordId)) - .get(PEER_PROBE_TIMEOUT_MS, TimeUnit.MILLISECONDS); + return future.get(timeoutMs, TimeUnit.MILLISECONDS); } catch (Exception e) { - return LeadChannel.MailboxState.absent(coordId); + future.cancel(true); // best-effort: don't leave a hung probe (and its channel) running forever + return LeadChannel.MailboxState.unknown(coordId); } } 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 7bdef7b..db16aa3 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadChannel.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadChannel.java @@ -44,9 +44,14 @@ public interface LeadChannel { * 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. * - *

Returns {@link MailboxState#absent(String)}, never throws, when {@code coordId}'s mailbox - * does not exist or the look otherwise fails (broker unreachable, timed out) — this is a - * best-effort fact-finding call, not an operation a caller must handle failing. + *

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 @@ -58,19 +63,52 @@ public interface LeadChannel { MailboxState inspect(String coordId); /** - * The result of {@link #inspect}. {@code exists} is {@code false} for a mailbox nobody has ever - * declared (or when the look could not be completed) — {@code pending}/{@code consumers} are - * meaningless in that case and always read {@code 0}. + * 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 exists whether the mailbox's queue is currently declared on the broker - * @param pending messages ready for delivery but not yet in a consumer's hands (0 when absent) - * @param consumers how many consumers are attached (0 when absent, or when nobody is reading it) + * @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, boolean exists, int pending, int consumers) { - /** The mailbox does not exist, or the look could not be completed. */ + 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, false, 0, 0); + 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 335c3e2..f9da75e 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; @@ -265,6 +266,20 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable { * 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. + *
*/ @Override public MailboxState inspect(String coordId) { @@ -272,17 +287,22 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable { Channel probe; try { probe = connection.createChannel(); - } catch (IOException e) { + } catch (IOException | RuntimeException e) { log.debug("lead mailbox inspect: cannot open a probe channel for {}: {}", coordId, e.toString()); - return MailboxState.absent(coordId); + return MailboxState.unknown(coordId); } try { AMQP.Queue.DeclareOk declared = probe.queueDeclarePassive(queue); - return new MailboxState(coordId, true, declared.getMessageCount(), declared.getConsumerCount()); + return MailboxState.exists(coordId, declared.getMessageCount(), declared.getConsumerCount()); } catch (IOException e) { - // Missing queue (404) or any other declare failure: the broker (or the client library) - // has already closed `probe` for us — report "does not exist" rather than throw. - return MailboxState.absent(coordId); + // 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()) { @@ -294,6 +314,22 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable { } } + /** + * {@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, no {@code ShutdownSignalException} cause, or none at all) is a declare failure of + * some other kind and must not be read as "confirmed absent". + */ + private 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/mcp/FleetMcpLeadCoordTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpLeadCoordTest.java index d0986ab..2384c1c 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpLeadCoordTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpLeadCoordTest.java @@ -103,7 +103,7 @@ class FleetMcpLeadCoordTest { @Test void warnsWhenThePeerMailboxHasNoConsumersButStillReportsSuccess() { var channel = new FakeLeadChannel(SELF) - .withMailbox(PEER, new LeadChannel.MailboxState(PEER, true, 0, 0)); + .withMailbox(PEER, LeadChannel.MailboxState.exists(PEER, 0, 0)); McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", null, null); @@ -116,7 +116,7 @@ class FleetMcpLeadCoordTest { @Test void staysQuietAboutConsumersWhenThePeerMailboxHasOne() { var channel = new FakeLeadChannel(SELF) - .withMailbox(PEER, new LeadChannel.MailboxState(PEER, true, 0, 1)); + .withMailbox(PEER, LeadChannel.MailboxState.exists(PEER, 0, 1)); McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", null, null); @@ -125,4 +125,25 @@ class FleetMcpLeadCoordTest { 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 1c90883..554d767 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java @@ -37,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.*; @@ -580,7 +582,7 @@ class FleetMcpTest { 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", new LeadChannel.MailboxState("mac-opus", true, 0, 1)); + .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, @@ -593,11 +595,36 @@ class FleetMcpTest { // 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\":{\"pending\":0,\"consumers\":1}"), 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 void listOmitsTheCoordinatorRowWhenLeadCoordinationIsOff() { FakeHerdr h = new FakeHerdr(); @@ -636,21 +663,82 @@ class FleetMcpTest { 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", new LeadChannel.MailboxState("fleet01-lead", true, 2, 1)); + .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 peer would report. + // 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"))); + new FleetMcp.CoordinationSource(channel, List.of("fleet01-lead", "fleet02-lead", "fleet03-lead"))); String out = textOf(res); - assertTrue(out.contains("\"coordId\":\"fleet01-lead\",\"reachable\":true,\"pending\":2,\"consumers\":1"), out); - assertTrue(out.contains("\"coordId\":\"fleet02-lead\",\"reachable\":false"), out); - assertFalse(out.contains("\"coordId\":\"fleet02-lead\",\"reachable\":false,\"pending\""), - "pending/consumers must be omitted, not faked as zero, for an unreachable peer: " + out); + 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 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 6bf4915..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,12 +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; @@ -181,9 +187,32 @@ class LeadMailboxTest { 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 @@ -220,6 +249,7 @@ class LeadMailboxTest { // 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(). @@ -236,6 +266,45 @@ class LeadMailboxTest { } } + /** + * 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 {