Compare commits
7 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| d2db8c7dc9 | |||
| 10ab58e4fc | |||
| c468953963 | |||
| 25d53e6ef7 | |||
| a9a37af957 | |||
| 7d497aa423 | |||
| 11eccc3a1b |
@@ -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.
|
||||
@@ -562,7 +562,8 @@ public final class FleetMcp {
|
||||
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) -> {
|
||||
@@ -759,6 +760,33 @@ public final class FleetMcp {
|
||||
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);
|
||||
@@ -972,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");
|
||||
}
|
||||
@@ -979,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);
|
||||
}
|
||||
|
||||
@@ -1158,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);
|
||||
}
|
||||
@@ -1174,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)");
|
||||
}
|
||||
@@ -1778,7 +1832,7 @@ public final class FleetMcp {
|
||||
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);
|
||||
Map.of(), false, coordination, false, true, true);
|
||||
}
|
||||
|
||||
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
|
||||
@@ -1788,7 +1842,7 @@ public final class FleetMcp {
|
||||
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);
|
||||
Map.of(), false, CoordinationSource.none(), false, true, true);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1806,7 +1860,7 @@ public final class FleetMcp {
|
||||
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);
|
||||
Map.of(), false, coordination, false, true, true);
|
||||
}
|
||||
|
||||
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
|
||||
@@ -1816,7 +1870,7 @@ public final class FleetMcp {
|
||||
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);
|
||||
Map.of(), false, coordination, false, true, true);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1839,7 +1893,7 @@ public final class FleetMcp {
|
||||
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);
|
||||
Map.of(), false, coordination, false, true, true);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1856,15 +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,
|
||||
Map.of(), false, coordination, callerIsPrimary);
|
||||
Map.of(), false, coordination, callerIsPrimary, leadsVisible, membersVisible);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1881,6 +1943,10 @@ public final class FleetMcp {
|
||||
* {@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,
|
||||
@@ -1890,28 +1956,41 @@ public final class FleetMcp {
|
||||
LeadConfigDirSource leadConfigDirs,
|
||||
Map<String, String> leads, String selfTerm,
|
||||
Map<String, String> collaborators, boolean collaboratorsVisible,
|
||||
CoordinationSource coordination, boolean callerIsPrimary) {
|
||||
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(),
|
||||
@@ -2445,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 "
|
||||
|
||||
@@ -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,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)
|
||||
|
||||
@@ -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,
|
||||
@@ -418,6 +418,91 @@ class FleetMcpAuthzTest {
|
||||
+ "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) ------------------------------------
|
||||
|
||||
/**
|
||||
|
||||
@@ -122,6 +122,6 @@ class FleetMcpLeadContextGaugeWiringTest {
|
||||
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);
|
||||
FleetMcp.CoordinationSource.none(), false, true, true);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -692,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
|
||||
@@ -721,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);
|
||||
@@ -780,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);
|
||||
}
|
||||
|
||||
@@ -808,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);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -842,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);
|
||||
@@ -851,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
|
||||
@@ -869,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);
|
||||
@@ -893,7 +961,7 @@ class FleetMcpTest {
|
||||
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);
|
||||
FleetMcp.CoordinationSource.none(), false, true, true);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -978,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);
|
||||
@@ -1007,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"),
|
||||
@@ -1076,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");
|
||||
|
||||
Reference in New Issue
Block a user