From 1e9b2c9b7eb9cdc8ce958643bdf94a72fad4abb6 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 10 Sep 2026 13:54:50 +0700 Subject: [PATCH] fleetd #421: let a lead peek its own held peer mail, primary-only fleet_list truncated held lead-to-lead messages to an 80-char preview with no way to read the full body, and fleet_poll{target} drained the wrong inbox (a worker's reply queue, not the coordinator mailbox) -- it silently returned []. fleet_ack would have destroyed the message unread. Add a non-destructive read: fleet_poll{coordId} peeks (never acks) this daemon's own held mail via LeadChannel.peek(). The coordId must equal the caller's own selfCoordId -- passing a peer's id is refused with a reason, instead of repeating the original silent-[] confusion. This is authorization-sensitive: mapping it to the existing READ action would let any worker read every peer lead's mail in full. READ's openness rests on "the roster carries no secrets" (Authz.java), which does not hold for lead-to-lead coordination bodies. Added Authz.Action.COORD_READ, primary-only (not even the architect, which holds READ today), and made pollAction's signature depend on both target and coordId so every call site states explicitly what it passes. Also fixes fleet_list's "pending: 0" trap: mailbox.pending only counts broker-ready messages, so a healthy held mailbox reads as empty. Added heldCount/heldDurable beside held[] so the durability fact isn't implied only by reading the code. Mutation-tested: pollAction's COORD_READ->READ mapping, the peek->ack substitution, the 80-char preview cap widened to 81, and the Authz case widened to include caller.isWorker() -- each breaks exactly its matching test and nothing else. The first attempt at the preview-cap test used a homogeneous "x"*200 body, which a widened cap slipped through unnoticed (contains() found a shifted match); replaced with a sentinel character at index 80 to actually pin the boundary. Updates CLAUDE.md's intent->tool table for fleet_poll's new coordId semantics, per this repo's own "prompt is part of the product" rule. wiki/ is a submodule and not committable from a worker's worktree -- wiki-bound content is in the PR body instead. --- CLAUDE.md | 1 + .../main/java/dev/ltms/fleet/auth/Authz.java | 15 +++ .../java/dev/ltms/fleet/mcp/FleetMcp.java | 122 +++++++++++++++--- .../java/dev/ltms/fleet/auth/AuthzTest.java | 14 ++ .../dev/ltms/fleet/mcp/FleetMcpAuthzTest.java | 60 +++++++-- .../java/dev/ltms/fleet/mcp/FleetMcpTest.java | 84 +++++++++++- 6 files changed, 271 insertions(+), 25 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index 29d90a1..3cbdf43 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -138,6 +138,7 @@ the merge — and merging on a reviewer's word is delegating it by proxy. | Message a **peer lead** on this host | `fleet_send{sessionId: , content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task | | Message a **peer lead** on another daemon or host | `fleet_send{coordId: , content}` — needs a `coordinator:` block; your own coord-id is in `fleet_list`. Coordination only, **never** a task | | Answer a peer lead that messaged you | `fleet_send{coordId}` — or `{sessionId}` if they are on this host. **Not** `fleet_reply`: it has no peer route and the publish is refused | +| Read your own held lead-to-lead mail (no ack) | `fleet_poll{coordId: }` — primary-only; never acks, so `fleet_list`'s `held[]` still shows it after. `fleet_list`'s `held[]` gives only a truncated preview — this is the only way to read the full body | | Collect a held reply | `fleet_poll{target}` · then `fleet_ack{target, msgId}` | | Tear down a member | `fleet_stop{paneId}` | diff --git a/fleetd/src/main/java/dev/ltms/fleet/auth/Authz.java b/fleetd/src/main/java/dev/ltms/fleet/auth/Authz.java index 1cb536c..1c186e0 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/auth/Authz.java +++ b/fleetd/src/main/java/dev/ltms/fleet/auth/Authz.java @@ -30,6 +30,16 @@ public final class Authz { DRAIN, /** Read-only observation: status, roster, profiles, task polling. */ READ, + /** + * Read (never ack) this daemon's own held lead-to-lead coordination mail (fleetd #421). + * + *

