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 6685105..16e19a8 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -1397,11 +1397,11 @@ public final class FleetMcp { * {@code mailbox.pending} counts only broker-ready messages; a held message is already * an unacked delivery sitting with this consumer, so the normal, healthy state of a blocked lead * is {@code "pending": 0} next to a non-empty {@code held[]} — which invites the false reading - * "these are only in memory, a restart will lose them". They are not: {@code LeadMailbox} - * consumes with manual ack, so held mail is a durable broker delivery. {@code heldCount} is the - * honest second number beside {@code pending} ({@code held.size()}, not left for the reader to - * count the array), and {@code heldDurable} states the fact in words rather than leaving - * {@code pending} as the only number next to {@code held[]}. + * "these are only in memory, a restart will lose them". {@code heldCount} is the honest second + * number beside {@code pending} ({@code held.size()}, not left for the reader to count the + * array). {@code heldDurable} comes straight from {@link LeadChannel#heldDurable}, which the + * channel implementation derives from what it actually did when it declared and consumed its own + * queue (fleetd #440) — this method never asserts the fact itself. */ private static Map coordinatorView(CoordinationSource coordination) { LeadChannel channel = coordination.leadChannel(); @@ -1415,7 +1415,7 @@ public final class FleetMcp { row.put("configured", true); row.put("mailbox", mailboxView(probe(channel, selfId))); row.put("heldCount", held.size()); - row.put("heldDurable", true); + row.put("heldDurable", channel.heldDurable()); row.put("held", held.stream().map(FleetMcp::heldView).toList()); row.put("peers", coordination.peers().stream().map(p -> peerView(channel, p)).toList()); return row; 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 d7c8931..653aa44 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadChannel.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadChannel.java @@ -42,6 +42,18 @@ public interface LeadChannel { /** This daemon's own lead coordination id — the mailbox it owns, and the {@code from} it sends as. */ String selfCoordId(); + /** + * Whether a message sitting in {@link #peek}'s held set (fetched but not yet {@link #ack}ed) is + * still safe if this daemon crashes or restarts right now — the conclusion of two independent + * facts about how this channel owns its own queue: the queue was declared durable, and + * the consumer that filled {@code held} uses manual ack, so an unacked delivery is still + * owned by the broker rather than only in this process's memory. Both must hold for {@code true}; + * an implementation must derive this from what it actually did when it declared and consumed its + * queue, never return a literal — fleetd #440 found {@code FleetMcp}'s {@code heldDurable} field + * doing exactly that, unable to ever report {@code false} even after the fact stopped being true. + */ + boolean heldDurable(); + /** * 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 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 a9283f1..e57a6ad 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java @@ -86,6 +86,12 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable { private final Object channelLock = new Object(); /** msgId → held delivery, for this mailbox's own queue only (there is exactly one). */ private final LinkedHashMap held = new LinkedHashMap<>(); + /** + * fleetd #440: the answer to {@link #heldDurable()}, set once by {@link #own()} from the exact + * booleans it passed to {@code queueDeclare}/{@code basicConsume} — never a separate literal that + * could drift from what those calls actually did. + */ + private boolean heldDurable; /** Successful broker acks on this connection, retained only to make a repeated caller ack quiet. */ private final LinkedHashMap recentlyAcked = new LinkedHashMap<>(); /** Bounds {@link #recentlyAcked}: it is only an idempotency aid, never delivery state. */ @@ -191,13 +197,22 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable { /** Declare + consume this daemon's own {@code lead..inbox}. Called once, at construction. */ private void own() throws IOException { String queue = queueName(selfCoordId); + boolean durableQueue = true; // durable, non-exclusive, keep on idle + boolean autoAck = false; // manual ack synchronized (channelLock) { - channel.queueDeclare(queue, true, false, false, null); // durable, non-exclusive, keep on idle - channel.basicConsume(queue, false, deliverCallback(), _ -> { }); // autoAck=false: manual ack + channel.queueDeclare(queue, durableQueue, false, false, null); + channel.basicConsume(queue, autoAck, deliverCallback(), _ -> { }); } + // fleetd #440: held mail is durable only while both hold — a durable queue AND manual ack. + this.heldDurable = durableQueue && !autoAck; log.debug("lead mailbox owns queue {} for coord-id {}", queue, selfCoordId); } + @Override + public boolean heldDurable() { + return heldDurable; + } + /** * Publish {@code msg} to {@code toCoordId}'s mailbox and block until the broker's publisher * confirm for it lands. Does not imply owning or consuming {@code toCoordId}'s queue. 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 e55bdbd..2186371 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java @@ -734,6 +734,32 @@ class FleetMcpTest { "must state the durability fact, not leave pending as the only number next to held[]: " + out); } + /** + * fleetd #440: {@code heldDurable} must be a derived fact, not a literal — so it can report + * {@code false} when the channel behind it says held mail is not durable (a non-durable queue, + * or a consumer running with {@code autoAck=true}). A test that only ever asserts {@code true} + * repeats the defect this ticket fixes. + */ + @Test + void listReportsHeldDurableFalseWhenTheChannelSaysMailIsNotDurable() { + 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)) + .withHeldDurable(false) + .hold(new LeadMessage("m1", "fleet01-lead", "mac-opus", "one")); + + 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("\"heldDurable\":false"), + "heldDurable must follow the channel, not a hardcoded true: " + out); + } + // ── fleetd #421: a lead reads (never consumes) its own held peer mail ────────────────────── @Test @@ -834,6 +860,9 @@ class FleetMcpTest { @Override public String selfCoordId() { return "mac-opus"; } + @Override + public boolean heldDurable() { return true; } + @Override public MailboxState inspect(String coordId) { started.countDown(); 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 b603a79..f631483 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/FakeLeadChannel.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/FakeLeadChannel.java @@ -28,11 +28,19 @@ public final class FakeLeadChannel implements LeadChannel { 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<>(); + /** fleetd #440: matches {@link LeadMailbox}'s real default (durable queue + manual ack) unless overridden. */ + private volatile boolean heldDurable = true; public FakeLeadChannel(String selfCoordId) { this.selfCoordId = selfCoordId; } + /** Make {@link #heldDurable()} report {@code durable} — the fleetd #440 seam for the false case. */ + public FakeLeadChannel withHeldDurable(boolean durable) { + this.heldDurable = durable; + return this; + } + /** Make {@link #inspect(String)} return {@code state} for {@code coordId} instead of "absent". */ public FakeLeadChannel withMailbox(String coordId, MailboxState state) { mailboxes.put(coordId, state); @@ -80,6 +88,11 @@ public final class FakeLeadChannel implements LeadChannel { return mailboxes.getOrDefault(coordId, MailboxState.absent(coordId)); } + @Override + public boolean heldDurable() { + return heldDurable; + } + public List published() { return List.copyOf(published); } 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 6998b47..e0c36ae 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/LeadMailboxTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/LeadMailboxTest.java @@ -207,6 +207,21 @@ class LeadMailboxTest { } } + /** + * fleetd #440: {@code heldDurable()} must be derived from what {@link LeadMailbox#own} actually + * did against the real broker — a durable queue declare plus a manual-ack consumer — not a + * hardcoded literal. This is the mutation-sensitive test: flip {@code own()}'s {@code autoAck} + * local to {@code true} (or its {@code durableQueue} local to {@code false}) and this must fail. + */ + @Test + void heldDurableReportsTrueBecauseTheQueueIsDurableAndTheConsumeIsManualAck() throws Exception { + String self = coordId("lead-held-durable"); + try (LeadMailbox mailbox = LeadMailbox.open(uri(), self)) { + assertTrue(mailbox.heldDurable(), + "own() declares a durable queue and consumes with autoAck=false, so held mail is durable"); + } + } + @Test void inspectReportsAMissingMailboxAsAbsentRatherThanThrowing() throws Exception { String nobody = coordId("lead-inspect-nobody");