Compare commits

...

11 Commits

Author SHA1 Message Date
Dai Ha d2db8c7dc9 fleetd #710 PR #717 review: leadsVisibleTo must also admit a collaborator
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Failing after 1m47s
A collaborator may SEND to a lead, and fleet_list's leads array is the
only place this tool gives it a lead's sessionId -- its own
fleet_whoami carries no lead address. leadsVisibleTo now returns true
for caller.isCollaborator() as well as primary and architect, matching
the existing rule for the sibling collaborators array (every role that
may SEND to a named peer). membersVisibleTo is unchanged: a
collaborator may never SEND to a spawned member.

Updated the truth-table tests for both predicates, added a behavioural
test proving a collaborator's fleet_list output contains leads and not
members, and updated fleet_list's tool description.
2026-10-04 07:37:55 +02:00
Dai Ha 10ab58e4fc fleetd #710: gate fleet_list's leads and members arrays by caller role
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 1m50s
Omit the leads and members keys entirely (never an empty array) from
fleet_list's result for a worker, matching the existing
coordinatorVisibleTo/collaboratorsVisibleTo pattern: two new named
predicates (leadsVisibleTo, membersVisibleTo) are consulted before
assembling either array, so a worker holding only READ can no longer
read every session on the daemon through this tool. Primary and
architect callers are unaffected.
2026-10-04 07:25:55 +02:00
Dai Ha c468953963 Merge PR #714: fleetd #711 — re-key two startup refusals onto the hazard that is still real
CI / shell-tests (push) Failing after 8s
CI / build (push) Failing after 1m53s
CI / contract (push) Successful in 50s
The roster is consulted ahead of every tab map, so a live registered member
is no longer read back as a lead. Both refusals still matter, for the narrower
case where the pane is alive and the roster holds no entry for it.

Behaviour unchanged. Verified in a throwaway worktree, not piped:
Tests run: 2008, FleetConfigTest 170, BUILD SUCCESS.
2026-10-04 07:04:21 +02:00
Dai Ha 25d53e6ef7 Merge PR #713: fleetd #710 — remove the unused CallerResolver.members() accessor
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 58s
CI / build (push) Failing after 2m3s
It returned architectTerminals, so its name contradicted its contents and
collided with fleet_list's members array. No production caller.

Verified in a throwaway worktree, not piped: Tests run: 2008, BUILD SUCCESS.
2026-10-04 07:00:35 +02:00
Dai Ha a9a37af957 Merge PR #712: fleetd #705 — gate fleet_poll's ticket lookup by the creating caller
Closes the MCP door only. The REST door (GET /tasks/{ticket}) still calls the
no-check poll overload, so #705 stays open.

Verified in a throwaway worktree, not piped: Tests run: 2008, BUILD SUCCESS.
2026-10-04 07:00:35 +02:00
Dai Ha 7d497aa423 fleetd #711: re-key two pane-tab hazard reasons onto the unregistered-pane condition
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 1m39s
Both refusals said a member landing in a lead's or collaborator's
labelled tab would be read back as that identity. CallerResolver
consults the spawned-member roster ahead of every tab map, so a live
registered member is never misread this way. State the real condition
instead: the hazard applies only while the pane is alive and carries
no entry in the spawned-member roster.
2026-10-04 06:54:43 +02:00
Dai Ha b205bcc2aa Remove CallerResolver members accessor
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 1m47s
2026-10-04 06:53:39 +02:00
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
Dai Ha 4939d40362 Merge PR #709: fleetd #703 — fleet_list reports collaborators, narrowed by role
CI / shell-tests (push) Failing after 12s
CI / contract (push) Successful in 48s
CI / build (push) Failing after 1m53s
A collaborator could find and message a lead, but a lead had nowhere to read a
collaborator's sessionId, so the channel only worked once the collaborator
spoke first. fleet_list now carries a collaborators array, and
CallerResolver.collaborators() has its first caller outside the resolver.