Deliberately not folded into {@link #READ}. {@code READ}'s grant + * rests on "the roster carries no secrets" (see its case below) — a lead-to-lead body is + * not the roster; it is where leads discuss host shapes, credentials and unmerged work. + * Mapping this to {@code READ} would let any worker read every peer lead's mail in full + * and would silently falsify that comment for every other {@code READ} caller. + */ + COORD_READ, /** Scrape the metrics endpoint. */ METRICS } @@ -68,6 +78,11 @@ public final class Authz { // Observation is open to every authenticated role: a worker legitimately polls its own // status, and the roster carries no secrets. case READ, METRICS -> caller.isPrimary() || caller.isWorker() || caller.isArchitect(); + + // fleetd #421: reading held lead-to-lead mail is the primary's alone. An architect + // holds READ today (CB-548), so "not primary" must mean not-architect here too — this + // is coordination between leads, not observation of the roster. + case COORD_READ -> caller.isPrimary(); }; } 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 792305a..52d8082 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -386,10 +386,11 @@ public final class FleetMcp { (exchange, req) -> { Map a = req.arguments(); String target = str(a, "target"); + String coordId = str(a, "coordId"); // The action depends on the ARGUMENTS, not on the tool name -- see pollAction. McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_poll", a), target); if (denied != null) return denied; - return poll(messages, str(a, "ticket"), target); + return poll(messages, leadChannel, str(a, "ticket"), target, coordId); }; // CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox). // Acking removes a reply from the inbox, so it is a drain, not a read. @@ -811,29 +812,42 @@ public final class FleetMcp { /** * Which authorization action a {@code fleet_poll} call needs, decided by its arguments - * (fleetd #272). + * (fleetd #272, widened by fleetd #421). * - *

{@code fleet_poll} is two operations behind one tool name. With {@code - * ticket} it observes an async delegation and changes nothing, which is a {@link + *

{@code fleet_poll} is now three operations behind one tool name. With + * {@code ticket} it observes an async delegation and changes nothing, which is a {@link * Authz.Action#READ}. With {@code target} it calls {@link MessageService#drainReplies} on that * session -- the replies are removed from the inbox and a second call returns nothing -- so it * is a {@link Authz.Action#DRAIN}, the same gate {@code fleet_ack} already uses for removing a - * single message, and the same one the REST path uses at {@code FleetApp.drainReplies}. + * single message, and the same one the REST path uses at {@code FleetApp.drainReplies}. With + * {@code coordId} it reads (never acks) this daemon's own held lead-to-lead mail, which is a + * {@link Authz.Action#COORD_READ} -- not {@code READ}, even though nothing is + * consumed: {@code READ}'s grant is open to every authenticated role on the premise that the + * roster carries no secrets, and a lead-to-lead body is not the roster. Mapping a non-destructive + * peer-mail read to {@code READ} would let any worker read every peer lead's mail in full. * - *

Until this method existed the handler passed a constant {@code READ} for both branches. - * {@code READ} is open to every authenticated role, so any worker could read a peer's id out of - * {@code fleet_list} and destroy the replies that peer had queued for the primary. The gate - * failed open, and it did so because the required action is a function of the arguments while - * the handler chose it before looking at them. + *

Before this method existed (fleetd #272) the handler passed a constant {@code READ} for + * both of the original branches. {@code READ} is open to every authenticated role, so any + * worker could read a peer's id out of {@code fleet_list} and destroy the replies that peer had + * queued for the primary. The gate failed open, and it did so because the required action is a + * function of the arguments while the handler chose it before looking at them. * *

The choice lives in this method, and not inline in the handler, so that a test can assert * the mapping the handler actually uses. {@code FleetMcpAuthzTest} already checked every * {@link Authz.Action} against every {@link Role} and passed throughout -- it tested the policy * table, which was correct, while the defect was in which action the caller handed it. * - * @param target the {@code target} argument of the call, or {@code null}/blank when absent + *

Checked first, and exclusively of {@code target}: a call naming {@code coordId} is reading + * a different inbox entirely (this daemon's own lead channel, never a worker's), so it takes + * priority over whatever {@code target} might also say. + * + * @param target the {@code target} argument of the call, or {@code null}/blank when absent + * @param coordId the {@code coordId} argument of the call, or {@code null}/blank when absent */ - static Authz.Action pollAction(String target) { + static Authz.Action pollAction(String target, String coordId) { + if (!isBlank(coordId)) { + return Authz.Action.COORD_READ; + } return isBlank(target) ? Authz.Action.READ : Authz.Action.DRAIN; } @@ -848,7 +862,7 @@ public final class FleetMcp { case "fleet_reply" -> Authz.Action.REPLY; case "fleet_ask" -> Authz.Action.ASK; case "fleet_status", "fleet_list", "fleet_profiles", "fleet_whoami" -> Authz.Action.READ; - case "fleet_poll" -> pollAction(str(arguments, "target")); + case "fleet_poll" -> pollAction(str(arguments, "target"), str(arguments, "coordId")); case "fleet_ack" -> Authz.Action.DRAIN; case "fleet_spawn" -> Authz.Action.SPAWN; case "fleet_stop" -> Authz.Action.STOP; @@ -858,6 +872,21 @@ public final class FleetMcp { /** {@code fleet_poll}: check an async delegation by ticket, or drain a worker's inbox by target. */ static McpSchema.CallToolResult poll(MessageService messages, String ticket, String target) { + return poll(messages, null, ticket, target, null); + } + + /** + * As above, plus (fleetd #421) a held-peer-mail read when {@code coordId} is present: returns + * this daemon's own held lead-to-lead messages, in full, without acking them. Checked first and + * exclusively of {@code ticket}/{@code target} — see {@link #pollAction}'s javadoc for why this + * is a different inbox (this daemon's own {@link LeadChannel}) that authorizes differently + * ({@link Authz.Action#COORD_READ}, primary-only) from either of the original two branches. + */ + static McpSchema.CallToolResult poll(MessageService messages, LeadChannel leadChannel, String ticket, + String target, String coordId) { + if (!isBlank(coordId)) { + return pollHeldPeerMail(leadChannel, coordId); + } if (!isBlank(target)) { var replies = messages.drainReplies(target); if (replies.isEmpty()) { @@ -885,6 +914,46 @@ public final class FleetMcp { }; } + /** + * fleetd #421: a lead's own held lead-to-lead mail, read without consuming it. + * + *

{@link LeadChannel#peek} is non-destructive, so calling this twice returns the same + * bodies, and {@code fleet_list}'s {@code coordinator.held[]} is unaffected — this adds a + * read, it never acks. Never let this short-circuit the delivery contract: {@code + * LeadCoordLoop}'s javadoc explains why a message must stay unacked until actually delivered, + * and that is unchanged here. + * + *

{@code coordId} must be THIS daemon's own coord-id ({@code fleet_list}'s + * {@code coordinator.selfId}) — there is no route here to read a PEER's outbound mail, only + * your own inbound mail. Requiring the caller to echo its own id catches the exact confusion + * that opened fleetd #421: the original failed attempt passed a PEER's id ("fleet01") as + * {@code fleet_poll}'s {@code target}, expecting to read that peer's messages, and got the same + * silent {@code []} as a genuinely empty worker inbox. This refuses the same mistake here with a + * reason, instead of a second silent wrong answer. + */ + static McpSchema.CallToolResult pollHeldPeerMail(LeadChannel leadChannel, String coordId) { + if (leadChannel == null) { + return error("lead coordination is not configured (no coordinator: block) — there is " + + "no held peer mail to read."); + } + String selfId = leadChannel.selfCoordId(); + if (!coordId.equals(selfId)) { + return error("coordId \"" + coordId + "\" is not this daemon's own coord-id (\"" + selfId + + "\"). fleet_poll reads only YOUR OWN held mail — pass your own coordId " + + "(fleet_list's coordinator.selfId), not a peer's."); + } + return text(json(leadChannel.peek().stream().map(FleetMcp::heldMailView).toList())); + } + + /** The full body of one held lead-to-lead message — never truncated, unlike {@link #heldView}. */ + private static Map heldMailView(LeadMessage m) { + Map row = new LinkedHashMap<>(); + row.put("msgId", m.msgId()); + row.put("from", m.from()); + row.put("content", m.content()); + return row; + } + /** * {@code fleet_reply}: the worker returns its structured answer, resolving the awaiting send * or — when no send is open — queueing the reply in the inbox for later drain (CB-307). @@ -1314,6 +1383,16 @@ public final class FleetMcp { * 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. + * + *

fleetd #421: {@code heldCount}/{@code heldDurable} fix the "pending: 0" trap. + * {@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[]}. */ private static Map coordinatorView(CoordinationSource coordination) { LeadChannel channel = coordination.leadChannel(); @@ -1321,11 +1400,14 @@ public final class FleetMcp { return null; } String selfId = channel.selfCoordId(); + List held = channel.peek(); 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("heldCount", held.size()); + row.put("heldDurable", true); + row.put("held", held.stream().map(FleetMcp::heldView).toList()); row.put("peers", coordination.peers().stream().map(p -> peerView(channel, p)).toList()); return row; } @@ -1654,10 +1736,18 @@ public final class FleetMcp { "Check an async delegation (a fleet_send with wait:false) by its ticket: " + "pending, done (with the worker's reply), or failed. When target (a worker " + "session id) is present instead of ticket, drain that worker's inbox of " - + "replies delivered when no send was open.", + + "replies delivered when no send was open. When coordId is present instead, " + + "read (never consume) your own held lead-to-lead mail — primary-only.", objectSchema(Map.of( "ticket", stringProp("The ticket returned by fleet_send wait:false"), - "target", stringProp("Worker session id to drain pending replies from (optional)")), + "target", stringProp("Worker session id to drain pending replies from (optional)"), + "coordId", stringProp("Your own coord-id (fleet_list's coordinator.selfId) — " + + "reads every message currently held[] for you in full, without " + + "acking. Read twice, get the same bodies both times; fleet_list's " + + "held[] still reports them afterward. Primary-only, and this can " + + "read only YOUR OWN mailbox — there is no route to a peer's outbound " + + "mail, so passing a peer's coordId here is refused rather than " + + "silently returning the wrong thing (or nothing).")), List.of())); } diff --git a/fleetd/src/test/java/dev/ltms/fleet/auth/AuthzTest.java b/fleetd/src/test/java/dev/ltms/fleet/auth/AuthzTest.java index 0abbb69..07bab44 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/auth/AuthzTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/auth/AuthzTest.java @@ -113,6 +113,20 @@ class AuthzTest { assertTrue(Authz.permits(WORKER_A, METRICS, null)); } + /** + * fleetd #421 — unlike READ (the case above), COORD_READ (a lead's own held peer mail) is the + * primary's alone. A worker or an architect reading it would disclose lead-to-lead + * coordination bodies, not the secret-free roster READ is open about. + */ + @Test + void coordReadIsThePrimarysAloneNotAWidenedRead() { + assertTrue(Authz.permits(PRIMARY, COORD_READ, null)); + assertFalse(Authz.permits(WORKER_A, COORD_READ, null), + "a worker must not read held lead-to-lead mail"); + assertFalse(Authz.permits(ARCH_DESIGN, COORD_READ, null), + "an architect holds READ today, but not-primary must mean not-architect here too"); + } + @Test void unauthenticatedIsDistinguishedFromMerelyForbidden() { // Drives the 401-vs-403 split: a missing credential is fixable by the caller, a wrong role diff --git a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpAuthzTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpAuthzTest.java index 191f0a7..016037c 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpAuthzTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpAuthzTest.java @@ -201,14 +201,34 @@ class FleetMcpAuthzTest { */ @Test void pollingByTargetIsADrainAndPollingByTicketIsARead() { - assertEquals(Authz.Action.DRAIN, FleetMcp.pollAction("term_b"), + assertEquals(Authz.Action.DRAIN, FleetMcp.pollAction("term_b", null), "poll by target removes the replies — that is a drain, not an observation"); - assertEquals(Authz.Action.READ, FleetMcp.pollAction(null), + assertEquals(Authz.Action.READ, FleetMcp.pollAction(null, null), "poll by ticket changes nothing"); - assertEquals(Authz.Action.READ, FleetMcp.pollAction(" "), + assertEquals(Authz.Action.READ, FleetMcp.pollAction(" ", null), "a blank target is an absent target"); } + /** + * fleetd #421: a coordId branch is a THIRD operation behind fleet_poll's one name, and it must + * map to {@link Authz.Action#COORD_READ} — never {@link Authz.Action#READ}, even though this + * branch also consumes nothing. READ's grant is open to every authenticated role on the premise + * that the roster carries no secrets; a lead-to-lead body is not the roster, so folding this + * branch into READ would let any worker read every peer lead's mail in full. coordId also takes + * priority over target when both happen to be present — it addresses a different inbox entirely. + */ + @Test + void pollingByCoordIdIsACoordReadNeverAPlainRead() { + assertEquals(Authz.Action.COORD_READ, FleetMcp.pollAction(null, "mac-opus"), + "reading held peer mail must not be mapped to the everyone-readable READ action"); + assertEquals(Authz.Action.COORD_READ, FleetMcp.pollAction(" ", "mac-opus"), + "a blank target must not fall through to READ/DRAIN when coordId is present"); + assertEquals(Authz.Action.READ, FleetMcp.pollAction(null, " "), + "a blank coordId is an absent coordId, same as target/ticket"); + assertEquals(Authz.Action.COORD_READ, FleetMcp.pollAction("term_b", "mac-opus"), + "coordId takes priority over target — this is a different inbox, not a drain"); + } + @Test void everyRegisteredToolHasItsHandlerActionPinned() { Set registered = toolsTheServerRegisters(); @@ -230,6 +250,8 @@ class FleetMcpAuthzTest { assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_whoami", Map.of())); assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_poll", Map.of("ticket", "task"))); assertEquals(Authz.Action.DRAIN, FleetMcp.toolAction("fleet_poll", Map.of("target", "term_b"))); + assertEquals(Authz.Action.COORD_READ, + FleetMcp.toolAction("fleet_poll", Map.of("coordId", "mac-opus"))); } private static Set toolsTheServerRegisters() { @@ -249,11 +271,11 @@ class FleetMcpAuthzTest { void aWorkerMayNotDrainAnotherSessionsInboxByPolling() { FleetMcp m = mcp(true); - assertNotNull(m.denyFor(WORKER_A, FleetMcp.pollAction("term_b"), "term_b"), + assertNotNull(m.denyFor(WORKER_A, FleetMcp.pollAction("term_b", null), "term_b"), "a worker draining a peer's inbox would destroy replies queued for the primary"); - assertNotNull(m.denyFor(ARCH_DESIGN, FleetMcp.pollAction("term_b"), "term_b"), + assertNotNull(m.denyFor(ARCH_DESIGN, FleetMcp.pollAction("term_b", null), "term_b"), "an architect has no lifecycle rights either — same gate as fleet_ack"); - assertNull(m.denyFor(PRIMARY, FleetMcp.pollAction("term_b"), "term_b"), + assertNull(m.denyFor(PRIMARY, FleetMcp.pollAction("term_b", null), "term_b"), "collecting a held reply is the primary's job"); } @@ -265,11 +287,33 @@ class FleetMcpAuthzTest { void pollingAnOwnTicketStaysOpenToWorkersAndArchitects() { FleetMcp m = mcp(true); - assertNull(m.denyFor(WORKER_A, FleetMcp.pollAction(null), null)); - assertNull(m.denyFor(ARCH_DESIGN, FleetMcp.pollAction(null), null), + assertNull(m.denyFor(WORKER_A, FleetMcp.pollAction(null, null), null)); + assertNull(m.denyFor(ARCH_DESIGN, FleetMcp.pollAction(null, null), null), "an architect delegates with wait:false, so it must be able to poll its ticket"); } + /** + * fleetd #421 acceptance: the whole ticket, pinned at the policy-table layer. A worker must be + * refused the coordId branch, the primary must be allowed, and (CB-548) an architect — which + * holds READ today — must be refused too, because "not primary" means not-architect here. + */ + @Test + void aWorkerAndAnArchitectMayNotReadHeldPeerMailOnlyThePrimaryMay() { + FleetMcp m = mcp(true); + Authz.Action coordRead = FleetMcp.pollAction(null, "mac-opus"); + + McpSchema.CallToolResult workerDenied = m.denyFor(WORKER_A, coordRead, null); + assertNotNull(workerDenied, "a worker must not read held lead-to-lead mail"); + assertTrue(workerDenied.isError()); + + McpSchema.CallToolResult archDenied = m.denyFor(ARCH_DESIGN, coordRead, null); + assertNotNull(archDenied, "an architect holds READ today, but not-primary must mean " + + "not-architect here too"); + assertTrue(archDenied.isError()); + + assertNull(m.denyFor(PRIMARY, coordRead, null), "reading its own held mail is the primary's job"); + } + // --- identity reconstruction from the transport context ------------------------------------ @Test 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 da985f9..e55bdbd 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java @@ -680,7 +680,11 @@ class FleetMcpTest { 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); + // A homogeneous "x".repeat(200) body would not pin the exact 80-char boundary: a preview + // widened by one (81 chars) still CONTAINS "x".repeat(80) + "…" as a substring one position + // later, because every character is 'x'. Put a sentinel ("Y") exactly at index 80 — the + // first character a widened cap would leak — so any preview past 80 chars is caught. + String longContent = "x".repeat(80) + "Y" + "z".repeat(119); FakeLeadChannel channel = new FakeLeadChannel("mac-opus") .hold(new LeadMessage("m1", "fleet01-lead", "mac-opus", longContent)); @@ -694,9 +698,87 @@ class FleetMcpTest { 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); + assertFalse(out.contains("Y"), + "the sentinel at index 80 must never appear — a preview past 80 chars leaked it: " + out); assertTrue(out.contains("x".repeat(80) + "…"), "expected an 80-char preview with an ellipsis: " + out); } + /** + * fleetd #421: {@code mailbox.pending} counts only broker-ready messages, so a blocked lead's + * normal, healthy state is {@code "pending": 0} next to a non-empty {@code held[]} — which + * invites the false reading that held mail is in-memory-only. {@code heldCount}/{@code + * heldDurable} put an honest second number and fact beside {@code pending} instead of leaving it + * as the only one. + */ + @Test + void listReportsAnHonestHeldCountAndDurabilityNotJustPendingZero() { + 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)) + .hold(new LeadMessage("m1", "fleet01-lead", "mac-opus", "one")) + .hold(new LeadMessage("m2", "fleet01-lead", "mac-opus", "two")) + .hold(new LeadMessage("m3", "fleet01-lead", "mac-opus", "three")); + + 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("\"pending\":0"), out); + assertTrue(out.contains("\"heldCount\":3"), + "the honest count beside pending: 0 — three messages really are held: " + out); + assertTrue(out.contains("\"heldDurable\":true"), + "must state the durability fact, not leave pending as the only number next to held[]: " + out); + } + + // ── fleetd #421: a lead reads (never consumes) its own held peer mail ────────────────────── + + @Test + void pollWithCoordIdReturnsTheFullBodyWithoutAckingAndLeavesItHeld() { + FakeLeadChannel channel = new FakeLeadChannel("mac-opus") + .hold(new LeadMessage("m1", "fleet01-lead", "mac-opus", "x".repeat(200))); + + McpSchema.CallToolResult first = FleetMcp.poll(messages, channel, null, null, "mac-opus"); + assertNotEquals(Boolean.TRUE, first.isError(), textOf(first)); + String out1 = textOf(first); + assertTrue(out1.contains("\"msgId\":\"m1\""), out1); + assertTrue(out1.contains("\"from\":\"fleet01-lead\""), out1); + assertTrue(out1.contains("\"content\":\"" + "x".repeat(200) + "\""), + "the coordId route must return the FULL body, unlike fleet_list's preview: " + out1); + + // Read again: identical bodies, and nothing was acked — peek() still holds it. + McpSchema.CallToolResult second = FleetMcp.poll(messages, channel, null, null, "mac-opus"); + assertEquals(out1, textOf(second), "peek is non-destructive — reading twice must return the same bodies"); + assertTrue(channel.acked().isEmpty(), "a read must never ack — that is the whole point of the ticket"); + assertEquals(1, channel.peek().size(), "the message must still be held after being read"); + } + + @Test + void pollWithCoordIdRefusesAPeersCoordIdInsteadOfReturningTheWrongMailOrNothing() { + FakeLeadChannel channel = new FakeLeadChannel("mac-opus") + .hold(new LeadMessage("m1", "fleet01-lead", "mac-opus", "secret coordination body")); + + // The ORIGINAL fleetd #421 confusion: passing a PEER's id where a self-address was meant. + McpSchema.CallToolResult res = FleetMcp.poll(messages, channel, null, null, "fleet01-lead"); + + assertEquals(Boolean.TRUE, res.isError()); + String out = textOf(res); + assertFalse(out.contains("secret coordination body"), + "a refused read must never leak the body it refused to return: " + out); + assertTrue(out.contains("mac-opus"), "the error must name this daemon's own coordId: " + out); + } + + @Test + void pollWithCoordIdErrorsHonestlyWhenLeadCoordinationIsNotConfigured() { + McpSchema.CallToolResult res = FleetMcp.poll(messages, null, null, null, "mac-opus"); + + assertEquals(Boolean.TRUE, res.isError()); + assertTrue(textOf(res).toLowerCase().contains("coordinat"), textOf(res)); + } + @Test void listReportsEachDeclaredPeersLiveReachability() { FakeHerdr h = new FakeHerdr(); -- 2.52.0