Compare commits

..

38 Commits

Author SHA1 Message Date
Dai Ha bf2d940c26 fleetd #669: fix fleetd.example.yaml's false scan-exclusion claims; cover the shipped shape
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 2m10s
fleetd.example.yaml claimed member spaces are excluded from the lead-tab scan and that
a lead's workspace default is "leads" and must not match a member workspace. Both are
false: production always passes an empty excludedWorkspaceLabels set (CB-558), and a
lead's workspace defaults to the same shared "fleet" space the members use. Rewrite
both paragraphs to name the defenses that actually exist: the startup tabPrefix refusal
and CallerResolver's roster-first precedence.

Add the test nobody wrote: a lead is still discovered when its workspace label equals
the member workspace, with an empty exclusion set. Correct
aTabInAWorkerSpaceIsNeverALeadEvenWhenItsLabelMatches's javadoc, which overclaimed that
the excludedWorkspaceLabels parameter is the live production guard. Drop the stale line
number from FleetdAssemblyLeadTabScannerExclusionTest's assertion message and javadoc.
2026-10-04 06:31:20 +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
Dai Ha 7e838ba8b9 fleetd #669 Unit E: route a collaborator's pane to the lead herdr daemon
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 45s
CI / build (pull_request) Failing after 1m57s
HerdrRouter.agentsFor picked the member daemon for any terminal the lead
predicate did not recognize, so a configured collaborator's pane (opened
by a person, exactly like a lead's) was routed to the member herdr
daemon instead of the lead one.

FleetdAssembly now combines the leads map and the collaborator-terminals
map into the predicate it hands HerdrRouter. HerdrRouter's isLead field
and constructor parameter are renamed to routeToLead, with its javadoc
naming the real contract: true for any terminal whose pane lives in the
lead daemon, lead or collaborator.
2026-10-04 02:11:17 +02:00
Dai Ha c3e3554bde CLAUDE.md: the collaborator role, after #669 Unit D
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 52s
CI / build (push) Failing after 1m50s
The block is the instruction surface this repo ships, and Unit D made
four of its statements false. A session can now resolve as
collaborator, read "Which role am I?", and find no section telling it
what it may do.

Each claim below was read from the merged tree at b5bc5d4.

  - fleet_whoami returns four roles, not three. FleetMcp.whoami returns
    before the lead branch, so a collaborator gets its registry name
    and its own sessionId, and no leader key.
  - The fallback ladder said only that it cannot separate a worker from
    an architect. It also fires for no collaborator at all: every rung
    detects a spawned member, and nothing launched a collaborator. It
    therefore falls to "act as a worker", which is the safe direction
    but leaves it unable to learn what it is without asking.
  - Invariant 3 said send is lead or architect. Authz SEND also allows
    a collaborator when the target passes knownLeadOrCollaborator, so
    it may reach a lead or another collaborator and no spawned member.
  - The invariants heading said "both roles"; there are now four.

New Collaborator section, placed after the member turn contract because
a collaborator must be told that contract is not its own. It owes no
fleet_reply: nothing delegates to it. Its two surprising limits are
written down rather than left to be discovered, and both were measured
in Authz: it cannot reach a worker, and it cannot read a ticket,
because TASK_READ is withheld where READ is not.

The intent table gains a row for messaging a collaborator, and that row
states the gap instead of implying discovery works: fleet_list emits
only "leads" and "members", and CallerResolver.collaborators() has no
caller outside the resolver, so a lead cannot find a collaborator
unless the collaborator tells it its sessionId.

The wiki template is byte-identical again, verified with the project's
own check in the main clone, which printed False before the propagation
and True after.
2026-10-04 01:58:03 +02:00
Dai Ha b5bc5d4ab5 Merge PR #701: fleetd #669 Unit D — the resolver, a live spawned member outranks every tab map
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 49s
CI / build (push) Failing after 1m55s
2026-10-04 01:43:17 +02:00
Dai Ha d83972ede2 fleetd #669 Unit D correction: confirm live slot role before granting a spawned architect
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Failing after 1m44s
The new spawned-member roster step returned Principal.architect from the roster's
own MemberRole.ARCHITECT alone, never reconfirming against memberSlotRoles. A slot
revoked after the bind kept granting ARCHITECT to the already-bound session,
regressing fleetd #424's "config governs what a bound slot still grants" half. The
roster now only answers that the pane is a live spawned member; a confirmed live
slot role still decides whether that grants ARCHITECT, falling through to WORKER
otherwise — mirroring the existing architect-slot step a few lines below.

Also corrects LeadTabScanner.buildTabIndex's javadoc: a colliding lead/collaborator
key resolves to the lead because the lead entry is put last (overwriting), not
because it is put first.
2026-10-04 01:37:13 +02:00
Dai Ha 02c6909546 fleetd #669 Unit D: a live spawned member outranks every tab map
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 56s
CI / build (pull_request) Failing after 1m43s
CallerResolver.resolve() now checks the spawned-member roster first, ahead
of every tab map, so a live member's own role wins over a lead or
collaborator tab naming the same terminal (closes #661 at the resolver
level). Adds collaborator resolution (Role.COLLABORATOR) and a real
knownLeadOrCollaborator() classifier, wired into both FleetMcp.denyFor and
FleetApp.allow in place of the inert NO_KNOWN_LEAD_OR_COLLABORATOR stand-in.

LeadTabScanner is generalized to match lead and collaborator tabs in one
pass, keeping every existing liveness/caching/grace-scan property for both
kinds. FleetdAssembly wires the collaborator tab map and a roster-backed
spawned-member lookup; a collaborator-only fleet (no leaders configured)
falls back to a 10s scan interval, matching FleetConfig.Leader's own default.
2026-10-04 01:13:25 +02:00
Dai Ha b92a669ddc Merge PR #699: fleetd #669 Unit C — the COLLABORATOR role and its authorization row
CI / shell-tests (push) Failing after 11s
CI / contract (push) Successful in 52s
CI / build (push) Failing after 1m49s
2026-10-04 00:34:13 +02:00
Dai Ha bb29b001e4 fleetd #669 Unit C: add the COLLABORATOR role and its authorization row
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 53s
CI / build (pull_request) Failing after 1m48s
Adds Role.COLLABORATOR, Principal.collaborator(), and the Authz.permits grant:
READ/METRICS open, REPLY/ASK via ownership, SEND limited to a configured lead
or collaborator via a new classifier parameter, everything else denied.
isSpawnedMember() stays WORKER || ARCHITECT. fleet_whoami reports role and
collaborator name with no leader key.

