Compare commits

..

16 Commits

Author SHA1 Message Date
Dai Ha 0f5985b419 fleetd #721: gate fleet_status's pending-ask block by the delegation's creator terminal
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Failing after 1m55s
fleet_status (MCP) and GET /sessions/{id}/status (REST) handed any TASK_READ
holder another session's open fleet_ask question, its turnId and its ticket,
with no check that the caller created that delegation. MessageService.pendingAsk
now takes the caller's terminal and reuses the existing ownsTicket comparison;
FleetMcp.status and FleetApp.sessionStatus both thread the resolved caller
terminal through. The base status line and REST's ready field are unaffected.
2026-10-04 08:46:38 +02:00
Dai Ha 6e06058b07 Merge PR #720: fleetd #718 — pin that no production caller uses MessageService.poll(String)
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 50s
CI / build (push) Failing after 1m44s
2026-10-04 08:23:58 +02:00
Dai Ha e5f4fb81ab fleetd #718: pin the naming convention the poll-usage scan depends on
CI / shell-tests (pull_request) Failing after 6s
CI / contract (pull_request) Successful in 1m15s
CI / build (pull_request) Failing after 2m22s
MessageServicePollUsageTest's messages.poll( receiver anchor only
covers a MessageService reached through a variable, field, or
parameter named "messages". Adds a second assertion in the same
class that every such declaration under src/main/java uses that
name, with its own file-walk and declaration-count controls, so a
future declaration under a different name turns this check red
instead of leaving the original scan silently blind to it.

Also tightens the Authz.permits(Principal, Action, String) javadoc
sentence to read as a plain contract statement.
2026-10-04 08:18:54 +02:00
Dai Ha d28ab0968b fleetd #718: pin that no production caller uses MessageService.poll(String)
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 52s
CI / build (pull_request) Failing after 1m47s
Adds a source-scrape test over src/main/java that fails if any caller
reaches the fail-open single-argument poll(String) overload instead of
poll(String, String). The scan anchors on the "messages.poll(" receiver
to avoid matching java.util.Queue.poll(), and balances parentheses to
avoid being fooled by a two-argument call whose first argument contains
nested parens.

Also documents Authz.permits(Principal, Action, String) as a test
convenience whose default classifier denies every collaborator.
2026-10-04 08:09:33 +02:00
Dai Ha 9a64d42599 Merge PR #716: fleetd #705 — gate the REST ticket routes by the creating caller
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 51s
CI / build (push) Failing after 1m57s
Three parts. GET /tasks/{ticket} passed no caller, so it used the overload that skips the ownership
check and any worker could read any ticket; it now passes the caller resolved from the same CALLER
attribute the authorization gate reads. The wait:false send path recorded no creator terminal, so a
REST-created ticket matched no terminal-bearing caller and its own creator was refused; it now
records one. Both handlers gained a scrape guard with its own control assertion.

Resolved one conflict in FleetMcpAuthzTest by keeping both sides: PR #717 and this branch each
appended tests at the same point. The test count is the check on that resolution — 2014 + 3 + 1 =
2018, so no test was dropped.

Verified in a throwaway worktree off main: Tests run: 2018, Failures: 0, BUILD SUCCESS. Three
mutations, each confirmed live with mvn -o compile before the suite ran. FleetMcp:527 is the one
that survived before this work and now kills theFleetPollHandlerActuallyThreadsCallerTerminalIntoPoll.
FleetApp:698 kills the new creator test. FleetApp:899 kills three, including the behavioural test.
Every file restored byte-identical.
2026-10-04 07:48:18 +02:00
Dai Ha f4e0ca41e6 Merge PR #717: fleetd #710 — gate fleet_list's leads and members arrays by caller role
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 55s
CI / build (push) Failing after 1m52s
A worker now gets neither array, and the key is absent rather than present-and-empty. leads stays
visible to the primary, an architect and a collaborator; members only to the primary and an
architect. That split follows the rule the collaborators array already states: you may list what
you could address. A collaborator may send to a lead, and leads is the only place the bridge gives
it that address, so hiding it would have left a shipped grant unusable.

Verified in a throwaway worktree off main: Tests run: 2014, Failures: 0, BUILD SUCCESS, against a
main baseline of 2008 that I measured myself. Dropping the collaborator clause from leadsVisibleTo
compiled green and then killed exactly two tests, the truth-table row and the behavioural test.
2026-10-04 07:41:21 +02:00
Dai Ha 29a2f97c25 fleetd: thread the creating caller's terminal into the REST fire-and-poll send path
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 53s
CI / build (pull_request) Failing after 1m53s
sendMessage's wait:false branch now records the resolved caller's own
terminal as the ticket's creatorTerminal, the same way taskStatus already
resolves its caller, so a REST-created ticket's own creator can still poll
it under the ownership check that now gates GET /tasks/{ticket}.
2026-10-04 07:41:12 +02:00
Dai Ha d2db8c7dc9 fleetd #710 PR #717 review: leadsVisibleTo must also admit a collaborator
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Failing after 1m47s
A collaborator may SEND to a lead, and fleet_list's leads array is the
only place this tool gives it a lead's sessionId -- its own
fleet_whoami carries no lead address. leadsVisibleTo now returns true
for caller.isCollaborator() as well as primary and architect, matching
the existing rule for the sibling collaborators array (every role that
may SEND to a named peer). membersVisibleTo is unchanged: a
collaborator may never SEND to a spawned member.

Updated the truth-table tests for both predicates, added a behavioural
test proving a collaborator's fleet_list output contains leads and not
members, and updated fleet_list's tool description.
2026-10-04 07:37:55 +02:00
Dai Ha 95311c6e8e Correct two stale claims in the canonical bridge block
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 1m5s
CI / build (push) Failing after 2m4s
The block is the instruction surface this repo ships, so a merged change that makes it false is an
incomplete change. Two had gone stale:

- the lead table said fleet_list does not report collaborators. #709 made it report them, visible
  to the primary, an architect and another collaborator, never a worker.
- the collaborator section justified the ticket refusal by saying ticket ids are a plain counter
  with no owner check. #712 added that owner check, so the stated reason no longer held. The role
  gate is what refuses a collaborator; the recorded creator terminal is the second line.

