Compare commits

..

23 Commits

Author SHA1 Message Date
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 b205bcc2aa Remove CallerResolver members accessor
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 1m47s
2026-10-04 06:53:39 +02:00
Dai Ha 11eccc3a1b fleetd #705: gate fleet_poll's ticket lookup by the creating caller's terminal
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 1m45s
A ticket id is a plain sequential counter, so any session holding
TASK_READ could walk task-1, task-2, ... and read another session's
delegation reply. sendAsync now records the creating caller's terminal
on the Task, and poll refuses a caller whose terminal differs from it.
A caller with no terminal (the unnamed primary) is still allowed
through regardless, since it never carries a herdr pane to compare.
2026-10-04 06:50:26 +02:00
Dai Ha 4939d40362 Merge PR #709: fleetd #703 — fleet_list reports collaborators, narrowed by role
CI / shell-tests (push) Failing after 12s
CI / contract (push) Successful in 48s
CI / build (push) Failing after 1m53s
A collaborator could find and message a lead, but a lead had nowhere to read a
collaborator's sessionId, so the channel only worked once the collaborator
spoke first. fleet_list now carries a collaborators array, and
CallerResolver.collaborators() has its first caller outside the resolver.

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

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

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

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

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

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

mvn clean install: exit 0, BUILD SUCCESS, Tests run: 2000, Failures: 0.
2026-10-04 06:35:49 +02:00
Dai Ha 3e8e314656 Merge PR #706 and PR #707: fleetd #669 — a collaborator is deliverable, and a collaborator config change is reported
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 1m0s
CI / build (push) Failing after 1m46s
Two follow-ups to #669, verified together in one throwaway worktree.

PR #707 (f820737): ConfigRef reports that a fleet.collaborators change needs a
restart. The map is read once at startup and a reload does not rebuild it, so
before this an operator saw a clean reload and no live effect.

PR #706 (20fc42b): deliverableTo opens for a collaborator. A collaborator is
never in MemberPresence and never found by the lead scan, so every send to one
was authorized, then held on the injector gate for ~60s and dropped without a
keystroke reaching the pane.