Visible to the primary, an architect and a collaborator; absent for a worker.
The rule is that the roles which may see a collaborator row are the roles that
can act on it, and a worker cannot send to a collaborator at all. Narrowed in
the payload the way coordinatorVisibleTo already does it, because fleet_list is
gated on READ and a worker holds READ.

Each row carries the registry name and the sessionId only. A collaborator is
never spawned and has no profile, so a lead row's context and seat fields have
no meaning for it.

mvn clean install: exit 0, BUILD SUCCESS, Tests run: 2005, Failures: 0.

Nothing yet pins the handler to collaboratorsVisibleTo, unlike the coordinator
flag, so a literal passed there would regress silently. Tracked in #710.
2026-10-04 06:42:05 +02:00
Dai Ha e7a7711a4e fleetd #703: fleet_list reports a collaborators array so a lead can discover one
CI / shell-tests (pull_request) Failing after 11s
CI / contract (pull_request) Successful in 58s
CI / build (pull_request) Failing after 1m54s
Threads CallerResolver#collaborators() into fleet_list's canonical listFleet
overload and adds a collaborators array (name, sessionId), visible only to
the primary, an architect, and a collaborator -- the same roles Authz grants
SEND to a named peer -- never a worker. Gives CallerResolver#collaborators()
its first real caller outside the resolver itself.
2026-10-04 06:37:48 +02:00
Dai Ha 1fd2cfa716 Merge PR #708: fleetd #669 — fleetd.example.yaml promised a defence that is not wired
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 1m10s
CI / build (push) Failing after 2m15s
The example file said two things stop the tab-name convention from becoming a
way to claim leadership, and named workspace exclusion as the first. Production
builds the scanner with an empty excludedWorkspaceLabels, so that defence does
not exist. It now names the two that do: the startup refusal of a colliding
tabPrefix/tabLabel, and CallerResolver asking the live spawned-member roster
before any tab map.

It also said a lead's workspace default is "leads" and must not be a member
workspace. The default is "fleet", the same space the members use, so that
advice was against the shipped shape.

Adds the test nobody wrote: a lead is still discovered when its workspace is
the member workspace, with an empty exclusion set. Corrects a test javadoc that
called itself the guard that matters while covering a parameter production never
passes.