Per a same-day ticket correction, the existing 3-argument Authz.permits is
kept (fail-closed default via the new NO_KNOWN_LEAD_OR_COLLABORATOR
classifier) rather than deleted, so the 47 pre-existing call sites in
AuthzTest/CallerResolverTest are untouched; both production gates
(FleetMcp#denyFor, FleetApp#allow via the new permitsFor seam) call the
4-argument form explicitly with the shared constant.

Mutation-verified: removing the classifier conjunct from the SEND arm kills
exactly 3 tests (AuthzTest, FleetMcpAuthzTest, FleetAppAuthTest), one per
gate, no cascade.

Full mvn clean install at this commit: 1974 tests, 0 failures (baseline at
f0ff252 was 1964; +10 are the new collaborator-matrix tests). Independently
counted from target/surefire-reports/*.xml after rm -rf: 173 report files,
aggregate tests=1974 failures=0 errors=0.
2026-10-04 00:22:27 +02:00
Dai Ha f0ff25221e Merge PR #697: fleetd #669 Unit B — recognise-only fleet.collaborators config block
CI / shell-tests (push) Failing after 14s
CI / contract (push) Successful in 55s
CI / build (push) Failing after 2m2s
2026-10-04 00:02:10 +02:00
Dai Ha 780cb342ad fleetd #669: pin the fleet-wide tabLabel-vs-collaborator-tab check
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Failing after 1m58s
Adds a test driving validateLeadTabPrefixes()'s fleet-wide
fleet.tabLabel branch directly against a collaborator tab, with no
profile override involved. Every existing collaborator test exercised
only the profile-override branch below it, leaving the fleet-wide
branch without a pinning test.
2026-10-03 23:55:47 +02:00
Dai Ha 5c563f02c8 fleetd #669: comments-only fixes — drop history/narration and premature behaviour claims
CI / shell-tests (pull_request) Failing after 6s
CI / contract (pull_request) Successful in 53s
CI / build (pull_request) Failing after 1m54s
The duplicate-collaborator test's javadoc narrated the before/after and a
timestamp; replaced with the present-tense behaviour it protects, per the
project's code-comment rule.

The Collaborator record javadoc and fleetd.example.yaml both claimed the
daemon "learns to address" a collaborator tab. Nothing does yet — that is
Units C-E. Dropped the claim from the record's contract and the example
now says the block is parsed and validated today, nothing more.

Also drops a confidence marker and a cross-class pointer from the
Collaborator record javadoc. No logic, assertion, or file outside this
PR's existing three changed.
2026-10-03 23:47:16 +02:00
Dai Ha ad593c9bb9 fleetd #669: add collaborators to the duplicate-slot-key guard
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 52s
CI / build (pull_request) Failing after 2m3s
FLEET_POOL_KEYS was missing "collaborators", so a duplicated
fleet.collaborators.<name> key took the skipValue branch in
rejectDuplicateSlotsInPools and silently collapsed last-wins, unlike
the other five pools. Measured: before this change, a test loading a
config with a duplicated collaborator name observed no exception;
after, it fails exactly like the existing architect-pool case.

Also updates the javadoc's "five pools" count to six.
2026-10-03 23:41:13 +02:00
Dai Ha 736fd9cf4b Merge PR #696: fleetd #692 — bound the unbounded userinfo mask in redact()
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 49s
CI / build (push) Failing after 1m46s
2026-10-03 23:33:37 +02:00
Dai Ha 05244a82b3 fleetd #669 Unit B: recognise-only fleet.collaborators config block
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 58s
CI / build (pull_request) Failing after 1m46s
Adds fleet.collaborators.<name>.tab (a Collaborator record with only a tab
field — no profile, no instances, no kind; never auto-launched) and widens
the two startup refusals that stop it being a privilege hole:

- validatePanePlacementAgainstLeadTabs() now fires on either a lead tab or
  a collaborator tab, and no longer early-returns when fleet.leaders is
  empty but fleet.collaborators is not.
- validateLeadTabPrefixes() now also refuses a tabLabel template able to
  render as a collaborator tab, two collaborators sharing one exact tab,
  and a collaborator tab equal to a lead tab (case-insensitively).

validateMembers() refuses a collaborator with no (or blank) tab, next to
the existing leader-must-name-a-tab check.

Does not touch Authz, CallerResolver, LeadTabScanner, MemberRole or
FleetMcp — those are later units per the ticket's plan.
2026-10-03 23:33:28 +02:00
Dai Ha ef4996a01e fleetd #692: bound the unbounded userinfo mask at redact()'s line 273
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 2m9s
redact()'s final fallthrough still used the unbounded [^@]* class #638
removed from the verdict-line masker, so a diff line with a URL that has
no userinfo plus a later @ elsewhere had its text between them silently
deleted. Share one bounded implementation, mask_url_userinfo(), between
redact() and mask_verdict_userinfo() so the bound lives in one place.
2026-10-03 23:27:15 +02:00
Dai Ha edbd8d816a Merge PR #694: fleetd #689 — check SEND before reading the request body in sendMessage
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 44s
CI / build (push) Failing after 1m53s
2026-10-03 22:58:17 +02:00
Dai Ha 9425a9b696 fleetd #689: pin the answerGatePasses call site via the audit trail
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 48s
CI / build (pull_request) Failing after 2m3s
A unit test on the extracted helper proves the helper, not the call
site in sendMessage. allow() logs an AuditLog.allowed() entry for
every granted non-READ/METRICS/TASK_READ action, so a granted turnId
request must log both SEND and ANSWER, and a granted plain request
must log SEND alone. Verified this goes red when the call site is
deleted from sendMessage, and restores to a clean diff.
2026-10-03 22:56:04 +02:00
Dai Ha d0688c8a60 Merge PR #695: fleetd #693 — pin the lead-tab guard's case-insensitivity; fleetd #676 — drop the last stale validator count
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 55s
CI / build (push) Failing after 1m47s
2026-10-03 22:54:04 +02:00
Dai Ha 2c467c2553 fleetd #693, #676: pin lead-tab case-insensitivity, drop stale validator count
CI / shell-tests (pull_request) Failing after 20s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 1m40s
#693: add a FleetConfigTest case for two fleet.leaders tabs differing only
in case — mutating equalsIgnoreCase to equals left this uncaught before.

#676: rename theSweepMechanismIsGenericNotHardcodedToFleetConfigsSixNames
to theSweepRunsEveryValidateMethodOnAnUnrelatedClass, since there are
eight validators now, not six, and the count had drifted into the name
and its class-javadoc {@link}.
2026-10-03 22:51:23 +02:00
Dai Ha c6430d8edd fleetd #689: check SEND before reading the request body in sendMessage
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Failing after 1m56s
Authorize twice: the coarse SEND grant first, with no body read, then
parse the body, then check ANSWER too when turnId is present. Restores
the pre-#687 ordering (no attacker-controlled body parse before the
gate) while keeping the SEND/ANSWER split #687 introduced.
2026-10-03 22:47:25 +02:00
Dai Ha bfee23acc3 Merge PR #691: fleetd #677 — refuse two leads sharing one exact tab (supersedes #686)
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 52s
CI / build (push) Failing after 1m37s
2026-10-03 22:45:23 +02:00
Dai Ha 7dec74f1b4 Merge PR #690: fleetd #638 — stop the verdict userinfo mask from crossing / or whitespace (supersedes #685)
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 54s
CI / build (push) Failing after 1m58s
2026-10-03 22:39:42 +02:00
Dai Ha 482598e2a6 fleetd #677: refuse two leads that share one exact tab
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 1m17s
CI / build (pull_request) Failing after 2m17s
validateLeadTabPrefixes() compared each lead's tab against the member
tabLabel template, but never against another lead's tab. Two leads
configured with the same exact tab loaded cleanly, even though exact
tab is the only thing lead identity is matched on, so only one of them
could ever be found.

Add an independent pass, case-insensitive, that refuses when two
fleet.leaders entries share one exact tab, with its own exception so
the message stays accurate for this relation.

Decided not to also refuse a lead's tab starting with a sibling's
tabPrefix: tabPrefix plays no role in identity resolution (only the
exact tab does), and the default tabPrefix is "lead:", the same
string used as the conventional lead tab prefix throughout this
codebase's own fixtures (e.g. "lead: opus" / "lead: sol"). Refusing
that case would reject the standard multi-lead setup with no matching
identity hazard.
2026-10-03 22:36:44 +02:00
Dai Ha 5f5d16fbd4 fleetd #638: stop the verdict userinfo mask from crossing / or whitespace
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 45s
CI / build (pull_request) Failing after 1m57s
mask_verdict_userinfo's character class [^@]* crossed a '/' or a space, so a
verdict line with a URI that has no userinfo plus a later @ (e.g. an email
address in diagnostic prose) had everything between them destroyed. Restrict
the class to [^@/[:space:]]* so the match stops at the end of the URI.

Adds the uncovered-direction test: a URI with no userinfo plus a later @ in
the same line must pass through byte for byte. Also drops the design-rationale
sentence from the helper's comment (now in the PR description).
2026-10-03 22:34:02 +02:00
Dai Ha 7df7985a16 Merge remote-tracking branch 'refs/remotes/pr/686' into worker/677-fix-lead-collision-f69073-12 2026-10-03 22:29:29 +02:00
Dai Ha 724b35b46e Merge remote-tracking branch 'refs/remotes/pr/685' into worker/638-fix-overmask-dbb1bf-11 2026-10-03 22:29:14 +02:00
Dai Ha 2eb2d6112e Merge PR #688: fleetd #675 — pin three unpinned FleetdAssembly constructor args
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 57s
CI / build (push) Failing after 1m35s
2026-10-03 22:27:32 +02:00
Dai Ha 7c458e8bf2 Merge PR #687: fleetd #669 Unit A — split SEND and READ into their real call shapes (closes #678) 2026-10-03 22:27:32 +02:00
Dai Ha 133f03e428 Merge PR #684: fleetd #683 — decouple the completion-fallback test's two 5s budgets 2026-10-03 22:27:27 +02:00
Dai Ha 804279175d fleetd #675: pin assembly loop timing defaults
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Failing after 1m42s
2026-10-03 22:12:46 +02:00
Dai Ha cb4a6869b9 fleetd #677: guard exact lead tab collisions
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Failing after 1m56s
2026-10-03 22:10:29 +02:00
Dai Ha 28b45d97e5 fleetd #638: mask userinfo in the daemon verdict line before it prints
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 1m7s
CI / build (pull_request) Failing after 2m4s
Apply the userinfo-only rewrite (sed -E 's#://[^@]*@#://<redacted>@#g') to every
path that prints config-edit.sh's $VERDICT_LINE: restore_and_confirm's two
prints, report_outcome's shared local copy (covering its clean/needs-restart/
refused branches), and check_mode's "last verdict in log" line via
last_verdict_line. Deliberately not routed through redact() — that function's
key:value masking does not match this line's prose, and the rest of the line
(e.g. the pattern quoted in a parse-failure refusal) is the detail an operator
needs to fix the refusal.

No current refusal message echoes a URI, token or password, so this is a guard
against a future validator doing so, not a fix for an observed leak.
2026-10-03 22:06:54 +02:00
Dai Ha 1a397e962e fleetd #683: decouple the completion-fallback test's send budget from its own setup clock
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 56s
CI / build (pull_request) Failing after 1m53s
completionFallbackResolvesATurnThatNeverCalledFleetReply gave messages.send a 5000 ms budget
that started ticking the instant sendAsync() ran, then raced that same clock against
awaitWaiting()'s own 2000 ms deadline plus several onStatus/readText calls before asserting with
send.get(5, SECONDS). On a loaded machine the setup could eat enough of the 5000 ms that the
production call expired first, returning TIMED_OUT_QUEUED instead of COMPLETED_UNREPLIED.

Give this one test's send a 30 000 ms budget (sendAsync(content, timeoutMillis)) so the setup can
never compete with it; send.get(5, SECONDS) stays the one clock the test depends on. A new
regression test injects a deterministic 5500 ms delay in the same spot and proves the budget is no
longer the binding constraint — reverting it to 5000 ms turns that test red with the same
TIMED_OUT_QUEUED mismatch, confirmed by mutation.
2026-10-03 22:06:26 +02:00
32 changed files with 2512 additions and 315 deletions
+39 -10
View File
@@ -27,10 +27,11 @@ through its `fleet_*` tools. No session addresses a peer, a broker, or the netwo
**Every role reads this file.** A member runs in a git worktree of this same repo, so it inherits
this `CLAUDE.md` verbatim, and every rule below is role-conditional.
**Call `fleet_whoami`.** It returns `primary`, `worker`, or `architect`, resolved by the daemon from
your connection — unforgeable, and the same resolution its authorization gate uses. A worker also
carries its `sessionId`, `profile`, `worktree` and `branch`; an architect carries the slot name it
was bound to. Don't infer what you can ask.
**Call `fleet_whoami`.** It returns `primary`, `worker`, `architect`, or `collaborator`, resolved by
the daemon from your connection — unforgeable, and the same resolution its authorization gate uses.
A worker also carries its `sessionId`, `profile`, `worktree` and `branch`; an architect carries the
slot name it was bound to; a collaborator carries its registry name and its own `sessionId`, and
**no `leader` key** — a collaborator is a named peer, not a primary. Don't infer what you can ask.
Only if that call is unavailable, fall back to these — each is one-way, so keep reading until one
fires: the reply charter in your system prompt (*"You are a spawned member in the
@@ -38,12 +39,17 @@ claude-bridge fleet"*) ⇒ **spawned member**; fleet tools prefixed `mcp__fleet_
member** (the launcher fixes that mount name; a primary's mount is named by whoever wrote its
`.mcp.json`, so it varies — and a member spawned before CB-632 still says `mcp__bridge__*`); `ANTHROPIC_BASE_URL` set ⇒ **spawned member** (Claude-model members run
on a clean env, so its *absence* proves nothing). None of these separate a worker from an architect —
only `fleet_whoami` does. **Still unsure ⇒ act as a worker**, the most restricted member role. The
only `fleet_whoami` does. **And none of them fires for a collaborator at all**: every signal in the
ladder detects a *spawned* member, while a collaborator is a tab a person opened by hand, so it has
no charter, no fixed mount name and a normal environment. A collaborator that cannot call
`fleet_whoami` therefore falls to the line below and acts as a worker. That is the safe direction —
it under-privileges, and the refusals are loud — but it means a collaborator has no way to learn
what it is except by asking. **Still unsure ⇒ act as a worker**, the most restricted member role. The
two mistakes are not symmetric: a primary acting as a worker is refused by the authorization gate —
loud and self-correcting — while a member acting as the primary ends its turn with no `fleet_reply`,
and the sender silently receives nothing. Fail toward the recoverable error.
### Invariants — both roles, no exceptions
### Invariants — every role, no exceptions
1. **Never set, export, or forward `ANTHROPIC_BASE_URL`** (or `ANTHROPIC_AUTH_TOKEN`). The primary
stays on subscription; only the bridge puts a member off it, at spawn. Mounting the bridge must
@@ -51,9 +57,10 @@ and the sender silently receives nothing. Fail toward the recoverable error.
2. **The bridge is the only channel.** Text you print in your terminal reaches nobody — the other
side cannot see your screen. An answer that isn't in a `fleet_*` call is silently discarded.
3. **Identity comes from the connection, never an argument.** Workers never pass a target; you
cannot act as another session. Spawn/stop/drain are lead-only; **send is lead or architect**;
reply/ask are only-as-itself — any peer may answer for its own pane, and for no other. A call
outside your role is refused, not queued.
cannot act as another session. Spawn/stop/drain are lead-only; **send is lead, architect, or
collaborator** — and a collaborator may send only to a lead or another collaborator, never to a
spawned member's terminal; reply/ask are only-as-itself — any peer may answer for its own pane,
and for no other. A call outside your role is refused, not queued.
4. **Delivery is status-gated: one message per turn.** Don't busy-poll a peer's terminal and don't
re-send because a call looks slow — the bridge delivers when the peer is `idle`, `blocked` or
`done`. A spawned member must **also** have mounted the bridge MCP: until it has, it is not
@@ -161,6 +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 |
| 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}` |
@@ -234,6 +242,27 @@ simply complies has thrown away the reason there are two of you.
7. **Never merge.** Stage files explicitly — never `git add -A` — and leave alone anything the
project marks as not-yours-to-commit.
### Collaborator — a named peer, not a member
`fleet_whoami` answered `collaborator`, so your pane's tab matches a `fleet.collaborators.<name>.tab`
entry. You are **not** a member: nothing delegates to you, you have no brief, no worktree and no
ticket, and **you owe no `fleet_reply`** — the turn contract above is for a session a lead spawned,
and it does not apply to you. Read it only to understand what the members around you are doing.
What you may do: observe the fleet (`fleet_list`, `fleet_profiles`, `fleet_whoami`), and send to a
lead or to another collaborator. What you may not: spawn, stop or drain anything, roll a lead's
session, answer a member's `fleet_ask`, poll a ticket, or send to a spawned member's terminal. Each
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.
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.
### Where each rule lives (don't duplicate — extend the right layer)
| Layer | Scope | Reaches |
@@ -359,7 +388,7 @@ Before you call any work done, check the row that matches what you touched:
| `ConnectionIdentity` / how a caller is resolved | the `fleet_whoami` paragraph and the fallback ladder |
| `REPLY_CHARTER`, or a launcher's mount/flags | the fallback ladder (`mcp__fleet__*`), and the layering table's top row |
| the injector / status gating | invariant 4 |
| worktree provisioning or the parity overlay | the "both roles read this file" premise — it rests on the worker's worktree being a checkout of this repo |
| worktree provisioning or the parity overlay | the "every role reads this file" premise — it rests on the worker's worktree being a checkout of this repo |
| `.claude/skills/**` | the addendum's skill list, and the "name the playbook" rule |
| a new peer kind (non-Claude adapter) | what that peer can read — anything it must obey belongs in its charter, not in the block |
| **anything an operator can use, configure, or observe** — an MCP tool, a `fleetd.yaml` knob, an endpoint, a visible behaviour | **[Features](wiki/11-Features.md)** — one entry: what it does · the knob that turns it on · **why it exists** · the gotcha |
+26 -8
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,6 +671,22 @@ fleet:
# kind: opencode
# model: openai/gpt-5.6-terra
# 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"
# architects:
# architect-1:
# profile: opus # a strong model, on the operator's subscription
+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);
}
/**
@@ -43,6 +43,7 @@ import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.msg.ReplyInbox;
import dev.ltms.fleet.msg.ReplyPushLoop;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
@@ -50,6 +51,7 @@ import dev.ltms.fleet.power.CaffeinateSleepAssertionMechanism;
import dev.ltms.fleet.power.IdleSleepGuard;
import dev.ltms.fleet.rest.FleetApp;
import dev.ltms.fleet.session.GitWorktrees;
import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.session.SessionReaper;
import io.javalin.Javalin;
@@ -141,8 +143,12 @@ final class FleetdAssembly {
? ports.connectHerdr(Path.of(cfg.memberHerdrSocket()))
: herdr;
AtomicReference<Supplier<Map<String, String>>> leadsRef = new AtomicReference<>(Map::of);
// fleetd #669 Unit E: a collaborator's pane is opened by a person, exactly like a lead's,
// so its terminal must also route to the lead herdr daemon rather than the member one.
AtomicReference<Supplier<Map<String, String>>> collaboratorTerminalsRef = new AtomicReference<>(Map::of);
HerdrRouter router = new HerdrRouter(herdr, memberHerdr,
target -> leadsRef.get().get().containsKey(target));
target -> leadsRef.get().get().containsKey(target)
|| collaboratorTerminalsRef.get().get().containsKey(target));
// CB-402: one adapter per configured peer kind, fronted by a composite router. A profile's
// `kind:` selects its adapter — claude-code (the default) and opencode partition the profile
// set — and the composite dispatches each SPI call to the adapter that owns the profile/pane.
@@ -250,26 +256,47 @@ final class FleetdAssembly {
log.info("leads: {} panes recognised {}", leadTerminals.size(), leadTerminals.values());
}
// CB-531/CB-579: discover leads by the tab labels the operator writes, one scanner per
// configured lead's own exact `tab:` label.
// configured lead's own exact `tab:` label. fleetd #669: the same scan also recognises a
// configured collaborator's tab, so one herdr pass answers both.
final Supplier<Map<String, String>> leads;
final Supplier<Map<String, String>> collaboratorTerminals;
var leaders = cfg.fleet().leaders();
if (!leaders.isEmpty()) {
var collaboratorsConfig = cfg.fleet().collaborators();
if (!leaders.isEmpty() || !collaboratorsConfig.isEmpty()) {
Map<String, String> tabToName = new LinkedHashMap<>();
leaders.forEach((name, leader) -> {
if (leader != null && leader.tab() != null && !leader.tab().isBlank()) {
tabToName.put(leader.tab(), name);
}
});
int scanIntervalSeconds = leaders.values().iterator().next().scanIntervalSeconds();
Map<String, String> collaboratorTabToName = new LinkedHashMap<>();
collaboratorsConfig.forEach((name, collaborator) -> {
if (collaborator != null && collaborator.tab() != null && !collaborator.tab().isBlank()) {
collaboratorTabToName.put(collaborator.tab(), name);
}
});
// A collaborator-only fleet configures no `leaders:` entry to read a scan interval from
// — FleetConfig.Collaborator carries no scanIntervalSeconds of its own. Falling back to
// FleetConfig.Leader's own compact-constructor default keeps a collaborator-only
// deployment on the same rescan cadence as the default lead cadence, instead of
// inventing a second number for the same kind of scan.
int scanIntervalSeconds = leaders.isEmpty()
? 10
: leaders.values().iterator().next().scanIntervalSeconds();
// This must use the lead daemon: scanning member tabs would demote the lead to a worker.
leads = new LeadTabScanner(herdr, tabToName, Set.of(),
LeadTabScanner scanner = new LeadTabScanner(herdr, tabToName, collaboratorTabToName, Set.of(),
TimeUnit.SECONDS.toNanos(scanIntervalSeconds), ports.nanoClock());
log.info("lead scan: tabs {} host a lead (rescan every {}s, shared fleet space)",
tabToName.keySet(), scanIntervalSeconds);
leads = scanner;
collaboratorTerminals = scanner::collaborators;
log.info("lead/collaborator scan: tabs {} host a lead, tabs {} host a collaborator "
+ "(rescan every {}s, shared fleet space)",
tabToName.keySet(), collaboratorTabToName.keySet(), scanIntervalSeconds);
} else {
leads = () -> leadTerminals;
collaboratorTerminals = Map::of;
}
leadsRef.set(leads);
collaboratorTerminalsRef.set(collaboratorTerminals);
// CB-558: start any declared lead that is not already running. After the scanner is built,
// and only when herdr answered — the launcher's whole safety property is that it can count
@@ -341,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.
@@ -452,6 +479,15 @@ final class FleetdAssembly {
ConnectionIdentity identity = new ConnectionIdentity(
new PaneLocator(herdr, memberHerdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
// fleetd #669 Unit D: a live spawned member resolves as its own role, whatever a tab map
// says about the same terminal — read from the roster meant for a hot path (SessionManager
// javadoc), never rosterResolved(), since resolve() runs on every request.
Function<String, MemberRole> spawnedMemberRole = terminal -> sessions.roster().stream()
.filter(s -> terminal.equals(s.terminalId()))
.map(MemberSession::role)
.findFirst()
.orElse(null);
// CB-501: one resolver behind both entry paths. Worker identity still comes from the
// connection and is never token-gated, so enabling token mode cannot lock the fleet out.
final CallerResolver callers;
@@ -461,11 +497,13 @@ final class FleetdAssembly {
throw new IllegalStateException("auth.mode=token but env var " + cfg.auth().tokenEnv()
+ " is unset or empty — export it before starting fleetd");
}
callers = CallerResolver.withLeadsAndMembers(identity, true, token, leads, members);
callers = CallerResolver.withLeadsAndMembers(identity, true, token, leads, members,
spawnedMemberRole, collaboratorTerminals);
log.info("auth: token mode (bearer required for non-worker callers, env {})",
cfg.auth().tokenEnv());
} else {
callers = CallerResolver.withLeadsAndMembers(identity, false, null, leads, members);
callers = CallerResolver.withLeadsAndMembers(identity, false, null, leads, members,
spawnedMemberRole, collaboratorTerminals);
log.info("auth: loopback-trust (any loopback non-worker caller is the primary)");
}
@@ -1,5 +1,7 @@
package dev.ltms.fleet.auth;
import java.util.function.Predicate;
/**
* The authorization table (CB-505), stated once and enforced on both entry paths.
*
@@ -60,58 +62,92 @@ public final class Authz {
}
/**
* Whether {@code caller} may perform {@code action} against {@code targetSession}.
*
* @param targetSession the session id in the request path; only consulted for the worker-scoped
* actions ({@code REPLY}, {@code ASK}), ignored otherwise, may be
* {@code null}
* The fail-closed classifier: answers no for every target, so a collaborator's {@code SEND}
* is refused unless a caller supplies a real one. {@code CallerResolver#knownLeadOrCollaborator()}
* is the real one, read from the same lead and collaborator maps {@code CallerResolver#resolve}
* consults, so a target that classifier calls known is one {@code resolve} would actually
* resolve as a lead or collaborator.
*/
public static final Predicate<String> NO_KNOWN_LEAD_OR_COLLABORATOR = target -> false;
/**
* Convenience form for a caller with no classifier to supply. Fails closed: a collaborator's
* {@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.
*/
public static boolean permits(Principal caller, Action action, String targetSession) {
return permits(caller, action, targetSession, NO_KNOWN_LEAD_OR_COLLABORATOR);
}
/**
* Whether {@code caller} may perform {@code action} against {@code targetSession}.
*
* @param targetSession the session id in the request path; only consulted for the
* worker-scoped actions ({@code REPLY}, {@code ASK}) and for a
* collaborator's {@code SEND}, ignored otherwise, may be
* {@code null}
* @param knownLeadOrCollaborator whether a terminal is a configured lead or collaborator —
* consulted only for a collaborator's {@code SEND}, to confine
* it to another named peer and never a spawned member's
* terminal
*/
public static boolean permits(Principal caller, Action action, String targetSession,
Predicate<String> knownLeadOrCollaborator) {
if (caller == null || caller.isAnonymous()) {
return false; // authenticated as nothing ⇒ authorized for nothing
}
return switch (action) {
// Fleet lifecycle is the primary's alone — spawn, stop, drain. An architect
// deliberately does NOT get these (CB-548), so it cannot tear down or stand up workers
// even though it coordinates them; and a worker driving any of these would be a worker
// escalating into the orchestrator role.
// Fleet lifecycle is the primary's alone — spawn, stop, drain. An architect and a
// collaborator deliberately do NOT get these, so neither can tear down or stand up
// workers even though one of them coordinates them; and a worker driving any of these
// would be a worker escalating into the orchestrator role.
case SPAWN, STOP, DRAIN, HANDOVER -> caller.isPrimary();
// Delivering a turn to a local session is open to the primary and the architect: an
// architect delegates to workers (that is the role's point) but still has no lifecycle
// rights. A worker is excluded — sending would be it escalating.
case SEND -> caller.isPrimary() || caller.isArchitect();
// Delivering a turn to a local session is open to the primary, the architect, and a
// collaborator whose target is itself a configured lead or collaborator: the architect
// delegates to workers (that is the role's point); a collaborator may reach only
// another named peer, never a spawned member's terminal. A worker is excluded —
// sending would be it escalating.
case SEND -> caller.isPrimary() || caller.isArchitect()
|| (caller.isCollaborator() && knownLeadOrCollaborator.test(targetSession));
// Same grant as SEND. Resolving a worker's blocked question is part of delegating to
// it, not a separate capability.
// Resolving a worker's blocked question is part of delegating to it, open to the same
// two roles that may stand up that delegation in the first place. Not a collaborator:
// resuming another session's turn is lifecycle-adjacent, not peer messaging.
case ANSWER -> caller.isPrimary() || caller.isArchitect();
// Same grant as SEND. This leaves the daemon over the coordination broker rather than
// addressing a local session, but the caller who may do one may do the other.
// Leaves the daemon over the coordination broker rather than addressing a local
// session, open to the same two roles as ANSWER. Not a collaborator: it is a
// local-tab peer with no cross-host route.
case COORD_SEND -> caller.isPrimary() || caller.isArchitect();
// The load-bearing rule: a caller acts only as the pane it occupies. CB-532 widened who
// that can be — a lead answering another lead is replying for its OWN terminal, which
// this already permits — while the rule itself is unchanged, and is what stops anyone
// forging a reply for a rendezvous someone else is waiting on. An architect's own pane
// passes through the same check, so it can answer a funnel that delegated to it. An
// unnamed primary (token/loopback, no pane) owns nothing and is still excluded.
// forging a reply for a rendezvous someone else is waiting on. An architect's or a
// collaborator's own pane passes through the same check, so each can answer a funnel
// that delegated to it. An unnamed primary (token/loopback, no pane) owns nothing and
// is still excluded.
case REPLY, ASK -> caller.ownsSession(targetSession);
// READ is roster, profile, and identity observation — fleet_list, fleet_profiles, and
// fleet_whoami — and carries no secrets: no ticket reply, no pending question, and no
// other session's turn state. Those live under TASK_READ. METRICS is the separate
// Prometheus scrape. Both stay open to every authenticated role.
case READ, METRICS -> caller.isPrimary() || caller.isWorker() || caller.isArchitect();
// Prometheus scrape. Both are open to every authenticated role, including a
// collaborator.
case READ, METRICS -> caller.isPrimary() || caller.isWorker() || caller.isArchitect()
|| caller.isCollaborator();
// Ticket polling and session status, open to every authenticated role the same as READ.
// Unlike READ, a holder may poll a ticket it did not create, or read another session's
// pending question and the turnId that answers it.
// Ticket polling and session status, open to every role READ is open to except a
// collaborator: ticket ids are a sequential counter with no owner check, so a holder
// could walk every ticket and read another session's delegation reply.
case TASK_READ -> caller.isPrimary() || caller.isWorker() || caller.isArchitect();
// fleetd #421: reading held lead-to-lead mail is the primary's alone. An architect
// holds READ today (CB-548), so "not primary" must mean not-architect here too — this
// is coordination between leads, not observation of the roster.
// is coordination between leads, not observation of the roster. The same reasoning
// excludes a collaborator.
case COORD_READ -> caller.isPrimary();
};
}
@@ -7,6 +7,7 @@ import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
import java.util.Map;
import java.util.function.Function;
import java.util.function.Predicate;
import java.util.function.Supplier;
/**
@@ -19,6 +20,13 @@ import java.util.function.Supplier;
*
* <p><strong>Resolution order</strong> — connection identity first, token second, nothing third:
* <ol>
* <li>A loopback peer PID that maps to a pane this gateway itself spawned ⇒ that member's own
* role: {@link Role#WORKER} for a dev, hunter, or reviewer; {@link Role#ARCHITECT} for an
* architect, but only while the live slot role still confirms it (fleetd #424 — a slot
* revoked from config demotes an already-bound session on its very next request, so the
* roster's own role is never granted on its word alone). No tab map is consulted — a live
* spawned member's identity comes from the registry that spawned it, never from a label a
* pane could also carry.</li>
* <li>A loopback peer PID that maps to a pane named by {@code leaders:}, by the legacy
* {@code primary.terminal} pin, or by an operator-labelled lead tab (CB-307, CB-530, CB-531)
* ⇒ {@link Role#PRIMARY}, carrying that lead's
@@ -28,8 +36,10 @@ import java.util.function.Supplier;
* so two leads can work as peers rather than one being demoted.</li>
* <li>A loopback peer PID that maps to a pane bound to a CB-548 architect slot ⇒
* {@link Role#ARCHITECT}, carrying the slot name. Just unforgeable as a worker's, and
* resolved from the <em>live</em> terminal→slot binding (never a request argument), before
* the generic worker fallback.</li>
* resolved from the <em>live</em> terminal→slot binding (never a request argument). This is
* the case the previous step does not catch: a binding with no live spawned-member session.</li>
* <li>A loopback peer PID that maps to an operator-labelled collaborator tab ⇒
* {@link Role#COLLABORATOR}, carrying that collaborator's name.</li>
* <li>A loopback peer PID that maps to any other herdr pane ⇒ {@link Role#WORKER}. This is
* unforgeable (the OS reports the PID, herdr owns the PID→pane map) and is honoured
* regardless of auth mode, so enabling auth never breaks the fleet.</li>
@@ -64,6 +74,20 @@ public final class CallerResolver {
private final Supplier<Map<String, String>> architectTerminals;
private final Function<String, MemberRole> memberSlotRoles;
private final Function<String, String> memberSlotNames;
/**
* terminal_id → the role of the live spawned member occupying it, or {@code null} for a
* terminal no spawned member occupies. Consulted first, ahead of every tab map: a live
* spawned member's identity is its own, whatever a tab map says about the same terminal.
* A function rather than the roster itself, so a resolve on the hot path never scans a list —
* the lookup strategy is the caller's to choose.
*/
private final Function<String, MemberRole> spawnedMemberRole;
/**
* terminal_id → collaborator name; empty when none are configured. A supplier for the same
* reason as {@link #leadTerminals}: a collaborator tab recognised after construction (the tab
* scan discovering a newly-labelled tab) takes effect without a restart.
*/
private final Supplier<Map<String, String>> collaboratorTerminals;
/** Loopback-trust resolver: no token required, historical behaviour. Test-only. */
CallerResolver(ConnectionIdentity identity) {
@@ -124,17 +148,42 @@ public final class CallerResolver {
/**
* Live registry form that can confirm a bound slot is an architect slot.
*
* <p>This is the only public construction path. It keeps terminal bindings and slot roles in
* the same {@link MemberRegistry}, so a configured architect can resolve as an architect.
* <p>It keeps terminal bindings and slot roles in the same {@link MemberRegistry}, so a
* configured architect can resolve as an architect. No spawned-member roster or collaborator
* registry is consulted — equivalent to {@link #withLeadsAndMembers(ConnectionIdentity,
* boolean, String, Supplier, MemberRegistry, Function, Supplier)} with both absent. Kept for
* every caller that has neither to offer, so adding them did not churn every construction site.
*/
public static CallerResolver withLeadsAndMembers(ConnectionIdentity identity,
boolean tokenMode, String token,
Supplier<Map<String, String>> leadTerminals,
MemberRegistry members) {
return withLeadsAndMembers(identity, tokenMode, token, leadTerminals, members, null, null);
}
/**
* Live registry form that also resolves a live spawned member to its own role, and a
* configured collaborator tab to {@link Role#COLLABORATOR}.
*
* <p>This is the only public construction path that exercises the full resolution order.
*
* @param spawnedMemberRole terminal_id → the role of the live spawned member occupying
* it, or {@code null} for a terminal no spawned member occupies.
* {@code null} here means no roster is consulted at all (every
* terminal falls through to the tab maps), not that none matches.
* @param collaboratorTerminals terminal_id → collaborator name, live like {@code leadTerminals}
*/
public static CallerResolver withLeadsAndMembers(ConnectionIdentity identity,
boolean tokenMode, String token,
Supplier<Map<String, String>> leadTerminals,
MemberRegistry members,
Function<String, MemberRole> spawnedMemberRole,
Supplier<Map<String, String>> collaboratorTerminals) {
return new CallerResolver(identity, tokenMode, token, leadTerminals,
members == null ? null : members::snapshot,
members == null ? null : members::roleForSlot,
members == null ? null : members::nameForSlot);
members == null ? null : members::nameForSlot,
spawnedMemberRole, collaboratorTerminals);
}
private static Supplier<Map<String, String>> fixed(Map<String, String> leadTerminals) {
@@ -160,6 +209,17 @@ public final class CallerResolver {
Supplier<Map<String, String>> architectTerminals,
Function<String, MemberRole> memberSlotRoles,
Function<String, String> memberSlotNames) {
this(identity, tokenMode, token, leadTerminals, architectTerminals, memberSlotRoles,
memberSlotNames, null, null);
}
private CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token,
Supplier<Map<String, String>> leadTerminals,
Supplier<Map<String, String>> architectTerminals,
Function<String, MemberRole> memberSlotRoles,
Function<String, String> memberSlotNames,
Function<String, MemberRole> spawnedMemberRole,
Supplier<Map<String, String>> collaboratorTerminals) {
if (tokenMode && (token == null || token.isBlank())) {
throw new IllegalArgumentException(
"auth.mode=token requires a non-empty token; check that the env var named by "
@@ -172,6 +232,8 @@ public final class CallerResolver {
this.architectTerminals = architectTerminals == null ? Map::of : architectTerminals;
this.memberSlotRoles = memberSlotRoles == null ? _ -> null : memberSlotRoles;
this.memberSlotNames = memberSlotNames == null ? Function.identity() : memberSlotNames;
this.spawnedMemberRole = spawnedMemberRole == null ? _ -> null : spawnedMemberRole;
this.collaboratorTerminals = collaboratorTerminals == null ? Map::of : collaboratorTerminals;
}
/**
@@ -198,6 +260,27 @@ public final class CallerResolver {
return architectTerminals.get();
}
/**
* The currently-recognised collaborator tabs, {@code terminal_id → name}.
*
* <p>Read from the same supplier {@link #resolve} consults, for the reason given in
* {@link #leads()}. Live for the same reason as {@link #leads()}.
*/
public Map<String, String> collaborators() {
return collaboratorTerminals.get();
}
/**
* Whether {@code target} names a terminal this resolver would resolve as a lead or a
* collaborator — the classifier a collaborator's {@code SEND} is checked against, read from the
* exact maps {@link #resolve} consults so a target that would resolve as a lead or collaborator
* is never the one a collaborator is refused to reach, or the reverse.
*/
public Predicate<String> knownLeadOrCollaborator() {
return target -> leadTerminals.get().containsKey(target)
|| collaboratorTerminals.get().containsKey(target);
}
/**
* Resolve the caller of a request.
*
@@ -208,6 +291,25 @@ public final class CallerResolver {
public Principal resolve(String remoteAddr, int remotePort, String authorizationHeader) {
ConnectionIdentity.Caller c = identity.resolve(remoteAddr, remotePort);
if (c.terminal() != null) {
MemberRole spawnedRole = spawnedMemberRole.apply(c.terminal());
if (spawnedRole != null) {
// A live spawned member occupies this pane. Its identity is its own, whatever a tab
// map says about the same terminal — checked before every tab map, consulting none
// of them, so a tab label can never override a roster entry for the same terminal.
if (spawnedRole == MemberRole.ARCHITECT) {
// The roster only answers THAT this pane is a live spawned member; config still
// decides WHAT that member's slot grants (fleetd #424). A slot revoked after the
// bind must still demote this session on its very next request, so the roster's
// own ARCHITECT role is confirmed against the live slot role, exactly as the
// architect-slot step below confirms a binding with no live member session.
String slot = architectTerminals.get().get(c.terminal());
if (slot != null && memberSlotRoles.apply(slot) == MemberRole.ARCHITECT) {
return Principal.architect(memberSlotNames.apply(slot), c.terminal(), c.pid());
}
return Principal.worker(c.terminal(), c.pid());
}
return Principal.worker(c.terminal(), c.pid());
}
String lead = leadTerminals.get().get(c.terminal());
if (lead != null) {
// The config names this pane as a lead's own. The pane mapping is exactly as
@@ -221,10 +323,17 @@ public final class CallerResolver {
// The config/live binding names this pane as an architect slot's own. Same
// unforgeable pane mapping; the live binding, never a request argument, decides.
// Check the slot role too: this defence in depth prevents a bad lifecycle bind from
// escalating a dev, hunter or reviewer into an architect. Checked before
// the worker fallback.
// escalating a dev, hunter or reviewer into an architect. This is the case the
// spawned-member step above does not catch: a binding with no live member session.
return Principal.architect(memberSlotNames.apply(slot), c.terminal(), c.pid());
}
String collaborator = collaboratorTerminals.get().get(c.terminal());
if (collaborator != null) {
// An operator-labelled collaborator tab, confirmed live by the same scan that
// confirms a lead tab. Checked last among the tab maps so a pane also matching one
// of the above keeps that stronger role.
return Principal.collaborator(collaborator, c.terminal(), c.pid());
}
return Principal.worker(c.terminal(), c.pid()); // unforgeable; never token-gated
}
@@ -11,7 +11,9 @@ package dev.ltms.fleet.auth;
* @param pid the connecting process id, or {@code -1} when not resolvable (audit context)
* @param name for a lead resolved from the CB-530 {@code leaders:} registry, which lead it is;
* for an architect resolved from the CB-548 {@code architects:} registry, which
* slot it occupies; {@code null} for every other caller, including an unnamed primary
* slot it occupies; for a collaborator resolved from the {@code collaborators:}
* registry, which collaborator it is; {@code null} for every other caller,
* including an unnamed primary
*/
public record Principal(Role role, String terminal, long pid, String name) {
@@ -73,6 +75,19 @@ public record Principal(Role role, String terminal, long pid, String name) {
return new Principal(Role.ARCHITECT, terminal, pid, slotName);
}
/**
* A collaborator: a human-opened tab recognised by its exact label in the
* {@code collaborators:} registry.
*
* <p>Carries {@link Role#COLLABORATOR}. {@code name} is reporting only — it lets
* {@code fleet_whoami} say which collaborator is asking. Identity is the {@code terminal}:
* like a worker's it comes from the connection, so {@code ownsSession} works exactly as it
* does for a worker — a collaborator acts as its own pane and no other.
*/
public static Principal collaborator(String name, String terminal, long pid) {
return new Principal(Role.COLLABORATOR, terminal, pid, name);
}
public boolean isPrimary() {
return role == Role.PRIMARY;
}
@@ -81,6 +96,10 @@ public record Principal(Role role, String terminal, long pid, String name) {
return role == Role.ARCHITECT;
}
public boolean isCollaborator() {
return role == Role.COLLABORATOR;
}
public boolean isWorker() {
return role == Role.WORKER;
}
@@ -120,6 +139,7 @@ public record Principal(Role role, String terminal, long pid, String name) {
return switch (role) {
case WORKER -> "worker:" + terminal;
case ARCHITECT -> "architect:" + name;
case COLLABORATOR -> "collaborator:" + name;
case PRIMARY -> name == null ? "primary" : "leader:" + name;
case ANONYMOUS -> "anonymous";
};
@@ -34,6 +34,17 @@ public enum Role {
*/
ARCHITECT,
/**
* A config-declared, human-opened tab recognised by its exact label (the {@code
* fleet.collaborators.<name>.tab} registry). Never spawned — identity comes from the
* connection, never a request argument, exactly like {@link #WORKER} and {@link #ARCHITECT}.
* May {@code SEND} only to a configured lead or collaborator, {@code REPLY}/{@code ASK} only
* as its own pane, and {@code READ}/{@code METRICS}; may not {@code SPAWN}/{@code STOP}/
* {@code DRAIN}/{@code HANDOVER}, poll a ticket ({@code TASK_READ}), or reach the
* coordination broker ({@code COORD_SEND}/{@code COORD_READ}).
*/
COLLABORATOR,
/** Authenticated as nothing. Authorized for nothing but {@code /healthz}. */
ANONYMOUS
}
@@ -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
@@ -1129,11 +1129,8 @@ public record FleetConfig(
* {@code tab} can never be discovered, launched or not
* @param instances how many of this lead should be live (default 1). The daemon
* launches only the shortfall, so a restart adopts rather than doubles
* @param tabPrefix no longer used to find a lead's tab — {@code tab} is matched
* exactly. Its only remaining job is the startup collision guard
* ({@link #validateLeadTabPrefixes()}), which still uses it to refuse
* a worker {@code tabLabel} template that could be misread as a lead.
* Default {@code "lead:"}
* @param tabPrefix lead-tab naming convention checked against member labels. Lead
* identity uses {@code tab}. Default {@code "lead:"}
* @param scanIntervalSeconds how long a tab scan is cached before herdr is asked again; also the
* worst case before a newly-labelled tab is recognised. Default 10
* @param kind which agent runs there ({@code claude}, {@code opencode}, …)
@@ -1179,6 +1176,25 @@ public record FleetConfig(
}
}
/**
* A tab fleetd recognises as a collaborator, keyed by name (fleetd #669).
*
* <p>Recognise-only: there is no {@code profile}, no {@code instances} and no {@code kind}.
* Nothing here ever launches a pane.
*
* <p>{@code tabPrefix} is absent. Identity is matched on the exact {@code tab} alone.
*
* @param tab the exact tab label hosting this collaborator, matched case-insensitively; the
* only field identity depends on. Required — an entry with no {@code tab} can
* never be discovered.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Collaborator(String tab) {
public Collaborator {
tab = (tab == null || tab.isBlank()) ? null : tab.strip();
}
}
/**
* One entry of a {@code fleet:} role pool — a role paired with the backend it runs on.
*
@@ -1216,15 +1232,18 @@ public record FleetConfig(
* is exactly compatible with that. The pool is also what replaced {@code defaultProfile:} — an
* unqualified spawn names a role, and the role's pool supplies the candidates.
*
* @param leaders panes that orchestrate rather than are orchestrated, keyed by lead name
* @param architects profiles the {@code architect} role may run on
* @param developers profiles the {@code dev} role may run on
* @param hunters profiles the {@code hunter} role may run on
* @param reviewers profiles the {@code reviewer} role may run on
* @param charters optional launch-charter text keyed by singular role wire name
* @param tabLabel template for a member tab's label; {@code {role}}, {@code {profile}},
* {@code {model}} and {@code {n}} (a per role+profile counter) are
* substituted. Default {@link #DEFAULT_TAB_LABEL}
* @param leaders panes that orchestrate rather than are orchestrated, keyed by lead name
* @param architects profiles the {@code architect} role may run on
* @param developers profiles the {@code dev} role may run on
* @param hunters profiles the {@code hunter} role may run on
* @param reviewers profiles the {@code reviewer} role may run on
* @param charters optional launch-charter text keyed by singular role wire name
* @param tabLabel template for a member tab's label; {@code {role}}, {@code {profile}},
* {@code {model}} and {@code {n}} (a per role+profile counter) are
* substituted. Default {@link #DEFAULT_TAB_LABEL}
* @param collaborators tabs fleetd recognises as collaborators (fleetd #669), keyed by name.
* Recognise-only, exactly like a {@code profile}-less {@link Leader}:
* nothing here is ever auto-launched.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Fleet(Map<String, Leader> leaders,
@@ -1233,13 +1252,11 @@ public record FleetConfig(
Map<String, Slot> hunters,
Map<String, Slot> reviewers,
Map<String, String> charters,
String tabLabel) {
String tabLabel,
Map<String, Collaborator> collaborators) {
/**
* Role first, so the tab bar reads as the fleet and so the label shares a namespace with a
* lead's {@code tabPrefix}. Because {@code {role}} comes from a closed enum, a generated
* member label can never begin with {@code "lead:"} — the clash that
* {@link #validateLeadTabPrefixes()} used to have to check for is unrepresentable here.
* Role first, so the tab bar identifies the member's fleet role.
*/
public static final String DEFAULT_TAB_LABEL = "{role}: {profile} #{n}";
@@ -1251,26 +1268,30 @@ public record FleetConfig(
reviewers = unmodifiableOrEmpty(reviewers);
charters = unmodifiableOrEmpty(charters);
tabLabel = (tabLabel == null || tabLabel.isBlank()) ? DEFAULT_TAB_LABEL : tabLabel;
collaborators = unmodifiableOrEmpty(collaborators);
}
/**
* A fleet with no configured launch charters — the shape every deployment had before
* CB-566, and what most tests want.
* A fleet with no configured launch charters and no collaborators — the shape every
* deployment had before CB-566, and what most tests want.
*
* <p>Kept deliberately, even though an overload that drops a new field is normally the
* shape to avoid. It is safe here because nothing <em>reads</em> a charter through a
* constructor: the launcher reads {@code fleet.charters()} from the live config. Jackson
* binds the canonical constructor, so this one cannot swallow an operator's YAML.
* {@code collaborators} is dropped the same way and for the same reason: no caller of
* this overload has ever needed to set it, so it defaults to empty here exactly as the
* canonical constructor would default an absent YAML key.
*/
public Fleet(Map<String, Leader> leaders, Map<String, Slot> architects,
Map<String, Slot> developers, Map<String, Slot> reviewers,
Map<String, String> charters, String tabLabel) {
this(leaders, architects, developers, null, reviewers, charters, tabLabel);
this(leaders, architects, developers, null, reviewers, charters, tabLabel, null);
}
public Fleet(Map<String, Leader> leaders, Map<String, Slot> architects,
Map<String, Slot> developers, Map<String, Slot> reviewers, String tabLabel) {
this(leaders, architects, developers, null, reviewers, null, tabLabel);
this(leaders, architects, developers, null, reviewers, null, tabLabel, null);
}
/**
@@ -1916,7 +1937,7 @@ public record FleetConfig(
/** The {@code fleet:} child blocks whose direct children are slot names. */
private static final Set<String> FLEET_POOL_KEYS =
Set.of("leaders", "architects", "developers", "hunters", "reviewers");
Set.of("leaders", "architects", "developers", "hunters", "reviewers", "collaborators");
/**
* Reject a {@code fleet:} role pool whose slot names repeat (CB-548, re-homed by CB-557).
@@ -1926,7 +1947,7 @@ public record FleetConfig(
* daemon would never know. Jackson's YAML parser does not fail on duplicate mapping keys by
* default, so duplicates are caught here, at parse time, before the map is built.
*
* <p>Only the five pools <em>directly under the top-level {@code fleet:}</em> are considered,
* <p>Only the six pools <em>directly under the top-level {@code fleet:}</em> are considered,
* and only their direct child keys (the slot names). A nested field elsewhere, even one also
* named {@code developers:}, is ignored, so parsing of the rest of the config is unaffected.
*
@@ -2685,31 +2706,21 @@ public record FleetConfig(
}
/**
* Reject a lead-scan convention that a worker tab would also satisfy (CB-531).
* Reject a member tab-label template that could render as a configured lead or collaborator
* tab or match a lead-tab naming convention, and reject two {@code fleet.leaders} or
* {@code fleet.collaborators} entries — across either registry — that share one exact tab.
*
* <p>The scan reads a tab label and concludes "a lead lives here". fleetd also <em>writes</em>
* tab labels — every member gets one rendered into its tab. Choose a lead {@code tabPrefix} that
* a member template matches and the daemon starts labelling its own members as leads, promoting
* the entire fleet to {@link dev.ltms.fleet.auth.Role#PRIMARY} with no message and no diff.
* {@link #validatePanePlacementAgainstLeadTabs()} is the check that stops a pane-placed member
* from landing inside a lead's tab in the first place; this check is a second, independent
* guard that catches the hazard even when every profile places members correctly, by refusing
* a label that a scan would still misread as a lead.
* <p>{@code fleet.collaborators} has no {@code tabPrefix}: identity is matched on the exact
* {@code tab} alone, so only the exact-render check applies there, not the prefix check.
*
* <p>CB-557 shrank this check rather than removing it. The default template is
* {@code "{role}: {profile} #{n}"} and {@code {role}} comes from a closed enum, so a
* <em>generated</em> label can no longer collide by construction. What remains checkable is what
* an operator still writes by hand: the {@code fleet.tabLabel} template and any per-profile
* {@code tabLabel} override.
*
* <p>Fatal rather than a warning, unlike {@link #warnUnknownTopLevelKeys}: an unknown key means
* a feature does nothing, while this means a feature does the opposite of what it says.
*
* @throws IllegalStateException when the fleet template or any profile's {@code tabLabel}
* override starts with a configured lead prefix
* @throws IllegalStateException when the fleet template or a profile {@code tabLabel} override
* can render as a configured lead or collaborator tab or match a
* lead-tab prefix, or when two entries — of either registry, or
* one of each — carry the same exact {@code tab}
* (case-insensitively)
*/
public void validateLeadTabPrefixes() {
if (fleet == null || fleet.leaders().isEmpty()) {
if (fleet == null) {
return;
}
List<String> bad = new ArrayList<>();
@@ -2717,51 +2728,173 @@ public record FleetConfig(
if (leader == null) {
return;
}
String tab = leader.tab();
String prefix = leader.tabPrefix();
// The fleet-wide template is checked once per prefix: it labels every member that has no
// override, so one bad template promotes the entire fleet, not one profile.
if (startsWithIgnoreCase(fleet.tabLabel(), prefix)) {
if (templateCanRenderAs(fleet.tabLabel(), tab)) {
bad.add("fleet.tabLabel=\"" + fleet.tabLabel() + "\" can render as the tab of "
+ "lead '" + leadName + "' (\"" + tab + "\")");
} else if (startsWithIgnoreCase(fleet.tabLabel(), prefix)) {
bad.add("fleet.tabLabel=\"" + fleet.tabLabel() + "\" starts with the tabPrefix of "
+ "lead '" + leadName + "' (\"" + prefix + "\")");
}
profiles().entrySet().stream()
.filter(e -> startsWithIgnoreCase(e.getValue().tabLabel(), prefix))
.map(Map.Entry::getKey)
.sorted()
.forEach(p -> bad.add("profile '" + p + "' overrides tabLabel with \""
+ profiles().get(p).tabLabel() + "\", which starts with the tabPrefix of "
+ "lead '" + leadName + "' (\"" + prefix + "\")"));
.forEach(p -> {
String label = profiles().get(p).tabLabel();
if (templateCanRenderAs(label, tab)) {
bad.add("profile '" + p + "' overrides tabLabel with \"" + label
+ "\", which can render as the tab of lead '" + leadName
+ "' (\"" + tab + "\")");
} else if (startsWithIgnoreCase(label, prefix)) {
bad.add("profile '" + p + "' overrides tabLabel with \"" + label
+ "\", which starts with the tabPrefix of lead '" + leadName
+ "' (\"" + prefix + "\")");
}
});
});
if (bad.isEmpty()) {
fleet.collaborators().forEach((collabName, collaborator) -> {
if (collaborator == null) {
return;
}
String tab = collaborator.tab();
if (templateCanRenderAs(fleet.tabLabel(), tab)) {
bad.add("fleet.tabLabel=\"" + fleet.tabLabel() + "\" can render as the tab of "
+ "collaborator '" + collabName + "' (\"" + tab + "\")");
}
profiles().entrySet().stream()
.map(Map.Entry::getKey)
.sorted()
.forEach(p -> {
String label = profiles().get(p).tabLabel();
if (templateCanRenderAs(label, tab)) {
bad.add("profile '" + p + "' overrides tabLabel with \"" + label
+ "\", which can render as the tab of collaborator '"
+ collabName + "' (\"" + tab + "\")");
}
});
});
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.");
}
List<String> collisions = new ArrayList<>();
List<String> leadNames = fleet.leaders().keySet().stream().sorted().toList();
for (int i = 0; i < leadNames.size(); i++) {
String nameA = leadNames.get(i);
Leader a = fleet.leaders().get(nameA);
if (a == null || a.tab() == null || a.tab().isBlank()) {
continue;
}
for (int j = i + 1; j < leadNames.size(); j++) {
String nameB = leadNames.get(j);
Leader b = fleet.leaders().get(nameB);
if (b == null || b.tab() == null || b.tab().isBlank()) {
continue;
}
if (a.tab().equalsIgnoreCase(b.tab())) {
collisions.add("lead '" + nameA + "' and lead '" + nameB + "' both use tab \""
+ a.tab() + "\"");
}
}
}
List<String> collabNames = fleet.collaborators().keySet().stream().sorted().toList();
for (int i = 0; i < collabNames.size(); i++) {
String nameA = collabNames.get(i);
Collaborator a = fleet.collaborators().get(nameA);
if (a == null || a.tab() == null || a.tab().isBlank()) {
continue;
}
for (int j = i + 1; j < collabNames.size(); j++) {
String nameB = collabNames.get(j);
Collaborator b = fleet.collaborators().get(nameB);
if (b == null || b.tab() == null || b.tab().isBlank()) {
continue;
}
if (a.tab().equalsIgnoreCase(b.tab())) {
collisions.add("collaborator '" + nameA + "' and collaborator '" + nameB
+ "' both use tab \"" + a.tab() + "\"");
}
}
}
for (String leadName : leadNames) {
Leader lead = fleet.leaders().get(leadName);
if (lead == null || lead.tab() == null || lead.tab().isBlank()) {
continue;
}
for (String collabName : collabNames) {
Collaborator collaborator = fleet.collaborators().get(collabName);
if (collaborator == null || collaborator.tab() == null
|| collaborator.tab().isBlank()) {
continue;
}
if (lead.tab().equalsIgnoreCase(collaborator.tab())) {
collisions.add("lead '" + leadName + "' and collaborator '" + collabName
+ "' both use tab \"" + lead.tab() + "\"");
}
}
}
if (collisions.isEmpty()) {
return;
}
throw new IllegalStateException("refusing to start: " + String.join("; ", bad)
+ ". Every member labelled that way would be read back as a lead and granted "
+ "spawn/stop/send on the whole fleet. Change one of the two so member tabs and "
+ "lead tabs cannot be confused.");
throw new IllegalStateException("refusing to start: " + String.join("; ", collisions)
+ ". Tab identity is matched exactly, so only one of two entries sharing a tab can "
+ "ever be found — the other is silently unreachable. Give each lead and "
+ "collaborator its own exact tab.");
}
private static boolean templateCanRenderAs(String template, String tab) {
if (template == null || template.isBlank() || tab == null || tab.isBlank()) {
return false;
}
var placeholders = Pattern.compile("\\{(?:role|profile|model|n)}").matcher(template);
StringBuilder expression = new StringBuilder("^");
int literalStart = 0;
while (placeholders.find()) {
expression.append(Pattern.quote(template.substring(literalStart, placeholders.start())));
expression.append(".*");
literalStart = placeholders.end();
}
expression.append(Pattern.quote(template.substring(literalStart))).append("$");
return Pattern.compile(expression.toString(), Pattern.CASE_INSENSITIVE).matcher(tab).matches();
}
/** Case-insensitive prefix test that tolerates a null or blank label. */
private static boolean startsWithIgnoreCase(String label, String prefix) {
if (label == null || prefix == null || prefix.isBlank()) {
return false;
}
String stripped = label.strip();
return stripped.regionMatches(true, 0, prefix, 0, prefix.length());
}
/**
* Reject a profile that places its members by {@code "pane"} while any {@code fleet.leaders}
* entry names a {@code tab}. A pane-placed member lands inside the focused tab rather than its
* own, so it can land inside a lead's own labelled tab. {@link
* dev.ltms.fleet.herdr.LeadTabScanner} identifies a lead 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 the lead
* and granted spawn/stop/send on the whole fleet.
* or {@code fleet.collaborators} entry names a {@code tab}. A pane-placed member lands inside
* 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.
*
* <p>Only a leader with a non-blank {@code tab} is in scope: one with no {@code tab} feeds
* <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.
*
* @throws IllegalStateException when any {@code profiles:} entry is pane-placed while any
* {@code fleet.leaders} entry names a non-blank {@code tab}
* {@code fleet.leaders} or {@code fleet.collaborators} entry
* names a non-blank {@code tab}
*/
public void validatePanePlacementAgainstLeadTabs() {
if (fleet == null || fleet.leaders().isEmpty()) {
if (fleet == null) {
return;
}
boolean anyLeaderHasTab = fleet.leaders().values().stream()
.anyMatch(leader -> leader != null && leader.tab() != null && !leader.tab().isBlank());
if (!anyLeaderHasTab) {
boolean anyCollaboratorHasTab = fleet.collaborators().values().stream()
.anyMatch(c -> c != null && c.tab() != null && !c.tab().isBlank());
if (!anyLeaderHasTab && !anyCollaboratorHasTab) {
return;
}
List<String> bad = new ArrayList<>();
@@ -2774,10 +2907,11 @@ public record FleetConfig(
return;
}
throw new IllegalStateException("refusing to start: profile(s) " + bad
+ " use placement: pane while fleet.leaders names a tab. A pane-placed member can "
+ "land inside a lead's labelled tab and be read back as the lead, granted "
+ "spawn/stop/send on the whole fleet. Set placement: tab for each named profile, "
+ "or remove the tab from every fleet.leaders entry.");
+ " 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 "
+ "each named profile, or remove the tab from every fleet.leaders and "
+ "fleet.collaborators entry.");
}
/**
@@ -2800,15 +2934,6 @@ public record FleetConfig(
}
}
/** Case-insensitive prefix test that tolerates a null/blank label. */
private static boolean startsWithIgnoreCase(String label, String prefix) {
if (label == null || prefix == null || prefix.isBlank()) {
return false;
}
String stripped = label.strip();
return stripped.regionMatches(true, 0, prefix, 0, prefix.length());
}
/**
* Reject a subscription profile whose {@code env:} block tries to reseat the Anthropic binding
* (CB-542).
@@ -2893,8 +3018,14 @@ public record FleetConfig(
* so duplicates are unrepresentable by construction once loaded — and {@link #load(Path)}
* already rejects a duplicated slot name at parse time, before the map collapses.
*
* @throws IllegalStateException when a slot names no profile or an unknown one, or when a lead
* can be neither found nor created, naming the offending entry
* <p>Also rejects a {@code fleet.collaborators} entry with no (or a blank) {@code tab}. A
* {@code profile}-less lead is still useful recognise-only — {@code tab} is the only field
* that matters to it either way. A collaborator carries no other field at all, so a blank
* {@code tab} leaves nothing for the entry to mean.
*
* @throws IllegalStateException when a slot names no profile or an unknown one, when a lead
* can be neither found nor created, or when a collaborator names
* no tab, naming the offending entry
*/
public void validateMembers() {
if (fleet == null) {
@@ -2932,6 +3063,16 @@ public record FleetConfig(
+ "auto-launched, labelled) purely by its tab, so every entry must name one.");
}
});
fleet.collaborators().forEach((name, collaborator) -> {
if (collaborator == null) {
return;
}
if (collaborator.tab() == null || collaborator.tab().isBlank()) {
bad.add("fleet.collaborators." + name + " has no tab: — a collaborator is "
+ "recognised purely by its tab, and carries no other field, so every "
+ "entry must name one.");
}
});
if (!bad.isEmpty()) {
throw new IllegalStateException("refusing to start: " + String.join(" ", bad));
}
@@ -11,12 +11,18 @@ public final class HerdrRouter implements AutoCloseable {
private final AgentControl memberAgents;
private final WorkspaceControl leadSpaces;
private final WorkspaceControl memberSpaces;
private final Predicate<String> isLead;
private final Predicate<String> routeToLead;
public HerdrRouter(HerdrClient lead, HerdrClient member, Predicate<String> isLead) {
/**
* @param routeToLead true for a terminal whose pane lives in the lead herdr daemon — a lead's
* own pane or a configured collaborator's, both opened by a person at a
* terminal rather than spawned, so both are found in the lead daemon rather
* than the member one
*/
public HerdrRouter(HerdrClient lead, HerdrClient member, Predicate<String> routeToLead) {
this.lead = Objects.requireNonNull(lead, "lead");
this.member = member != null ? member : lead;
this.isLead = Objects.requireNonNull(isLead, "isLead");
this.routeToLead = Objects.requireNonNull(routeToLead, "routeToLead");
leadAgents = new AgentControl(this.lead);
memberAgents = this.member == this.lead ? leadAgents : new AgentControl(this.member);
leadSpaces = new WorkspaceControl(this.lead);
@@ -27,7 +33,7 @@ public final class HerdrRouter implements AutoCloseable {
public WorkspaceControl leadSpaces() { return leadSpaces; }
public AgentControl memberAgents() { return memberAgents; }
public WorkspaceControl memberSpaces() { return memberSpaces; }
public AgentControl agentsFor(String targetId) { return isLead.test(targetId) ? leadAgents : memberAgents; }
public AgentControl agentsFor(String targetId) { return routeToLead.test(targetId) ? leadAgents : memberAgents; }
HerdrClient leadClient() { return lead; }
HerdrClient memberClient() { return member; }
@@ -96,13 +96,19 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
private static final Logger log = LoggerFactory.getLogger(LeadTabScanner.class);
/** What a matched tab names: a lead or a collaborator. */
private enum Kind { LEAD, COLLABORATOR }
/** One matched tab's name and what it names. */
private record Entry(String name, Kind kind) {}
private final HerdrClient herdr;
private final Map<String, String> tabToName;
private final Map<String, Entry> tabToEntry;
private final Set<String> excludedWorkspaceLabels;
private final long ttlNanos;
private final LongSupplier clock;
private Map<String, String> cached = Map.of();
private Map<String, Entry> cached = Map.of();
private long scannedAtNanos;
private boolean everScanned;
@@ -126,26 +132,53 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
*/
public LeadTabScanner(HerdrClient herdr, Map<String, String> tabToName,
Set<String> excludedWorkspaceLabels, long ttlNanos, LongSupplier clock) {
this(herdr, tabToName, Map.of(), excludedWorkspaceLabels, ttlNanos, clock);
}
/**
* As {@link #LeadTabScanner(HerdrClient, Map, Set, long, LongSupplier)}, additionally scanning
* for configured collaborator tabs in the same pass.
*
* @param collaboratorTabToName every configured collaborator's exact tab label → its name
* ({@code fleet.collaborators.<name>.tab}), matched the same way as
* {@code tabToName}
*/
public LeadTabScanner(HerdrClient herdr, Map<String, String> tabToName,
Map<String, String> collaboratorTabToName,
Set<String> excludedWorkspaceLabels, long ttlNanos, LongSupplier clock) {
this.herdr = herdr;
this.tabToName = normalize(tabToName);
this.tabToEntry = buildTabIndex(tabToName, collaboratorTabToName);
this.excludedWorkspaceLabels = excludedWorkspaceLabels == null
? Set.of() : Set.copyOf(excludedWorkspaceLabels);
this.ttlNanos = ttlNanos;
this.clock = clock;
}
/** Keys stripped and lower-cased once, so every lookup is a plain map hit. */
private static Map<String, String> normalize(Map<String, String> tabToName) {
if (tabToName == null || tabToName.isEmpty()) {
return Map.of();
/**
* Keys stripped and lower-cased once, so every lookup is a plain map hit. Leads and
* collaborators merge into a single index, so {@link #scan()} matches both kinds in one pass
* over the tab list; a label naming both a lead and a collaborator takes the lead entry —
* leads are put last, so a colliding key's lead entry is the one that overwrites — since a lead
* can already do everything a collaborator can. Config validation already refuses a lead and a
* collaborator sharing one exact tab, so this ordering is defence in depth, not the control.
*/
private static Map<String, Entry> buildTabIndex(Map<String, String> tabToName,
Map<String, String> collaboratorTabToName) {
Map<String, Entry> out = new LinkedHashMap<>();
putNormalized(out, collaboratorTabToName, Kind.COLLABORATOR);
putNormalized(out, tabToName, Kind.LEAD);
return Collections.unmodifiableMap(out);
}
private static void putNormalized(Map<String, Entry> out, Map<String, String> tabToName, Kind kind) {
if (tabToName == null) {
return;
}
Map<String, String> out = new LinkedHashMap<>();
tabToName.forEach((tab, name) -> {
if (tab != null && !tab.isBlank() && name != null && !name.isBlank()) {
out.put(tab.strip().toLowerCase(Locale.ROOT), name);
out.put(tab.strip().toLowerCase(Locale.ROOT), new Entry(name, kind));
}
});
return Collections.unmodifiableMap(out);
}
/**
@@ -156,6 +189,29 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
*/
@Override
public synchronized Map<String, String> get() {
return byKind(refresh(), Kind.LEAD);
}
/**
* The current {@code terminal_id → collaborator name} map, sharing the same scan and cache as
* {@link #get()} — both kinds are matched in one pass, so this never costs a second herdr call.
*/
public synchronized Map<String, String> collaborators() {
return byKind(refresh(), Kind.COLLABORATOR);
}
private static Map<String, String> byKind(Map<String, Entry> entries, Kind kind) {
Map<String, String> out = new LinkedHashMap<>();
entries.forEach((terminal, entry) -> {
if (entry.kind() == kind) {
out.put(terminal, entry.name());
}
});
return Collections.unmodifiableMap(out);
}
/** Rescans if the cache has expired, otherwise returns the cached answer. */
private Map<String, Entry> refresh() {
long now = clock.getAsLong();
if (everScanned && now - scannedAtNanos < ttlNanos) {
return cached;
@@ -165,21 +221,21 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
scannedAtNanos = now;
everScanned = true;
try {
Map<String, String> fresh = scan();
Map<String, Entry> fresh = scan();
if (!fresh.equals(cached)) {
log.info("lead panes: {}", fresh);
log.info("lead/collaborator panes: {}", fresh);
}
cached = fresh;
} catch (HerdrException e) {
log.warn("lead-tab scan failed, keeping the {} lead(s) already known: {}",
log.warn("lead-tab scan failed, keeping the {} entr(y/ies) already known: {}",
cached.size(), e.getMessage());
}
return cached;
}
/** One full pass: labelled tabs → live agents in them → those panes' terminals. */
private Map<String, String> scan() {
Map<String, String> nameByTab = new LinkedHashMap<>();
private Map<String, Entry> scan() {
Map<String, Entry> entryByTab = new LinkedHashMap<>();
for (JsonNode w : herdr.call("workspace.list").path("workspaces")) {
Workspace ws = Workspace.from(w);
if (ws.workspaceId() == null || excludedWorkspaceLabels.contains(ws.label())) {
@@ -187,21 +243,22 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
}
for (JsonNode t : herdr.call("tab.list", Map.of("workspace_id", ws.workspaceId())).path("tabs")) {
Tab tab = Tab.from(t);
String name = leadNameOf(tab.label());
if (name != null && tab.tabId() != null) {
nameByTab.put(tab.tabId(), name);
Entry entry = entryOf(tab.label());
if (entry != null && tab.tabId() != null) {
entryByTab.put(tab.tabId(), entry);
}
}
}
if (nameByTab.isEmpty()) {
if (entryByTab.isEmpty()) {
gracedTerminals = Set.of();
return Map.of();
}
// fleetd #359: a labelled tab is only a lead when herdr also reports a running agent in
// it — the same liveness signal LeadLauncher.countLeads trusts for the identical purpose.
// Without this, a tab left behind by a session that has since died reads as live forever.
// fleetd #359: a labelled tab is only a lead (or collaborator) when herdr also reports a
// running agent in it — the same liveness signal LeadLauncher.countLeads trusts for the
// identical purpose. Without this, a tab left behind by a session that has since died reads
// as live forever.
Set<String> tabsWithAgent = new HashSet<>();
for (JsonNode a : herdr.call("agent.list").path("agents")) {
String tabId = a.path("tab_id").asText(null);
@@ -210,18 +267,18 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
}
}
Map<String, String> byTerminal = new LinkedHashMap<>();
Map<String, Entry> byTerminal = new LinkedHashMap<>();
Set<String> stillGraced = new HashSet<>();
// One pane.list for every tab: panes carry tab_id, so the join is local.
for (JsonNode p : herdr.call("pane.list", Map.of()).path("panes")) {
String tabId = p.path("tab_id").asText(null);
String name = nameByTab.get(tabId);
Entry entry = entryByTab.get(tabId);
String terminal = p.path("terminal_id").asText(null);
if (name == null || terminal == null || terminal.isBlank()) {
if (entry == null || terminal == null || terminal.isBlank()) {
continue;
}
if (tabsWithAgent.contains(tabId)) {
byTerminal.put(terminal, name);
byTerminal.put(terminal, entry);
continue;
}
// No agent reported for this tab, but its tab/pane are still here — this is the
@@ -230,7 +287,7 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
// reported as live; a terminal we never reported live gets none, so the original #359
// fix (a genuinely dead tab is never reported) is unaffected for the common case.
if (cached.containsKey(terminal) && !gracedTerminals.contains(terminal)) {
byTerminal.put(terminal, name);
byTerminal.put(terminal, entry);
stillGraced.add(terminal);
}
}
@@ -239,18 +296,19 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
}
/**
* The lead name a tab label declares, or {@code null} if it names none of the configured leads.
* The entry a tab label declares, or {@code null} if it names neither a configured lead nor a
* configured collaborator.
*
* <p>Exact match (case-insensitive, ends stripped) against {@link #tabToName} — no prefix
* <p>Exact match (case-insensitive, ends stripped) against {@link #tabToEntry} — no prefix
* stripping, so an operator's {@code "lead: something-else"} tab is never mistaken for a
* configured lead just because it shares a prefix. The match strips a trailing
* {@link PendingCloseMarker} first, so a tab {@code LeadLauncher} has flagged as maybe-dead but
* not yet closed keeps resolving normally while that reconcile is pending.
*/
private String leadNameOf(String label) {
private Entry entryOf(String label) {
if (label == null) {
return null;
}
return tabToName.get(PendingCloseMarker.strip(label).toLowerCase(Locale.ROOT));
return tabToEntry.get(PendingCloseMarker.strip(label).toLowerCase(Locale.ROOT));
}
}
@@ -114,6 +114,13 @@ public final class FleetMcp {
* dev.ltms.fleet.herdr.PaneLocator} — that a real assembly wired up. See {@link #identity()}.
*/
private final ConnectionIdentity identity;
/**
* Kept as a field (rather than only captured by the {@code contextExtractor} closure) so
* {@link #denyFor} can read {@link CallerResolver#knownLeadOrCollaborator()} — the classifier a
* collaborator's {@code SEND} is checked against, built from the same lead and collaborator
* maps {@link #identity}-based resolution reads.
*/
private final CallerResolver callers;
private final Metrics metrics; // CB-502: null → auth failures not counted
private final CapacitySource capacity;
private final HealthCoverageSource healthCoverage;
@@ -409,6 +416,7 @@ public final class FleetMcp {
this.authorizationEnforced = Objects.requireNonNull(authorizationMode, "authorizationMode")
== AuthorizationMode.ENFORCED;
this.identity = identity;
this.callers = callers;
this.leadChannel = leadChannel;
this.peers = peers == null ? List.of() : List.copyOf(peers);
this.capacity = capacity;
@@ -689,7 +697,7 @@ public final class FleetMcp {
if (!authorizationEnforced) {
return null; // AuthorizationMode.UNENFORCED: authorization not enforced (fleetd #518)
}
if (Authz.permits(caller, action, target)) {
if (Authz.permits(caller, action, target, callers.knownLeadOrCollaborator())) {
if (action != Authz.Action.READ && action != Authz.Action.TASK_READ) {
AuditLog.allowed(caller, action, target); // reads would drown the trail
}
@@ -1309,6 +1317,18 @@ public final class FleetMcp {
}
return text(json(m));
}
if (caller.isCollaborator()) {
// A collaborator's name is its slot in the collaborators: registry; sessionId is its
// pane so a peer knows where to reach it. No leader key: a collaborator is not a
// primary for authorization, unlike a lead.
if (caller.name() != null) {
m.put("collaborator", caller.name());
}
if (caller.terminal() != null) {
m.put("sessionId", caller.terminal());
}
return text(json(m));
}
if (!caller.isWorker()) {
// CB-530: which lead, once more than one pane is configured as one. `role` deliberately
// still reads "primary" — the fallback ladder in CLAUDE.md keys on it, and a lead IS a
@@ -83,6 +83,18 @@ public final class FleetApp {
};
}
/**
* The second gate for {@code POST /sessions/{id}/message}: checked only when {@code turnId}
* is present and non-blank, against {@link Authz.Action#ANSWER}. A request with no {@code
* turnId} passes this gate unconditionally, without consulting {@code permit} at all, having
* already cleared the coarse {@link Authz.Action#SEND} grant checked ahead of it.
*
* @param permit reports whether the caller holds the named grant
*/
static boolean answerGatePasses(String turnId, Predicate<Authz.Action> permit) {
return turnId == null || turnId.isBlank() || permit.test(Authz.Action.ANSWER);
}
/** Default blocking window for a message; kept under typical HTTP idle timeouts. */
private static final long DEFAULT_MESSAGE_TIMEOUT_MS = 25_000;
private static final long MAX_MESSAGE_TIMEOUT_MS = 120_000;
@@ -254,6 +266,21 @@ public final class FleetApp {
return app;
}
/**
* The authorization decision behind {@link #allow}, taking the caller directly rather than
* pulling it from a servlet {@link Context} — unit-testable without fabricating a live
* request, the same reason {@code FleetMcp#denyFor} is split from {@code FleetMcp#deny}.
*
* @param knownLeadOrCollaborator the classifier a collaborator's {@code SEND} is checked
* against; pass {@link #auth}'s own {@code
* knownLeadOrCollaborator()} to exercise the real production
* gate, as {@link #allow} does
*/
static boolean permitsFor(Principal caller, Authz.Action action, String target,
Predicate<String> knownLeadOrCollaborator) {
return Authz.permits(caller, action, target, knownLeadOrCollaborator);
}
/**
* Gate a handler on the CB-505 authorization table. Returns {@code true} when the request may
* proceed; otherwise writes the error response and returns {@code false}.
@@ -267,7 +294,7 @@ public final class FleetApp {
return true; // legacy: authorization not enforced
}
Principal caller = ctx.attribute(CALLER);
if (Authz.permits(caller, action, target)) {
if (permitsFor(caller, action, target, auth.knownLeadOrCollaborator())) {
if (action != Authz.Action.READ && action != Authz.Action.METRICS
&& action != Authz.Action.TASK_READ) {
AuditLog.allowed(caller, action, target); // reads would drown the trail
@@ -625,27 +652,31 @@ public final class FleetApp {
*
* <p>Two call shapes share this route, exactly as {@code fleet_send} does over MCP (see
* {@code FleetMcp#sendAction}): a plain delivery to {@code id}, and -- when the body carries
* {@code turnId} -- resolving a worker's blocked question. The body is parsed before the
* authorization check so the right one of {@link Authz.Action#SEND}/{@link Authz.Action#ANSWER}
* reaches the gate; a body that fails to parse is treated as the plain shape for that check
* alone, and is rejected afterward exactly as before.
* {@code turnId} -- resolving a worker's blocked question. The coarse {@link
* Authz.Action#SEND} grant is checked first, before the body is read at all; only once that
* passes is the body parsed, and a present {@code turnId} is then checked again against
* {@link Authz.Action#ANSWER}. A body that fails to parse is rejected with 400 and reaches
* neither {@code messages.answer} nor {@code messages.send}.
*/
private void sendMessage(Context ctx) {
String id = ctx.pathParam("id");
if (!allow(ctx, routeAction("POST /sessions/{id}/message"), id)) {
return;
}
JsonNode body;
try {
body = mapper.readTree(ctx.body());
} catch (Exception e) {
body = null;
}
String turnId = body == null ? null : body.path("turnId").asText(null);
if (!allow(ctx, routeAction("POST /sessions/{id}/message", turnId), id)) {
return;
}
if (body == null) {
ctx.status(400).json(Map.of("error", "bad_request", "detail", "body must be JSON"));
return;
}
String turnId = body.path("turnId").asText(null);
if (!answerGatePasses(turnId, action -> allow(ctx, action, id))) {
return;
}
String content = body.path("content").asText("");
long timeout = body.path("timeoutMs").asLong(DEFAULT_MESSAGE_TIMEOUT_MS);
boolean wait = body.path("wait").asBoolean(true); // default: block for the reply (CB-104)
@@ -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"));
}
}
@@ -0,0 +1,165 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.msg.ReplyInbox;
import io.javalin.Javalin;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertSame;
/**
* fleetd #669 Unit E. Reaches the real {@link dev.ltms.fleet.herdr.HerdrRouter} that {@link
* FleetdAssembly#assembleAndStart} builds and wires — not a copy built for this test — and proves
* that a configured collaborator's terminal routes to the LEAD herdr daemon.
*
* <p>Two distinct {@link FakeHerdr} instances are required, the same pattern {@code
* FleetdAssemblyConnectionIdentityTest} and {@code FleetdLeadRolloverAssemblyTest} already use:
* with one client shared between {@code herdrSocket} and {@code memberHerdrSocket},
* {@code HerdrRouter} folds {@code leadAgents} and {@code memberAgents} into the same instance
* (see its constructor), and {@code agentsFor} would return that one object regardless of whether
* the collaborator map was ever consulted — invisible to a mutation of the predicate this ticket
* fixes. This test's two sockets resolve to two different fakes, so the assertion only passes when
* the collaborator's terminal is actually recognised and routed to the lead one.
*/
class FleetdAssemblyCollaboratorHerdrRoutingTest {
private static final Path LEAD_SOCKET = Path.of("/fake/lead-herdr.sock");
private static final Path MEMBER_SOCKET = Path.of("/fake/member-herdr.sock");
private static final class RecordingResourcePorts implements ResourcePorts {
final Map<Path, HerdrClient> herdrsBySocket = new LinkedHashMap<>();
Runnable shutdownHook;
@Override
public Map<String, String> environment() {
return Map.of();
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
HerdrClient client = herdrsBySocket.get(socketPath);
if (client == null) {
throw new IllegalStateException("no fake herdr registered for socket " + socketPath);
}
return client;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
return (uri, prefetch) -> new ReplyInbox() {
@Override public void own(String target) { }
@Override public void release(String target) { }
@Override public void publish(String target, String msgId, String content) { }
@Override public List<InboxMessage> peek(String target) { return List.of(); }
@Override public boolean ack(String target, String msgId) { return false; }
};
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> {
throw new UnsupportedOperationException(
"leadMailboxOpener must not be called — no coordinator: block is configured");
};
}
@Override
public LongSupplier nanoClock() {
return System::nanoTime;
}
@Override
public LongSupplier wallClockNanos() {
return System::nanoTime;
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor();
}
@Override
public void addShutdownHook(Runnable hook) {
shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
// Do not bind a real port in this assembly test.
}
@Override
public Runnable herdrPollWait() {
return () -> {
throw new UnsupportedOperationException("FakeHerdr is healthy; no poll wait is expected");
};
}
}
private static FleetConfig writeConfig(Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: "%s"
memberHerdrSocket: "%s"
idleSleepGuard:
enabled: false
fleet:
collaborators:
reviewer-alex:
tab: "collab: alex"
profiles:
sonnet:
subscription: true
argv: ["ccs", "sonnet"]
""".formatted(LEAD_SOCKET, MEMBER_SOCKET));
return FleetConfig.load(file);
}
@Test
void assembledRouterRoutesACollaboratorTerminalToTheLeadDaemon(@TempDir Path dir) throws Exception {
FleetConfig cfg = writeConfig(dir);
RecordingResourcePorts ports = new RecordingResourcePorts();
// The fixed FakeHerdr fixture already ties terminal "term_a" to a live agent on tab
// "w2:t7" (pane "w2:p7") — seeding only the tab LABEL to match the configured collaborator
// is enough to make LeadTabScanner resolve "term_a" as that collaborator. Seeded on the
// LEAD fake only: a collaborator's pane lives in the lead daemon, exactly like a lead's.
FakeHerdr lead = new FakeHerdr().withTab("w2", "w2:t7", "collab: alex");
FakeHerdr member = new FakeHerdr();
ports.herdrsBySocket.put(LEAD_SOCKET, lead);
ports.herdrsBySocket.put(MEMBER_SOCKET, member);
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg,
new ConfigRef(dir.resolve("fleetd.yaml"), cfg), new SubscriptionGuard(cfg.guard().hostSet())), ports);
try {
assertSame(runtime.router().leadAgents(), runtime.router().agentsFor("term_a"),
"a configured collaborator's terminal must route to the LEAD daemon — "
+ "FleetdAssembly must wire the collaborator map into the router's "
+ "predicate, not just LeadTabScanner.get()");
assertSame(runtime.router().memberAgents(), runtime.router().agentsFor("term_shell"),
"control: a terminal naming neither a lead nor a collaborator (term_shell, on "
+ "the unlabelled tab w2:t8) must still route to the member daemon");
} finally {
assertNotNull(ports.shutdownHook, "control: assembly must capture its shutdown hook");
ports.shutdownHook.run();
}
}
}
@@ -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");
@@ -0,0 +1,190 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.inject.StatusPoller;
import dev.ltms.fleet.msg.LeadChannelHandle;
import dev.ltms.fleet.msg.LeadCoordLoop;
import dev.ltms.fleet.msg.LeadMessage;
import dev.ltms.fleet.msg.ReplyInbox;
import dev.ltms.fleet.msg.ReplyPushLoop;
import io.javalin.Javalin;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.lang.reflect.Field;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
/**
* Asserts that the assembled loops use the production reminder, coordination, and delivery timing
* defaults when no {@code primary:} block configures the reply-push values.
*/
class FleetdAssemblyTimingDefaultsTest {
private static final class FakeLeadChannel implements LeadChannelHandle {
@Override
public void publish(String toCoordId, LeadMessage message) {
}
@Override
public List<LeadMessage> peek() {
return List.of();
}
@Override
public void ack(String msgId) {
}
@Override
public String selfCoordId() {
return "test-lead";
}
@Override
public boolean heldDurable() {
return true;
}
@Override
public MailboxState inspect(String coordId) {
return MailboxState.unknown(coordId);
}
@Override
public void close() {
}
}
private static final class TestResourcePorts implements ResourcePorts {
final FakeHerdr herdr = new FakeHerdr();
Runnable shutdownHook;
@Override
public Map<String, String> environment() {
return Map.of();
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
return herdr;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
return (uri, prefetch) -> new ReplyInbox() {
@Override public void own(String target) { }
@Override public void release(String target) { }
@Override public void publish(String target, String msgId, String content) { }
@Override public List<InboxMessage> peek(String target) { return List.of(); }
@Override public boolean ack(String target, String msgId) { return false; }
};
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> new FakeLeadChannel();
}
@Override
public LongSupplier nanoClock() {
return System::nanoTime;
}
@Override
public LongSupplier wallClockNanos() {
return System::nanoTime;
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor();
}
@Override
public void addShutdownHook(Runnable hook) {
shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
}
@Override
public Runnable herdrPollWait() {
return () -> {
throw new UnsupportedOperationException("FakeHerdr is healthy; no poll wait is expected");
};
}
}
private TestResourcePorts ports;
@AfterEach
void tearDown() {
if (ports != null && ports.shutdownHook != null) {
ports.shutdownHook.run();
}
}
private static FleetConfig writeConfig(Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, """
bind:
host: 127.0.0.1
port: 8765
idleSleepGuard:
enabled: false
coordinator:
uri: "amqp://fake-lead-broker/vh"
selfId: "test-lead"
""");
return FleetConfig.load(file);
}
private FleetdRuntime assemble(Path dir) throws Exception {
FleetConfig cfg = writeConfig(dir);
ports = new TestResourcePorts();
return FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg,
new ConfigRef(dir.resolve("fleetd.yaml"), cfg), new SubscriptionGuard(cfg.guard().hostSet())), ports);
}
private static long longField(Object target, String name) throws Exception {
Field field = target.getClass().getDeclaredField(name);
field.setAccessible(true);
return field.getLong(target);
}
@Test
void productionBootPathUsesTheExpectedLoopTimingDefaults(@TempDir Path dir) throws Exception {
FleetdRuntime runtime = assemble(dir);
ReplyPushLoop pushLoop = runtime.pushLoop();
assertEquals(5, longField(pushLoop, "maxReminders"),
"without primary:, ReplyPushLoop must stop after five reminder attempts");
assertEquals(15_000L, longField(pushLoop, "backoffMs"),
"without primary:, ReplyPushLoop must wait fifteen seconds before the next reminder");
LeadCoordLoop leadCoordLoop = runtime.leadCoordLoop();
assertNotNull(leadCoordLoop, "control: coordinator: must build LeadCoordLoop");
assertEquals(3_000L, longField(leadCoordLoop, "intervalMs"),
"LeadCoordLoop must poll for peer-lead mail every three seconds");
StatusPoller poller = runtime.poller();
assertEquals(Injector.POLL_INTERVAL_MILLIS, longField(poller, "intervalMillis"),
"StatusPoller must use Injector's delivery poll interval");
}
}
@@ -14,6 +14,7 @@ class AuthzTest {
private static final Principal ANON = Principal.anonymous();
private static final Principal ARCH_DESIGN = Principal.architect("lead-designer", "term_design", 400);
private static final Principal ARCH_OTHER = Principal.architect("reviewer", "term_review", 500);
private static final Principal COLLABORATOR = Principal.collaborator("ops", "term_collab", 600);
@Test
void anonymousIsAuthorizedForNothing() {
@@ -166,4 +167,87 @@ class AuthzTest {
assertFalse(Authz.isUnauthenticated(WORKER_A));
assertFalse(Authz.isUnauthenticated(PRIMARY));
}
// ── the collaborator matrix ─────────────────────────────────────────────────────────────────
/**
* {@code SEND} for a collaborator is the one grant that is conditional rather than fixed:
* flipping only the classifier's answer for the target flips only this outcome.
*/
@Test
void aCollaboratorMaySendOnlyWhenTheClassifierAcceptsTheTarget() {
assertTrue(Authz.permits(COLLABORATOR, SEND, "term_lead", target -> true),
"the classifier accepting the target must grant SEND");
assertFalse(Authz.permits(COLLABORATOR, SEND, "term_lead", target -> false),
"the classifier refusing the target must deny SEND");
assertFalse(Authz.permits(COLLABORATOR, SEND, "term_lead"),
"the real production classifier recognises no terminal yet, so SEND is refused today");
}
/**
* Control for the test above: every other action's result for a collaborator does not move
* when the classifier does. Only {@code SEND} is wired to it.
*/
@Test
void theClassifierMovesOnlySendForACollaborator() {
for (Authz.Action a : Authz.Action.values()) {
if (a == SEND) {
continue;
}
assertEquals(
Authz.permits(COLLABORATOR, a, "term_lead"),
Authz.permits(COLLABORATOR, a, "term_lead", target -> true),
a + " must not depend on the classifier at all");
}
}
@Test
void aCollaboratorMayReadAndScrapeMetrics() {
assertTrue(Authz.permits(COLLABORATOR, READ, null));
assertTrue(Authz.permits(COLLABORATOR, METRICS, null));
}
@Test
void aCollaboratorMayReplyAndAskOnlyAsItsOwnPane() {
assertTrue(Authz.permits(COLLABORATOR, REPLY, "term_collab"),
"its own pane is its own");
assertTrue(Authz.permits(COLLABORATOR, ASK, "term_collab"));
assertFalse(Authz.permits(COLLABORATOR, REPLY, "term_design"),
"a collaborator must not reply on another pane");
assertFalse(Authz.permits(COLLABORATOR, REPLY, null),
"an absent target must not pass the own-session rule");
}
/**
* Every action denied to a collaborator, asserted denied even when the classifier would
* accept any target — proving none of these is actually gated on the classifier at all.
*/
@Test
void aCollaboratorIsDeniedLifecycleCoordinationAndTicketPolling() {
for (Authz.Action a : new Authz.Action[]{SPAWN, STOP, DRAIN, HANDOVER, ANSWER, COORD_SEND,
COORD_READ, TASK_READ}) {
assertFalse(Authz.permits(COLLABORATOR, a, "term_lead", target -> true),
"a collaborator must not " + a + " even when the classifier accepts every target");
}
}
@Test
void aCollaboratorIsNotCountedAsPrimaryWorkerOrArchitect() {
assertFalse(COLLABORATOR.isPrimary());
assertFalse(COLLABORATOR.isWorker());
assertFalse(COLLABORATOR.isArchitect());
assertTrue(COLLABORATOR.isCollaborator());
}
/**
* A collaborator is never spawned, so it must not be enrolled in the presence map as an
* available member. Control: both a worker and an architect — which ARE spawned — still are.
*/
@Test
void isSpawnedMemberIsFalseForACollaboratorButTrueForAWorkerAndAnArchitect() {
assertFalse(COLLABORATOR.isSpawnedMember());
assertTrue(WORKER_A.isSpawnedMember());
assertTrue(ARCH_DESIGN.isSpawnedMember());
}
}
@@ -523,4 +523,129 @@ class CallerResolverTest {
}
}
// ── fleetd #669 Unit D: a live spawned member outranks every tab map ────────────────────────
/**
* Criterion 1: a terminal present in BOTH the spawned-member roster AND the lead tab map
* resolves as its member role, not as a lead — the roster is checked first, consulting no tab
* map at all when it matches.
*/
@Test
void aSpawnedMemberWinsOverALeadTabForTheSamePane() {
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
() -> Map.of("term_a", "opus-5.0"), new MemberRegistry(null),
t -> "term_a".equals(t) ? MemberRole.DEV : null, Map::of)
.resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, p.role(),
"a live spawned member's own identity must win over a tab map naming the same pane a lead");
assertEquals("term_a", p.terminal());
}
/** A spawned architect in the roster resolves ARCHITECT, carrying its bound slot's name. */
@Test
void aSpawnedArchitectInTheRosterResolvesArchitectWithItsSlotName() {
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
boundMembers("architect:lead-designer", MemberRole.ARCHITECT),
t -> "term_a".equals(t) ? MemberRole.ARCHITECT : null, Map::of)
.resolve("127.0.0.1", 42, null);
assertEquals(Role.ARCHITECT, p.role());
assertEquals("lead-designer", p.name());
assertEquals("term_a", p.terminal());
}
/**
* fleetd #424 regression: the roster only answers THAT a pane is a live spawned member; config
* still decides WHAT that member's slot grants. A slot revoked after the bind must still demote
* the session on its very next request, exactly as it would for a pane with no roster entry at
* all — the roster's own ARCHITECT role must never be granted on its word alone.
*
* <p>{@code bind} refuses an unconfigured slot, so the revoked state can only be reached by
* binding while the slot is configured and then swapping the config out from under it, the way
* a live reload does.
*/
@Test
void aRevokedArchitectSlotDemotesALiveSpawnedArchitectToWorker() {
FleetConfig.Fleet configured = new FleetConfig.Fleet(Map.of(),
Map.of("lead-designer", new FleetConfig.Slot("sonnet")), Map.of(), Map.of(), null);
java.util.concurrent.atomic.AtomicReference<FleetConfig.Fleet> live =
new java.util.concurrent.atomic.AtomicReference<>(configured);
MemberRegistry members = MemberRegistry.live(live::get);
assertTrue(members.bind("architect:lead-designer", "term_a"));
live.set(new FleetConfig.Fleet(Map.of(), Map.of(), Map.of(), Map.of(), null)); // slot revoked
// Setup controls: the slot is really gone from config, but the occupancy is still there —
// otherwise this test would pass for the wrong reason.
assertNull(members.roleForSlot("architect:lead-designer"), "setup control: the slot must be gone from config");
assertEquals("architect:lead-designer", members.snapshot().get("term_a"),
"setup control: the binding itself must still be there");
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
members, t -> "term_a".equals(t) ? MemberRole.ARCHITECT : null, Map::of)
.resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, p.role(),
"a revoked slot must demote a live spawned architect on its very next request");
}
/** Criterion 3: a configured collaborator tab that is not a spawned member resolves COLLABORATOR. */
@Test
void aConfiguredCollaboratorTabResolvesToCollaboratorCarryingItsName() {
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
new MemberRegistry(null), t -> null, () -> Map.of("term_a", "ops"))
.resolve("127.0.0.1", 42, null);
assertEquals(Role.COLLABORATOR, p.role());
assertEquals("ops", p.name());
assertEquals("term_a", p.terminal());
}
/** Regression: an empty collaborator registry leaves every pane exactly as before. */
@Test
void anEmptyCollaboratorRegistryLeavesEveryPaneAsBefore() {
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
new MemberRegistry(null), t -> null, Map::of)
.resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, p.role());
assertNull(p.name());
}
@Test
void describeNamesTheCollaborator() {
assertEquals("collaborator:ops", Principal.collaborator("ops", "term_a", 1).describe());
}
// ── fleetd #669 Unit D: knownLeadOrCollaborator() reads the same maps resolve() does ───────────
@Test
void knownLeadOrCollaboratorIsTrueForALeadTerminal() {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
() -> Map.of("term_lead", "opus-5.0"), new MemberRegistry(null), t -> null, Map::of);
assertTrue(r.knownLeadOrCollaborator().test("term_lead"));
assertFalse(r.knownLeadOrCollaborator().test("term_other"));
}
@Test
void knownLeadOrCollaboratorIsTrueForACollaboratorTerminal() {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
new MemberRegistry(null), t -> null, () -> Map.of("term_collab", "ops"));
assertTrue(r.knownLeadOrCollaborator().test("term_collab"));
assertFalse(r.knownLeadOrCollaborator().test("term_other"));
}
@Test
void knownLeadOrCollaboratorIsFalseForASpawnedMembersTerminal() {
// The exact scenario a collaborator's SEND must never reach: a live spawned member's own
// terminal, which is neither a configured lead nor a configured collaborator.
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
new MemberRegistry(null), t -> "term_a".equals(t) ? MemberRole.DEV : null, Map::of);
assertFalse(r.knownLeadOrCollaborator().test("term_a"));
}
}
@@ -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
@@ -676,12 +676,8 @@ class FleetConfigTest {
assertEquals(5, hb.quietNudgeCap());
}
/**
* The hazard the guard exists for: fleetd writes worker tab labels and reads lead tab labels.
* Overlap the two and every worker it spawns is read back as a lead.
*/
@Test
void aLeadPrefixThatAProfileTabLabelOverrideAlsoMatchesRefusesToStart(@TempDir Path dir)
void aProfileTabLabelOverrideMatchingALeadTabRefusesToStart(@TempDir Path dir)
throws Exception {
Path f = dir.resolve("collide.yaml");
Files.writeString(f, """
@@ -689,32 +685,33 @@ class FleetConfigTest {
port: 8080
profiles:
gx10:
tabLabel: "lead: {profile} #{n}"
tabLabel: "alpha"
fleet:
leaders:
opus:
tab: "lead: opus"
tabPrefix: "lead:"
tab: "alpha"
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e =
assertThrows(IllegalStateException.class, cfg::validateLeadTabPrefixes);
assertTrue(e.getMessage().contains("gx10"), "the message must name the offending profile");
assertTrue(e.getMessage().contains("alpha"), "the message must name the offending label");
}
/** A bad fleet-wide template promotes every member, not one profile — so it is checked too. */
@Test
void aFleetTabLabelThatMatchesALeadPrefixRefusesToStart(@TempDir Path dir) throws Exception {
void aFleetTabLabelTemplateThatCanRenderAsALeadTabRefusesToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("collide-template.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
pha: {}
fleet:
tabLabel: "lead: {role} {profile}"
tabLabel: "al{profile}"
leaders:
opus:
tab: "lead: opus"
tab: "alpha"
""");
FleetConfig cfg = FleetConfig.load(f);
@@ -723,12 +720,28 @@ class FleetConfigTest {
assertTrue(e.getMessage().contains("fleet.tabLabel"));
}
/**
* The point of making role the label's first field: {@code {role}} comes from a closed enum, so
* a generated label cannot begin with {@code "lead:"} however the fleet is configured.
*/
@Test
void theDefaultTabLabelCannotCollideWithTheDefaultLeadPrefix(@TempDir Path dir) throws Exception {
void anExactFleetTabLabelCollisionRefusesToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("exact-tab-collision.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
tabLabel: "alpha"
leaders:
alpha:
tab: "alpha"
""");
IllegalStateException e = assertThrows(IllegalStateException.class,
() -> FleetConfig.load(f).validateAll());
assertTrue(e.getMessage().contains("fleet.tabLabel"),
"the message must name the offending label");
assertTrue(e.getMessage().contains("alpha"), "the message must name the colliding lead tab");
}
@Test
void aFleetTabLabelTemplateThatCannotRenderAsALeadTabIsAllowed(@TempDir Path dir) throws Exception {
Path f = dir.resolve("ok.yaml");
Files.writeString(f, """
bind:
@@ -737,17 +750,13 @@ class FleetConfigTest {
gx10:
baseUrl: http://gx00.gw:8000
fleet:
tabLabel: "worker-{profile}"
leaders:
opus:
tab: "lead: opus"
tab: "alpha"
""");
assertDoesNotThrow(() -> FleetConfig.load(f).validateLeadTabPrefixes());
for (MemberRole role : MemberRole.values()) {
assertFalse(FleetConfig.Fleet.DEFAULT_TAB_LABEL
.replace("{role}", role.wireName()).startsWith("lead:"),
"no role renders a label that reads as a lead");
}
assertDoesNotThrow(() -> FleetConfig.load(f).validateAll());
}
@Test
@@ -765,6 +774,78 @@ class FleetConfigTest {
"a label that collides with a convention nobody reads is not a problem");
}
/**
* fleetd #677: identity is matched on a lead's exact {@code tab} alone, so two leads sharing
* one tab means only one of them is ever found — the guard must catch this independently of
* the member-template checks above.
*/
@Test
void twoLeadsSharingTheSameExactTabRefusesToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("shared-tab.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
leaders:
opus:
tab: "shared tab"
sonnet:
tab: "shared tab"
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e =
assertThrows(IllegalStateException.class, cfg::validateLeadTabPrefixes);
assertTrue(e.getMessage().contains("opus"), "the message must name one offending lead");
assertTrue(e.getMessage().contains("sonnet"), "the message must name the other offending lead");
}
/**
* fleetd #693: the guard matches tabs case-insensitively, because
* {@code LeadTabScanner} keys its tab map on a lowercased label — two tabs differing only in
* case collide there too, and the guard must catch that independently of the exact-match case
* above.
*/
@Test
void twoLeadsSharingTheSameTabInDifferentCaseRefusesToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("shared-tab-case.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
leaders:
opus:
tab: "Shared Tab"
sonnet:
tab: "shared tab"
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e =
assertThrows(IllegalStateException.class, cfg::validateLeadTabPrefixes);
assertTrue(e.getMessage().contains("opus"), "the message must name one offending lead");
assertTrue(e.getMessage().contains("sonnet"), "the message must name the other offending lead");
}
/** Control for {@link #twoLeadsSharingTheSameExactTabRefusesToStart}: distinct tabs load cleanly. */
@Test
void twoLeadsWithDistinctExactTabsAreAllowed(@TempDir Path dir) throws Exception {
Path f = dir.resolve("distinct-tabs.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
leaders:
opus:
tab: "opus tab"
sonnet:
tab: "sonnet tab"
""");
FleetConfig cfg = FleetConfig.load(f);
assertDoesNotThrow(cfg::validateLeadTabPrefixes);
}
// ── validatePanePlacementAgainstLeadTabs ────────────────────────────────────────────────────
/**
@@ -919,6 +1000,273 @@ class FleetConfigTest {
"no primary.terminal pin ⇒ nothing registered, even with fleet.leaders configured");
}
// ── fleetd #669: the collaborators registry ────────────────────────────────────────────────
/**
* {@code Fleet} is {@code @JsonIgnoreProperties(ignoreUnknown = true)}, so a config naming
* {@code fleet.collaborators.<name>.tab} loads with no exception whether or not the key is
* ever read into the object model. Asserting only "no exception" would pass both before and
* after the real fix, so this asserts the parsed value is actually reachable from the loaded
* {@code FleetConfig} — the one thing a vacuous "no exception" test cannot tell apart.
*/
@Test
void collaboratorsBlockIsActuallyParsedNotSilentlyDropped(@TempDir Path dir) throws Exception {
Path f = dir.resolve("collaborators.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
collaborators:
reviewer-alex:
tab: "collab: alex"
""");
FleetConfig cfg = FleetConfig.load(f);
assertEquals("collab: alex", cfg.fleet().collaborators().get("reviewer-alex").tab());
}
/** A collaborator carries no field other than {@code tab}, so a blank one is meaningless. */
@Test
void aCollaboratorWithNoTabRefusesToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("useless-collaborator.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
collaborators:
ghost: {}
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e = assertThrows(IllegalStateException.class, cfg::validateMembers);
assertTrue(e.getMessage().contains("ghost"), "the message must name the useless entry");
assertTrue(e.getMessage().contains("tab:"), "the message must say what is missing");
}
/** Control for {@link #aCollaboratorWithNoTabRefusesToStart}: a named tab loads cleanly. */
@Test
void aCollaboratorWithATabIsAllowed(@TempDir Path dir) throws Exception {
Path f = dir.resolve("named-collaborator.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
collaborators:
reviewer-alex:
tab: "collab: alex"
""");
FleetConfig cfg = FleetConfig.load(f);
assertDoesNotThrow(cfg::validateMembers);
}
/**
* fleetd #669: the fleet-wide {@code tabLabel} template can render as a collaborator tab, the
* same hazard {@link #aFleetTabLabelTemplateThatCanRenderAsALeadTabRefusesToStart} covers on
* the lead side. Drives the fleet-wide branch directly, with no profile override involved.
*/
@Test
void aFleetTabLabelTemplateThatCanRenderAsACollaboratorTabRefusesToStart(@TempDir Path dir)
throws Exception {
Path f = dir.resolve("collide-template-collaborator.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
tabLabel: "al{profile}"
collaborators:
alex:
tab: "alpha"
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e =
assertThrows(IllegalStateException.class, cfg::validateLeadTabPrefixes);
assertTrue(e.getMessage().contains("fleet.tabLabel"));
assertTrue(e.getMessage().contains("alex"), "the message must name the offending collaborator");
}
/**
* fleetd #669: a member tabLabel that can render as a configured collaborator tab is the same
* hazard as the lead case above — a member labelled that way is read back as the collaborator.
*/
@Test
void aProfileTabLabelOverrideMatchingACollaboratorTabRefusesToStart(@TempDir Path dir)
throws Exception {
Path f = dir.resolve("collide-collaborator.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
gx10:
tabLabel: "collab-tab"
fleet:
collaborators:
alex:
tab: "collab-tab"
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e =
assertThrows(IllegalStateException.class, cfg::validateLeadTabPrefixes);
assertTrue(e.getMessage().contains("gx10"), "the message must name the offending profile");
assertTrue(e.getMessage().contains("collab-tab"), "the message must name the offending label");
}
/** Control: a profile tabLabel that cannot render as the collaborator tab is allowed. */
@Test
void aProfileTabLabelThatCannotRenderAsACollaboratorTabIsAllowed(@TempDir Path dir)
throws Exception {
Path f = dir.resolve("ok-collaborator.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
gx10:
tabLabel: "worker-{profile}"
fleet:
collaborators:
alex:
tab: "collab-tab"
""");
FleetConfig cfg = FleetConfig.load(f);
assertDoesNotThrow(cfg::validateLeadTabPrefixes);
}
/** fleetd #669: identity is matched on a collaborator's exact tab, so two sharing one are unreachable. */
@Test
void twoCollaboratorsSharingTheSameExactTabRefusesToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("shared-collaborator-tab.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
collaborators:
alex:
tab: "shared tab"
sam:
tab: "Shared Tab"
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e =
assertThrows(IllegalStateException.class, cfg::validateLeadTabPrefixes);
assertTrue(e.getMessage().contains("alex"), "the message must name one offending collaborator");
assertTrue(e.getMessage().contains("sam"), "the message must name the other offending collaborator");
}
/** Control for {@link #twoCollaboratorsSharingTheSameExactTabRefusesToStart}: distinct tabs load cleanly. */
@Test
void twoCollaboratorsWithDistinctExactTabsAreAllowed(@TempDir Path dir) throws Exception {
Path f = dir.resolve("distinct-collaborator-tabs.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
collaborators:
alex:
tab: "alex tab"
sam:
tab: "sam tab"
""");
FleetConfig cfg = FleetConfig.load(f);
assertDoesNotThrow(cfg::validateLeadTabPrefixes);
}
/**
* fleetd #669: a collaborator tab equal to a lead tab crosses a privilege boundary — the worst
* of the three new collisions, since only one of the two identities is ever found.
*/
@Test
void aCollaboratorTabEqualToALeadTabRefusesToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("lead-collaborator-collision.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
leaders:
opus:
tab: "shared tab"
collaborators:
alex:
tab: "Shared Tab"
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e =
assertThrows(IllegalStateException.class, cfg::validateLeadTabPrefixes);
assertTrue(e.getMessage().contains("opus"), "the message must name the offending lead");
assertTrue(e.getMessage().contains("alex"), "the message must name the offending collaborator");
}
/** Control for {@link #aCollaboratorTabEqualToALeadTabRefusesToStart}: distinct tabs load cleanly. */
@Test
void aLeadAndACollaboratorWithDistinctTabsAreAllowed(@TempDir Path dir) throws Exception {
Path f = dir.resolve("lead-collaborator-ok.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
leaders:
opus:
tab: "lead tab"
collaborators:
alex:
tab: "collab tab"
""");
FleetConfig cfg = FleetConfig.load(f);
assertDoesNotThrow(cfg::validateLeadTabPrefixes);
}
/**
* fleetd #669: a pane-placed member can land in a collaborator's labelled tab exactly as it
* can land in a lead's — {@code validatePanePlacementAgainstLeadTabs()} must fire even when
* {@code fleet.leaders} is empty, which is the early-return the brief flagged as the bug.
*/
@Test
void aPanePlacedProfileWithACollaboratorTabRefusesToStartEvenWithNoLeaders(@TempDir Path dir)
throws Exception {
Path f = dir.resolve("pane-hazard-collaborator.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
gx10:
placement: pane
fleet:
collaborators:
alex:
tab: "collab: alex"
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e = assertThrows(IllegalStateException.class,
cfg::validatePanePlacementAgainstLeadTabs);
assertTrue(e.getMessage().contains("gx10"), "the message must name the offending profile");
}
/** Control: a pane-placed profile with no lead or collaborator tab configured is allowed. */
@Test
void aPanePlacedProfileWithNoLeaderOrCollaboratorTabIsAllowed(@TempDir Path dir) throws Exception {
Path f = dir.resolve("pane-no-tab-at-all.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
gx10:
placement: pane
fleet:
collaborators:
alex: {}
""");
assertDoesNotThrow(() -> FleetConfig.load(f).validatePanePlacementAgainstLeadTabs(),
"a collaborator with no tab feeds nothing into the scanner, so pane placement is safe");
}
// ── CB-548: the architects registry ────────────────────────────────────────────────────────
@Test
@@ -1191,6 +1539,60 @@ class FleetConfigTest {
assertEquals(Set.of("sonnet"), cfg.fleet().pool(MemberRole.REVIEWER).keySet());
}
/**
* A duplicated name in {@code fleet.collaborators} is refused at parse time, like any other
* {@code fleet:} pool. See {@link #duplicateSlotNamesInOnePoolAreRejectedAtParseTime}.
*/
@Test
void duplicateCollaboratorNamesAreRejectedAtParseTime(@TempDir Path dir) throws Exception {
Path f = dir.resolve("collaborator-dup.yaml");
Files.writeString(f, """
bind:
port: 8080
fleet:
collaborators:
alex:
tab: "collab: alex"
alex:
tab: "collab: alex, second"
""");
IllegalStateException e =
assertThrows(IllegalStateException.class, () -> FleetConfig.load(f));
assertTrue(e.getMessage().contains("alex"),
"the refusal names the duplicated entry, was: " + e.getMessage());
assertTrue(e.getMessage().contains("fleet.collaborators"),
"the refusal names the pool the duplicate is in, was: " + e.getMessage());
}
/**
* Control for {@link #duplicateCollaboratorNamesAreRejectedAtParseTime}: the same name reused
* across the collaborators registry and a member role pool is the role × profile matrix doing
* its job in the other pool, not a mistake — only a repeat within one pool loses an entry.
*/
@Test
void theSameNameInCollaboratorsAndAnotherPoolIsNotADuplicate(@TempDir Path dir) throws Exception {
Path f = dir.resolve("collaborator-cross-pool.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
sonnet:
baseUrl: http://gx10.gw:8000
fleet:
architects:
alex:
profile: sonnet
collaborators:
alex:
tab: "collab: alex"
""");
FleetConfig cfg = assertDoesNotThrow(() -> FleetConfig.load(f));
assertEquals(Set.of("alex"), cfg.fleet().pool(MemberRole.ARCHITECT).keySet());
assertEquals(Set.of("alex"), cfg.fleet().collaborators().keySet());
}
@Test
void duplicateKeysOutsideTheFleetPoolsAreUnaffected(@TempDir Path dir) throws Exception {
// The duplicate check is scoped to the fleet pools — a duplicate elsewhere is not this
@@ -17,63 +17,21 @@ import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* The gap this class exists to close: mutation testing on the fleetd ticket "central allow-list
* of usable models" found that although {@link FleetConfig#validateModels()}'s own logic was well
* pinned, nothing proved either real caller ({@code Fleetd.main} and {@link ConfigRef#reload()})
* still invoked it — deleting the call site left the full suite green (1478/0/0/0). A follow-up
* measurement (same technique — remove one call site, run the suite, not read the code) found the
* SAME gap for all five of {@link FleetConfig}'s other validators at startup, and for four of the
* six inside {@link ConfigRef#reload()}. This is a class of gap, not one line's mistake: every one
* of those thirteen tests called the validator itself directly, never the real caller that was
* supposed to.
* Tests the reflective validator sweep and {@link FleetConfig#validateAll()} reachability.
*
* <p>The fix replaces the six individual {@code cfg.validateXxx()} calls at each of the two real
* call sites with one {@link FleetConfig#validateAll()}, which reaches every validator by
* reflection rather than by a hand-maintained list of names. A hand-maintained list of six names
* would have exactly the defect it replaces: the seventh validator someone adds next month has no
* reason to be added to it, and nothing would say so. This class proves TWO separate claims, and
* keeps them separate on purpose:
*
* <ol>
* <li>{@link #theSweepMechanismIsGenericNotHardcodedToFleetConfigsSixNames()} and its neighbours
* prove the reflective sweep itself ({@link FleetConfig#invokeAllValidators}) is a general
* mechanism — it runs whatever public, no-arg, void {@code validateXxx()} methods a class
* happens to declare today, including a class with more of them than {@link FleetConfig}
* has right now. This is the proof that a future, real seventh validator on {@link
* FleetConfig} would be swept automatically, without needing to add a real (unwanted)
* seventh validator just to exercise the claim.</li>
* <li>{@link #validateAllReachesEveryOneOfTodaysRealValidators()} proves {@link
* FleetConfig#validateAll()} itself is wired to that same generic mechanism and genuinely
* reaches every one of today's real validators — reusing the exact minimal failing
* configurations {@code FleetConfigTest} already established for each one directly, plus a
* dedicated fixture for {@link FleetConfig#validateLeadRollover()}, which no other test
* drives through {@code validateAll()} — so a single call to {@code validateAll()} is shown
* to reproduce every one of those failures.</li>
* </ol>
*
* <p>Together with the direct-{@code Fleetd.main}-invocation tests in {@code
* FleetdStartupValidationTest} (which prove the real startup call site still calls {@code
* validateAll()}) and the {@code ConfigRefTest} reload tests (which prove the same for {@link
* ConfigRef#reload()}), removing {@code cfg.validateAll();} from either real call site now fails
* a test in this module.
*
* <p><b>What is NOT pinned, measured rather than assumed.</b> Reverting {@link
* FleetConfig#validateAll()} to a hardcoded list of today's method calls leaves the whole
* suite green. Nothing ties {@code validateAll()} to
* the generic sweep — claim 1 proves {@link FleetConfig#invokeAllValidators} is generic, and claim
* 2 proves {@code validateAll()} reaches today's validators, and a hardcoded list satisfies both. So the
* reflective sweep is a convenience, not the guarantee. The guarantee is {@link
* #fleetConfigDeclaresExactlyTheseValidatorsToday()}: it fails the moment any validator is added
* or removed, which forces whoever changes the set to look at this file.
* <p>{@link #theSweepRunsEveryValidateMethodOnAnUnrelatedClass()} and its neighbours
* prove that {@link FleetConfig#invokeAllValidators} runs each public, no-arg, void
* {@code validateXxx()} method on its target. {@link #fleetConfigDeclaresExactlyTheseValidatorsToday()}
* is the canary for the validator set. {@link #validateAllReachesEveryOneOfTodaysRealValidators()}
* is the reachability check for that set.
*/
class FleetConfigValidateAllTest {
// ── Claim 1: the reflective sweep is a general mechanism, not six names in disguise ──────────
// ── Claim 1: the reflective sweep is a general mechanism ─────────────────────────────────────
/**
* A throwaway fixture class, unrelated to {@link FleetConfig} in every way except shape: three
* public, no-arg, void methods named {@code validateXxx}. Proves the sweep works on ANY class
* with this shape, not on something special-cased to {@link FleetConfig}.
* Fixture with public, no-arg, void methods named {@code validateXxx}. It proves the sweep uses
* the target's method shape rather than special handling for {@link FleetConfig}.
*/
static class ThreeValidators {
final List<String> ran = new ArrayList<>();
@@ -92,7 +50,7 @@ class FleetConfigValidateAllTest {
}
@Test
void theSweepMechanismIsGenericNotHardcodedToFleetConfigsSixNames() {
void theSweepRunsEveryValidateMethodOnAnUnrelatedClass() {
ThreeValidators target = new ThreeValidators();
FleetConfig.invokeAllValidators(target);
assertEquals(List.of("validateAlpha", "validateBeta", "validateGamma"), target.ran,
@@ -102,12 +60,8 @@ class FleetConfigValidateAllTest {
}
/**
* The core of the "self-maintaining" requirement: the exact same class shape as {@link
* ThreeValidators}, plus one more method — standing in for "a developer adds a validator next
* month". Nothing about the sweep changes to pick it up; the new method is invoked purely
* because it exists and matches the shape. This is what makes adding a seventh real validator
* to {@link FleetConfig} safe without touching {@link FleetConfig#validateAll()} or either
* call site — there is no "wire it in" step left to forget.
* Fixture with an added valid method. It proves the sweep reaches a method because it matches
* the validator shape.
*/
static class FourValidators {
final List<String> ran = new ArrayList<>();
@@ -135,7 +89,7 @@ class FleetConfigValidateAllTest {
FleetConfig.invokeAllValidators(target);
assertEquals(List.of("validateAlpha", "validateBeta", "validateDelta", "validateGamma"),
sorted(target.ran),
"the fourth method must be reached automatically — proving a class can grow the "
"the added method must be reached automatically — proving a class can grow the "
+ "set of things it validates with no change to the sweep itself");
}
@@ -297,16 +251,16 @@ class FleetConfigValidateAllTest {
port: 8765
""", "auth.mode: token");
// validateLeadTabPrefixes: a fleet-wide tabLabel that starts with a lead's own tabPrefix.
// validateLeadTabPrefixes: a fleet-wide tabLabel that equals a lead tab.
assertValidateAllRefuses(dir, "lead-tab-prefixes.yaml", """
bind:
host: 127.0.0.1
port: 8765
fleet:
tabLabel: "lead: {role} {profile}"
tabLabel: "alpha"
leaders:
opus:
tab: "lead: opus"
tab: "alpha"
""", "fleet.tabLabel");
// validateSubscriptionProfiles: subscription: true with env: reseating ANTHROPIC_BASE_URL.
@@ -2,6 +2,8 @@ package dev.ltms.fleet.herdr;
import org.junit.jupiter.api.Test;
import java.util.Map;
import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertNotSame;
@@ -28,4 +30,44 @@ class HerdrRouterTest {
assertSame(router.leadAgents(), router.agentsFor("lead"));
assertSame(router.memberAgents(), router.agentsFor("member"));
}
/**
* fleetd #669 Unit E. The predicate shape here is exactly what {@code FleetdAssembly} builds:
* true when the terminal is a known lead OR a known collaborator. Two distinct clients are
* required — with one shared client {@code agentsFor} would return the same object regardless
* of the predicate's answer, and this assertion would pass whether or not the collaborator map
* was ever consulted.
*/
@Test
void collaboratorTerminalRoutesToTheLeadDaemon() {
FakeHerdr lead = new FakeHerdr();
FakeHerdr member = new FakeHerdr();
Map<String, String> leads = Map.of("term_lead", "primary");
Map<String, String> collaborators = Map.of("term_collab", "reviewer-alex");
HerdrRouter router = new HerdrRouter(lead, member,
id -> leads.containsKey(id) || collaborators.containsKey(id));
assertSame(router.leadAgents(), router.agentsFor("term_collab"),
"a configured collaborator's terminal must route to the LEAD daemon, not the member "
+ "one — its pane is opened by a person, exactly like a lead's");
}
/**
* Companion to {@link #collaboratorTerminalRoutesToTheLeadDaemon}: a terminal that is neither a
* known lead nor a known collaborator must still route to the member daemon. Without this, a
* predicate of {@code _ -> true} would also pass the test above.
*/
@Test
void terminalInNeitherMapStillRoutesToTheMemberDaemon() {
FakeHerdr lead = new FakeHerdr();
FakeHerdr member = new FakeHerdr();
Map<String, String> leads = Map.of("term_lead", "primary");
Map<String, String> collaborators = Map.of("term_collab", "reviewer-alex");
HerdrRouter router = new HerdrRouter(lead, member,
id -> leads.containsKey(id) || collaborators.containsKey(id));
assertSame(router.memberAgents(), router.agentsFor("term_worker"),
"a terminal absent from both maps must stay on the member daemon — the fix widens "
+ "the predicate, it does not make it unconditionally true");
}
}
@@ -163,6 +163,13 @@ class LeadTabScannerTest {
return new LeadTabScanner(herdr, tabToName, Set.of("fleetd-workers"), TTL, clock::get);
}
private LeadTabScanner scannerWithCollaborators(TopologyHerdr herdr, Map<String, String> tabToName,
Map<String, String> collaboratorTabToName,
AtomicLong clock) {
return new LeadTabScanner(herdr, tabToName, collaboratorTabToName,
Set.of("fleetd-workers"), TTL, clock::get);
}
@Test
void everyConfiguredTabBecomesALeadNamedByItsEntry() {
Map<String, String> leads = scanner(twoLeads(), twoLeadsConfigured(), new AtomicLong()).get();
@@ -198,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() {
@@ -212,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")
@@ -464,4 +497,75 @@ class LeadTabScannerTest {
assertEquals(afterFirst, herdr.calls, "the failure path must be rate-limited too");
}
// ── fleetd #669 Unit D: a collaborator tab is matched the same way as a lead tab, one pass ─────
@Test
void aConfiguredCollaboratorTabIsReportedByCollaboratorsNotByGet() {
TopologyHerdr herdr = new TopologyHerdr()
.workspace("w1", "main")
.tab("w1:t1", "w1", "collab: ops")
.pane("w1:p1", "w1:t1", "term_ops");
LeadTabScanner s = scannerWithCollaborators(herdr, Map.of(), Map.of("collab: ops", "ops"),
new AtomicLong());
assertEquals(Map.of("term_ops", "ops"), s.collaborators(),
"a collaborator tab is matched exactly like a lead tab");
assertEquals(Map.of(), s.get(), "a collaborator tab must never also appear as a lead");
}
/**
* Criterion 4, scanner level: a collaborator tab that is labelled but runs no agent is not
* reported — the same #359 liveness cross-check a lead tab gets.
*/
@Test
void aDeadCollaboratorTabIsNotReported() {
TopologyHerdr herdr = new TopologyHerdr()
.workspace("w1", "main")
.tab("w1:t1", "w1", "collab: ops")
.pane("w1:p1", "w1:t1", "term_ops")
.deadAgent("w1:t1");
LeadTabScanner s = scannerWithCollaborators(herdr, Map.of(), Map.of("collab: ops", "ops"),
new AtomicLong());
assertFalse(s.collaborators().containsKey("term_ops"),
"a dead collaborator tab must never resolve as a live collaborator");
}
/**
* Both kinds are matched in a single pass over the same tab list — not a second scanner, not a
* second scan. Proven by herdr call count: scanning one lead tab and one collaborator tab in the
* same instance costs exactly as many calls as scanning two lead tabs in {@link #twoLeads()}.
*/
@Test
void leadsAndCollaboratorsAreMatchedInOnePassOverTheSameScan() {
TopologyHerdr oneOfEach = new TopologyHerdr()
.workspace("w1", "main")
.tab("w1:t1", "w1", "lead: opus-5.0")
.tab("w1:t2", "w1", "collab: ops")
.pane("w1:p1", "w1:t1", "term_opus")
.pane("w1:p2", "w1:t2", "term_ops");
LeadTabScanner s = scannerWithCollaborators(oneOfEach, Map.of("lead: opus-5.0", "opus-5.0"),
Map.of("collab: ops", "ops"), new AtomicLong());
assertEquals(Map.of("term_opus", "opus-5.0"), s.get());
assertEquals(Map.of("term_ops", "ops"), s.collaborators());
TopologyHerdr twoLeadsBaseline = twoLeads();
scanner(twoLeadsBaseline, twoLeadsConfigured(), new AtomicLong()).get();
assertEquals(twoLeadsBaseline.calls, oneOfEach.calls,
"one lead tab + one collaborator tab must cost exactly as many herdr calls as two "
+ "lead tabs — proof this is one pass, not a second scan");
}
/** Regression: with no collaborators configured, every existing lead-only behaviour is unchanged. */
@Test
void anEmptyCollaboratorMapLeavesCollaboratorsEmptyAndGetUnaffected() {
LeadTabScanner s = scannerWithCollaborators(twoLeads(), twoLeadsConfigured(), Map.of(),
new AtomicLong());
assertEquals(Map.of(), s.collaborators());
assertEquals(Map.of("term_opus", "opus-5.0", "term_gpt", "gpt-sol-5.6"), s.get());
}
}
@@ -63,6 +63,17 @@ class FleetMcpAuthzTest {
/** A fully wired FleetMcp on fakes — constructing it is itself part of what is under test. */
private FleetMcp mcp(boolean enforce) {
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> 999_999);
return mcp(enforce, CallerResolver.withLeadsAndMembers(identity, false, null,
Map::of, new MemberRegistry(null)));
}
/**
* As {@link #mcp(boolean)}, with an explicit {@link CallerResolver} — so a test can wire known
* leads/collaborators and drive {@code denyFor}'s real {@code knownLeadOrCollaborator()}
* classifier instead of the default empty one.
*/
private FleetMcp mcp(boolean enforce, CallerResolver callers) {
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
@@ -83,9 +94,7 @@ class FleetMcpAuthzTest {
// to be omitted to reach "legacy" is now always real, and AuthorizationMode is the
// separate, explicit choice that governs enforcement.
mcp = new FleetMcp(messages, workers, sessions, identity, sessions.asPresence(),
new PrimaryRegistry(null),
CallerResolver.withLeadsAndMembers(identity, false, null,
Map::of, new MemberRegistry(null)),
new PrimaryRegistry(null), callers,
enforce ? FleetMcp.AuthorizationMode.ENFORCED : FleetMcp.AuthorizationMode.UNENFORCED,
metrics, FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), null, FleetMcp.OutageSource.none(),
@@ -97,6 +106,7 @@ class FleetMcpAuthzTest {
private static final Principal WORKER_A = Principal.worker("term_a", 200);
private static final Principal ANON = Principal.anonymous();
private static final Principal ARCH_DESIGN = Principal.architect("lead-designer", "term_design", 400);
private static final Principal COLLABORATOR = Principal.collaborator("ops", "term_collab", 600);
// --- the table, enforced on THIS path too ---------------------------------------------------
@@ -226,6 +236,39 @@ class FleetMcpAuthzTest {
"the caller IS authenticated — it is just not the right role");
}
/**
* {@code denyFor} passes the real production classifier, not a test-supplied one — no
* terminal is recognised as a configured lead or collaborator, so a collaborator's SEND is
* refused over MCP.
*/
@Test
void aCollaboratorMayNotSendOverMcpWithTheRealProductionClassifier() {
FleetMcp m = mcp(true);
McpSchema.CallToolResult denied = m.denyFor(COLLABORATOR, Authz.Action.SEND, "term_lead");
assertNotNull(denied, "no terminal is recognised as a lead or collaborator yet");
assertTrue(denied.isError());
}
/**
* fleetd #669 Unit D: wires a real {@link CallerResolver} with a known lead and a known
* collaborator tab, and leaves a spawned member's own terminal recognised by neither map — so a
* collaborator's SEND reaches both named peers and is refused for the spawned member's terminal,
* over MCP's {@code denyFor}.
*/
@Test
void aCollaboratorMaySendToAKnownLeadOrCollaboratorButNotToASpawnedMembersTerminal() {
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> 999_999);
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null,
() -> Map.of("term_lead_known", "lead-x"), new MemberRegistry(null),
t -> null, () -> Map.of("term_collab_known", "ops2"));
FleetMcp m = mcp(true, callers);
assertNull(m.denyFor(COLLABORATOR, Authz.Action.SEND, "term_lead_known"));
assertNull(m.denyFor(COLLABORATOR, Authz.Action.SEND, "term_collab_known"));
assertNotNull(m.denyFor(COLLABORATOR, Authz.Action.SEND, "term_a"),
"a spawned member's own terminal must stay unreachable, even once the classifier is real");
}
@Test
void theLegacyConstructorLeavesTheGateOpen() {
// The 22 pre-existing FleetMcpTest cases rely on no authorization being enforced.
@@ -1837,6 +1837,32 @@ class FleetMcpTest {
assertTrue(out.contains("\"sessionId\":\"term_design\""), out);
}
/**
* A collaborator reports its own role and name, never the {@code leader} key a lead gets —
* {@code role} already reads {@code "collaborator"}, so a {@code leader} key alongside it
* would be self-contradicting. Control: the same call shape fed a named lead must still carry
* {@code leader}, so this is not passing because the key stopped being emitted for everyone.
*/
@Test
void whoamiReportsACollaboratorWithNoLeaderKeyButALeadStillGetsOne() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = sessionManager(h, "http://gx00.gw:8000", Set.of("gx00.gw"));
McpSchema.CallToolResult collabRes = FleetMcp.whoami(
Principal.collaborator("ops", "term_collab", 700), sessions);
assertNotEquals(Boolean.TRUE, collabRes.isError());
String collabOut = textOf(collabRes);
assertTrue(collabOut.contains("\"role\":\"collaborator\""), collabOut);
assertTrue(collabOut.contains("\"collaborator\":\"ops\""), collabOut);
assertTrue(collabOut.contains("\"sessionId\":\"term_collab\""), collabOut);
assertFalse(collabOut.contains("leader"), collabOut);
McpSchema.CallToolResult leadRes = FleetMcp.whoami(
Principal.leader("opus", "term_lead", 100), sessions);
String leadOut = textOf(leadRes);
assertTrue(leadOut.contains("\"leader\":\"opus\""), leadOut);
}
/**
* CB-548: an architect SEND delegates as its own pane (recording the per-target delegation) but
* must NEVER become the legacy singleton "primary" fallback — the per-target map does not cure
@@ -60,13 +60,25 @@ class MessageServiceTest {
inbox.own(T);
}
/**
* A send budget large enough that a test's own setup — {@link #awaitWaiting()} plus whatever
* status transitions it drives afterward — can never compete with it for the same clock. A test
* that needs {@code send.get(...)}'s own window to be the only timing bound it depends on uses
* {@link #sendAsync(String, long)} with this value instead of the default 5000 ms.
*/
private static final long GENEROUS_SEND_BUDGET_MILLIS = 30_000;
/** Run {@code send} on a background thread; the current thread drives the worker's turn. */
private CompletableFuture<MessageService.Reply> sendAsync() {
return sendAsync("do the task");
}
private CompletableFuture<MessageService.Reply> sendAsync(String content) {
return CompletableFuture.supplyAsync(() -> messages.send(T, content, 5000));
return sendAsync(content, 5000);
}
private CompletableFuture<MessageService.Reply> sendAsync(String content, long timeoutMillis) {
return CompletableFuture.supplyAsync(() -> messages.send(T, content, timeoutMillis));
}
private void awaitWaiting() throws InterruptedException {
@@ -80,7 +92,7 @@ class MessageServiceTest {
@Test
void completionFallbackResolvesATurnThatNeverCalledFleetReply() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
CompletableFuture<MessageService.Reply> send = sendAsync("do the task", GENEROUS_SEND_BUDGET_MILLIS);
awaitWaiting();
herdr.readText("$ prompt"); // pre-turn pane: no answer yet (baseline reference)
@@ -96,6 +108,32 @@ class MessageServiceTest {
assertTrue(reply.completed(), "a scraped completion still counts as completed");
}
/**
* Pins {@link #GENEROUS_SEND_BUDGET_MILLIS} as the budget {@link
* #completionFallbackResolvesATurnThatNeverCalledFleetReply} depends on. A 5500 ms delay between
* {@link #awaitWaiting()} and the status transitions that drive completion stands in for a loaded
* machine's setup overhead — comfortably past the 5000 ms budget this send no longer uses, and
* still well inside this method's own 30 000 ms budget. The only clock this test depends on is
* {@code send.get}'s own 10 s window.
*/
@Test
void completionFallbackSurvivesASlowHarnessBecauseItsSendBudgetIsNotTheBindingClock() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync("do the task", GENEROUS_SEND_BUDGET_MILLIS);
awaitWaiting();
Thread.sleep(5500);
herdr.readText("$ prompt");
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
herdr.readText("BUILD GREEN: 391 files");
injector.onStatus(T, AgentStatus.IDLE);
MessageService.Reply reply = send.get(10, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.COMPLETED_UNREPLIED, reply.outcome(),
"a slow harness must not be mistaken for a timed-out delivery");
}
@Test
void completionFallbackReplacesAnEchoedInjectedBriefWithNoReportOutcome() throws Exception {
String brief = "Implement the requested change. ".repeat(20);
@@ -1,8 +1,12 @@
package dev.ltms.fleet.rest;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.spi.ILoggingEvent;
import com.fasterxml.jackson.databind.ObjectMapper;
import dev.ltms.fleet.auth.CallerResolver;
import dev.ltms.fleet.auth.Authz;
import dev.ltms.fleet.auth.MemberRegistry;
import dev.ltms.fleet.auth.Principal;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
@@ -18,6 +22,7 @@ import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.session.FakeWorktrees;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.testing.CapturedLog;
import io.javalin.Javalin;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
@@ -28,10 +33,13 @@ import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Set;
import java.util.function.Predicate;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
import java.util.stream.Collectors;
@@ -144,6 +152,38 @@ class FleetAppAuthTest {
assertEquals(Authz.Action.ANSWER, FleetApp.routeAction("POST /sessions/{id}/message", "turn-1"));
}
/**
* fleetd #689: {@code answerGatePasses} is the second, conditional gate behind {@code
* sendMessage}'s coarse {@link Authz.Action#SEND} check. With the {@code ANSWER} grant denied,
* a {@code turnId}-bearing request is refused while a plain one still passes — and the denied
* permit is queried only for the {@code turnId} case, never for the plain one, which is what
* proves this is a genuinely separate, conditional check rather than the {@code SEND} check
* renamed or an unconditional call whose result is ignored. Flipping only the {@code ANSWER}
* grant to allowed then flips only the {@code turnId} shape's outcome.
*/
@Test
void answerGatePassesOnlyWhenTurnIdAbsentOrAnswerGranted() {
List<Authz.Action> queried = new ArrayList<>();
Predicate<Authz.Action> denyAnswer = action -> {
queried.add(action);
return false;
};
assertFalse(FleetApp.answerGatePasses("turn-1", denyAnswer),
"ANSWER denied ⇒ the turnId shape is refused");
assertEquals(List.of(Authz.Action.ANSWER), queried,
"the ANSWER grant, specifically, must be the one consulted");
queried.clear();
assertTrue(FleetApp.answerGatePasses(null, denyAnswer),
"no turnId ⇒ the plain shape passes even though ANSWER is denied");
assertTrue(FleetApp.answerGatePasses(" ", denyAnswer), "a blank turnId is treated as absent");
assertEquals(List.of(), queried, "the plain shape must never consult the permit at all");
assertTrue(FleetApp.answerGatePasses("turn-1", action -> true),
"flipping only the ANSWER grant to allowed flips only the turnId shape's outcome");
}
private static Set<String> routesTheServerRegisters() {
try {
String source = Files.readString(REST_SOURCE).lines()
@@ -163,6 +203,74 @@ class FleetAppAuthTest {
}
}
/**
* {@code permitsFor} is the exact decision {@link FleetApp#allow} makes, passing the real
* production classifier rather than a test-supplied one — built from an empty {@link
* CallerResolver}, so no terminal is recognised as a configured lead or collaborator and a
* collaborator's SEND is refused through the REST gate.
*/
@Test
void aCollaboratorMayNotSendOverRestWithTheRealProductionClassifier() {
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(new FakeHerdr()), _ -> 700L);
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null,
Map::of, new MemberRegistry(null));
Principal collaborator = Principal.collaborator("ops", "term_collab", 700);
assertFalse(FleetApp.permitsFor(collaborator, Authz.Action.SEND, "term_lead",
callers.knownLeadOrCollaborator()),
"no terminal is recognised as a lead or collaborator yet");
}
/**
* fleetd #669 Unit D: wires a real {@link CallerResolver} with a known lead and a known
* collaborator tab, and a spawned member's own terminal recognised by neither map. A
* collaborator's SEND reaches the known lead and the known collaborator, and is refused for the
* spawned member's terminal — over the REST route, not just the unit-level classifier, so a
* test covering only MCP cannot leave this route open.
*/
@Test
void aCollaboratorMaySendToAKnownLeadOrCollaboratorButNotToASpawnedMembersTerminalOverRest() throws Exception {
int port = startWithRealClassifier(FakeHerdr.WORKER_PID,
Map.of("term_lead_known", "lead-x"), Map.of("term_a", "ops2"));
HttpResponse<String> toLead = send(port, "POST", "/sessions/term_lead_known/message",
"{\"content\":\"hi\",\"wait\":false}", null);
assertEquals(202, toLead.statusCode(), toLead.body());
HttpResponse<String> toSpawnedMembersTerminal = send(port, "POST", "/sessions/term_worker/message",
"{\"content\":\"hi\",\"wait\":false}", null);
assertEquals(403, toSpawnedMembersTerminal.statusCode(), toSpawnedMembersTerminal.body());
}
/**
* As {@link #start}, but with explicit lead/collaborator maps and no spawned-member roster, so
* a test can wire the real {@link CallerResolver#knownLeadOrCollaborator()} classifier instead
* of the default empty one.
*/
private int startWithRealClassifier(long pid, Map<String, String> leadTerminals,
Map<String, String> collaboratorTerminals) {
FakeHerdr herdr = new FakeHerdr();
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());
Injector injector = new Injector(agents);
MessageService messages = new MessageService(agents, injector, new Rendezvous());
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> pid);
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null,
() -> leadTerminals, new MemberRegistry(null), t -> null, () -> collaboratorTerminals);
metrics = FleetMetrics.create(sessions, new dev.ltms.fleet.msg.InMemoryReplyInbox());
app = new FleetApp(herdr, workers, sessions, messages, sessions.asPresence(), null,
callers, metrics).build().start("127.0.0.1", 0);
return app.port();
}
// --- loopback-trust: the caller is the primary -------------------------------------------
@Test
@@ -223,6 +331,92 @@ class FleetAppAuthTest {
"resolving another session's blocked question would be a worker escalating too");
}
/**
* fleetd #689: a caller refused the coarse {@link Authz.Action#SEND} grant is refused on
* {@code SEND} specifically, even on the {@code turnId}-bearing shape that otherwise raises
* the check to {@link Authz.Action#ANSWER} — proving {@code turnId} was never read from the
* body before the refusal (reading it would have changed which action is named in the 403).
* The same caller refused with no body at all gets the identical detail, which could not hold
* if the decision depended on anything read from the body. Control: a caller who IS granted
* reaches past the gate and the body is used normally.
*/
@Test
void aDeniedCallerIsRefusedOnSendEvenWithATurnIdBodyAndNeverReadsTheBody() throws Exception {
int workerPort = start(FakeHerdr.WORKER_PID, false, null); // denied: not primary/architect
HttpResponse<String> withTurnId = send(workerPort, "POST", "/sessions/term_b/message",
"{\"turnId\":\"turn-1\",\"content\":\"hi\"}", null);
assertEquals(403, withTurnId.statusCode());
assertTrue(withTurnId.body().contains("may not SEND"),
"the SEND check must be the one that fired, not ANSWER — ANSWER would only be "
+ "reachable by having already read turnId out of the body");
HttpResponse<String> noBody = send(workerPort, "POST", "/sessions/term_b/message", null, null);
assertEquals(403, noBody.statusCode());
assertTrue(noBody.body().contains("may not SEND"),
"refused identically with no body at all — the refusal cannot depend on body content");
// Control: a primary IS granted SEND, so the same turnId body is read and acted on —
// reaching messages.answer, which reports this unknown turnId as a stale one.
int primaryPort = start(999_999, false, null);
HttpResponse<String> granted = send(primaryPort, "POST", "/sessions/term_b/message",
"{\"turnId\":\"turn-1\",\"content\":\"hi\"}", null);
assertEquals(409, granted.statusCode());
assertTrue(granted.body().contains("stale_turn"), "a granted caller's body IS read and acted on");
}
/**
* fleetd #689 (ticket comment 18353): the only place {@code sendMessage}'s call to {@code
* answerGatePasses} is observable is the audit trail — {@code allow()} logs an {@code
* "allowed"} entry for every granted action except {@code READ}/{@code METRICS}/{@code
* TASK_READ}, and {@code ANSWER} is none of those. A granted {@code turnId} request must
* therefore log both a {@code SEND} and an {@code ANSWER} entry; a granted plain request must
* log {@code SEND} alone. A unit test of the extracted helper pins the helper; this pins the
* call site — deleting the {@code answerGatePasses} call from {@code sendMessage} leaves the
* helper's own test green but turns this one red.
*/
@Test
void aGrantedTurnIdRequestAuditsBothSendAndAnswerButAPlainRequestAuditsSendAlone() throws Exception {
int port = start(999_999, false, null); // primary: granted both SEND and ANSWER
ObjectMapper mapper = new ObjectMapper();
try (CapturedLog audit = CapturedLog.at("audit", Level.INFO)) {
send(port, "POST", "/sessions/term_b/message",
"{\"turnId\":\"turn-1\",\"content\":\"hi\"}", null);
List<String> allowed = allowedActions(audit, mapper);
assertTrue(allowed.contains("SEND"),
"a turnId request must still clear the coarse SEND grant first");
assertTrue(allowed.contains("ANSWER"),
"a turnId request must ALSO clear the ANSWER grant — this is the call site itself");
}
try (CapturedLog audit = CapturedLog.at("audit", Level.INFO)) {
send(port, "POST", "/sessions/term_b/message",
"{\"content\":\"hi\",\"timeoutMs\":50}", null);
List<String> allowed = allowedActions(audit, mapper);
assertEquals(List.of("SEND"), allowed,
"a plain request must log SEND and nothing else — ANSWER is conditional on "
+ "turnId, not something every request happens to log");
}
}
private static List<String> allowedActions(CapturedLog audit, ObjectMapper mapper) {
return audit.events().stream()
.map(ILoggingEvent::getFormattedMessage)
.map(line -> {
try {
return mapper.readTree(line);
} catch (Exception e) {
throw new AssertionError("audit line is not valid JSON: " + line, e);
}
})
.filter(n -> "allowed".equals(n.path("outcome").asText()))
.map(n -> n.path("action").asText())
.toList();
}
// --- token mode ---------------------------------------------------------------------------
@Test
@@ -675,6 +675,26 @@ class FleetAppTest {
assertEquals(400, postMessage(port, "{}").statusCode());
}
/**
* fleetd #689: a body that fails to parse is rejected with 400 before {@code turnId} is ever
* read from it, so it reaches neither {@code messages.answer} (which needs a {@code turnId})
* nor {@code messages.send} — confirmed here for {@code send} by the fake agent's idle status,
* which would otherwise make an immediate {@code agent.prompt} delivery observable.
*/
@Test
void malformedBodyReturns400AndNeverReachesSendOrAnswer() throws Exception {
FakeHerdr herdr = new FakeHerdr().agentStatus("idle"); // idle ⇒ send would deliver right away if reached
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
HttpResponse<String> res = postMessage(port, "not json at all");
assertEquals(400, res.statusCode());
JsonNode err = mapper.readTree(res.body());
assertEquals("bad_request", err.get("error").asText());
assertEquals("body must be JSON", err.get("detail").asText());
assertFalse(herdr.called("agent.prompt"),
"a malformed body must never reach messages.send's delivery");
}
@Test
void sessionStatusReportsLiveAgentStatus() throws Exception {
FakeHerdr herdr = new FakeHerdr().agentStatus("blocked");
+22 -5
View File
@@ -213,6 +213,17 @@ map_masked_lines() {
done < "$file"
}
# Masks every `scheme://user:pass@host` userinfo on one line of text, replacing just that
# userinfo with `<redacted>` and leaving the rest of the line untouched, byte for byte. The
# pattern stops at the first `/`, whitespace, or `@` reached after `://` — a URI's userinfo
# component cannot contain any of those three characters — so a URI with no userinfo, followed
# later on the same line by an unrelated `@`, never matches. The `g` flag matters: a line can
# carry more than one URI. Shared by every caller that prints a line which may hold a
# credentialed URI, so the bound lives in exactly one place.
mask_url_userinfo() {
printf '%s\n' "$1" | sed -E 's#://[^@/[:space:]]*@#://<redacted>@#g'
}
redact() {
local old_file="$1" new_file="$2"
local line prefix content indent lead key old_line=0 new_line=0 in_hunk=0
@@ -270,7 +281,7 @@ redact() {
continue
fi
fi
printf '%s\n' "$line" | sed -E 's#://[^@]*@#://<redacted>@#g'
mask_url_userinfo "$line"
done
[ "$saved_nocasematch" = 1 ] || shopt -u nocasematch
}
@@ -580,6 +591,12 @@ install_candidate() {
# -------------------------------------------------------------------------------- the report path
#
# Masks basic-auth userinfo (scheme://user:pass@host) in a daemon verdict line before it reaches
# the terminal.
mask_verdict_userinfo() {
mask_url_userinfo "$1"
}
# Prints the literal command the operator (or a test) can run to restore the backup by hand — the
# absolute path to THIS script plus the overrides actually in force, so it works from any cwd.
restore_command_line() {
@@ -599,8 +616,8 @@ restore_and_confirm() {
ok "restored from $backup"
if wait_for_verdict "$LOG" "$mark2" "$WAIT_SECONDS"; then
case "$VERDICT_KIND" in
refused) warn "the RESTORE was also refused by the daemon: $VERDICT_LINE" ;;
*) ok "restore confirmed: $VERDICT_LINE" ;;
refused) warn "the RESTORE was also refused by the daemon: $(mask_verdict_userinfo "$VERDICT_LINE")" ;;
*) ok "restore confirmed: $(mask_verdict_userinfo "$VERDICT_LINE")" ;;
esac
else
warn "the restore is on disk, but no confirming verdict line appeared within ${WAIT_SECONDS}s"
@@ -616,7 +633,7 @@ report_outcome() {
say "waiting for the daemon's verdict (up to ${WAIT_SECONDS}s)"
if wait_for_verdict "$LOG" "$mark" "$WAIT_SECONDS"; then
kind="$VERDICT_KIND"; line="$VERDICT_LINE"
kind="$VERDICT_KIND"; line="$(mask_verdict_userinfo "$VERDICT_LINE")"
else
kind="none"
fi
@@ -673,7 +690,7 @@ check_mode() {
local verdict
verdict="$(last_verdict_line "$LOG")"
if [ -n "$verdict" ]; then
ok "last verdict in log: $verdict"
ok "last verdict in log: $(mask_verdict_userinfo "$verdict")"
else
warn "no reload verdict line found in $LOG"
fi
+131
View File
@@ -233,6 +233,66 @@ test_redaction_holds() {
assert_contains "weight" "$RUN_OUTPUT" "a diff must have been demonstrably printed at all"
}
# redact()'s key-name filter only inspects the KEY, so a diff line whose key does not match
# TOKEN|SECRET|PASSWORD|PASSWD|PASSPHRASE|CREDENTIAL|URI|_KEY still reaches the final userinfo
# sed even when its VALUE holds a credentialed URI. "note" is not a sensitive key name, so this
# line must fall all the way through to that sed, not the earlier whole-value branch. The
# trailing prose on both sides of the userinfo is a positive control: it proves the line reached
# the userinfo sed (which touches only the userinfo) rather than the earlier branch (which would
# have replaced the whole value with a bare "<redacted>" and dropped the prose).
test_diff_line_userinfo_is_masked_with_positive_control() {
local dir
dir="$(new_fixture)"
start_run "$dir" 5 --set '.profiles.sonnet.note=see amqp://alice:wonderland@rabbit.local:5672/vhost for details'
sleep 1
printf 'config reloaded\n' >> "$dir/fleetd.out"
collect_run "$dir"
assert_equals 0 "$RUN_RC" "diff-userinfo-case reload exit code"
assert_not_contains "alice:wonderland" "$RUN_OUTPUT" "the userinfo must never reach the output"
assert_contains "amqp://<redacted>@rabbit.local:5672/vhost" "$RUN_OUTPUT" \
"the userinfo must be MASKED, not deleted — the rest of the value must survive"
assert_contains "note:" "$RUN_OUTPUT" "the key name must still reach the output"
assert_contains "see " "$RUN_OUTPUT" "prose BEFORE the userinfo must still reach the output"
assert_contains "for details" "$RUN_OUTPUT" "prose AFTER the userinfo must still reach the output"
}
# A diff line can hold a URL with no userinfo, followed later on the same line by an unrelated @
# (free text in a string value, for example an email address). The line must pass through the
# userinfo sed byte for byte: the match must stop at the end of the URL and must not treat the
# later @ as a second userinfo delimiter.
test_diff_line_uri_without_userinfo_survives_a_later_at_sign() {
local dir
dir="$(new_fixture)"
start_run "$dir" 5 --set '.profiles.sonnet.note2=see https://docs.local/guide and mail ops@example.com'
sleep 1
printf 'config reloaded\n' >> "$dir/fleetd.out"
collect_run "$dir"
assert_equals 0 "$RUN_RC" "diff-no-userinfo-with-later-at-sign reload exit code"
assert_contains "note2: see https://docs.local/guide and mail ops@example.com" "$RUN_OUTPUT" \
"a URL with no userinfo plus a later @ on the same line must pass through byte for byte"
}
# Two credentialed URIs on one diff line must both be masked — the g flag matters.
test_diff_line_masks_multiple_userinfo_with_g_flag() {
local dir
dir="$(new_fixture)"
start_run "$dir" 5 --set '.profiles.sonnet.note3=amqp://u1:p1@host1/vhost1 and amqp://u2:p2@host2/vhost2'
sleep 1
printf 'config reloaded\n' >> "$dir/fleetd.out"
collect_run "$dir"
assert_equals 0 "$RUN_RC" "diff-two-userinfo-on-one-line reload exit code"
assert_not_contains "u1:p1" "$RUN_OUTPUT" "the first userinfo must never reach the output"
assert_not_contains "u2:p2" "$RUN_OUTPUT" "the second userinfo must never reach the output"
assert_contains "amqp://<redacted>@host1/vhost1" "$RUN_OUTPUT" "the first URI must be masked"
assert_contains "amqp://<redacted>@host2/vhost2" "$RUN_OUTPUT" "the second URI must be masked"
}
# ------------------------------------------------------- acceptance criterion 9: forgotten value
# `--set .a.b=` is a plausible typo (the value simply forgotten), and it must be refused outright
# rather than silently nulling the field — a null numeric field falls back to its default, which
@@ -629,6 +689,65 @@ test_refusal_shape_from_parse_failure_wording_is_recognised() {
assert_equals 4 "$RUN_RC" "the parse-failure refusal shape must also exit 4, not be read as silence"
}
# A verdict line carrying a credentialed URI has its userinfo masked, with a positive control
# proving the rest of the line still reaches the output unchanged.
test_verdict_userinfo_is_masked_with_positive_control() {
local dir
dir="$(new_fixture)"
start_run "$dir" 5 --set '.profiles.sonnet.weight=4'
sleep 1
printf 'config reload from %s refused, keeping the running config: refusing to start: malformed pattern — profiles.local.errorPattern ("amqp://user:hunter2@host/vhost"): Unclosed character class near index 8\n' \
"$dir/fleetd.yaml" >> "$dir/fleetd.out"
collect_run "$dir"
assert_equals 4 "$RUN_RC" "refusal-with-userinfo exit code"
assert_not_contains "user:hunter2" "$RUN_OUTPUT" "the userinfo must never reach the output"
assert_contains "amqp://<redacted>@host/vhost" "$RUN_OUTPUT" \
"the userinfo must be MASKED, not deleted — the rest of the quoted value must survive"
# Positive control: the diagnostic prose on both sides of the userinfo must still reach the
# output. Without this, a mutant that drops the whole verdict line would pass identically.
assert_contains "malformed pattern" "$RUN_OUTPUT" "prose BEFORE the userinfo must still reach the output"
assert_contains "Unclosed character class near index 8" "$RUN_OUTPUT" \
"prose AFTER the userinfo must still reach the output"
}
# An ordinary refusal line quotes the offending pattern, not a credential, and must survive byte
# for byte: the rewrite is scoped to userinfo only, and the quoted pattern is the detail an
# operator needs to fix the refusal.
test_ordinary_refusal_line_passes_through_unchanged() {
local dir real_line
dir="$(new_fixture)"
real_line="config reload from $dir/fleetd.yaml refused, keeping the running config: refusing to start: malformed pattern — profiles.local.errorPattern (\"[unclosed\"): Unclosed character class near index 8"
start_run "$dir" 5 --set '.profiles.sonnet.weight=4'
sleep 1
printf '%s\n' "$real_line" >> "$dir/fleetd.out"
collect_run "$dir"
assert_equals 4 "$RUN_RC" "ordinary refusal exit code"
assert_contains "$real_line" "$RUN_OUTPUT" \
"an ordinary refusal with no userinfo must pass through byte for byte, unchanged"
}
# A verdict line can hold a URI with NO userinfo and a later, unrelated @ further on in the same
# line (an email address in diagnostic prose, for example). The rewrite must stop at the end of
# the URI and must not treat the later @ as a second userinfo delimiter.
test_uri_without_userinfo_survives_a_later_at_sign() {
local dir real_line
dir="$(new_fixture)"
real_line="config reload from $dir/fleetd.yaml refused, keeping the running config: broker.uri amqp://broker.local/vhost unreachable, contact ops@example.com"
start_run "$dir" 5 --set '.profiles.sonnet.weight=4'
sleep 1
printf '%s\n' "$real_line" >> "$dir/fleetd.out"
collect_run "$dir"
assert_equals 4 "$RUN_RC" "no-userinfo-with-later-at-sign exit code"
assert_contains "$real_line" "$RUN_OUTPUT" \
"a URI with no userinfo plus a later @ in the same line must pass through byte for byte"
}
# --set runs yq over the whole candidate. It warns when that changes more lines than the requested
# pairs, but a simple file with only the intended changed line must stay quiet.
new_fixture_reformat_sensitive() {
@@ -694,6 +813,12 @@ echo "== acceptance criterion 6: the marker works =="
test_marker_skips_lines_before_it
echo "== acceptance criterion 7 (+13: redaction is proven to have run) =="
test_redaction_holds
echo "== fleetd #692: a diff line's userinfo is masked, rest of the value survives =="
test_diff_line_userinfo_is_masked_with_positive_control
echo "== fleetd #692: a diff line's URI with no userinfo survives a later @ in the line =="
test_diff_line_uri_without_userinfo_survives_a_later_at_sign
echo "== fleetd #692: two userinfo URIs on one diff line are both masked =="
test_diff_line_masks_multiple_userinfo_with_g_flag
echo "== acceptance criterion 9: a forgotten value refuses and installs nothing =="
test_forgotten_value_refuses_and_installs_nothing
echo "== acceptance criterion 10: an explicit clear writes a bare null =="
@@ -720,6 +845,12 @@ echo "== extra: --check is read-only and always exits 0 =="
test_check_is_read_only_and_exits_zero
echo "== extra: the parse-failure refusal shape is also recognised =="
test_refusal_shape_from_parse_failure_wording_is_recognised
echo "== verdict-redaction criteria 2+3: verdict userinfo is masked, rest of line survives =="
test_verdict_userinfo_is_masked_with_positive_control
echo "== verdict-redaction criterion 4: an ordinary refusal passes through unchanged =="
test_ordinary_refusal_line_passes_through_unchanged
echo "== fleetd #638: a URI with no userinfo survives a later @ in the same line =="
test_uri_without_userinfo_survives_a_later_at_sign
echo "== acceptance criterion 17: --set warns about yq formatting churn =="
test_set_warns_when_yq_reformats_extra_lines
echo "== acceptance criterion 18: --set stays quiet without formatting churn =="