mvn clean install on the merged tree: exit 0, BUILD SUCCESS, Tests run: 1999,
Failures: 0, Errors: 0. ConfigRefTest 32, FleetDeliverabilityTest 9.
2026-10-04 06:19:47 +02:00
Dai Ha 20fc42b572 fleetd #669 follow-up: deliverableTo also opens for a collaborator
CI / shell-tests (pull_request) Failing after 11s
CI / contract (pull_request) Successful in 59s
CI / build (pull_request) Failing after 1m48s
A collaborator's terminal was authorized at Authz but never deliverable: it
is never enrolled in MemberPresence and never discovered by the lead scan,
so a send to a collaborator sat on the injector's readiness gate for
~60s and failed, never typed into the pane. deliverableTo now takes a
third collaborators supplier, read through on each call like the lead
supplier, and FleetdAssembly wires the existing collaboratorTerminals
supplier into it.
2026-10-04 06:13:03 +02:00
Dai Ha f82073717a fleetd #669: report collaborator reload restart
CI / shell-tests (pull_request) Failing after 11s
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Failing after 2m1s
2026-10-04 06:12:37 +02:00
Dai Ha 16b52fac6c fleetd.example.yaml: the collaborator block works now (#669)
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 52s
CI / build (push) Failing after 1m49s
The comment above the collaborators: example said the block 'is parsed and
validated today; nothing yet recognises or addresses the tab it names'. That
was true when Unit B landed the config shape. Units C, D and E have since
merged and deployed, so the tab is recognised, the role is authorized, and
the pane routes to the lead herdr daemon.

This is the worst place for that sentence to go stale: it sits in the file an
operator copies to turn the feature on, and it tells them the block does
nothing. Replaced with what the role may and may not do, taken from the rows
in auth/Authz.java rather than from the wiki, plus the #703 limit that a lead
cannot discover a collaborator.

FleetConfigTest (which loads this file) green at 170 tests.
2026-10-04 02:30:28 +02:00
Dai Ha 13b6ae2628 Merge PR #704: fleetd #669 Unit E — route a collaborator's pane to the lead herdr daemon
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 58s
CI / build (push) Failing after 2m5s
A configured collaborator's terminal now routes to the lead herdr daemon
rather than the member one. FleetdAssembly ORs a second terminal map into
the predicate it hands HerdrRouter, and the predicate's field is renamed
isLead -> routeToLead so its name matches the widened contract.

Verified in the lead's own throwaway worktree: mvn clean install green at
1992 tests, 0 failures (174 surefire report files, fresh worktree). The
fix rests on one test, so I reproduced the mutation myself rather than
taking the worker's word: dropping the collaborator clause from the
predicate turns FleetdAssemblyCollaboratorHerdrRoutingTest red at runtime
with two distinct AgentControl identities, while HerdrRouterTest's 4 tests
stay green. That confirms both the kill and that the two-distinct-fakes
setup is actually discriminating.

Also checked by hand: both terminal maps come from LeadTabScanner's one
byKind helper and are keyed terminal_id -> name, so containsKey(target) is
correct for both; no agentsFor call runs between router construction and
the ref being set, so the publication window is harmless and matches the
one leadsRef already had.

One reviewer fanned out against the diff, briefed from the diff rather
than the implementer's rationale. It reported no issue.
2026-10-04 02:19:44 +02:00
23 changed files with 1300 additions and 119 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.
+22 -13
View File
@@ -62,11 +62,12 @@ bind:
#
# Leads are configured under `fleet.leaders:` — see THE FLEET further down.
#
# Two things stop the tab-name convention from becoming a way to claim leadership: the configured
# member spaces are excluded from the scan, so nothing fleetd places can land in a matching tab;
# and startup REFUSES a `tabPrefix` that the fleet tabLabel template, or any per-profile `tabLabel`
# override, also matches — so the two namespaces cannot overlap by accident. The label is a NAME,
# never a capability: what a pane may do is decided by the role the daemon resolves for it.
# Two things stop the tab-name convention from becoming a way to claim leadership: startup REFUSES
# a `tabPrefix` that the fleet tabLabel template, or any per-profile `tabLabel` override, also
# matches, so the two namespaces cannot overlap by accident; and the CallerResolver asks the live
# spawned-member roster BEFORE any tab map, so a live member is never mistaken for a lead no matter
# what its tab says. The label is a NAME, never a capability: what a pane may do is decided by the
# role the daemon resolves for it.
# CB-551: IDLE-LEAD HEARTBEAT — nudge the single lead back to work when it has been continuously
# idle (no open fleet_send driving it) past the quiet period. The fleet is one lead + architects +
@@ -659,9 +660,10 @@ fleet:
# tabPrefix: "lead:" # only used to guard against a worker tabLabel colliding with
# # this convention at startup; plays no part in matching a lead
# scanIntervalSeconds: 10 # rescan cadence, and the worst case before a new tab is seen
# workspace: leads # where a launched lead's tab is created (default "leads").
# # MUST NOT be a member workspace — those are excluded from the
# # scan, so a lead placed in one is never found again.
# workspace: leads # where a launched lead's tab is created (default "fleet",
# # the same shared space the members use). Sharing that space
# # with members is the normal shipped shape: the scanner tells
# # a lead from a member by the exact tab label, not by workspace.
# cwd: /path/to/repo # the launched lead's working directory (default: fleetd's own)
# kind: claude # descriptive; reported by fleet_whoami
# gpt-sol-5.6:
@@ -669,11 +671,18 @@ fleet:
# kind: opencode
# model: openai/gpt-5.6-terra
# A collaborator tab, keyed by name (fleetd #669). This block is parsed and validated today;
# nothing yet recognises or addresses the tab it names. Recognise-only, like a profile-less
# `leaders:` entry above: there is no `profile:`, no `instances:` and no `kind:`. `tab:` is
# REQUIRED and is the only field identity depends on, matched case-insensitively — the same
# GET-THE-VALUE-RIGHT warning above the `leaders:` block applies here too.
# A collaborator tab, keyed by name (fleetd #669). A pane whose tab matches resolves as the
# COLLABORATOR role: `fleet_whoami` answers `collaborator`, and the session may observe the fleet
# (`fleet_list`, `fleet_profiles`, `fleet_whoami`) and `fleet_send` to a lead or to another
# collaborator. It may NOT spawn, stop or drain anything, roll a lead's session, answer a
# member's `fleet_ask`, poll a ticket, send across hosts, or send to a spawned member's terminal.
# Each of those is refused at the gate rather than queued.
# Recognise-only, like a profile-less `leaders:` entry above: there is no `profile:`, no
# `instances:` and no `kind:`. `tab:` is REQUIRED and is the only field identity depends on,
# matched case-insensitively — the same GET-THE-VALUE-RIGHT warning above the `leaders:` block
# applies here too.
# A lead cannot discover a collaborator yet (#703): `fleet_list` has no `collaborators` key, so
# the collaborator must speak first, or pass on the `sessionId` its own `fleet_whoami` reports.
# collaborators:
# reviewer-alex:
# tab: "collab: alex"
+15 -12
View File
@@ -217,25 +217,28 @@ public final class Fleetd {
/**
* The {@link Injector}'s readiness gate (CB-534): a target is deliverable if it is a spawned
* member whose agent has connected the bridge MCP, <em>or</em> a lead.
* member whose agent has connected the bridge MCP, a lead, <em>or</em> a collaborator.
*
* <p>The gate exists for one reason — to hold a delivery out of a <em>spawned</em> member's boot
* window, where herdr already reports {@code idle} but the TUI would drop an injected paste. That
* hazard is a property of spawning. A lead is never spawned: the operator started it and named it
* (or labelled its tab) only once it was up, so there is no boot window to guard.
* (or labelled its tab) only once it was up, so there is no boot window to guard. A collaborator
* is the same way — a person's own tab, matched to a configured name, never spawned.
*
* <p>A lead is also never enrolled in {@link MemberPresence} — {@code FleetMcp} marks presence
* for every spawned member (worker and architect), deliberately, since that map doubles as the
* member roster's availability signal and a lead counted there would show up as an available
* member. So without the second disjunct a lead is permanently un-deliverable: every
* lead→lead send sat on the gate for {@code READINESS_GRACE_POLLS} (~60s) and then failed
* having never been typed into the pane.
* <p>Neither a lead nor a collaborator is ever enrolled in {@link MemberPresence} — {@code
* FleetMcp} marks presence for every spawned member (worker and architect), deliberately, since
* that map doubles as the member roster's availability signal and a lead or collaborator counted
* there would show up as an available member. So without the second and third disjuncts a lead or
* collaborator is permanently un-deliverable: every send to one sat on the gate for
* {@code READINESS_GRACE_POLLS} (~60s) and then failed having never been typed into the pane.
*
* <p>The lead set is read through the supplier on each call rather than snapshotted, so a lead
* discovered by {@code leadScan} after startup becomes deliverable without a restart.
* <p>Both sets are read through their supplier on each call rather than snapshotted, so a lead or
* collaborator discovered by {@code leadScan} after startup becomes deliverable without a restart.
*/
static Predicate<String> deliverableTo(MemberPresence presence, Supplier<Map<String, String>> leads) {
return target -> presence.isPresent(target) || leads.get().containsKey(target);
static Predicate<String> deliverableTo(MemberPresence presence, Supplier<Map<String, String>> leads,
Supplier<Map<String, String>> collaborators) {
return target -> presence.isPresent(target) || leads.get().containsKey(target)
|| collaborators.get().containsKey(target);
}
/**
@@ -368,7 +368,7 @@ final class FleetdAssembly {
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
MemberPresence presence = sessions.asPresence();
TurnListener turnListener = Fleetd.turnListener(completion, sessions);
Predicate<String> deliverable = Fleetd.deliverableTo(presence, leads);
Predicate<String> deliverable = Fleetd.deliverableTo(presence, leads, collaboratorTerminals);
// fleetd #556: registration is wired directly to `completion`, not folded into the
// `turnListener` fan-out above — so it survives `sessions.onDelivered` (or any future
// listener) throwing, regardless of call order.
@@ -75,6 +75,9 @@ public final class Authz {
* {@code SEND} is refused, as if no terminal were a configured lead or collaborator — the
* same decision {@link #NO_KNOWN_LEAD_OR_COLLABORATOR} gives explicitly. Every other action's
* result is identical to the four-argument form's, since none of them consult the classifier.
*
* <p>Its default classifier denies every collaborator, so a caller enforcing authorization
* must use the four-argument form instead.
*/
public static boolean permits(Principal caller, Action action, String targetSession) {
return permits(caller, action, targetSession, NO_KNOWN_LEAD_OR_COLLABORATOR);
@@ -249,17 +249,6 @@ public final class CallerResolver {
return leadTerminals.get();
}
/**
* The currently-recognised architect slots, {@code terminal_id → slot name} (CB-548).
*
* <p>Read from the same supplier {@link #resolve} consults, so a slot that is <em>listed</em>
* here but would not <em>resolve</em> (or the reverse) cannot drift apart. Live for the same
* reason as {@link #leads()}.
*/
public Map<String, String> members() {
return architectTerminals.get();
}
/**
* The currently-recognised collaborator tabs, {@code terminal_id → name}.
*
@@ -615,7 +615,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
// may bind to, AND what a slot already bound still grants) through its own instance of that
// same supplier shape — see MemberRegistry.live and its class doc for the binding rule:
// removing a slot revokes ARCHITECT on the bound pane's very next request, and only the slot
// OCCUPANCY survives, so the demoted session keeps its slot key until it unbinds. Only
// OCCUPANCY survives, so the demoted session keeps its slot key until it unbinds.
// fleet.leaders is frozen (Fleetd.java:281 reads cfg.fleet().leaders() off the startup
// snapshot to build both the LeadTabScanner's tab-label-to-name map, wired into
// CallerResolver.withLeadsAndMembers at Fleetd.java:620/624, and — when herdr answered —
@@ -635,6 +635,11 @@ public final class ConfigRef implements Supplier<FleetConfig> {
+ "live through that same supplier for placement AND through a separate supplier "
+ "on MemberRegistry for spawn-time identity — both already applied");
}
if (!Objects.equals(collaboratorsOf(old), collaboratorsOf(fresh))) {
changed.add("fleet: fleet.collaborators (each collaborator's tab) is read once at "
+ "startup to build the LeadTabScanner's identity map, which is not rebuilt on "
+ "reload, so a collaborator added, removed, or given a new tab: needs a restart");
}
// Kept in step with SPLIT_KEYS the same way changedColdKeys is kept in step with COLD_KEYS —
// every message here must be traceable to one of the split keys the class doc documents.
// NOTE what this does NOT prove, per the javadoc above: it does not catch a SPLIT_KEYS
@@ -650,6 +655,11 @@ public final class ConfigRef implements Supplier<FleetConfig> {
return cfg.fleet() == null ? Map.of() : cfg.fleet().leaders();
}
/** {@code cfg.fleet().collaborators()}, defensively, in case a caller hands in a non-defaulted config. */
private static Map<String, FleetConfig.Collaborator> collaboratorsOf(FleetConfig cfg) {
return cfg.fleet() == null ? Map.of() : cfg.fleet().collaborators();
}
/**
* {@link FleetConfig.Profile} record components deliberately left out of
* {@link #sameLaunchSettings} because they are read <em>live</em>, not baked in at spawn — see
@@ -2776,9 +2776,10 @@ public record FleetConfig(
});
if (!bad.isEmpty()) {
throw new IllegalStateException("refusing to start: " + String.join("; ", bad)
+ ". Every member labelled that way would be read back as a lead or "
+ "collaborator and granted that identity's authority. Change one of the two "
+ "so member tabs cannot be confused with a lead's or collaborator's tab.");
+ ". A member labelled that way, while its pane carries no entry in the "
+ "spawned-member roster, is read back as a lead or collaborator and granted "
+ "that identity's authority. Change one of the two so member tabs cannot be "
+ "confused with a lead's or collaborator's tab.");
}
List<String> collisions = new ArrayList<>();
@@ -2877,7 +2878,8 @@ public record FleetConfig(
* the focused tab rather than its own, so it can land inside a lead's or collaborator's own
* labelled tab. {@link dev.ltms.fleet.herdr.LeadTabScanner} identifies a lead or collaborator
* purely by that tab's label — it does not exclude the member space — so a member that ends up
* there would be read back as that lead or collaborator and granted that identity's authority.
* there, while its pane carries no entry in the spawned-member roster, is read back as that
* lead or collaborator and granted that identity's authority.
*
* <p>Only an entry with a non-blank {@code tab} is in scope: one with no {@code tab} feeds
* nothing into {@link dev.ltms.fleet.herdr.LeadTabScanner}, so it creates no hazard here.
@@ -2908,8 +2910,9 @@ public record FleetConfig(
}
throw new IllegalStateException("refusing to start: profile(s) " + bad
+ " use placement: pane while fleet.leaders or fleet.collaborators names a tab. A "
+ "pane-placed member can land inside that labelled tab and be read back as the "
+ "lead or collaborator, granted that identity's authority. Set placement: tab for "
+ "pane-placed member can land inside that labelled tab, and while its pane "
+ "carries no entry in the spawned-member roster, it is read back as the lead or "
+ "collaborator and granted that identity's authority. Set placement: tab for "
+ "each named profile, or remove the tab from every fleet.leaders and "
+ "fleet.collaborators entry.");
}
@@ -490,7 +490,7 @@ public final class FleetMcp {
Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, caller);
// wait defaults to true (block for the reply); wait:false is fire-and-poll.
return Boolean.FALSE.equals(a.get("wait"))
? sendAsync(messages, target, content, onAccepted, workers.profiles())
? sendAsync(messages, target, content, onAccepted, workers.profiles(), caller)
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles());
};
// fleet_reply's identity is the CONNECTION, never an argument — so the authz check
@@ -524,7 +524,7 @@ public final class FleetMcp {
// The action depends on the ARGUMENTS, not on the tool name -- see pollAction.
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_poll", a), target);
if (denied != null) return denied;
return poll(messages, leadChannel, str(a, "ticket"), target, coordId);
return poll(messages, leadChannel, str(a, "ticket"), target, coordId, callerTerminal(exchange));
};
// CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox).
// Acking removes a reply from the inbox, so it is a drain, not a read.
@@ -560,8 +560,10 @@ public final class FleetMcp {
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage,
leadSeats, leadContextGauge, leadConfigDirs, callers.leads(),
callerTerminal(exchange),
callers.collaborators(), collaboratorsVisibleTo(principal(exchange)),
new CoordinationSource(leadChannel, peers),
coordinatorVisibleTo(principal(exchange)));
coordinatorVisibleTo(principal(exchange)),
leadsVisibleTo(principal(exchange)), membersVisibleTo(principal(exchange)));
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
(exchange, req) -> {
@@ -743,6 +745,48 @@ public final class FleetMcp {
return caller.isPrimary();
}
/**
* fleetd #703: who may see {@code fleet_list}'s {@code collaborators} array — a roster of the
* human-opened tabs this daemon recognises as named peers. Visible to exactly the roles that
* may {@link Authz.Action#SEND} to a named peer ({@code Authz.java}'s {@code SEND} case):
* the primary, an architect, and a collaborator (reaching another collaborator or a lead). A
* worker holds {@code READ} but can never {@code SEND} to a collaborator, so listing them to a
* worker would expose which human tabs exist on the host with no use to that caller. Split out
* for the same reason as {@link #coordinatorVisibleTo}: the decision must be unit-testable
* without fabricating an SDK {@code McpSyncServerExchange}, and the handler must call this
* named predicate rather than inlining the check.
*/
static boolean collaboratorsVisibleTo(Principal caller) {
return caller.isPrimary() || caller.isArchitect() || caller.isCollaborator();
}
/**
* Who may see {@code fleet_list}'s {@code leads} array — exactly the roles that may
* {@link Authz.Action#SEND} to a lead: the primary, an architect, and a collaborator. A
* collaborator's own {@code fleet_whoami} carries no lead address, and {@code leads} is the
* only place this tool gives one, so a collaborator needs this array to use the send it
* already holds. A worker can never {@code SEND} at all, so it still sees neither this array
* nor {@code members}; a worker's own facts come from {@code fleet_whoami} instead. Split
* out for the same reason as {@link #coordinatorVisibleTo} and
* {@link #collaboratorsVisibleTo}: the decision must be unit-testable without fabricating an
* SDK {@code McpSyncServerExchange}, and the handler must call this named predicate rather
* than inlining the check.
*/
static boolean leadsVisibleTo(Principal caller) {
return caller.isPrimary() || caller.isArchitect() || caller.isCollaborator();
}
/**
* Who may see {@code fleet_list}'s {@code members} array — the primary and an architect,
* which may {@link Authz.Action#SEND} to a member. A worker can never {@code SEND} at all,
* and a collaborator may {@code SEND} only to a lead or another collaborator, never to a
* spawned member, so neither sees this array even though {@link #leadsVisibleTo} grants a
* collaborator the sibling one.
*/
static boolean membersVisibleTo(Principal caller) {
return caller.isPrimary() || caller.isArchitect();
}
/** The worker identity resolved from this call's connection, or {@code null} if the primary. */
private static String callerTerminal(McpSyncServerExchange exchange) {
Object v = exchange.transportContext().get(CALLER_TERMINAL);
@@ -956,6 +1000,17 @@ public final class FleetMcp {
*/
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
Runnable onAccepted, Set<String> profiles) {
return sendAsync(messages, sessionId, content, onAccepted, profiles, null);
}
/**
* As above, recording {@code creatorTerminal} as this ticket's owner (the caller's own
* terminal, resolved from the connection) so a later {@code fleet_poll{ticket}} only hands the
* result back to that same caller — see {@link MessageService#poll(String, String)}.
*/
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
Runnable onAccepted, Set<String> profiles,
String creatorTerminal) {
if (isBlank(sessionId) || isBlank(content)) {
return error("sessionId and content are required");
}
@@ -963,7 +1018,7 @@ public final class FleetMcp {
if (targetError != null) {
return targetError;
}
String ticket = messages.sendAsync(sessionId, content, onAccepted);
String ticket = messages.sendAsync(sessionId, content, onAccepted, creatorTerminal);
return text("accepted — task delegated. Poll fleet_poll with ticket=" + ticket);
}
@@ -1142,9 +1197,24 @@ public final class FleetMcp {
* exclusively of {@code ticket}/{@code target} — see {@link #pollAction}'s javadoc for why this
* is a different inbox (this daemon's own {@link LeadChannel}) that authorizes differently
* ({@link Authz.Action#COORD_READ}, primary-only) from either of the original two branches.
*
* <p>Does not check who owns {@code ticket} — see the overload that takes {@code
* callerTerminal} for that. Callers that do not resolve a caller terminal (tests, or a surface
* with no connection-based identity) use this one.
*/
static McpSchema.CallToolResult poll(MessageService messages, LeadChannel leadChannel, String ticket,
String target, String coordId) {
return poll(messages, leadChannel, ticket, target, coordId, null);
}
/**
* As above, refusing a ticket lookup whose caller's terminal differs from the terminal that
* created it — see {@link MessageService#poll(String, String)}. {@code callerTerminal} is the
* CALLING session's terminal id, resolved by the MCP layer from the connection, never a
* client-supplied value.
*/
static McpSchema.CallToolResult poll(MessageService messages, LeadChannel leadChannel, String ticket,
String target, String coordId, String callerTerminal) {
if (!isBlank(coordId)) {
return pollHeldPeerMail(leadChannel, coordId);
}
@@ -1158,7 +1228,7 @@ public final class FleetMcp {
if (isBlank(ticket)) {
return error("ticket (or target) is required");
}
MessageService.TaskView v = messages.poll(ticket);
MessageService.TaskView v = messages.poll(ticket, callerTerminal);
if (v == null) {
return error("unknown ticket: " + ticket + " (never issued, or expired)");
}
@@ -1761,7 +1831,8 @@ public final class FleetMcp {
Map<String, String> leads, String selfTerm,
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine,
OutageSource.none(), LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
OutageSource.none(), LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, false, true, true);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
@@ -1770,7 +1841,8 @@ public final class FleetMcp {
QuarantineSource quarantine, OutageSource outage,
Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, CoordinationSource.none(), false);
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, CoordinationSource.none(), false, true, true);
}
/**
@@ -1787,7 +1859,8 @@ public final class FleetMcp {
QuarantineSource quarantine, Map<String, String> leads, String selfTerm,
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, OutageSource.none(),
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, false, true, true);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
@@ -1796,7 +1869,8 @@ public final class FleetMcp {
QuarantineSource quarantine, OutageSource outage,
Map<String, String> leads, String selfTerm, CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, false, true, true);
}
/**
@@ -1818,7 +1892,8 @@ public final class FleetMcp {
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, false, true, true);
}
/**
@@ -1835,14 +1910,23 @@ public final class FleetMcp {
* {@code false} (fleetd #463: a forgotten argument fails closed, not
* open), so a test that wants the {@code coordinator} row must pass
* an explicit {@code true}
* @param leadsVisible whether this caller may see the {@code leads} array (see
* {@link #leadsVisibleTo}) — unlike {@code callerIsPrimary}, this is
* also {@code true} for an architect or a collaborator, so it cannot
* be derived from {@code callerIsPrimary} alone
* @param membersVisible whether this caller may see the {@code members} array (see
* {@link #membersVisibleTo}); also {@code true} for an architect, but
* unlike {@code leadsVisible}, never for a collaborator
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, OutageSource outage,
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
CoordinationSource coordination, boolean callerIsPrimary) {
CoordinationSource coordination, boolean callerIsPrimary,
boolean leadsVisible, boolean membersVisible) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine,
outage, leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, callerIsPrimary);
outage, leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm,
Map.of(), false, coordination, callerIsPrimary, leadsVisible, membersVisible);
}
/**
@@ -1852,6 +1936,17 @@ public final class FleetMcp {
* matter); the one caller that matters for caching, {@code fleet_list}'s MCP handler, passes
* its own single long-lived instance instead (see {@code FleetMcp}'s {@code leadContextGauge}
* field).
*
* @param collaborators terminal_id → collaborator name, live from the resolver
* ({@link dev.ltms.fleet.auth.CallerResolver#collaborators()})
* @param collaboratorsVisible whether this caller may see the {@code collaborators} array (see
* {@link #collaboratorsVisibleTo}); every wrapper overload above
* passes {@code false}, so a test that wants the row must call this
* overload with an explicit {@code true}
* @param leadsVisible whether this caller may see the {@code leads} array (see
* {@link #leadsVisibleTo})
* @param membersVisible whether this caller may see the {@code members} array (see
* {@link #membersVisibleTo})
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
@@ -1860,32 +1955,56 @@ public final class FleetMcp {
LeadSeatSource leadSeats, LeadContextGauge contextGauge,
LeadConfigDirSource leadConfigDirs,
Map<String, String> leads, String selfTerm,
CoordinationSource coordination, boolean callerIsPrimary) {
Map<String, String> collaborators, boolean collaboratorsVisible,
CoordinationSource coordination, boolean callerIsPrimary,
boolean leadsVisible, boolean membersVisible) {
try {
Map<String, Agent> live = workers.list().stream()
.map(Agent.class::cast)
.filter(a -> a.terminalId() != null)
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
List<Map<String, Object>> leadRows = leads.entrySet().stream()
.sorted(Map.Entry.comparingByValue())
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm, contextGauge,
leadConfigDirs))
.toList();
// Neither row's assembly (leadView/memberCapacityView probing herdr for live status)
// runs unless at least one of them needs the live-agent lookup backing it.
Map<String, Agent> live = (leadsVisible || membersVisible)
? workers.list().stream()
.map(Agent.class::cast)
.filter(a -> a.terminalId() != null)
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b))
: Map.of();
// fleetd #209: this is the caller-driven fleet_list read that actually reports
// agentSessionId (via memberCapacityView -> SessionManager.rosterView), so it uses the
// resolving roster; the heartbeat/health/metrics timers stay on the plain sessions.roster().
List<MemberSession> roster = sessions.rosterResolved();
List<Map<String, Object>> out = roster.stream()
.map(s -> memberCapacityView(s, live.get(s.terminalId()), messages, capacity.clock().getAsLong()))
.toList();
Set<String> profiles = new java.util.TreeSet<>(capacity.configuredProfiles().get());
roster.stream().map(MemberSession::profile).forEach(profiles::add);
Map<String, Object> result = new LinkedHashMap<>();
result.put("leads", leadRows); result.put("members", out);
// READ is permission to enter this tool, not permission to receive every field it can
// build -- gate BEFORE assembling each row, so the key is absent rather than
// present-and-empty; a caller without either row gets its own facts from fleet_whoami.
if (leadsVisible) {
List<Map<String, Object>> leadRows = leads.entrySet().stream()
.sorted(Map.Entry.comparingByValue())
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm, contextGauge,
leadConfigDirs))
.toList();
result.put("leads", leadRows);
}
if (membersVisible) {
List<Map<String, Object>> out = roster.stream()
.map(s -> memberCapacityView(s, live.get(s.terminalId()), messages, capacity.clock().getAsLong()))
.toList();
result.put("members", out);
}
result.put("healthCoverage", healthCoverage.value().get());
result.put("loopHealth", Map.of(
"statusPoller", loopHealth.statusPoller().get().name(),
"sessionReaper", loopHealth.sessionReaper().get().name()));
// fleetd #703: a collaborator tab is a person's own tab, so the row is assembled and
// included only for the roles that may SEND to a named peer -- gate BEFORE assembling
// it, same reason as the coordinator row just below: the key must be absent for a
// worker, never present-and-empty.
if (collaboratorsVisible) {
result.put("collaborators", collaborators.entrySet().stream()
.sorted(Map.Entry.comparingByValue())
.map(e -> collaboratorRow(e.getKey(), e.getValue()))
.toList());
}
// fleetd #439: coordinator/coordinatorView is lead-to-lead coordination state and must
// never reach a worker or an architect -- gate BEFORE assembling it, not after, so the
// key is absent rather than present-and-empty.
@@ -2197,6 +2316,20 @@ public final class FleetMcp {
return m;
}
/**
* fleetd #703: one row of the {@code collaborators} array — a named peer's registry name and
* the herdr {@code terminal_id} a {@code fleet_send} must target to reach it. No status, no
* context gauge, no profile: a collaborator is never spawned and carries no profile, so those
* fields have no meaning for it, and the scan behind {@code terminal} only reports a tab it
* actually found in herdr, so a tab nobody has open does not appear here at all.
*/
private static Map<String, Object> collaboratorRow(String terminal, String name) {
Map<String, Object> m = new LinkedHashMap<>();
m.put("name", name);
m.put("sessionId", terminal);
return m;
}
/**
* Reads the lead context gauge for one lead. {@code configDir} is this lead's configured
* {@code CLAUDE_CONFIG_DIR} override (see {@link LeadConfigDirSource}), derived from
@@ -2391,7 +2524,14 @@ public final class FleetMcp {
private static McpSchema.Tool listTool() {
return tool(FleetTool.LIST.wireName(),
"List the whole fleet the bridge tracks, in two parts. 'leads' are your PEERS — other "
"List the whole fleet the bridge tracks, in two parts. 'members' is visible to "
+ "the primary and an architect only. 'leads' is visible to those two AND a "
+ "collaborator — exactly the roles that may fleet_send to a lead, so a "
+ "collaborator can learn a lead's sessionId before using the send it already "
+ "holds. A worker holds READ to call this tool at all, but gets neither "
+ "array, never an empty one; a worker reads its own session, "
+ "profile, state, worktree, branch and owner from fleet_whoami instead. "
+ "'leads' are your PEERS — other "
+ "orchestrators, each with its sessionId (the address to fleet_send to), "
+ "name, live status, and 'self': true on your own row; this is how you "
+ "discover a peer lead without being told its address. 'members' are the "
@@ -2422,7 +2562,12 @@ public final class FleetMcp {
+ "msgId/from/preview, never the full body), and one row per coordinator.peers "
+ "coord-id ('peers': coordId/reachable, plus pending/consumers when reachable) — "
+ "this is peer DISCOVERY for cross-host leads, distinct from the local 'leads' "
+ "array above. It is omitted entirely when no coordinator is configured.",
+ "array above. It is omitted entirely when no coordinator is configured. A "
+ "'collaborators' array, visible only to the primary, an architect, and a "
+ "collaborator (never a worker), reports every other named peer tab this daemon "
+ "recognises: each row is 'name' (its registry name) and 'sessionId' (the "
+ "terminal id a fleet_send targets to reach it). Absent entirely for a worker, "
+ "whatever is configured.",
objectSchema(Map.of(), List.of()));
}
@@ -275,10 +275,19 @@ public final class MessageService {
* {@link #abandon}) can never match again regardless of this flag's value.
*/
private volatile boolean askTimedOut;
/**
* The terminal of the caller whose {@code fleet_send{wait:false}} created this ticket, or
* {@code null} when that caller had no terminal (the unnamed primary) or the ticket was
* created through an overload that does not record one. {@link #poll(String, String)}
* compares a polling caller's own terminal against this field before handing back the
* ticket's state.
*/
private final String creatorTerminal;
private Task(String ticket, String target, LongSupplier nowNanos) {
private Task(String ticket, String target, LongSupplier nowNanos, String creatorTerminal) {
this.ticket = ticket;
this.target = target;
this.creatorTerminal = creatorTerminal;
this.createdNanos = nowNanos.getAsLong();
future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong());
}
@@ -1279,8 +1288,20 @@ public final class MessageService {
* @return the ticket to poll for the eventual result
*/
public String sendAsync(String target, String content, Runnable onAccepted) {
return sendAsync(target, content, onAccepted, null);
}
/**
* As {@link #sendAsync(String, String, Runnable)}, recording {@code creatorTerminal} as this
* ticket's owner. {@link #poll(String, String)} refuses a later caller whose own terminal
* differs from this one; {@code null} records no owner (a caller with no terminal — the
* unnamed primary — is always allowed to poll the result regardless).
*
* @return the ticket to poll for the eventual result
*/
public String sendAsync(String target, String content, Runnable onAccepted, String creatorTerminal) {
String ticket = "task-" + ticketSeq.incrementAndGet();
Task task = new Task(ticket, target, nowNanos);
Task task = new Task(ticket, target, nowNanos, creatorTerminal);
tasks.put(ticket, task);
if (pushLoop != null) {
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a
@@ -1343,15 +1364,31 @@ public final class MessageService {
}
/**
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket;
* otherwise a {@link Phase#PENDING} view (with the live worker status as detail), a
* {@link Phase#DONE} view carrying the reply, or a {@link Phase#FAILED} view with the reason.
* As {@link #poll(String, String)}, with no caller terminal — the ticket's ownership is never
* checked, so this overload must only be used where the caller's identity is otherwise
* irrelevant (a test, or a surface that does not resolve a caller terminal at all).
*/
public TaskView poll(String ticket) {
return poll(ticket, null);
}
/**
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket.
* Refuses a {@code callerTerminal} that differs from the terminal that created the ticket (see
* {@link #sendAsync(String, String, Runnable, String)}) with a {@link Phase#FAILED} view that
* carries no reply text — a caller with no terminal (the unnamed primary) is never refused.
* Otherwise returns a {@link Phase#PENDING} view (with the live worker status as detail), a
* {@link Phase#DONE} view carrying the reply, or a {@link Phase#FAILED} view with the reason.
*/
public TaskView poll(String ticket, String callerTerminal) {
Task task = tasks.get(ticket);
if (task == null) {
return null;
}
if (!ownsTicket(task, callerTerminal)) {
return new TaskView(ticket, Phase.FAILED, null, null,
"forbidden: this ticket was created by a different session", null);
}
CompletableFuture<Reply> f = task.future;
if (!f.isDone()) {
Reply question = task.question;
@@ -1387,6 +1424,17 @@ public final class MessageService {
return new TaskView(ticket, Phase.FAILED, null, null, detail, null);
}
/**
* Whether {@code callerTerminal} may read {@code task}'s state. A caller with no terminal
* always may — that is the unnamed primary, resolved by token or loopback trust, which never
* carries a herdr pane and must keep reading every ticket. Otherwise the caller's terminal must
* equal the terminal recorded on the task; a task with no recorded terminal matches no
* terminal-bearing caller.
*/
private static boolean ownsTicket(Task task, String callerTerminal) {
return callerTerminal == null || callerTerminal.equals(task.creatorTerminal);
}
/**
* Test seam only — carries no production behaviour, and nothing in this class calls it;
* {@link #pruneTerminalTickets} still reads {@link Task#completedNanos} directly.
@@ -694,7 +694,8 @@ public final class FleetApp {
if (!wait) {
// Fire-and-poll (CB-107): return a ticket immediately; the caller polls GET /tasks/{ticket}.
String ticket = messages.sendAsync(id, content);
Principal caller = ctx.attribute(CALLER);
String ticket = messages.sendAsync(id, content, null, caller == null ? null : caller.terminal());
ctx.status(202).json(Map.of("sessionId", id, "ticket", ticket, "status", "accepted"));
return;
}
@@ -894,7 +895,8 @@ public final class FleetApp {
if (!allow(ctx, routeAction("GET /tasks/{ticket}"), null)) {
return;
}
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"));
Principal caller = ctx.attribute(CALLER);
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"), caller == null ? null : caller.terminal());
if (v == null) {
ctx.status(404).json(Map.of("error", "unknown_ticket", "detail", "no such task (or it has expired)"));
return;
@@ -14,10 +14,12 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* CB-534: the injector's readiness gate must open for a lead as well as for a present worker.
* fleetd #669 follow-up: the same gate must also open for a collaborator, which — like a lead —
* is never enrolled in {@link MemberPresence} and never discovered by the lead scan.
*
* <p>The bug these cover was silent and slow: a lead was never marked present (only workers are), so
* every lead→lead delivery sat on the gate for the full readiness grace and failed ~60s later without
* a keystroke ever reaching the pane.
* <p>The bug these cover was silent and slow: a lead (and later a collaborator) was never marked
* present (only workers are) and never counted as a lead, so every send to one sat on the gate for
* the full readiness grace and failed ~60s later without a keystroke ever reaching the pane.
*/
class FleetDeliverabilityTest {
@@ -25,19 +27,25 @@ class FleetDeliverabilityTest {
return () -> m;
}
private static Supplier<Map<String, String>> collaborators(Map<String, String> m) {
return () -> m;
}
@Test
@DisplayName("a worker that has connected its MCP is deliverable")
void presentWorkerIsDeliverable() {
MemberPresence presence = new MemberPresence();
presence.markPresent("term_worker");
assertTrue(Fleetd.deliverableTo(presence, leads(Map.of())).test("term_worker"));
assertTrue(Fleetd.deliverableTo(presence, leads(Map.of()), collaborators(Map.of()))
.test("term_worker"));
}
@Test
@DisplayName("a worker still in its boot window is held back")
void absentWorkerIsNotDeliverable() {
assertFalse(Fleetd.deliverableTo(new MemberPresence(), leads(Map.of())).test("term_booting"));
assertFalse(Fleetd.deliverableTo(new MemberPresence(), leads(Map.of()), collaborators(Map.of()))
.test("term_booting"));
}
@Test
@@ -45,19 +53,31 @@ class FleetDeliverabilityTest {
void leadIsDeliverableWithoutPresence() {
MemberPresence presence = new MemberPresence();
Predicate<String> deliverable =
Fleetd.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")));
Fleetd.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")), collaborators(Map.of()));
assertFalse(presence.isPresent("term_lead"), "a lead is never enrolled in worker presence");
assertTrue(deliverable.test("term_lead"), "…and must be deliverable anyway");
}
@Test
@DisplayName("an unknown terminal is deliverable to neither")
@DisplayName("a collaborator is deliverable without ever being marked present or scanned as a lead")
void collaboratorIsDeliverableWithoutPresenceOrLeadStatus() {
MemberPresence presence = new MemberPresence();
Predicate<String> deliverable = Fleetd.deliverableTo(presence, leads(Map.of()),
collaborators(Map.of("term_collab", "kevin")));
assertFalse(presence.isPresent("term_collab"), "a collaborator is never enrolled in worker presence");
assertTrue(deliverable.test("term_collab"), "…and must be deliverable anyway");
}
@Test
@DisplayName("an unknown terminal is deliverable to none of presence, leads, or collaborators")
void strangerIsNotDeliverable() {
MemberPresence presence = new MemberPresence();
presence.markPresent("term_worker");
assertFalse(Fleetd.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")))
assertFalse(Fleetd.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")),
collaborators(Map.of("term_collab", "kevin")))
.test("term_stranger"));
}
@@ -65,22 +85,47 @@ class FleetDeliverabilityTest {
@DisplayName("a lead discovered after startup becomes deliverable with no restart")
void leadSetIsReadThroughOnEveryCall() {
Map<String, String> discovered = new HashMap<>();
Predicate<String> deliverable = Fleetd.deliverableTo(new MemberPresence(), leads(discovered));
Predicate<String> deliverable =
Fleetd.deliverableTo(new MemberPresence(), leads(discovered), collaborators(Map.of()));
assertFalse(deliverable.test("term_late"));
discovered.put("term_late", "gpt-sol-5.6"); // leadScan picks up a newly labelled tab
assertTrue(deliverable.test("term_late"), "the supplier must be re-read, not snapshotted");
}
@Test
@DisplayName("a collaborator discovered after startup becomes deliverable with no restart")
void collaboratorSetIsReadThroughOnEveryCall() {
Map<String, String> discovered = new HashMap<>();
Predicate<String> deliverable =
Fleetd.deliverableTo(new MemberPresence(), leads(Map.of()), collaborators(discovered));
assertFalse(deliverable.test("term_late_collab"));
discovered.put("term_late_collab", "kevin"); // the same tab scan picks up a newly labelled collaborator tab
assertTrue(deliverable.test("term_late_collab"), "the supplier must be re-read, not snapshotted");
}
@Test
@DisplayName("forgetting a torn-down worker does not strip a lead of its deliverability")
void forgetDoesNotDisarmALead() {
MemberPresence presence = new MemberPresence();
Predicate<String> deliverable =
Fleetd.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")));
Fleetd.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")), collaborators(Map.of()));
presence.forget("term_lead"); // the injector's cleanup path runs against every target
assertTrue(deliverable.test("term_lead"));
}
@Test
@DisplayName("forgetting a torn-down worker does not strip a collaborator of its deliverability")
void forgetDoesNotDisarmACollaborator() {
MemberPresence presence = new MemberPresence();
Predicate<String> deliverable = Fleetd.deliverableTo(presence, leads(Map.of()),
collaborators(Map.of("term_collab", "kevin")));
presence.forget("term_collab"); // the injector's cleanup path runs against every target
assertTrue(deliverable.test("term_collab"));
}
}
@@ -31,8 +31,7 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #670 — pins the {@code excludedWorkspaceLabels} argument {@link FleetdAssembly}'s
* production boot path passes to {@link LeadTabScanner} at {@code FleetdAssembly.java:265}
* ({@code Set.of()}).
* production boot path passes to {@link LeadTabScanner} ({@code Set.of()}).
*
* <p>{@code LeadTabScannerTest} already covers this constructor parameter, but it builds its own
* {@link LeadTabScanner} with its own set, so it tests the seam and proves nothing about the
@@ -191,7 +190,7 @@ class FleetdAssemblyLeadTabScannerExclusionTest {
Set<?> excluded = (Set<?>) excludedField.get(leads);
assertTrue(excluded.isEmpty(),
"FleetdAssembly.java:265 must pass an empty excludedWorkspaceLabels to "
"FleetdAssembly must pass an empty excludedWorkspaceLabels to "
+ "LeadTabScanner — scanning member tabs would demote the lead to a worker");
} finally {
assertNotNull(ports.shutdownHook, "control: assembly must capture its shutdown hook");
@@ -460,7 +460,6 @@ class CallerResolverTest {
assertTrue(members.bind("architect:lead-designer", "term_a")); // the later lifecycle binds the slot
assertEquals(Role.ARCHITECT, r.resolve("127.0.0.1", 42, null).role());
assertEquals("architect:lead-designer", r.members().get("term_a"));
}
@Test
@@ -107,6 +107,8 @@ class ConfigRefTest {
assertTrue(out.applied());
assertTrue(out.deferred().isEmpty());
assertTrue(out.split().stream().noneMatch(s -> s.contains("fleet.collaborators")),
out.split().toString());
assertEquals("new charter", ref.get().fleet().charterFor(
dev.ltms.fleet.peer.MemberRole.ARCHITECT));
}
@@ -786,14 +788,100 @@ class ConfigRefTest {
assertTrue(out.split().getFirst().startsWith("fleet:"), out.split().toString());
assertTrue(out.split().getFirst().contains("restart"), out.split().toString());
assertTrue(out.split().getFirst().contains("live"), out.split().toString());
assertTrue(out.split().stream().noneMatch(s -> s.contains("fleet.collaborators")),
out.split().toString());
assertTrue(out.summary().contains("partially live"), out.summary());
// The snapshot still carries the new value — the LeadTabScanner's identity map and
// LeadLauncher's auto-launch are what wait for a restart; a reload rebuilds neither.
assertEquals("lead: opus-b", ref.get().fleet().leaders().get("opus").tab());
}
@Test
void addingAFleetCollaboratorIsReportedAsSplit(@TempDir Path dir) throws Exception {
assertCollaboratorChangeIsReported(dir, """
fleet:
collaborators:
alex:
tab: "collaborator: alex"
""", "add a collaborator");
}
@Test
void removingAFleetCollaboratorIsReportedAsSplit(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml("""
fleet:
collaborators:
alex:
tab: "collaborator: alex"
"""));
ConfigRef ref = refFor(f);
Files.writeString(f, yaml("fleet: {}\n"));
assertCollaboratorSplit(ref.reload(), "remove a collaborator");
}
@Test
void changingAFleetCollaboratorTabIsReportedAsSplit(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml("""
fleet:
collaborators:
alex:
tab: "collaborator: alex-a"
"""));
ConfigRef ref = refFor(f);
Files.writeString(f, yaml("""
fleet:
collaborators:
alex:
tab: "collaborator: alex-b"
"""));
assertCollaboratorSplit(ref.reload(), "change a collaborator tab");
}
@Test
void unchangedFleetCollaboratorsProduceNoSplitReport(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
String config = yaml("""
fleet:
collaborators:
alex:
tab: "collaborator: alex"
""");
Files.writeString(f, config);
ConfigRef ref = refFor(f);
Files.writeString(f, config);
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertTrue(out.split().isEmpty(), out.split().toString());
}
private static void assertCollaboratorChangeIsReported(Path dir, String changed, String action)
throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml("fleet: {}\n"));
ConfigRef ref = refFor(f);
Files.writeString(f, yaml(changed));
assertCollaboratorSplit(ref.reload(), action);
}
private static void assertCollaboratorSplit(ConfigRef.Outcome out, String action) {
assertTrue(out.applied(), action);
assertTrue(out.deferred().isEmpty(), out.deferred().toString());
assertEquals(1, out.split().size(), out.split().toString());
String report = out.split().getFirst();
assertTrue(report.startsWith("fleet: fleet.collaborators"), report);
assertTrue(report.contains("read once at startup"), report);
assertTrue(report.contains("restart"), report);
}
/**
* fleetd #333: {@code fleet.leaders} is the ONLY frozen part of {@code fleet:}. A reload that
* fleetd #333: {@code fleet.leaders} is a frozen part of {@code fleet:}. A reload that
* changes {@code tabLabel} (or charters, or a role pool) without touching {@code fleet.leaders}
* must stay fully hot with nothing reported — proving {@link ConfigRef#changedSplitKeys}
* compares {@code fleet.leaders} specifically rather than the whole {@code Fleet} record, which
@@ -850,7 +850,8 @@ class FleetConfigTest {
/**
* The hazard this guard closes: a pane-placed member lands inside the focused tab rather than
* its own, so it can land inside a lead's labelled tab and be read back as that lead.
* its own, so it can land inside a lead's labelled tab and, while its pane carries no entry in
* the spawned-member roster, be read back as that lead.
*/
@Test
void aPanePlacedProfileWithALeadTabRefusesToStart(@TempDir Path dir) throws Exception {
@@ -1087,8 +1088,9 @@ class FleetConfigTest {
}
/**
* fleetd #669: a member tabLabel that can render as a configured collaborator tab is the same
* hazard as the lead case above — a member labelled that way is read back as the collaborator.
* A member tabLabel that can render as a configured collaborator tab is the same hazard as the
* lead case above — while its pane carries no entry in the spawned-member roster, a member
* labelled that way is read back as the collaborator.
*/
@Test
void aProfileTabLabelOverrideMatchingACollaboratorTabRefusesToStart(@TempDir Path dir)
@@ -205,9 +205,12 @@ class LeadTabScannerTest {
}
/**
* The guard that matters: fleetd labels its own worker tabs, so if a worker space were scanned
* a naming accident would promote the fleet. The exclusion is by workspace, not by hoping the
* worker template never collides.
* Covers the {@code excludedWorkspaceLabels} parameter: a tab in an excluded workspace is never
* matched, whatever its label. Production always constructs this class with an empty set (CB-558,
* {@code FleetdAssembly}), so this parameter plays no part in the live guard against a worker
* tab being mistaken for a lead — that guard is {@code CallerResolver} asking the live
* spawned-member roster before any tab map. This test exists because the parameter still exists
* and is worth covering on its own terms.
*/
@Test
void aTabInAWorkerSpaceIsNeverALeadEvenWhenItsLabelMatches() {
@@ -219,6 +222,29 @@ class LeadTabScannerTest {
assertFalse(scanner(herdr, tabToName, new AtomicLong()).get().containsKey("term_impostor"));
}
/**
* The shipped shape (CB-558, {@code FleetdAssembly}): production always constructs this class
* with an empty {@code excludedWorkspaceLabels}, and a lead's {@code workspace:} default is the
* same shared {@code "fleet"} space the members use. A lead tab is still found when it sits in
* the exact same workspace as a member-labelled tab — the scanner tells them apart by the exact
* tab label, not by which workspace either one is in.
*/
@Test
void aLeadIsDiscoveredWhenItsWorkspaceIsTheSameAsTheMemberWorkspace() {
TopologyHerdr herdr = new TopologyHerdr()
.workspace("w1", "fleet")
.tab("w1:t1", "w1", "lead: opus-5.0")
.tab("w1:t2", "w1", "worker: gx10 #1")
.pane("w1:p1", "w1:t1", "term_opus")
.pane("w1:p2", "w1:t2", "term_worker");
LeadTabScanner s = new LeadTabScanner(herdr, Map.of("lead: opus-5.0", "opus-5.0"),
Set.of(), TTL, new AtomicLong()::get);
assertEquals("opus-5.0", s.get().get("term_opus"),
"a lead sharing the members' workspace is still discovered — the label, not the "
+ "workspace, is what matches it");
}
@Test
void aLabelWithNoConfiguredEntryIsIgnored() {
TopologyHerdr herdr = new TopologyHerdr().workspace("w1", "main")
@@ -355,10 +355,10 @@ class FleetMcpAuthzTest {
+ "the anchors have drifted, this test is not testing what it claims to");
Pattern trailingArg = Pattern.compile(
"listFleet\\([^;]*?,\\s*(coordinatorVisibleTo\\(principal\\(exchange\\)\\)|true|false)\\s*\\)\\s*;",
"listFleet\\([^;]*?,\\s*(coordinatorVisibleTo\\(principal\\(exchange\\)\\)|true|false)\\s*[,)]",
Pattern.DOTALL);
Matcher m = trailingArg.matcher(handlerBlock);
assertTrue(m.find(), "could not locate listFleet(...)'s trailing boolean argument in the "
assertTrue(m.find(), "could not locate listFleet(...)'s coordinator boolean argument in the "
+ "listHandler block -- the call shape changed, update this test's anchor: " + handlerBlock);
String trailing = m.group(1);
assertEquals("coordinatorVisibleTo(principal(exchange))", trailing,
@@ -366,6 +366,170 @@ class FleetMcpAuthzTest {
+ "calling, not pass a literal boolean -- found: " + trailing);
}
// --- fleetd #703: who may see fleet_list's collaborators array -------------------------------
/**
* fleetd #703: {@link FleetMcp#collaboratorsVisibleTo} is the whole policy decision for
* {@code fleet_list}'s {@code collaborators} array. Visible to exactly the roles that may
* {@code SEND} to a named peer -- the primary, an architect, and a collaborator itself -- never
* a worker, which holds {@code READ} but can never {@code SEND} to a collaborator, and never an
* anonymous caller.
*/
@Test
void onlyPrimaryArchitectAndCollaboratorMaySeeTheCollaboratorsArray() {
assertTrue(FleetMcp.collaboratorsVisibleTo(PRIMARY), "the primary must see the collaborators array");
assertTrue(FleetMcp.collaboratorsVisibleTo(ARCH_DESIGN), "an architect must see the collaborators array");
assertTrue(FleetMcp.collaboratorsVisibleTo(COLLABORATOR), "a collaborator must see its own peer roster");
assertFalse(FleetMcp.collaboratorsVisibleTo(WORKER_A),
"a worker holds READ but can never SEND to a collaborator, so it must not see the array");
assertFalse(FleetMcp.collaboratorsVisibleTo(ANON), "authenticated as nothing must not see it either");
}
/**
* fleetd #703, same reasoning as {@link #theFleetListHandlerActuallyConsultsCoordinatorVisibleTo}:
* the predicate above can be perfectly correct while the one production call site never asks it.
* This reads {@code FleetMcp.java}'s own source and asserts the {@code fleet_list} handler's
* {@code listFleet(...)} call both threads {@code callers.collaborators()} into the payload and
* asks {@code collaboratorsVisibleTo(principal(exchange))} for the visibility flag, rather than a
* literal boolean or an empty map.
*/
@Test
void theFleetListHandlerActuallyConsultsCollaboratorsVisibleTo() throws Exception {
String source = Files.readString(MCP_SOURCE);
int start = source.indexOf("listHandler =");
assertTrue(start >= 0, "could not find the fleet_list handler (listHandler) in " + MCP_SOURCE
+ " -- the scrape has stopped matching, fix the anchor before trusting this test");
int end = source.indexOf("stopHandler =", start);
assertTrue(end > start, "could not find the handler declared after listHandler to bound the scrape");
String handlerBlock = source.substring(start, end);
// CONTROL: the block we scraped really does contain a call to listFleet(...) -- if this
// fails, the anchors above moved and the assertions below would otherwise pass on nothing.
assertTrue(handlerBlock.contains("listFleet("),
"control failed: the scraped listHandler block contains no listFleet( call at all -- "
+ "the anchors have drifted, this test is not testing what it claims to");
assertTrue(handlerBlock.contains("callers.collaborators()"),
"the fleet_list handler must thread callers.collaborators() into listFleet(...), not an "
+ "empty or literal map -- block: " + handlerBlock);
assertTrue(handlerBlock.contains("collaboratorsVisibleTo(principal(exchange))"),
"the fleet_list handler must ask collaboratorsVisibleTo(principal(exchange)) who is "
+ "calling, not pass a literal boolean -- block: " + handlerBlock);
}
// --- who may see fleet_list's leads and members arrays ---------------------------------------
/**
* {@link FleetMcp#leadsVisibleTo} is the whole policy decision for {@code fleet_list}'s
* {@code leads} array: visible to exactly the roles that may {@code SEND} to a lead -- the
* primary, an architect, and a collaborator -- never a worker, which holds {@code READ} but
* can never {@code SEND} at all, and never an anonymous caller.
*/
@Test
void primaryArchitectAndCollaboratorMaySeeTheLeadsArray() {
assertTrue(FleetMcp.leadsVisibleTo(PRIMARY), "the primary must see the leads array");
assertTrue(FleetMcp.leadsVisibleTo(ARCH_DESIGN), "an architect must see the leads array");
assertTrue(FleetMcp.leadsVisibleTo(COLLABORATOR),
"a collaborator may SEND to a lead, so it must see the leads array to learn where");
assertFalse(FleetMcp.leadsVisibleTo(WORKER_A),
"a worker holds READ but can never SEND, so it must not see the leads array");
assertFalse(FleetMcp.leadsVisibleTo(ANON), "authenticated as nothing must not see it either");
}
/**
* As {@link #primaryArchitectAndCollaboratorMaySeeTheLeadsArray}, for the {@code members}
* array -- but a collaborator may {@code SEND} only to a lead or another collaborator, never
* to a spawned member, so it must not see this one.
*/
@Test
void onlyPrimaryAndArchitectMaySeeTheMembersArray() {
assertTrue(FleetMcp.membersVisibleTo(PRIMARY), "the primary must see the members array");
assertTrue(FleetMcp.membersVisibleTo(ARCH_DESIGN), "an architect must see the members array");
assertFalse(FleetMcp.membersVisibleTo(WORKER_A),
"a worker holds READ but can never SEND, so it must not see the members array");
assertFalse(FleetMcp.membersVisibleTo(COLLABORATOR),
"a collaborator may SEND to a lead, never to a spawned member, so it must not see the members array");
assertFalse(FleetMcp.membersVisibleTo(ANON), "authenticated as nothing must not see it either");
}
/**
* Same reasoning as {@link #theFleetListHandlerActuallyConsultsCoordinatorVisibleTo}: the
* predicate above can be perfectly correct while the one production call site never asks it.
* This reads {@code FleetMcp.java}'s own source and asserts the {@code fleet_list} handler's
* {@code listFleet(...)} call asks {@code leadsVisibleTo(principal(exchange))} for the leads
* visibility flag, rather than a literal boolean.
*/
@Test
void theFleetListHandlerActuallyConsultsLeadsVisibleTo() throws Exception {
String source = Files.readString(MCP_SOURCE);
int start = source.indexOf("listHandler =");
assertTrue(start >= 0, "could not find the fleet_list handler (listHandler) in " + MCP_SOURCE
+ " -- the scrape has stopped matching, fix the anchor before trusting this test");
int end = source.indexOf("stopHandler =", start);
assertTrue(end > start, "could not find the handler declared after listHandler to bound the scrape");
String handlerBlock = source.substring(start, end);
// CONTROL: the block we scraped really does contain a call to listFleet(...) -- if this
// fails, the anchors above moved and the assertion below would otherwise pass on nothing.
assertTrue(handlerBlock.contains("listFleet("),
"control failed: the scraped listHandler block contains no listFleet( call at all -- "
+ "the anchors have drifted, this test is not testing what it claims to");
assertTrue(handlerBlock.contains("leadsVisibleTo(principal(exchange))"),
"the fleet_list handler must ask leadsVisibleTo(principal(exchange)) who is calling, "
+ "not pass a literal boolean -- block: " + handlerBlock);
}
/** As {@link #theFleetListHandlerActuallyConsultsLeadsVisibleTo}, for {@code membersVisibleTo}. */
@Test
void theFleetListHandlerActuallyConsultsMembersVisibleTo() throws Exception {
String source = Files.readString(MCP_SOURCE);
int start = source.indexOf("listHandler =");
assertTrue(start >= 0, "could not find the fleet_list handler (listHandler) in " + MCP_SOURCE
+ " -- the scrape has stopped matching, fix the anchor before trusting this test");
int end = source.indexOf("stopHandler =", start);
assertTrue(end > start, "could not find the handler declared after listHandler to bound the scrape");
String handlerBlock = source.substring(start, end);
assertTrue(handlerBlock.contains("listFleet("),
"control failed: the scraped listHandler block contains no listFleet( call at all -- "
+ "the anchors have drifted, this test is not testing what it claims to");
assertTrue(handlerBlock.contains("membersVisibleTo(principal(exchange))"),
"the fleet_list handler must ask membersVisibleTo(principal(exchange)) who is calling, "
+ "not pass a literal boolean -- block: " + handlerBlock);
}
/**
* {@code fleet_poll{ticket}} must thread the calling connection's own terminal into
* {@link MessageService#poll(String, String)}, so a worker cannot read a ticket a different
* session created.
*/
@Test
void theFleetPollHandlerActuallyThreadsCallerTerminalIntoPoll() throws Exception {
String source = Files.readString(MCP_SOURCE);
int start = source.indexOf("pollHandler =");
assertTrue(start >= 0, "could not find the fleet_poll handler (pollHandler) in " + MCP_SOURCE
+ " -- the scrape has stopped matching, fix the anchor before trusting this test");
int end = source.indexOf("ackHandler =", start);
assertTrue(end > start, "could not find the handler declared after pollHandler to bound the scrape");
String handlerBlock = source.substring(start, end);
// CONTROL: the block we scraped really does call poll(...) -- if this fails, the anchors
// above moved and the assertions below would otherwise pass on nothing.
assertTrue(handlerBlock.contains("poll(messages,"),
"control failed: the scraped pollHandler block contains no poll(messages, ...) call "
+ "at all -- the anchors have drifted, this test is not testing what it claims to");
assertTrue(handlerBlock.contains("callerTerminal(exchange)"),
"the fleet_poll handler must thread callerTerminal(exchange) into poll(...), not omit "
+ "it or pass a literal null -- block: " + handlerBlock);
}
// --- which action each tool hands the gate (fleetd #272) ------------------------------------
/**
@@ -121,6 +121,7 @@ class FleetMcpLeadContextGaugeWiringTest {
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.LoopHealthSource.none(), FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
FleetMcp.LeadSeatSource.none(), contextGauge, leadConfigDirs,
Map.of(LEAD_TERMINAL, LEAD_NAME), LEAD_TERMINAL, FleetMcp.CoordinationSource.none(), false);
Map.of(LEAD_TERMINAL, LEAD_NAME), LEAD_TERMINAL, Map.of(), false,
FleetMcp.CoordinationSource.none(), false, true, true);
}
}
@@ -10,6 +10,7 @@ import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.inject.LoopWatchdog;
import dev.ltms.fleet.lead.LeadContextGauge;
import dev.ltms.fleet.herdr.PaneLocator;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.Injector;
@@ -691,7 +692,7 @@ class FleetMcpTest {
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
String out = textOf(res);
// fleetd #361: reports both which coord-id a peer must use to reach ME, and this daemon's
@@ -720,7 +721,7 @@ class FleetMcpTest {
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
String out = textOf(res);
assertTrue(out.contains("\"mailbox\":{\"status\":\"unknown\"}"), out);
@@ -779,19 +780,82 @@ class FleetMcpTest {
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
.withMailbox("mac-opus", LeadChannel.MailboxState.exists("mac-opus", 0, 1));
boolean callerIsPrimary = Principal.worker("term_a", 1).isPrimary();
Principal worker = Principal.worker("term_a", 1);
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), callerIsPrimary);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), worker.isPrimary(),
FleetMcp.leadsVisibleTo(worker), FleetMcp.membersVisibleTo(worker));
String out = textOf(res);
assertFalse(out.contains("\"coordinator\""), "a worker must never see the coordinator key at all: " + out);
assertFalse(out.contains("mac-opus"), "no fragment of the coordinator row may leak either: " + out);
assertTrue(out.contains("\"leads\""), "the rest of the result must still be present: " + out);
assertTrue(out.contains("\"members\""), out);
assertFalse(out.contains("\"leads\""), "a worker must never see the leads key at all: " + out);
assertFalse(out.contains("\"members\""), "a worker must never see the members key at all: " + out);
assertTrue(out.contains("\"healthCoverage\""), "the rest of the result must still be present: " + out);
}
/**
* A worker's result must carry neither the {@code leads} nor the {@code members} key, and no
* fragment of either row leaks even though both are fully populated for this call -- a key
* check alone would pass on an implementation that still built the rows and only renamed or
* nested them.
*/
@Test
void listLeaksNoLeadOrMemberRowFragmentToAWorker() {
FakeHerdr h = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", "term_owner",
new WorktreeRequest("cb-304", null));
Principal worker = Principal.worker("term_a", 1);
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of("term_lead_x", "mac-opus"), "", FleetMcp.CoordinationSource.none(), worker.isPrimary(),
FleetMcp.leadsVisibleTo(worker), FleetMcp.membersVisibleTo(worker));
String out = textOf(res);
assertFalse(out.contains("\"leads\""), out);
assertFalse(out.contains("\"members\""), out);
assertFalse(out.contains(s.terminalId()), "no member row fragment may leak to a worker: " + out);
assertFalse(out.contains("term_owner"), "no member owner fragment may leak to a worker: " + out);
assertFalse(out.contains("mac-opus"), "no lead row fragment may leak to a worker: " + out);
assertTrue(out.contains("\"healthCoverage\""), out);
}
/**
* A collaborator may {@code SEND} to a lead, so it must see the {@code leads} array -- the
* only place {@code fleet_whoami} does not already give it a lead's address. It may never
* {@code SEND} to a spawned member, so the {@code members} key must stay absent for it, with
* no fragment of a populated member row leaking either.
*/
@Test
void listShowsLeadsButNotMembersToACollaborator() {
FakeHerdr h = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", "term_owner",
new WorktreeRequest("cb-304", null));
Principal collaborator = Principal.collaborator("ops", "term_collab", 600);
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of("term_lead_x", "mac-opus"), "", FleetMcp.CoordinationSource.none(), collaborator.isPrimary(),
FleetMcp.leadsVisibleTo(collaborator), FleetMcp.membersVisibleTo(collaborator));
String out = textOf(res);
assertTrue(out.contains("\"leads\""), "a collaborator must see the leads array: " + out);
assertTrue(out.contains("mac-opus"), "a collaborator must see the lead's name/address: " + out);
assertFalse(out.contains("\"members\""), "a collaborator must never see the members key: " + out);
assertFalse(out.contains(s.terminalId()), "no member row fragment may leak to a collaborator: " + out);
assertFalse(out.contains("term_owner"), "no member owner fragment may leak to a collaborator: " + out);
assertTrue(out.contains("\"healthCoverage\""), out);
}
@@ -807,16 +871,19 @@ class FleetMcpTest {
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
.withMailbox("mac-opus", LeadChannel.MailboxState.exists("mac-opus", 0, 1));
boolean callerIsPrimary = Principal.architect("lead-designer", "term_design", 400).isPrimary();
Principal architect = Principal.architect("lead-designer", "term_design", 400);
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), callerIsPrimary);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), architect.isPrimary(),
FleetMcp.leadsVisibleTo(architect), FleetMcp.membersVisibleTo(architect));
String out = textOf(res);
assertFalse(out.contains("\"coordinator\""), "an architect must never see the coordinator key either: " + out);
assertTrue(out.contains("\"leads\""), "an architect must still see the leads array: " + out);
assertTrue(out.contains("\"members\""), "an architect must still see the members array: " + out);
}
/**
@@ -841,7 +908,7 @@ class FleetMcpTest {
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true));
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true));
assertTrue(gatedAsPrimary.contains("\"coordinator\""), gatedAsPrimary);
assertTrue(gatedAsPrimary.contains("\"selfId\":\"mac-opus\""), gatedAsPrimary);
@@ -850,6 +917,8 @@ class FleetMcpTest {
assertTrue(gatedAsPrimary.contains("\"heldDurable\""), gatedAsPrimary);
assertTrue(gatedAsPrimary.contains("\"held\""), gatedAsPrimary);
assertTrue(gatedAsPrimary.contains("\"peers\""), gatedAsPrimary);
assertTrue(gatedAsPrimary.contains("\"leads\""), "the primary must still see the leads array: " + gatedAsPrimary);
assertTrue(gatedAsPrimary.contains("\"members\""), "the primary must still see the members array: " + gatedAsPrimary);
}
@Test
@@ -868,7 +937,7 @@ class FleetMcpTest {
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
String out = textOf(res);
assertTrue(out.contains("\"msgId\":\"m1\""), out);
@@ -879,6 +948,83 @@ class FleetMcpTest {
assertTrue(out.contains("x".repeat(80) + "…"), "expected an 80-char preview with an ellipsis: " + out);
}
// --- fleetd #703: fleet_list's collaborators array -------------------------------------------
/** Calls the canonical {@code listFleet} overload directly, so a test can set the collaborators
* payload and its visibility independently of a real {@code Principal} / MCP exchange. */
private static McpSchema.CallToolResult listFleetWithCollaborators(FakeHerdr h,
Map<String, String> collaborators, boolean collaboratorsVisible) {
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
return FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.LoopHealthSource.none(), FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
FleetMcp.LeadSeatSource.none(), new LeadContextGauge(), FleetMcp.LeadConfigDirSource.none(),
Map.of(), "", collaborators, collaboratorsVisible,
FleetMcp.CoordinationSource.none(), false, true, true);
}
/**
* fleetd #703 acceptance A: a visible caller with one configured collaborator gets a
* {@code collaborators} row whose key ({@code CallerResolver.collaborators()}'s
* {@code terminal_id -> name} entry) lands as that row's {@code sessionId}, and {@code leads}/
* {@code members} are unaffected by the new key.
*/
@Test
void listReportsACollaboratorsRowKeyedByTheCollaboratorsTerminalId() {
FakeHerdr h = new FakeHerdr();
String out = textOf(listFleetWithCollaborators(h, Map.of("term_collab", "ops"), true));
assertTrue(out.contains("\"collaborators\":["), out);
assertTrue(out.contains("\"name\":\"ops\""), out);
assertTrue(out.contains("\"sessionId\":\"term_collab\""), out);
assertTrue(out.contains("\"leads\":[]"), out);
assertTrue(out.contains("\"members\":[]"), out);
}
/**
* fleetd #703 acceptance A, control half: with no collaborator configured, the key is absent
* (caller cannot see it) or an empty array (caller can), and {@code leads}/{@code members} are
* unchanged either way.
*/
@Test
void listOmitsOrEmptiesCollaboratorsWhenNoneAreConfigured() {
FakeHerdr h = new FakeHerdr();
String visibleButEmpty = textOf(listFleetWithCollaborators(h, Map.of(), true));
assertTrue(visibleButEmpty.contains("\"collaborators\":[]"), visibleButEmpty);
assertTrue(visibleButEmpty.contains("\"leads\":[]"), visibleButEmpty);
assertTrue(visibleButEmpty.contains("\"members\":[]"), visibleButEmpty);
String notVisible = textOf(listFleetWithCollaborators(h, Map.of(), false));
assertFalse(notVisible.contains("\"collaborators\""), notVisible);
assertTrue(notVisible.contains("\"leads\":[]"), notVisible);
assertTrue(notVisible.contains("\"members\":[]"), notVisible);
}
/**
* fleetd #703 acceptance B: a worker must not see the {@code collaborators} array at all, while
* an architect -- a role that also holds READ, same as a worker -- does see it. Both halves are
* asserted: a test that only checked the worker-hidden half would pass even if the feature were
* never wired up for anyone.
*/
@Test
void listOmitsCollaboratorsForAWorkerAndIncludesThemForAnArchitect() {
FakeHerdr h = new FakeHerdr();
Map<String, String> collaborators = Map.of("term_collab", "ops");
String asWorker = textOf(listFleetWithCollaborators(h, collaborators,
FleetMcp.collaboratorsVisibleTo(Principal.worker("term_w", 1))));
assertFalse(asWorker.contains("\"collaborators\""),
"a worker must not see the collaborators array: " + asWorker);
String asArchitect = textOf(listFleetWithCollaborators(h, collaborators,
FleetMcp.collaboratorsVisibleTo(Principal.architect("design", "term_arch", 2))));
assertTrue(asArchitect.contains("\"collaborators\":["),
"an architect must see the collaborators array: " + asArchitect);
assertTrue(asArchitect.contains("\"sessionId\":\"term_collab\""), asArchitect);
}
/**
* fleetd #421: {@code mailbox.pending} counts only broker-ready messages, so a blocked lead's
* normal, healthy state is {@code "pending": 0} next to a non-empty {@code held[]} — which
@@ -900,7 +1046,7 @@ class FleetMcpTest {
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
String out = textOf(res);
assertTrue(out.contains("\"pending\":0"), out);
@@ -929,7 +1075,7 @@ class FleetMcpTest {
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true, true, true);
String out = textOf(res);
assertTrue(out.contains("\"heldDurable\":false"),
@@ -998,7 +1144,7 @@ class FleetMcpTest {
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "",
new FleetMcp.CoordinationSource(channel, List.of("fleet01-lead", "fleet02-lead", "fleet03-lead")), true);
new FleetMcp.CoordinationSource(channel, List.of("fleet01-lead", "fleet02-lead", "fleet03-lead")), true, true, true);
String out = textOf(res);
assertTrue(out.contains("\"coordId\":\"fleet01-lead\",\"status\":\"exists\",\"pending\":2,\"consumers\":1"), out);
@@ -0,0 +1,280 @@
package dev.ltms.fleet.msg;
import org.junit.jupiter.api.Test;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.List;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
import java.util.stream.Stream;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* Pins that no file under {@code src/main/java} calls the fail-open
* {@link MessageService#poll(String)} overload. That overload skips the ownership check in
* {@code MessageService}'s {@code ownsTicket} entirely, so a caller of it can read any session's
* ticket. Every production caller must go through {@link MessageService#poll(String, String)}
* and pass a {@code callerTerminal} explicitly, even when it is {@code null}.
*
* <p>This reads each file's own source text rather than reflecting on compiled bytecode, because
* the risk is a future one-word edit at a call site, not a missing overload.
*
* <p>The scan below finds a violation by its receiver, {@code messages.poll(}, rather than the
* bare method name, so it does not mistake {@link java.util.Queue#poll()} for a violation. That
* anchor only covers a {@code MessageService} reached through a variable or field named
* {@code messages}, so {@link #everyMessageServiceDeclarationIsNamedMessages} pins the naming
* convention the anchor depends on: a declaration under any other name would be invisible to the
* scan above, and must turn this second check red instead of passing silently.
*/
class MessageServicePollUsageTest {
private static final Path PRODUCTION_SOURCE = Path.of("src/main/java");
@Test
void noProductionFileCallsTheSingleArgumentPollOverload() throws IOException {
List<String> violations = new ArrayList<>();
List<String> twoArgSites = new ArrayList<>();
int filesScanned = scanForPollCalls(PRODUCTION_SOURCE, violations, twoArgSites);
// CONTROL: the scan actually walked files -- a wrong root would otherwise report "no
// violations found" having looked at nothing.
assertTrue(filesScanned > 0, "control failed: the scan under " + PRODUCTION_SOURCE
+ " visited zero .java files -- the path is wrong, so the absence of violations "
+ "below proves nothing");
assertTrue(violations.isEmpty(), "found a call to the fail-open MessageService.poll(String) "
+ "overload, which skips the ownership check entirely -- pass a callerTerminal "
+ "explicitly (even if null) through poll(String, String) instead: " + violations);
// CONTROL: the arity parser actually finds the two genuine two-argument call sites (the
// MCP handler in FleetMcp and the REST handler in FleetApp). If this drops, the parser
// itself is broken, not the production code -- a broken parser (or a scan root that
// reaches no real source) must fail loudly here rather than pass vacuously above.
assertEquals(2, twoArgSites.size(), "control failed: expected exactly the two known "
+ "two-argument messages.poll(...) call sites, found: " + twoArgSites);
assertTrue(twoArgSites.stream().anyMatch(s -> s.contains("FleetMcp.java")),
"control failed: did not find the FleetMcp.java messages.poll(ticket, callerTerminal) "
+ "site among: " + twoArgSites);
assertTrue(twoArgSites.stream().anyMatch(s -> s.contains("FleetApp.java")),
"control failed: did not find the FleetApp.java messages.poll(...) site among: "
+ twoArgSites);
}
/**
* The {@code messages.poll(} anchor above only sees a {@code MessageService} reached through
* a variable, field, or parameter named {@code messages}. This asserts that every such
* declaration under {@code src/main/java} uses that name, so a differently named declaration
* -- invisible to the scan above -- fails loudly here instead of letting that scan pass on a
* call site it never looked at.
*/
@Test
void everyMessageServiceDeclarationIsNamedMessages() throws IOException {
List<String> names = new ArrayList<>();
int filesScanned = scanForDeclarationNames(PRODUCTION_SOURCE, names);
// CONTROL: the scan actually walked files -- a wrong root would otherwise report "every
// declaration is named messages" having looked at nothing.
assertTrue(filesScanned > 0, "control failed: the scan under " + PRODUCTION_SOURCE
+ " visited zero .java files -- the path is wrong, so the result below proves nothing");
// CONTROL: the declaration pattern actually finds real declarations. Zero means the
// pattern is broken, not that every MessageService variable, field, or parameter vanished.
assertTrue(names.size() > 0, "control failed: found zero MessageService declarations under "
+ PRODUCTION_SOURCE + " -- the declaration pattern is broken, update it before "
+ "trusting the naming check below");
List<String> other = names.stream().filter(n -> !n.equals("messages")).distinct().toList();
assertTrue(other.isEmpty(), "found a MessageService declaration not named \"messages\": "
+ other + " -- the messages.poll( scan above only looks for that name, so a call "
+ "through a differently named variable or field is invisible to it; either rename "
+ "the declaration or widen that scan's anchor to cover it");
}
private static int scanForPollCalls(Path root, List<String> violations, List<String> twoArgSites)
throws IOException {
List<Path> files = javaFiles(root);
for (Path file : files) {
scanFileForPollCalls(file, violations, twoArgSites);
}
return files.size();
}
private static int scanForDeclarationNames(Path root, List<String> names) throws IOException {
List<Path> files = javaFiles(root);
Pattern declaration = Pattern.compile("MessageService\\s+([A-Za-z_][A-Za-z0-9_]*)");
for (Path file : files) {
if (file.getFileName().toString().equals("MessageService.java")) {
continue; // the type's own declaration, not a caller holding a reference to it
}
String stripped = stripComments(Files.readString(file));
Matcher m = declaration.matcher(stripped);
while (m.find()) {
int j = m.end();
while (j < stripped.length() && Character.isWhitespace(stripped.charAt(j))) j++;
if (j < stripped.length() && stripped.charAt(j) == '(') {
continue; // a method named like the convention, e.g. "MessageService messages()"
}
names.add(m.group(1));
}
}
return files.size();
}
private static List<Path> javaFiles(Path root) throws IOException {
try (Stream<Path> paths = Files.walk(root)) {
return paths.filter(p -> p.toString().endsWith(".java")).toList();
}
}
/**
* Replaces {@code //} and {@code /* *}{@code /} comment text with nothing, leaving code,
* string/char literals and line breaks untouched -- so a comment that merely mentions
* {@code MessageService} in prose can never be read as a declaration.
*/
private static String stripComments(String source) {
StringBuilder out = new StringBuilder(source.length());
boolean inString = false;
boolean inChar = false;
int i = 0;
while (i < source.length()) {
char c = source.charAt(i);
if (inString) {
out.append(c);
if (c == '\\' && i + 1 < source.length()) { out.append(source.charAt(i + 1)); i += 2; continue; }
if (c == '"') inString = false;
i++;
continue;
}
if (inChar) {
out.append(c);
if (c == '\\' && i + 1 < source.length()) { out.append(source.charAt(i + 1)); i += 2; continue; }
if (c == '\'') inChar = false;
i++;
continue;
}
if (c == '"') { inString = true; out.append(c); i++; continue; }
if (c == '\'') { inChar = true; out.append(c); i++; continue; }
if (c == '/' && i + 1 < source.length() && source.charAt(i + 1) == '/') {
while (i < source.length() && source.charAt(i) != '\n') i++;
continue; // leaves the newline itself for the next iteration to append
}
if (c == '/' && i + 1 < source.length() && source.charAt(i + 1) == '*') {
i += 2;
while (i < source.length() && !(source.charAt(i) == '*' && i + 1 < source.length()
&& source.charAt(i + 1) == '/')) {
if (source.charAt(i) == '\n') out.append('\n');
i++;
}
i += 2;
continue;
}
out.append(c);
i++;
}
return out.toString();
}
private static void scanFileForPollCalls(Path file, List<String> violations, List<String> twoArgSites)
throws IOException {
String source = Files.readString(file);
String needle = "messages.poll(";
int from = 0;
int idx;
while ((idx = source.indexOf(needle, from)) >= 0) {
int argsStart = idx + needle.length();
String args = extractBalancedArgs(source, argsStart, file, idx);
int closeParenIndex = argsStart + args.length();
from = closeParenIndex + 1;
if (args.isBlank()) {
continue; // MessageService has no zero-argument poll() -- nothing to classify
}
String site = file + ":" + lineOf(source, idx);
if (topLevelCommaCount(args) == 0) {
violations.add(site + " -- messages.poll(" + args.trim() + ")");
} else {
twoArgSites.add(site);
}
}
}
/**
* The text between {@code messages.poll(} and its matching close paren: balanced over nested
* calls, and never split by a paren or comma sitting inside a string or char literal.
*/
private static String extractBalancedArgs(String source, int start, Path file, int callIndex) {
int depth = 1;
boolean inString = false;
boolean inChar = false;
int i = start;
while (i < source.length()) {
char c = source.charAt(i);
if (inString) {
if (c == '\\') { i += 2; continue; }
if (c == '"') inString = false;
} else if (inChar) {
if (c == '\\') { i += 2; continue; }
if (c == '\'') inChar = false;
} else if (c == '"') {
inString = true;
} else if (c == '\'') {
inChar = true;
} else if (c == '(') {
depth++;
} else if (c == ')') {
depth--;
if (depth == 0) return source.substring(start, i);
}
i++;
}
throw new IllegalStateException(
"unbalanced parentheses scanning " + file + ":" + lineOf(source, callIndex));
}
/**
* Commas at paren/bracket/brace depth zero, skipping string and char literals -- the argument
* separators a human reader would see, not every comma character in the text.
*/
private static int topLevelCommaCount(String args) {
int depth = 0;
int commas = 0;
boolean inString = false;
boolean inChar = false;
int i = 0;
while (i < args.length()) {
char c = args.charAt(i);
if (inString) {
if (c == '\\') { i += 2; continue; }
if (c == '"') inString = false;
} else if (inChar) {
if (c == '\\') { i += 2; continue; }
if (c == '\'') inChar = false;
} else if (c == '"') {
inString = true;
} else if (c == '\'') {
inChar = true;
} else if (c == '(' || c == '[' || c == '{') {
depth++;
} else if (c == ')' || c == ']' || c == '}') {
depth--;
} else if (c == ',' && depth == 0) {
commas++;
}
i++;
}
return commas;
}
private static int lineOf(String source, int index) {
int line = 1;
for (int i = 0; i < index; i++) {
if (source.charAt(i) == '\n') line++;
}
return line;
}
}
@@ -865,11 +865,82 @@ class MessageServiceTest {
+ "follow-up send would come back BUSY instead of timing out on its own work");
}
/**
* Drive an async send on {@code T} to a resolved reply, polling as {@code owner} until the
* ticket reports {@link MessageService.Phase#DONE} (or the 2s deadline runs out). {@code owner}
* must be a terminal this ticket's creator check actually accepts, or this loops to the
* deadline and returns a non-{@code DONE} view.
*/
private MessageService.TaskView driveAsyncTicketToDone(String ticket, String owner) throws InterruptedException {
MessageService.TaskView view = null;
long deadline = System.currentTimeMillis() + 2000;
while (view == null || view.phase() != MessageService.Phase.DONE) {
if (System.currentTimeMillis() >= deadline) break;
view = messages.poll(ticket, owner);
//noinspection BusyWait
Thread.sleep(5);
}
return view;
}
@Test
void pollReturnsNullForAnUnknownTicket() {
assertNull(messages.poll("task-999999"), "a ticket that was never minted is unknown");
}
// --- fleetd #705: a ticket's creator terminal gates who may poll it -------------------------
@Test
void pollByAnotherTerminalIsRefused() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, "term_a");
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker works
assertTrue(rendezvous.resolve(T, "secret async result"), "a reply resolves the async send");
MessageService.TaskView owner = driveAsyncTicketToDone(ticket, "term_a");
assertNotNull(owner, "the creator must still be able to read its own ticket");
assertEquals(MessageService.Phase.DONE, owner.phase());
MessageService.TaskView refused = messages.poll(ticket, "term_b");
assertNotNull(refused, "a different terminal gets a refusal, not silence");
assertNotEquals(MessageService.Phase.DONE, refused.phase(),
"a different terminal must never see the ticket as DONE");
assertNull(refused.reply(), "a refusal must never carry the reply text");
assertFalse(String.valueOf(refused).contains("secret async result"),
"the reply text must not appear anywhere in the refused view");
}
@Test
void unnamedPrimaryStillReadsAnyTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, "term_lead");
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker works
assertTrue(rendezvous.resolve(T, "primary-visible result"), "a reply resolves the async send");
// callerTerminal == null is the unnamed primary (resolved by token or loopback trust, with
// no herdr pane) — it must read a ticket a terminal-bearing lead created.
MessageService.TaskView view = driveAsyncTicketToDone(ticket, null);
assertNotNull(view, "the unnamed primary must be able to read any ticket");
assertEquals(MessageService.Phase.DONE, view.phase());
assertEquals("primary-visible result", view.reply());
}
@Test
void creatorReadsItsOwnTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, "term_creator");
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker works
assertTrue(rendezvous.resolve(T, "own result"), "a reply resolves the async send");
MessageService.TaskView view = driveAsyncTicketToDone(ticket, "term_creator");
assertNotNull(view, "the ticket's own creator must be able to read it");
assertEquals(MessageService.Phase.DONE, view.phase());
assertEquals("own result", view.reply());
}
@Test
void pollReportsACompletedTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task");
@@ -139,6 +139,154 @@ class FleetAppAuthTest {
assertThrows(IllegalArgumentException.class, () -> FleetApp.routeAction("GET /healthz"));
}
// --- GET /tasks/{ticket} must not be the no-check overload ----------------------------------
/**
* {@code GET /tasks/{ticket}} must resolve its caller the same way {@code allow(...)} does
* and thread that terminal into {@link MessageService#poll(String, String)}, not the
* no-check overload that ignores who is asking.
*/
@Test
void theTaskStatusRouteActuallyThreadsTheCallersTerminalIntoPoll() throws Exception {
String source = Files.readString(REST_SOURCE);
int start = source.indexOf("private void taskStatus(Context ctx) {");
assertTrue(start >= 0, "could not find taskStatus in " + REST_SOURCE
+ " -- the scrape has stopped matching, fix the anchor before trusting this test");
int end = source.indexOf("private static void herdrError(Context ctx, HerdrException e) {", start);
assertTrue(end > start, "could not find the method declared after taskStatus to bound the scrape");
String handlerBlock = source.substring(start, end);
// CONTROL: the block we scraped really does call messages.poll(...) -- if this fails, the
// anchors above moved and the assertions below would otherwise pass on nothing.
assertTrue(handlerBlock.contains("messages.poll("),
"control failed: the scraped taskStatus block contains no messages.poll( call at all "
+ "-- the anchors have drifted, this test is not testing what it claims to");
assertTrue(handlerBlock.contains("caller.terminal()"),
"the taskStatus route must thread the resolved caller's terminal into messages.poll(...), "
+ "not the no-check overload -- block: " + handlerBlock);
assertTrue(handlerBlock.contains("ctx.attribute(CALLER)"),
"the taskStatus route must resolve its caller the same way allow(...) does, not via a "
+ "second, separate resolution path -- block: " + handlerBlock);
}
/**
* {@code GET /tasks/{ticket}} refuses a worker whose terminal did not create the ticket,
* while the creating worker and the unnamed primary both still read it. The ticket is minted
* directly on the shared {@link MessageService}, the same way {@code MessageServiceTest}
* drives {@link MessageService#poll(String, String)}, so this exercises only the REST poll
* route's own handling of the ownership already recorded on the ticket.
*/
@Test
void restPollRefusesADifferentWorkerButAllowsTheCreatorAndTheUnnamedPrimary() throws Exception {
FakeHerdr herdr = new FakeHerdr();
AgentControl agents = new AgentControl(herdr);
Injector injector = new Injector(agents);
MessageService messages = new MessageService(agents, injector, new Rendezvous());
Javalin creatorApp = startOnSharedService(messages, herdr, FakeHerdr.WORKER_PID); // -> term_a
Javalin otherWorkerApp = startOnSharedService(messages, herdr, 9001L); // -> term_shell
Javalin primaryApp = startOnSharedService(messages, herdr, 999_999L); // no pane -> primary
try {
String ticket = messages.sendAsync("term_a", "long task", null, "term_a");
HttpResponse<String> refused = send(otherWorkerApp.port(), "GET", "/tasks/" + ticket, null, null);
assertEquals(200, refused.statusCode());
assertTrue(refused.body().contains("forbidden"),
"a different worker's terminal must be refused, not shown the ticket: " + refused.body());
assertFalse(refused.body().contains("\"reply\""),
"a refusal must never carry reply text: " + refused.body());
HttpResponse<String> own = send(creatorApp.port(), "GET", "/tasks/" + ticket, null, null);
assertEquals(200, own.statusCode());
assertFalse(own.body().contains("forbidden"),
"the creating worker must read its own ticket: " + own.body());
HttpResponse<String> primary = send(primaryApp.port(), "GET", "/tasks/" + ticket, null, null);
assertEquals(200, primary.statusCode());
assertFalse(primary.body().contains("forbidden"),
"the unnamed primary must read any ticket: " + primary.body());
} finally {
creatorApp.stop();
otherWorkerApp.stop();
primaryApp.stop();
}
}
/**
* {@code POST /sessions/{id}/message} with {@code wait:false} must record the creating
* caller's own terminal on the ticket it returns, so that caller can still poll its own
* ticket over REST, while a different terminal is refused.
*/
@Test
void restSendAsyncRecordsTheCreatingCallersTerminalSoItCanStillPollItsOwnTicket() throws Exception {
FakeHerdr herdr = new FakeHerdr();
AgentControl agents = new AgentControl(herdr);
Injector injector = new Injector(agents);
MessageService messages = new MessageService(agents, injector, new Rendezvous());
Javalin leadApp = startOnSharedService(messages, herdr, FakeHerdr.WORKER_PID, Map.of("term_a", "lead-x"));
Javalin otherWorkerApp = startOnSharedService(messages, herdr, 9001L); // -> term_shell
try {
ObjectMapper mapper = new ObjectMapper();
HttpResponse<String> created = send(leadApp.port(), "POST", "/sessions/term_a/message",
"{\"content\":\"long task\",\"wait\":false}", null);
assertEquals(202, created.statusCode(), created.body());
String ticket = mapper.readTree(created.body()).path("ticket").asText(null);
assertNotNull(ticket, "the accepted response carried no ticket: " + created.body());
HttpResponse<String> own = send(leadApp.port(), "GET", "/tasks/" + ticket, null, null);
assertEquals(200, own.statusCode());
assertFalse(own.body().contains("forbidden"),
"the session that created the ticket over REST must be able to poll it: " + own.body());
HttpResponse<String> refused = send(otherWorkerApp.port(), "GET", "/tasks/" + ticket, null, null);
assertEquals(200, refused.statusCode());
assertTrue(refused.body().contains("forbidden: this ticket was created by a different session"),
"a different terminal must still be refused with the ownership detail, not some "
+ "other rejection: " + refused.body());
} finally {
leadApp.stop();
otherWorkerApp.stop();
}
}
/**
* As {@link #start}, but shares {@code messages} and {@code herdr} across several app
* instances bound to different pids, each returned as its own started {@link Javalin} rather
* than through the shared {@code app} field, so several differently-resolved callers can
* poll the same ticket.
*/
private Javalin startOnSharedService(MessageService messages, FakeHerdr herdr, long pid) {
return startOnSharedService(messages, herdr, pid, Map.of());
}
/**
* As {@link #startOnSharedService(MessageService, FakeHerdr, long)}, but {@code leadTerminals}
* resolves the given pid's terminal to a named lead (a caller with SEND permission) instead of
* a plain worker, for a test that needs a terminal-bearing caller able to create a ticket.
*/
private Javalin startOnSharedService(MessageService messages, FakeHerdr herdr, long pid,
Map<String, String> leadTerminals) {
FleetConfig.Profile wcfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
AgentControl agents = new AgentControl(herdr);
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(
agents, new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
Map.of(wcfg.profile(), wcfg), wcfg.profile(),
k -> "FLEETD_WORKER_TOKEN".equals(k) ? "tok-abc" : null);
SessionManager sessions = new SessionManager(workers, new FakeWorktrees());
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> pid);
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null,
() -> leadTerminals, new MemberRegistry(null));
Metrics appMetrics = FleetMetrics.create(sessions, new dev.ltms.fleet.msg.InMemoryReplyInbox());
return new FleetApp(herdr, workers, sessions, messages, sessions.asPresence(), null,
callers, appMetrics).build().start("127.0.0.1", 0);
}
/**
* fleetd #669 Unit A: {@code POST /sessions/{id}/message} is two call shapes behind one route,
* mirroring {@code fleet_send}'s MCP-side split into {@link Authz.Action#SEND} and {@link