Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 11eccc3a1b |
@@ -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");
|
||||
|
||||
Reference in New Issue
Block a user