wiki/7-Use-Cases.md carries the same edit and was pushed to its own remote, verified by ref. The
sync check prints True.
2026-10-04 07:34:23 +02:00
Dai Ha 10ab58e4fc fleetd #710: gate fleet_list's leads and members arrays by caller role
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 1m50s
Omit the leads and members keys entirely (never an empty array) from
fleet_list's result for a worker, matching the existing
coordinatorVisibleTo/collaboratorsVisibleTo pattern: two new named
predicates (leadsVisibleTo, membersVisibleTo) are consulted before
assembling either array, so a worker holding only READ can no longer
read every session on the daemon through this tool. Primary and
architect callers are unaffected.
2026-10-04 07:25:55 +02:00
Dai Ha 0ba597e394 fleetd #705: close the REST ticket-poll door and pin both handlers' caller-terminal wiring
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 1m47s
GET /tasks/{ticket} now resolves the caller the same way allow(...) does and
threads that terminal into MessageService.poll(ticket, callerTerminal) instead
of the no-check overload, so a worker can no longer read a ticket a different
session created over REST. Adds a source-scrape guard (with its own control
assertion) for both the fleet_poll MCP handler and this REST route, plus a
behavioural test driving GET /tasks/{ticket} with three differently-resolved
callers against one shared MessageService.
2026-10-04 07:23:29 +02:00
Dai Ha c468953963 Merge PR #714: fleetd #711 — re-key two startup refusals onto the hazard that is still real
CI / shell-tests (push) Failing after 8s
CI / build (push) Failing after 1m53s
CI / contract (push) Successful in 50s
The roster is consulted ahead of every tab map, so a live registered member
is no longer read back as a lead. Both refusals still matter, for the narrower
case where the pane is alive and the roster holds no entry for it.

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

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

