Compare commits
23 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| e5f4fb81ab | |||
| d28ab0968b | |||
| 9a64d42599 | |||
| f4e0ca41e6 | |||
| 29a2f97c25 | |||
| d2db8c7dc9 | |||
| 95311c6e8e | |||
| 10ab58e4fc | |||
| 0ba597e394 | |||
| c468953963 | |||
| 25d53e6ef7 | |||
| a9a37af957 | |||
| 7d497aa423 | |||
| b205bcc2aa | |||
| 11eccc3a1b | |||
| 4939d40362 | |||
| e7a7711a4e | |||
| 1fd2cfa716 | |||
| 3e8e314656 | |||
| 20fc42b572 | |||
| f82073717a | |||
| 16b52fac6c | |||
| 13b6ae2628 |
@@ -168,7 +168,7 @@ you decide.
|
||||
| Answer a member's `fleet_ask` | `fleet_send{turnId, content}` — **not** `sessionId` |
|
||||
| Message a **peer lead** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task |
|
||||
| Message a **peer lead** on another daemon or host | `fleet_send{coordId: <their coord-id>, content}` — needs a `coordinator:` block; your own coord-id is in `fleet_list`. Coordination only, **never** a task |
|
||||
| Message a **collaborator** on this host | `fleet_send{sessionId: <their terminal>, content}` — but **`fleet_list` does not report collaborators**, so you cannot discover one: it must tell you its `sessionId`, which its own `fleet_whoami` gives it. Coordination only, **never** a task |
|
||||
| Message a **collaborator** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` reports a `collaborators` array, and each row carries that peer's `name` and the `sessionId` you send to. It is visible to you, to an architect and to another collaborator, never to a worker. Coordination only, **never** a task |
|
||||
| Answer a peer lead that messaged you | `fleet_send{coordId}` — or `{sessionId}` if they are on this host. **Not** `fleet_reply`: it has no peer route and the publish is refused |
|
||||
| Read your own held lead-to-lead mail (no ack) | `fleet_poll{coordId: <your own coord-id, from fleet_list's coordinator.selfId>}` — primary-only; never acks, so `fleet_list`'s `held[]` still shows it after. `fleet_list`'s `held[]` gives only a truncated preview — this is the only way to read the full body |
|
||||
| Collect a held reply | `fleet_poll{target}` · then `fleet_ack{target, msgId}` |
|
||||
@@ -257,8 +257,8 @@ of those is refused at the gate, not queued.
|
||||
Two limits worth knowing before you hit them. **You cannot reach a worker** — not even to help one —
|
||||
because a worker belongs to the lead that spawned it, and routing around that would make you a
|
||||
second orchestrator with no plan. Send to the lead instead. And **you cannot read a ticket**, so you
|
||||
cannot collect a delegation's reply; ticket ids are a plain counter with no owner check, so holding
|
||||
one would let you walk every other session's answers.
|
||||
cannot collect a delegation's reply: `fleet_poll` refuses you at the role gate, and a ticket also
|
||||
records the terminal that created it, so even a leaked id reads nothing.
|
||||
|
||||
Being named buys you a channel, not authority. Your `fleet_send` to a lead is coordination between
|
||||
peers: the lead owes you no obedience, and you owe it none.
|
||||
|
||||
+22
-13
@@ -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:
|
||||
@@ -669,11 +671,18 @@ fleet:
|
||||
# kind: opencode
|
||||
# model: openai/gpt-5.6-terra
|
||||
|
||||
# A collaborator tab, keyed by name (fleetd #669). This block is parsed and validated today;
|
||||
# nothing yet recognises or addresses the tab it names. Recognise-only, like a profile-less
|
||||
# `leaders:` entry above: there is no `profile:`, no `instances:` and no `kind:`. `tab:` is
|
||||
# REQUIRED and is the only field identity depends on, matched case-insensitively — the same
|
||||
# GET-THE-VALUE-RIGHT warning above the `leaders:` block applies here too.
|
||||
# A collaborator tab, keyed by name (fleetd #669). A pane whose tab matches resolves as the
|
||||
# COLLABORATOR role: `fleet_whoami` answers `collaborator`, and the session may observe the fleet
|
||||
# (`fleet_list`, `fleet_profiles`, `fleet_whoami`) and `fleet_send` to a lead or to another
|
||||
# collaborator. It may NOT spawn, stop or drain anything, roll a lead's session, answer a
|
||||
# member's `fleet_ask`, poll a ticket, send across hosts, or send to a spawned member's terminal.
|
||||
# Each of those is refused at the gate rather than queued.
|
||||
# Recognise-only, like a profile-less `leaders:` entry above: there is no `profile:`, no
|
||||
# `instances:` and no `kind:`. `tab:` is REQUIRED and is the only field identity depends on,
|
||||
# matched case-insensitively — the same GET-THE-VALUE-RIGHT warning above the `leaders:` block
|
||||
# applies here too.
|
||||
# A lead cannot discover a collaborator yet (#703): `fleet_list` has no `collaborators` key, so
|
||||
# the collaborator must speak first, or pass on the `sessionId` its own `fleet_whoami` reports.
|
||||
# collaborators:
|
||||
# reviewer-alex:
|
||||
# tab: "collab: alex"
|
||||
|
||||
@@ -217,25 +217,28 @@ public final class Fleetd {
|
||||
|
||||
/**
|
||||
* The {@link Injector}'s readiness gate (CB-534): a target is deliverable if it is a spawned
|
||||
* member whose agent has connected the bridge MCP, <em>or</em> a lead.
|
||||
* member whose agent has connected the bridge MCP, a lead, <em>or</em> a collaborator.
|
||||
*
|
||||
* <p>The gate exists for one reason — to hold a delivery out of a <em>spawned</em> member's boot
|
||||
* window, where herdr already reports {@code idle} but the TUI would drop an injected paste. That
|
||||
* hazard is a property of spawning. A lead is never spawned: the operator started it and named it
|
||||
* (or labelled its tab) only once it was up, so there is no boot window to guard.
|
||||
* (or labelled its tab) only once it was up, so there is no boot window to guard. A collaborator
|
||||
* is the same way — a person's own tab, matched to a configured name, never spawned.
|
||||
*
|
||||
* <p>A lead is also never enrolled in {@link MemberPresence} — {@code FleetMcp} marks presence
|
||||
* for every spawned member (worker and architect), deliberately, since that map doubles as the
|
||||
* member roster's availability signal and a lead counted there would show up as an available
|
||||
* member. So without the second disjunct a lead is permanently un-deliverable: every
|
||||
* lead→lead send sat on the gate for {@code READINESS_GRACE_POLLS} (~60s) and then failed
|
||||
* having never been typed into the pane.
|
||||
* <p>Neither a lead nor a collaborator is ever enrolled in {@link MemberPresence} — {@code
|
||||
* FleetMcp} marks presence for every spawned member (worker and architect), deliberately, since
|
||||
* that map doubles as the member roster's availability signal and a lead or collaborator counted
|
||||
* there would show up as an available member. So without the second and third disjuncts a lead or
|
||||
* collaborator is permanently un-deliverable: every send to one sat on the gate for
|
||||
* {@code READINESS_GRACE_POLLS} (~60s) and then failed having never been typed into the pane.
|
||||
*
|
||||
* <p>The lead set is read through the supplier on each call rather than snapshotted, so a lead
|
||||
* discovered by {@code leadScan} after startup becomes deliverable without a restart.
|
||||
* <p>Both sets are read through their supplier on each call rather than snapshotted, so a lead or
|
||||
* collaborator discovered by {@code leadScan} after startup becomes deliverable without a restart.
|
||||
*/
|
||||
static Predicate<String> deliverableTo(MemberPresence presence, Supplier<Map<String, String>> leads) {
|
||||
return target -> presence.isPresent(target) || leads.get().containsKey(target);
|
||||
static Predicate<String> deliverableTo(MemberPresence presence, Supplier<Map<String, String>> leads,
|
||||
Supplier<Map<String, String>> collaborators) {
|
||||
return target -> presence.isPresent(target) || leads.get().containsKey(target)
|
||||
|| collaborators.get().containsKey(target);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -368,7 +368,7 @@ final class FleetdAssembly {
|
||||
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
|
||||
MemberPresence presence = sessions.asPresence();
|
||||
TurnListener turnListener = Fleetd.turnListener(completion, sessions);
|
||||
Predicate<String> deliverable = Fleetd.deliverableTo(presence, leads);
|
||||
Predicate<String> deliverable = Fleetd.deliverableTo(presence, leads, collaboratorTerminals);
|
||||
// fleetd #556: registration is wired directly to `completion`, not folded into the
|
||||
// `turnListener` fan-out above — so it survives `sessions.onDelivered` (or any future
|
||||
// listener) throwing, regardless of call order.
|
||||
|
||||
@@ -75,6 +75,9 @@ public final class Authz {
|
||||
* {@code SEND} is refused, as if no terminal were a configured lead or collaborator — the
|
||||
* same decision {@link #NO_KNOWN_LEAD_OR_COLLABORATOR} gives explicitly. Every other action's
|
||||
* result is identical to the four-argument form's, since none of them consult the classifier.
|
||||
*
|
||||
* <p>Its default classifier denies every collaborator, so a caller enforcing authorization
|
||||
* must use the four-argument form instead.
|
||||
*/
|
||||
public static boolean permits(Principal caller, Action action, String targetSession) {
|
||||
return permits(caller, action, targetSession, NO_KNOWN_LEAD_OR_COLLABORATOR);
|
||||
|
||||
@@ -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}.
|
||||
*
|
||||
|
||||
@@ -615,7 +615,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
// may bind to, AND what a slot already bound still grants) through its own instance of that
|
||||
// same supplier shape — see MemberRegistry.live and its class doc for the binding rule:
|
||||
// removing a slot revokes ARCHITECT on the bound pane's very next request, and only the slot
|
||||
// OCCUPANCY survives, so the demoted session keeps its slot key until it unbinds. Only
|
||||
// OCCUPANCY survives, so the demoted session keeps its slot key until it unbinds.
|
||||
// fleet.leaders is frozen (Fleetd.java:281 reads cfg.fleet().leaders() off the startup
|
||||
// snapshot to build both the LeadTabScanner's tab-label-to-name map, wired into
|
||||
// CallerResolver.withLeadsAndMembers at Fleetd.java:620/624, and — when herdr answered —
|
||||
@@ -635,6 +635,11 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
+ "live through that same supplier for placement AND through a separate supplier "
|
||||
+ "on MemberRegistry for spawn-time identity — both already applied");
|
||||
}
|
||||
if (!Objects.equals(collaboratorsOf(old), collaboratorsOf(fresh))) {
|
||||
changed.add("fleet: fleet.collaborators (each collaborator's tab) is read once at "
|
||||
+ "startup to build the LeadTabScanner's identity map, which is not rebuilt on "
|
||||
+ "reload, so a collaborator added, removed, or given a new tab: needs a restart");
|
||||
}
|
||||
// Kept in step with SPLIT_KEYS the same way changedColdKeys is kept in step with COLD_KEYS —
|
||||
// every message here must be traceable to one of the split keys the class doc documents.
|
||||
// NOTE what this does NOT prove, per the javadoc above: it does not catch a SPLIT_KEYS
|
||||
@@ -650,6 +655,11 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
return cfg.fleet() == null ? Map.of() : cfg.fleet().leaders();
|
||||
}
|
||||
|
||||
/** {@code cfg.fleet().collaborators()}, defensively, in case a caller hands in a non-defaulted config. */
|
||||
private static Map<String, FleetConfig.Collaborator> collaboratorsOf(FleetConfig cfg) {
|
||||
return cfg.fleet() == null ? Map.of() : cfg.fleet().collaborators();
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link FleetConfig.Profile} record components deliberately left out of
|
||||
* {@link #sameLaunchSettings} because they are read <em>live</em>, not baked in at spawn — see
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -694,7 +694,8 @@ public final class FleetApp {
|
||||
|
||||
if (!wait) {
|
||||
// Fire-and-poll (CB-107): return a ticket immediately; the caller polls GET /tasks/{ticket}.
|
||||
String ticket = messages.sendAsync(id, content);
|
||||
Principal caller = ctx.attribute(CALLER);
|
||||
String ticket = messages.sendAsync(id, content, null, caller == null ? null : caller.terminal());
|
||||
ctx.status(202).json(Map.of("sessionId", id, "ticket", ticket, "status", "accepted"));
|
||||
return;
|
||||
}
|
||||
@@ -894,7 +895,8 @@ public final class FleetApp {
|
||||
if (!allow(ctx, routeAction("GET /tasks/{ticket}"), null)) {
|
||||
return;
|
||||
}
|
||||
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"));
|
||||
Principal caller = ctx.attribute(CALLER);
|
||||
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"), caller == null ? null : caller.terminal());
|
||||
if (v == null) {
|
||||
ctx.status(404).json(Map.of("error", "unknown_ticket", "detail", "no such task (or it has expired)"));
|
||||
return;
|
||||
|
||||
@@ -14,10 +14,12 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* CB-534: the injector's readiness gate must open for a lead as well as for a present worker.
|
||||
* fleetd #669 follow-up: the same gate must also open for a collaborator, which — like a lead —
|
||||
* is never enrolled in {@link MemberPresence} and never discovered by the lead scan.
|
||||
*
|
||||
* <p>The bug these cover was silent and slow: a lead was never marked present (only workers are), so
|
||||
* every lead→lead delivery sat on the gate for the full readiness grace and failed ~60s later without
|
||||
* a keystroke ever reaching the pane.
|
||||
* <p>The bug these cover was silent and slow: a lead (and later a collaborator) was never marked
|
||||
* present (only workers are) and never counted as a lead, so every send to one sat on the gate for
|
||||
* the full readiness grace and failed ~60s later without a keystroke ever reaching the pane.
|
||||
*/
|
||||
class FleetDeliverabilityTest {
|
||||
|
||||
@@ -25,19 +27,25 @@ class FleetDeliverabilityTest {
|
||||
return () -> m;
|
||||
}
|
||||
|
||||
private static Supplier<Map<String, String>> collaborators(Map<String, String> m) {
|
||||
return () -> m;
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a worker that has connected its MCP is deliverable")
|
||||
void presentWorkerIsDeliverable() {
|
||||
MemberPresence presence = new MemberPresence();
|
||||
presence.markPresent("term_worker");
|
||||
|
||||
assertTrue(Fleetd.deliverableTo(presence, leads(Map.of())).test("term_worker"));
|
||||
assertTrue(Fleetd.deliverableTo(presence, leads(Map.of()), collaborators(Map.of()))
|
||||
.test("term_worker"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a worker still in its boot window is held back")
|
||||
void absentWorkerIsNotDeliverable() {
|
||||
assertFalse(Fleetd.deliverableTo(new MemberPresence(), leads(Map.of())).test("term_booting"));
|
||||
assertFalse(Fleetd.deliverableTo(new MemberPresence(), leads(Map.of()), collaborators(Map.of()))
|
||||
.test("term_booting"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -45,19 +53,31 @@ class FleetDeliverabilityTest {
|
||||
void leadIsDeliverableWithoutPresence() {
|
||||
MemberPresence presence = new MemberPresence();
|
||||
Predicate<String> deliverable =
|
||||
Fleetd.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")));
|
||||
Fleetd.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")), collaborators(Map.of()));
|
||||
|
||||
assertFalse(presence.isPresent("term_lead"), "a lead is never enrolled in worker presence");
|
||||
assertTrue(deliverable.test("term_lead"), "…and must be deliverable anyway");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("an unknown terminal is deliverable to neither")
|
||||
@DisplayName("a collaborator is deliverable without ever being marked present or scanned as a lead")
|
||||
void collaboratorIsDeliverableWithoutPresenceOrLeadStatus() {
|
||||
MemberPresence presence = new MemberPresence();
|
||||
Predicate<String> deliverable = Fleetd.deliverableTo(presence, leads(Map.of()),
|
||||
collaborators(Map.of("term_collab", "kevin")));
|
||||
|
||||
assertFalse(presence.isPresent("term_collab"), "a collaborator is never enrolled in worker presence");
|
||||
assertTrue(deliverable.test("term_collab"), "…and must be deliverable anyway");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("an unknown terminal is deliverable to none of presence, leads, or collaborators")
|
||||
void strangerIsNotDeliverable() {
|
||||
MemberPresence presence = new MemberPresence();
|
||||
presence.markPresent("term_worker");
|
||||
|
||||
assertFalse(Fleetd.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")))
|
||||
assertFalse(Fleetd.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")),
|
||||
collaborators(Map.of("term_collab", "kevin")))
|
||||
.test("term_stranger"));
|
||||
}
|
||||
|
||||
@@ -65,22 +85,47 @@ class FleetDeliverabilityTest {
|
||||
@DisplayName("a lead discovered after startup becomes deliverable with no restart")
|
||||
void leadSetIsReadThroughOnEveryCall() {
|
||||
Map<String, String> discovered = new HashMap<>();
|
||||
Predicate<String> deliverable = Fleetd.deliverableTo(new MemberPresence(), leads(discovered));
|
||||
Predicate<String> deliverable =
|
||||
Fleetd.deliverableTo(new MemberPresence(), leads(discovered), collaborators(Map.of()));
|
||||
|
||||
assertFalse(deliverable.test("term_late"));
|
||||
discovered.put("term_late", "gpt-sol-5.6"); // leadScan picks up a newly labelled tab
|
||||
assertTrue(deliverable.test("term_late"), "the supplier must be re-read, not snapshotted");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a collaborator discovered after startup becomes deliverable with no restart")
|
||||
void collaboratorSetIsReadThroughOnEveryCall() {
|
||||
Map<String, String> discovered = new HashMap<>();
|
||||
Predicate<String> deliverable =
|
||||
Fleetd.deliverableTo(new MemberPresence(), leads(Map.of()), collaborators(discovered));
|
||||
|
||||
assertFalse(deliverable.test("term_late_collab"));
|
||||
discovered.put("term_late_collab", "kevin"); // the same tab scan picks up a newly labelled collaborator tab
|
||||
assertTrue(deliverable.test("term_late_collab"), "the supplier must be re-read, not snapshotted");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("forgetting a torn-down worker does not strip a lead of its deliverability")
|
||||
void forgetDoesNotDisarmALead() {
|
||||
MemberPresence presence = new MemberPresence();
|
||||
Predicate<String> deliverable =
|
||||
Fleetd.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")));
|
||||
Fleetd.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")), collaborators(Map.of()));
|
||||
|
||||
presence.forget("term_lead"); // the injector's cleanup path runs against every target
|
||||
|
||||
assertTrue(deliverable.test("term_lead"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("forgetting a torn-down worker does not strip a collaborator of its deliverability")
|
||||
void forgetDoesNotDisarmACollaborator() {
|
||||
MemberPresence presence = new MemberPresence();
|
||||
Predicate<String> deliverable = Fleetd.deliverableTo(presence, leads(Map.of()),
|
||||
collaborators(Map.of("term_collab", "kevin")));
|
||||
|
||||
presence.forget("term_collab"); // the injector's cleanup path runs against every target
|
||||
|
||||
assertTrue(deliverable.test("term_collab"));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -107,6 +107,8 @@ class ConfigRefTest {
|
||||
|
||||
assertTrue(out.applied());
|
||||
assertTrue(out.deferred().isEmpty());
|
||||
assertTrue(out.split().stream().noneMatch(s -> s.contains("fleet.collaborators")),
|
||||
out.split().toString());
|
||||
assertEquals("new charter", ref.get().fleet().charterFor(
|
||||
dev.ltms.fleet.peer.MemberRole.ARCHITECT));
|
||||
}
|
||||
@@ -786,14 +788,100 @@ class ConfigRefTest {
|
||||
assertTrue(out.split().getFirst().startsWith("fleet:"), out.split().toString());
|
||||
assertTrue(out.split().getFirst().contains("restart"), out.split().toString());
|
||||
assertTrue(out.split().getFirst().contains("live"), out.split().toString());
|
||||
assertTrue(out.split().stream().noneMatch(s -> s.contains("fleet.collaborators")),
|
||||
out.split().toString());
|
||||
assertTrue(out.summary().contains("partially live"), out.summary());
|
||||
// The snapshot still carries the new value — the LeadTabScanner's identity map and
|
||||
// LeadLauncher's auto-launch are what wait for a restart; a reload rebuilds neither.
|
||||
assertEquals("lead: opus-b", ref.get().fleet().leaders().get("opus").tab());
|
||||
}
|
||||
|
||||
@Test
|
||||
void addingAFleetCollaboratorIsReportedAsSplit(@TempDir Path dir) throws Exception {
|
||||
assertCollaboratorChangeIsReported(dir, """
|
||||
fleet:
|
||||
collaborators:
|
||||
alex:
|
||||
tab: "collaborator: alex"
|
||||
""", "add a collaborator");
|
||||
}
|
||||
|
||||
@Test
|
||||
void removingAFleetCollaboratorIsReportedAsSplit(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml("""
|
||||
fleet:
|
||||
collaborators:
|
||||
alex:
|
||||
tab: "collaborator: alex"
|
||||
"""));
|
||||
ConfigRef ref = refFor(f);
|
||||
|
||||
Files.writeString(f, yaml("fleet: {}\n"));
|
||||
assertCollaboratorSplit(ref.reload(), "remove a collaborator");
|
||||
}
|
||||
|
||||
@Test
|
||||
void changingAFleetCollaboratorTabIsReportedAsSplit(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml("""
|
||||
fleet:
|
||||
collaborators:
|
||||
alex:
|
||||
tab: "collaborator: alex-a"
|
||||
"""));
|
||||
ConfigRef ref = refFor(f);
|
||||
|
||||
Files.writeString(f, yaml("""
|
||||
fleet:
|
||||
collaborators:
|
||||
alex:
|
||||
tab: "collaborator: alex-b"
|
||||
"""));
|
||||
assertCollaboratorSplit(ref.reload(), "change a collaborator tab");
|
||||
}
|
||||
|
||||
@Test
|
||||
void unchangedFleetCollaboratorsProduceNoSplitReport(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
String config = yaml("""
|
||||
fleet:
|
||||
collaborators:
|
||||
alex:
|
||||
tab: "collaborator: alex"
|
||||
""");
|
||||
Files.writeString(f, config);
|
||||
ConfigRef ref = refFor(f);
|
||||
|
||||
Files.writeString(f, config);
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
|
||||
assertTrue(out.applied());
|
||||
assertTrue(out.split().isEmpty(), out.split().toString());
|
||||
}
|
||||
|
||||
private static void assertCollaboratorChangeIsReported(Path dir, String changed, String action)
|
||||
throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml("fleet: {}\n"));
|
||||
ConfigRef ref = refFor(f);
|
||||
|
||||
Files.writeString(f, yaml(changed));
|
||||
assertCollaboratorSplit(ref.reload(), action);
|
||||
}
|
||||
|
||||
private static void assertCollaboratorSplit(ConfigRef.Outcome out, String action) {
|
||||
assertTrue(out.applied(), action);
|
||||
assertTrue(out.deferred().isEmpty(), out.deferred().toString());
|
||||
assertEquals(1, out.split().size(), out.split().toString());
|
||||
String report = out.split().getFirst();
|
||||
assertTrue(report.startsWith("fleet: fleet.collaborators"), report);
|
||||
assertTrue(report.contains("read once at startup"), report);
|
||||
assertTrue(report.contains("restart"), report);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #333: {@code fleet.leaders} is the ONLY frozen part of {@code fleet:}. A reload that
|
||||
* fleetd #333: {@code fleet.leaders} is a frozen part of {@code fleet:}. A reload that
|
||||
* changes {@code tabLabel} (or charters, or a role pool) without touching {@code fleet.leaders}
|
||||
* must stay fully hot with nothing reported — proving {@link ConfigRef#changedSplitKeys}
|
||||
* compares {@code fleet.leaders} specifically rather than the whole {@code Fleet} record, which
|
||||
|
||||
@@ -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,170 @@ 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);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code fleet_poll{ticket}} must thread the calling connection's own terminal into
|
||||
* {@link MessageService#poll(String, String)}, so a worker cannot read a ticket a different
|
||||
* session created.
|
||||
*/
|
||||
@Test
|
||||
void theFleetPollHandlerActuallyThreadsCallerTerminalIntoPoll() throws Exception {
|
||||
String source = Files.readString(MCP_SOURCE);
|
||||
|
||||
int start = source.indexOf("pollHandler =");
|
||||
assertTrue(start >= 0, "could not find the fleet_poll handler (pollHandler) in " + MCP_SOURCE
|
||||
+ " -- the scrape has stopped matching, fix the anchor before trusting this test");
|
||||
int end = source.indexOf("ackHandler =", start);
|
||||
assertTrue(end > start, "could not find the handler declared after pollHandler to bound the scrape");
|
||||
String handlerBlock = source.substring(start, end);
|
||||
|
||||
// CONTROL: the block we scraped really does call poll(...) -- if this fails, the anchors
|
||||
// above moved and the assertions below would otherwise pass on nothing.
|
||||
assertTrue(handlerBlock.contains("poll(messages,"),
|
||||
"control failed: the scraped pollHandler block contains no poll(messages, ...) call "
|
||||
+ "at all -- the anchors have drifted, this test is not testing what it claims to");
|
||||
|
||||
assertTrue(handlerBlock.contains("callerTerminal(exchange)"),
|
||||
"the fleet_poll handler must thread callerTerminal(exchange) into poll(...), not omit "
|
||||
+ "it or pass a literal null -- 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);
|
||||
|
||||
@@ -0,0 +1,280 @@
|
||||
package dev.ltms.fleet.msg;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.regex.Matcher;
|
||||
import java.util.regex.Pattern;
|
||||
import java.util.stream.Stream;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* Pins that no file under {@code src/main/java} calls the fail-open
|
||||
* {@link MessageService#poll(String)} overload. That overload skips the ownership check in
|
||||
* {@code MessageService}'s {@code ownsTicket} entirely, so a caller of it can read any session's
|
||||
* ticket. Every production caller must go through {@link MessageService#poll(String, String)}
|
||||
* and pass a {@code callerTerminal} explicitly, even when it is {@code null}.
|
||||
*
|
||||
* <p>This reads each file's own source text rather than reflecting on compiled bytecode, because
|
||||
* the risk is a future one-word edit at a call site, not a missing overload.
|
||||
*
|
||||
* <p>The scan below finds a violation by its receiver, {@code messages.poll(}, rather than the
|
||||
* bare method name, so it does not mistake {@link java.util.Queue#poll()} for a violation. That
|
||||
* anchor only covers a {@code MessageService} reached through a variable or field named
|
||||
* {@code messages}, so {@link #everyMessageServiceDeclarationIsNamedMessages} pins the naming
|
||||
* convention the anchor depends on: a declaration under any other name would be invisible to the
|
||||
* scan above, and must turn this second check red instead of passing silently.
|
||||
*/
|
||||
class MessageServicePollUsageTest {
|
||||
|
||||
private static final Path PRODUCTION_SOURCE = Path.of("src/main/java");
|
||||
|
||||
@Test
|
||||
void noProductionFileCallsTheSingleArgumentPollOverload() throws IOException {
|
||||
List<String> violations = new ArrayList<>();
|
||||
List<String> twoArgSites = new ArrayList<>();
|
||||
int filesScanned = scanForPollCalls(PRODUCTION_SOURCE, violations, twoArgSites);
|
||||
|
||||
// CONTROL: the scan actually walked files -- a wrong root would otherwise report "no
|
||||
// violations found" having looked at nothing.
|
||||
assertTrue(filesScanned > 0, "control failed: the scan under " + PRODUCTION_SOURCE
|
||||
+ " visited zero .java files -- the path is wrong, so the absence of violations "
|
||||
+ "below proves nothing");
|
||||
|
||||
assertTrue(violations.isEmpty(), "found a call to the fail-open MessageService.poll(String) "
|
||||
+ "overload, which skips the ownership check entirely -- pass a callerTerminal "
|
||||
+ "explicitly (even if null) through poll(String, String) instead: " + violations);
|
||||
|
||||
// CONTROL: the arity parser actually finds the two genuine two-argument call sites (the
|
||||
// MCP handler in FleetMcp and the REST handler in FleetApp). If this drops, the parser
|
||||
// itself is broken, not the production code -- a broken parser (or a scan root that
|
||||
// reaches no real source) must fail loudly here rather than pass vacuously above.
|
||||
assertEquals(2, twoArgSites.size(), "control failed: expected exactly the two known "
|
||||
+ "two-argument messages.poll(...) call sites, found: " + twoArgSites);
|
||||
assertTrue(twoArgSites.stream().anyMatch(s -> s.contains("FleetMcp.java")),
|
||||
"control failed: did not find the FleetMcp.java messages.poll(ticket, callerTerminal) "
|
||||
+ "site among: " + twoArgSites);
|
||||
assertTrue(twoArgSites.stream().anyMatch(s -> s.contains("FleetApp.java")),
|
||||
"control failed: did not find the FleetApp.java messages.poll(...) site among: "
|
||||
+ twoArgSites);
|
||||
}
|
||||
|
||||
/**
|
||||
* The {@code messages.poll(} anchor above only sees a {@code MessageService} reached through
|
||||
* a variable, field, or parameter named {@code messages}. This asserts that every such
|
||||
* declaration under {@code src/main/java} uses that name, so a differently named declaration
|
||||
* -- invisible to the scan above -- fails loudly here instead of letting that scan pass on a
|
||||
* call site it never looked at.
|
||||
*/
|
||||
@Test
|
||||
void everyMessageServiceDeclarationIsNamedMessages() throws IOException {
|
||||
List<String> names = new ArrayList<>();
|
||||
int filesScanned = scanForDeclarationNames(PRODUCTION_SOURCE, names);
|
||||
|
||||
// CONTROL: the scan actually walked files -- a wrong root would otherwise report "every
|
||||
// declaration is named messages" having looked at nothing.
|
||||
assertTrue(filesScanned > 0, "control failed: the scan under " + PRODUCTION_SOURCE
|
||||
+ " visited zero .java files -- the path is wrong, so the result below proves nothing");
|
||||
|
||||
// CONTROL: the declaration pattern actually finds real declarations. Zero means the
|
||||
// pattern is broken, not that every MessageService variable, field, or parameter vanished.
|
||||
assertTrue(names.size() > 0, "control failed: found zero MessageService declarations under "
|
||||
+ PRODUCTION_SOURCE + " -- the declaration pattern is broken, update it before "
|
||||
+ "trusting the naming check below");
|
||||
|
||||
List<String> other = names.stream().filter(n -> !n.equals("messages")).distinct().toList();
|
||||
assertTrue(other.isEmpty(), "found a MessageService declaration not named \"messages\": "
|
||||
+ other + " -- the messages.poll( scan above only looks for that name, so a call "
|
||||
+ "through a differently named variable or field is invisible to it; either rename "
|
||||
+ "the declaration or widen that scan's anchor to cover it");
|
||||
}
|
||||
|
||||
private static int scanForPollCalls(Path root, List<String> violations, List<String> twoArgSites)
|
||||
throws IOException {
|
||||
List<Path> files = javaFiles(root);
|
||||
for (Path file : files) {
|
||||
scanFileForPollCalls(file, violations, twoArgSites);
|
||||
}
|
||||
return files.size();
|
||||
}
|
||||
|
||||
private static int scanForDeclarationNames(Path root, List<String> names) throws IOException {
|
||||
List<Path> files = javaFiles(root);
|
||||
Pattern declaration = Pattern.compile("MessageService\\s+([A-Za-z_][A-Za-z0-9_]*)");
|
||||
for (Path file : files) {
|
||||
if (file.getFileName().toString().equals("MessageService.java")) {
|
||||
continue; // the type's own declaration, not a caller holding a reference to it
|
||||
}
|
||||
String stripped = stripComments(Files.readString(file));
|
||||
Matcher m = declaration.matcher(stripped);
|
||||
while (m.find()) {
|
||||
int j = m.end();
|
||||
while (j < stripped.length() && Character.isWhitespace(stripped.charAt(j))) j++;
|
||||
if (j < stripped.length() && stripped.charAt(j) == '(') {
|
||||
continue; // a method named like the convention, e.g. "MessageService messages()"
|
||||
}
|
||||
names.add(m.group(1));
|
||||
}
|
||||
}
|
||||
return files.size();
|
||||
}
|
||||
|
||||
private static List<Path> javaFiles(Path root) throws IOException {
|
||||
try (Stream<Path> paths = Files.walk(root)) {
|
||||
return paths.filter(p -> p.toString().endsWith(".java")).toList();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Replaces {@code //} and {@code /* *}{@code /} comment text with nothing, leaving code,
|
||||
* string/char literals and line breaks untouched -- so a comment that merely mentions
|
||||
* {@code MessageService} in prose can never be read as a declaration.
|
||||
*/
|
||||
private static String stripComments(String source) {
|
||||
StringBuilder out = new StringBuilder(source.length());
|
||||
boolean inString = false;
|
||||
boolean inChar = false;
|
||||
int i = 0;
|
||||
while (i < source.length()) {
|
||||
char c = source.charAt(i);
|
||||
if (inString) {
|
||||
out.append(c);
|
||||
if (c == '\\' && i + 1 < source.length()) { out.append(source.charAt(i + 1)); i += 2; continue; }
|
||||
if (c == '"') inString = false;
|
||||
i++;
|
||||
continue;
|
||||
}
|
||||
if (inChar) {
|
||||
out.append(c);
|
||||
if (c == '\\' && i + 1 < source.length()) { out.append(source.charAt(i + 1)); i += 2; continue; }
|
||||
if (c == '\'') inChar = false;
|
||||
i++;
|
||||
continue;
|
||||
}
|
||||
if (c == '"') { inString = true; out.append(c); i++; continue; }
|
||||
if (c == '\'') { inChar = true; out.append(c); i++; continue; }
|
||||
if (c == '/' && i + 1 < source.length() && source.charAt(i + 1) == '/') {
|
||||
while (i < source.length() && source.charAt(i) != '\n') i++;
|
||||
continue; // leaves the newline itself for the next iteration to append
|
||||
}
|
||||
if (c == '/' && i + 1 < source.length() && source.charAt(i + 1) == '*') {
|
||||
i += 2;
|
||||
while (i < source.length() && !(source.charAt(i) == '*' && i + 1 < source.length()
|
||||
&& source.charAt(i + 1) == '/')) {
|
||||
if (source.charAt(i) == '\n') out.append('\n');
|
||||
i++;
|
||||
}
|
||||
i += 2;
|
||||
continue;
|
||||
}
|
||||
out.append(c);
|
||||
i++;
|
||||
}
|
||||
return out.toString();
|
||||
}
|
||||
|
||||
private static void scanFileForPollCalls(Path file, List<String> violations, List<String> twoArgSites)
|
||||
throws IOException {
|
||||
String source = Files.readString(file);
|
||||
String needle = "messages.poll(";
|
||||
int from = 0;
|
||||
int idx;
|
||||
while ((idx = source.indexOf(needle, from)) >= 0) {
|
||||
int argsStart = idx + needle.length();
|
||||
String args = extractBalancedArgs(source, argsStart, file, idx);
|
||||
int closeParenIndex = argsStart + args.length();
|
||||
from = closeParenIndex + 1;
|
||||
|
||||
if (args.isBlank()) {
|
||||
continue; // MessageService has no zero-argument poll() -- nothing to classify
|
||||
}
|
||||
String site = file + ":" + lineOf(source, idx);
|
||||
if (topLevelCommaCount(args) == 0) {
|
||||
violations.add(site + " -- messages.poll(" + args.trim() + ")");
|
||||
} else {
|
||||
twoArgSites.add(site);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The text between {@code messages.poll(} and its matching close paren: balanced over nested
|
||||
* calls, and never split by a paren or comma sitting inside a string or char literal.
|
||||
*/
|
||||
private static String extractBalancedArgs(String source, int start, Path file, int callIndex) {
|
||||
int depth = 1;
|
||||
boolean inString = false;
|
||||
boolean inChar = false;
|
||||
int i = start;
|
||||
while (i < source.length()) {
|
||||
char c = source.charAt(i);
|
||||
if (inString) {
|
||||
if (c == '\\') { i += 2; continue; }
|
||||
if (c == '"') inString = false;
|
||||
} else if (inChar) {
|
||||
if (c == '\\') { i += 2; continue; }
|
||||
if (c == '\'') inChar = false;
|
||||
} else if (c == '"') {
|
||||
inString = true;
|
||||
} else if (c == '\'') {
|
||||
inChar = true;
|
||||
} else if (c == '(') {
|
||||
depth++;
|
||||
} else if (c == ')') {
|
||||
depth--;
|
||||
if (depth == 0) return source.substring(start, i);
|
||||
}
|
||||
i++;
|
||||
}
|
||||
throw new IllegalStateException(
|
||||
"unbalanced parentheses scanning " + file + ":" + lineOf(source, callIndex));
|
||||
}
|
||||
|
||||
/**
|
||||
* Commas at paren/bracket/brace depth zero, skipping string and char literals -- the argument
|
||||
* separators a human reader would see, not every comma character in the text.
|
||||
*/
|
||||
private static int topLevelCommaCount(String args) {
|
||||
int depth = 0;
|
||||
int commas = 0;
|
||||
boolean inString = false;
|
||||
boolean inChar = false;
|
||||
int i = 0;
|
||||
while (i < args.length()) {
|
||||
char c = args.charAt(i);
|
||||
if (inString) {
|
||||
if (c == '\\') { i += 2; continue; }
|
||||
if (c == '"') inString = false;
|
||||
} else if (inChar) {
|
||||
if (c == '\\') { i += 2; continue; }
|
||||
if (c == '\'') inChar = false;
|
||||
} else if (c == '"') {
|
||||
inString = true;
|
||||
} else if (c == '\'') {
|
||||
inChar = true;
|
||||
} else if (c == '(' || c == '[' || c == '{') {
|
||||
depth++;
|
||||
} else if (c == ')' || c == ']' || c == '}') {
|
||||
depth--;
|
||||
} else if (c == ',' && depth == 0) {
|
||||
commas++;
|
||||
}
|
||||
i++;
|
||||
}
|
||||
return commas;
|
||||
}
|
||||
|
||||
private static int lineOf(String source, int index) {
|
||||
int line = 1;
|
||||
for (int i = 0; i < index; i++) {
|
||||
if (source.charAt(i) == '\n') line++;
|
||||
}
|
||||
return line;
|
||||
}
|
||||
}
|
||||
@@ -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");
|
||||
|
||||
@@ -139,6 +139,154 @@ class FleetAppAuthTest {
|
||||
assertThrows(IllegalArgumentException.class, () -> FleetApp.routeAction("GET /healthz"));
|
||||
}
|
||||
|
||||
// --- GET /tasks/{ticket} must not be the no-check overload ----------------------------------
|
||||
|
||||
/**
|
||||
* {@code GET /tasks/{ticket}} must resolve its caller the same way {@code allow(...)} does
|
||||
* and thread that terminal into {@link MessageService#poll(String, String)}, not the
|
||||
* no-check overload that ignores who is asking.
|
||||
*/
|
||||
@Test
|
||||
void theTaskStatusRouteActuallyThreadsTheCallersTerminalIntoPoll() throws Exception {
|
||||
String source = Files.readString(REST_SOURCE);
|
||||
|
||||
int start = source.indexOf("private void taskStatus(Context ctx) {");
|
||||
assertTrue(start >= 0, "could not find taskStatus in " + REST_SOURCE
|
||||
+ " -- the scrape has stopped matching, fix the anchor before trusting this test");
|
||||
int end = source.indexOf("private static void herdrError(Context ctx, HerdrException e) {", start);
|
||||
assertTrue(end > start, "could not find the method declared after taskStatus to bound the scrape");
|
||||
String handlerBlock = source.substring(start, end);
|
||||
|
||||
// CONTROL: the block we scraped really does call messages.poll(...) -- if this fails, the
|
||||
// anchors above moved and the assertions below would otherwise pass on nothing.
|
||||
assertTrue(handlerBlock.contains("messages.poll("),
|
||||
"control failed: the scraped taskStatus block contains no messages.poll( call at all "
|
||||
+ "-- the anchors have drifted, this test is not testing what it claims to");
|
||||
|
||||
assertTrue(handlerBlock.contains("caller.terminal()"),
|
||||
"the taskStatus route must thread the resolved caller's terminal into messages.poll(...), "
|
||||
+ "not the no-check overload -- block: " + handlerBlock);
|
||||
assertTrue(handlerBlock.contains("ctx.attribute(CALLER)"),
|
||||
"the taskStatus route must resolve its caller the same way allow(...) does, not via a "
|
||||
+ "second, separate resolution path -- block: " + handlerBlock);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code GET /tasks/{ticket}} refuses a worker whose terminal did not create the ticket,
|
||||
* while the creating worker and the unnamed primary both still read it. The ticket is minted
|
||||
* directly on the shared {@link MessageService}, the same way {@code MessageServiceTest}
|
||||
* drives {@link MessageService#poll(String, String)}, so this exercises only the REST poll
|
||||
* route's own handling of the ownership already recorded on the ticket.
|
||||
*/
|
||||
@Test
|
||||
void restPollRefusesADifferentWorkerButAllowsTheCreatorAndTheUnnamedPrimary() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
Injector injector = new Injector(agents);
|
||||
MessageService messages = new MessageService(agents, injector, new Rendezvous());
|
||||
|
||||
Javalin creatorApp = startOnSharedService(messages, herdr, FakeHerdr.WORKER_PID); // -> term_a
|
||||
Javalin otherWorkerApp = startOnSharedService(messages, herdr, 9001L); // -> term_shell
|
||||
Javalin primaryApp = startOnSharedService(messages, herdr, 999_999L); // no pane -> primary
|
||||
try {
|
||||
String ticket = messages.sendAsync("term_a", "long task", null, "term_a");
|
||||
|
||||
HttpResponse<String> refused = send(otherWorkerApp.port(), "GET", "/tasks/" + ticket, null, null);
|
||||
assertEquals(200, refused.statusCode());
|
||||
assertTrue(refused.body().contains("forbidden"),
|
||||
"a different worker's terminal must be refused, not shown the ticket: " + refused.body());
|
||||
assertFalse(refused.body().contains("\"reply\""),
|
||||
"a refusal must never carry reply text: " + refused.body());
|
||||
|
||||
HttpResponse<String> own = send(creatorApp.port(), "GET", "/tasks/" + ticket, null, null);
|
||||
assertEquals(200, own.statusCode());
|
||||
assertFalse(own.body().contains("forbidden"),
|
||||
"the creating worker must read its own ticket: " + own.body());
|
||||
|
||||
HttpResponse<String> primary = send(primaryApp.port(), "GET", "/tasks/" + ticket, null, null);
|
||||
assertEquals(200, primary.statusCode());
|
||||
assertFalse(primary.body().contains("forbidden"),
|
||||
"the unnamed primary must read any ticket: " + primary.body());
|
||||
} finally {
|
||||
creatorApp.stop();
|
||||
otherWorkerApp.stop();
|
||||
primaryApp.stop();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code POST /sessions/{id}/message} with {@code wait:false} must record the creating
|
||||
* caller's own terminal on the ticket it returns, so that caller can still poll its own
|
||||
* ticket over REST, while a different terminal is refused.
|
||||
*/
|
||||
@Test
|
||||
void restSendAsyncRecordsTheCreatingCallersTerminalSoItCanStillPollItsOwnTicket() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
Injector injector = new Injector(agents);
|
||||
MessageService messages = new MessageService(agents, injector, new Rendezvous());
|
||||
|
||||
Javalin leadApp = startOnSharedService(messages, herdr, FakeHerdr.WORKER_PID, Map.of("term_a", "lead-x"));
|
||||
Javalin otherWorkerApp = startOnSharedService(messages, herdr, 9001L); // -> term_shell
|
||||
try {
|
||||
ObjectMapper mapper = new ObjectMapper();
|
||||
HttpResponse<String> created = send(leadApp.port(), "POST", "/sessions/term_a/message",
|
||||
"{\"content\":\"long task\",\"wait\":false}", null);
|
||||
assertEquals(202, created.statusCode(), created.body());
|
||||
String ticket = mapper.readTree(created.body()).path("ticket").asText(null);
|
||||
assertNotNull(ticket, "the accepted response carried no ticket: " + created.body());
|
||||
|
||||
HttpResponse<String> own = send(leadApp.port(), "GET", "/tasks/" + ticket, null, null);
|
||||
assertEquals(200, own.statusCode());
|
||||
assertFalse(own.body().contains("forbidden"),
|
||||
"the session that created the ticket over REST must be able to poll it: " + own.body());
|
||||
|
||||
HttpResponse<String> refused = send(otherWorkerApp.port(), "GET", "/tasks/" + ticket, null, null);
|
||||
assertEquals(200, refused.statusCode());
|
||||
assertTrue(refused.body().contains("forbidden: this ticket was created by a different session"),
|
||||
"a different terminal must still be refused with the ownership detail, not some "
|
||||
+ "other rejection: " + refused.body());
|
||||
} finally {
|
||||
leadApp.stop();
|
||||
otherWorkerApp.stop();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #start}, but shares {@code messages} and {@code herdr} across several app
|
||||
* instances bound to different pids, each returned as its own started {@link Javalin} rather
|
||||
* than through the shared {@code app} field, so several differently-resolved callers can
|
||||
* poll the same ticket.
|
||||
*/
|
||||
private Javalin startOnSharedService(MessageService messages, FakeHerdr herdr, long pid) {
|
||||
return startOnSharedService(messages, herdr, pid, Map.of());
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #startOnSharedService(MessageService, FakeHerdr, long)}, but {@code leadTerminals}
|
||||
* resolves the given pid's terminal to a named lead (a caller with SEND permission) instead of
|
||||
* a plain worker, for a test that needs a terminal-bearing caller able to create a ticket.
|
||||
*/
|
||||
private Javalin startOnSharedService(MessageService messages, FakeHerdr herdr, long pid,
|
||||
Map<String, String> leadTerminals) {
|
||||
FleetConfig.Profile wcfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
|
||||
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(
|
||||
agents, new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
Map.of(wcfg.profile(), wcfg), wcfg.profile(),
|
||||
k -> "FLEETD_WORKER_TOKEN".equals(k) ? "tok-abc" : null);
|
||||
SessionManager sessions = new SessionManager(workers, new FakeWorktrees());
|
||||
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> pid);
|
||||
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null,
|
||||
() -> leadTerminals, new MemberRegistry(null));
|
||||
Metrics appMetrics = FleetMetrics.create(sessions, new dev.ltms.fleet.msg.InMemoryReplyInbox());
|
||||
|
||||
return new FleetApp(herdr, workers, sessions, messages, sessions.asPresence(), null,
|
||||
callers, appMetrics).build().start("127.0.0.1", 0);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #669 Unit A: {@code POST /sessions/{id}/message} is two call shapes behind one route,
|
||||
* mirroring {@code fleet_send}'s MCP-side split into {@link Authz.Action#SEND} and {@link
|
||||
|
||||
Reference in New Issue
Block a user