Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha 11eccc3a1b fleetd #705: gate fleet_poll's ticket lookup by the creating caller's terminal
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 1m45s
A ticket id is a plain sequential counter, so any session holding
TASK_READ could walk task-1, task-2, ... and read another session's
delegation reply. sendAsync now records the creating caller's terminal
on the Task, and poll refuses a caller whose terminal differs from it.
A caller with no terminal (the unnamed primary) is still allowed
through regardless, since it never carries a herdr pane to compare.
2026-10-04 06:50:26 +02:00
8 changed files with 171 additions and 221 deletions
@@ -2776,10 +2776,9 @@ public record FleetConfig(
});
if (!bad.isEmpty()) {
throw new IllegalStateException("refusing to start: " + String.join("; ", bad)
+ ". A member labelled that way, while its pane carries no entry in the "
+ "spawned-member roster, is read back as a lead or collaborator and granted "
+ "that identity's authority. Change one of the two so member tabs cannot be "
+ "confused with a lead's or collaborator's tab.");
+ ". Every member labelled that way would be read back as a lead or "
+ "collaborator and granted that identity's authority. Change one of the two "
+ "so member tabs cannot be confused with a lead's or collaborator's tab.");
}
List<String> collisions = new ArrayList<>();
@@ -2878,8 +2877,7 @@ public record FleetConfig(
* the focused tab rather than its own, so it can land inside a lead's or collaborator's own
* labelled tab. {@link dev.ltms.fleet.herdr.LeadTabScanner} identifies a lead or collaborator
* purely by that tab's label — it does not exclude the member space — so a member that ends up
* there, while its pane carries no entry in the spawned-member roster, is read back as that
* lead or collaborator and granted that identity's authority.
* there would be read back as that lead or collaborator and granted that identity's authority.
*
* <p>Only an entry with a non-blank {@code tab} is in scope: one with no {@code tab} feeds
* nothing into {@link dev.ltms.fleet.herdr.LeadTabScanner}, so it creates no hazard here.
@@ -2910,9 +2908,8 @@ public record FleetConfig(
}
throw new IllegalStateException("refusing to start: profile(s) " + bad
+ " use placement: pane while fleet.leaders or fleet.collaborators names a tab. A "
+ "pane-placed member can land inside that labelled tab, and while its pane "
+ "carries no entry in the spawned-member roster, it is read back as the lead or "
+ "collaborator and granted that identity's authority. Set placement: tab for "
+ "pane-placed member can land inside that labelled tab and be read back as the "
+ "lead or collaborator, granted that identity's authority. Set placement: tab for "
+ "each named profile, or remove the tab from every fleet.leaders and "
+ "fleet.collaborators entry.");
}
@@ -490,7 +490,7 @@ public final class FleetMcp {
Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, caller);
// wait defaults to true (block for the reply); wait:false is fire-and-poll.
return Boolean.FALSE.equals(a.get("wait"))
? sendAsync(messages, target, content, onAccepted, workers.profiles())
? sendAsync(messages, target, content, onAccepted, workers.profiles(), caller)
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles());
};
// fleet_reply's identity is the CONNECTION, never an argument — so the authz check
@@ -524,7 +524,7 @@ public final class FleetMcp {
// 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, leadChannel, str(a, "ticket"), target, coordId);
return poll(messages, leadChannel, str(a, "ticket"), target, coordId, callerTerminal(exchange));
};
// 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.
@@ -560,7 +560,6 @@ public final class FleetMcp {
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage,
leadSeats, leadContextGauge, leadConfigDirs, callers.leads(),
callerTerminal(exchange),
callers.collaborators(), collaboratorsVisibleTo(principal(exchange)),
new CoordinationSource(leadChannel, peers),
coordinatorVisibleTo(principal(exchange)));
};
@@ -744,21 +743,6 @@ public final class FleetMcp {
return caller.isPrimary();
}
/**
* fleetd #703: who may see {@code fleet_list}'s {@code collaborators} array — a roster of the
* human-opened tabs this daemon recognises as named peers. Visible to exactly the roles that
* may {@link Authz.Action#SEND} to a named peer ({@code Authz.java}'s {@code SEND} case):
* the primary, an architect, and a collaborator (reaching another collaborator or a lead). A
* worker holds {@code READ} but can never {@code SEND} to a collaborator, so listing them to a
* worker would expose which human tabs exist on the host with no use to that caller. Split out
* for the same reason as {@link #coordinatorVisibleTo}: the decision must be unit-testable
* without fabricating an SDK {@code McpSyncServerExchange}, and the handler must call this
* named predicate rather than inlining the check.
*/
static boolean collaboratorsVisibleTo(Principal caller) {
return caller.isPrimary() || caller.isArchitect() || caller.isCollaborator();
}
/** The worker identity resolved from this call's connection, or {@code null} if the primary. */
private static String callerTerminal(McpSyncServerExchange exchange) {
Object v = exchange.transportContext().get(CALLER_TERMINAL);
@@ -972,6 +956,17 @@ public final class FleetMcp {
*/
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
Runnable onAccepted, Set<String> profiles) {
return sendAsync(messages, sessionId, content, onAccepted, profiles, null);
}
/**
* As above, recording {@code creatorTerminal} as this ticket's owner (the caller's own
* terminal, resolved from the connection) so a later {@code fleet_poll{ticket}} only hands the
* result back to that same caller — see {@link MessageService#poll(String, String)}.
*/
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
Runnable onAccepted, Set<String> profiles,
String creatorTerminal) {
if (isBlank(sessionId) || isBlank(content)) {
return error("sessionId and content are required");
}
@@ -979,7 +974,7 @@ public final class FleetMcp {
if (targetError != null) {
return targetError;
}
String ticket = messages.sendAsync(sessionId, content, onAccepted);
String ticket = messages.sendAsync(sessionId, content, onAccepted, creatorTerminal);
return text("accepted — task delegated. Poll fleet_poll with ticket=" + ticket);
}
@@ -1158,9 +1153,24 @@ public final class FleetMcp {
* 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.
*
* <p>Does not check who owns {@code ticket} — see the overload that takes {@code
* callerTerminal} for that. Callers that do not resolve a caller terminal (tests, or a surface
* with no connection-based identity) use this one.
*/
static McpSchema.CallToolResult poll(MessageService messages, LeadChannel leadChannel, String ticket,
String target, String coordId) {
return poll(messages, leadChannel, ticket, target, coordId, null);
}
/**
* As above, refusing a ticket lookup whose caller's terminal differs from the terminal that
* created it — see {@link MessageService#poll(String, String)}. {@code callerTerminal} is the
* CALLING session's terminal id, resolved by the MCP layer from the connection, never a
* client-supplied value.
*/
static McpSchema.CallToolResult poll(MessageService messages, LeadChannel leadChannel, String ticket,
String target, String coordId, String callerTerminal) {
if (!isBlank(coordId)) {
return pollHeldPeerMail(leadChannel, coordId);
}
@@ -1174,7 +1184,7 @@ public final class FleetMcp {
if (isBlank(ticket)) {
return error("ticket (or target) is required");
}
MessageService.TaskView v = messages.poll(ticket);
MessageService.TaskView v = messages.poll(ticket, callerTerminal);
if (v == null) {
return error("unknown ticket: " + ticket + " (never issued, or expired)");
}
@@ -1777,8 +1787,7 @@ public final class FleetMcp {
Map<String, String> leads, String selfTerm,
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine,
OutageSource.none(), LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, false);
OutageSource.none(), LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
@@ -1787,8 +1796,7 @@ public final class FleetMcp {
QuarantineSource quarantine, OutageSource outage,
Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, CoordinationSource.none(), false);
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, CoordinationSource.none(), false);
}
/**
@@ -1805,8 +1813,7 @@ public final class FleetMcp {
QuarantineSource quarantine, Map<String, String> leads, String selfTerm,
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, OutageSource.none(),
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, false);
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
@@ -1815,8 +1822,7 @@ public final class FleetMcp {
QuarantineSource quarantine, OutageSource outage,
Map<String, String> leads, String selfTerm, CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, false);
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
}
/**
@@ -1838,8 +1844,7 @@ public final class FleetMcp {
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, false);
leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
}
/**
@@ -1863,8 +1868,7 @@ public final class FleetMcp {
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
CoordinationSource coordination, boolean callerIsPrimary) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine,
outage, leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, callerIsPrimary);
outage, leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, callerIsPrimary);
}
/**
@@ -1874,13 +1878,6 @@ public final class FleetMcp {
* matter); the one caller that matters for caching, {@code fleet_list}'s MCP handler, passes
* its own single long-lived instance instead (see {@code FleetMcp}'s {@code leadContextGauge}
* field).
*
* @param collaborators terminal_id → collaborator name, live from the resolver
* ({@link dev.ltms.fleet.auth.CallerResolver#collaborators()})
* @param collaboratorsVisible whether this caller may see the {@code collaborators} array (see
* {@link #collaboratorsVisibleTo}); every wrapper overload above
* passes {@code false}, so a test that wants the row must call this
* overload with an explicit {@code true}
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
@@ -1889,7 +1886,6 @@ public final class FleetMcp {
LeadSeatSource leadSeats, LeadContextGauge contextGauge,
LeadConfigDirSource leadConfigDirs,
Map<String, String> leads, String selfTerm,
Map<String, String> collaborators, boolean collaboratorsVisible,
CoordinationSource coordination, boolean callerIsPrimary) {
try {
Map<String, Agent> live = workers.list().stream()
@@ -1916,16 +1912,6 @@ public final class FleetMcp {
result.put("loopHealth", Map.of(
"statusPoller", loopHealth.statusPoller().get().name(),
"sessionReaper", loopHealth.sessionReaper().get().name()));
// fleetd #703: a collaborator tab is a person's own tab, so the row is assembled and
// included only for the roles that may SEND to a named peer -- gate BEFORE assembling
// it, same reason as the coordinator row just below: the key must be absent for a
// worker, never present-and-empty.
if (collaboratorsVisible) {
result.put("collaborators", collaborators.entrySet().stream()
.sorted(Map.Entry.comparingByValue())
.map(e -> collaboratorRow(e.getKey(), e.getValue()))
.toList());
}
// fleetd #439: coordinator/coordinatorView is lead-to-lead coordination state and must
// never reach a worker or an architect -- gate BEFORE assembling it, not after, so the
// key is absent rather than present-and-empty.
@@ -2237,20 +2223,6 @@ public final class FleetMcp {
return m;
}
/**
* fleetd #703: one row of the {@code collaborators} array — a named peer's registry name and
* the herdr {@code terminal_id} a {@code fleet_send} must target to reach it. No status, no
* context gauge, no profile: a collaborator is never spawned and carries no profile, so those
* fields have no meaning for it, and the scan behind {@code terminal} only reports a tab it
* actually found in herdr, so a tab nobody has open does not appear here at all.
*/
private static Map<String, Object> collaboratorRow(String terminal, String name) {
Map<String, Object> m = new LinkedHashMap<>();
m.put("name", name);
m.put("sessionId", terminal);
return m;
}
/**
* Reads the lead context gauge for one lead. {@code configDir} is this lead's configured
* {@code CLAUDE_CONFIG_DIR} override (see {@link LeadConfigDirSource}), derived from
@@ -2476,12 +2448,7 @@ public final class FleetMcp {
+ "msgId/from/preview, never the full body), and one row per coordinator.peers "
+ "coord-id ('peers': coordId/reachable, plus pending/consumers when reachable) — "
+ "this is peer DISCOVERY for cross-host leads, distinct from the local 'leads' "
+ "array above. It is omitted entirely when no coordinator is configured. A "
+ "'collaborators' array, visible only to the primary, an architect, and a "
+ "collaborator (never a worker), reports every other named peer tab this daemon "
+ "recognises: each row is 'name' (its registry name) and 'sessionId' (the "
+ "terminal id a fleet_send targets to reach it). Absent entirely for a worker, "
+ "whatever is configured.",
+ "array above. It is omitted entirely when no coordinator is configured.",
objectSchema(Map.of(), List.of()));
}
@@ -275,10 +275,19 @@ public final class MessageService {
* {@link #abandon}) can never match again regardless of this flag's value.
*/
private volatile boolean askTimedOut;
/**
* The terminal of the caller whose {@code fleet_send{wait:false}} created this ticket, or
* {@code null} when that caller had no terminal (the unnamed primary) or the ticket was
* created through an overload that does not record one. {@link #poll(String, String)}
* compares a polling caller's own terminal against this field before handing back the
* ticket's state.
*/
private final String creatorTerminal;
private Task(String ticket, String target, LongSupplier nowNanos) {
private Task(String ticket, String target, LongSupplier nowNanos, String creatorTerminal) {
this.ticket = ticket;
this.target = target;
this.creatorTerminal = creatorTerminal;
this.createdNanos = nowNanos.getAsLong();
future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong());
}
@@ -1279,8 +1288,20 @@ public final class MessageService {
* @return the ticket to poll for the eventual result
*/
public String sendAsync(String target, String content, Runnable onAccepted) {
return sendAsync(target, content, onAccepted, null);
}
/**
* As {@link #sendAsync(String, String, Runnable)}, recording {@code creatorTerminal} as this
* ticket's owner. {@link #poll(String, String)} refuses a later caller whose own terminal
* differs from this one; {@code null} records no owner (a caller with no terminal — the
* unnamed primary — is always allowed to poll the result regardless).
*
* @return the ticket to poll for the eventual result
*/
public String sendAsync(String target, String content, Runnable onAccepted, String creatorTerminal) {
String ticket = "task-" + ticketSeq.incrementAndGet();
Task task = new Task(ticket, target, nowNanos);
Task task = new Task(ticket, target, nowNanos, creatorTerminal);
tasks.put(ticket, task);
if (pushLoop != null) {
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a
@@ -1343,15 +1364,31 @@ public final class MessageService {
}
/**
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket;
* otherwise a {@link Phase#PENDING} view (with the live worker status as detail), a
* {@link Phase#DONE} view carrying the reply, or a {@link Phase#FAILED} view with the reason.
* As {@link #poll(String, String)}, with no caller terminal — the ticket's ownership is never
* checked, so this overload must only be used where the caller's identity is otherwise
* irrelevant (a test, or a surface that does not resolve a caller terminal at all).
*/
public TaskView poll(String ticket) {
return poll(ticket, null);
}
/**
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket.
* Refuses a {@code callerTerminal} that differs from the terminal that created the ticket (see
* {@link #sendAsync(String, String, Runnable, String)}) with a {@link Phase#FAILED} view that
* carries no reply text — a caller with no terminal (the unnamed primary) is never refused.
* Otherwise returns a {@link Phase#PENDING} view (with the live worker status as detail), a
* {@link Phase#DONE} view carrying the reply, or a {@link Phase#FAILED} view with the reason.
*/
public TaskView poll(String ticket, String callerTerminal) {
Task task = tasks.get(ticket);
if (task == null) {
return null;
}
if (!ownsTicket(task, callerTerminal)) {
return new TaskView(ticket, Phase.FAILED, null, null,
"forbidden: this ticket was created by a different session", null);
}
CompletableFuture<Reply> f = task.future;
if (!f.isDone()) {
Reply question = task.question;
@@ -1387,6 +1424,17 @@ public final class MessageService {
return new TaskView(ticket, Phase.FAILED, null, null, detail, null);
}
/**
* Whether {@code callerTerminal} may read {@code task}'s state. A caller with no terminal
* always may — that is the unnamed primary, resolved by token or loopback trust, which never
* carries a herdr pane and must keep reading every ticket. Otherwise the caller's terminal must
* equal the terminal recorded on the task; a task with no recorded terminal matches no
* terminal-bearing caller.
*/
private static boolean ownsTicket(Task task, String callerTerminal) {
return callerTerminal == null || callerTerminal.equals(task.creatorTerminal);
}
/**
* Test seam only — carries no production behaviour, and nothing in this class calls it;
* {@link #pruneTerminalTickets} still reads {@link Task#completedNanos} directly.
@@ -850,8 +850,7 @@ class FleetConfigTest {
/**
* The hazard this guard closes: a pane-placed member lands inside the focused tab rather than
* its own, so it can land inside a lead's labelled tab and, while its pane carries no entry in
* the spawned-member roster, be read back as that lead.
* its own, so it can land inside a lead's labelled tab and be read back as that lead.
*/
@Test
void aPanePlacedProfileWithALeadTabRefusesToStart(@TempDir Path dir) throws Exception {
@@ -1088,9 +1087,8 @@ class FleetConfigTest {
}
/**
* A member tabLabel that can render as a configured collaborator tab is the same hazard as the
* lead case above — while its pane carries no entry in the spawned-member roster, a member
* labelled that way is read back as the collaborator.
* fleetd #669: a member tabLabel that can render as a configured collaborator tab is the same
* hazard as the lead case above — a member labelled that way is read back as the collaborator.
*/
@Test
void aProfileTabLabelOverrideMatchingACollaboratorTabRefusesToStart(@TempDir Path dir)
@@ -366,58 +366,6 @@ class FleetMcpAuthzTest {
+ "calling, not pass a literal boolean -- found: " + trailing);
}
// --- fleetd #703: who may see fleet_list's collaborators array -------------------------------
/**
* fleetd #703: {@link FleetMcp#collaboratorsVisibleTo} is the whole policy decision for
* {@code fleet_list}'s {@code collaborators} array. Visible to exactly the roles that may
* {@code SEND} to a named peer -- the primary, an architect, and a collaborator itself -- never
* a worker, which holds {@code READ} but can never {@code SEND} to a collaborator, and never an
* anonymous caller.
*/
@Test
void onlyPrimaryArchitectAndCollaboratorMaySeeTheCollaboratorsArray() {
assertTrue(FleetMcp.collaboratorsVisibleTo(PRIMARY), "the primary must see the collaborators array");
assertTrue(FleetMcp.collaboratorsVisibleTo(ARCH_DESIGN), "an architect must see the collaborators array");
assertTrue(FleetMcp.collaboratorsVisibleTo(COLLABORATOR), "a collaborator must see its own peer roster");
assertFalse(FleetMcp.collaboratorsVisibleTo(WORKER_A),
"a worker holds READ but can never SEND to a collaborator, so it must not see the array");
assertFalse(FleetMcp.collaboratorsVisibleTo(ANON), "authenticated as nothing must not see it either");
}
/**
* fleetd #703, same reasoning as {@link #theFleetListHandlerActuallyConsultsCoordinatorVisibleTo}:
* the predicate above can be perfectly correct while the one production call site never asks it.
* This reads {@code FleetMcp.java}'s own source and asserts the {@code fleet_list} handler's
* {@code listFleet(...)} call both threads {@code callers.collaborators()} into the payload and
* asks {@code collaboratorsVisibleTo(principal(exchange))} for the visibility flag, rather than a
* literal boolean or an empty map.
*/
@Test
void theFleetListHandlerActuallyConsultsCollaboratorsVisibleTo() throws Exception {
String source = Files.readString(MCP_SOURCE);
int start = source.indexOf("listHandler =");
assertTrue(start >= 0, "could not find the fleet_list handler (listHandler) in " + MCP_SOURCE
+ " -- the scrape has stopped matching, fix the anchor before trusting this test");
int end = source.indexOf("stopHandler =", start);
assertTrue(end > start, "could not find the handler declared after listHandler to bound the scrape");
String handlerBlock = source.substring(start, end);
// CONTROL: the block we scraped really does contain a call to listFleet(...) -- if this
// fails, the anchors above moved and the assertions below would otherwise pass on nothing.
assertTrue(handlerBlock.contains("listFleet("),
"control failed: the scraped listHandler block contains no listFleet( call at all -- "
+ "the anchors have drifted, this test is not testing what it claims to");
assertTrue(handlerBlock.contains("callers.collaborators()"),
"the fleet_list handler must thread callers.collaborators() into listFleet(...), not an "
+ "empty or literal map -- block: " + handlerBlock);
assertTrue(handlerBlock.contains("collaboratorsVisibleTo(principal(exchange))"),
"the fleet_list handler must ask collaboratorsVisibleTo(principal(exchange)) who is "
+ "calling, not pass a literal boolean -- block: " + handlerBlock);
}
// --- which action each tool hands the gate (fleetd #272) ------------------------------------
/**
@@ -121,7 +121,6 @@ class FleetMcpLeadContextGaugeWiringTest {
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.LoopHealthSource.none(), FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
FleetMcp.LeadSeatSource.none(), contextGauge, leadConfigDirs,
Map.of(LEAD_TERMINAL, LEAD_NAME), LEAD_TERMINAL, Map.of(), false,
FleetMcp.CoordinationSource.none(), false);
Map.of(LEAD_TERMINAL, LEAD_NAME), LEAD_TERMINAL, FleetMcp.CoordinationSource.none(), false);
}
}
@@ -10,7 +10,6 @@ import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.inject.LoopWatchdog;
import dev.ltms.fleet.lead.LeadContextGauge;
import dev.ltms.fleet.herdr.PaneLocator;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.Injector;
@@ -880,83 +879,6 @@ class FleetMcpTest {
assertTrue(out.contains("x".repeat(80) + "…"), "expected an 80-char preview with an ellipsis: " + out);
}
// --- fleetd #703: fleet_list's collaborators array -------------------------------------------
/** Calls the canonical {@code listFleet} overload directly, so a test can set the collaborators
* payload and its visibility independently of a real {@code Principal} / MCP exchange. */
private static McpSchema.CallToolResult listFleetWithCollaborators(FakeHerdr h,
Map<String, String> collaborators, boolean collaboratorsVisible) {
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
return FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.LoopHealthSource.none(), FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
FleetMcp.LeadSeatSource.none(), new LeadContextGauge(), FleetMcp.LeadConfigDirSource.none(),
Map.of(), "", collaborators, collaboratorsVisible,
FleetMcp.CoordinationSource.none(), false);
}
/**
* fleetd #703 acceptance A: a visible caller with one configured collaborator gets a
* {@code collaborators} row whose key ({@code CallerResolver.collaborators()}'s
* {@code terminal_id -> name} entry) lands as that row's {@code sessionId}, and {@code leads}/
* {@code members} are unaffected by the new key.
*/
@Test
void listReportsACollaboratorsRowKeyedByTheCollaboratorsTerminalId() {
FakeHerdr h = new FakeHerdr();
String out = textOf(listFleetWithCollaborators(h, Map.of("term_collab", "ops"), true));
assertTrue(out.contains("\"collaborators\":["), out);
assertTrue(out.contains("\"name\":\"ops\""), out);
assertTrue(out.contains("\"sessionId\":\"term_collab\""), out);
assertTrue(out.contains("\"leads\":[]"), out);
assertTrue(out.contains("\"members\":[]"), out);
}
/**
* fleetd #703 acceptance A, control half: with no collaborator configured, the key is absent
* (caller cannot see it) or an empty array (caller can), and {@code leads}/{@code members} are
* unchanged either way.
*/
@Test
void listOmitsOrEmptiesCollaboratorsWhenNoneAreConfigured() {
FakeHerdr h = new FakeHerdr();
String visibleButEmpty = textOf(listFleetWithCollaborators(h, Map.of(), true));
assertTrue(visibleButEmpty.contains("\"collaborators\":[]"), visibleButEmpty);
assertTrue(visibleButEmpty.contains("\"leads\":[]"), visibleButEmpty);
assertTrue(visibleButEmpty.contains("\"members\":[]"), visibleButEmpty);
String notVisible = textOf(listFleetWithCollaborators(h, Map.of(), false));
assertFalse(notVisible.contains("\"collaborators\""), notVisible);
assertTrue(notVisible.contains("\"leads\":[]"), notVisible);
assertTrue(notVisible.contains("\"members\":[]"), notVisible);
}
/**
* fleetd #703 acceptance B: a worker must not see the {@code collaborators} array at all, while
* an architect -- a role that also holds READ, same as a worker -- does see it. Both halves are
* asserted: a test that only checked the worker-hidden half would pass even if the feature were
* never wired up for anyone.
*/
@Test
void listOmitsCollaboratorsForAWorkerAndIncludesThemForAnArchitect() {
FakeHerdr h = new FakeHerdr();
Map<String, String> collaborators = Map.of("term_collab", "ops");
String asWorker = textOf(listFleetWithCollaborators(h, collaborators,
FleetMcp.collaboratorsVisibleTo(Principal.worker("term_w", 1))));
assertFalse(asWorker.contains("\"collaborators\""),
"a worker must not see the collaborators array: " + asWorker);
String asArchitect = textOf(listFleetWithCollaborators(h, collaborators,
FleetMcp.collaboratorsVisibleTo(Principal.architect("design", "term_arch", 2))));
assertTrue(asArchitect.contains("\"collaborators\":["),
"an architect must see the collaborators array: " + asArchitect);
assertTrue(asArchitect.contains("\"sessionId\":\"term_collab\""), asArchitect);
}
/**
* 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
@@ -865,11 +865,82 @@ class MessageServiceTest {
+ "follow-up send would come back BUSY instead of timing out on its own work");
}
/**
* Drive an async send on {@code T} to a resolved reply, polling as {@code owner} until the
* ticket reports {@link MessageService.Phase#DONE} (or the 2s deadline runs out). {@code owner}
* must be a terminal this ticket's creator check actually accepts, or this loops to the
* deadline and returns a non-{@code DONE} view.
*/
private MessageService.TaskView driveAsyncTicketToDone(String ticket, String owner) throws InterruptedException {
MessageService.TaskView view = null;
long deadline = System.currentTimeMillis() + 2000;
while (view == null || view.phase() != MessageService.Phase.DONE) {
if (System.currentTimeMillis() >= deadline) break;
view = messages.poll(ticket, owner);
//noinspection BusyWait
Thread.sleep(5);
}
return view;
}
@Test
void pollReturnsNullForAnUnknownTicket() {
assertNull(messages.poll("task-999999"), "a ticket that was never minted is unknown");
}
// --- fleetd #705: a ticket's creator terminal gates who may poll it -------------------------
@Test
void pollByAnotherTerminalIsRefused() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, "term_a");
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker works
assertTrue(rendezvous.resolve(T, "secret async result"), "a reply resolves the async send");
MessageService.TaskView owner = driveAsyncTicketToDone(ticket, "term_a");
assertNotNull(owner, "the creator must still be able to read its own ticket");
assertEquals(MessageService.Phase.DONE, owner.phase());
MessageService.TaskView refused = messages.poll(ticket, "term_b");
assertNotNull(refused, "a different terminal gets a refusal, not silence");
assertNotEquals(MessageService.Phase.DONE, refused.phase(),
"a different terminal must never see the ticket as DONE");
assertNull(refused.reply(), "a refusal must never carry the reply text");
assertFalse(String.valueOf(refused).contains("secret async result"),
"the reply text must not appear anywhere in the refused view");
}
@Test
void unnamedPrimaryStillReadsAnyTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, "term_lead");
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker works
assertTrue(rendezvous.resolve(T, "primary-visible result"), "a reply resolves the async send");
// callerTerminal == null is the unnamed primary (resolved by token or loopback trust, with
// no herdr pane) — it must read a ticket a terminal-bearing lead created.
MessageService.TaskView view = driveAsyncTicketToDone(ticket, null);
assertNotNull(view, "the unnamed primary must be able to read any ticket");
assertEquals(MessageService.Phase.DONE, view.phase());
assertEquals("primary-visible result", view.reply());
}
@Test
void creatorReadsItsOwnTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, "term_creator");
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker works
assertTrue(rendezvous.resolve(T, "own result"), "a reply resolves the async send");
MessageService.TaskView view = driveAsyncTicketToDone(ticket, "term_creator");
assertNotNull(view, "the ticket's own creator must be able to read it");
assertEquals(MessageService.Phase.DONE, view.phase());
assertEquals("own result", view.reply());
}
@Test
void pollReportsACompletedTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task");