Verified in a throwaway worktree, not piped: Tests run: 2008, BUILD SUCCESS.
2026-10-04 07:00:35 +02:00
Dai Ha 7d497aa423 fleetd #711: re-key two pane-tab hazard reasons onto the unregistered-pane condition
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 1m39s
Both refusals said a member landing in a lead's or collaborator's
labelled tab would be read back as that identity. CallerResolver
consults the spawned-member roster ahead of every tab map, so a live
registered member is never misread this way. State the real condition
instead: the hazard applies only while the pane is alive and carries
no entry in the spawned-member roster.
2026-10-04 06:54:43 +02:00
Dai Ha 11eccc3a1b fleetd #705: gate fleet_poll's ticket lookup by the creating caller's terminal
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 1m45s
A ticket id is a plain sequential counter, so any session holding
TASK_READ could walk task-1, task-2, ... and read another session's
delegation reply. sendAsync now records the creating caller's terminal
on the Task, and poll refuses a caller whose terminal differs from it.
A caller with no terminal (the unnamed primary) is still allowed
through regardless, since it never carries a herdr pane to compare.
2026-10-04 06:50:26 +02:00
13 changed files with 1171 additions and 89 deletions
+3 -3
View File
@@ -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.
@@ -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);
@@ -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
@@ -514,7 +514,7 @@ public final class FleetMcp {
(exchange, req) -> {
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_status", req.arguments()), null);
if (denied != null) return denied;
return status(messages, str(req.arguments(), "sessionId"));
return status(messages, str(req.arguments(), "sessionId"), callerTerminal(exchange));
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> pollHandler =
(exchange, req) -> {
@@ -524,7 +524,7 @@ public final class FleetMcp {
// The action depends on the ARGUMENTS, not on the tool name -- see pollAction.
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_poll", a), target);
if (denied != null) return denied;
return poll(messages, leadChannel, str(a, "ticket"), target, coordId);
return poll(messages, leadChannel, str(a, "ticket"), target, coordId, callerTerminal(exchange));
};
// CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox).
// Acking removes a reply from the inbox, so it is a drain, not a read.
@@ -562,7 +562,8 @@ public final class FleetMcp {
callerTerminal(exchange),
callers.collaborators(), collaboratorsVisibleTo(principal(exchange)),
new CoordinationSource(leadChannel, peers),
coordinatorVisibleTo(principal(exchange)));
coordinatorVisibleTo(principal(exchange)),
leadsVisibleTo(principal(exchange)), membersVisibleTo(principal(exchange)));
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
(exchange, req) -> {
@@ -759,6 +760,33 @@ public final class FleetMcp {
return caller.isPrimary() || caller.isArchitect() || caller.isCollaborator();
}
/**
* Who may see {@code fleet_list}'s {@code leads} array — exactly the roles that may
* {@link Authz.Action#SEND} to a lead: the primary, an architect, and a collaborator. A
* collaborator's own {@code fleet_whoami} carries no lead address, and {@code leads} is the
* only place this tool gives one, so a collaborator needs this array to use the send it
* already holds. A worker can never {@code SEND} at all, so it still sees neither this array
* nor {@code members}; a worker's own facts come from {@code fleet_whoami} instead. Split
* out for the same reason as {@link #coordinatorVisibleTo} and
* {@link #collaboratorsVisibleTo}: the decision must be unit-testable without fabricating an
* SDK {@code McpSyncServerExchange}, and the handler must call this named predicate rather
* than inlining the check.
*/
static boolean leadsVisibleTo(Principal caller) {
return caller.isPrimary() || caller.isArchitect() || caller.isCollaborator();
}
/**
* Who may see {@code fleet_list}'s {@code members} array — the primary and an architect,
* which may {@link Authz.Action#SEND} to a member. A worker can never {@code SEND} at all,
* and a collaborator may {@code SEND} only to a lead or another collaborator, never to a
* spawned member, so neither sees this array even though {@link #leadsVisibleTo} grants a
* collaborator the sibling one.
*/
static boolean membersVisibleTo(Principal caller) {
return caller.isPrimary() || caller.isArchitect();
}
/** The worker identity resolved from this call's connection, or {@code null} if the primary. */
private static String callerTerminal(McpSyncServerExchange exchange) {
Object v = exchange.transportContext().get(CALLER_TERMINAL);
@@ -972,6 +1000,17 @@ public final class FleetMcp {
*/
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
Runnable onAccepted, Set<String> profiles) {
return sendAsync(messages, sessionId, content, onAccepted, profiles, null);
}
/**
* As above, recording {@code creatorTerminal} as this ticket's owner (the caller's own
* terminal, resolved from the connection) so a later {@code fleet_poll{ticket}} only hands the
* result back to that same caller — see {@link MessageService#poll(String, String)}.
*/
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
Runnable onAccepted, Set<String> profiles,
String creatorTerminal) {
if (isBlank(sessionId) || isBlank(content)) {
return error("sessionId and content are required");
}
@@ -979,7 +1018,7 @@ public final class FleetMcp {
if (targetError != null) {
return targetError;
}
String ticket = messages.sendAsync(sessionId, content, onAccepted);
String ticket = messages.sendAsync(sessionId, content, onAccepted, creatorTerminal);
return text("accepted — task delegated. Poll fleet_poll with ticket=" + ticket);
}
@@ -1158,9 +1197,24 @@ public final class FleetMcp {
* exclusively of {@code ticket}/{@code target} — see {@link #pollAction}'s javadoc for why this
* is a different inbox (this daemon's own {@link LeadChannel}) that authorizes differently
* ({@link Authz.Action#COORD_READ}, primary-only) from either of the original two branches.
*
* <p>Does not check who owns {@code ticket} — see the overload that takes {@code
* callerTerminal} for that. Callers that do not resolve a caller terminal (tests, or a surface
* with no connection-based identity) use this one.
*/
static McpSchema.CallToolResult poll(MessageService messages, LeadChannel leadChannel, String ticket,
String target, String coordId) {
return poll(messages, leadChannel, ticket, target, coordId, null);
}
/**
* As above, refusing a ticket lookup whose caller's terminal differs from the terminal that
* created it — see {@link MessageService#poll(String, String)}. {@code callerTerminal} is the
* CALLING session's terminal id, resolved by the MCP layer from the connection, never a
* client-supplied value.
*/
static McpSchema.CallToolResult poll(MessageService messages, LeadChannel leadChannel, String ticket,
String target, String coordId, String callerTerminal) {
if (!isBlank(coordId)) {
return pollHeldPeerMail(leadChannel, coordId);
}
@@ -1174,7 +1228,7 @@ public final class FleetMcp {
if (isBlank(ticket)) {
return error("ticket (or target) is required");
}
MessageService.TaskView v = messages.poll(ticket);
MessageService.TaskView v = messages.poll(ticket, callerTerminal);
if (v == null) {
return error("unknown ticket: " + ticket + " (never issued, or expired)");
}
@@ -1281,16 +1335,20 @@ public final class FleetMcp {
/**
* {@code fleet_status}: the live lifecycle status of a worker session, plus — when the worker
* is paused mid-turn in an async {@code fleet_ask} (CB-582) — the open question and how to
* answer it, so a lead on its normal poll cadence does not need the ticket to notice.
* is paused mid-turn in an async {@code fleet_ask} — the open question and how to answer it, so
* a lead on its normal poll cadence does not need the ticket to notice. The question, its
* {@code turnId} and its ticket id are shown only to the caller whose terminal created that
* delegation, or to a caller with no terminal at all (the unnamed primary); any other caller
* still sees the base status. {@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 status(MessageService messages, String sessionId) {
static McpSchema.CallToolResult status(MessageService messages, String sessionId, String callerTerminal) {
if (isBlank(sessionId)) {
return error("sessionId is required");
}
try {
String base = messages.status(sessionId).name().toLowerCase();
MessageService.PendingAsk ask = messages.pendingAsk(sessionId);
MessageService.PendingAsk ask = messages.pendingAsk(sessionId, callerTerminal);
if (ask == null) {
return text(base);
}
@@ -1778,7 +1836,7 @@ public final class FleetMcp {
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine,
OutageSource.none(), LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, false);
Map.of(), false, coordination, false, true, true);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
@@ -1788,7 +1846,7 @@ public final class FleetMcp {
Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, CoordinationSource.none(), false);
Map.of(), false, CoordinationSource.none(), false, true, true);
}
/**
@@ -1806,7 +1864,7 @@ public final class FleetMcp {
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, OutageSource.none(),
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, false);
Map.of(), false, coordination, false, true, true);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
@@ -1816,7 +1874,7 @@ public final class FleetMcp {
Map<String, String> leads, String selfTerm, CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, false);
Map.of(), false, coordination, false, true, true);
}
/**
@@ -1839,7 +1897,7 @@ public final class FleetMcp {
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, false);
Map.of(), false, coordination, false, true, true);
}
/**
@@ -1856,15 +1914,23 @@ public final class FleetMcp {
* {@code false} (fleetd #463: a forgotten argument fails closed, not
* open), so a test that wants the {@code coordinator} row must pass
* an explicit {@code true}
* @param leadsVisible whether this caller may see the {@code leads} array (see
* {@link #leadsVisibleTo}) — unlike {@code callerIsPrimary}, this is
* also {@code true} for an architect or a collaborator, so it cannot
* be derived from {@code callerIsPrimary} alone
* @param membersVisible whether this caller may see the {@code members} array (see
* {@link #membersVisibleTo}); also {@code true} for an architect, but
* unlike {@code leadsVisible}, never for a collaborator
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, OutageSource outage,
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
CoordinationSource coordination, boolean callerIsPrimary) {
CoordinationSource coordination, boolean callerIsPrimary,
boolean leadsVisible, boolean membersVisible) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine,
outage, leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, callerIsPrimary);
Map.of(), false, coordination, callerIsPrimary, leadsVisible, membersVisible);
}
/**
@@ -1881,6 +1947,10 @@ public final class FleetMcp {
* {@link #collaboratorsVisibleTo}); every wrapper overload above
* passes {@code false}, so a test that wants the row must call this
* overload with an explicit {@code true}
* @param leadsVisible whether this caller may see the {@code leads} array (see
* {@link #leadsVisibleTo})
* @param membersVisible whether this caller may see the {@code members} array (see
* {@link #membersVisibleTo})
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
@@ -1890,28 +1960,41 @@ public final class FleetMcp {
LeadConfigDirSource leadConfigDirs,
Map<String, String> leads, String selfTerm,
Map<String, String> collaborators, boolean collaboratorsVisible,
CoordinationSource coordination, boolean callerIsPrimary) {
CoordinationSource coordination, boolean callerIsPrimary,
boolean leadsVisible, boolean membersVisible) {
try {
Map<String, Agent> live = workers.list().stream()
.map(Agent.class::cast)
.filter(a -> a.terminalId() != null)
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
List<Map<String, Object>> leadRows = leads.entrySet().stream()
.sorted(Map.Entry.comparingByValue())
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm, contextGauge,
leadConfigDirs))
.toList();
// Neither row's assembly (leadView/memberCapacityView probing herdr for live status)
// runs unless at least one of them needs the live-agent lookup backing it.
Map<String, Agent> live = (leadsVisible || membersVisible)
? workers.list().stream()
.map(Agent.class::cast)
.filter(a -> a.terminalId() != null)
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b))
: Map.of();
// fleetd #209: this is the caller-driven fleet_list read that actually reports
// agentSessionId (via memberCapacityView -> SessionManager.rosterView), so it uses the
// resolving roster; the heartbeat/health/metrics timers stay on the plain sessions.roster().
List<MemberSession> roster = sessions.rosterResolved();
List<Map<String, Object>> out = roster.stream()
.map(s -> memberCapacityView(s, live.get(s.terminalId()), messages, capacity.clock().getAsLong()))
.toList();
Set<String> profiles = new java.util.TreeSet<>(capacity.configuredProfiles().get());
roster.stream().map(MemberSession::profile).forEach(profiles::add);
Map<String, Object> result = new LinkedHashMap<>();
result.put("leads", leadRows); result.put("members", out);
// READ is permission to enter this tool, not permission to receive every field it can
// build -- gate BEFORE assembling each row, so the key is absent rather than
// present-and-empty; a caller without either row gets its own facts from fleet_whoami.
if (leadsVisible) {
List<Map<String, Object>> leadRows = leads.entrySet().stream()
.sorted(Map.Entry.comparingByValue())
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm, contextGauge,
leadConfigDirs))
.toList();
result.put("leads", leadRows);
}
if (membersVisible) {
List<Map<String, Object>> out = roster.stream()
.map(s -> memberCapacityView(s, live.get(s.terminalId()), messages, capacity.clock().getAsLong()))
.toList();
result.put("members", out);
}
result.put("healthCoverage", healthCoverage.value().get());
result.put("loopHealth", Map.of(
"statusPoller", loopHealth.statusPoller().get().name(),
@@ -2445,7 +2528,14 @@ public final class FleetMcp {
private static McpSchema.Tool listTool() {
return tool(FleetTool.LIST.wireName(),
"List the whole fleet the bridge tracks, in two parts. 'leads' are your PEERS — other "
"List the whole fleet the bridge tracks, in two parts. 'members' is visible to "
+ "the primary and an architect only. 'leads' is visible to those two AND a "
+ "collaborator — exactly the roles that may fleet_send to a lead, so a "
+ "collaborator can learn a lead's sessionId before using the send it already "
+ "holds. A worker holds READ to call this tool at all, but gets neither "
+ "array, never an empty one; a worker reads its own session, "
+ "profile, state, worktree, branch and owner from fleet_whoami instead. "
+ "'leads' are your PEERS — other "
+ "orchestrators, each with its sessionId (the address to fleet_send to), "
+ "name, live status, and 'self': true on your own row; this is how you "
+ "discover a peer lead without being told its address. 'members' are the "
@@ -275,10 +275,19 @@ public final class MessageService {
* {@link #abandon}) can never match again regardless of this flag's value.
*/
private volatile boolean askTimedOut;
/**
* The terminal of the caller whose {@code fleet_send{wait:false}} created this ticket, or
* {@code null} when that caller had no terminal (the unnamed primary) or the ticket was
* created through an overload that does not record one. {@link #poll(String, String)}
* compares a polling caller's own terminal against this field before handing back the
* ticket's state.
*/
private final String creatorTerminal;
private Task(String ticket, String target, LongSupplier nowNanos) {
private Task(String ticket, String target, LongSupplier nowNanos, String creatorTerminal) {
this.ticket = ticket;
this.target = target;
this.creatorTerminal = creatorTerminal;
this.createdNanos = nowNanos.getAsLong();
future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong());
}
@@ -1279,8 +1288,20 @@ public final class MessageService {
* @return the ticket to poll for the eventual result
*/
public String sendAsync(String target, String content, Runnable onAccepted) {
return sendAsync(target, content, onAccepted, null);
}
/**
* As {@link #sendAsync(String, String, Runnable)}, recording {@code creatorTerminal} as this
* ticket's owner. {@link #poll(String, String)} refuses a later caller whose own terminal
* differs from this one; {@code null} records no owner (a caller with no terminal — the
* unnamed primary — is always allowed to poll the result regardless).
*
* @return the ticket to poll for the eventual result
*/
public String sendAsync(String target, String content, Runnable onAccepted, String creatorTerminal) {
String ticket = "task-" + ticketSeq.incrementAndGet();
Task task = new Task(ticket, target, nowNanos);
Task task = new Task(ticket, target, nowNanos, creatorTerminal);
tasks.put(ticket, task);
if (pushLoop != null) {
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a
@@ -1343,15 +1364,31 @@ public final class MessageService {
}
/**
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket;
* otherwise a {@link Phase#PENDING} view (with the live worker status as detail), a
* {@link Phase#DONE} view carrying the reply, or a {@link Phase#FAILED} view with the reason.
* As {@link #poll(String, String)}, with no caller terminal — the ticket's ownership is never
* checked, so this overload must only be used where the caller's identity is otherwise
* irrelevant (a test, or a surface that does not resolve a caller terminal at all).
*/
public TaskView poll(String ticket) {
return poll(ticket, null);
}
/**
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket.
* Refuses a {@code callerTerminal} that differs from the terminal that created the ticket (see
* {@link #sendAsync(String, String, Runnable, String)}) with a {@link Phase#FAILED} view that
* carries no reply text — a caller with no terminal (the unnamed primary) is never refused.
* Otherwise returns a {@link Phase#PENDING} view (with the live worker status as detail), a
* {@link Phase#DONE} view carrying the reply, or a {@link Phase#FAILED} view with the reason.
*/
public TaskView poll(String ticket, String callerTerminal) {
Task task = tasks.get(ticket);
if (task == null) {
return null;
}
if (!ownsTicket(task, callerTerminal)) {
return new TaskView(ticket, Phase.FAILED, null, null,
"forbidden: this ticket was created by a different session", null);
}
CompletableFuture<Reply> f = task.future;
if (!f.isDone()) {
Reply question = task.question;
@@ -1387,6 +1424,17 @@ public final class MessageService {
return new TaskView(ticket, Phase.FAILED, null, null, detail, null);
}
/**
* Whether {@code callerTerminal} may read {@code task}'s state. A caller with no terminal
* always may — that is the unnamed primary, resolved by token or loopback trust, which never
* carries a herdr pane and must keep reading every ticket. Otherwise the caller's terminal must
* equal the terminal recorded on the task; a task with no recorded terminal matches no
* terminal-bearing caller.
*/
private static boolean ownsTicket(Task task, String callerTerminal) {
return callerTerminal == null || callerTerminal.equals(task.creatorTerminal);
}
/**
* Test seam only — carries no production behaviour, and nothing in this class calls it;
* {@link #pruneTerminalTickets} still reads {@link Task#completedNanos} directly.
@@ -1706,16 +1754,17 @@ public final class MessageService {
}
/**
* The question {@code workerSession} is currently paused on via {@code fleet_ask}, if any
* (CB-582) — {@code fleet_status} uses this to show a pending question without the caller
* needing the ticket. {@code null} when the session has no open async question (including a
* session mid a <em>blocking</em> {@code fleet_ask}, which has no {@link Task} to look up — see
* {@link PendingAsk}).
* The question {@code workerSession} is currently paused on via {@code fleet_ask}, if any —
* {@code fleet_status} uses this to show a pending question without the caller needing the
* ticket. {@code null} when the session has no open async question (including a session mid a
* <em>blocking</em> {@code fleet_ask}, which has no {@link Task} to look up — see
* {@link PendingAsk}), or when {@code callerTerminal} does not own the task the question
* belongs to (see {@link #ownsTicket(Task, String)}).
*/
public PendingAsk pendingAsk(String workerSession) {
public PendingAsk pendingAsk(String workerSession, String callerTerminal) {
for (Task task : tasks.values()) {
Reply q = task.question;
if (q != null && workerSession.equals(task.target)) {
if (q != null && workerSession.equals(task.target) && ownsTicket(task, callerTerminal)) {
return new PendingAsk(task.ticket, q.text(), q.turnId());
}
}
@@ -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;
}
@@ -874,10 +875,12 @@ public final class FleetApp {
body.put("sessionId", id);
body.put("status", messages.status(id).name().toLowerCase());
body.put("ready", deliverable.test(id));
// CB-582: a worker paused mid-turn in an async fleet_ask is otherwise invisible to a
// status poll — surface the open question and how to answer it, same as fleet_poll's
// Phase.ASKING view.
MessageService.PendingAsk ask = messages.pendingAsk(id);
// A worker paused mid-turn in an async fleet_ask is otherwise invisible to a status
// poll — surface the open question and how to answer it, same as fleet_poll's
// Phase.ASKING view, but only to the caller whose terminal created that delegation, or
// to a caller with no terminal at all (the unnamed primary).
Principal caller = ctx.attribute(CALLER);
MessageService.PendingAsk ask = messages.pendingAsk(id, caller == null ? null : caller.terminal());
if (ask != null) {
body.put("question", ask.question());
body.put("turnId", ask.turnId());
@@ -894,7 +897,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;
@@ -850,7 +850,8 @@ class FleetConfigTest {
/**
* The hazard this guard closes: a pane-placed member lands inside the focused tab rather than
* its own, so it can land inside a lead's labelled tab and be read back as that lead.
* its own, so it can land inside a lead's labelled tab and, while its pane carries no entry in
* the spawned-member roster, be read back as that lead.
*/
@Test
void aPanePlacedProfileWithALeadTabRefusesToStart(@TempDir Path dir) throws Exception {
@@ -1087,8 +1088,9 @@ class FleetConfigTest {
}
/**
* fleetd #669: a member tabLabel that can render as a configured collaborator tab is the same
* hazard as the lead case above — a member labelled that way is read back as the collaborator.
* A member tabLabel that can render as a configured collaborator tab is the same hazard as the
* lead case above — while its pane carries no entry in the spawned-member roster, a member
* labelled that way is read back as the collaborator.
*/
@Test
void aProfileTabLabelOverrideMatchingACollaboratorTabRefusesToStart(@TempDir Path dir)
@@ -355,10 +355,10 @@ class FleetMcpAuthzTest {
+ "the anchors have drifted, this test is not testing what it claims to");
Pattern trailingArg = Pattern.compile(
"listFleet\\([^;]*?,\\s*(coordinatorVisibleTo\\(principal\\(exchange\\)\\)|true|false)\\s*\\)\\s*;",
"listFleet\\([^;]*?,\\s*(coordinatorVisibleTo\\(principal\\(exchange\\)\\)|true|false)\\s*[,)]",
Pattern.DOTALL);
Matcher m = trailingArg.matcher(handlerBlock);
assertTrue(m.find(), "could not locate listFleet(...)'s trailing boolean argument in the "
assertTrue(m.find(), "could not locate listFleet(...)'s coordinator boolean argument in the "
+ "listHandler block -- the call shape changed, update this test's anchor: " + handlerBlock);
String trailing = m.group(1);
assertEquals("coordinatorVisibleTo(principal(exchange))", trailing,
@@ -418,6 +418,145 @@ class FleetMcpAuthzTest {
+ "calling, not pass a literal boolean -- block: " + handlerBlock);
}
// --- who may see fleet_list's leads and members arrays ---------------------------------------
/**
* {@link FleetMcp#leadsVisibleTo} is the whole policy decision for {@code fleet_list}'s
* {@code leads} array: visible to exactly the roles that may {@code SEND} to a lead -- the
* primary, an architect, and a collaborator -- never a worker, which holds {@code READ} but
* can never {@code SEND} at all, and never an anonymous caller.
*/
@Test
void primaryArchitectAndCollaboratorMaySeeTheLeadsArray() {
assertTrue(FleetMcp.leadsVisibleTo(PRIMARY), "the primary must see the leads array");
assertTrue(FleetMcp.leadsVisibleTo(ARCH_DESIGN), "an architect must see the leads array");
assertTrue(FleetMcp.leadsVisibleTo(COLLABORATOR),
"a collaborator may SEND to a lead, so it must see the leads array to learn where");
assertFalse(FleetMcp.leadsVisibleTo(WORKER_A),
"a worker holds READ but can never SEND, so it must not see the leads array");
assertFalse(FleetMcp.leadsVisibleTo(ANON), "authenticated as nothing must not see it either");
}
/**
* As {@link #primaryArchitectAndCollaboratorMaySeeTheLeadsArray}, for the {@code members}
* array -- but a collaborator may {@code SEND} only to a lead or another collaborator, never
* to a spawned member, so it must not see this one.
*/
@Test
void onlyPrimaryAndArchitectMaySeeTheMembersArray() {
assertTrue(FleetMcp.membersVisibleTo(PRIMARY), "the primary must see the members array");
assertTrue(FleetMcp.membersVisibleTo(ARCH_DESIGN), "an architect must see the members array");
assertFalse(FleetMcp.membersVisibleTo(WORKER_A),
"a worker holds READ but can never SEND, so it must not see the members array");
assertFalse(FleetMcp.membersVisibleTo(COLLABORATOR),
"a collaborator may SEND to a lead, never to a spawned member, so it must not see the members array");
assertFalse(FleetMcp.membersVisibleTo(ANON), "authenticated as nothing must not see it either");
}
/**
* Same reasoning as {@link #theFleetListHandlerActuallyConsultsCoordinatorVisibleTo}: the
* predicate above can be perfectly correct while the one production call site never asks it.
* This reads {@code FleetMcp.java}'s own source and asserts the {@code fleet_list} handler's
* {@code listFleet(...)} call asks {@code leadsVisibleTo(principal(exchange))} for the leads
* visibility flag, rather than a literal boolean.
*/
@Test
void theFleetListHandlerActuallyConsultsLeadsVisibleTo() throws Exception {
String source = Files.readString(MCP_SOURCE);
int start = source.indexOf("listHandler =");
assertTrue(start >= 0, "could not find the fleet_list handler (listHandler) in " + MCP_SOURCE
+ " -- the scrape has stopped matching, fix the anchor before trusting this test");
int end = source.indexOf("stopHandler =", start);
assertTrue(end > start, "could not find the handler declared after listHandler to bound the scrape");
String handlerBlock = source.substring(start, end);
// CONTROL: the block we scraped really does contain a call to listFleet(...) -- if this
// fails, the anchors above moved and the assertion below would otherwise pass on nothing.
assertTrue(handlerBlock.contains("listFleet("),
"control failed: the scraped listHandler block contains no listFleet( call at all -- "
+ "the anchors have drifted, this test is not testing what it claims to");
assertTrue(handlerBlock.contains("leadsVisibleTo(principal(exchange))"),
"the fleet_list handler must ask leadsVisibleTo(principal(exchange)) who is calling, "
+ "not pass a literal boolean -- block: " + handlerBlock);
}
/** As {@link #theFleetListHandlerActuallyConsultsLeadsVisibleTo}, for {@code membersVisibleTo}. */
@Test
void theFleetListHandlerActuallyConsultsMembersVisibleTo() throws Exception {
String source = Files.readString(MCP_SOURCE);
int start = source.indexOf("listHandler =");
assertTrue(start >= 0, "could not find the fleet_list handler (listHandler) in " + MCP_SOURCE
+ " -- the scrape has stopped matching, fix the anchor before trusting this test");
int end = source.indexOf("stopHandler =", start);
assertTrue(end > start, "could not find the handler declared after listHandler to bound the scrape");
String handlerBlock = source.substring(start, end);
assertTrue(handlerBlock.contains("listFleet("),
"control failed: the scraped listHandler block contains no listFleet( call at all -- "
+ "the anchors have drifted, this test is not testing what it claims to");
assertTrue(handlerBlock.contains("membersVisibleTo(principal(exchange))"),
"the fleet_list handler must ask membersVisibleTo(principal(exchange)) who is calling, "
+ "not pass a literal boolean -- block: " + handlerBlock);
}
/**
* {@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);
}
/**
* {@code fleet_status} must thread the calling connection's own terminal into
* {@link FleetMcp#status(MessageService, String, String)}, so a caller that did not create a
* worker's open delegation cannot read its pending question through the status handler either.
*/
@Test
void theFleetStatusHandlerActuallyThreadsCallerTerminalIntoStatus() throws Exception {
String source = Files.readString(MCP_SOURCE);
int start = source.indexOf("statusHandler =");
assertTrue(start >= 0, "could not find the fleet_status handler (statusHandler) in " + MCP_SOURCE
+ " -- the scrape has stopped matching, fix the anchor before trusting this test");
int end = source.indexOf("pollHandler =", start);
assertTrue(end > start, "could not find the handler declared after statusHandler to bound the scrape");
String handlerBlock = source.substring(start, end);
// CONTROL: the block we scraped really does call status(...) -- if this fails, the anchors
// above moved and the assertions below would otherwise pass on nothing.
assertTrue(handlerBlock.contains("status(messages,"),
"control failed: the scraped statusHandler block contains no status(messages, ...) "
+ "call at all -- the anchors have drifted, this test is not testing what it claims to");
assertTrue(handlerBlock.contains("callerTerminal(exchange)"),
"the fleet_status handler must thread callerTerminal(exchange) into status(...), not "
+ "omit it or pass a literal null -- block: " + handlerBlock);
}
// --- which action each tool hands the gate (fleetd #272) ------------------------------------
/**
@@ -122,6 +122,6 @@ class FleetMcpLeadContextGaugeWiringTest {
FleetMcp.LoopHealthSource.none(), FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
FleetMcp.LeadSeatSource.none(), contextGauge, leadConfigDirs,
Map.of(LEAD_TERMINAL, LEAD_NAME), LEAD_TERMINAL, Map.of(), false,
FleetMcp.CoordinationSource.none(), false);
FleetMcp.CoordinationSource.none(), false, true, true);
}
}
@@ -692,7 +692,7 @@ class FleetMcpTest {
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
String out = textOf(res);
// fleetd #361: reports both which coord-id a peer must use to reach ME, and this daemon's
@@ -721,7 +721,7 @@ class FleetMcpTest {
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
String out = textOf(res);
assertTrue(out.contains("\"mailbox\":{\"status\":\"unknown\"}"), out);
@@ -780,19 +780,82 @@ class FleetMcpTest {
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
.withMailbox("mac-opus", LeadChannel.MailboxState.exists("mac-opus", 0, 1));
boolean callerIsPrimary = Principal.worker("term_a", 1).isPrimary();
Principal worker = Principal.worker("term_a", 1);
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), callerIsPrimary);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), worker.isPrimary(),
FleetMcp.leadsVisibleTo(worker), FleetMcp.membersVisibleTo(worker));
String out = textOf(res);
assertFalse(out.contains("\"coordinator\""), "a worker must never see the coordinator key at all: " + out);
assertFalse(out.contains("mac-opus"), "no fragment of the coordinator row may leak either: " + out);
assertTrue(out.contains("\"leads\""), "the rest of the result must still be present: " + out);
assertTrue(out.contains("\"members\""), out);
assertFalse(out.contains("\"leads\""), "a worker must never see the leads key at all: " + out);
assertFalse(out.contains("\"members\""), "a worker must never see the members key at all: " + out);
assertTrue(out.contains("\"healthCoverage\""), "the rest of the result must still be present: " + out);
}
/**
* A worker's result must carry neither the {@code leads} nor the {@code members} key, and no
* fragment of either row leaks even though both are fully populated for this call -- a key
* check alone would pass on an implementation that still built the rows and only renamed or
* nested them.
*/
@Test
void listLeaksNoLeadOrMemberRowFragmentToAWorker() {
FakeHerdr h = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", "term_owner",
new WorktreeRequest("cb-304", null));
Principal worker = Principal.worker("term_a", 1);
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of("term_lead_x", "mac-opus"), "", FleetMcp.CoordinationSource.none(), worker.isPrimary(),
FleetMcp.leadsVisibleTo(worker), FleetMcp.membersVisibleTo(worker));
String out = textOf(res);
assertFalse(out.contains("\"leads\""), out);
assertFalse(out.contains("\"members\""), out);
assertFalse(out.contains(s.terminalId()), "no member row fragment may leak to a worker: " + out);
assertFalse(out.contains("term_owner"), "no member owner fragment may leak to a worker: " + out);
assertFalse(out.contains("mac-opus"), "no lead row fragment may leak to a worker: " + out);
assertTrue(out.contains("\"healthCoverage\""), out);
}
/**
* A collaborator may {@code SEND} to a lead, so it must see the {@code leads} array -- the
* only place {@code fleet_whoami} does not already give it a lead's address. It may never
* {@code SEND} to a spawned member, so the {@code members} key must stay absent for it, with
* no fragment of a populated member row leaking either.
*/
@Test
void listShowsLeadsButNotMembersToACollaborator() {
FakeHerdr h = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", "term_owner",
new WorktreeRequest("cb-304", null));
Principal collaborator = Principal.collaborator("ops", "term_collab", 600);
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of("term_lead_x", "mac-opus"), "", FleetMcp.CoordinationSource.none(), collaborator.isPrimary(),
FleetMcp.leadsVisibleTo(collaborator), FleetMcp.membersVisibleTo(collaborator));
String out = textOf(res);
assertTrue(out.contains("\"leads\""), "a collaborator must see the leads array: " + out);
assertTrue(out.contains("mac-opus"), "a collaborator must see the lead's name/address: " + out);
assertFalse(out.contains("\"members\""), "a collaborator must never see the members key: " + out);
assertFalse(out.contains(s.terminalId()), "no member row fragment may leak to a collaborator: " + out);
assertFalse(out.contains("term_owner"), "no member owner fragment may leak to a collaborator: " + out);
assertTrue(out.contains("\"healthCoverage\""), out);
}
@@ -808,16 +871,19 @@ class FleetMcpTest {
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
.withMailbox("mac-opus", LeadChannel.MailboxState.exists("mac-opus", 0, 1));
boolean callerIsPrimary = Principal.architect("lead-designer", "term_design", 400).isPrimary();
Principal architect = Principal.architect("lead-designer", "term_design", 400);
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), callerIsPrimary);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), architect.isPrimary(),
FleetMcp.leadsVisibleTo(architect), FleetMcp.membersVisibleTo(architect));
String out = textOf(res);
assertFalse(out.contains("\"coordinator\""), "an architect must never see the coordinator key either: " + out);
assertTrue(out.contains("\"leads\""), "an architect must still see the leads array: " + out);
assertTrue(out.contains("\"members\""), "an architect must still see the members array: " + out);
}
/**
@@ -842,7 +908,7 @@ class FleetMcpTest {
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true));
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true));
assertTrue(gatedAsPrimary.contains("\"coordinator\""), gatedAsPrimary);
assertTrue(gatedAsPrimary.contains("\"selfId\":\"mac-opus\""), gatedAsPrimary);
@@ -851,6 +917,8 @@ class FleetMcpTest {
assertTrue(gatedAsPrimary.contains("\"heldDurable\""), gatedAsPrimary);
assertTrue(gatedAsPrimary.contains("\"held\""), gatedAsPrimary);
assertTrue(gatedAsPrimary.contains("\"peers\""), gatedAsPrimary);
assertTrue(gatedAsPrimary.contains("\"leads\""), "the primary must still see the leads array: " + gatedAsPrimary);
assertTrue(gatedAsPrimary.contains("\"members\""), "the primary must still see the members array: " + gatedAsPrimary);
}
@Test
@@ -869,7 +937,7 @@ class FleetMcpTest {
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
String out = textOf(res);
assertTrue(out.contains("\"msgId\":\"m1\""), out);
@@ -893,7 +961,7 @@ class FleetMcpTest {
FleetMcp.LoopHealthSource.none(), FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
FleetMcp.LeadSeatSource.none(), new LeadContextGauge(), FleetMcp.LeadConfigDirSource.none(),
Map.of(), "", collaborators, collaboratorsVisible,
FleetMcp.CoordinationSource.none(), false);
FleetMcp.CoordinationSource.none(), false, true, true);
}
/**
@@ -978,7 +1046,7 @@ class FleetMcpTest {
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
String out = textOf(res);
assertTrue(out.contains("\"pending\":0"), out);
@@ -1007,7 +1075,7 @@ class FleetMcpTest {
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
String out = textOf(res);
assertTrue(out.contains("\"heldDurable\":false"),
@@ -1076,7 +1144,7 @@ class FleetMcpTest {
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "",
new FleetMcp.CoordinationSource(channel, List.of("fleet01-lead", "fleet02-lead", "fleet03-lead")), true);
new FleetMcp.CoordinationSource(channel, List.of("fleet01-lead", "fleet02-lead", "fleet03-lead")), true, true, true);
String out = textOf(res);
assertTrue(out.contains("\"coordId\":\"fleet01-lead\",\"status\":\"exists\",\"pending\":2,\"consumers\":1"), out);
@@ -1792,14 +1860,14 @@ class FleetMcpTest {
FakeHerdr blocked = new FakeHerdr().agentStatus("blocked");
AgentControl blockedAgents = new AgentControl(blocked);
McpSchema.CallToolResult res = FleetMcp.status(
new MessageService(blockedAgents, new Injector(blockedAgents), rendezvous), "term_a");
new MessageService(blockedAgents, new Injector(blockedAgents), rendezvous), "term_a", null);
assertNotEquals(Boolean.TRUE, res.isError());
assertEquals("blocked", textOf(res));
}
/**
* CB-582: a lead polling {@code fleet_status} on its normal cadence — not {@code fleet_poll}
* — must also see a worker's open async {@code fleet_ask} question, since the reverse-rendezvous
* A lead polling {@code fleet_status} on its normal cadence — not {@code fleet_poll} — must
* also see a worker's open async {@code fleet_ask} question, since the reverse-rendezvous
* window it opened with is far shorter than that cadence.
*/
@Test
@@ -1822,7 +1890,7 @@ class FleetMcpTest {
} while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline);
assertEquals(MessageService.Phase.ASKING, asking.phase());
McpSchema.CallToolResult res = FleetMcp.status(messages, T);
McpSchema.CallToolResult res = FleetMcp.status(messages, T, null);
assertNotEquals(Boolean.TRUE, res.isError());
String out = textOf(res);
assertTrue(out.startsWith("idle"), "the live status must still lead the text: " + out);
@@ -1843,6 +1911,63 @@ class FleetMcpTest {
answer.get(5, TimeUnit.SECONDS);
}
/**
* {@code fleet_status}'s pending-ask block (the question, its {@code turnId} and its ticket)
* is shown only to the caller whose terminal created the delegation, or to a caller with no
* terminal at all (the unnamed primary) — a different terminal-bearing caller still sees the
* base status line, but none of the pending-ask fields.
*/
@Test
void statusGatesThePendingAskFieldsByTheDelegationsCreatorTerminal() throws Exception {
String ticket = messages.sendAsync(T, "task that asks", null, "term_creator");
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting(T), "sendAsync should have opened its rendezvous waiter");
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.TaskView asking;
deadline = System.currentTimeMillis() + 3000;
do {
asking = messages.poll(ticket);
Thread.sleep(5);
} while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline);
assertEquals(MessageService.Phase.ASKING, asking.phase());
String other = textOf(FleetMcp.status(messages, T, "term_other"));
assertTrue(other.startsWith("idle"), "the base status must still be shown: " + other);
assertFalse(other.contains("which config file?"),
"a non-creating caller must not see the question text: " + other);
assertFalse(other.contains(asking.turnId()),
"a non-creating caller must not see the turnId: " + other);
assertFalse(other.contains(ticket),
"a non-creating caller must not see the ticket: " + other);
String creator = textOf(FleetMcp.status(messages, T, "term_creator"));
assertTrue(creator.contains("which config file?"), "the creator must see the question: " + creator);
assertTrue(creator.contains(asking.turnId()), "the creator must see the turnId: " + creator);
assertTrue(creator.contains(ticket), "the creator must see the ticket: " + creator);
String unnamed = textOf(FleetMcp.status(messages, T, null));
assertTrue(unnamed.contains("which config file?"),
"a caller with no terminal (the unnamed primary) must see the question: " + unnamed);
// Clean up the still-open ask so the background thread does not linger past the test.
String turnId = asking.turnId();
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(turnId, "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.resolve(T, "done"));
answer.get(5, TimeUnit.SECONDS);
}
// --- fleet_whoami: the caller's own identity, so an agent never has to guess its role -------
@Test
@@ -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");
@@ -1804,20 +1875,20 @@ class MessageServiceTest {
}
}
// --- CB-582: fleet_status pendingAsk() ------------------------------------------------------
// --- fleet_status pendingAsk() -----------------------------------------------------------
@Test
void pendingAskReturnsNullWhenNoQuestionIsOpen() throws Exception {
assertNull(messages.pendingAsk(T), "no async ticket at all -> no pending ask");
assertNull(messages.pendingAsk(T, null), "no async ticket at all -> no pending ask");
String ticket = messages.sendAsync(T, "long task");
awaitWaiting();
assertNull(messages.pendingAsk(T), "a plain pending delegation is not a question");
assertNull(messages.pendingAsk(T, null), "a plain pending delegation is not a question");
injectDelivery();
assertTrue(rendezvous.resolve(T, "done"));
awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertNull(messages.pendingAsk(T), "a finished ticket carries no open question either");
assertNull(messages.pendingAsk(T, null), "a finished ticket carries no open question either");
}
@Test
@@ -1830,7 +1901,7 @@ class MessageServiceTest {
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
MessageService.PendingAsk pending = messages.pendingAsk(T);
MessageService.PendingAsk pending = messages.pendingAsk(T, null);
assertNotNull(pending, "fleet_status should see the open question");
assertEquals(ticket, pending.ticket());
assertEquals("which config file?", pending.question());
@@ -1839,7 +1910,69 @@ class MessageServiceTest {
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
assertNull(messages.pendingAsk(T), "an answered question is no longer pending");
assertNull(messages.pendingAsk(T, null), "an answered question is no longer pending");
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
}
/**
* A caller's own terminal must match the terminal that created the delegation to see its
* pending question; a different terminal-bearing caller sees nothing, and a caller with no
* terminal at all (the unnamed primary) always sees it.
*/
@Test
void pendingAskGatesTheQuestionByTheDelegationsCreatorTerminal() throws Exception {
String ticket = messages.sendAsync(T, "task that asks", null, "term_creator");
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
assertNull(messages.pendingAsk(T, "term_other"),
"a caller whose terminal did not create the delegation must not see the question");
MessageService.PendingAsk own = messages.pendingAsk(T, "term_creator");
assertNotNull(own, "the creating caller must see its own open question");
assertEquals("which config file?", own.question());
MessageService.PendingAsk unnamed = messages.pendingAsk(T, null);
assertNotNull(unnamed, "a caller with no terminal (the unnamed primary) must always see the question");
assertEquals("which config file?", unnamed.question());
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
}
/**
* A task created with no recorded creator terminal (a short {@code sendAsync} overload) must
* not hand its open question to any caller that does have a terminal — only a caller with no
* terminal at all may still see it.
*/
@Test
void pendingAskDeniesATerminalBearingCallerWhenTheTaskRecordsNoCreator() throws Exception {
String ticket = messages.sendAsync(T, "task that asks"); // no creatorTerminal recorded
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
assertNull(messages.pendingAsk(T, "term_someone"),
"a terminal-bearing caller must not see a question whose task records no creator");
assertNotNull(messages.pendingAsk(T, null),
"the unnamed primary must still see it even with no recorded creator");
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
@@ -2,6 +2,7 @@ package dev.ltms.fleet.rest;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.spi.ILoggingEvent;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import dev.ltms.fleet.auth.CallerResolver;
import dev.ltms.fleet.auth.Authz;
@@ -39,6 +40,8 @@ import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
import java.util.function.Predicate;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
@@ -139,6 +142,257 @@ 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 /sessions/{id}/status} must resolve its caller the same way {@code allow(...)}
* does and thread that terminal into {@link MessageService#pendingAsk(String, String)}, not
* the no-check overload that ignores who is asking.
*/
@Test
void theSessionStatusRouteActuallyThreadsTheCallersTerminalIntoPendingAsk() throws Exception {
String source = Files.readString(REST_SOURCE);
int start = source.indexOf("private void sessionStatus(Context ctx) {");
assertTrue(start >= 0, "could not find sessionStatus in " + REST_SOURCE
+ " -- the scrape has stopped matching, fix the anchor before trusting this test");
int end = source.indexOf("private void taskStatus(Context ctx) {", start);
assertTrue(end > start, "could not find the method declared after sessionStatus to bound the scrape");
String handlerBlock = source.substring(start, end);
// CONTROL: the block we scraped really does call messages.pendingAsk(...) -- if this fails,
// the anchors above moved and the assertions below would otherwise pass on nothing.
assertTrue(handlerBlock.contains("messages.pendingAsk("),
"control failed: the scraped sessionStatus block contains no messages.pendingAsk( "
+ "call at all -- the anchors have drifted, this test is not testing what it claims to");
assertTrue(handlerBlock.contains("caller.terminal()"),
"the sessionStatus route must thread the resolved caller's terminal into "
+ "messages.pendingAsk(...), not the no-check overload -- block: " + handlerBlock);
assertTrue(handlerBlock.contains("ctx.attribute(CALLER)"),
"the sessionStatus 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();
}
}
/**
* {@code GET /sessions/{id}/status} shows a worker's pending {@code fleet_ask} question, its
* {@code turnId} and its ticket only to the caller whose terminal created that delegation, or
* to a caller with no terminal at all (the unnamed primary) — a different terminal-bearing
* caller still sees the base status line, but none of the pending-ask fields.
*/
@Test
void restStatusGatesThePendingAskFieldsByTheDelegationsCreatorTerminal() throws Exception {
FakeHerdr herdr = new FakeHerdr();
AgentControl agents = new AgentControl(herdr);
Injector injector = new Injector(agents);
Rendezvous rendezvous = new Rendezvous();
MessageService messages = new MessageService(agents, injector, 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 {
ObjectMapper mapper = new ObjectMapper();
String ticket = messages.sendAsync("term_target", "task that asks", null, "term_a");
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_target") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting("term_target"), "sendAsync should have opened its rendezvous waiter");
CompletableFuture<MessageService.AskResult> ask = CompletableFuture.supplyAsync(
() -> messages.ask("term_target", "which config file?", 5000));
MessageService.TaskView asking;
deadline = System.currentTimeMillis() + 3000;
do {
asking = messages.poll(ticket, null);
Thread.sleep(5);
} while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline);
assertEquals(MessageService.Phase.ASKING, asking.phase());
String turnId = asking.turnId();
JsonNode other = mapper.readTree(
send(otherWorkerApp.port(), "GET", "/sessions/term_target/status", null, null).body());
assertEquals("idle", other.get("status").asText(), "the base status must still be shown");
assertFalse(other.has("question"), "a non-creating caller must not see the question: " + other);
assertFalse(other.has("turnId"), "a non-creating caller must not see the turnId: " + other);
assertFalse(other.has("ticket"), "a non-creating caller must not see the ticket: " + other);
JsonNode own = mapper.readTree(
send(creatorApp.port(), "GET", "/sessions/term_target/status", null, null).body());
assertEquals("which config file?", own.get("question").asText(), "the creator must see the question");
assertEquals(turnId, own.get("turnId").asText(), "the creator must see the turnId");
assertEquals(ticket, own.get("ticket").asText(), "the creator must see the ticket");
JsonNode primary = mapper.readTree(
send(primaryApp.port(), "GET", "/sessions/term_target/status", null, null).body());
assertEquals("which config file?", primary.get("question").asText(),
"a caller with no terminal (the unnamed primary) must see the question");
// Clean up the still-open ask so the background thread does not linger past the test.
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(turnId, "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_target") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.resolve("term_target", "done"));
answer.get(5, TimeUnit.SECONDS);
} finally {
creatorApp.stop();
otherWorkerApp.stop();
primaryApp.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