fleetd #361 review: fix false-negative absent, uncancelled probes, and a throw contract gap
Three findings from review of #364, fixed on the same branch: 1. LeadChannel.MailboxState.absent() was returned both for a genuinely absent mailbox AND for "the probe could not determine anything" (timeout, unreachable broker, other declare failure) -- exactly the overstatement #361 exists to fix, one level down. MailboxState now carries a Presence enum (EXISTS/ABSENT/UNKNOWN) with exists()/known() accessors; LeadMailbox.inspect classifies a real AMQP 404 (measured against a live broker, not assumed: an IOException wrapping a ShutdownSignalException whose Channel.Close reply code is 404) as ABSENT and everything else as UNKNOWN. fleet_list's mailbox/peer rows now render a "status" of exists/absent/unknown and only include pending/consumers when status is "exists", so an unresolved self- or peer-probe can never render as a measured zero. 2. FleetMcp.probe's get(timeoutMs) left a timed-out inspect() task running forever on its own virtual thread, holding the AMQP channel it had already opened -- against a hung (not down) broker this would orphan one channel per fleet_list call until the connection's channel-max was exhausted, breaking publish() too. probe() now holds the Future and calls cancel(true) on timeout/failure so the orphaned task is interrupted instead of abandoned, and now returns MailboxState.unknown() (never absent()) on timeout/exception. 3. LeadMailbox.inspect only caught IOException, but createChannel() on an already-closed connection throws AlreadyClosedException, an unchecked RuntimeException (measured against a live broker) -- so it could escape the "never throws" contract. Both places in inspect now also catch RuntimeException and report unknown(). Tests: MailboxState.exists()/absent()/unknown() call sites updated across FleetMcpTest/FleetMcpLeadCoordTest; new hermetic tests cover the tri-state fleet_list rendering (self-probe unknown, a peer that's absent vs. one that's unknown) and probe cancellation (a LeadChannel fake that blocks until interrupted, proving probe() doesn't just give up on it); new @Tag("contract") LeadMailboxTest cases pin the real exception shapes for both the 404 and the already-closed-connection paths and prove inspect() reports unknown (never throws) when the connection is already closed.
This commit is contained in:
@@ -1273,16 +1273,32 @@ public final class FleetMcp {
|
||||
Map<String, Object> row = new LinkedHashMap<>();
|
||||
row.put("selfId", selfId);
|
||||
row.put("configured", true);
|
||||
LeadChannel.MailboxState self = probe(channel, selfId);
|
||||
Map<String, Object> 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<String, Object> mailboxView(LeadChannel.MailboxState state) {
|
||||
Map<String, Object> 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<String, Object> heldView(LeadMessage m) {
|
||||
Map<String, Object> 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<String, Object> peerView(LeadChannel channel, String coordId) {
|
||||
LeadChannel.MailboxState state = probe(channel, coordId);
|
||||
Map<String, Object> 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 <em>bounded</em> 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.
|
||||
*
|
||||
* <p><strong>A timeout cancels the orphaned task</strong> 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<LeadChannel.MailboxState> 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);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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.
|
||||
*
|
||||
* <p>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.
|
||||
* <p><strong>Never throws</strong> — 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".
|
||||
*
|
||||
* <p><strong>Must never share fate with {@link #publish} or {@link #peek}/{@link #ack}.</strong>
|
||||
* 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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.
|
||||
*
|
||||
* <p>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:
|
||||
* <ul>
|
||||
* <li>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;
|
||||
* <li>{@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.
|
||||
* </ul>
|
||||
*/
|
||||
@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<LeadMessage> drain() {
|
||||
List<LeadMessage> snapshot = peek();
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<LeadMessage> 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
|
||||
|
||||
@@ -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<LeadMessage> awaitPeek(LeadMailbox inbox) throws InterruptedException {
|
||||
|
||||
Reference in New Issue
Block a user