fleet_poll{coordId} peeks this daemon's own held lead-to-lead mail and
returns full bodies without acking. New Authz.Action COORD_READ, primary
only — not the architect, which holds READ today. pollAction is now
argument-derived over both target and coordId.
heldView and HELD_PREVIEW_MAX_CHARS are untouched: fleet_list stays a
cheap always-safe scan, and the full read is a separately authorized call.
Verified by me, not taken from the report: merged tree builds 1562 green
(main 1555 + 7 new methods), 0 compile errors, no merge conflicts.
Mutation battery on lines the worker did NOT mutate, control 111 green:
- COORD_READ widened to the architect KILLED (AuthzTest + FleetMcpAuthzTest)
- self-coord-id guard removed KILLED (FleetMcpTest)
- full body swapped for the preview KILLED (FleetMcpTest)
The third mutation targets the ticket's own deliverable, and it is pinned.
Authz.permits has no default, so a new action is a compile error rather
than a silently unhandled case.
Follow-up filed as #439: fleet_list's coordinator row is still READ-gated,
so a worker sees peer coord-ids and 80-char previews of lead-to-lead
bodies. Pre-existing; the implementer flagged it and left it alone.
This commit was merged in pull request #438.
This commit is contained in:
@@ -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: <their terminal>, content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task |
|
||||
| Message a **peer lead** on another daemon or host | `fleet_send{coordId: <their coord-id>, 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: <your own coord-id, from fleet_list's coordinator.selfId>}` — 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}` |
|
||||
|
||||
|
||||
@@ -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).
|
||||
*
|
||||
* <p>Deliberately <strong>not</strong> 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();
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
@@ -386,10 +386,11 @@ public final class FleetMcp {
|
||||
(exchange, req) -> {
|
||||
Map<String, Object> 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).
|
||||
*
|
||||
* <p>{@code fleet_poll} is <strong>two operations behind one tool name</strong>. With {@code
|
||||
* ticket} it observes an async delegation and changes nothing, which is a {@link
|
||||
* <p>{@code fleet_poll} is now <strong>three operations behind one tool name</strong>. 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} -- <strong>not</strong> {@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.
|
||||
*
|
||||
* <p>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.
|
||||
* <p>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.
|
||||
*
|
||||
* <p>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
|
||||
* <p>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.
|
||||
*
|
||||
* <p>{@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.
|
||||
*
|
||||
* <p>{@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<String, Object> heldMailView(LeadMessage m) {
|
||||
Map<String, Object> 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).
|
||||
@@ -1323,6 +1392,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.
|
||||
*
|
||||
* <p><strong>fleetd #421: {@code heldCount}/{@code heldDurable} fix the "pending: 0" trap.</strong>
|
||||
* {@code mailbox.pending} counts only broker-<em>ready</em> 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<String, Object> coordinatorView(CoordinationSource coordination) {
|
||||
LeadChannel channel = coordination.leadChannel();
|
||||
@@ -1330,11 +1409,14 @@ public final class FleetMcp {
|
||||
return null;
|
||||
}
|
||||
String selfId = channel.selfCoordId();
|
||||
List<LeadMessage> held = channel.peek();
|
||||
Map<String, Object> 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;
|
||||
}
|
||||
@@ -1663,10 +1745,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()));
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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<String> 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<String> 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
|
||||
|
||||
@@ -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();
|
||||
|
||||
Reference in New Issue
Block a user