Compare commits

...

10 Commits

Author SHA1 Message Date
Dai Ha 11eccc3a1b fleetd #705: gate fleet_poll's ticket lookup by the creating caller's terminal
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 1m45s
A ticket id is a plain sequential counter, so any session holding
TASK_READ could walk task-1, task-2, ... and read another session's
delegation reply. sendAsync now records the creating caller's terminal
on the Task, and poll refuses a caller whose terminal differs from it.
A caller with no terminal (the unnamed primary) is still allowed
through regardless, since it never carries a herdr pane to compare.
2026-10-04 06:50:26 +02:00
Dai Ha 1fd2cfa716 Merge PR #708: fleetd #669 — fleetd.example.yaml promised a defence that is not wired
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 1m10s
CI / build (push) Failing after 2m15s
The example file said two things stop the tab-name convention from becoming a
way to claim leadership, and named workspace exclusion as the first. Production
builds the scanner with an empty excludedWorkspaceLabels, so that defence does
not exist. It now names the two that do: the startup refusal of a colliding
tabPrefix/tabLabel, and CallerResolver asking the live spawned-member roster
before any tab map.

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

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

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

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

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

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

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

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

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

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

One reviewer fanned out against the diff, briefed from the diff rather
than the implementer's rationale. It reported no issue.
2026-10-04 02:19:44 +02:00
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
15 changed files with 640 additions and 68 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 |
+22 -13
View File
@@ -62,11 +62,12 @@ bind:
#
# Leads are configured under `fleet.leaders:` — see THE FLEET further down.
#
# Two things stop the tab-name convention from becoming a way to claim leadership: the configured
# member spaces are excluded from the scan, so nothing fleetd places can land in a matching tab;
# and startup REFUSES a `tabPrefix` that the fleet tabLabel template, or any per-profile `tabLabel`
# override, also matches — so the two namespaces cannot overlap by accident. The label is a NAME,
# never a capability: what a pane may do is decided by the role the daemon resolves for it.
# Two things stop the tab-name convention from becoming a way to claim leadership: startup REFUSES
# a `tabPrefix` that the fleet tabLabel template, or any per-profile `tabLabel` override, also
# matches, so the two namespaces cannot overlap by accident; and the CallerResolver asks the live
# spawned-member roster BEFORE any tab map, so a live member is never mistaken for a lead no matter
# what its tab says. The label is a NAME, never a capability: what a pane may do is decided by the
# role the daemon resolves for it.
# CB-551: IDLE-LEAD HEARTBEAT — nudge the single lead back to work when it has been continuously
# idle (no open fleet_send driving it) past the quiet period. The fleet is one lead + architects +
@@ -659,9 +660,10 @@ fleet:
# tabPrefix: "lead:" # only used to guard against a worker tabLabel colliding with
# # this convention at startup; plays no part in matching a lead
# scanIntervalSeconds: 10 # rescan cadence, and the worst case before a new tab is seen
# workspace: leads # where a launched lead's tab is created (default "leads").
# # MUST NOT be a member workspace — those are excluded from the
# # scan, so a lead placed in one is never found again.
# workspace: leads # where a launched lead's tab is created (default "fleet",
# # the same shared space the members use). Sharing that space
# # with members is the normal shipped shape: the scanner tells
# # a lead from a member by the exact tab label, not by workspace.
# cwd: /path/to/repo # the launched lead's working directory (default: fleetd's own)
# kind: claude # descriptive; reported by fleet_whoami
# gpt-sol-5.6:
@@ -669,11 +671,18 @@ fleet:
# kind: opencode
# model: openai/gpt-5.6-terra
# A collaborator tab, keyed by name (fleetd #669). This block is parsed and validated today;
# nothing yet recognises or addresses the tab it names. Recognise-only, like a profile-less
# `leaders:` entry above: there is no `profile:`, no `instances:` and no `kind:`. `tab:` is
# REQUIRED and is the only field identity depends on, matched case-insensitively — the same
# GET-THE-VALUE-RIGHT warning above the `leaders:` block applies here too.
# A collaborator tab, keyed by name (fleetd #669). A pane whose tab matches resolves as the
# COLLABORATOR role: `fleet_whoami` answers `collaborator`, and the session may observe the fleet
# (`fleet_list`, `fleet_profiles`, `fleet_whoami`) and `fleet_send` to a lead or to another
# collaborator. It may NOT spawn, stop or drain anything, roll a lead's session, answer a
# member's `fleet_ask`, poll a ticket, send across hosts, or send to a spawned member's terminal.
# Each of those is refused at the gate rather than queued.
# Recognise-only, like a profile-less `leaders:` entry above: there is no `profile:`, no
# `instances:` and no `kind:`. `tab:` is REQUIRED and is the only field identity depends on,
# matched case-insensitively — the same GET-THE-VALUE-RIGHT warning above the `leaders:` block
# applies here too.
# A lead cannot discover a collaborator yet (#703): `fleet_list` has no `collaborators` key, so
# the collaborator must speak first, or pass on the `sessionId` its own `fleet_whoami` reports.
# collaborators:
# reviewer-alex:
# tab: "collab: alex"
+15 -12
View File
@@ -217,25 +217,28 @@ public final class Fleetd {
/**
* The {@link Injector}'s readiness gate (CB-534): a target is deliverable if it is a spawned
* member whose agent has connected the bridge MCP, <em>or</em> a lead.
* member whose agent has connected the bridge MCP, a lead, <em>or</em> a collaborator.
*
* <p>The gate exists for one reason — to hold a delivery out of a <em>spawned</em> member's boot
* window, where herdr already reports {@code idle} but the TUI would drop an injected paste. That
* hazard is a property of spawning. A lead is never spawned: the operator started it and named it
* (or labelled its tab) only once it was up, so there is no boot window to guard.
* (or labelled its tab) only once it was up, so there is no boot window to guard. A collaborator
* is the same way — a person's own tab, matched to a configured name, never spawned.
*
* <p>A lead is also never enrolled in {@link MemberPresence} — {@code FleetMcp} marks presence
* for every spawned member (worker and architect), deliberately, since that map doubles as the
* member roster's availability signal and a lead counted there would show up as an available
* member. So without the second disjunct a lead is permanently un-deliverable: every
* lead→lead send sat on the gate for {@code READINESS_GRACE_POLLS} (~60s) and then failed
* having never been typed into the pane.
* <p>Neither a lead nor a collaborator is ever enrolled in {@link MemberPresence} — {@code
* FleetMcp} marks presence for every spawned member (worker and architect), deliberately, since
* that map doubles as the member roster's availability signal and a lead or collaborator counted
* there would show up as an available member. So without the second and third disjuncts a lead or
* collaborator is permanently un-deliverable: every send to one sat on the gate for
* {@code READINESS_GRACE_POLLS} (~60s) and then failed having never been typed into the pane.
*
* <p>The lead set is read through the supplier on each call rather than snapshotted, so a lead
* discovered by {@code leadScan} after startup becomes deliverable without a restart.
* <p>Both sets are read through their supplier on each call rather than snapshotted, so a lead or
* collaborator discovered by {@code leadScan} after startup becomes deliverable without a restart.
*/
static Predicate<String> deliverableTo(MemberPresence presence, Supplier<Map<String, String>> leads) {
return target -> presence.isPresent(target) || leads.get().containsKey(target);
static Predicate<String> deliverableTo(MemberPresence presence, Supplier<Map<String, String>> leads,
Supplier<Map<String, String>> collaborators) {
return target -> presence.isPresent(target) || leads.get().containsKey(target)
|| collaborators.get().containsKey(target);
}
/**
@@ -143,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.
@@ -292,6 +296,7 @@ final class FleetdAssembly {
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
@@ -363,7 +368,7 @@ final class FleetdAssembly {
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
MemberPresence presence = sessions.asPresence();
TurnListener turnListener = Fleetd.turnListener(completion, sessions);
Predicate<String> deliverable = Fleetd.deliverableTo(presence, leads);
Predicate<String> deliverable = Fleetd.deliverableTo(presence, leads, collaboratorTerminals);
// fleetd #556: registration is wired directly to `completion`, not folded into the
// `turnListener` fan-out above — so it survives `sessions.onDelivered` (or any future
// listener) throwing, regardless of call order.
@@ -615,7 +615,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
// may bind to, AND what a slot already bound still grants) through its own instance of that
// same supplier shape — see MemberRegistry.live and its class doc for the binding rule:
// removing a slot revokes ARCHITECT on the bound pane's very next request, and only the slot
// OCCUPANCY survives, so the demoted session keeps its slot key until it unbinds. Only
// OCCUPANCY survives, so the demoted session keeps its slot key until it unbinds.
// fleet.leaders is frozen (Fleetd.java:281 reads cfg.fleet().leaders() off the startup
// snapshot to build both the LeadTabScanner's tab-label-to-name map, wired into
// CallerResolver.withLeadsAndMembers at Fleetd.java:620/624, and — when herdr answered —
@@ -635,6 +635,11 @@ public final class ConfigRef implements Supplier<FleetConfig> {
+ "live through that same supplier for placement AND through a separate supplier "
+ "on MemberRegistry for spawn-time identity — both already applied");
}
if (!Objects.equals(collaboratorsOf(old), collaboratorsOf(fresh))) {
changed.add("fleet: fleet.collaborators (each collaborator's tab) is read once at "
+ "startup to build the LeadTabScanner's identity map, which is not rebuilt on "
+ "reload, so a collaborator added, removed, or given a new tab: needs a restart");
}
// Kept in step with SPLIT_KEYS the same way changedColdKeys is kept in step with COLD_KEYS —
// every message here must be traceable to one of the split keys the class doc documents.
// NOTE what this does NOT prove, per the javadoc above: it does not catch a SPLIT_KEYS
@@ -650,6 +655,11 @@ public final class ConfigRef implements Supplier<FleetConfig> {
return cfg.fleet() == null ? Map.of() : cfg.fleet().leaders();
}
/** {@code cfg.fleet().collaborators()}, defensively, in case a caller hands in a non-defaulted config. */
private static Map<String, FleetConfig.Collaborator> collaboratorsOf(FleetConfig cfg) {
return cfg.fleet() == null ? Map.of() : cfg.fleet().collaborators();
}
/**
* {@link FleetConfig.Profile} record components deliberately left out of
* {@link #sameLaunchSettings} because they are read <em>live</em>, not baked in at spawn — see
@@ -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; }
@@ -490,7 +490,7 @@ public final class FleetMcp {
Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, caller);
// wait defaults to true (block for the reply); wait:false is fire-and-poll.
return Boolean.FALSE.equals(a.get("wait"))
? sendAsync(messages, target, content, onAccepted, workers.profiles())
? sendAsync(messages, target, content, onAccepted, workers.profiles(), caller)
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles());
};
// fleet_reply's identity is the CONNECTION, never an argument — so the authz check
@@ -524,7 +524,7 @@ public final class FleetMcp {
// The action depends on the ARGUMENTS, not on the tool name -- see pollAction.
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_poll", a), target);
if (denied != null) return denied;
return poll(messages, leadChannel, str(a, "ticket"), target, coordId);
return poll(messages, leadChannel, str(a, "ticket"), target, coordId, callerTerminal(exchange));
};
// CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox).
// Acking removes a reply from the inbox, so it is a drain, not a read.
@@ -956,6 +956,17 @@ public final class FleetMcp {
*/
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
Runnable onAccepted, Set<String> profiles) {
return sendAsync(messages, sessionId, content, onAccepted, profiles, null);
}
/**
* As above, recording {@code creatorTerminal} as this ticket's owner (the caller's own
* terminal, resolved from the connection) so a later {@code fleet_poll{ticket}} only hands the
* result back to that same caller — see {@link MessageService#poll(String, String)}.
*/
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
Runnable onAccepted, Set<String> profiles,
String creatorTerminal) {
if (isBlank(sessionId) || isBlank(content)) {
return error("sessionId and content are required");
}
@@ -963,7 +974,7 @@ public final class FleetMcp {
if (targetError != null) {
return targetError;
}
String ticket = messages.sendAsync(sessionId, content, onAccepted);
String ticket = messages.sendAsync(sessionId, content, onAccepted, creatorTerminal);
return text("accepted — task delegated. Poll fleet_poll with ticket=" + ticket);
}
@@ -1142,9 +1153,24 @@ public final class FleetMcp {
* exclusively of {@code ticket}/{@code target} — see {@link #pollAction}'s javadoc for why this
* is a different inbox (this daemon's own {@link LeadChannel}) that authorizes differently
* ({@link Authz.Action#COORD_READ}, primary-only) from either of the original two branches.
*
* <p>Does not check who owns {@code ticket} — see the overload that takes {@code
* callerTerminal} for that. Callers that do not resolve a caller terminal (tests, or a surface
* with no connection-based identity) use this one.
*/
static McpSchema.CallToolResult poll(MessageService messages, LeadChannel leadChannel, String ticket,
String target, String coordId) {
return poll(messages, leadChannel, ticket, target, coordId, null);
}
/**
* As above, refusing a ticket lookup whose caller's terminal differs from the terminal that
* created it — see {@link MessageService#poll(String, String)}. {@code callerTerminal} is the
* CALLING session's terminal id, resolved by the MCP layer from the connection, never a
* client-supplied value.
*/
static McpSchema.CallToolResult poll(MessageService messages, LeadChannel leadChannel, String ticket,
String target, String coordId, String callerTerminal) {
if (!isBlank(coordId)) {
return pollHeldPeerMail(leadChannel, coordId);
}
@@ -1158,7 +1184,7 @@ public final class FleetMcp {
if (isBlank(ticket)) {
return error("ticket (or target) is required");
}
MessageService.TaskView v = messages.poll(ticket);
MessageService.TaskView v = messages.poll(ticket, callerTerminal);
if (v == null) {
return error("unknown ticket: " + ticket + " (never issued, or expired)");
}
@@ -275,10 +275,19 @@ public final class MessageService {
* {@link #abandon}) can never match again regardless of this flag's value.
*/
private volatile boolean askTimedOut;
/**
* The terminal of the caller whose {@code fleet_send{wait:false}} created this ticket, or
* {@code null} when that caller had no terminal (the unnamed primary) or the ticket was
* created through an overload that does not record one. {@link #poll(String, String)}
* compares a polling caller's own terminal against this field before handing back the
* ticket's state.
*/
private final String creatorTerminal;
private Task(String ticket, String target, LongSupplier nowNanos) {
private Task(String ticket, String target, LongSupplier nowNanos, String creatorTerminal) {
this.ticket = ticket;
this.target = target;
this.creatorTerminal = creatorTerminal;
this.createdNanos = nowNanos.getAsLong();
future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong());
}
@@ -1279,8 +1288,20 @@ public final class MessageService {
* @return the ticket to poll for the eventual result
*/
public String sendAsync(String target, String content, Runnable onAccepted) {
return sendAsync(target, content, onAccepted, null);
}
/**
* As {@link #sendAsync(String, String, Runnable)}, recording {@code creatorTerminal} as this
* ticket's owner. {@link #poll(String, String)} refuses a later caller whose own terminal
* differs from this one; {@code null} records no owner (a caller with no terminal — the
* unnamed primary — is always allowed to poll the result regardless).
*
* @return the ticket to poll for the eventual result
*/
public String sendAsync(String target, String content, Runnable onAccepted, String creatorTerminal) {
String ticket = "task-" + ticketSeq.incrementAndGet();
Task task = new Task(ticket, target, nowNanos);
Task task = new Task(ticket, target, nowNanos, creatorTerminal);
tasks.put(ticket, task);
if (pushLoop != null) {
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a
@@ -1343,15 +1364,31 @@ public final class MessageService {
}
/**
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket;
* otherwise a {@link Phase#PENDING} view (with the live worker status as detail), a
* {@link Phase#DONE} view carrying the reply, or a {@link Phase#FAILED} view with the reason.
* As {@link #poll(String, String)}, with no caller terminal — the ticket's ownership is never
* checked, so this overload must only be used where the caller's identity is otherwise
* irrelevant (a test, or a surface that does not resolve a caller terminal at all).
*/
public TaskView poll(String ticket) {
return poll(ticket, null);
}
/**
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket.
* Refuses a {@code callerTerminal} that differs from the terminal that created the ticket (see
* {@link #sendAsync(String, String, Runnable, String)}) with a {@link Phase#FAILED} view that
* carries no reply text — a caller with no terminal (the unnamed primary) is never refused.
* Otherwise returns a {@link Phase#PENDING} view (with the live worker status as detail), a
* {@link Phase#DONE} view carrying the reply, or a {@link Phase#FAILED} view with the reason.
*/
public TaskView poll(String ticket, String callerTerminal) {
Task task = tasks.get(ticket);
if (task == null) {
return null;
}
if (!ownsTicket(task, callerTerminal)) {
return new TaskView(ticket, Phase.FAILED, null, null,
"forbidden: this ticket was created by a different session", null);
}
CompletableFuture<Reply> f = task.future;
if (!f.isDone()) {
Reply question = task.question;
@@ -1387,6 +1424,17 @@ public final class MessageService {
return new TaskView(ticket, Phase.FAILED, null, null, detail, null);
}
/**
* Whether {@code callerTerminal} may read {@code task}'s state. A caller with no terminal
* always may — that is the unnamed primary, resolved by token or loopback trust, which never
* carries a herdr pane and must keep reading every ticket. Otherwise the caller's terminal must
* equal the terminal recorded on the task; a task with no recorded terminal matches no
* terminal-bearing caller.
*/
private static boolean ownsTicket(Task task, String callerTerminal) {
return callerTerminal == null || callerTerminal.equals(task.creatorTerminal);
}
/**
* Test seam only — carries no production behaviour, and nothing in this class calls it;
* {@link #pruneTerminalTickets} still reads {@link Task#completedNanos} directly.
@@ -14,10 +14,12 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* CB-534: the injector's readiness gate must open for a lead as well as for a present worker.
* fleetd #669 follow-up: the same gate must also open for a collaborator, which — like a lead —
* is never enrolled in {@link MemberPresence} and never discovered by the lead scan.
*
* <p>The bug these cover was silent and slow: a lead was never marked present (only workers are), so
* every lead→lead delivery sat on the gate for the full readiness grace and failed ~60s later without
* a keystroke ever reaching the pane.
* <p>The bug these cover was silent and slow: a lead (and later a collaborator) was never marked
* present (only workers are) and never counted as a lead, so every send to one sat on the gate for
* the full readiness grace and failed ~60s later without a keystroke ever reaching the pane.
*/
class FleetDeliverabilityTest {
@@ -25,19 +27,25 @@ class FleetDeliverabilityTest {
return () -> m;
}
private static Supplier<Map<String, String>> collaborators(Map<String, String> m) {
return () -> m;
}
@Test
@DisplayName("a worker that has connected its MCP is deliverable")
void presentWorkerIsDeliverable() {
MemberPresence presence = new MemberPresence();
presence.markPresent("term_worker");
assertTrue(Fleetd.deliverableTo(presence, leads(Map.of())).test("term_worker"));
assertTrue(Fleetd.deliverableTo(presence, leads(Map.of()), collaborators(Map.of()))
.test("term_worker"));
}
@Test
@DisplayName("a worker still in its boot window is held back")
void absentWorkerIsNotDeliverable() {
assertFalse(Fleetd.deliverableTo(new MemberPresence(), leads(Map.of())).test("term_booting"));
assertFalse(Fleetd.deliverableTo(new MemberPresence(), leads(Map.of()), collaborators(Map.of()))
.test("term_booting"));
}
@Test
@@ -45,19 +53,31 @@ class FleetDeliverabilityTest {
void leadIsDeliverableWithoutPresence() {
MemberPresence presence = new MemberPresence();
Predicate<String> deliverable =
Fleetd.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")));
Fleetd.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")), collaborators(Map.of()));
assertFalse(presence.isPresent("term_lead"), "a lead is never enrolled in worker presence");
assertTrue(deliverable.test("term_lead"), "…and must be deliverable anyway");
}
@Test
@DisplayName("an unknown terminal is deliverable to neither")
@DisplayName("a collaborator is deliverable without ever being marked present or scanned as a lead")
void collaboratorIsDeliverableWithoutPresenceOrLeadStatus() {
MemberPresence presence = new MemberPresence();
Predicate<String> deliverable = Fleetd.deliverableTo(presence, leads(Map.of()),
collaborators(Map.of("term_collab", "kevin")));
assertFalse(presence.isPresent("term_collab"), "a collaborator is never enrolled in worker presence");
assertTrue(deliverable.test("term_collab"), "…and must be deliverable anyway");
}
@Test
@DisplayName("an unknown terminal is deliverable to none of presence, leads, or collaborators")
void strangerIsNotDeliverable() {
MemberPresence presence = new MemberPresence();
presence.markPresent("term_worker");
assertFalse(Fleetd.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")))
assertFalse(Fleetd.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")),
collaborators(Map.of("term_collab", "kevin")))
.test("term_stranger"));
}
@@ -65,22 +85,47 @@ class FleetDeliverabilityTest {
@DisplayName("a lead discovered after startup becomes deliverable with no restart")
void leadSetIsReadThroughOnEveryCall() {
Map<String, String> discovered = new HashMap<>();
Predicate<String> deliverable = Fleetd.deliverableTo(new MemberPresence(), leads(discovered));
Predicate<String> deliverable =
Fleetd.deliverableTo(new MemberPresence(), leads(discovered), collaborators(Map.of()));
assertFalse(deliverable.test("term_late"));
discovered.put("term_late", "gpt-sol-5.6"); // leadScan picks up a newly labelled tab
assertTrue(deliverable.test("term_late"), "the supplier must be re-read, not snapshotted");
}
@Test
@DisplayName("a collaborator discovered after startup becomes deliverable with no restart")
void collaboratorSetIsReadThroughOnEveryCall() {
Map<String, String> discovered = new HashMap<>();
Predicate<String> deliverable =
Fleetd.deliverableTo(new MemberPresence(), leads(Map.of()), collaborators(discovered));
assertFalse(deliverable.test("term_late_collab"));
discovered.put("term_late_collab", "kevin"); // the same tab scan picks up a newly labelled collaborator tab
assertTrue(deliverable.test("term_late_collab"), "the supplier must be re-read, not snapshotted");
}
@Test
@DisplayName("forgetting a torn-down worker does not strip a lead of its deliverability")
void forgetDoesNotDisarmALead() {
MemberPresence presence = new MemberPresence();
Predicate<String> deliverable =
Fleetd.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")));
Fleetd.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")), collaborators(Map.of()));
presence.forget("term_lead"); // the injector's cleanup path runs against every target
assertTrue(deliverable.test("term_lead"));
}
@Test
@DisplayName("forgetting a torn-down worker does not strip a collaborator of its deliverability")
void forgetDoesNotDisarmACollaborator() {
MemberPresence presence = new MemberPresence();
Predicate<String> deliverable = Fleetd.deliverableTo(presence, leads(Map.of()),
collaborators(Map.of("term_collab", "kevin")));
presence.forget("term_collab"); // the injector's cleanup path runs against every target
assertTrue(deliverable.test("term_collab"));
}
}
@@ -0,0 +1,165 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.msg.ReplyInbox;
import io.javalin.Javalin;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertSame;
/**
* fleetd #669 Unit E. Reaches the real {@link dev.ltms.fleet.herdr.HerdrRouter} that {@link
* FleetdAssembly#assembleAndStart} builds and wires — not a copy built for this test — and proves
* that a configured collaborator's terminal routes to the LEAD herdr daemon.
*
* <p>Two distinct {@link FakeHerdr} instances are required, the same pattern {@code
* FleetdAssemblyConnectionIdentityTest} and {@code FleetdLeadRolloverAssemblyTest} already use:
* with one client shared between {@code herdrSocket} and {@code memberHerdrSocket},
* {@code HerdrRouter} folds {@code leadAgents} and {@code memberAgents} into the same instance
* (see its constructor), and {@code agentsFor} would return that one object regardless of whether
* the collaborator map was ever consulted — invisible to a mutation of the predicate this ticket
* fixes. This test's two sockets resolve to two different fakes, so the assertion only passes when
* the collaborator's terminal is actually recognised and routed to the lead one.
*/
class FleetdAssemblyCollaboratorHerdrRoutingTest {
private static final Path LEAD_SOCKET = Path.of("/fake/lead-herdr.sock");
private static final Path MEMBER_SOCKET = Path.of("/fake/member-herdr.sock");
private static final class RecordingResourcePorts implements ResourcePorts {
final Map<Path, HerdrClient> herdrsBySocket = new LinkedHashMap<>();
Runnable shutdownHook;
@Override
public Map<String, String> environment() {
return Map.of();
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
HerdrClient client = herdrsBySocket.get(socketPath);
if (client == null) {
throw new IllegalStateException("no fake herdr registered for socket " + socketPath);
}
return client;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
return (uri, prefetch) -> new ReplyInbox() {
@Override public void own(String target) { }
@Override public void release(String target) { }
@Override public void publish(String target, String msgId, String content) { }
@Override public List<InboxMessage> peek(String target) { return List.of(); }
@Override public boolean ack(String target, String msgId) { return false; }
};
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> {
throw new UnsupportedOperationException(
"leadMailboxOpener must not be called — no coordinator: block is configured");
};
}
@Override
public LongSupplier nanoClock() {
return System::nanoTime;
}
@Override
public LongSupplier wallClockNanos() {
return System::nanoTime;
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor();
}
@Override
public void addShutdownHook(Runnable hook) {
shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
// Do not bind a real port in this assembly test.
}
@Override
public Runnable herdrPollWait() {
return () -> {
throw new UnsupportedOperationException("FakeHerdr is healthy; no poll wait is expected");
};
}
}
private static FleetConfig writeConfig(Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: "%s"
memberHerdrSocket: "%s"
idleSleepGuard:
enabled: false
fleet:
collaborators:
reviewer-alex:
tab: "collab: alex"
profiles:
sonnet:
subscription: true
argv: ["ccs", "sonnet"]
""".formatted(LEAD_SOCKET, MEMBER_SOCKET));
return FleetConfig.load(file);
}
@Test
void assembledRouterRoutesACollaboratorTerminalToTheLeadDaemon(@TempDir Path dir) throws Exception {
FleetConfig cfg = writeConfig(dir);
RecordingResourcePorts ports = new RecordingResourcePorts();
// The fixed FakeHerdr fixture already ties terminal "term_a" to a live agent on tab
// "w2:t7" (pane "w2:p7") — seeding only the tab LABEL to match the configured collaborator
// is enough to make LeadTabScanner resolve "term_a" as that collaborator. Seeded on the
// LEAD fake only: a collaborator's pane lives in the lead daemon, exactly like a lead's.
FakeHerdr lead = new FakeHerdr().withTab("w2", "w2:t7", "collab: alex");
FakeHerdr member = new FakeHerdr();
ports.herdrsBySocket.put(LEAD_SOCKET, lead);
ports.herdrsBySocket.put(MEMBER_SOCKET, member);
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg,
new ConfigRef(dir.resolve("fleetd.yaml"), cfg), new SubscriptionGuard(cfg.guard().hostSet())), ports);
try {
assertSame(runtime.router().leadAgents(), runtime.router().agentsFor("term_a"),
"a configured collaborator's terminal must route to the LEAD daemon — "
+ "FleetdAssembly must wire the collaborator map into the router's "
+ "predicate, not just LeadTabScanner.get()");
assertSame(runtime.router().memberAgents(), runtime.router().agentsFor("term_shell"),
"control: a terminal naming neither a lead nor a collaborator (term_shell, on "
+ "the unlabelled tab w2:t8) must still route to the member daemon");
} finally {
assertNotNull(ports.shutdownHook, "control: assembly must capture its shutdown hook");
ports.shutdownHook.run();
}
}
}
@@ -31,8 +31,7 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #670 — pins the {@code excludedWorkspaceLabels} argument {@link FleetdAssembly}'s
* production boot path passes to {@link LeadTabScanner} at {@code FleetdAssembly.java:265}
* ({@code Set.of()}).
* production boot path passes to {@link LeadTabScanner} ({@code Set.of()}).
*
* <p>{@code LeadTabScannerTest} already covers this constructor parameter, but it builds its own
* {@link LeadTabScanner} with its own set, so it tests the seam and proves nothing about the
@@ -191,7 +190,7 @@ class FleetdAssemblyLeadTabScannerExclusionTest {
Set<?> excluded = (Set<?>) excludedField.get(leads);
assertTrue(excluded.isEmpty(),
"FleetdAssembly.java:265 must pass an empty excludedWorkspaceLabels to "
"FleetdAssembly must pass an empty excludedWorkspaceLabels to "
+ "LeadTabScanner — scanning member tabs would demote the lead to a worker");
} finally {
assertNotNull(ports.shutdownHook, "control: assembly must capture its shutdown hook");
@@ -107,6 +107,8 @@ class ConfigRefTest {
assertTrue(out.applied());
assertTrue(out.deferred().isEmpty());
assertTrue(out.split().stream().noneMatch(s -> s.contains("fleet.collaborators")),
out.split().toString());
assertEquals("new charter", ref.get().fleet().charterFor(
dev.ltms.fleet.peer.MemberRole.ARCHITECT));
}
@@ -786,14 +788,100 @@ class ConfigRefTest {
assertTrue(out.split().getFirst().startsWith("fleet:"), out.split().toString());
assertTrue(out.split().getFirst().contains("restart"), out.split().toString());
assertTrue(out.split().getFirst().contains("live"), out.split().toString());
assertTrue(out.split().stream().noneMatch(s -> s.contains("fleet.collaborators")),
out.split().toString());
assertTrue(out.summary().contains("partially live"), out.summary());
// The snapshot still carries the new value — the LeadTabScanner's identity map and
// LeadLauncher's auto-launch are what wait for a restart; a reload rebuilds neither.
assertEquals("lead: opus-b", ref.get().fleet().leaders().get("opus").tab());
}
@Test
void addingAFleetCollaboratorIsReportedAsSplit(@TempDir Path dir) throws Exception {
assertCollaboratorChangeIsReported(dir, """
fleet:
collaborators:
alex:
tab: "collaborator: alex"
""", "add a collaborator");
}
@Test
void removingAFleetCollaboratorIsReportedAsSplit(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml("""
fleet:
collaborators:
alex:
tab: "collaborator: alex"
"""));
ConfigRef ref = refFor(f);
Files.writeString(f, yaml("fleet: {}\n"));
assertCollaboratorSplit(ref.reload(), "remove a collaborator");
}
@Test
void changingAFleetCollaboratorTabIsReportedAsSplit(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml("""
fleet:
collaborators:
alex:
tab: "collaborator: alex-a"
"""));
ConfigRef ref = refFor(f);
Files.writeString(f, yaml("""
fleet:
collaborators:
alex:
tab: "collaborator: alex-b"
"""));
assertCollaboratorSplit(ref.reload(), "change a collaborator tab");
}
@Test
void unchangedFleetCollaboratorsProduceNoSplitReport(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
String config = yaml("""
fleet:
collaborators:
alex:
tab: "collaborator: alex"
""");
Files.writeString(f, config);
ConfigRef ref = refFor(f);
Files.writeString(f, config);
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertTrue(out.split().isEmpty(), out.split().toString());
}
private static void assertCollaboratorChangeIsReported(Path dir, String changed, String action)
throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml("fleet: {}\n"));
ConfigRef ref = refFor(f);
Files.writeString(f, yaml(changed));
assertCollaboratorSplit(ref.reload(), action);
}
private static void assertCollaboratorSplit(ConfigRef.Outcome out, String action) {
assertTrue(out.applied(), action);
assertTrue(out.deferred().isEmpty(), out.deferred().toString());
assertEquals(1, out.split().size(), out.split().toString());
String report = out.split().getFirst();
assertTrue(report.startsWith("fleet: fleet.collaborators"), report);
assertTrue(report.contains("read once at startup"), report);
assertTrue(report.contains("restart"), report);
}
/**
* fleetd #333: {@code fleet.leaders} is the ONLY frozen part of {@code fleet:}. A reload that
* fleetd #333: {@code fleet.leaders} is a frozen part of {@code fleet:}. A reload that
* changes {@code tabLabel} (or charters, or a role pool) without touching {@code fleet.leaders}
* must stay fully hot with nothing reported — proving {@link ConfigRef#changedSplitKeys}
* compares {@code fleet.leaders} specifically rather than the whole {@code Fleet} record, which
@@ -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");
}
}
@@ -205,9 +205,12 @@ class LeadTabScannerTest {
}
/**
* The guard that matters: fleetd labels its own worker tabs, so if a worker space were scanned
* a naming accident would promote the fleet. The exclusion is by workspace, not by hoping the
* worker template never collides.
* Covers the {@code excludedWorkspaceLabels} parameter: a tab in an excluded workspace is never
* matched, whatever its label. Production always constructs this class with an empty set (CB-558,
* {@code FleetdAssembly}), so this parameter plays no part in the live guard against a worker
* tab being mistaken for a lead — that guard is {@code CallerResolver} asking the live
* spawned-member roster before any tab map. This test exists because the parameter still exists
* and is worth covering on its own terms.
*/
@Test
void aTabInAWorkerSpaceIsNeverALeadEvenWhenItsLabelMatches() {
@@ -219,6 +222,29 @@ class LeadTabScannerTest {
assertFalse(scanner(herdr, tabToName, new AtomicLong()).get().containsKey("term_impostor"));
}
/**
* The shipped shape (CB-558, {@code FleetdAssembly}): production always constructs this class
* with an empty {@code excludedWorkspaceLabels}, and a lead's {@code workspace:} default is the
* same shared {@code "fleet"} space the members use. A lead tab is still found when it sits in
* the exact same workspace as a member-labelled tab — the scanner tells them apart by the exact
* tab label, not by which workspace either one is in.
*/
@Test
void aLeadIsDiscoveredWhenItsWorkspaceIsTheSameAsTheMemberWorkspace() {
TopologyHerdr herdr = new TopologyHerdr()
.workspace("w1", "fleet")
.tab("w1:t1", "w1", "lead: opus-5.0")
.tab("w1:t2", "w1", "worker: gx10 #1")
.pane("w1:p1", "w1:t1", "term_opus")
.pane("w1:p2", "w1:t2", "term_worker");
LeadTabScanner s = new LeadTabScanner(herdr, Map.of("lead: opus-5.0", "opus-5.0"),
Set.of(), TTL, new AtomicLong()::get);
assertEquals("opus-5.0", s.get().get("term_opus"),
"a lead sharing the members' workspace is still discovered — the label, not the "
+ "workspace, is what matches it");
}
@Test
void aLabelWithNoConfiguredEntryIsIgnored() {
TopologyHerdr herdr = new TopologyHerdr().workspace("w1", "main")
@@ -865,11 +865,82 @@ class MessageServiceTest {
+ "follow-up send would come back BUSY instead of timing out on its own work");
}
/**
* Drive an async send on {@code T} to a resolved reply, polling as {@code owner} until the
* ticket reports {@link MessageService.Phase#DONE} (or the 2s deadline runs out). {@code owner}
* must be a terminal this ticket's creator check actually accepts, or this loops to the
* deadline and returns a non-{@code DONE} view.
*/
private MessageService.TaskView driveAsyncTicketToDone(String ticket, String owner) throws InterruptedException {
MessageService.TaskView view = null;
long deadline = System.currentTimeMillis() + 2000;
while (view == null || view.phase() != MessageService.Phase.DONE) {
if (System.currentTimeMillis() >= deadline) break;
view = messages.poll(ticket, owner);
//noinspection BusyWait
Thread.sleep(5);
}
return view;
}
@Test
void pollReturnsNullForAnUnknownTicket() {
assertNull(messages.poll("task-999999"), "a ticket that was never minted is unknown");
}
// --- fleetd #705: a ticket's creator terminal gates who may poll it -------------------------
@Test
void pollByAnotherTerminalIsRefused() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, "term_a");
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker works
assertTrue(rendezvous.resolve(T, "secret async result"), "a reply resolves the async send");
MessageService.TaskView owner = driveAsyncTicketToDone(ticket, "term_a");
assertNotNull(owner, "the creator must still be able to read its own ticket");
assertEquals(MessageService.Phase.DONE, owner.phase());
MessageService.TaskView refused = messages.poll(ticket, "term_b");
assertNotNull(refused, "a different terminal gets a refusal, not silence");
assertNotEquals(MessageService.Phase.DONE, refused.phase(),
"a different terminal must never see the ticket as DONE");
assertNull(refused.reply(), "a refusal must never carry the reply text");
assertFalse(String.valueOf(refused).contains("secret async result"),
"the reply text must not appear anywhere in the refused view");
}
@Test
void unnamedPrimaryStillReadsAnyTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, "term_lead");
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker works
assertTrue(rendezvous.resolve(T, "primary-visible result"), "a reply resolves the async send");
// callerTerminal == null is the unnamed primary (resolved by token or loopback trust, with
// no herdr pane) — it must read a ticket a terminal-bearing lead created.
MessageService.TaskView view = driveAsyncTicketToDone(ticket, null);
assertNotNull(view, "the unnamed primary must be able to read any ticket");
assertEquals(MessageService.Phase.DONE, view.phase());
assertEquals("primary-visible result", view.reply());
}
@Test
void creatorReadsItsOwnTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, "term_creator");
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker works
assertTrue(rendezvous.resolve(T, "own result"), "a reply resolves the async send");
MessageService.TaskView view = driveAsyncTicketToDone(ticket, "term_creator");
assertNotNull(view, "the ticket's own creator must be able to read it");
assertEquals(MessageService.Phase.DONE, view.phase());
assertEquals("own result", view.reply());
}
@Test
void pollReportsACompletedTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task");