Compare commits

...

10 Commits

Author SHA1 Message Date
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 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 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
20 changed files with 1146 additions and 95 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 |
@@ -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
@@ -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
}
@@ -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
@@ -266,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}.
@@ -279,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
@@ -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();
}
}
}
@@ -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"));
}
}
@@ -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();
@@ -464,4 +471,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
@@ -6,6 +6,7 @@ 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;
@@ -202,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
+13 -2
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
}
@@ -583,7 +594,7 @@ install_candidate() {
# Masks basic-auth userinfo (scheme://user:pass@host) in a daemon verdict line before it reaches
# the terminal.
mask_verdict_userinfo() {
printf '%s\n' "$1" | sed -E 's#://[^@/[:space:]]*@#://<redacted>@#g'
mask_url_userinfo "$1"
}
# Prints the literal command the operator (or a test) can run to restore the backup by hand — the
+66
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
@@ -753,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 =="