mvn clean install: exit 0, BUILD SUCCESS, Tests run: 2000, Failures: 0.
2026-10-04 06:35:49 +02:00
13 changed files with 652 additions and 84 deletions
+10 -8
View File
@@ -62,11 +62,12 @@ bind:
#
# Leads are configured under `fleet.leaders:` — see THE FLEET further down.
#
# Two things stop the tab-name convention from becoming a way to claim leadership: the configured
# member spaces are excluded from the scan, so nothing fleetd places can land in a matching tab;
# and startup REFUSES a `tabPrefix` that the fleet tabLabel template, or any per-profile `tabLabel`
# override, also matches — so the two namespaces cannot overlap by accident. The label is a NAME,
# never a capability: what a pane may do is decided by the role the daemon resolves for it.
# Two things stop the tab-name convention from becoming a way to claim leadership: startup REFUSES
# a `tabPrefix` that the fleet tabLabel template, or any per-profile `tabLabel` override, also
# matches, so the two namespaces cannot overlap by accident; and the CallerResolver asks the live
# spawned-member roster BEFORE any tab map, so a live member is never mistaken for a lead no matter
# what its tab says. The label is a NAME, never a capability: what a pane may do is decided by the
# role the daemon resolves for it.
# CB-551: IDLE-LEAD HEARTBEAT — nudge the single lead back to work when it has been continuously
# idle (no open fleet_send driving it) past the quiet period. The fleet is one lead + architects +
@@ -659,9 +660,10 @@ fleet:
# tabPrefix: "lead:" # only used to guard against a worker tabLabel colliding with
# # this convention at startup; plays no part in matching a lead
# scanIntervalSeconds: 10 # rescan cadence, and the worst case before a new tab is seen
# workspace: leads # where a launched lead's tab is created (default "leads").
# # MUST NOT be a member workspace — those are excluded from the
# # scan, so a lead placed in one is never found again.
# workspace: leads # where a launched lead's tab is created (default "fleet",
# # the same shared space the members use). Sharing that space
# # with members is the normal shipped shape: the scanner tells
# # a lead from a member by the exact tab label, not by workspace.
# cwd: /path/to/repo # the launched lead's working directory (default: fleetd's own)
# kind: claude # descriptive; reported by fleet_whoami
# gpt-sol-5.6:
@@ -249,17 +249,6 @@ public final class CallerResolver {
return leadTerminals.get();
}
/**
* The currently-recognised architect slots, {@code terminal_id → slot name} (CB-548).
*
* <p>Read from the same supplier {@link #resolve} consults, so a slot that is <em>listed</em>
* here but would not <em>resolve</em> (or the reverse) cannot drift apart. Live for the same
* reason as {@link #leads()}.
*/
public Map<String, String> members() {
return architectTerminals.get();
}
/**
* The currently-recognised collaborator tabs, {@code terminal_id → name}.
*
@@ -2776,9 +2776,10 @@ public record FleetConfig(
});
if (!bad.isEmpty()) {
throw new IllegalStateException("refusing to start: " + String.join("; ", bad)
+ ". 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.");
+ ". 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.");
}
List<String> collisions = new ArrayList<>();
@@ -2877,7 +2878,8 @@ 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 would be read back as that lead or collaborator and granted that identity's authority.
* 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.
*
* <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.
@@ -2908,8 +2910,9 @@ 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 be read back as the "
+ "lead or collaborator, granted that identity's authority. Set placement: tab for "
+ "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 "
+ "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,8 +560,10 @@ 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)));
coordinatorVisibleTo(principal(exchange)),
leadsVisibleTo(principal(exchange)), membersVisibleTo(principal(exchange)));
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
(exchange, req) -> {
@@ -743,6 +745,48 @@ 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();
}
/**
* Who may see {@code fleet_list}'s {@code leads} array — exactly the roles that may
* {@link Authz.Action#SEND} to a lead: the primary, an architect, and a collaborator. A
* collaborator's own {@code fleet_whoami} carries no lead address, and {@code leads} is the
* only place this tool gives one, so a collaborator needs this array to use the send it
* already holds. A worker can never {@code SEND} at all, so it still sees neither this array
* nor {@code members}; a worker's own facts come from {@code fleet_whoami} instead. Split
* out for the same reason as {@link #coordinatorVisibleTo} and
* {@link #collaboratorsVisibleTo}: 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 leadsVisibleTo(Principal caller) {
return caller.isPrimary() || caller.isArchitect() || caller.isCollaborator();
}
/**
* Who may see {@code fleet_list}'s {@code members} array — the primary and an architect,
* which may {@link Authz.Action#SEND} to a member. A worker can never {@code SEND} at all,
* and a collaborator may {@code SEND} only to a lead or another collaborator, never to a
* spawned member, so neither sees this array even though {@link #leadsVisibleTo} grants a
* collaborator the sibling one.
*/
static boolean membersVisibleTo(Principal caller) {
return caller.isPrimary() || caller.isArchitect();
}
/** 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);
@@ -956,6 +1000,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");
}
@@ -963,7 +1018,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);
}
@@ -1142,9 +1197,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);
}
@@ -1158,7 +1228,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)");
}
@@ -1761,7 +1831,8 @@ 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, coordination, false);
OutageSource.none(), LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, false, true, true);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
@@ -1770,7 +1841,8 @@ 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, CoordinationSource.none(), false);
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, CoordinationSource.none(), false, true, true);
}
/**
@@ -1787,7 +1859,8 @@ 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, coordination, false);
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, false, true, true);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
@@ -1796,7 +1869,8 @@ 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, coordination, false);
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, false, true, true);
}
/**
@@ -1818,7 +1892,8 @@ 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, coordination, false);
leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, false, true, true);
}
/**
@@ -1835,14 +1910,23 @@ public final class FleetMcp {
* {@code false} (fleetd #463: a forgotten argument fails closed, not
* open), so a test that wants the {@code coordinator} row must pass
* an explicit {@code true}
* @param leadsVisible whether this caller may see the {@code leads} array (see
* {@link #leadsVisibleTo}) — unlike {@code callerIsPrimary}, this is
* also {@code true} for an architect or a collaborator, so it cannot
* be derived from {@code callerIsPrimary} alone
* @param membersVisible whether this caller may see the {@code members} array (see
* {@link #membersVisibleTo}); also {@code true} for an architect, but
* unlike {@code leadsVisible}, never for a collaborator
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, OutageSource outage,
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
CoordinationSource coordination, boolean callerIsPrimary) {
CoordinationSource coordination, boolean callerIsPrimary,
boolean leadsVisible, boolean membersVisible) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine,
outage, leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, callerIsPrimary);
outage, leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, callerIsPrimary, leadsVisible, membersVisible);
}
/**
@@ -1852,6 +1936,17 @@ 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}
* @param leadsVisible whether this caller may see the {@code leads} array (see
* {@link #leadsVisibleTo})
* @param membersVisible whether this caller may see the {@code members} array (see
* {@link #membersVisibleTo})
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
@@ -1860,32 +1955,56 @@ public final class FleetMcp {
LeadSeatSource leadSeats, LeadContextGauge contextGauge,
LeadConfigDirSource leadConfigDirs,
Map<String, String> leads, String selfTerm,
CoordinationSource coordination, boolean callerIsPrimary) {
Map<String, String> collaborators, boolean collaboratorsVisible,
CoordinationSource coordination, boolean callerIsPrimary,
boolean leadsVisible, boolean membersVisible) {
try {
Map<String, Agent> live = workers.list().stream()
.map(Agent.class::cast)
.filter(a -> a.terminalId() != null)
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
List<Map<String, Object>> leadRows = leads.entrySet().stream()
.sorted(Map.Entry.comparingByValue())
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm, contextGauge,
leadConfigDirs))
.toList();
// Neither row's assembly (leadView/memberCapacityView probing herdr for live status)
// runs unless at least one of them needs the live-agent lookup backing it.
Map<String, Agent> live = (leadsVisible || membersVisible)
? workers.list().stream()
.map(Agent.class::cast)
.filter(a -> a.terminalId() != null)
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b))
: Map.of();
// fleetd #209: this is the caller-driven fleet_list read that actually reports
// agentSessionId (via memberCapacityView -> SessionManager.rosterView), so it uses the
// resolving roster; the heartbeat/health/metrics timers stay on the plain sessions.roster().
List<MemberSession> roster = sessions.rosterResolved();
List<Map<String, Object>> out = roster.stream()
.map(s -> memberCapacityView(s, live.get(s.terminalId()), messages, capacity.clock().getAsLong()))
.toList();
Set<String> profiles = new java.util.TreeSet<>(capacity.configuredProfiles().get());
roster.stream().map(MemberSession::profile).forEach(profiles::add);
Map<String, Object> result = new LinkedHashMap<>();
result.put("leads", leadRows); result.put("members", out);
// READ is permission to enter this tool, not permission to receive every field it can
// build -- gate BEFORE assembling each row, so the key is absent rather than
// present-and-empty; a caller without either row gets its own facts from fleet_whoami.
if (leadsVisible) {
List<Map<String, Object>> leadRows = leads.entrySet().stream()
.sorted(Map.Entry.comparingByValue())
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm, contextGauge,
leadConfigDirs))
.toList();
result.put("leads", leadRows);
}
if (membersVisible) {
List<Map<String, Object>> out = roster.stream()
.map(s -> memberCapacityView(s, live.get(s.terminalId()), messages, capacity.clock().getAsLong()))
.toList();
result.put("members", out);
}
result.put("healthCoverage", healthCoverage.value().get());
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.
@@ -2197,6 +2316,20 @@ 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
@@ -2391,7 +2524,14 @@ public final class FleetMcp {
private static McpSchema.Tool listTool() {
return tool(FleetTool.LIST.wireName(),
"List the whole fleet the bridge tracks, in two parts. 'leads' are your PEERS — other "
"List the whole fleet the bridge tracks, in two parts. 'members' is visible to "
+ "the primary and an architect only. 'leads' is visible to those two AND a "
+ "collaborator — exactly the roles that may fleet_send to a lead, so a "
+ "collaborator can learn a lead's sessionId before using the send it already "
+ "holds. A worker holds READ to call this tool at all, but gets neither "
+ "array, never an empty one; a worker reads its own session, "
+ "profile, state, worktree, branch and owner from fleet_whoami instead. "
+ "'leads' are your PEERS — other "
+ "orchestrators, each with its sessionId (the address to fleet_send to), "
+ "name, live status, and 'self': true on your own row; this is how you "
+ "discover a peer lead without being told its address. 'members' are the "
@@ -2422,7 +2562,12 @@ 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.",
+ "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.",
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.
@@ -31,8 +31,7 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #670 — pins the {@code excludedWorkspaceLabels} argument {@link FleetdAssembly}'s
* production boot path passes to {@link LeadTabScanner} at {@code FleetdAssembly.java:265}
* ({@code Set.of()}).
* production boot path passes to {@link LeadTabScanner} ({@code Set.of()}).
*
* <p>{@code LeadTabScannerTest} already covers this constructor parameter, but it builds its own
* {@link LeadTabScanner} with its own set, so it tests the seam and proves nothing about the
@@ -191,7 +190,7 @@ class FleetdAssemblyLeadTabScannerExclusionTest {
Set<?> excluded = (Set<?>) excludedField.get(leads);
assertTrue(excluded.isEmpty(),
"FleetdAssembly.java:265 must pass an empty excludedWorkspaceLabels to "
"FleetdAssembly must pass an empty excludedWorkspaceLabels to "
+ "LeadTabScanner — scanning member tabs would demote the lead to a worker");
} finally {
assertNotNull(ports.shutdownHook, "control: assembly must capture its shutdown hook");
@@ -460,7 +460,6 @@ class CallerResolverTest {
assertTrue(members.bind("architect:lead-designer", "term_a")); // the later lifecycle binds the slot
assertEquals(Role.ARCHITECT, r.resolve("127.0.0.1", 42, null).role());
assertEquals("architect:lead-designer", r.members().get("term_a"));
}
@Test
@@ -850,7 +850,8 @@ 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 be read back as that lead.
* 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.
*/
@Test
void aPanePlacedProfileWithALeadTabRefusesToStart(@TempDir Path dir) throws Exception {
@@ -1087,8 +1088,9 @@ class FleetConfigTest {
}
/**
* 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.
* 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.
*/
@Test
void aProfileTabLabelOverrideMatchingACollaboratorTabRefusesToStart(@TempDir Path dir)
@@ -205,9 +205,12 @@ class LeadTabScannerTest {
}
/**
* The guard that matters: fleetd labels its own worker tabs, so if a worker space were scanned
* a naming accident would promote the fleet. The exclusion is by workspace, not by hoping the
* worker template never collides.
* Covers the {@code excludedWorkspaceLabels} parameter: a tab in an excluded workspace is never
* matched, whatever its label. Production always constructs this class with an empty set (CB-558,
* {@code FleetdAssembly}), so this parameter plays no part in the live guard against a worker
* tab being mistaken for a lead — that guard is {@code CallerResolver} asking the live
* spawned-member roster before any tab map. This test exists because the parameter still exists
* and is worth covering on its own terms.
*/
@Test
void aTabInAWorkerSpaceIsNeverALeadEvenWhenItsLabelMatches() {
@@ -219,6 +222,29 @@ class LeadTabScannerTest {
assertFalse(scanner(herdr, tabToName, new AtomicLong()).get().containsKey("term_impostor"));
}
/**
* The shipped shape (CB-558, {@code FleetdAssembly}): production always constructs this class
* with an empty {@code excludedWorkspaceLabels}, and a lead's {@code workspace:} default is the
* same shared {@code "fleet"} space the members use. A lead tab is still found when it sits in
* the exact same workspace as a member-labelled tab — the scanner tells them apart by the exact
* tab label, not by which workspace either one is in.
*/
@Test
void aLeadIsDiscoveredWhenItsWorkspaceIsTheSameAsTheMemberWorkspace() {
TopologyHerdr herdr = new TopologyHerdr()
.workspace("w1", "fleet")
.tab("w1:t1", "w1", "lead: opus-5.0")
.tab("w1:t2", "w1", "worker: gx10 #1")
.pane("w1:p1", "w1:t1", "term_opus")
.pane("w1:p2", "w1:t2", "term_worker");
LeadTabScanner s = new LeadTabScanner(herdr, Map.of("lead: opus-5.0", "opus-5.0"),
Set.of(), TTL, new AtomicLong()::get);
assertEquals("opus-5.0", s.get().get("term_opus"),
"a lead sharing the members' workspace is still discovered — the label, not the "
+ "workspace, is what matches it");
}
@Test
void aLabelWithNoConfiguredEntryIsIgnored() {
TopologyHerdr herdr = new TopologyHerdr().workspace("w1", "main")
@@ -355,10 +355,10 @@ class FleetMcpAuthzTest {
+ "the anchors have drifted, this test is not testing what it claims to");
Pattern trailingArg = Pattern.compile(
"listFleet\\([^;]*?,\\s*(coordinatorVisibleTo\\(principal\\(exchange\\)\\)|true|false)\\s*\\)\\s*;",
"listFleet\\([^;]*?,\\s*(coordinatorVisibleTo\\(principal\\(exchange\\)\\)|true|false)\\s*[,)]",
Pattern.DOTALL);
Matcher m = trailingArg.matcher(handlerBlock);
assertTrue(m.find(), "could not locate listFleet(...)'s trailing boolean argument in the "
assertTrue(m.find(), "could not locate listFleet(...)'s coordinator boolean argument in the "
+ "listHandler block -- the call shape changed, update this test's anchor: " + handlerBlock);
String trailing = m.group(1);
assertEquals("coordinatorVisibleTo(principal(exchange))", trailing,
@@ -366,6 +366,143 @@ 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);
}
// --- who may see fleet_list's leads and members arrays ---------------------------------------
/**
* {@link FleetMcp#leadsVisibleTo} is the whole policy decision for {@code fleet_list}'s
* {@code leads} array: visible to exactly the roles that may {@code SEND} to a lead -- the
* primary, an architect, and a collaborator -- never a worker, which holds {@code READ} but
* can never {@code SEND} at all, and never an anonymous caller.
*/
@Test
void primaryArchitectAndCollaboratorMaySeeTheLeadsArray() {
assertTrue(FleetMcp.leadsVisibleTo(PRIMARY), "the primary must see the leads array");
assertTrue(FleetMcp.leadsVisibleTo(ARCH_DESIGN), "an architect must see the leads array");
assertTrue(FleetMcp.leadsVisibleTo(COLLABORATOR),
"a collaborator may SEND to a lead, so it must see the leads array to learn where");
assertFalse(FleetMcp.leadsVisibleTo(WORKER_A),
"a worker holds READ but can never SEND, so it must not see the leads array");
assertFalse(FleetMcp.leadsVisibleTo(ANON), "authenticated as nothing must not see it either");
}
/**
* As {@link #primaryArchitectAndCollaboratorMaySeeTheLeadsArray}, for the {@code members}
* array -- but a collaborator may {@code SEND} only to a lead or another collaborator, never
* to a spawned member, so it must not see this one.
*/
@Test
void onlyPrimaryAndArchitectMaySeeTheMembersArray() {
assertTrue(FleetMcp.membersVisibleTo(PRIMARY), "the primary must see the members array");
assertTrue(FleetMcp.membersVisibleTo(ARCH_DESIGN), "an architect must see the members array");
assertFalse(FleetMcp.membersVisibleTo(WORKER_A),
"a worker holds READ but can never SEND, so it must not see the members array");
assertFalse(FleetMcp.membersVisibleTo(COLLABORATOR),
"a collaborator may SEND to a lead, never to a spawned member, so it must not see the members array");
assertFalse(FleetMcp.membersVisibleTo(ANON), "authenticated as nothing must not see it either");
}
/**
* 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 asks {@code leadsVisibleTo(principal(exchange))} for the leads
* visibility flag, rather than a literal boolean.
*/
@Test
void theFleetListHandlerActuallyConsultsLeadsVisibleTo() 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 assertion 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("leadsVisibleTo(principal(exchange))"),
"the fleet_list handler must ask leadsVisibleTo(principal(exchange)) who is calling, "
+ "not pass a literal boolean -- block: " + handlerBlock);
}
/** As {@link #theFleetListHandlerActuallyConsultsLeadsVisibleTo}, for {@code membersVisibleTo}. */
@Test
void theFleetListHandlerActuallyConsultsMembersVisibleTo() 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);
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("membersVisibleTo(principal(exchange))"),
"the fleet_list handler must ask membersVisibleTo(principal(exchange)) who is calling, "
+ "not pass a literal boolean -- block: " + handlerBlock);
}
// --- which action each tool hands the gate (fleetd #272) ------------------------------------
/**
@@ -121,6 +121,7 @@ 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, FleetMcp.CoordinationSource.none(), false);
Map.of(LEAD_TERMINAL, LEAD_NAME), LEAD_TERMINAL, Map.of(), false,
FleetMcp.CoordinationSource.none(), false, true, true);
}
}
@@ -10,6 +10,7 @@ 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;
@@ -691,7 +692,7 @@ class FleetMcpTest {
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
String out = textOf(res);
// fleetd #361: reports both which coord-id a peer must use to reach ME, and this daemon's
@@ -720,7 +721,7 @@ class FleetMcpTest {
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
String out = textOf(res);
assertTrue(out.contains("\"mailbox\":{\"status\":\"unknown\"}"), out);
@@ -779,19 +780,82 @@ class FleetMcpTest {
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));
boolean callerIsPrimary = Principal.worker("term_a", 1).isPrimary();
Principal worker = Principal.worker("term_a", 1);
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(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), callerIsPrimary);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), worker.isPrimary(),
FleetMcp.leadsVisibleTo(worker), FleetMcp.membersVisibleTo(worker));
String out = textOf(res);
assertFalse(out.contains("\"coordinator\""), "a worker must never see the coordinator key at all: " + out);
assertFalse(out.contains("mac-opus"), "no fragment of the coordinator row may leak either: " + out);
assertTrue(out.contains("\"leads\""), "the rest of the result must still be present: " + out);
assertTrue(out.contains("\"members\""), out);
assertFalse(out.contains("\"leads\""), "a worker must never see the leads key at all: " + out);
assertFalse(out.contains("\"members\""), "a worker must never see the members key at all: " + out);
assertTrue(out.contains("\"healthCoverage\""), "the rest of the result must still be present: " + out);
}
/**
* A worker's result must carry neither the {@code leads} nor the {@code members} key, and no
* fragment of either row leaks even though both are fully populated for this call -- a key
* check alone would pass on an implementation that still built the rows and only renamed or
* nested them.
*/
@Test
void listLeaksNoLeadOrMemberRowFragmentToAWorker() {
FakeHerdr h = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", "term_owner",
new WorktreeRequest("cb-304", null));
Principal worker = Principal.worker("term_a", 1);
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(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of("term_lead_x", "mac-opus"), "", FleetMcp.CoordinationSource.none(), worker.isPrimary(),
FleetMcp.leadsVisibleTo(worker), FleetMcp.membersVisibleTo(worker));
String out = textOf(res);
assertFalse(out.contains("\"leads\""), out);
assertFalse(out.contains("\"members\""), out);
assertFalse(out.contains(s.terminalId()), "no member row fragment may leak to a worker: " + out);
assertFalse(out.contains("term_owner"), "no member owner fragment may leak to a worker: " + out);
assertFalse(out.contains("mac-opus"), "no lead row fragment may leak to a worker: " + out);
assertTrue(out.contains("\"healthCoverage\""), out);
}
/**
* A collaborator may {@code SEND} to a lead, so it must see the {@code leads} array -- the
* only place {@code fleet_whoami} does not already give it a lead's address. It may never
* {@code SEND} to a spawned member, so the {@code members} key must stay absent for it, with
* no fragment of a populated member row leaking either.
*/
@Test
void listShowsLeadsButNotMembersToACollaborator() {
FakeHerdr h = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", "term_owner",
new WorktreeRequest("cb-304", null));
Principal collaborator = Principal.collaborator("ops", "term_collab", 600);
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(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of("term_lead_x", "mac-opus"), "", FleetMcp.CoordinationSource.none(), collaborator.isPrimary(),
FleetMcp.leadsVisibleTo(collaborator), FleetMcp.membersVisibleTo(collaborator));
String out = textOf(res);
assertTrue(out.contains("\"leads\""), "a collaborator must see the leads array: " + out);
assertTrue(out.contains("mac-opus"), "a collaborator must see the lead's name/address: " + out);
assertFalse(out.contains("\"members\""), "a collaborator must never see the members key: " + out);
assertFalse(out.contains(s.terminalId()), "no member row fragment may leak to a collaborator: " + out);
assertFalse(out.contains("term_owner"), "no member owner fragment may leak to a collaborator: " + out);
assertTrue(out.contains("\"healthCoverage\""), out);
}
@@ -807,16 +871,19 @@ class FleetMcpTest {
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));
boolean callerIsPrimary = Principal.architect("lead-designer", "term_design", 400).isPrimary();
Principal architect = Principal.architect("lead-designer", "term_design", 400);
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(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), callerIsPrimary);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), architect.isPrimary(),
FleetMcp.leadsVisibleTo(architect), FleetMcp.membersVisibleTo(architect));
String out = textOf(res);
assertFalse(out.contains("\"coordinator\""), "an architect must never see the coordinator key either: " + out);
assertTrue(out.contains("\"leads\""), "an architect must still see the leads array: " + out);
assertTrue(out.contains("\"members\""), "an architect must still see the members array: " + out);
}
/**
@@ -841,7 +908,7 @@ class FleetMcpTest {
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true));
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true));
assertTrue(gatedAsPrimary.contains("\"coordinator\""), gatedAsPrimary);
assertTrue(gatedAsPrimary.contains("\"selfId\":\"mac-opus\""), gatedAsPrimary);
@@ -850,6 +917,8 @@ class FleetMcpTest {
assertTrue(gatedAsPrimary.contains("\"heldDurable\""), gatedAsPrimary);
assertTrue(gatedAsPrimary.contains("\"held\""), gatedAsPrimary);
assertTrue(gatedAsPrimary.contains("\"peers\""), gatedAsPrimary);
assertTrue(gatedAsPrimary.contains("\"leads\""), "the primary must still see the leads array: " + gatedAsPrimary);
assertTrue(gatedAsPrimary.contains("\"members\""), "the primary must still see the members array: " + gatedAsPrimary);
}
@Test
@@ -868,7 +937,7 @@ class FleetMcpTest {
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
String out = textOf(res);
assertTrue(out.contains("\"msgId\":\"m1\""), out);
@@ -879,6 +948,83 @@ 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, true, true);
}
/**
* 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
@@ -900,7 +1046,7 @@ class FleetMcpTest {
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
String out = textOf(res);
assertTrue(out.contains("\"pending\":0"), out);
@@ -929,7 +1075,7 @@ class FleetMcpTest {
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
String out = textOf(res);
assertTrue(out.contains("\"heldDurable\":false"),
@@ -998,7 +1144,7 @@ class FleetMcpTest {
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "",
new FleetMcp.CoordinationSource(channel, List.of("fleet01-lead", "fleet02-lead", "fleet03-lead")), true);
new FleetMcp.CoordinationSource(channel, List.of("fleet01-lead", "fleet02-lead", "fleet03-lead")), true, true, true);
String out = textOf(res);
assertTrue(out.contains("\"coordId\":\"fleet01-lead\",\"status\":\"exists\",\"pending\":2,\"consumers\":1"), out);
@@ -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");