Compare commits

..

18 Commits

Author SHA1 Message Date
Dai Ha a37acd5ee3 fleetd #369 review round 2: pin the factory's behaviour, not its call count
CI / contract (pull_request) Successful in 41s
CI / build (pull_request) Successful in 1m26s
everyGitSubprocessGoesThroughTheHermeticFactory counts ProcessBuilder("git",
...) call sites, so it catches a new helper built the old way, but a
reviewer proved it does not catch gitProcessBuilder itself being gutted:
removing pb.environment().putAll(hermeticEnv()) from inside the factory
leaves every call site unchanged, the count stays 2, and the whole
unpoisoned suite stays green.

Add gitProcessBuilderCarriesTheFullHermeticEnvironment, which inspects what
the factory actually hands to ProcessBuilder#start(): every hermetic key
present with the isolating value, and XDG_CONFIG_HOME pointed inside the
class's own throwaway directory rather than left unset or pointing at the
operator's real one. This fails the moment the hermetic environment stops
being applied, on any machine, with no poison needed. Keep the call-site
count check too — the two catch different regressions.
2026-09-06 20:09:32 +07:00
Dai Ha dd2efd8541 fleetd #369: make GitWorktreesTest hermetic against the machine's real git config
CI / contract (pull_request) Successful in 47s
CI / build (pull_request) Failing after 1m28s
gitOutput set GIT_CONFIG_GLOBAL/SYSTEM/TERMINAL_PROMPT but not XDG_CONFIG_HOME,
and status/fullStatus (plus every other raw git subprocess in this class) set no
isolation at all — inheriting the JVM's real environment, including the
operator's real ~/.gitconfig and default excludes file. Measured: with a
poisoned XDG_CONFIG_HOME, 56 of 59 tests failed.

Centralize every git subprocess this test starts through one factory,
gitProcessBuilder, which always applies the existing hermeticGitEnv isolation
(extended with XDG_CONFIG_HOME, the same fix #366 already applied to the
production-instance seam). Add a self-check test that counts direct
ProcessBuilder("git", ...) constructions in this file's own source and fails
if a future helper bypasses the factory, so the omission that caused this
ticket is caught by name instead of rediscovered on a poisoned machine.
2026-09-06 19:54:35 +07:00
Dai Ha 6c61355f8f Merge #367: dead lead tabs stop accumulating, without destroying a live one
CI / contract (push) Successful in 1m16s
CI / build (push) Successful in 2m7s
fleetd #359. LeadTabScanner used to join labelled tabs straight to terminals
with no liveness check, and its javadoc excused that ("a stale name costs
nothing here"). It cost plenty: LeadCoordLoop reads that map to pick which
pane a peer message goes into, so a dead tab was a candidate it could pick.
LeadLauncher, meanwhile, had no cleanup path at all -- every reconcile that
found 0 live created another tab and left the old one.

Both now cross-check agent.list, and neither trusts a single reading of it.
That matters because the daemon's own evidence on fleet01 was agent.list
reporting 0 live while ps showed one real claude. A first cut of this fix
closed tabs on that single reading, which would have closed the operator's
live lead instead of leaving a spare tab. So: a dead tab is flagged, not
closed, and only closed when a later reconcile still finds it dead; and the
scanner grants one grace scan to a terminal it already knew was live.

Verified on this merge, not taken from the worker's report:
  mvn clean install -> Tests run: 1412, Failures: 0, Errors: 0, BUILD SUCCESS

Mutation run on merge (PendingCloseMarker.strip made identity, so a flagged
tab stops matching its configured label): 4 failures, BUILD FAILURE.

Known and accepted: ensureLeads() runs at startup, so the second reading
arrives at the next restart. A tab that dies mid-session stays flagged and
open until then. Deliberate -- the bug is about repeated restarts, and one
leftover tab is cheaper than closing a live session on unverified evidence.
2026-09-05 16:00:04 +07:00
Dai Ha 96c406b968 #359 review: require two independent readings before destroying a lead's tab
CI / contract (pull_request) Successful in 45s
CI / build (pull_request) Successful in 1m33s
Finding 1 (LeadLauncher): closing a labelled tab on a single agent.list
miss could destroy a live lead's session — the ticket's own evidence
showed that exact signal missing a genuinely running agent. A dead
reading now only flags the tab (PendingCloseMarker); it is closed only
if a later, independently-connected reconcile still finds it dead
while flagged. A tab found live again has its flag cleared instead.

Finding 2 (LeadTabScanner): the new agent.list cross-check in scan()
was not covered by get()'s "keep the cache on a failed scan" contract,
which only fires on a thrown HerdrException. A successful-but-short
agent.list could silently drop a lead CallerResolver had already
resolved, demoting it to Role.WORKER. A terminal already reported live
now gets one grace scan before being dropped; a terminal never
reported live gets none, so the original #359 exclusion is unaffected.

Both mechanisms were mutation-tested: reverting either change turns
exactly its own new tests red and nothing else.
2026-09-05 15:54:52 +07:00
Dai Ha a5d81c3f70 Merge remote-tracking branch 'origin/main' into worker/359-dead-lead-tabs-f1253b-4 2026-09-05 15:28:49 +07:00
Dai Ha 92c0f164f1 Merge #366: seed the bridge's worker skills into every provisioned worktree
CI / contract (push) Successful in 46s
CI / build (push) Successful in 1m46s
fleetd #362 item 3. A member spawned against a repo that does not ship its
own .claude/skills/ could not load implementer, reviewer or hunter at all.
Every brief starts with "Load the <name> skill", and outside this repo that
line was silently a no-op. memberSkills: <dir> now copies those folders into
each provisioned worktree, skipping any name the target repo already ships.

Two review rounds, both about the same hazard: core.excludesFile is
single-valued, so pointing it at fleetd's own file would SHADOW the
operator's. It now composes instead of replacing, and the XDG default
excludes file is carried forward too.

Verified on this merge, not taken from the worker's report:
  mvn clean install -> Tests run: 1402, Failures: 0, Errors: 0, BUILD SUCCESS

Two mutations run on merge:
  drop the XDG fallback          -> 1 failure, BUILD FAILURE
  remove the composition itself  -> 2 failures, BUILD FAILURE
Both directions are pinned.
2026-09-05 13:32:16 +07:00
Dai Ha 9e813ec179 fleetd #362 review fix 2: route the XDG excludesFile fallback through gitEnv too
CI / contract (pull_request) Successful in 1m16s
CI / build (pull_request) Successful in 1m24s
Finding 1 (lead): the XDG fallback branch of previouslyEffectiveExcludesFileContent was
unpinned — deleting it left the suite green (Tests run: 1386, Failures: 0). Added
seedSkillsComposesWithTheXdgDefaultExcludesFileWhenNoneIsConfigured to pin it: isolates
XDG_CONFIG_HOME via the gitEnv seam at a temp dir carrying a synthetic git/ignore, points
GIT_CONFIG_GLOBAL at an empty file so core.excludesFile is genuinely unset (forcing the
fallback branch), seeds a skill, and asserts a file matching the XDG-default pattern still
reads as clean. Reverting the fix (mutating the fallback to resolve to "") turns this test
red with a real pasted failure (see PR body): "expected: <> but was: <?? xdg-fallback-marker>".

Finding 2 (lead, the one that actually needed a code fix): the fallback read XDG_CONFIG_HOME
and HOME straight from the JVM's own environment, not through the gitEnv seam every git
subprocess in this class already honours — so no test could isolate it, and on a machine
carrying a real ~/.config/git/ignore (this dev machine does), every seeding test silently
composed with that real file. Added resolveEnv/resolveHome, which check gitEnv first and
fall back to the JVM's real environment only when the seam doesn't supply a value (production
behaviour, where gitEnv is always Map.of(), is unchanged). Added a hermeticGitEnv() test
helper and routed every seeding test in GitWorktreesTest through it, so no test in the class
can reach the real machine's home directory for this fallback.

Also documents two non-defects the lead asked for one javadoc line each on: the composed
excludesFile is a snapshot taken at seed time, not a live reference to the operator's file;
and excludeSeededSkillsFromGitStatus assumes a fresh worktree (not idempotent, but the
double-seed path does not exist today, so no guard was added for it).
2026-09-05 13:23:39 +07:00
Dai Ha 395b3b5c46 #359: LeadTabScanner drops dead lead tabs; LeadLauncher closes them on relaunch
CI / contract (pull_request) Successful in 58s
CI / build (pull_request) Successful in 1m49s
LeadTabScanner.scan() joined labelled tabs to panes with no liveness check at
all, so a tab left behind by a crashed/relaunched lead read as a live lead
forever -- exactly the hazard its own javadoc predicted but excused. It now
cross-checks agent.list, the same signal LeadLauncher already trusted, and
drops any labelled tab with no agent running in it. That alone removes the
duplicate candidates LeadCoordLoop.resolveLocalLead() could pick from,
including the dead one its own WARN's advice (name a lead after
coordinator.selfId) could land a message in.

LeadLauncher never actually stopped the accumulation: relaunching on "0 live"
always created a brand new tab and left the old dead-labelled one right where
it was, so any restart that found 0 live for any reason (a real crash, or a
herdr read that missed a still-running agent) added one more dead tab,
forever. ensureLeads() now closes every dead-labelled tab for a lead as part
of the same reconcile that decides to relaunch, so at most one tab survives
per configured lead once a restart's reconcile has run.
2026-09-05 13:20:37 +07:00
Dai Ha 6e9e464d62 Merge #364: fleet_list reports lead-coordination state instead of guessing at it
CI / contract (push) Successful in 53s
CI / build (push) Successful in 1m56s
fleetd #361. fleet_list's coordination block now carries this daemon's own
mailbox state and one row per operator-declared peer, each as a tri-state
status (exists / absent / unknown) rather than a boolean. pending and
consumers appear only when status is "exists", so an unresolved probe can
never render as a measured zero.

Verified on this merge, not taken from the worker's report:
  mvn clean install -> Tests run: 1394, Failures: 0, Errors: 0, BUILD SUCCESS

Mutation run on merge (isMissingQueue always returns true, which restores
the exact defect the ticket fixes): 4 failures, BUILD FAILURE. The
discriminator's false branch is pinned.
2026-09-05 13:19:46 +07:00
Dai Ha 105c065615 Merge origin/main into worker/362-worktree-skills-c03e51-3 (picks up #363) 2026-09-05 13:18:32 +07:00
Dai Ha 01492059d4 fleetd #361 review round 2: pin isMissingQueue's false branch
CI / contract (pull_request) Successful in 1m5s
CI / build (pull_request) Successful in 1m34s
The reviewer's mutation (isMissingQueue always returns true) restored
the exact overstatement fleetd #361 exists to fix -- every declare
failure reading as a confirmed absence -- and still left mvn clean
install green (1389/1389), because no test drove a non-404 shape
through inspect(). The false branch was the whole discriminator
between MailboxState.absent() and MailboxState.unknown(), unpinned.

Widened LeadMailbox.isMissingQueue from private to package-private and
added LeadMailboxIsMissingQueueTest: five hermetic tests (no broker)
covering the true case and all three false shapes isMissingQueue's own
javadoc lists -- a different reply code, a ShutdownSignalException
whose reason isn't a Channel.Close, and an IOException with no such
cause at all (plus an IOException wrapping an unrelated exception
type). Re-ran the reviewer's exact mutation locally: 4 of 5 new tests
went red with the expected assertion messages; reverted, and mvn clean
install is green again at 1394/1394 (1389 + 5 new).

Also added a one-line javadoc note on LeadMailbox.inspect being honest
about which of its two RuntimeException catches is proven by a test
(the createChannel() one, end-to-end against a real broker) and which
stays purely defensive (the declare-site one, for a connection-drops-
mid-call race no test drives on purpose).
2026-09-05 13:12:24 +07:00
Dai Ha f84824ee29 fleetd #362 review fix: compose skill-seeding excludes with the operator's own excludesFile
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 1m54s
core.excludesFile is single-valued, so pointing it at fleetd's own seeded-skill exclude file
with --replace-all at worktree scope was SHADOWING whatever excludesFile the worktree already
resolved (an operator's global config, most commonly) instead of adding to it. This repo's own
.gitignore does not ignore target/ — only an operator's global excludesFile does — so every
worker's `mvn clean install` would make target/ show up as untracked, and CB-576's deliberately
untracked-inclusive hasUncommitted would then read every such worktree as dirty forever, so it
is never cleaned up.

excludeSeededSkillsFromGitStatus now reads whatever core.excludesFile resolves to BEFORE writing
anything (falling back to git's own $XDG_CONFIG_HOME/git/ignore default when the key is unset
entirely, per gitignore(5)), and writes that content into fleetd's own exclude file ahead of the
seeded skill patterns, so every operator-configured pattern keeps applying inside the seeded
worktree. Proven with a new test, seedSkillsComposesWithAnAlreadyEffectiveGlobalExcludesFile,
which isolates a synthetic "operator's global config" via a new gitEnv test seam on GitWorktrees
(GIT_CONFIG_GLOBAL pointed at a throwaway temp file, never the real machine's config) and drives
the real add() path end to end.

Also documents (FleetConfig javadoc + fleetd.example.yaml) that memberSkills copies every
non-hidden subdirectory of its source wholesale, with no per-file allowlist.
2026-09-05 13:10:27 +07:00
Dai Ha c4d40fbc2b fleetd #361 review: fix false-negative absent, uncancelled probes, and a throw contract gap
CI / contract (pull_request) Successful in 47s
CI / build (pull_request) Successful in 1m45s
Three findings from review of #364, fixed on the same branch:

1. LeadChannel.MailboxState.absent() was returned both for a genuinely
   absent mailbox AND for "the probe could not determine anything"
   (timeout, unreachable broker, other declare failure) -- exactly the
   overstatement #361 exists to fix, one level down. MailboxState now
   carries a Presence enum (EXISTS/ABSENT/UNKNOWN) with exists()/known()
   accessors; LeadMailbox.inspect classifies a real AMQP 404 (measured
   against a live broker, not assumed: an IOException wrapping a
   ShutdownSignalException whose Channel.Close reply code is 404) as
   ABSENT and everything else as UNKNOWN. fleet_list's mailbox/peer rows
   now render a "status" of exists/absent/unknown and only include
   pending/consumers when status is "exists", so an unresolved self- or
   peer-probe can never render as a measured zero.

2. FleetMcp.probe's get(timeoutMs) left a timed-out inspect() task
   running forever on its own virtual thread, holding the AMQP channel
   it had already opened -- against a hung (not down) broker this would
   orphan one channel per fleet_list call until the connection's
   channel-max was exhausted, breaking publish() too. probe() now holds
   the Future and calls cancel(true) on timeout/failure so the orphaned
   task is interrupted instead of abandoned, and now returns
   MailboxState.unknown() (never absent()) on timeout/exception.

3. LeadMailbox.inspect only caught IOException, but createChannel() on
   an already-closed connection throws AlreadyClosedException, an
   unchecked RuntimeException (measured against a live broker) -- so it
   could escape the "never throws" contract. Both places in inspect now
   also catch RuntimeException and report unknown().

Tests: MailboxState.exists()/absent()/unknown() call sites updated
across FleetMcpTest/FleetMcpLeadCoordTest; new hermetic tests cover the
tri-state fleet_list rendering (self-probe unknown, a peer that's
absent vs. one that's unknown) and probe cancellation (a LeadChannel
fake that blocks until interrupted, proving probe() doesn't just give
up on it); new @Tag("contract") LeadMailboxTest cases pin the real
exception shapes for both the 404 and the already-closed-connection
paths and prove inspect() reports unknown (never throws) when the
connection is already closed.
2026-09-05 13:01:48 +07:00
Dai Ha 7c684e40d3 fleetd #362 (item 3): seed .claude/skills/ into provisioned worktrees
CI / contract (pull_request) Successful in 1m20s
CI / build (pull_request) Successful in 1m46s
Add memberSkills: <dir> to FleetConfig. GitWorktrees#add copies each
skill folder from that directory into <worktree>/.claude/skills/ so a
member spawned against ANY repo — not only one that already ships its
own skills — can load a bridge skill (e.g. implementer). A skill the
target repo already carries is never overwritten.

Every seeded path is hidden from `git status` in that worktree ONLY,
via a --worktree-scoped core.excludesFile pointing at a file under the
worktree's own private git dir (outside the working tree, so it can
never be committed) — not the shared .git/info/exclude, which a linked
worktree resolves to the repo's common git dir and would otherwise leak
visibility changes into the primary checkout and every sibling
worktree. Proven with a real `git status --porcelain` in
GitWorktreesTest, not by reasoning.

Seeding is best-effort like the existing overlayParity/isolateToolSurface
steps: a missing/unreadable source or a copy/exclude failure is logged
and skipped, never fails the spawn. memberSkills is triaged as a
DEFERRED config key in ConfigRef (baked once into GitWorktrees at
startup, like worktreeGroup), with its own changedDeferredKeys branch
and coverage-test entries.
2026-09-05 12:54:01 +07:00
ltms 6938f52155 Merge #363: make the plugin visible, and fix the drift that made it unusable
CI / contract (push) Successful in 1m12s
CI / build (push) Successful in 1m28s
2026-09-05 07:50:19 +02:00
Dai Ha 0df34f3220 fleetd #361: close the lead-coordination visibility gap
CI / contract (pull_request) Successful in 1m10s
CI / build (pull_request) Successful in 1m34s
Lead-to-lead AMQP coordination had a send half with tools and a receive
half without. This closes three blind spots:

- LeadChannel gains inspect(coordId) -> MailboxState(exists, pending,
  consumers), implemented in LeadMailbox with a throwaway probe channel
  (never the long-lived publish/consume channels) so a passive-declare
  404 on a missing queue can never take down publish() on the same
  instance.
- FleetConfig.Coordinator gains peers: List<String> (defaults to empty,
  blank entries dropped) so a daemon can declare which peer coord-ids
  it expects to reach.
- fleet_list reports coordination state via a new CoordinationSource
  (own coord-id, own mailbox state, held messages as msgId/from/preview
  only, and one row per configured peer with reachability/pending/
  consumers), following the existing OutageSource/QuarantineSource
  "Source record with none()" idiom instead of growing listFleet's
  overload chain by another positional parameter. Every peer probe is
  bounded by a 1.5s timeout on a virtual-thread pool and degrades to
  absent rather than ever slowing or failing fleet_list.
- fleet_send{coordId}'s success text now says "durably confirmed by the
  broker" instead of "delivered", and warns (while still reporting
  success) when the target mailbox has zero consumers attached.

Tests: hermetic unit tests for exists/absent/zero-consumer/old-config-
no-peers-key/new fleet_list shape using FakeLeadChannel, plus a
@Tag("contract") LeadMailboxTest.inspectingAMissingMailboxNeverBreaks
PublishOnTheSameInstance proving the invariant against a real broker.
2026-09-05 12:44:29 +07:00
Dai Ha 457458437f #362: make the plugin visible, and fix the drift that made it unusable
CI / contract (pull_request) Successful in 1m12s
CI / build (pull_request) Successful in 1m31s
CB-527 shipped a Claude Code plugin and a marketplace in this repo. Nothing in
CLAUDE.md or docs/ ever named it, so a later session planned the same feature
from scratch. The wiki Features entry existed and was correct, but wiki/ is a
submodule whose pointer is never advanced, so no session reads it.

Visibility:
- CLAUDE.md addendum now names plugin/ and both structural limits, so every
  session sees it. This is the change that stops the rebuild happening again.
- wiki/11-Features.md records the rename and why the entry alone was not enough.

Drift (each measured against the code, not assumed):
- mount name fleetd -> fleet, matching PeerLauncher.MCP_MOUNT_NAME. The old name
  gave a lead with both a project .mcp.json and the plugin two mounts of one
  daemon and a duplicated fleet_* tool set.
- url is now ${FLEETD_MCP_URL} instead of a hardcoded address, so one plugin can
  serve hosts running the daemon on different ports. Plain ${VAR}, the form
  kb-alms proves works here; ${VAR:-default} is untested and not used.
- plugin claude-bridge -> fleet, marketplace claude-bridge -> fleetd, version
  0.2.0. Breaking for a 0.1.0 install: mcp__fleetd__* becomes mcp__fleet__*.
- README install path ltms/claude-bridge -> the fleet/fleetd remote.
- the setup skill's §5 told operators to pin primary.terminal:. CB-579 replaced
  that with fleet.leaders.*.tab. Replaced, with the duplicate-tab warning (#359).

Scope: the plugin is lead-side only, and cannot be otherwise. The launcher adds
--agent only when <worktree>/.claude/agents/<role>.md exists in the member's own
tree (ClaudeCodeLauncher.java:371,391), and a member's CLAUDE_CONFIG_DIR points
at its profile's config dir (ClaudeCodeLauncher.java:285), so a member never
reads the operator's plugin store. On this Mac all four Claude profiles set
configDir, and the four ccs instances hold four separate copies of the plugin
store -- same md5, different inodes. Seeding member skills through the worktree
is #362 scope item 3, implemented separately.

Note for anyone verifying a plugin: `claude plugin validate` does NOT read
.mcp.json. Replacing it with `{ this is not json at all` still passes, exit 0.

Refs #362, #359
2026-09-05 12:42:20 +07:00
Dai Ha 3759c41f99 Merge #354: the redeploy health gate classifies AMQP errors instead of counting them
CI / contract (push) Successful in 1m14s
CI / build (push) Successful in 1m36s
The gate counted ERROR lines since RESTART_MARK. On a laptop that idle-sleeps after one
minute on battery that meant 6 ERROR lines for an AMQP link that recovered every time,
and a gate that cries wolf is a gate nobody reads.

It now reports three states: no errors; only errors proven to have recovered (quiet, and
the gate passes); anything else (the old warning, unchanged). Attribution is per
connection, using the names #356 put into the log -- a lead-mailbox recovery can no
longer clear an unrecovered reply-inbox reset. A candidate carrying neither name is
unattributable and stays LOUD.

Two earlier rounds were rejected. Round 1 was inert: it matched nothing in the real log,
because the layout abbreviates the logger and 'Connection reset' sits in the stack trace,
not on the ERROR line -- my brief had pointed the worker at fleetd.out, which is untracked
and so absent from its worktree. Round 2 was correct and honest but could not attribute
anything, which is what motivated #356.

Verified on merge beyond the worker's own mutations:
 - ran the classifier against the REAL log, which is still in the pre-#356 format: 6 total,
   0 recovered, 6 unexplained. Old-format lines carry no connection name, so they stay loud
   -- the safe direction, on genuine data rather than a fixture.
 - adversarial fixture the worker did not write: a lead recovery BEFORE any failure banks
   no credit; 2 inbox resets with 1 recovery leaves 1 unexplained; a non-AMQP ERROR stays
   loud. total=3 recovered=1 unexplained=2, as intended.
 - RESTART_MARK still anchors the scanned region.

Caveat carried from the PR: the patterns are source-derived. The daemon has not been
redeployed, so they are not yet confirmed against a live log.
2026-09-05 06:09:29 +07:00
31 changed files with 2723 additions and 166 deletions
+4 -4
View File
@@ -1,15 +1,15 @@
{
"name": "claude-bridge",
"name": "fleetd",
"description": "Tooling for orchestrating a fleet of delegated coding agents through the fleetd MCP gateway.",
"owner": {
"name": "LTMS"
},
"plugins": [
{
"name": "claude-bridge",
"name": "fleet",
"source": "./plugin",
"description": "Make a project bridge-ready: mount the fleetd MCP gateway and apply standard Claude Code settings so a session can orchestrate delegated workers. Ships no credentials.",
"version": "0.1.0",
"description": "Mount the fleetd MCP gateway and apply standard Claude Code settings so a session can orchestrate delegated workers. Ships no credentials.",
"version": "0.2.0",
"author": {
"name": "LTMS"
}
+11
View File
@@ -210,6 +210,17 @@ must obey belongs in the charter, not here.
- **Primary-side skills** (not delegation playbooks — a worker cannot use them):
`port-to-opencode` (make an OpenCode session a participant in this workspace) and
`fleets-status` (report every fleet that shares one LavinMQ instance).
- **This repo is also a Claude Code marketplace, and ships a plugin.** `.claude-plugin/marketplace.json`
points at `plugin/`, which carries the MCP mount and the `setup` skill
(`/claude-bridge:setup` — make any project bridge-ready). It was added in CB-527 and then went
unmentioned by every instruction file, so it drifted and a later session planned it from scratch
(#362). **Read `plugin/` before designing anything about onboarding a project.** Two limits are
structural, not bugs: a plugin cannot carry the role agent files, because
`ClaudeCodeLauncher.java:371` requires `<cwd>/.claude/agents/<role>.md` in the member's own
worktree; and a plugin cannot deliver anything to members at all, because
`ClaudeCodeLauncher.java:285` exports `CLAUDE_CONFIG_DIR` and every Claude profile here sets it,
so a member never reads the operator's plugin store. **The plugin is the lead-side surface;
member-facing assets travel in the worktree.**
- **Never commit** `.mcp.json` (the primary's local copy, flagged `--skip-worktree`) or `wiki/`
(a submodule with its own remote).
- **A provisioned worktree neutralizes `.mcp.json`, `opencode.json` and `.autoenv`** — the repo's
+20
View File
@@ -774,6 +774,19 @@ guard:
# so fleetd falls back to the weaker CB-596 sentinel overlay instead (a WARN names the gap).
# worktreeGroup: fleet-workers
# fleetd #362: a directory of skill folders (each a subdirectory holding a SKILL.md, the same
# shape as this repo's own .claude/skills/) copied into every PROVISIONED worktree's
# .claude/skills/, so a member spawned against ANY repo — not only one that already ships its own
# copy — can load a bridge skill (e.g. implementer). Unset (the default): no worktree is touched
# beyond today's behaviour. A skill folder the target repo already carries under
# .claude/skills/<name> is never overwritten — the repo's own copy always wins. Claude Code
# members only; an opencode member reads a different path (.opencode/agent) this key does not
# touch. Best-effort like worktreeGroup above: a missing/unreadable directory here is logged and
# skipped, never a failed spawn. Every non-hidden subdirectory of this directory is copied
# wholesale, with no per-file allowlist — don't park scratch files or drafts alongside the real
# skill folders, they will be copied into every provisioned worktree too.
# memberSkills: /path/to/fleetd/checkout/.claude/skills
# Session lifecycle limits (CB-303). All knobs are opt-in; omit or set to null to keep
# the feature disabled. By default the daemon never reaps, caps, or drains sessions.
# idleTtlSeconds → reap READY/DONE sessions idle longer than this (never BUSY/SPAWNING)
@@ -821,10 +834,17 @@ guard:
# across every daemon sharing this vhost.
# prefetch → consumer basicQos, capping how many unacked messages the mailbox holds in-heap.
# Default 32 when omitted.
# peers → fleetd #361: the coord-ids of the OTHER daemons on this vhost, declared by the
# operator (the daemon never guesses). fleet_list reports each one's live reachability
# (a passive queue check, never a presence protocol) alongside this daemon's own
# mailbox state. Omit, or leave empty, for a daemon with no known peers yet — an
# undeclared peer can still reach you and be reached by fleet_send, it just will not
# show up as a row in fleet_list.
# coordinator:
# uriEnv: LEAD_COORD_URI
# selfId: mac-opus
# prefetch: 32
# peers: [fleet01-lead]
# Active push-to-primary (CB-307 Stage 3). When a worker reply lands with no open fleet_send,
# the ReplyPushLoop injects a *drain nudge* (never the payload) into the primary's own herdr
@@ -248,7 +248,8 @@ public final class Fleetd {
contextCap = cfg.lifecycle().contextCap();
}
boolean clearAfterTurn = cfg.lifecycle() != null && cfg.lifecycle().clearAfterTurn();
SessionManager sessions = new SessionManager(workers, new GitWorktrees(cfg.worktreeRoot(), cfg.worktreeGroup()),
SessionManager sessions = new SessionManager(workers,
new GitWorktrees(cfg.worktreeRoot(), cfg.worktreeGroup(), cfg.memberSkills()),
System::nanoTime, contextCap, clearAfterTurn);
liveCountRef.set(profileName -> liveSessionCount(sessions.roster(), profileName));
@@ -652,7 +653,12 @@ public final class Fleetd {
quarantineSource,
leadMailbox,
outageSource,
new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)));
new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)),
// fleetd #361: the operator-declared peers this daemon's fleet_list should try to
// reach. Read from the SAME snapshot leadMailbox itself opened from (cfg.coordinator()),
// not the live config.get() — coordinator wiring is already boot-time-fixed (see
// leadMailbox above), so peers follows the same rule rather than half hot-reloading.
cfg.coordinator() == null ? List.of() : cfg.coordinator().peers());
// CB-637: the receive half. Only constructed when a lead mailbox actually opened — with no
// coordinator (or an unreachable one) there is nothing to deliver, so no scheduler is
@@ -39,10 +39,11 @@ import java.util.function.Supplier;
* keeps the old value until a restart: {@code lifecycle:}, {@code leadHeartbeat:},
* {@code spawnReadyTimeoutMs} / {@code spawnReadyPollMs}, {@code quarantineCooldownSeconds}
* (CB-578 stage B — baked once into the {@code BackendQuarantine} built at startup),
* {@code guard:}, {@code worktreeRoot:} and {@code worktreeGroup:} (both baked once into the
* {@code GitWorktrees} built at {@code Fleetd.java:251} and never rebuilt — fleetd #323
* instance 2 found {@code worktreeGroup} missing from this list and from
* {@link #changedDeferredKeys}), {@code primary:} (fleetd #326 — {@code Fleetd.java:506, 519,
* {@code guard:}, {@code worktreeRoot:}, {@code worktreeGroup:} and {@code memberSkills:}
* (all three of the latter baked once into the {@code GitWorktrees} built at
* {@code Fleetd.java:251} and never rebuilt — fleetd #323 instance 2 found
* {@code worktreeGroup} missing from this list and from {@link #changedDeferredKeys};
* {@code memberSkills} (fleetd #362) followed the same shape), {@code primary:} (fleetd #326 — {@code Fleetd.java:506, 519,
* 520} read {@code cfg.primary()} only off the startup snapshot to build {@code
* PrimaryRegistry} and size {@code ReplyPushLoop}'s reminder cap/backoff, and neither is
* rebuilt on reload. Say the consequence exactly: {@code primary.terminal} is DEPRECATED
@@ -129,9 +130,10 @@ import java.util.function.Supplier;
* five of COLD_KEYS" rather than re-listing them, so prose and set cannot drift again.</li>
* </ul>
*
* <p><strong>The denominator, measured on 2026-09-04 (fleetd #330; recounted for fleetd #333).</strong>
* {@code FleetConfig} has 22 top-level record components: 5 cold, 11 deferred, 3 split, 3
* hot-excluded. Three of them are named nowhere in this file, and the reason is the same for all
* <p><strong>The denominator, measured on 2026-09-04 (fleetd #330; recounted for fleetd #333);
* recounted again for fleetd #362.</strong> {@code FleetConfig} has 23 top-level record components:
* 5 cold, 12 deferred, 3 split, 3 hot-excluded. Three of them are named nowhere in this file, and
* the reason is the same for all
* three: {@code placement}, {@code memberCredentials} and {@code memberLoginShell} are
* <strong>hot</strong> and correctly absent — all three are read live off {@code config.get()}
* (placement through the {@code CompositePeerLauncher} supplier the Hot bullet names;
@@ -211,7 +213,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
* read {@link #COLD_KEYS} and {@link #SPLIT_KEYS}.
*/
static final Set<String> DEFERRED_KEYS = Set.of(
"guard", "worktreeRoot", "worktreeGroup", "primary", "configReload",
"guard", "worktreeRoot", "worktreeGroup", "memberSkills", "primary", "configReload",
"leadHeartbeat", "lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs",
"quarantineCooldownSeconds", "profiles");
@@ -402,6 +404,13 @@ public final class ConfigRef implements Supplier<FleetConfig> {
if (!Objects.equals(old.worktreeGroup(), fresh.worktreeGroup())) {
changed.add("worktreeGroup");
}
// fleetd #362: baked into the same GitWorktrees as worktreeRoot/worktreeGroup
// (Fleetd.java:251) and never rebuilt either — a reload that changes only memberSkills
// must be reported the same way, or a newly provisioned worktree keeps seeding from (or
// skipping) the old source directory with nothing telling the operator why.
if (!Objects.equals(old.memberSkills(), fresh.memberSkills())) {
changed.add("memberSkills");
}
// fleetd #326: Fleetd.java:506, 519, 520 read cfg.primary() only off the startup snapshot
// (PrimaryRegistry's pinned terminal, ReplyPushLoop's reminder cap and backoff) — neither is
// rebuilt on reload, so a changed value needs a restart. Note what it does NOT mean:
@@ -106,6 +106,19 @@ import java.util.regex.PatternSyntaxException;
* When {@code memberHerdrSocket} is NOT configured this field is never
* consulted at all; fleetd keeps reading its own {@code $SHELL}, exactly as
* before this field existed.
* @param memberSkills fleetd #362: nullable directory of skill folders (each a subdirectory
* holding a {@code SKILL.md}, the same shape as this repo's own {@code
* .claude/skills/}) copied into every provisioned worktree's {@code
* .claude/skills/}, so a member spawned against ANY repo — not only one that
* already ships its own copy — can load a bridge skill such as {@code
* implementer}. {@code null}/blank ⇒ off: no worktree is touched beyond
* today's behaviour. A skill folder the target repo already carries is never
* overwritten — see {@link dev.ltms.fleet.session.GitWorktrees}. Claude Code
* members only; an opencode member's equivalent lives under a different path
* ({@code .opencode/agent}) and is not covered by this key. Every non-hidden
* subdirectory of this directory is copied wholesale, with no per-file
* allowlist — do not park scratch files or drafts alongside the real skill
* folders, they will be copied into every provisioned worktree too.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record FleetConfig(
@@ -130,7 +143,22 @@ public record FleetConfig(
MemberCredentials memberCredentials,
Coordinator coordinator,
String worktreeGroup,
String memberLoginShell) {
String memberLoginShell,
String memberSkills) {
/** Back-compat form before the {@code memberSkills} key was added. */
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
Guard guard, String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
ConfigReload configReload, Integer quarantineCooldownSeconds,
MemberCredentials memberCredentials, Coordinator coordinator, String worktreeGroup,
String memberLoginShell) {
this(bind, herdrSocket, memberHerdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, worktreeGroup,
memberLoginShell, null);
}
/** Back-compat form before the {@code memberLoginShell} key was added. */
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
@@ -888,12 +916,18 @@ public record FleetConfig(
* {@code null} ⇒ kept as {@code null} (no self id configured).
* @param prefetch the consumer's {@code basicQos} prefetch count. {@code null}/non-positive ⇒
* {@link LeadMailbox#DEFAULT_PREFETCH}.
* @param peers fleetd #361: the coord-ids the operator declares as this daemon's peers — the
* daemon never guesses who else exists. {@code fleet_list} reports each one's
* live reachability. Blank entries are dropped; {@code null} ⇒ an empty list, so
* a config written before this field existed still parses unchanged.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Coordinator(String uri, String uriEnv, String selfId, Integer prefetch) {
public record Coordinator(String uri, String uriEnv, String selfId, Integer prefetch, List<String> peers) {
public Coordinator {
selfId = (selfId == null || selfId.isBlank()) ? null : selfId;
peers = peers == null ? List.of()
: peers.stream().filter(p -> p != null && !p.isBlank()).toList();
}
/** True when a {@code uriEnv} is configured by name, whether or not its variable resolves. */
@@ -1495,7 +1529,7 @@ public record FleetConfig(
"bind", "herdrSocket", "memberHerdrSocket", "profiles", "guard", "worktreeRoot",
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds",
"memberCredentials", "coordinator", "worktreeGroup", "memberLoginShell");
"memberCredentials", "coordinator", "worktreeGroup", "memberLoginShell", "memberSkills");
/** Load and validate config from {@code path}. */
public static FleetConfig load(Path path) {
@@ -2172,9 +2206,12 @@ public record FleetConfig(
// memberLoginShell is left as-is (fleetd #213), like worktreeGroup: null/blank is "not
// configured", and there is no sane non-null default — a member's login shell is
// operator-specific and only meaningful when memberHerdrSocket is also set.
// memberSkills is left as-is (fleetd #362), like worktreeGroup/memberLoginShell: null/blank
// is "off", and there is no sane non-null default — the daemon may not even run from a
// checkout that ships its own .claude/skills/.
return new FleetConfig(b, herdrSocket, memberHerdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
quarantineCooldown, mc, coordinator, worktreeGroup, memberLoginShell);
quarantineCooldown, mc, coordinator, worktreeGroup, memberLoginShell, memberSkills);
}
/**
@@ -5,6 +5,7 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.Collections;
import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.Locale;
import java.util.Map;
@@ -53,17 +54,40 @@ import java.util.function.Supplier;
* daemon and found again by this scan. The trust direction above is unaffected — fleetd writing a
* name for a lead it just started is not a pane promoting itself — but <em>staleness</em> becomes
* real: a label left behind by a session that has since died would read as a live lead forever.
* This scanner does not solve that (its job is naming, and a stale name costs nothing here); the
* launcher does, by requiring a running agent in the tab before it counts the lead as live. If you
* ever make a decision that <em>removes</em> something based on this map, add the same check.
* The remaining hazard is an <em>operator</em> one — a worker {@code tabLabel} template that
* happens to start with the same prefix would promote the whole fleet — and that is refused at
* startup by {@code FleetConfig.validateLeadTabPrefixes} rather than documented here.
*
* <p><strong>fleetd #359 — the staleness check this class used to skip.</strong> This used to say
* "a stale name costs nothing here" and leave liveness to {@code LeadLauncher}, on the theory that
* naming and removing are different decisions. That was wrong: {@code dev.ltms.fleet.msg.LeadCoordLoop}
* makes exactly the kind of removal decision the old javadoc warned about, by reading this map to
* pick which pane a peer message goes into — and a stale entry there is not free. On a host where
* {@code fleetd} had restarted more than once, a labelled-but-dead tab from a previous life was
* reported right alongside the live one; {@code LeadCoordLoop.resolveLocalLead()} saw more than one
* candidate and refused to guess (safe), but the fix the daemon's own WARN suggests — name a lead
* after {@code coordinator.selfId} — stops being safe once two tabs can share a label: step 1 of
* that resolution picks whichever matching entry it finds first, which can be the dead one, and
* typing a peer's message into a dead shell does not fail — it is silently gone instead of merely
* held. {@link #scan()} now cross-checks every labelled tab against {@code agent.list} (the same
* signal {@code LeadLauncher.countLeads} already trusts for the same purpose) and drops any tab
* with no agent running in it, so a dead tab is never in the map for a caller to pick at all.
*
* <p><strong>Caching.</strong> {@link #get()} is on the request path (every resolve), so the scan
* is TTL-cached and a stale-but-valid map is preferred to a herdr round-trip. A failed scan keeps
* the previous answer instead of emptying it — a herdr hiccup must not silently demote a live lead
* mid-session.
* is TTL-cached and a stale-but-valid map is preferred to a herdr round-trip. A failed scan (herdr
* throws) keeps the previous answer instead of emptying it — a herdr hiccup must not silently
* demote a live lead mid-session.
*
* <p><strong>fleetd #359 review, finding 2 — a successful-but-wrong scan is the same hazard.</strong>
* The catch above only fires when a call throws. It does nothing for a call that returns 200 with an
* incomplete answer — exactly what the ticket's own evidence showed {@code agent.list} can do. Once
* this class started trusting that signal, an empty read would otherwise get cached as fact and
* silently drop a lead {@code CallerResolver} had, until then, correctly resolved — turning it into a
* {@code Role.WORKER}, which refuses every orchestration call. So a terminal this class already
* reported as live is not dropped the first time {@code agent.list} loses it: {@link #scan()} grants
* it one grace scan (see {@code gracedTerminals}) and only drops it if a <em>later</em> scan still
* finds no agent. A terminal never reported live before gets no grace — that would weaken the
* original #359 fix itself, which this class's own test suite already pins.
*/
public final class LeadTabScanner implements Supplier<Map<String, String>> {
@@ -79,6 +103,15 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
private long scannedAtNanos;
private boolean everScanned;
/**
* Terminals currently on their one grace scan: {@code cached} reported them live, the most
* recent {@link #scan()} found no agent for them, and they were re-included anyway. Cleared for
* a terminal the instant it is seen live again; a terminal still here on the <em>next</em> scan
* is finally dropped. Scoped separately from {@link #cached} so a graced terminal cannot renew
* its own grace forever just by staying in the exposed map (fleetd #359 review, finding 2).
*/
private Set<String> gracedTerminals = Set.of();
/**
* @param herdr the herdr client to query ({@code workspace.list},
* {@code tab.list}, {@code pane.list} — all read-only)
@@ -141,7 +174,7 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
return cached;
}
/** One full pass: labelled tabs → their panes → those panes' terminals. */
/** One full pass: labelled tabs → live agents in them → those panes' terminals. */
private Map<String, String> scan() {
Map<String, String> nameByTab = new LinkedHashMap<>();
for (JsonNode w : herdr.call("workspace.list").path("workspaces")) {
@@ -158,17 +191,47 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
}
}
Map<String, String> byTerminal = new LinkedHashMap<>();
if (!nameByTab.isEmpty()) {
// 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 name = nameByTab.get(p.path("tab_id").asText(null));
String terminal = p.path("terminal_id").asText(null);
if (name != null && terminal != null && !terminal.isBlank()) {
byTerminal.put(terminal, name);
}
if (nameByTab.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.
Set<String> tabsWithAgent = new HashSet<>();
for (JsonNode a : herdr.call("agent.list").path("agents")) {
String tabId = a.path("tab_id").asText(null);
if (tabId != null) {
tabsWithAgent.add(tabId);
}
}
Map<String, String> 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);
String terminal = p.path("terminal_id").asText(null);
if (name == null || terminal == null || terminal.isBlank()) {
continue;
}
if (tabsWithAgent.contains(tabId)) {
byTerminal.put(terminal, name);
continue;
}
// No agent reported for this tab, but its tab/pane are still here — this is the
// ambiguous case review finding 2 named: a successful agent.list that came back short
// does not prove the lead is dead. Grant one grace scan to a terminal we had already
// 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);
stillGraced.add(terminal);
}
}
gracedTerminals = stillGraced;
return Collections.unmodifiableMap(byTerminal);
}
@@ -177,12 +240,14 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
*
* <p>Exact match (case-insensitive, ends stripped) against {@link #tabToName} — 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.
* 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) {
if (label == null) {
return null;
}
return tabToName.get(label.strip().toLowerCase(Locale.ROOT));
return tabToName.get(PendingCloseMarker.strip(label).toLowerCase(Locale.ROOT));
}
}
@@ -0,0 +1,42 @@
package dev.ltms.fleet.herdr;
/**
* The suffix {@code dev.ltms.fleet.lead.LeadLauncher} appends to a lead tab's label the first time a
* reconcile finds no running agent in it, before it is sure enough to close the tab outright.
*
* <p><strong>fleetd #359 review, finding 1.</strong> The daemon's own evidence showed
* {@code agent.list} can read "no agent" for a tab that genuinely has one running — so a single such
* reading must never be treated as proof a tab is dead. {@code LeadLauncher} now writes this marker
* on the first miss, and only closes the tab if a <em>later</em>, independent reconcile still finds
* it dead while the marker is still there. Two consecutive misses, one restart apart, is a much
* stronger claim than one.
*
* <p>{@link LeadTabScanner} strips the same suffix before matching a label against a configured
* lead's {@code tab}, so a flagged-but-actually-still-live tab keeps resolving normally — the marker
* changes nothing about which pane {@code LeadCoordLoop} can reach while the flag is pending. Both
* classes must use exactly this suffix, which is why it lives here rather than as a private constant
* on either.
*/
public final class PendingCloseMarker {
public static final String SUFFIX = " [fleetd:pending-close]";
private PendingCloseMarker() {
}
/** The label with any trailing pending-close marker removed, for name matching. */
public static String strip(String label) {
if (label == null) {
return null;
}
String stripped = label.strip();
return stripped.endsWith(SUFFIX)
? stripped.substring(0, stripped.length() - SUFFIX.length()).strip()
: stripped;
}
/** Whether a label currently carries the marker. */
public static boolean isFlagged(String label) {
return label != null && label.strip().endsWith(SUFFIX);
}
}
@@ -4,6 +4,7 @@ import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.herdr.PendingCloseMarker;
import dev.ltms.fleet.herdr.Tab;
import dev.ltms.fleet.herdr.Workspace;
import dev.ltms.fleet.herdr.WorkspaceControl;
@@ -12,9 +13,11 @@ import org.slf4j.LoggerFactory;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import dev.ltms.fleet.peer.PeerLauncher;
/**
@@ -44,8 +47,33 @@ import dev.ltms.fleet.peer.PeerLauncher;
* privilege escalation (the tab label never granted anything a pane could take for itself; see that
* class's javadoc), but <em>staleness</em>: a label left behind by a crashed session would otherwise
* read as a live lead forever, and the lead would never be relaunched. So a lead counts as live only
* when herdr also reports a <em>running agent</em> in that tab — see {@link #liveLeads}. A labelled
* when herdr also reports a <em>running agent</em> in that tab — see {@link #countLeads}. A labelled
* tab with no agent in it is not a lead.
*
* <p><strong>fleetd #359 — a stale label used to pile up, not just mislead.</strong> Finding "not
* live" here used to mean only one thing: launch another. The old label was left exactly where it
* was, so a daemon that restarted enough times — or hit one herdr read that missed a genuinely
* running agent — accumulated one more identically-labelled dead tab per occurrence, and
* {@code LeadTabScanner} (before its own #359 fix) reported every one of them as a lead.
*
* <p><strong>Review finding 1 — closing on one reading is worse than the bug.</strong> The first
* version of this fix closed a name's dead tabs the moment a single {@link #countLeads} reading
* called them dead. The ticket's own live evidence rules that out: on a real host, {@code
* agent.list} was seen reporting "0 live" for a tab that a plain {@code ps} confirmed was running a
* real session. Closing on that reading would have destroyed the operator's actual lead — a worse
* failure than the extra tab it replaces. So {@link #ensureLeads()} now needs the same dead reading
* <em>twice</em>, one restart apart, before it closes anything: the first time a labelled tab reads
* dead, it is only flagged ({@link PendingCloseMarker}), left running, untouched; it is closed only
* if a <em>later</em>, independently-connected reconcile still finds it dead while the flag is
* still there. A transient miss self-heals — the next reconcile sees the agent again and clears the
* flag (see {@code toUnflag} below) — so the worst case for a single bad reading is one extra tab
* surviving one more restart, never a live session destroyed. Two alternatives were considered and
* rejected: corroborating {@code agent.list} against a second, truly independent signal was dropped
* because nothing else herdr exposes proves "is a process attached to this pane" any better — a
* second call to the same unreliable source is not independent evidence; capping the close to "all
* but the most recent dead tab" was dropped because "most recent" has no reliable ordering across
* tab ids and would leave the true failure mode (a name that is <em>never</em> reconfirmed) growing
* by one tab per bad reading forever, which is the exact defect this ticket exists to fix.
*/
public final class LeadLauncher {
@@ -79,9 +107,9 @@ public final class LeadLauncher {
return 0;
}
Map<String, Integer> live;
Map<String, LeadCount> live;
try {
live = liveLeads(leaders);
live = countLeads(leaders);
} catch (HerdrException e) {
// Counting is the whole safety mechanism against double-spawning. If we cannot count, we
// must not guess — spawning a second orchestrator is worse than starting none.
@@ -93,9 +121,53 @@ public final class LeadLauncher {
for (Map.Entry<String, FleetConfig.Leader> e : leaders.entrySet()) {
String name = e.getKey();
FleetConfig.Leader lead = e.getValue();
int running = live.getOrDefault(name, 0);
LeadCount state = live.getOrDefault(name, LeadCount.NONE);
int running = state.running();
int wanted = lead.instances();
// fleetd #359 review finding 1: a tab already flagged pending-close, still labelled for
// this lead, and STILL hosting no agent on this separate reconcile — two independent
// readings agree, so close it. A tab found dead for the first time is only flagged below,
// never closed on the spot.
for (String tabId : state.toClose()) {
log.info("lead '{}': closing tab {} — flagged pending-close on a previous reconcile "
+ "and still no agent running in it", name, tabId);
try {
spaces.closeTab(tabId);
} catch (RuntimeException cleanup) {
log.warn("could not close stale tab {} for lead '{}': {}",
tabId, name, cleanup.getMessage());
}
}
// A tab labelled for this lead, with no agent running in it, seen dead for the first
// time — flag it rather than closing it. One reading of `agent.list` is not enough
// evidence to destroy a tab that might genuinely be live (see the class javadoc).
for (String tabId : state.toFlag()) {
String flagged = lead.tabLabel() + PendingCloseMarker.SUFFIX;
log.info("lead '{}': tab {} has no agent running in it this reconcile — flagging it "
+ "'{}' rather than closing; it is only closed if a later reconcile still "
+ "finds it dead", name, tabId, flagged);
try {
spaces.renameTab(tabId, flagged);
} catch (RuntimeException cleanup) {
log.warn("could not flag stale tab {} for lead '{}': {}",
tabId, name, cleanup.getMessage());
}
}
// A previously-flagged tab that is running an agent again — the miss that flagged it was
// transient. Clear the flag so a future, unrelated miss starts its own two-reading count
// rather than closing on the strength of this one's already-spent flag.
for (String tabId : state.toUnflag()) {
log.info("lead '{}': tab {} is running an agent again — clearing its pending-close flag",
name, tabId);
try {
spaces.renameTab(tabId, lead.tabLabel());
} catch (RuntimeException cleanup) {
log.warn("could not clear the pending-close flag on tab {} for lead '{}': {}",
tabId, name, cleanup.getMessage());
}
}
if (running >= wanted) {
log.info("lead '{}': {} live, {} wanted — nothing to start", name, running, wanted);
continue;
@@ -126,18 +198,29 @@ public final class LeadLauncher {
}
/**
* How many live leads exist per configured name: a running agent in a tab labelled with that
* lead's exact {@code tab} (CB-579). Member workspaces are excluded, exactly as the scanner
* excludes them: a member must not be counted as a lead because it happens to sit in a matching
* tab.
* How many live leads exist per configured name, and which of that name's labelled tabs are
* <em>not</em> live: a running agent in a tab labelled with that lead's exact {@code tab}
* (CB-579). Member workspaces are excluded, exactly as the scanner excludes them: a member must
* not be counted as a lead because it happens to sit in a matching tab.
*
* <p>There used to be a second path here — a running agent on the terminal a
* {@code fleet.leaders.<name>.terminal} pin named, for a lead opened and pinned by hand. That
* pin is retired: {@code tab} is now the only field identity depends on, and {@link Agent}
* already carries {@link Agent#tabId()} directly, so a hand-opened lead is found the same way an
* auto-launched one is — by labelling its tab to match.
*
* <p>fleetd #359 review finding 1: a labelled tab with nothing running in it is split into
* {@code toClose} (already flagged pending-close by a previous reconcile, and still dead — two
* independent readings agree) and {@code toFlag} (dead for the first time — not enough evidence
* to close yet). {@code toUnflag} is the reverse: a tab flagged pending-close that is running an
* agent again, so the flag it carries no longer means anything and {@link #ensureLeads()} clears
* it.
*/
private Map<String, Integer> liveLeads(Map<String, FleetConfig.Leader> leaders) {
private record LeadCount(int running, List<String> toClose, List<String> toFlag, List<String> toUnflag) {
static final LeadCount NONE = new LeadCount(0, List.of(), List.of(), List.of());
}
private Map<String, LeadCount> countLeads(Map<String, FleetConfig.Leader> leaders) {
// A lead and the members share ONE workspace now (the operator asked for a single "session"
// with many tabs), so a workspace can no longer be excluded wholesale — the lead lives in the
// member workspace by design. The sole discriminator is the exact tab label: a lead carries
@@ -145,6 +228,7 @@ public final class LeadLauncher {
// profile's `worker: {profile} #{n}` template. These never collide, so an exact-label match
// separates them without needing to know which workspace anyone is in.
Map<String, String> nameByTab = new LinkedHashMap<>();
Set<String> flaggedTabIds = new LinkedHashSet<>();
for (Workspace ws : spaces.listWorkspaces()) {
if (ws.workspaceId() == null) {
continue;
@@ -153,31 +237,65 @@ public final class LeadLauncher {
String declared = leadNameOf(tab.label(), leaders);
if (declared != null && tab.tabId() != null) {
nameByTab.put(tab.tabId(), declared);
if (PendingCloseMarker.isFlagged(tab.label())) {
flaggedTabIds.add(tab.tabId());
}
}
}
}
Set<String> liveTabIds = new LinkedHashSet<>();
Map<String, Integer> counts = new LinkedHashMap<>();
for (Agent a : agents.list()) {
String name = nameByTab.get(a.tabId());
if (name != null) {
counts.merge(name, 1, Integer::sum);
liveTabIds.add(a.tabId());
}
}
return counts;
Map<String, List<String>> toCloseByName = new LinkedHashMap<>();
Map<String, List<String>> toFlagByName = new LinkedHashMap<>();
Map<String, List<String>> toUnflagByName = new LinkedHashMap<>();
nameByTab.forEach((tabId, name) -> {
boolean live = liveTabIds.contains(tabId);
boolean flagged = flaggedTabIds.contains(tabId);
if (live) {
if (flagged) {
toUnflagByName.computeIfAbsent(name, k -> new ArrayList<>()).add(tabId);
}
return;
}
if (flagged) {
toCloseByName.computeIfAbsent(name, k -> new ArrayList<>()).add(tabId);
} else {
toFlagByName.computeIfAbsent(name, k -> new ArrayList<>()).add(tabId);
}
});
Map<String, LeadCount> out = new LinkedHashMap<>();
for (String name : leaders.keySet()) {
out.put(name, new LeadCount(counts.getOrDefault(name, 0),
toCloseByName.getOrDefault(name, List.of()),
toFlagByName.getOrDefault(name, List.of()),
toUnflagByName.getOrDefault(name, List.of())));
}
return out;
}
/**
* The configured lead a tab label names, or {@code null} for a label that names none.
*
* <p>Matched exactly (case-insensitively) against each lead's configured {@code tab}, so an
* operator's {@code "lead: something-else"} tab is not mistaken for a configured lead.
* operator's {@code "lead: something-else"} tab is not mistaken for a configured lead. A
* trailing {@link PendingCloseMarker} is stripped first, so a tab this class flagged on a
* previous reconcile is still recognised as the same lead's tab on this one.
*/
private String leadNameOf(String label, Map<String, FleetConfig.Leader> leaders) {
if (label == null) {
return null;
}
String l = label.strip();
String l = PendingCloseMarker.strip(label);
for (Map.Entry<String, FleetConfig.Leader> e : leaders.entrySet()) {
String tab = e.getValue().tabLabel();
if (tab != null && l.equalsIgnoreCase(tab.strip())) {
@@ -44,6 +44,7 @@ import java.util.Objects;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.function.BiFunction;
import java.util.function.Function;
import java.util.function.LongSupplier;
@@ -101,6 +102,8 @@ public final class FleetMcp {
private final LeadSeatSource leadSeats;
/** CB-637: this daemon's lead-to-lead channel; {@code null} when no coordinator is configured. */
private final LeadChannel leadChannel;
/** fleetd #361: {@code coordinator.peers} — see {@link CoordinationSource}. Empty when unset. */
private final List<String> peers;
/** Capacity facts used by {@code fleet_list}; production must supply the placement live count. */
public record CapacitySource(Function<String, Integer> liveCount, Function<String, Integer> maxLoad,
@@ -164,6 +167,32 @@ public final class FleetMcp {
public static LeadSeatSource none() { return new LeadSeatSource(_ -> 0); }
}
/**
* fleetd #361: peer-visibility facts for {@code fleet_list}'s {@code coordinator} row — this
* daemon's own {@link LeadChannel} (for its self mailbox state and held messages) plus the
* coord-ids the operator has declared as peers ({@code coordinator.peers}). Bundled as its own
* Source, the same idiom as {@link OutageSource}/{@link QuarantineSource}/{@link LeadSeatSource},
* so {@code listFleet}'s already-long overload chain gains exactly one new required parameter
* instead of a further bare positional argument.
*
* <p>Every mailbox look this triggers goes through {@link LeadChannel#inspect}, which is
* specified to run on its own disposable channel — never the channel {@link LeadChannel#publish}
* or the consume loop depends on — so a peer that happens to be down, or the coordination broker
* itself being unreachable, can never take {@code fleet_send}/{@code LeadCoordLoop}'s own path
* down with it. See {@code FleetMcp.probe} for the additional timeout bound on top of that.
*
* @param leadChannel this daemon's own channel, or {@code null} when no coordinator is configured
* @param peers the coord-ids declared under {@code coordinator.peers}, or empty
*/
public record CoordinationSource(LeadChannel leadChannel, List<String> peers) {
public CoordinationSource {
peers = peers == null ? List.of() : List.copyOf(peers);
}
/** Inert source — no coordinator row is ever reported. */
public static CoordinationSource none() { return new CoordinationSource(null, List.of()); }
}
/**
* @param callers resolves each call's {@link Principal}; {@code null} disables authorization.
* This surface needs its own enforcement: {@code /mcp} is a raw servlet on
@@ -211,19 +240,34 @@ public final class FleetMcp {
}
/**
* As above, with fleetd #176 lead-seat facts (see {@link LeadSeatSource}). This is what
* {@code Fleetd.main} actually wires up.
*
* @param leadSeats required — pass {@link LeadSeatSource#none()} for a caller that does not want
* the feature, never a defaulting overload (the same rule {@code quarantine} and
* {@code outage} follow).
* As above, with fleetd #176 lead-seat facts (see {@link LeadSeatSource}).
*/
public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage,
LeadSeatSource leadSeats) {
this(messages, workers, sessions, identity, presence, primaryRegistry, callers, metrics, capacity,
healthCoverage, quarantine, leadChannel, outage, leadSeats, List.of());
}
/**
* As above, with fleetd #361 {@code coordinator.peers} (see {@link CoordinationSource}). This is
* what {@code Fleetd.main} actually wires up.
*
* @param leadSeats required — pass {@link LeadSeatSource#none()} for a caller that does not want
* the feature, never a defaulting overload (the same rule {@code quarantine} and
* {@code outage} follow).
* @param peers the coord-ids declared under {@code coordinator.peers}; empty when unset or
* when {@code leadChannel} is {@code null}.
*/
public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage,
LeadSeatSource leadSeats, List<String> peers) {
this.leadChannel = leadChannel;
this.peers = peers == null ? List.of() : List.copyOf(peers);
this.capacity = capacity;
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
this.outage = Objects.requireNonNull(outage, "outage");
@@ -361,7 +405,7 @@ public final class FleetMcp {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
leadSeats, callers == null ? Map.of() : callers.leads(),
callerTerminal(exchange),
leadChannel == null ? null : leadChannel.selfCoordId());
new CoordinationSource(leadChannel, peers));
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
(exchange, req) -> {
@@ -697,6 +741,15 @@ public final class FleetMcp {
* turned into a tool error naming the coord-id. It is never allowed to escape as a crash: an
* unreachable peer is an ordinary outcome of addressing a fleet you do not control.
*
* <p><strong>fleetd #361: the success text is honest about what "durably confirmed" does and
* does not mean.</strong> The broker's publisher confirm proves the message is durably queued —
* it says nothing about whether the peer's pane has, or ever will, receive it. After a
* successful publish this looks at the target mailbox's consumer count (via
* {@link LeadChannel#inspect}, bounded and never allowed to fail the call — see {@link #probe})
* and appends a warning when it is zero: that is the observable form of "nobody is reading this
* right now". A zero-consumer publish is still reported as a SUCCESS, never an error — the
* message is safely queued and will be read once a daemon owning that coord-id connects.
*
* @param leadChannel this daemon's channel, or {@code null} when no coordinator is configured
*/
static McpSchema.CallToolResult sendToLead(LeadChannel leadChannel, String coordId, String content,
@@ -724,7 +777,16 @@ public final class FleetMcp {
+ ". Check that a daemon is running with coordinator.selfId=\"" + coordId
+ "\" and is connected to the same coordination broker.");
}
return text("delivered to peer lead " + coordId + " (msgId " + msg.msgId() + ")");
String result = "published to peer lead \"" + coordId + "\"'s mailbox and durably confirmed "
+ "by the broker (msgId " + msg.msgId() + ").";
LeadChannel.MailboxState state = probe(leadChannel, coordId);
if (state.exists() && state.consumers() == 0) {
result += " Warning: that mailbox currently has NO consumers attached — nobody is reading "
+ "it right now. The message is safely queued and will be delivered once a daemon "
+ "with coordinator.selfId=\"" + coordId + "\" is running and connected; until then "
+ "it will not reach that lead's pane.";
}
return text(result);
}
/**
@@ -1110,7 +1172,8 @@ public final class FleetMcp {
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, leads, selfTerm, null);
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, leads, selfTerm,
CoordinationSource.none());
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
@@ -1119,33 +1182,33 @@ public final class FleetMcp {
QuarantineSource quarantine, OutageSource outage,
Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
LeadSeatSource.none(), leads, selfTerm, null);
LeadSeatSource.none(), leads, selfTerm, CoordinationSource.none());
}
/**
* As above, additionally reporting this daemon's own lead coordination id (CB-637) when one is
* configured and its channel opened. There is no peer-discovery surface yet — a lead addresses a
* peer by a coord-id it was told — so this row exists to answer the one question the operator
* cannot answer any other way: what is MY coord-id, the one a peer must use to reach me. It is
* omitted entirely when no coordinator is configured, so an ordinary fleet's output is unchanged.
* As above, additionally reporting this daemon's own lead coordination state (CB-637, fleetd
* #361) when a coordinator is configured and its channel opened — see {@link #coordinatorView}
* for the shape. Omitted entirely when no coordinator is configured, so an ordinary fleet's
* output is unchanged.
*
* @param selfCoordId this daemon's coord-id, or {@code null} when lead coordination is off
* @param coordination this daemon's lead channel plus its declared peers, or
* {@link CoordinationSource#none()} when lead coordination is off
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, Map<String, String> leads, String selfTerm,
String selfCoordId) {
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, OutageSource.none(),
LeadSeatSource.none(), leads, selfTerm, selfCoordId);
LeadSeatSource.none(), leads, selfTerm, coordination);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, OutageSource outage,
Map<String, String> leads, String selfTerm, String selfCoordId) {
Map<String, String> leads, String selfTerm, CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
LeadSeatSource.none(), leads, selfTerm, selfCoordId);
LeadSeatSource.none(), leads, selfTerm, coordination);
}
/** As above, plus fleetd #176 lead-seat facts (see {@link LeadSeatSource}). */
@@ -1153,7 +1216,7 @@ public final class FleetMcp {
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, OutageSource outage,
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
String selfCoordId) {
CoordinationSource coordination) {
try {
Map<String, Agent> live = workers.list().stream()
.map(Agent.class::cast)
@@ -1175,8 +1238,9 @@ public final class FleetMcp {
Map<String, Object> result = new LinkedHashMap<>();
result.put("leads", leadRows); result.put("members", out);
result.put("healthCoverage", healthCoverage.value().get());
if (selfCoordId != null && !selfCoordId.isBlank()) {
result.put("coordinator", Map.of("selfId", selfCoordId, "configured", true));
Map<String, Object> coordinatorRow = coordinatorView(coordination);
if (coordinatorRow != null) {
result.put("coordinator", coordinatorRow);
}
if (capacity.available()) result.put("capacity", profiles.stream()
.map(profile -> capacityView(profile, capacity.liveCount(), capacity.maxLoad(), roster, messages,
@@ -1187,6 +1251,140 @@ public final class FleetMcp {
}
}
/**
* fleetd #361: the {@code coordinator} row — this daemon's own coord-id and mailbox state, the
* messages currently held for it, and the live reachability of every operator-declared peer.
* {@code null} (the row is then omitted entirely) when lead coordination is off, so an ordinary
* fleet's {@code fleet_list} output is byte-identical to before this feature existed.
*
* <p>Every peer/self mailbox look goes through {@link #probe}, which bounds each
* {@link LeadChannel#inspect} call to {@link #PEER_PROBE_TIMEOUT_MS} and never lets it throw —
* a coordination broker that is down or slow degrades this row toward "unreachable"/"unknown"
* counts, it can never make {@code fleet_list} itself slow or fail. {@code held} comes from
* {@link LeadChannel#peek}, a pure in-memory read with no broker round trip, so it is never
* subject to that bound.
*/
private static Map<String, Object> coordinatorView(CoordinationSource coordination) {
LeadChannel channel = coordination.leadChannel();
if (channel == null) {
return null;
}
String selfId = channel.selfCoordId();
Map<String, Object> row = new LinkedHashMap<>();
row.put("selfId", selfId);
row.put("configured", true);
row.put("mailbox", mailboxView(probe(channel, selfId)));
row.put("held", channel.peek().stream().map(FleetMcp::heldView).toList());
row.put("peers", coordination.peers().stream().map(p -> peerView(channel, p)).toList());
return row;
}
/**
* fleetd #361: render a {@link LeadChannel.MailboxState} without ever presenting an unmeasured
* fact as a measured one. {@code status} is the tri-state itself — {@code "exists"},
* {@code "absent"} (the broker positively confirmed no such queue), or {@code "unknown"} (the
* probe could not determine either way: down, unreachable, or timed out). {@code pending}/
* {@code consumers} are included ONLY when {@code status == "exists"} — a reader must never see
* them default to {@code 0} for a mailbox this call never actually measured. This is the fix for
* the review finding that a collapsed {@code absent()} rendered a self-probe timeout as
* "pending: 0, consumers: 0", indistinguishable from an actually-empty, actually-unread mailbox.
*/
private static Map<String, Object> mailboxView(LeadChannel.MailboxState state) {
Map<String, Object> row = new LinkedHashMap<>();
row.put("status", state.exists() ? "exists" : state.known() ? "absent" : "unknown");
if (state.exists()) {
row.put("pending", state.pending());
row.put("consumers", state.consumers());
}
return row;
}
/** One held-for-me message: enough to identify it and see roughly what it says, never the whole body. */
private static Map<String, Object> heldView(LeadMessage m) {
Map<String, Object> row = new LinkedHashMap<>();
row.put("msgId", m.msgId());
row.put("from", m.from());
row.put("preview", preview(m.content()));
return row;
}
/** Cap a held message's content to a short preview — {@code fleet_list} must never dump a full body. */
private static final int HELD_PREVIEW_MAX_CHARS = 80;
private static String preview(String content) {
if (content == null) {
return "";
}
return content.length() <= HELD_PREVIEW_MAX_CHARS
? content
: content.substring(0, HELD_PREVIEW_MAX_CHARS) + "…";
}
/**
* One declared peer's row: its coord-id, then the same tri-state {@link #mailboxView} shape.
* Deliberately no boolean "reachable" field — that collapsed "confirmed gone" and "could not
* check" into the same {@code false}, which is exactly the review finding this row now avoids:
* an operator reading {@code status} can tell "fleet01 is down" (a {@code coordinator.selfId}
* nobody has ever run) apart from "my own broker is slow or unreachable right now".
*/
private static Map<String, Object> peerView(LeadChannel channel, String coordId) {
Map<String, Object> row = new LinkedHashMap<>();
row.put("coordId", coordId);
row.putAll(mailboxView(probe(channel, coordId)));
return row;
}
/**
* fleetd #361: how long {@code fleet_list} waits on any single {@link LeadChannel#inspect} call
* before giving up on it — see {@link #probe}.
*/
private static final long PEER_PROBE_TIMEOUT_MS = 1_500L;
/**
* Dedicated pool for {@link LeadChannel#inspect} calls so a slow one blocks only its own virtual
* thread, never the MCP request thread calling {@code fleet_list}. Not a <em>bounded</em> pool —
* {@code newThreadPerTaskExecutor} starts a fresh virtual thread per call with no cap on how many
* run at once; virtual threads make that cheap, not bounded. What actually keeps a hung probe
* from accumulating forever is the {@link Future#cancel} in {@link #probe}, not a pool limit.
*/
private static final java.util.concurrent.ExecutorService PEER_PROBE_POOL =
java.util.concurrent.Executors.newThreadPerTaskExecutor(Thread.ofVirtual().name("fleet-peer-probe-", 0).factory());
/**
* fleetd #361: {@link LeadChannel#inspect}, bounded to {@link #PEER_PROBE_TIMEOUT_MS} and never
* allowed to throw or hang the caller — a coordination broker that is unreachable or slow
* degrades to {@link LeadChannel.MailboxState#unknown} (never {@code absent}: a timeout proves
* nothing about whether the mailbox exists) rather than making {@code fleet_list} slow or
* failing it. {@code inspect} itself is already specified to never throw, but this is the seam
* that also survives an implementation that does, or one that blocks indefinitely on a dead
* connection.
*
* <p><strong>A timeout cancels the orphaned task</strong> rather than abandoning it. Before this,
* {@code get(timeout)} on a hung {@code inspect} left the submitted task running forever on its
* own virtual thread, holding the AMQP channel it had already opened — against a broker that
* hangs rather than fails fast, every {@code fleet_list} call would orphan one more channel until
* the connection's channel-max (2047 by default) was exhausted, which would break {@link
* LeadChannel#publish} too. {@link Future#cancel(boolean) cancel(true)} interrupts the orphaned
* task's thread; {@link LeadMailbox#inspect} has no interruptible wait of its own to catch that,
* but the underlying AMQP RPC continuation does block on one, so the interrupt reaches it and the
* task's {@code finally} still closes the probe channel it opened rather than leaking it forever.
*/
private static LeadChannel.MailboxState probe(LeadChannel channel, String coordId) {
return probe(channel, coordId, PEER_PROBE_TIMEOUT_MS);
}
/** As {@link #probe(LeadChannel, String)}, with an explicit timeout — a seam for tests. */
static LeadChannel.MailboxState probe(LeadChannel channel, String coordId, long timeoutMs) {
java.util.concurrent.Future<LeadChannel.MailboxState> future =
PEER_PROBE_POOL.submit(() -> channel.inspect(coordId));
try {
return future.get(timeoutMs, TimeUnit.MILLISECONDS);
} catch (Exception e) {
future.cancel(true); // best-effort: don't leave a hung probe (and its channel) running forever
return LeadChannel.MailboxState.unknown(coordId);
}
}
/**
* Capacity is advisory only. {@code reclaimable} says there is no bridge work, not that fleetd
* may stop the member: the bridge has capacity facts but no work list, and choosing work needs
@@ -1377,8 +1575,12 @@ public final class FleetMcp {
+ "routes your answer back into the same turn (omit for a normal delegation)"),
"coordId", stringProp("A peer LEAD's coordination id — delivers content to that "
+ "lead's durable mailbox on the shared coordination broker, which works "
+ "across hosts. Mutually exclusive with sessionId and turnId. Your own "
+ "coordId is reported by fleet_list.")),
+ "across hosts. Mutually exclusive with sessionId and turnId. Success means "
+ "the message is durably queued and confirmed by the broker, with a warning "
+ "if that mailbox has no consumers attached right now (queued, but nobody is "
+ "reading it yet) — it does not mean the peer's pane has seen it. Your own "
+ "coordId, this daemon's mailbox state, and every coordinator.peers entry's "
+ "live reachability are reported by fleet_list's coordinator row.")),
List.of("content")));
}
@@ -1486,7 +1688,13 @@ public final class FleetMcp {
+ "for a quarantined profile's credential (see fleet_profiles), whatever its "
+ "maxLoad/live — with credentialId and quarantinedForSeconds naming the "
+ "quarantine, so 'free: 0, busy' can be told apart from 'free: 0, refusing "
+ "for N seconds'.",
+ "for N seconds'. When lead-to-lead coordination is configured, a 'coordinator' "
+ "object reports this daemon's own coord-id ('selfId') and mailbox state "
+ "('mailbox': pending/consumers), the messages currently held for it ('held': "
+ "msgId/from/preview, never the full body), and one row per coordinator.peers "
+ "coord-id ('peers': coordId/reachable, plus pending/consumers when reachable) — "
+ "this is peer DISCOVERY for cross-host leads, distinct from the local 'leads' "
+ "array above. It is omitted entirely when no coordinator is configured.",
objectSchema(Map.of(), List.of()));
}
@@ -4,8 +4,8 @@ import java.util.List;
/**
* The lead-to-lead message channel this daemon speaks, as its callers need it — one lead's own
* mailbox: publish to a peer's coord-id, look at what has arrived for me, and ack what I have
* delivered.
* mailbox: publish to a peer's coord-id, look at what has arrived for me, ack what I have
* delivered, and (fleetd #361) inspect any coord-id's mailbox from the outside without owning it.
*
* <p>Extracted from {@link LeadMailbox} purely as a seam. {@code LeadMailbox} is the one production
* implementation and owns a live AMQP connection, so a test that wanted to exercise the routing in
@@ -37,4 +37,78 @@ public interface LeadChannel {
/** This daemon's own lead coordination id — the mailbox it owns, and the {@code from} it sends as. */
String selfCoordId();
/**
* A non-destructive look at {@code coordId}'s mailbox — does it exist, how many messages are
* waiting on it, and how many consumers are attached — without owning, consuming, or otherwise
* changing it. {@code consumers == 0} on an existing mailbox is the observable form of "nobody
* is reading this right now": a publish to it will sit queued rather than reach a pane.
*
* <p><strong>Never throws</strong> — this is a best-effort fact-finding call, not an operation a
* caller must handle failing. But it must never turn "I could not check" into a false negative:
* {@link MailboxState#absent(String)} means the broker positively confirmed there is no such
* queue, and {@link MailboxState#unknown(String)} — a distinct value — means the look could not
* be completed at all (broker unreachable, timed out, connection closed). A caller that
* collapses those two into one, as fleetd #361 initially did, cannot tell "that peer is down"
* from "I could not check", and a reader of {@code pending}/{@code consumers} cannot tell a
* measured zero from a zero standing in for "not measured".
*
* <p><strong>Must never share fate with {@link #publish} or {@link #peek}/{@link #ack}.</strong>
* fleetd #361: in AMQP 0-9-1 a passive queue declare of a queue that does not exist closes the
* channel it was declared on with a 404. An implementation backed by a real broker connection
* must inspect on a channel it can afford to lose — never the channel {@link #publish} or the
* consume loop depends on — so that looking at a peer that happens to be down can never break
* this daemon's own send or receive path.
*/
MailboxState inspect(String coordId);
/**
* The result of {@link #inspect}. {@code presence} tells apart three states a caller must not
* conflate: a confirmed-existing mailbox ({@link Presence#EXISTS}, the only case where
* {@code pending}/{@code consumers} are measured facts), a confirmed-absent one
* ({@link Presence#ABSENT} — the broker positively said "no such queue"), and one this call
* simply could not determine ({@link Presence#UNKNOWN} — broker unreachable, timed out,
* connection closed). {@code pending}/{@code consumers} are always {@code 0} and meaningless
* outside {@link Presence#EXISTS}; a renderer must gate on {@link #exists()} (or {@code
* presence} directly), never present them as measured otherwise.
*
* @param coordId the coord-id inspected
* @param presence whether the mailbox is confirmed to exist, confirmed absent, or unknown
* @param pending messages ready for delivery but not yet in a consumer's hands (0 unless EXISTS)
* @param consumers how many consumers are attached (0 unless EXISTS)
*/
record MailboxState(String coordId, Presence presence, int pending, int consumers) {
/** Whether {@link #inspect} was able to reach a definite answer, of either kind. */
public enum Presence { EXISTS, ABSENT, UNKNOWN }
/** {@code true} only when the broker confirmed this exact queue is currently declared. */
public boolean exists() {
return presence == Presence.EXISTS;
}
/**
* {@code true} when {@link #inspect} reached a definite answer (exists or confirmed
* absent); {@code false} when it could not determine either way. A caller must never treat
* {@code !known()} the same as a confirmed absence — the mailbox may well exist.
*/
public boolean known() {
return presence != Presence.UNKNOWN;
}
/** The broker confirmed this queue exists, with these measured counts. */
public static MailboxState exists(String coordId, int pending, int consumers) {
return new MailboxState(coordId, Presence.EXISTS, pending, consumers);
}
/** The broker positively confirmed there is no such queue (e.g. a 404 on passive declare). */
public static MailboxState absent(String coordId) {
return new MailboxState(coordId, Presence.ABSENT, 0, 0);
}
/** The look could not be completed — broker unreachable, timed out, or connection closed. */
public static MailboxState unknown(String coordId) {
return new MailboxState(coordId, Presence.UNKNOWN, 0, 0);
}
}
}
@@ -10,6 +10,7 @@ import com.rabbitmq.client.DeliverCallback;
import com.rabbitmq.client.Recoverable;
import com.rabbitmq.client.RecoveryListener;
import com.rabbitmq.client.Return;
import com.rabbitmq.client.ShutdownSignalException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -259,6 +260,89 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
}
}
/**
* fleetd #361: look at {@code coordId}'s mailbox on a fresh, immediately-closed throwaway
* channel — never {@link #channel} (consume/ack) or {@link #publishChannel} (publish). A
* passive queue declare of a queue that does not exist closes the channel it was declared on
* with a 404; using a disposable probe channel means that closure can never touch either
* long-lived channel this instance depends on for {@link #publish} or the consume loop.
*
* <p>Classifies failures rather than collapsing them, both measured against a real broker in
* {@code LeadMailboxTest} rather than assumed from the AMQP 0-9-1 spec text:
* <ul>
* <li>a genuine 404 — an {@link IOException} wrapping a {@link ShutdownSignalException} whose
* {@link AMQP.Channel.Close#getReplyCode()} is {@code 404} — reports
* {@link MailboxState#absent}; every other declare failure reports
* {@link MailboxState#unknown} instead of quietly becoming the same "absent" value;
* <li>{@code catch (RuntimeException e)} on both attempts matters as much as the checked
* catches: a connection that is already closed makes {@link Connection#createChannel()}
* throw {@link com.rabbitmq.client.AlreadyClosedException} (a {@link RuntimeException},
* not an {@link IOException}) — an {@code inspect} that only caught {@code IOException}
* would let that escape, breaking the "never throws" contract this method promises.
* </ul>
*
* <p><strong>Honesty about which catch is measured and which is defensive:</strong> the
* {@code createChannel()} catch above is exercised end-to-end against a real broker by
* {@code LeadMailboxTest.inspectReportsUnknownRatherThanThrowingWhenTheConnectionIsAlreadyClosed}.
* The second {@code catch (RuntimeException e)}, around the passive declare itself — for the
* narrower race where the connection drops <em>between</em> {@code createChannel()} succeeding
* and the declare landing — has no such test; reaching it needs a connection that dies at that
* exact instant, which is not a scenario this suite drives on purpose. It stays purely
* defensive: correct by the same reasoning as the first catch, but unproven the way the first
* one is proven.
*/
@Override
public MailboxState inspect(String coordId) {
String queue = queueName(coordId);
Channel probe;
try {
probe = connection.createChannel();
} catch (IOException | RuntimeException e) {
log.debug("lead mailbox inspect: cannot open a probe channel for {}: {}", coordId, e.toString());
return MailboxState.unknown(coordId);
}
try {
AMQP.Queue.DeclareOk declared = probe.queueDeclarePassive(queue);
return MailboxState.exists(coordId, declared.getMessageCount(), declared.getConsumerCount());
} catch (IOException e) {
// The broker (or the client library) has already closed `probe` for us either way; only
// a confirmed 404 means "no such queue" — anything else (a different declare failure) is
// "could not determine", never silently reported as the same value as a genuine absence.
return isMissingQueue(e) ? MailboxState.absent(coordId) : MailboxState.unknown(coordId);
} catch (RuntimeException e) {
// E.g. the connection dropped between createChannel() and the declare landing.
log.debug("lead mailbox inspect: declare failed unexpectedly for {}: {}", coordId, e.toString());
return MailboxState.unknown(coordId);
} finally {
try {
if (probe.isOpen()) {
probe.close();
}
} catch (Exception e) {
log.debug("lead mailbox inspect: probe channel close for {}: {}", coordId, e.toString());
}
}
}
/**
* {@code true} only for the specific shape a missing-queue passive declare actually produces —
* measured against a real broker, not assumed from the spec text (see {@code
* LeadMailboxTest.passiveDeclareOfAMissingQueueThrowsAnIOExceptionWrappingA404ShutdownSignal}):
* an {@link IOException} whose cause is a {@link ShutdownSignalException} carrying an
* {@link AMQP.Channel.Close} reason with {@code replyCode == 404}. Any other shape (a different
* reply code, a {@code ShutdownSignalException} cause whose reason is not a
* {@code Channel.Close}, or no cause at all) is a declare failure of some other kind and must
* not be read as "confirmed absent" — pinned hermetically, with no broker needed, by
* {@code LeadMailboxIsMissingQueueTest} for exactly those three false shapes. Package-private
* (not {@code private}) so that test can call it directly.
*/
static boolean isMissingQueue(IOException e) {
if (!(e.getCause() instanceof ShutdownSignalException sse)) {
return false;
}
return sse.getReason() instanceof AMQP.Channel.Close close && close.getReplyCode() == AMQP.NOT_FOUND;
}
/** Convenience: {@link #peek} the current snapshot, then {@link #ack} every message in it. */
public List<LeadMessage> drain() {
List<LeadMessage> snapshot = peek();
@@ -92,6 +92,15 @@ public final class GitWorktrees implements Worktrees {
private final String configuredRoot;
/** OS group name for {@link #shareWithGroup} (fleetd #185 stage 3); {@code null} ⇒ feature off. */
private final String group;
/** Source directory of skill folders for {@link #seedSkills} (fleetd #362, {@code memberSkills:}
* in config); {@code null} ⇒ feature off, no worktree is touched beyond today's behaviour. */
private final String memberSkillsSource;
/** Extra environment merged into every {@code git} subprocess this instance runs. Always {@code
* Map.of()} from every production constructor. Test seam only (fleetd #362 review fix): lets
* {@code GitWorktreesTest} point {@code GIT_CONFIG_GLOBAL} at an isolated temp file so it can
* drive the real {@link #add} path against a controlled "operator's global git config" and
* prove the excludesFile composition below without ever touching the real machine's config. */
private final Map<String, String> gitEnv;
private final Consumer<String> afterWorktreeAdded;
/** How the initial {@code git worktree add} command runs. Package-private test seam for an
* interrupted command after Git has made worktree state. */
@@ -125,7 +134,20 @@ public final class GitWorktrees implements Worktrees {
* config); null/blank ⇒ {@link #shareWithGroup} is a no-op.
*/
public GitWorktrees(String configuredRoot, String group) {
this(configuredRoot, group, _ -> {});
this(configuredRoot, group, (String) null);
}
/**
* @param configuredRoot nullable absolute or relative path; null/blank derives a sibling of
* the repo root.
* @param group optional OS group name (fleetd #185 stage 3, {@code worktreeGroup:}
* in config); null/blank ⇒ {@link #shareWithGroup} is a no-op.
* @param memberSkillsSource fleetd #362: optional directory of skill folders ({@code
* memberSkills:} in config) copied into every provisioned worktree's
* {@code .claude/skills/}; null/blank ⇒ {@link #seedSkills} is a no-op.
*/
public GitWorktrees(String configuredRoot, String group, String memberSkillsSource) {
this(configuredRoot, group, _ -> {}, null, null, memberSkillsSource);
}
/** Test seam for changing a real worktree between its creation and its security check. */
@@ -135,13 +157,20 @@ public final class GitWorktrees implements Worktrees {
/** Test seam combining a configurable {@code group} with {@link #afterWorktreeAdded}. */
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded) {
this(configuredRoot, group, afterWorktreeAdded, null, null);
this(configuredRoot, group, afterWorktreeAdded, null, null, null);
}
/** Test seam for changing how {@link #shareWithGroup}'s processes run. */
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded,
Function<String[], String> shareGroupRunner) {
this(configuredRoot, group, afterWorktreeAdded, shareGroupRunner, null);
this(configuredRoot, group, afterWorktreeAdded, shareGroupRunner, null, null);
}
/** Test seam for changing how the initial {@code git worktree add} command runs, with no
* {@code memberSkillsSource} configured. */
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded,
Function<String[], String> shareGroupRunner, Function<String[], String> worktreeAddRunner) {
this(configuredRoot, group, afterWorktreeAdded, shareGroupRunner, worktreeAddRunner, null);
}
/**
@@ -151,11 +180,30 @@ public final class GitWorktrees implements Worktrees {
*
* @param shareGroupRunner {@code null} ⇒ the real {@link #exec(String...)}.
* @param worktreeAddRunner {@code null} ⇒ the real {@link #exec(String...)}.
* @param memberSkillsSource {@code null}/blank ⇒ {@link #seedSkills} is a no-op.
*/
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded,
Function<String[], String> shareGroupRunner, Function<String[], String> worktreeAddRunner) {
Function<String[], String> shareGroupRunner, Function<String[], String> worktreeAddRunner,
String memberSkillsSource) {
this(configuredRoot, group, afterWorktreeAdded, shareGroupRunner, worktreeAddRunner,
memberSkillsSource, Map.of());
}
/**
* Full test seam, plus {@code gitEnv} (fleetd #362 review fix, verification only): extra
* environment merged into every {@code git} subprocess this instance runs, so a test can isolate
* something like {@code GIT_CONFIG_GLOBAL} from the real machine while still driving the real
* {@link #add} path end to end. Every production constructor above delegates here with {@code
* Map.of()}.
*/
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded,
Function<String[], String> shareGroupRunner, Function<String[], String> worktreeAddRunner,
String memberSkillsSource, Map<String, String> gitEnv) {
this.configuredRoot = configuredRoot;
this.group = (group == null || group.isBlank()) ? null : group;
this.memberSkillsSource = (memberSkillsSource == null || memberSkillsSource.isBlank())
? null : memberSkillsSource;
this.gitEnv = gitEnv == null ? Map.of() : gitEnv;
this.afterWorktreeAdded = afterWorktreeAdded == null ? _ -> {} : afterWorktreeAdded;
this.shareGroupRunner = shareGroupRunner != null ? shareGroupRunner : this::exec;
this.worktreeAddRunner = worktreeAddRunner != null ? worktreeAddRunner : this::exec;
@@ -184,6 +232,7 @@ public final class GitWorktrees implements Worktrees {
configureEnvironmentCredentialHelper(repoRoot, wt);
configureHttpsUrlRewriteForSshOrigin(repoRoot, wt);
isolateToolSurface(wt);
seedSkills(wt);
} catch (RuntimeException e) {
cleanupAfterAddFailure(repoRoot, wt, branch, e);
throw e;
@@ -574,6 +623,310 @@ public final class GitWorktrees implements Worktrees {
return true;
}
/**
* fleetd #362: copy each skill folder from the configured {@link #memberSkillsSource} directory
* into {@code <worktreePath>/.claude/skills/}, so a member spawned against ANY repo — not only
* one that already ships its own {@code .claude/skills/} — can load a bridge skill such as
* {@code implementer}. Every brief this fleet sends starts with {@code "Load the <name>
* skill."}; outside a repo carrying its own copy that line was previously a no-op.
*
* <p><b>No-op — nothing read, nothing written, nothing logged</b> — when {@link
* #memberSkillsSource} is null/blank (today's default), the same off-switch shape as
* {@link #shareWithGroup}.
*
* <p><b>Invariant 1 — a repo's own skill wins.</b> A skill folder already present at
* {@code <worktreePath>/.claude/skills/<name>} — because the just-checked-out branch commits its
* own copy — is left completely untouched: never overwritten, and never even opened.
*
* <p><b>Invariant 2 — a seeded skill can never end up in a worker's commit.</b> Every path this
* writes is untracked in the target repo (that is the whole reason it is being seeded), so
* {@code git status} would otherwise show each one as a new, addable, committable path. The
* repo-wide {@code .git/info/exclude} is NOT used for this: measured against a real linked
* worktree, that file resolves to the repository's COMMON git dir even from a worktree (the
* same file {@link dev.ltms.fleet.member.ClaudeCodeLauncher#writeIdeOverlay writeIdeOverlay}
* appends {@code CLAUDE.local.md} to), so an entry written there would hide the seeded skill
* from {@code git status} in the PRIMARY's own checkout and every sibling worktree too — not
* only this one. Instead, {@link #excludeSeededSkillsFromGitStatus} points {@code
* core.excludesFile} at a file scoped {@code --worktree} (the same {@code
* extensions.worktreeConfig} mechanism {@link #configureEnvironmentCredentialHelper} already
* relies on) that itself lives under this worktree's own private git dir ({@code
* .git/worktrees/<nonce>/}, OUTSIDE the working tree) — invisible to this worktree's {@code git
* status} and structurally impossible for this worktree to commit, with no effect on any other
* worktree or the primary checkout. Proven with a real {@code git status --porcelain} in
* {@code GitWorktreesTest}, not by reasoning.
*
* <p><b>Compose, don't replace.</b> {@code core.excludesFile} is single-valued: the first cut of
* this method pointed it at fleetd's own file with {@code --replace-all}, which SHADOWS whatever
* the operator's own (global, or repo-local) {@code core.excludesFile} was already resolving to
* inside this worktree, rather than adding to it. Measured concretely: this repo's own {@code
* .gitignore} does not ignore {@code target/} — only an operator's global excludesFile does — so
* every worker's {@code mvn clean install} would otherwise make {@code target/} appear as
* untracked, and {@link #hasUncommitted}'s deliberately-untracked-inclusive {@code git status
* --porcelain} (CB-576) would then read every such worktree as dirty forever, so it is never
* cleaned up. {@link #excludeSeededSkillsFromGitStatus} now reads whatever {@code
* core.excludesFile} resolves to BEFORE writing anything (falling back to git's own documented
* default, {@code $XDG_CONFIG_HOME/git/ignore} or {@code $HOME/.config/git/ignore}, when the key
* is unset entirely — see {@code gitignore(5)}), and writes that content into its OWN exclude
* file ahead of the seeded skill patterns, so every operator-configured pattern keeps applying
* inside the seeded worktree exactly as it did before seeding ran.
*
* <p>Instead of using worktree-scoped-config as an add-then-append (a second key does not exist
* for {@code core.excludesFile} — it takes exactly one value), an actual second exclude source
* was ruled out because git resolves only ONE {@code core.excludesFile}; concatenating the prior
* content into fleetd's own file is what "compose" reduces to for a single-valued key.
*
* <p>Proven the same way as invariant 2's own leak check: {@code
* GitWorktreesTest#seedSkillsComposesWithAnAlreadyEffectiveGlobalExcludesFile} isolates a
* synthetic "operator's global config" via {@code GIT_CONFIG_GLOBAL} (never the real machine's),
* seeds a skill, and asserts {@code git status --porcelain} is still empty for a file matching
* that global config's own ignore pattern.
*
* <p><b>Invariant 3 — best-effort.</b> A missing/unreadable {@link #memberSkillsSource}, or a
* copy/exclude failure, is logged and skipped — it must never fail the spawn, the same contract
* {@link #overlayParity} and {@link #isolateToolSurface} already hold.
*
* <p>Recorded for the worker itself the same way fleetd #134 records {@code
* fleet.neutralizedConfig}: {@code fleet.seededSkills} (one value per seeded skill folder) and
* {@code fleet.seededSkillsNote}, readable with {@code git config --worktree --get-all
* fleet.seededSkills}.
*
* <p><b>Claude Code specific by construction, not by a backend check here.</b> Only {@code
* .claude/skills/<name>/SKILL.md} is a path any launcher reads today (opencode's equivalent is a
* different shape under {@code .opencode/agent}, out of scope — see issue #362). This method
* only copies files; like {@link #isolateToolSurface} — which neutralizes BOTH {@code .mcp.json}
* and {@code opencode.json} unconditionally — it runs the same for every worktree regardless of
* which backend ultimately spawns into it, because the backend is not yet chosen at {@link #add}
* time. A seeded {@code .claude/skills/} directory in an opencode member's worktree is simply
* never read by that launcher.
*/
private void seedSkills(String worktreePath) {
if (memberSkillsSource == null) {
return;
}
Path source = Path.of(memberSkillsSource).toAbsolutePath().normalize();
if (!Files.isDirectory(source)) {
log.warn("memberSkills source '{}' is not a directory — skipping skill seeding for worktree {}",
source, worktreePath);
return;
}
Path skillsRoot = Path.of(worktreePath).resolve(".claude").resolve("skills");
List<String> seeded = new ArrayList<>();
List<String> kept = new ArrayList<>();
try (var candidates = Files.list(source)) {
for (Path candidate : candidates
.filter(Files::isDirectory)
.filter(p -> !p.getFileName().toString().startsWith("."))
.sorted()
.toList()) {
String name = candidate.getFileName().toString();
Path dst = skillsRoot.resolve(name);
if (Files.exists(dst)) {
kept.add(name);
continue;
}
copySkillDirectory(candidate, dst);
seeded.add(name);
}
} catch (IOException | RuntimeException e) {
log.warn("failed to seed skills into worktree {} from memberSkills source '{}': {}",
worktreePath, source, e.getMessage());
return;
}
String detail = seeded.isEmpty() ? "" : "seeded: " + String.join(", ", seeded);
if (!kept.isEmpty()) {
detail += (detail.isEmpty() ? "" : "; ") + "kept the repo's own copy of: " + String.join(", ", kept);
}
if (detail.isEmpty()) {
detail = "no skill folders found under " + source;
}
log.info("skill seeding: {} of {} candidate(s) from {} into {}/.claude/skills — {}",
seeded.size(), seeded.size() + kept.size(), source, worktreePath, detail);
if (seeded.isEmpty()) {
return;
}
try {
excludeSeededSkillsFromGitStatus(worktreePath, seeded);
recordSeededSkillsForWorker(worktreePath, seeded);
} catch (RuntimeException e) {
log.warn("seeded skill(s) {} into {} but could not hide them from git status: {} — "
+ "they may show as untracked; never commit them", seeded, worktreePath, e.getMessage());
}
}
/** Recursively copy a skill folder ({@code src}) into a fresh destination ({@code dst}) that
* {@link #seedSkills} has already confirmed does not exist, preserving the directory structure
* (e.g. {@code implementer/SKILL.md}, {@code implementer/references/...}). */
private static void copySkillDirectory(Path src, Path dst) {
try (var walk = Files.walk(src)) {
for (Path path : walk.sorted().toList()) {
Path target = dst.resolve(src.relativize(path).toString());
if (Files.isDirectory(path)) {
Files.createDirectories(target);
} else {
Files.createDirectories(target.getParent());
Files.copy(path, target, StandardCopyOption.COPY_ATTRIBUTES);
}
}
} catch (IOException e) {
throw new WorktreeException("cannot copy skill directory " + src + " -> " + dst + ": "
+ e.getMessage(), e);
}
}
/**
* Make every path in {@code seededSkillNames} (each a name under {@code .claude/skills/})
* invisible to {@code git status} in THIS worktree only — see the invariant-2 discussion on
* {@link #seedSkills}. Sets {@code core.excludesFile} scoped {@code --worktree} to a file
* written under this worktree's own private git dir ({@code git rev-parse
* --absolute-git-dir}), which lives outside the working tree, so the exclude file itself can
* never be committed either.
*
* <p><b>Compose, don't replace.</b> {@code core.excludesFile} is single-valued, so pointing it at
* fleetd's own file would otherwise SHADOW whatever excludesFile this worktree was already
* resolving (an operator's global config, most commonly) rather than add to it — see the
* "Compose, don't replace" discussion on {@link #seedSkills}. {@link
* #previouslyEffectiveExcludesFileContent} is read BEFORE this method's own {@code --worktree}
* write below, so it still sees whatever was effective beforehand; that content is written into
* fleetd's own exclude file ahead of the seeded skill patterns, and the worktree-scoped override
* then points at that combined file — so every pattern the operator's own configuration already
* applied keeps applying, plus the seeded skill paths.
*
* <p><b>Assumes a fresh worktree — not idempotent.</b> {@link #seedSkills} only ever calls this
* from {@link #add}, which always creates a brand-new worktree, so {@code core.excludesFile} is
* never already worktree-scoped-set to fleetd's own file when this runs. A hypothetical second
* call on the SAME worktree would read fleetd's own already-composed file back as "previously
* effective" (worktree scope now wins) and append the seeded patterns a second time — harmless
* to {@code git status} (duplicate exclude lines are a no-op), but not something to rely on. No
* guard is added for this because the path does not exist today; if a future caller ever seeds
* the same worktree twice, it will need one.
*/
private void excludeSeededSkillsFromGitStatus(String worktreePath, List<String> seededSkillNames) {
exec("git", "-C", worktreePath, "config", "extensions.worktreeConfig", "true");
String previouslyEffective = previouslyEffectiveExcludesFileContent(worktreePath);
String gitDir = exec("git", "-C", worktreePath, "rev-parse", "--absolute-git-dir").trim();
Path excludeFile = Path.of(gitDir, "fleet-seeded-skills-exclude");
StringBuilder patterns = new StringBuilder();
if (!previouslyEffective.isEmpty()) {
patterns.append(previouslyEffective);
}
for (String name : seededSkillNames) {
patterns.append("/.claude/skills/").append(name).append('/').append(System.lineSeparator());
}
try {
Files.writeString(excludeFile, patterns.toString());
} catch (IOException e) {
throw new WorktreeException("cannot write skills exclude file " + excludeFile + ": "
+ e.getMessage(), e);
}
exec("git", "-C", worktreePath, "config", "--worktree", "--replace-all", "core.excludesFile",
excludeFile.toString());
}
/**
* The content of whatever {@code core.excludesFile} resolves to for {@code worktreePath} right
* now — BEFORE {@link #excludeSeededSkillsFromGitStatus} points that key at fleetd's own file —
* so it can be carried forward instead of shadowed. {@code --type=path} makes git itself perform
* {@code ~}/{@code ~user} expansion the same way it would when actually reading the key to build
* exclude rules, rather than handing back a raw, unexpanded config string.
*
* <p>When the key is unset entirely (exit code non-zero), falls back to git's own documented
* default excludes file — {@code $XDG_CONFIG_HOME/git/ignore}, or {@code
* $HOME/.config/git/ignore} when that variable is unset — per {@code gitignore(5)}: git applies
* that file even with no {@code core.excludesFile} configured at all, so skipping it here would
* silently drop patterns an operator never had to configure to get.
*
* <p>Never throws: a missing, unreadable, or unresolvable file is treated as "nothing to carry
* forward" (empty string) — this is a best-effort read in service of {@link #seedSkills}'s own
* invariant 3, not a new way for skill seeding to fail a spawn.
*
* <p><b>Review fix, finding 2.</b> The XDG-fallback branch below does not go through {@code git}
* at all, so a first cut of it read {@code XDG_CONFIG_HOME}/{@code HOME} straight from the JVM's
* own environment ({@link System#getenv} / {@code user.home}) — unlike every other value this
* class resolves, which goes through a {@code git} subprocess and therefore already honours
* {@link #gitEnv}. That meant no test could make this branch hermetic, and on any machine
* carrying a real {@code ~/.config/git/ignore} (this repo's own dev machine does), every
* skill-seeding test silently composed with that real file — correct in production, but
* machine-dependent in the test suite, and a future broader pattern in that real file could
* silently change what a seeded worktree's {@code git status} reports depending on whose home
* directory ran the test. {@link #resolveEnv} now checks {@link #gitEnv} first for both
* variables, falling back to the JVM's real environment only when the seam does not supply
* them — production behaviour (empty {@link #gitEnv}) is unchanged, and a test can now isolate
* this branch exactly as it already isolates every {@code git} subprocess call.
*
* <p><b>Snapshot, not a reference.</b> The content below is read once, at seeding time, and
* copied into fleetd's own exclude file. If the operator edits their global excludesFile
* afterward, an already-seeded worktree keeps the old copy — acceptable for a worktree's
* expected lifetime, but worth knowing before reading a stale pattern as a bug.
*
* @return the file's content, trailing-newline-normalized, or {@code ""} when there is nothing
* to compose with.
*/
private String previouslyEffectiveExcludesFileContent(String worktreePath) {
String resolvedPath;
if (exitCode("git", "-C", worktreePath, "config", "--get", "--type=path", "core.excludesFile") == 0) {
resolvedPath = exec("git", "-C", worktreePath, "config", "--get", "--type=path",
"core.excludesFile").trim();
} else {
String xdgConfigHome = resolveEnv("XDG_CONFIG_HOME");
Path fallback = (xdgConfigHome != null && !xdgConfigHome.isBlank())
? Path.of(xdgConfigHome, "git", "ignore")
: Path.of(resolveHome(), ".config", "git", "ignore");
resolvedPath = fallback.toString();
}
if (resolvedPath.isBlank()) {
return "";
}
Path file = Path.of(resolvedPath);
if (!Files.isRegularFile(file) || !Files.isReadable(file)) {
return "";
}
try {
String content = Files.readString(file);
return content.isBlank() ? "" : content.stripTrailing() + System.lineSeparator();
} catch (IOException e) {
log.warn("could not read previously-effective excludesFile {} while seeding skills into "
+ "{}: {} — its patterns will not carry forward into the seeded worktree",
file, worktreePath, e.getMessage());
return "";
}
}
/**
* Resolve environment variable {@code name} for {@link #previouslyEffectiveExcludesFileContent}'s
* XDG fallback, checking {@link #gitEnv} FIRST so a test can isolate this the same way it
* already isolates every {@code git} subprocess this class runs, and falling back to the JVM's
* real environment only when the seam does not supply it (always the case in production, where
* {@link #gitEnv} is {@code Map.of()}).
*/
private String resolveEnv(String name) {
String fromSeam = gitEnv.get(name);
return fromSeam != null ? fromSeam : System.getenv(name);
}
/** Same as {@link #resolveEnv(String)}, for {@code HOME} — falls back to {@code user.home}
* (rather than {@code System.getenv("HOME")}) when the seam does not supply it, matching this
* class's pre-existing behaviour for every other home-directory resolution. */
private String resolveHome() {
String fromSeam = gitEnv.get("HOME");
return fromSeam != null ? fromSeam : System.getProperty("user.home");
}
/**
* The worker-readable half of fleetd #362, mirroring {@link #recordNeutralizedConfigForWorker}:
* record which skill folders were seeded where the worker itself can read it, without a
* working-tree file that would show up in {@code git status}.
*/
private void recordSeededSkillsForWorker(String worktreePath, List<String> seeded) {
exec("git", "-C", worktreePath, "config", "extensions.worktreeConfig", "true");
for (String name : seeded) {
exec("git", "-C", worktreePath, "config", "--worktree", "--add", "fleet.seededSkills", name);
}
exec("git", "-C", worktreePath, "config", "--worktree", "fleet.seededSkillsNote",
"each fleet.seededSkills value names a skill folder fleetd copied into "
+ ".claude/skills/ because this repo did not already ship it; it is excluded "
+ "from git status (core.excludesFile, worktree-scoped) and must never be committed");
}
@Override
public void remove(String repoRoot, String worktreePath) {
Path p = Path.of(worktreePath);
@@ -1084,6 +1437,9 @@ public final class GitWorktrees implements Worktrees {
Process p;
try {
ProcessBuilder pb = new ProcessBuilder(command).redirectErrorStream(true);
if (!gitEnv.isEmpty()) {
pb.environment().putAll(gitEnv);
}
if (extraEnv != null && !extraEnv.isEmpty()) {
pb.environment().putAll(extraEnv);
}
@@ -1119,7 +1475,11 @@ public final class GitWorktrees implements Worktrees {
private int exitCode(String... command) {
Process p;
try {
p = new ProcessBuilder(command).redirectErrorStream(true).start();
ProcessBuilder pb = new ProcessBuilder(command).redirectErrorStream(true);
if (!gitEnv.isEmpty()) {
pb.environment().putAll(gitEnv);
}
p = pb.start();
} catch (IOException e) {
throw new WorktreeException("failed to start " + command[0] + ": " + e.getMessage(), e);
}
@@ -80,7 +80,7 @@ class FleetdLeadMailboxSelectionTest {
@Test
void opensTheMailboxWhenAUriAndSelfIdAreConfigured() {
var opener = new RecordingOpener();
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null);
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null, null);
Fleetd.openLeadMailbox(coordinator, Map.of(), opener);
@@ -94,7 +94,7 @@ class FleetdLeadMailboxSelectionTest {
void honoursUriEnvOverALiteralUri() {
var opener = new RecordingOpener();
var coordinator = new FleetConfig.Coordinator("amqp://stale:stale@old:5672/x", "COORD_URI",
"mac-opus", 8);
"mac-opus", 8, null);
Fleetd.openLeadMailbox(coordinator, Map.of("COORD_URI", RESOLVED_URI), opener);
@@ -105,7 +105,7 @@ class FleetdLeadMailboxSelectionTest {
@Test
void turnsOffWhenUriEnvDoesNotResolve() {
var opener = new RecordingOpener();
var coordinator = new FleetConfig.Coordinator(null, "COORD_URI", "mac-opus", null);
var coordinator = new FleetConfig.Coordinator(null, "COORD_URI", "mac-opus", null, null);
assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener));
@@ -116,7 +116,7 @@ class FleetdLeadMailboxSelectionTest {
void warnsAndStaysOffWhenSelfIdIsMissing() {
var appender = captureFleetdLogs();
var opener = new RecordingOpener();
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, null, null);
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, null, null, null);
assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener));
@@ -131,7 +131,7 @@ class FleetdLeadMailboxSelectionTest {
var appender = captureFleetdLogs();
var opener = new RecordingOpener();
opener.unreachable = true;
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null);
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null, null);
assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener),
"a down coordination broker turns the feature off; it must never take the daemon down");
@@ -104,9 +104,10 @@ class ConfigRefTopLevelReportingCoverageTest {
v.put("configReload", new FleetConfig.ConfigReload(true, 10));
v.put("quarantineCooldownSeconds", 1800);
v.put("memberCredentials", null);
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-a", null, "self-a", 1));
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-a", null, "self-a", 1, null));
v.put("worktreeGroup", "group-a");
v.put("memberLoginShell", null);
v.put("memberSkills", "/skills/a");
assertNamesMatchComponents(v);
return v;
}
@@ -144,9 +145,10 @@ class ConfigRefTopLevelReportingCoverageTest {
v.put("configReload", new FleetConfig.ConfigReload(false, 20));
v.put("quarantineCooldownSeconds", 3600);
v.put("memberCredentials", null);
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-b", null, "self-b", 2));
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-b", null, "self-b", 2, null));
v.put("worktreeGroup", "group-b");
v.put("memberLoginShell", null);
v.put("memberSkills", "/skills/b");
assertNamesMatchComponents(v);
return v;
}
@@ -1189,7 +1189,7 @@ class FleetConfigTest {
@Test
void coordinatorEffectiveUriHonorsUriEnv() {
FleetConfig.Coordinator withEnv = new FleetConfig.Coordinator(
"amqp://stale-clear-text@127.0.0.1:5672/coord", "LEAD_COORD_URI", "fleet01-lead", null);
"amqp://stale-clear-text@127.0.0.1:5672/coord", "LEAD_COORD_URI", "fleet01-lead", null, null);
assertEquals("amqp://from-env@127.0.0.1:5672/coord",
withEnv.effectiveUri(Map.of("LEAD_COORD_URI", "amqp://from-env@127.0.0.1:5672/coord")),
@@ -1200,12 +1200,61 @@ class FleetConfigTest {
"a blank uriEnv variable must not fall back to the literal uri");
FleetConfig.Coordinator noEnv = new FleetConfig.Coordinator(
"amqp://guest:guest@127.0.0.1:5672/coord", null, null, null);
"amqp://guest:guest@127.0.0.1:5672/coord", null, null, null, null);
assertEquals("amqp://guest:guest@127.0.0.1:5672/coord", noEnv.effectiveUri(Map.of()),
"the literal uri is used when no uriEnv is configured");
assertEquals(LeadMailbox.DEFAULT_PREFETCH, noEnv.prefetchOrDefault());
}
/**
* fleetd #361: the live block ships as just {@code uriEnv} + {@code selfId} (see
* {@code Fleetd.example.yaml} / the operator's real {@code fleetd.yaml}, gitignored). That exact
* shape, with no {@code peers:} key at all, must keep parsing unchanged after this field is added.
*/
@Test
void coordinatorBlockWithNoPeersKeyStillParses(@TempDir Path dir) throws Exception {
Path f = dir.resolve("coordinator-no-peers.yaml");
Files.writeString(f, """
bind:
port: 8080
coordinator:
uriEnv: COORD_AMQP_URI
selfId: mac
""");
FleetConfig cfg = FleetConfig.load(f);
assertNotNull(cfg.coordinator());
assertEquals("mac", cfg.coordinator().selfId());
assertEquals(List.of(), cfg.coordinator().peers(), "no peers: key means no configured peers, never null");
}
@Test
void coordinatorPeersParsesAndDropsBlankEntries(@TempDir Path dir) throws Exception {
Path f = dir.resolve("coordinator-peers.yaml");
Files.writeString(f, """
bind:
port: 8080
coordinator:
selfId: mac
uri: amqp://guest:guest@127.0.0.1:5672/coord
peers:
- fleet01
- ""
- fleet02
""");
FleetConfig cfg = FleetConfig.load(f);
assertEquals(List.of("fleet01", "fleet02"), cfg.coordinator().peers(),
"a blank peer entry must be dropped, never kept as an empty coord-id");
}
@Test
void coordinatorPeersDefaultsToEmptyWhenConstructedWithNull() {
FleetConfig.Coordinator c = new FleetConfig.Coordinator(
"amqp://guest:guest@127.0.0.1:5672/coord", null, "mac", null, null);
assertEquals(List.of(), c.peers(), "a null peers list must default to empty, never NPE downstream");
}
@Test
void absentWorktreeGroupLeavesItNull(@TempDir Path dir) throws Exception {
Path f = dir.resolve("no-worktree-group.yaml");
@@ -45,8 +45,8 @@ import static org.junit.jupiter.api.Assertions.assertEquals;
* <p>Why this is a valid check for every component, not just some: {@link #withDefaults()}'s own
* comments document that it only ever REPLACES a component when the incoming value is {@code null}
* (or blank, for {@code placement}) — {@code broker}/{@code primary}/{@code leadHeartbeat}/
* {@code configReload}/{@code coordinator}/{@code worktreeGroup}/{@code memberLoginShell} are left
* as-is unconditionally, and {@code bind}/{@code guard}/{@code lifecycle}/{@code auth}/
* {@code configReload}/{@code coordinator}/{@code worktreeGroup}/{@code memberLoginShell}/
* {@code memberSkills} are left as-is unconditionally, and {@code bind}/{@code guard}/{@code lifecycle}/{@code auth}/
* {@code fleet}/{@code quarantineCooldownSeconds}/{@code memberCredentials}/{@code placement} are
* replaced only on null/blank input. A value that is never null or blank going in must therefore
* never change coming out, for every current component. No exclusion is needed today.
@@ -92,9 +92,10 @@ class FleetConfigWithDefaultsPreservesEveryComponentTest {
v.put("memberCredentials", new FleetConfig.MemberCredentials(
FleetConfig.MemberCredentials.POLICY_DENY_BY_DEFAULT,
List.of("git"), List.of("git", "ssh"), null));
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-guard", null, "self-guard", 3));
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-guard", null, "self-guard", 3, null));
v.put("worktreeGroup", "group-guard");
v.put("memberLoginShell", "/bin/zsh");
v.put("memberSkills", "/skills/guard");
assertNamesMatchComponents(v);
return v;
}
@@ -6,6 +6,7 @@ import org.junit.jupiter.api.Test;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
@@ -35,6 +36,12 @@ class LeadTabScannerTest {
final Map<String, String[]> tabs = new LinkedHashMap<>();
/** pane_id → [tab_id, terminal_id]. */
final Map<String, String[]> panes = new LinkedHashMap<>();
/**
* tab_id → whether herdr reports a running agent there. Defaults to {@code true} for every
* tab that has a pane, so every existing fixture keeps meaning "a live lead" unless a test
* says otherwise via {@link #deadAgent}.
*/
final Set<String> deadTabs = new LinkedHashSet<>();
int calls;
boolean failing;
@@ -53,6 +60,18 @@ class LeadTabScannerTest {
return this;
}
/** Mark {@code tabId} as labelled but agent-less — a dead lead's leftover tab (fleetd #359). */
TopologyHerdr deadAgent(String tabId) {
deadTabs.add(tabId);
return this;
}
/** Undo {@link #deadAgent} — models {@code agent.list} reporting the tab live again. */
TopologyHerdr reviveAgent(String tabId) {
deadTabs.remove(tabId);
return this;
}
@Override
public JsonNode call(String method, Object params) {
calls++;
@@ -83,6 +102,19 @@ class LeadTabScannerTest {
.formatted(id, p[0], p[1])));
return read("{\"panes\":[%s]}".formatted(String.join(",", items)));
}
case "agent.list" -> {
// One agent per distinct tab that has a pane and isn't marked dead — mirrors
// AgentControl.list()'s "agents" shape closely enough for the scanner's join,
// which only reads tab_id off each entry.
Set<String> seen = new LinkedHashSet<>();
panes.forEach((paneId, p) -> {
String tabId = p[0];
if (!deadTabs.contains(tabId) && seen.add(tabId)) {
items.add("{\"tab_id\":\"%s\"}".formatted(tabId));
}
});
return read("{\"agents\":[%s]}".formatted(String.join(",", items)));
}
default -> throw new AssertionError("unexpected herdr call: " + method);
}
}
@@ -208,6 +240,114 @@ class LeadTabScannerTest {
scanner(herdr, twoLeadsConfigured(), new AtomicLong()).get().get("term_opus_split"));
}
// ── fleetd #359: a labelled tab is only a lead when something is running in it ─────────────────
/**
* The core invariant this ticket restores: {@code get()} must never report a terminal for a
* lead whose pane no longer runs an agent. Before this fix the scanner joined labelled tabs to
* panes with no liveness check at all, so a tab left behind by a crashed/relaunched lead (see
* {@code LeadLauncher}'s own staleness handling) was reported as live forever — which is exactly
* what let duplicate lead tabs make {@code LeadCoordLoop.resolveLocalLead()} permanently unable
* to pick one. Mutate this away (drop the {@code agent.list} cross-check in {@link
* LeadTabScanner#scan()}) and this test must fail.
*/
@Test
void aLabelledTabWithNoRunningAgentIsNotReported() {
TopologyHerdr herdr = twoLeads().deadAgent("w1:t1"); // opus-5.0's tab is labelled but dead
Map<String, String> leads = scanner(herdr, twoLeadsConfigured(), new AtomicLong()).get();
assertFalse(leads.containsKey("term_opus"),
"a labelled tab with no running agent must never be reported as a live lead");
assertEquals("gpt-sol-5.6", leads.get("term_gpt"),
"the other, genuinely live lead must be unaffected");
}
/**
* The other direction, pinned separately so a fix cannot satisfy the test above by simply
* returning nothing: a labelled tab that DOES have a running agent must still be reported. A
* scanner that always comes back empty is worse than the bug it fixes.
*/
@Test
void aLabelledTabWithARunningAgentIsStillReported() {
Map<String, String> leads = scanner(twoLeads(), twoLeadsConfigured(), new AtomicLong()).get();
assertEquals("opus-5.0", leads.get("term_opus"));
assertEquals("gpt-sol-5.6", leads.get("term_gpt"));
}
/**
* fleetd #359 review, finding 2 — the exact scenario the ticket's own evidence showed:
* {@code agent.list} can come back successfully but short, without the herdr call ever throwing.
* A lead this class already reported as live must not be dropped on the strength of one such
* read: {@link LeadTabScanner#get()}'s "keep the cache on failure" contract only fires on an
* exception, so without a fix a single short {@code agent.list} silently empties the cached map —
* which would resolve that lead's pane as {@code Role.WORKER} downstream, refusing every
* orchestration call. Mutate this away (drop the one-scan grace in {@link
* LeadTabScanner#scan()}) and this test must fail.
*/
@Test
void aTransientAgentListMissDoesNotDemoteALeadAlreadyKnownLive() {
TopologyHerdr herdr = twoLeads();
AtomicLong clock = new AtomicLong();
LeadTabScanner s = scanner(herdr, twoLeadsConfigured(), clock);
assertTrue(s.get().containsKey("term_opus"), "opus must be known live before the miss");
// One scan where agent.list comes back without opus's tab, even though the tab and pane are
// completely unchanged — the tab/pane are still there, only the liveness read is short.
herdr.deadAgent("w1:t1");
clock.addAndGet(TTL);
assertTrue(s.get().containsKey("term_opus"),
"a single missed detection must not empty the cached map for a lead already known live");
assertEquals("gpt-sol-5.6", s.get().get("term_gpt"), "the unaffected lead is unchanged");
}
/**
* The other half of finding 2, so the grace above cannot be mistaken for permanent amnesty: the
* original #359 invariant (a genuinely dead tab is not reported forever) must still hold once a
* SECOND, independent scan agrees the agent is gone.
*/
@Test
void aLeadMissingFromAgentListOnTwoConsecutiveScansIsFinallyDropped() {
TopologyHerdr herdr = twoLeads();
AtomicLong clock = new AtomicLong();
LeadTabScanner s = scanner(herdr, twoLeadsConfigured(), clock);
assertTrue(s.get().containsKey("term_opus"));
herdr.deadAgent("w1:t1");
clock.addAndGet(TTL);
assertTrue(s.get().containsKey("term_opus"), "first miss is a grace period, not a verdict");
clock.addAndGet(TTL); // a second, independent scan — still no agent
assertFalse(s.get().containsKey("term_opus"),
"a second consecutive miss for the same terminal must finally drop it");
}
/** A lead that recovers between the two misses keeps its grace spent, not renewed for free. */
@Test
void aLeadThatRecoversBetweenMissesIsReportedNormallyAndResetsItsGrace() {
TopologyHerdr herdr = twoLeads();
AtomicLong clock = new AtomicLong();
LeadTabScanner s = scanner(herdr, twoLeadsConfigured(), clock);
assertTrue(s.get().containsKey("term_opus"));
herdr.deadAgent("w1:t1");
clock.addAndGet(TTL);
assertTrue(s.get().containsKey("term_opus"), "graced on the first miss");
herdr.reviveAgent("w1:t1"); // the miss really was transient
clock.addAndGet(TTL);
assertTrue(s.get().containsKey("term_opus"), "found live again — reported normally");
// A later, unrelated miss must get its own fresh grace scan rather than being dropped
// immediately because the earlier miss had already "used up" a slot for this terminal.
herdr.deadAgent("w1:t1");
clock.addAndGet(TTL);
assertTrue(s.get().containsKey("term_opus"),
"a miss after a genuine recovery is a new event and deserves its own grace scan");
}
/**
* CB-579 acceptance (6): this is the bug the ticket closes. A stale pin used to be merged back
* over every scan and never expire; now a scan is the whole answer, so a lead whose tab is gone
@@ -9,6 +9,7 @@ import org.junit.jupiter.api.Test;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import static org.junit.jupiter.api.Assertions.*;
@@ -106,6 +107,7 @@ class LeadLauncherTest {
assertEquals(0, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
assertFalse(herdr.called("agent.start"), "the live lead must not be duplicated");
assertFalse(herdr.called("tab.close"), "a labelled tab WITH a live agent must never be closed");
}
/**
@@ -122,6 +124,121 @@ class LeadLauncherTest {
"a stale label is not a lead; the lead must be relaunched");
}
/**
* fleetd #359 review finding 1 — the exact scenario the ticket's own live evidence produced: on a
* real host, {@code agent.list} reported "0 live" for a tab a plain {@code ps} confirmed was
* running a real session. A single such reading must never close the tab outright — that would
* destroy the operator's actual lead, a worse failure than the stale-tab bug this ticket exists
* to fix. The first dead reading only flags the tab; mutate this away (make the first reading
* close instead of flag) and this test must fail.
*/
@Test
void aStaleLabelledTabIsFlaggedRatherThanClosedOnTheFirstReconcile() {
FakeHerdr herdr = new FakeHerdr()
.withWorkspace("wL", "fleet")
.withTab("wL", "wL:t1", "lead: opus"); // label only, first look — could be a live session agent.list missed
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
assertFalse(herdr.called("tab.close"),
"one missed reading must never close a tab that might be hosting a live session");
assertTrue(flaggedTabIds(herdr).contains("wL:t1"),
"the stale tab must be flagged pending-close so a later reconcile can confirm it");
}
/**
* The operator's own trace on fleet01 (#359): a daemon that restarted several times, each
* occasion finding "0 live" for whatever reason, had left several identically-labelled dead
* tabs sitting side by side. Every one of them is flagged on its first dead reading, not closed —
* none is more or less trustworthy than another.
*/
@Test
void allStaleLabelledTabsAreFlaggedRatherThanClosedOnTheFirstReconcile() {
FakeHerdr herdr = new FakeHerdr()
.withWorkspace("wL", "fleet")
.withTab("wL", "wL:t1", "lead: opus")
.withTab("wL", "wL:t2", "lead: opus")
.withTab("wL", "wL:t3", "lead: opus"); // three restarts' worth of debris, none confirmed twice yet
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
assertFalse(herdr.called("tab.close"), "no tab may be closed on its first dead reading");
assertEquals(Set.of("wL:t1", "wL:t2", "wL:t3"), flaggedTabIds(herdr),
"every dead labelled tab must be flagged, not just the first one found");
}
/**
* fleetd #359 review finding 1, the other half — a tab already flagged pending-close by an
* earlier reconcile, and STILL dead on this one, has now been read dead on two independent,
* separately-connected reconciles. That is strong enough evidence to actually close it. Mutate
* this away (never close a flagged tab) and this test must fail — the original #359 growth bug
* would come back for good.
*/
@Test
void aTabAlreadyFlaggedPendingCloseIsClosedWhenStillDeadOnALaterReconcile() {
FakeHerdr herdr = new FakeHerdr()
.withWorkspace("wL", "fleet")
.withTab("wL", "wL:t1", "lead: opus [fleetd:pending-close]"); // flagged last reconcile, still dead
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
assertTrue(herdr.called("tab.close"),
"a tab dead on two independent reconciles must finally be closed");
assertEquals("wL:t1", ((Map<?, ?>) herdr.lastCall("tab.close").params()).get("tab_id"));
}
/** All of several already-flagged, still-dead tabs are closed — not just the first found. */
@Test
void allTabsAlreadyFlaggedPendingCloseAreClosedWhenStillDead() {
FakeHerdr herdr = new FakeHerdr()
.withWorkspace("wL", "fleet")
.withTab("wL", "wL:t1", "lead: opus [fleetd:pending-close]")
.withTab("wL", "wL:t2", "lead: opus [fleetd:pending-close]")
.withTab("wL", "wL:t3", "lead: opus [fleetd:pending-close]");
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
List<Object> closedTabIds = herdr.calls.stream()
.filter(c -> c.method().equals("tab.close"))
.<Object>map(c -> ((Map<?, ?>) c.params()).get("tab_id"))
.toList();
assertEquals(3, closedTabIds.size(),
"every confirmed-dead labelled tab must be closed, not just the first one found");
assertEquals(Set.of("wL:t1", "wL:t2", "wL:t3"), Set.copyOf(closedTabIds));
}
/**
* A tab flagged pending-close on a previous reconcile that is running an agent again — the miss
* that flagged it was transient. It must never be closed, and its flag must be cleared so a
* future, unrelated miss starts its own two-reading count from zero.
*/
@Test
void aFlaggedTabRunningAnAgentAgainHasItsFlagClearedInsteadOfBeingClosed() {
FakeHerdr herdr = new FakeHerdr()
.withWorkspace("wL", "fleet")
.withTab("wL", "wL:t1", "lead: opus [fleetd:pending-close]")
.withAgent("lead-opus", "term_lead", "wL:p1", "wL:t1"); // it recovered — really alive now
assertEquals(0, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
"the lead is live again — nothing to relaunch");
assertFalse(herdr.called("tab.close"), "a tab running an agent again must never be closed");
assertFalse(herdr.called("agent.start"), "the lead is live again — nothing to relaunch");
assertEquals("wL:t1", ((Map<?, ?>) herdr.lastCall("tab.rename").params()).get("tab_id"));
assertEquals("lead: opus", ((Map<?, ?>) herdr.lastCall("tab.rename").params()).get("label"),
"the pending-close flag must be cleared once the tab is confirmed live again");
}
/** Every {@code tab.rename} call whose label carries the pending-close marker, by tab id. */
private static Set<String> flaggedTabIds(FakeHerdr herdr) {
return herdr.calls.stream()
.filter(c -> c.method().equals("tab.rename"))
.filter(c -> String.valueOf(((Map<?, ?>) c.params()).get("label"))
.endsWith("[fleetd:pending-close]"))
.map(c -> String.valueOf(((Map<?, ?>) c.params()).get("tab_id")))
.collect(java.util.stream.Collectors.toUnmodifiableSet());
}
/**
* A lead the operator opened by hand is live once its tab carries the configured `tab:` label —
* CB-579 retired the `terminal:` pin, so a hand-opened lead is found the same way an
@@ -1,6 +1,7 @@
package dev.ltms.fleet.mcp;
import dev.ltms.fleet.msg.FakeLeadChannel;
import dev.ltms.fleet.msg.LeadChannel;
import dev.ltms.fleet.msg.LeadMessage;
import io.modelcontextprotocol.spec.McpSchema;
import org.junit.jupiter.api.Test;
@@ -93,4 +94,56 @@ class FleetMcpLeadCoordTest {
assertTrue(FleetMcp.sendToLead(channel, PEER, " ", null, null).isError());
assertEquals(0, channel.published().size());
}
/**
* fleetd #361: "delivered" overstated what publish actually proves — the broker's confirm means
* durably queued, not read. A zero-consumer target is the observable form of "this will not
* reach a pane right now", so the (still-successful) result must say so.
*/
@Test
void warnsWhenThePeerMailboxHasNoConsumersButStillReportsSuccess() {
var channel = new FakeLeadChannel(SELF)
.withMailbox(PEER, LeadChannel.MailboxState.exists(PEER, 0, 0));
McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", null, null);
assertFalse(res.isError(), "a zero-consumer mailbox is still a successful, durably-queued publish");
String out = textOf(res);
assertTrue(out.contains("durably confirmed"), out);
assertTrue(out.toLowerCase().contains("no consumers"), () -> "must warn nobody is reading it: " + out);
}
@Test
void staysQuietAboutConsumersWhenThePeerMailboxHasOne() {
var channel = new FakeLeadChannel(SELF)
.withMailbox(PEER, LeadChannel.MailboxState.exists(PEER, 0, 1));
McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", null, null);
assertFalse(res.isError());
String out = textOf(res);
assertTrue(out.contains("durably confirmed"), out);
assertFalse(out.toLowerCase().contains("no consumers"), () -> "a consumer IS attached: " + out);
}
/**
* fleetd #361 review finding 1: an unmeasured fact must never render as a definite one. When
* the post-publish probe could not determine the mailbox's consumer count at all (broker slow,
* unreachable, or the probe timed out — {@link LeadChannel.MailboxState#unknown}), the result
* must stay just as quiet as the has-a-consumer case — never assert "no consumers" for a mailbox
* this call never actually measured.
*/
@Test
void staysQuietAboutConsumersWhenThePeerMailboxStateIsUnknown() {
var channel = new FakeLeadChannel(SELF)
.withMailbox(PEER, LeadChannel.MailboxState.unknown(PEER));
McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", null, null);
assertFalse(res.isError(), "an unresolved post-publish probe must never turn a durably-confirmed publish into an error");
String out = textOf(res);
assertTrue(out.contains("durably confirmed"), out);
assertFalse(out.toLowerCase().contains("no consumers"),
() -> "an unmeasured fact must never be reported as a definite zero-consumer mailbox: " + out);
}
}
@@ -10,6 +10,9 @@ import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.PaneLocator;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.msg.FakeLeadChannel;
import dev.ltms.fleet.msg.LeadChannel;
import dev.ltms.fleet.msg.LeadMessage;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.session.FakeWorktrees;
@@ -34,7 +37,9 @@ import java.util.Map;
import java.util.EnumSet;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Function;
import static org.junit.jupiter.api.Assertions.*;
@@ -576,17 +581,48 @@ class FleetMcpTest {
void listReportsThisDaemonsOwnCoordIdWhenLeadCoordinationIsOn() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
.withMailbox("mac-opus", LeadChannel.MailboxState.exists("mac-opus", 0, 1));
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), Map.of(), "", "mac-opus");
FleetMcp.QuarantineSource.none(), Map.of(), "",
new FleetMcp.CoordinationSource(channel, List.of()));
String out = textOf(res);
// There is no peer-discovery surface yet, so this row answers the one question an operator
// cannot answer any other way: which coord-id a peer must use to reach ME.
// fleetd #361: reports both which coord-id a peer must use to reach ME, and this daemon's
// own mailbox state (a self-diagnosis: is my own consumer actually attached?).
assertTrue(out.contains("\"coordinator\""), out);
assertTrue(out.contains("\"selfId\":\"mac-opus\""), out);
assertTrue(out.contains("\"mailbox\":{\"status\":\"exists\",\"pending\":0,\"consumers\":1}"), out);
assertTrue(out.contains("\"held\":[]"), out);
assertTrue(out.contains("\"peers\":[]"), out);
}
/**
* fleetd #361 review finding 1: a self-probe that could not complete (broker unreachable, timed
* out) must never render the same as a measured "0 pending, 0 consumers" — that was exactly the
* bug: a reader could not tell "my mailbox is empty and idle" from "I could not check", and the
* second one is the far more alarming state.
*/
@Test
void listReportsAnUnresolvedSelfProbeAsUnknownNeverAsAMeasuredZero() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
.withMailbox("mac-opus", LeadChannel.MailboxState.unknown("mac-opus"));
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), Map.of(), "",
new FleetMcp.CoordinationSource(channel, List.of()));
String out = textOf(res);
assertTrue(out.contains("\"mailbox\":{\"status\":\"unknown\"}"), out);
assertFalse(out.contains("\"pending\""), "an unresolved probe must never carry a pending count at all: " + out);
assertFalse(out.contains("\"consumers\""), "an unresolved probe must never carry a consumers count at all: " + out);
}
@Test
@@ -601,6 +637,110 @@ class FleetMcpTest {
"an ordinary fleet's output must be unchanged by this feature");
}
@Test
void listReportsHeldMessagesWithATruncatedPreviewNeverTheFullBody() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
String longContent = "x".repeat(200);
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
.hold(new LeadMessage("m1", "fleet01-lead", "mac-opus", longContent));
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), Map.of(), "",
new FleetMcp.CoordinationSource(channel, List.of()));
String out = textOf(res);
assertTrue(out.contains("\"msgId\":\"m1\""), out);
assertTrue(out.contains("\"from\":\"fleet01-lead\""), out);
assertFalse(out.contains(longContent), "fleet_list must never dump a held message's full body: " + out);
assertTrue(out.contains("x".repeat(80) + "…"), "expected an 80-char preview with an ellipsis: " + out);
}
@Test
void listReportsEachDeclaredPeersLiveReachability() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
.withMailbox("fleet01-lead", LeadChannel.MailboxState.exists("fleet01-lead", 2, 1))
.withMailbox("fleet03-lead", LeadChannel.MailboxState.unknown("fleet03-lead"));
// "fleet02-lead" is declared as a peer but never configured on the fake — inspect() falls
// back to MailboxState.absent, exactly as a real down (never-run) peer would report.
// "fleet03-lead" IS configured, as unknown — a broker that could not be reached in time,
// which review finding 1 says must render distinctly from "fleet02-lead"'s confirmed absence.
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), Map.of(), "",
new FleetMcp.CoordinationSource(channel, List.of("fleet01-lead", "fleet02-lead", "fleet03-lead")));
String out = textOf(res);
assertTrue(out.contains("\"coordId\":\"fleet01-lead\",\"status\":\"exists\",\"pending\":2,\"consumers\":1"), out);
assertTrue(out.contains("\"coordId\":\"fleet02-lead\",\"status\":\"absent\""), out);
assertTrue(out.contains("\"coordId\":\"fleet03-lead\",\"status\":\"unknown\""), out);
assertFalse(out.contains("\"coordId\":\"fleet02-lead\",\"status\":\"absent\",\"pending\""),
"pending/consumers must be omitted, not faked as zero, for a confirmed-absent peer: " + out);
assertFalse(out.contains("\"coordId\":\"fleet03-lead\",\"status\":\"unknown\",\"pending\""),
"pending/consumers must be omitted, not faked as zero, for an unresolved peer probe: " + out);
}
/**
* fleetd #361 review finding 2: {@code get(timeout)} alone times out the CALLER but leaves the
* submitted {@link LeadChannel#inspect} task running forever on its own virtual thread — against
* a hung (not down) broker every probe would orphan one more thread holding an AMQP channel
* until the connection's channel-max is exhausted, which would break {@code publish} too. This
* proves {@link FleetMcp#probe(LeadChannel, String, long)} does not merely give up on a slow
* task: it interrupts it, so the task does not go on running unbounded after the caller has
* already moved on. No hung broker needed — a {@link LeadChannel} fake that blocks until
* interrupted is enough to observe the same mechanism.
*/
@Test
void aTimedOutProbeInterruptsTheOrphanedTaskRatherThanAbandoningIt() throws Exception {
CountDownLatch started = new CountDownLatch(1);
AtomicBoolean wasInterrupted = new AtomicBoolean(false);
LeadChannel hangs = new LeadChannel() {
@Override
public void publish(String toCoordId, LeadMessage m) { }
@Override
public List<LeadMessage> peek() { return List.of(); }
@Override
public void ack(String msgId) { }
@Override
public String selfCoordId() { return "mac-opus"; }
@Override
public MailboxState inspect(String coordId) {
started.countDown();
try {
Thread.sleep(60_000);
} catch (InterruptedException e) {
wasInterrupted.set(true);
Thread.currentThread().interrupt();
}
return MailboxState.unknown(coordId);
}
};
LeadChannel.MailboxState result = FleetMcp.probe(hangs, "fleet01-lead", 100L);
assertFalse(result.exists(), "a timed-out probe must never claim the mailbox exists");
assertFalse(result.known(), "a timed-out probe proves nothing either way — it must report unknown");
assertTrue(started.await(2, TimeUnit.SECONDS), "the probe task must actually have started");
// The interrupt is delivered asynchronously to the orphaned task's own thread — poll briefly
// rather than assume it has already landed the instant probe() returns.
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(2);
while (!wasInterrupted.get() && System.nanoTime() < deadline) {
Thread.sleep(20);
}
assertTrue(wasInterrupted.get(),
"probe() must cancel the orphaned task (interrupt it) instead of leaving it to run forever");
}
@Test
void capacityUsesThePlacementLiveCount() {
FakeHerdr h = new FakeHerdr();
@@ -911,7 +1051,8 @@ class FleetMcpTest {
String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
sessions, null, new FleetMcp.CapacitySource(profile -> 2, profile -> 3,
() -> Set.of("sonnet"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", null));
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "",
FleetMcp.CoordinationSource.none()));
assertTrue(out.contains("\"maxLoad\":3"), "maxLoad itself must be left untouched: " + out);
assertTrue(out.contains("\"live\":2"), out);
@@ -934,7 +1075,8 @@ class FleetMcpTest {
String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
sessions, null, new FleetMcp.CapacitySource(profile -> 0, profile -> 3,
() -> Set.of("sonnet"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", null));
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "",
FleetMcp.CoordinationSource.none()));
assertTrue(out.contains("\"live\":0"), out);
assertTrue(out.contains("\"free\":3"), "the real gate never subtracts the lead's seat: " + out);
@@ -988,7 +1130,7 @@ class FleetMcpTest {
String out = textOf(FleetMcp.listFleet(composite, sm, null,
new FleetMcp.CapacitySource(liveCount, p -> profiles.get(p).maxLoad(), profiles::keySet, () -> 0),
new FleetMcp.HealthCoverageSource(() -> "off"), FleetMcp.QuarantineSource.none(),
FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", null));
FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", FleetMcp.CoordinationSource.none()));
int reportedFree = extractInt(out, "free");
for (int i = 0; i < reportedFree; i++) {
@@ -1017,7 +1159,7 @@ class FleetMcpTest {
sessions, null, new FleetMcp.CapacitySource(profile -> 0, profile -> 2,
() -> Set.of("terra"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", null));
Map.of(), "", FleetMcp.CoordinationSource.none()));
assertTrue(out.contains("\"free\":2"), out);
assertFalse(out.contains("leadSeats"), "no lead shares this profile's credential: " + out);
@@ -146,7 +146,7 @@ class MemberEnvAllowListTest {
return new FleetConfig(null, null, null, Map.of(), null, null, null, null, null,
new FleetConfig.Broker(null, brokerUriEnv, null), null, null, null, null, null,
null, null, null, null,
new FleetConfig.Coordinator(null, coordinatorUriEnv, null, null)).withDefaults();
new FleetConfig.Coordinator(null, coordinatorUriEnv, null, null, null)).withDefaults();
}
/** {@code LC_*} categories are infrastructure by prefix; everything else needs an exact match. */
@@ -3,6 +3,8 @@ package dev.ltms.fleet.msg;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
/**
* Hermetic stand-in for {@link LeadChannel}: an in-memory mailbox that records what was published
@@ -24,11 +26,19 @@ public final class FakeLeadChannel implements LeadChannel {
private final List<String> acked = Collections.synchronizedList(new ArrayList<>());
/** When set, every {@link #publish} throws it — the unroutable/nacked/timed-out peer. */
private volatile IllegalStateException publishFailure;
/** Canned {@link #inspect} results by coord-id — absent for any coord-id not configured here. */
private final Map<String, MailboxState> mailboxes = new ConcurrentHashMap<>();
public FakeLeadChannel(String selfCoordId) {
this.selfCoordId = selfCoordId;
}
/** Make {@link #inspect(String)} return {@code state} for {@code coordId} instead of "absent". */
public FakeLeadChannel withMailbox(String coordId, MailboxState state) {
mailboxes.put(coordId, state);
return this;
}
/** Make every publish fail as an unreachable peer would. */
public FakeLeadChannel failPublishWith(String message) {
this.publishFailure = new IllegalStateException(message);
@@ -65,6 +75,11 @@ public final class FakeLeadChannel implements LeadChannel {
return selfCoordId;
}
@Override
public MailboxState inspect(String coordId) {
return mailboxes.getOrDefault(coordId, MailboxState.absent(coordId));
}
public List<LeadMessage> published() {
return List.copyOf(published);
}
@@ -0,0 +1,77 @@
package dev.ltms.fleet.msg;
import com.rabbitmq.client.ShutdownSignalException;
import com.rabbitmq.client.impl.AMQImpl;
import org.junit.jupiter.api.Test;
import java.io.IOException;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #361 review round 2: a mutation that made {@link LeadMailbox#isMissingQueue} return
* {@code true} unconditionally still left {@code mvn clean install} green — 1389 tests, 0
* failures — because nothing exercised its false branch. That branch is the whole discriminator
* between {@link LeadChannel.MailboxState#absent} and {@link LeadChannel.MailboxState#unknown};
* without a test pinning it, a future refactor that widens it back to "always true" (restoring the
* exact overstatement fleetd #361 exists to fix) would pass this suite.
*
* <p>Hermetic — no broker needed, per the review's own suggestion. {@code isMissingQueue} takes a
* plain {@link IOException}, so every input here is constructed directly rather than provoked from
* a live connection. The real 404 shape itself is still pinned against a real broker, in
* {@code LeadMailboxTest.passiveDeclareOfAMissingQueueThrowsAnIOExceptionWrappingA404ShutdownSignal}
* — this class covers the three false shapes {@link LeadMailbox#isMissingQueue}'s own javadoc
* lists, so both directions of the discriminator are proven somewhere.
*/
class LeadMailboxIsMissingQueueTest {
@Test
void aConfirmedMissingQueueIsRecognized() {
ShutdownSignalException sse = new ShutdownSignalException(true, false,
new AMQImpl.Channel.Close(404, "NOT_FOUND - no queue 'lead.x.inbox' in vhost '/'", 50, 10), null);
IOException e = new IOException("channel error", sse);
assertTrue(LeadMailbox.isMissingQueue(e), "a genuine 404 Channel.Close must be recognized as a missing queue");
}
@Test
void aDifferentReplyCodeIsNotAMissingQueue() {
// E.g. 403 ACCESS_REFUSED — the queue may well exist; this call was simply refused.
ShutdownSignalException sse = new ShutdownSignalException(true, false,
new AMQImpl.Channel.Close(403, "ACCESS_REFUSED", 50, 10), null);
IOException e = new IOException("channel error", sse);
assertFalse(LeadMailbox.isMissingQueue(e),
"a non-404 reply code must never be read as a confirmed absence — the mailbox's real state is unknown");
}
@Test
void aShutdownSignalWhoseReasonIsNotAChannelCloseIsNotAMissingQueue() {
// A Connection.Close (a whole different broker-level shutdown) is still a ShutdownSignalException,
// but its reason is not a Channel.Close at all — must not be misread as "no such queue".
ShutdownSignalException sse = new ShutdownSignalException(true, false,
new AMQImpl.Connection.Close(404, "coincidentally 404, but this is a CONNECTION close", 10, 50), null);
IOException e = new IOException("connection error", sse);
assertFalse(LeadMailbox.isMissingQueue(e),
"a ShutdownSignalException whose reason is not a Channel.Close must never be read as a missing queue,"
+ " even if its reply code happens to be 404");
}
@Test
void anIOExceptionWithNoCauseAtAllIsNotAMissingQueue() {
IOException e = new IOException("some other declare failure, no cause attached");
assertFalse(LeadMailbox.isMissingQueue(e),
"an IOException with no ShutdownSignalException cause must never be read as a confirmed absence");
}
@Test
void anIOExceptionWithAnUnrelatedCauseIsNotAMissingQueue() {
IOException e = new IOException("wrapped something else entirely", new RuntimeException("boom"));
assertFalse(LeadMailbox.isMissingQueue(e),
"a cause that isn't even a ShutdownSignalException must never be read as a confirmed absence");
}
}
@@ -1,5 +1,9 @@
package dev.ltms.fleet.msg;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ShutdownSignalException;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Tag;
import org.junit.jupiter.api.Test;
@@ -7,11 +11,14 @@ import org.testcontainers.containers.RabbitMQContainer;
import org.testcontainers.junit.jupiter.Testcontainers;
import org.testcontainers.utility.DockerImageName;
import java.io.IOException;
import java.util.List;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -159,6 +166,145 @@ class LeadMailboxTest {
}
}
@Test
void inspectReportsAnOwnedMailboxAsExistingWithItsOwnConsumer() throws Exception {
// A LeadMailbox declares AND consumes its own queue the moment open() returns (see own()),
// so inspecting a coord-id this same process owns must always find exactly one consumer.
String self = coordId("lead-inspect-self");
try (LeadMailbox mailbox = LeadMailbox.open(uri(), self)) {
LeadChannel.MailboxState state = mailbox.inspect(self);
assertEquals(self, state.coordId());
assertTrue(state.exists(), "this daemon owns and has declared this exact queue");
assertEquals(0, state.pending(), "nothing has been published to it yet");
assertEquals(1, state.consumers(), "the mailbox's own constructor already attached a consumer");
}
}
@Test
void inspectReportsAMissingMailboxAsAbsentRatherThanThrowing() throws Exception {
String nobody = coordId("lead-inspect-nobody");
try (LeadMailbox mailbox = LeadMailbox.open(uri(), coordId("lead-inspect-caller"))) {
LeadChannel.MailboxState state = mailbox.inspect(nobody);
assertEquals(LeadChannel.MailboxState.absent(nobody), state,
"a queue nobody has ever declared must report absent, never throw");
assertTrue(state.known(), "a confirmed 404 IS a definite answer — this is not the unknown case");
}
}
/**
* fleetd #361 review finding 3: {@code inspect} is specified to never throw, but the original
* implementation caught only {@link IOException} — and {@link Connection#createChannel()} on an
* already-closed connection throws {@link com.rabbitmq.client.AlreadyClosedException}, an
* unchecked {@link RuntimeException} (pinned by {@code
* createChannelOnAnAlreadyClosedConnectionThrowsAnUncheckedException} above). This drives that
* exact scenario through the real {@link LeadMailbox#inspect} — not the raw client call — and
* checks both halves of finding 1 and finding 3 at once: no exception escapes, and the result is
* {@code UNKNOWN} rather than the wrong-but-plausible-looking {@code ABSENT}.
*/
@Test
void inspectReportsUnknownRatherThanThrowingWhenTheConnectionIsAlreadyClosed() throws Exception {
LeadMailbox mailbox = LeadMailbox.open(uri(), coordId("lead-inspect-dead-connection"));
mailbox.close(); // tears down the connection `inspect` will try to open a probe channel on
LeadChannel.MailboxState state = mailbox.inspect(coordId("lead-inspect-irrelevant-target"));
assertFalse(state.exists());
assertFalse(state.known(), "a dead connection proves nothing about the target mailbox — it must be unknown, not absent");
assertEquals(LeadChannel.MailboxState.Presence.UNKNOWN, state.presence());
}
@Test
void inspectReportsPendingMessagesAndZeroConsumersWhenNobodyIsReadingAnymore() throws Exception {
// Publish into a mailbox this test owns, then never consume from it, to prove `pending` and
// `consumers` really come off the broker rather than off this process's own in-memory state.
String to = coordId("lead-inspect-pending");
String observerId = coordId("lead-inspect-observer");
try (LeadMailbox owner = LeadMailbox.open(uri(), to);
LeadMailbox observer = LeadMailbox.open(uri(), observerId)) {
owner.publish(to, new LeadMessage("m1", "lead-from", to, "sitting in the queue"));
awaitPeek(owner); // make sure the broker has actually enqueued it before inspecting
} // `owner` closes here: its consumer disconnects, but the durable, unacked message stays queued.
try (LeadMailbox observer = LeadMailbox.open(uri(), coordId("lead-inspect-observer-2"))) {
// The broker requeues `owner`'s unacked delivery asynchronously once its connection drops,
// so poll rather than assume the very first passive declare already sees the settled state.
LeadChannel.MailboxState state = awaitInspect(observer, to, s -> s.consumers() == 0);
assertTrue(state.exists());
assertEquals(0, state.consumers(), "the only owner just closed — nobody is reading this anymore");
assertEquals(1, state.pending(), "the unacked message must be requeued, never dropped");
}
}
/**
* fleetd #361's central invariant, proved rather than assumed: a passive queue declare of a
* missing queue closes ITS channel with a 404 in AMQP 0-9-1. {@link LeadMailbox#inspect} is
* specified to run on its own disposable channel for exactly this reason — this test is the one
* that actually exercises the failure mode and shows {@link LeadMailbox#publish} on the SAME
* instance is unaffected by it.
*/
@Test
void inspectingAMissingMailboxNeverBreaksPublishOnTheSameInstance() throws Exception {
String self = coordId("lead-invariant-self");
try (LeadMailbox mailbox = LeadMailbox.open(uri(), self)) {
// Miss on a queue that has never existed — this is exactly the 404-closes-the-channel case.
LeadChannel.MailboxState missed = mailbox.inspect(coordId("lead-invariant-nobody-home"));
assertFalse(missed.exists());
assertTrue(missed.known(), "a genuine 404 on a queue that never existed is a confirmed fact, not an unknown");
// publish() must still work on THIS SAME instance: if inspect() had reused `publishChannel`
// (or `channel`), the broker's 404 would have closed it out from underneath publish().
LeadMessage sent = new LeadMessage("after-miss", "lead-from", self, "still alive");
mailbox.publish(self, sent);
List<LeadMessage> got = awaitPeek(mailbox);
assertEquals(1, got.size(), "publish must still reach this mailbox's own queue after a missed inspect");
assertEquals("after-miss", got.getFirst().msgId());
// And a second inspect() — of a mailbox that DOES exist this time — must also still work,
// proving the miss did not wedge inspect() itself either.
LeadChannel.MailboxState self2 = mailbox.inspect(self);
assertTrue(self2.exists());
}
}
/**
* Pins the exact exception shape {@link LeadMailbox#inspect} relies on to tell a genuine 404
* (mailbox confirmed absent) apart from everything else (mailbox state unknown) — measured
* against a real broker rather than assumed from the AMQP 0-9-1 spec text. If this ever fails,
* the classification in {@code inspect} is reading the wrong shape and must be revisited.
*/
@Test
void passiveDeclareOfAMissingQueueThrowsAnIOExceptionWrappingA404ShutdownSignal() throws Exception {
try (Connection conn = LeadMailbox.connectionFactory(uri()).newConnection()) {
Channel probe = conn.createChannel();
String missing = LeadMailbox.queueName(coordId("lead-404-shape"));
IOException thrown = assertThrows(IOException.class, () -> probe.queueDeclarePassive(missing));
assertInstanceOf(ShutdownSignalException.class, thrown.getCause(),
() -> "expected the IOException to wrap a ShutdownSignalException, got: " + thrown);
ShutdownSignalException sse = (ShutdownSignalException) thrown.getCause();
assertInstanceOf(AMQP.Channel.Close.class, sse.getReason(),
() -> "expected a Channel.Close reason: " + sse);
AMQP.Channel.Close close = (AMQP.Channel.Close) sse.getReason();
assertEquals(404, close.getReplyCode(), () -> "expected AMQP NOT_FOUND (404): " + close);
assertFalse(probe.isOpen(), "the 404 must have closed the channel the declare ran on");
}
}
/**
* The other half of the same measurement: calling {@code createChannel()} on an
* already-closed connection — the shape {@link LeadMailbox#inspect} hits when the broker
* connection itself is gone — throws {@link com.rabbitmq.client.AlreadyClosedException}, an
* unchecked {@link RuntimeException}, not an {@link IOException}. An {@code inspect} that only
* caught {@code IOException} here would let this escape instead of reporting "unknown".
*/
@Test
void createChannelOnAnAlreadyClosedConnectionThrowsAnUncheckedException() throws Exception {
Connection conn = LeadMailbox.connectionFactory(uri()).newConnection();
conn.close();
RuntimeException thrown = assertThrows(RuntimeException.class, conn::createChannel);
assertInstanceOf(com.rabbitmq.client.AlreadyClosedException.class, thrown,
() -> "expected AlreadyClosedException, got: " + thrown);
}
/** Poll peek until at least one message is held, or ~10s elapse (broker delivery is async). */
@SuppressWarnings("BusyWait")
private static List<LeadMessage> awaitPeek(LeadMailbox inbox) throws InterruptedException {
@@ -170,4 +316,18 @@ class LeadMailboxTest {
}
return msgs;
}
/** Poll inspect(coordId) until it satisfies {@code done}, or ~10s elapse (broker state settles async). */
@SuppressWarnings("BusyWait")
private static LeadChannel.MailboxState awaitInspect(
LeadMailbox observer, String coordId, java.util.function.Predicate<LeadChannel.MailboxState> done)
throws InterruptedException {
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
LeadChannel.MailboxState state = observer.inspect(coordId);
while (!done.test(state) && System.nanoTime() < deadline) {
Thread.sleep(50);
state = observer.inspect(coordId);
}
return state;
}
}
@@ -12,6 +12,7 @@ import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
@@ -53,6 +54,56 @@ class GitWorktreesTest {
/** A non-empty autoenv file — the form that would prompt for authorization in a worktree. */
private static final String AUTOENV_WITH_DIRECTIVE = "export HELLO=world\n";
/**
* fleetd #369. A throwaway directory that lives for the whole class (JUnit 5.4+ supports a
* static {@code @TempDir} field, created once and removed once every test in this class has
* run) — backing every raw {@code git} subprocess's {@code XDG_CONFIG_HOME} below. It only
* ever needs to exist and be guaranteed free of a {@code git/ignore} file; nothing writes
* inside it.
*/
@TempDir
private static Path CLASS_TMP;
/**
* fleetd #369 — the leak measured: {@code XDG_CONFIG_HOME=<dir with a `*` git/ignore> mvn test
* -Dtest=GitWorktreesTest} failed 56 of 59 tests on an unpatched checkout, because {@link
* #gitOutput} set {@code GIT_CONFIG_GLOBAL}/{@code GIT_CONFIG_SYSTEM}/{@code
* GIT_TERMINAL_PROMPT} but not {@code XDG_CONFIG_HOME}, and {@link #status}/{@link
* #fullStatus} (plus every other raw {@code git} subprocess this class started) set NOTHING at
* all — inheriting the JVM's whole real environment, including the operator's real {@code
* ~/.gitconfig} and real default excludes file ({@code $XDG_CONFIG_HOME/git/ignore} or {@code
* $HOME/.config/git/ignore}, applied by git with no {@code core.excludesFile} configured at
* all — see {@code gitignore(5)}). {@code GIT_CONFIG_GLOBAL=/dev/null} does not stop that
* default from applying; only setting {@code XDG_CONFIG_HOME} to a directory that provably
* carries no {@code git/ignore} does.
*
* <p>This is the same isolation {@link #hermeticGitEnv} already gives {@link
* #seedingGitWorktrees}'s production {@link GitWorktrees} instances (fleetd #362 review fix,
* finding 2), reused here for every subprocess the TEST ITSELF starts to drive and inspect
* those fixture repos.
*/
private static Map<String, String> hermeticEnv() {
return hermeticGitEnv(CLASS_TMP);
}
/**
* The one seam every git subprocess in this class is built through — see criterion 4's
* self-check, {@link #everyGitSubprocessGoesThroughTheHermeticFactory}, which fails the moment
* a future helper builds its own {@code git} subprocess directly instead of calling this, so
* the omission that caused fleetd #369 gets caught by name rather than rediscovered by a
* poisoned machine. The one deliberate exception is {@link
* #worktreeCredentialHelperCompletesWithoutUsingAnInheritedHelper}, which needs a
* non-hermetic, test-controlled global config to prove the credential helper ignores it — see
* the comment on that test.
*/
private static ProcessBuilder gitProcessBuilder(Path cwd, String... args) {
List<String> cmd = new java.util.ArrayList<>(List.of("git"));
cmd.addAll(List.of(args));
ProcessBuilder pb = new ProcessBuilder(cmd).directory(cwd.toFile()).redirectErrorStream(true);
pb.environment().putAll(hermeticEnv());
return pb;
}
private static Path initRepo(Path dir) throws Exception {
Files.createDirectories(dir);
git(dir, "init", "-q", "-b", "main");
@@ -70,23 +121,16 @@ class GitWorktreesTest {
}
private static String gitOutput(Path cwd, String... args) throws Exception {
List<String> cmd = new java.util.ArrayList<>(List.of("git"));
cmd.addAll(List.of(args));
ProcessBuilder pb = new ProcessBuilder(cmd).directory(cwd.toFile()).redirectErrorStream(true);
pb.environment().put("GIT_CONFIG_GLOBAL", "/dev/null");
pb.environment().put("GIT_CONFIG_SYSTEM", "/dev/null");
pb.environment().put("GIT_TERMINAL_PROMPT", "0");
Process p = pb.start();
Process p = gitProcessBuilder(cwd, args).start();
String out = new String(p.getInputStream().readAllBytes());
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git timed out: " + String.join(" ", cmd));
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git timed out: git " + String.join(" ", args));
assertEquals(0, p.exitValue(), "git " + String.join(" ", args) + " failed:\n" + out);
return out;
}
/** Pending changes to {@code file} in {@code cwd}, empty when git considers it unmodified. */
private static String status(Path cwd, String file) throws Exception {
Process p = new ProcessBuilder("git", "status", "--porcelain", "--", file)
.directory(cwd.toFile()).redirectErrorStream(true).start();
Process p = gitProcessBuilder(cwd, "status", "--porcelain", "--", file).start();
String out = new String(p.getInputStream().readAllBytes());
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git status timed out");
return out;
@@ -94,16 +138,14 @@ class GitWorktreesTest {
/** Every pending change in {@code cwd} — the whole-tree porcelain status, unlike {@link #status}. */
private static String fullStatus(Path cwd) throws Exception {
Process p = new ProcessBuilder("git", "status", "--porcelain")
.directory(cwd.toFile()).redirectErrorStream(true).start();
Process p = gitProcessBuilder(cwd, "status", "--porcelain").start();
String out = new String(p.getInputStream().readAllBytes());
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git status timed out");
return out;
}
private static String revParse(Path cwd, String ref) throws Exception {
Process p = new ProcessBuilder("git", "-C", cwd.toString(), "rev-parse", ref)
.redirectErrorStream(true).start();
Process p = gitProcessBuilder(cwd, "rev-parse", ref).start();
String out = new String(p.getInputStream().readAllBytes()).trim();
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git rev-parse timed out");
assertEquals(0, p.exitValue(), "git rev-parse " + ref + " failed:\n" + out);
@@ -112,8 +154,7 @@ class GitWorktreesTest {
/** The recursive file list of a commit's tree — used to check what a snapshot actually committed. */
private static String lsTree(Path cwd, String ref) throws Exception {
Process p = new ProcessBuilder("git", "-C", cwd.toString(), "ls-tree", "-r", "--name-only", ref)
.redirectErrorStream(true).start();
Process p = gitProcessBuilder(cwd, "ls-tree", "-r", "--name-only", ref).start();
String out = new String(p.getInputStream().readAllBytes());
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git ls-tree timed out");
assertEquals(0, p.exitValue(), "git ls-tree " + ref + " failed:\n" + out);
@@ -124,8 +165,7 @@ class GitWorktreesTest {
* snapshot's tree changed relative to its parent, the same shape {@code git status --porcelain}
* reports for the worktree it was taken from. */
private static Set<String> diffNameOnly(Path cwd, String from, String to) throws Exception {
Process p = new ProcessBuilder("git", "-C", cwd.toString(), "diff", "--name-only", from, to)
.redirectErrorStream(true).start();
Process p = gitProcessBuilder(cwd, "diff", "--name-only", from, to).start();
String out = new String(p.getInputStream().readAllBytes());
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git diff timed out");
assertEquals(0, p.exitValue(), "git diff " + from + ".." + to + " failed:\n" + out);
@@ -151,8 +191,7 @@ class GitWorktreesTest {
}
private static String forEachRef(Path cwd, String pattern) throws Exception {
Process p = new ProcessBuilder("git", "-C", cwd.toString(), "for-each-ref", pattern)
.redirectErrorStream(true).start();
Process p = gitProcessBuilder(cwd, "for-each-ref", pattern).start();
String out = new String(p.getInputStream().readAllBytes());
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git for-each-ref timed out");
assertEquals(0, p.exitValue(), "git for-each-ref " + pattern + " failed:\n" + out);
@@ -161,8 +200,7 @@ class GitWorktreesTest {
/** Write {@code content} as a blob into the object database; returns its sha. */
private static String blobOf(Path cwd, String content) throws Exception {
Process p = new ProcessBuilder("git", "-C", cwd.toString(), "hash-object", "-w", "--stdin")
.redirectErrorStream(true).start();
Process p = gitProcessBuilder(cwd, "hash-object", "-w", "--stdin").start();
p.getOutputStream().write(content.getBytes(StandardCharsets.UTF_8));
p.getOutputStream().close();
String out = new String(p.getInputStream().readAllBytes()).trim();
@@ -173,8 +211,7 @@ class GitWorktreesTest {
/** Build a single-file tree object from {@code blob}; returns the tree's sha. */
private static String treeOf(Path cwd, String path, String blob) throws Exception {
Process p = new ProcessBuilder("git", "-C", cwd.toString(), "mktree")
.redirectErrorStream(true).start();
Process p = gitProcessBuilder(cwd, "mktree").start();
p.getOutputStream().write(("100644 blob " + blob + "\t" + path + "\n").getBytes(StandardCharsets.UTF_8));
p.getOutputStream().close();
String out = new String(p.getInputStream().readAllBytes()).trim();
@@ -186,10 +223,9 @@ class GitWorktreesTest {
/** {@code git commit-tree} rooted at {@code tree} with a chosen committer date; returns the sha. */
private static String commitTree(Path cwd, String tree, String parent, String committerDate,
String message) throws Exception {
ProcessBuilder pb = new ProcessBuilder("git", "-C", cwd.toString(), "commit-tree",
tree, "-p", parent, "-m", message);
ProcessBuilder pb = gitProcessBuilder(cwd, "commit-tree", tree, "-p", parent, "-m", message);
pb.environment().put("GIT_COMMITTER_DATE", committerDate);
Process p = pb.redirectErrorStream(true).start();
Process p = pb.start();
String out = new String(p.getInputStream().readAllBytes()).trim();
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git commit-tree timed out");
assertEquals(0, p.exitValue(), "git commit-tree failed:\n" + out);
@@ -262,6 +298,11 @@ class GitWorktreesTest {
helper = !f() { printf 'username=%s\\npassword=%s\\n\\n' operator operator-secret; }; f
""");
// fleetd #369: the one deliberate exception to gitProcessBuilder. This test's whole point is
// that git must resolve `globalConfig` (a synthetic "operator's global config", never the
// real machine's) and then IGNORE it — so it cannot use the shared hermetic env, which would
// point GIT_CONFIG_GLOBAL at /dev/null and defeat the very thing under test. It never runs
// `git status`, so it does not need XDG_CONFIG_HOME isolation either.
ProcessBuilder pb = new ProcessBuilder("git", "credential", "fill")
.directory(Path.of(wt).toFile()).redirectErrorStream(true);
pb.environment().put("GIT_CONFIG_GLOBAL", globalConfig.toString());
@@ -328,20 +369,16 @@ class GitWorktreesTest {
Path worktree = Path.of(wt);
assertEquals("https://git.ltms.dev/akb/kb.git",
gitOutput(worktree, "remote", "get-url", "origin").trim());
assertEquals(1, exitCode("git", "-C", wt, "config", "--worktree", "--get-regexp", "^url\\."),
assertEquals(1, gitExitCode(worktree, "config", "--worktree", "--get-regexp", "^url\\."),
"no url.*.insteadOf rewrite should be added for an already-HTTPS origin");
}
/** Test-local exit-code probe, mirroring {@link GitWorktrees#exitCode} for an assertion the
* production class does not expose. */
private static int exitCode(String... command) throws Exception {
ProcessBuilder pb = new ProcessBuilder(command).redirectErrorStream(true);
pb.environment().put("GIT_CONFIG_GLOBAL", "/dev/null");
pb.environment().put("GIT_CONFIG_SYSTEM", "/dev/null");
pb.environment().put("GIT_TERMINAL_PROMPT", "0");
Process p = pb.start();
private static int gitExitCode(Path cwd, String... args) throws Exception {
Process p = gitProcessBuilder(cwd, args).start();
p.getInputStream().readAllBytes();
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "command timed out: " + String.join(" ", command));
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "command timed out: git " + String.join(" ", args));
return p.exitValue();
}
@@ -1469,4 +1506,346 @@ class GitWorktreesTest {
assertTrue(reportingAppender.list.isEmpty(),
"a null/empty overlay must log nothing, got:\n" + capturedMessages());
}
// ---- fleetd #362: seedSkills. Drives GitWorktrees#add end-to-end (not a bare worktree) so the
// real memberSkillsSource wiring is exercised, exactly like the credential-helper/origin tests
// above do for their own seams. ----
/** Write {@code content} as {@code <dir>/<skillName>/SKILL.md}, creating {@code dir} first. */
private static void writeSkill(Path dir, String skillName, String content) throws IOException {
Path skillFile = dir.resolve(skillName).resolve("SKILL.md");
Files.createDirectories(skillFile.getParent());
Files.writeString(skillFile, content);
}
/**
* fleetd #362 review fix, finding 2. {@code core.excludesFile}'s XDG-fallback branch
* ({@link GitWorktrees#previouslyEffectiveExcludesFileContent}) does not go through a {@code
* git} subprocess, so a first cut of it read {@code XDG_CONFIG_HOME}/{@code HOME} straight from
* the JVM's real environment — no test could isolate it, and on any machine carrying a real
* {@code ~/.config/git/ignore} (this repo's own dev machine does — measured, not assumed), every
* seeding test below silently composed with that real file instead of a controlled fixture.
* Every test that seeds at least one skill now constructs its {@link GitWorktrees} with this —
* an empty, machine-independent {@code XDG_CONFIG_HOME} (so the fallback resolves to a file that
* provably does not exist) plus the same {@code GIT_CONFIG_GLOBAL}/{@code GIT_CONFIG_SYSTEM}/
* {@code GIT_TERMINAL_PROMPT} isolation the {@link #git}/{@link #gitOutput} helpers already use
* for repo setup — so no test in this class can reach the real machine's home directory.
*/
private static Map<String, String> hermeticGitEnv(Path tmp) {
return Map.of(
"GIT_CONFIG_GLOBAL", "/dev/null",
"GIT_CONFIG_SYSTEM", "/dev/null",
"GIT_TERMINAL_PROMPT", "0",
"XDG_CONFIG_HOME", tmp.resolve("hermetic-xdg-config-home-" + System.nanoTime()).toString());
}
/** {@link GitWorktrees}'s full test seam, with a {@code memberSkillsSource} and no other
* overrides — the shape every seeding test below needs, isolated via {@link #hermeticGitEnv}. */
private static GitWorktrees seedingGitWorktrees(Path root, String memberSkillsSource, Path tmp) {
return new GitWorktrees(root.toString(), null, _ -> {}, null, null, memberSkillsSource,
hermeticGitEnv(tmp));
}
/** Acceptance criterion 2 (part 1): a worktree with no {@code .claude/} at all gets the skill
* copied in from the configured {@code memberSkillsSource}, structure and content intact. */
@Test
void seedSkillsCopiesIntoAWorktreeWithNoClaudeDirAtAll(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
Path skillsSource = tmp.resolve("skills-src");
writeSkill(skillsSource, "implementer", "IMPLEMENTER SKILL\n");
String wt = seedingGitWorktrees(tmp.resolve("wts"), skillsSource.toString(), tmp)
.add(repo.toString(), "cb-362-fresh", "HEAD");
assertEquals("IMPLEMENTER SKILL\n",
Files.readString(Path.of(wt, ".claude", "skills", "implementer", "SKILL.md")));
}
/** Acceptance criterion 2 (part 2) / invariant 1: a repo that already ships its own {@code
* implementer} skill keeps it byte-for-byte — fleetd's copy is never written over it, even
* though the configured source also carries a same-named skill with different content. */
@Test
void seedSkillsNeverOverwritesAReposOwnSkill(@TempDir Path tmp) throws Exception {
Path repo = tmp.resolve("repo");
Files.createDirectories(repo);
git(repo, "init", "-q", "-b", "main");
git(repo, "config", "user.email", "test@example.invalid");
git(repo, "config", "user.name", "Test");
writeSkill(repo.resolve(".claude/skills"), "implementer", "REPO OWN SKILL\n");
git(repo, "add", ".claude");
git(repo, "commit", "-q", "-m", "repo ships its own implementer skill");
Path skillsSource = tmp.resolve("skills-src");
writeSkill(skillsSource, "implementer", "FLEETD SKILL — must never land here\n");
String wt = new GitWorktrees(tmp.resolve("wts").toString(), null, skillsSource.toString())
.add(repo.toString(), "cb-362-repo-own", "HEAD");
assertEquals("REPO OWN SKILL\n",
Files.readString(Path.of(wt, ".claude", "skills", "implementer", "SKILL.md")),
"the repo's own committed skill must survive untouched");
}
/** Acceptance criterion 2 (part 3) / invariant 3: a misconfigured or missing {@code
* memberSkillsSource} must never fail the spawn — the worktree is still created. */
@Test
void seedSkillsIsBestEffortWhenSourceDoesNotExist(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
String missingSource = tmp.resolve("does-not-exist").toString();
String wt = new GitWorktrees(tmp.resolve("wts").toString(), null, missingSource)
.add(repo.toString(), "cb-362-missing-src", "HEAD");
assertTrue(Files.isDirectory(Path.of(wt)), "the spawn must still produce a worktree");
assertFalse(Files.exists(Path.of(wt, ".claude", "skills")),
"nothing should be seeded when the source directory does not exist");
assertTrue(capturedMessages().stream().anyMatch(m -> m.contains("is not a directory")),
"expected a warning naming the bad memberSkills source, got:\n" + capturedMessages());
}
/** Acceptance criterion 3: prove invariant 2 with a real git command — a freshly seeded skill
* must not appear in {@code git status --porcelain} for the worktree it was seeded into. */
@Test
void seedSkillsHidesSeededPathsFromGitStatus(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
Path skillsSource = tmp.resolve("skills-src");
writeSkill(skillsSource, "implementer", "IMPLEMENTER SKILL\n");
String wt = seedingGitWorktrees(tmp.resolve("wts"), skillsSource.toString(), tmp)
.add(repo.toString(), "cb-362-status", "HEAD");
assertEquals("", fullStatus(Path.of(wt)),
"a seeded skill must be invisible to git status, so it can never be staged or committed");
}
/** Invariant 2, the other direction: the exclude {@link #seedSkillsHidesSeededPathsFromGitStatus}
* proves is scoped to ONE worktree, not the whole repo. A second worktree of the same repo,
* provisioned with no {@code memberSkillsSource}, still reports an untracked {@code
* .claude/skills/} the ordinary way — proving the exclude did not leak in via the shared
* {@code .git/info/exclude} (which a linked worktree resolves to the repo's COMMON git dir). */
@Test
void seedSkillsExcludeDoesNotLeakIntoASiblingWorktree(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
Path skillsSource = tmp.resolve("skills-src");
writeSkill(skillsSource, "implementer", "IMPLEMENTER SKILL\n");
GitWorktrees seeding = seedingGitWorktrees(tmp.resolve("wts"), skillsSource.toString(), tmp);
// `plain` never seeds anything (memberSkillsSource is null, so seedSkills no-ops before it
// ever touches core.excludesFile), so it does not need the hermetic gitEnv seam.
GitWorktrees plain = new GitWorktrees(tmp.resolve("wts").toString());
String seededWt = seeding.add(repo.toString(), "cb-362-scope-a", "HEAD");
String plainWt = plain.add(repo.toString(), "cb-362-scope-b", "HEAD");
// Simulate the same untracked shape landing in the sibling worktree by hand, since `plain`
// was never configured with a memberSkillsSource to seed it itself.
writeSkill(Path.of(plainWt, ".claude", "skills"), "implementer", "unrelated untracked content\n");
assertEquals("", fullStatus(Path.of(seededWt)), "seeded worktree stays clean");
assertTrue(porcelainPaths(fullStatus(Path.of(plainWt))).contains(".claude/"),
"an unrelated worktree's own untracked .claude/ must still show up in its status — "
+ "the seeded worktree's exclude must not have leaked into it, got:\n"
+ fullStatus(Path.of(plainWt)));
}
/** Criterion 2's log shape, mirroring the {@code overlayParity} log assertions above: the
* denominator, what was seeded, and what was kept because the repo already had it. */
@Test
void seedSkillsLogsSeededAndKept(@TempDir Path tmp) throws Exception {
reportingLogger.setLevel(Level.INFO);
Path repo = tmp.resolve("repo");
Files.createDirectories(repo);
git(repo, "init", "-q", "-b", "main");
git(repo, "config", "user.email", "test@example.invalid");
git(repo, "config", "user.name", "Test");
writeSkill(repo.resolve(".claude/skills"), "hunter", "REPO OWN HUNTER\n");
git(repo, "add", ".");
git(repo, "commit", "-q", "-m", "repo ships hunter only");
Path skillsSource = tmp.resolve("skills-src");
writeSkill(skillsSource, "hunter", "FLEETD HUNTER\n");
writeSkill(skillsSource, "implementer", "FLEETD IMPLEMENTER\n");
seedingGitWorktrees(tmp.resolve("wts"), skillsSource.toString(), tmp)
.add(repo.toString(), "cb-362-log", "HEAD");
assertTrue(capturedMessages().stream().anyMatch(m ->
m.contains("seeded: implementer") && m.contains("kept the repo's own copy of: hunter")),
"expected a summary naming both the seeded and kept skills, got:\n" + capturedMessages());
}
/**
* fleetd #362 review fix — the "compose, don't replace" invariant, a third direction alongside
* {@link #seedSkillsHidesSeededPathsFromGitStatus} and
* {@link #seedSkillsExcludeDoesNotLeakIntoASiblingWorktree}. {@code core.excludesFile} is
* single-valued: the first cut of {@code excludeSeededSkillsFromGitStatus} pointed it at
* fleetd's own exclude file with {@code --replace-all}, which SHADOWS whatever excludesFile the
* worktree was already resolving (an operator's global config, most commonly) instead of adding
* to it. Concretely, this repo's own {@code .gitignore} does not ignore {@code target/} — only an
* operator's global excludesFile does — so every worker's {@code mvn clean install} would
* otherwise make {@code target/} appear as untracked, and CB-576's deliberately
* untracked-inclusive {@code hasUncommitted} would then read every such worktree as dirty
* forever, so {@code SessionManager} never cleans it up.
*
* <p>A synthetic "operator's global git config" is isolated via {@code GIT_CONFIG_GLOBAL}
* pointed at a throwaway temp file, passed to {@link GitWorktrees} through its {@code gitEnv}
* test seam — never the real machine's own git config. That global config ignores {@code
* target}. A skill is then seeded through the real {@link GitWorktrees#add} path, and a file
* named {@code target} is written into the worktree afterward: {@code git status --porcelain}
* must still be empty, proving the operator's own global pattern kept applying after seeding.
*/
@Test
void seedSkillsComposesWithAnAlreadyEffectiveGlobalExcludesFile(@TempDir Path tmp) throws Exception {
Path globalExcludes = tmp.resolve("operator-global-ignore");
Files.writeString(globalExcludes, "target\n");
Path globalConfig = tmp.resolve("operator-global.gitconfig");
Files.writeString(globalConfig, "[core]\n\texcludesFile = " + globalExcludes + "\n");
Map<String, String> gitEnv = Map.of(
"GIT_CONFIG_GLOBAL", globalConfig.toString(),
"GIT_CONFIG_SYSTEM", "/dev/null",
"GIT_TERMINAL_PROMPT", "0",
// core.excludesFile is explicitly set above, so the XDG fallback branch is never
// reached here — this is belt-and-braces so the test stays hermetic even if that
// ever changes, matching every other seeding test in this file.
"XDG_CONFIG_HOME", tmp.resolve("unused-xdg-config-home").toString());
Path repo = initRepo(tmp.resolve("repo"));
Path skillsSource = tmp.resolve("skills-src");
writeSkill(skillsSource, "implementer", "IMPLEMENTER SKILL\n");
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString(), null, _ -> {},
null, null, skillsSource.toString(), gitEnv);
String wt = gitWorktrees.add(repo.toString(), "cb-362-global-compose", "HEAD");
assertEquals("IMPLEMENTER SKILL\n",
Files.readString(Path.of(wt, ".claude", "skills", "implementer", "SKILL.md")),
"fixture check — the skill really was seeded");
Files.writeString(Path.of(wt, "target"), "build output the operator's global config ignores\n");
String porcelain = fullStatus(Path.of(wt));
assertEquals("", porcelain,
"the operator's own global excludesFile pattern ('target') must still apply after "
+ "skill seeding ran — got:\n" + porcelain);
}
/**
* fleetd #362 review fix, finding 2: pins the XDG-fallback branch of {@link
* GitWorktrees#previouslyEffectiveExcludesFileContent}, exercised when {@code core.excludesFile}
* is unset entirely (no global, local, or worktree-scoped value at all) — the branch that used to
* read {@code XDG_CONFIG_HOME} straight from the JVM's own environment, unreachable by any test
* seam, and would silently compose with whatever real {@code ~/.config/git/ignore} happened to
* exist on the machine running the suite. {@code GIT_CONFIG_GLOBAL} points at an empty file (so
* {@code core.excludesFile} is genuinely unset, forcing the fallback branch to fire — not the
* "already configured" branch {@link #seedSkillsComposesWithAnAlreadyEffectiveGlobalExcludesFile}
* covers), and {@code XDG_CONFIG_HOME} is isolated through the {@code gitEnv} seam at a throwaway
* temp dir carrying a synthetic {@code git/ignore} that ignores {@code xdg-fallback-marker}. A
* skill is seeded through the real {@link GitWorktrees#add} path, and a file named {@code
* xdg-fallback-marker} is written into the worktree afterward: {@code git status --porcelain}
* must still be empty, proving the XDG-default pattern kept applying after seeding.
*
* <p>Deleting the fallback (so an unset key composes with {@code ""}) turns this test red with:
* {@code expected: <> but was: <?? xdg-fallback-marker\n>} — see the PR body for the pasted
* failure from actually running that mutation.
*/
@Test
void seedSkillsComposesWithTheXdgDefaultExcludesFileWhenNoneIsConfigured(@TempDir Path tmp) throws Exception {
Path xdgConfigHome = tmp.resolve("xdg-config-home");
Files.createDirectories(xdgConfigHome.resolve("git"));
Files.writeString(xdgConfigHome.resolve("git").resolve("ignore"), "xdg-fallback-marker\n");
Path emptyGlobalConfig = tmp.resolve("empty-global.gitconfig");
Files.writeString(emptyGlobalConfig, "");
Map<String, String> gitEnv = Map.of(
"GIT_CONFIG_GLOBAL", emptyGlobalConfig.toString(),
"GIT_CONFIG_SYSTEM", "/dev/null",
"GIT_TERMINAL_PROMPT", "0",
"XDG_CONFIG_HOME", xdgConfigHome.toString());
Path repo = initRepo(tmp.resolve("repo"));
Path skillsSource = tmp.resolve("skills-src");
writeSkill(skillsSource, "implementer", "IMPLEMENTER SKILL\n");
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString(), null, _ -> {},
null, null, skillsSource.toString(), gitEnv);
String wt = gitWorktrees.add(repo.toString(), "cb-362-xdg-fallback", "HEAD");
assertEquals("IMPLEMENTER SKILL\n",
Files.readString(Path.of(wt, ".claude", "skills", "implementer", "SKILL.md")),
"fixture check — the skill really was seeded");
Files.writeString(Path.of(wt, "xdg-fallback-marker"),
"build output the XDG default ignore file (not core.excludesFile) covers\n");
String porcelain = fullStatus(Path.of(wt));
assertEquals("", porcelain,
"the XDG default excludesFile pattern ('xdg-fallback-marker') must still apply "
+ "after skill seeding ran — got:\n" + porcelain);
}
/**
* fleetd #369, acceptance criterion 4 — make the fix hard to undo by accident. Every git
* subprocess this class starts is required to go through {@link #gitProcessBuilder}, the one
* place {@link #hermeticEnv} is applied; a helper built directly, the way the original leak in
* {@link #status}/{@link #fullStatus} was, is now a source-level fact this test can catch by
* name instead of a machine-dependent failure someone has to rediscover.
*
* <p>This counts a literal marker in this very file's own source, split into three
* concatenated pieces below so the count is not thrown off by this method's own text — a
* plain, unsplit occurrence of the marker anywhere in this file (a helper's construction, or a
* comment that happens to spell it out contiguously) adds to the count the same way. Today
* there are exactly two: the factory itself, and the one documented exception in {@link
* #worktreeCredentialHelperCompletesWithoutUsingAnInheritedHelper}, which needs a
* non-hermetic, test-controlled global config to prove the credential helper ignores it. A
* third means a new helper was added the old, leak-prone way — route it through {@link
* #gitProcessBuilder} instead, or explain the new exception here and bump this number.
*/
@Test
void everyGitSubprocessGoesThroughTheHermeticFactory() throws Exception {
Path source = Path.of("src/test/java/dev/ltms/fleet/session/GitWorktreesTest.java");
String text = Files.readString(source);
String marker = "new " + "ProcessBuilder" + "(";
int count = 0;
for (int from = text.indexOf(marker); from >= 0; from = text.indexOf(marker, from + marker.length())) {
count++;
}
assertEquals(2, count,
"expected exactly 2 direct git-subprocess constructions in this file (the "
+ "gitProcessBuilder factory itself, plus the one documented exception in "
+ "worktreeCredentialHelperCompletesWithoutUsingAnInheritedHelper) — a "
+ "different count means a helper now bypasses the hermetic factory; route "
+ "it through gitProcessBuilder or document the new exception here");
}
/**
* fleetd #369 review round 2. {@link #everyGitSubprocessGoesThroughTheHermeticFactory} counts
* call sites, not behaviour — it catches a NEW helper built the old, leak-prone way, but it
* cannot catch {@link #gitProcessBuilder} itself being gutted: deleting {@code
* pb.environment().putAll(hermeticEnv())} from inside the factory leaves every call site
* unchanged, the count stays 2, and the whole unpoisoned suite stays green — the exact leak
* this ticket fixed would come back silently, with nothing but a human remembering to re-run
* the poison command to catch it. This test instead inspects what the factory actually hands
* to {@link ProcessBuilder#start()}, so it fails the moment the hermetic environment stops
* being applied, on any machine, with no poison needed.
*
* <p>The property under test: every git subprocess this class starts must run with an
* environment that cannot see the operator's real git configuration. A call-site count is a
* proxy for that; this is the thing itself.
*/
@Test
void gitProcessBuilderCarriesTheFullHermeticEnvironment(@TempDir Path tmp) {
Map<String, String> env = gitProcessBuilder(tmp, "status", "--porcelain").environment();
assertEquals("/dev/null", env.get("GIT_CONFIG_GLOBAL"),
"GIT_CONFIG_GLOBAL must be neutralized, or the operator's real ~/.gitconfig applies");
assertEquals("/dev/null", env.get("GIT_CONFIG_SYSTEM"),
"GIT_CONFIG_SYSTEM must be neutralized, or the machine's real /etc/gitconfig applies");
assertEquals("0", env.get("GIT_TERMINAL_PROMPT"),
"GIT_TERMINAL_PROMPT must be disabled, or a credential prompt can hang the subprocess");
String xdg = env.get("XDG_CONFIG_HOME");
assertNotNull(xdg,
"XDG_CONFIG_HOME must be set — left unset, git falls back to the operator's real "
+ "$HOME/.config/git/ignore (gitignore(5)), exactly fleetd #369's leak");
assertFalse(xdg.isBlank(), "XDG_CONFIG_HOME must not be blank — blank behaves like unset");
assertTrue(Path.of(xdg).startsWith(CLASS_TMP),
"XDG_CONFIG_HOME must point inside this test class's own throwaway directory, "
+ "never the operator's real one or the JVM's inherited value — got: " + xdg);
assertFalse(Files.exists(Path.of(xdg, "git", "ignore")),
"the resolved XDG default excludes file must provably not exist, or its contents "
+ "would silently apply to every git status this test class runs");
}
}
+337
View File
@@ -0,0 +1,337 @@
# Fleet as a Claude Code plugin — plan
Status: draft for architect review. Not implemented.
Author: primary (lead `opus`, Mac fleet). Date: 2026-09-05.
## 0. Correction — this already exists, and that changes the plan
I wrote sections below as if the plugin were new work. It is not. **This repo is already a Claude
Code marketplace and already ships a plugin**, added in `ef1e014` (CB-527) and last touched in
`2e138a1` (CB-634):
```
.claude-plugin/marketplace.json -> name "claude-bridge", plugins: [ ./plugin ]
plugin/.claude-plugin/plugin.json -> name "claude-bridge", version 0.1.0
plugin/.mcp.json -> mounts "fleetd" at http://127.0.0.1:8765/mcp
plugin/skills/setup/SKILL.md -> a full onboarding skill
plugin/README.md
```
The `setup` skill is good and covers most of what section 5 proposes: preflight, merge-not-clobber
into `.mcp.json`, read-only permissions only, credentials by env-var name, and a verify step that
insists on a **real spawn** because a green `/healthz` proves nothing.
So the operator's question — "can we pack things into plugins?" — is already answered *yes, and it
was built*. The real question is why it did nothing for the kb session. The answer is drift plus
invisibility.
### The drift, measured
| # | Finding | Evidence |
|---|---|---|
| 1 | **Mount name mismatch.** The plugin mounts the server as `fleetd`; the daemon's own constant is `fleet` | `plugin/.mcp.json` vs `PeerLauncher.java:34` `String MCP_MOUNT_NAME = "fleet"` |
| 2 | **URL hardcoded**, no env indirection, so one plugin cannot serve two hosts or ports | `plugin/.mcp.json` |
| 3 | **Ships no worker skills and no agents** | `plugin/` has 1 skill (`setup`); `.claude/skills/` has 5 and `.claude/agents/` has 3, none of them in `plugin/` |
| 4 | **Stale identity advice.** `setup` §5 tells the operator to pin `primary.terminal:` | CB-579 replaced that with `fleet.leaders.*.tab`. `record Primary` still exists (`FleetConfig.java:954`), so the advice is not dead — but it is no longer the mechanism |
| 5 | **Stale install path.** README says `/plugin marketplace add ltms/claude-bridge` | the repo is `fleet/fleetd` since CB-623 |
| 6 | **Stale names.** Plugin and marketplace are both `claude-bridge` | the project renamed to `fleetd` in CB-634 |
| 7 | **Nothing references it.** `grep -rn "plugin/" CLAUDE.md docs/*.md` returns nothing | so no session is ever told the plugin exists — which is exactly why I planned it from scratch |
Finding 7 is the root cause of the other six. A shipped capability that no instruction file
mentions gets no maintenance, and the next person rebuilds it. That is the same failure the
`CLAUDE.md` "Features" rule was written to stop.
Finding 3 is the one that matters most for the operator's actual problem. A worker spawned into a
**kb** worktree has no `implementer` skill, because only `claude-bridge` carries one in
`.claude/skills/`. Every brief that says "Load the implementer skill" is a no-op outside this repo.
The plugin is the right home for those skills and does not carry them yet.
## 1. The problem, restated
A project is "fleet-enabled" today by hand-edits nobody wrote down in one place:
- a `fleet.leaders.<name>` entry in a host's gitignored `fleetd.yaml`;
- the repo must carry `.claude/skills/*` for a worker to load `implementer` or `reviewer`;
- the repo must carry the canonical bridge block in its `CLAUDE.md`;
- the MCP mount arrives only because `LeadLauncher` and `ClaudeCodeLauncher` add `--mcp-config`
to the argv they build.
Shown live on 2026-09-05: an operator opened `claude` by hand in `/home/ltms/LTMS/kb` on fleet01
and the session had **no `fleet_*` tools at all**, because a hand-started agent never gets the
launcher's `--mcp-config`. The plugin would have fixed that — if it had been installed, and if it
had been mentioned anywhere.
## 2. What was verified, and how
| Claim | Evidence |
|---|---|
| A plugin can install globally | `~/.claude/plugins/installed_plugins.json` — scopes in use are `project` (5), `local` (4), **`user` (1)** |
| A plugin can mount an MCP server | `~/.claude/plugins/marketplaces/kb-alms/.mcp.json` mounts `memory` at `"url": "${KB_MEMORY_URL}"` |
| A plugin can carry skills, agents, commands, hooks | `kb-alms` ships `skills/` + `hooks/hooks.json`; `umputun-cc-thingz/plugins/planning` ships `agents/` + `skills/` |
| Env vars interpolate in a plugin's `.mcp.json` | same `kb-alms` file: `${KB_MEMORY_URL}`, `${MEMORY_MCP_TOKEN}` |
| A marketplace can be a plain git repo | `known_marketplaces.json` — `mgnl-code-review` has `"source": "git", "url": "https://..."` |
| **This repo is already such a marketplace** | `.claude-plugin/marketplace.json`, committed in `ef1e014` |
| A lead already binds to a project directory | `FleetConfig.java:1017` `record Leader(..., String workspace, String cwd)`; used at `LeadLauncher.java:193-199` |
| The lead's mount comes from argv, not config | `LeadLauncher.java:253` adds `--mcp-config` |
## 2b. Architect review + one measurement changed the design
The architect verified the plan against the code and returned **build it with these changes**. Two
of its findings are load-bearing. I checked both myself.
### A plugin cannot carry the agent definitions — confirmed
`ClaudeCodeLauncher.java:371` calls
`agentDefinitionFile(spec.cwd(), spec.role(), ".claude", "agents")`, and
`HerdrPeerLauncher.java:359-365` returns a path only when
`<cwd>/.claude/agents/<role>.md` `Files.isRegularFile`. `ClaudeCodeLauncher.java:391-393` adds
`--agent` **only** when that returns non-null. `OpenCodeLauncher` does the same for
`.opencode/agent`.
So the file must exist **in the member's worktree**. Moving `.claude/agents/*.md` into the plugin
would silently stop every member getting `--agent`. **The agents stay in the repo.** My plan had
this wrong.
### The plugin does not reach members at all — confirmed, and worse than the architect could see
`ClaudeCodeLauncher.java:285` does `putIfPresent(workerEnv, "CLAUDE_CONFIG_DIR", cfg.configDir())`.
A member with `configDir` set reads that directory, not the operator's `~/.claude`.
The architect could not check how far that goes, because `fleetd.yaml` is gitignored. I measured it:
- **All four Claude profiles on the live fleet set `configDir`** — `local`, `local-direct`, `opus`,
`sonnet` (`grep -c configDir fleetd/fleetd.yaml` = 4).
- Each `ccs` instance has its **own** `plugins/` directory: `gx10` (8 entries), `ltms` (11),
`ollama` (9), `work` (9).
- Those directories are **four separate real directories with four separate inodes**, and
`installed_plugins.json` in each is a **separate inode with an identical md5**
(`51c6e1c853e32e656b817e123fbbfcc5`). They are *copies made once*, not links.
So a plugin installed at user scope lands in exactly one instance's store. It would have to be
installed once per `CLAUDE_CONFIG_DIR`, and each copy would then drift. **The plugin is not a
delivery mechanism for member-facing assets on this host.**
A side effect worth recording: my own session's `CLAUDE_CONFIG_DIR` is set, so the
`~/.claude/plugins/*` evidence in section 2 is not even this session's store. The claims about what
a plugin *can* do still hold — they were read from real manifests — but the directory I read them
from is the wrong one for any conclusion about *this* session.
### The design that follows
Split by audience, not by mechanism:
| Audience | Delivered by | Carries |
|---|---|---|
| operator / lead (a human opening any project) | **the plugin**, per config dir | the MCP mount, `setup`, the bridge charter |
| member (a worker in a provisioned worktree) | **worktree provisioning** | `.claude/skills/*`, `.claude/agents/*` |
The second row is not a new idea — it is what the code already does for agents, and it is why
`agentDefinitionFile` looks in the worktree. Extending worktree provisioning to seed
`.claude/skills/` from a fleetd-owned source is the consistent move, and it is what actually fixes
"a worker in kb has no `implementer` skill". The plugin never could.
### Other review findings I accepted
- **Mount-name collision is real.** `PeerLauncher.java:34` is `fleet`; the plugin mounts `fleetd`.
A lead with both gets two mounts of one daemon and duplicate `fleet_*` tools. Rename the
plugin's server to `fleet`.
- **Keep `--mcp-config` in `LeadLauncher`.** It is config→argv from `profile.mcpUrl()`, not drift,
and it is the only path that works for a member with its own `configDir`.
- **I overstated #359.** `LeadCoordLoop.java:174-197` returns null and logs a warning that names the
fix; `tick()` leaves the message unacked, so the broker holds it and delivers once a lead is
named. It **stalls loudly and recovers** — it is not silent, and it is not data loss. My wording
in issue #361 needs the same correction.
- **Stage 1 ends with tools that mostly cannot be used until stage 3**, because authority still
comes from the tab. Section 3 already said this; the stage table did not.
## 3. What a plugin can and cannot do
This is the part that decides the design, so it is stated before the design.
**A plugin gives tools. It does not give authority.**
`fleet_whoami` resolves a caller's role from the connection, not from what is mounted. A
hand-started session in kb that mounts `fleet_*` through a plugin will be resolved as a **worker**
and refused on every orchestration call, because its pane is not in a tab matching
`fleet.leaders.*.tab`.
So the plugin alone does not make a project fleet-enabled. It makes it *tool*-enabled. Registering
the lead stays fleetd's job. Any plan that forgets this ships a plugin that looks installed and
does nothing.
```mermaid
flowchart TB
P["fleet plugin<br/>(user scope, every session)"] --> T["fleet_* tools mounted"]
D["fleetd.yaml<br/>leaders.kb {tab, cwd}"] --> A["role = primary"]
T --> W["can call fleet_*"]
A --> W2["calls are authorized"]
W --> OK["working lead"]
W2 --> OK
T --> NO["tools mounted, every call refused"]
classDef good fill:#2f855a,stroke:#22543d,color:#ffffff;
classDef bad fill:#9b2c2c,stroke:#63171b,color:#ffffff;
class OK good
class NO bad
```
*Both halves are needed. The plugin is the left half only.*
## 4. Proposed architecture — three layers
### Layer 1: the plugin — lead-side only, fix the one that exists
Keep it at `plugin/`, keep the marketplace at `.claude-plugin/marketplace.json`. Do not create a
second one, and do not put member-facing assets in it (see 2b).
```
.claude-plugin/marketplace.json -> rename to "fleetd"; keep source ./plugin
plugin/
.claude-plugin/plugin.json -> rename to "fleet"; bump version
.mcp.json -> mount name "fleet" (match PeerLauncher.MCP_MOUNT_NAME),
url "${FLEETD_MCP_URL}", 8765 default documented
skills/setup/SKILL.md -> EXISTS. fix the stale primary.terminal advice (§5)
skills/bridge-charter/SKILL.md -> NEW: the canonical CLAUDE.md block
README.md -> fix the install path (fleet/fleetd, not ltms/claude-bridge)
```
**Not in the plugin:** `agents/*.md` (the launcher requires them in the member's worktree —
`ClaudeCodeLauncher.java:371,391`) and the three worker skills (a member with `configDir` set never
reads the operator's plugin store — measured in 2b). Those belong to layer 1b.
**A rename is a breaking change for anyone who installed 0.1.0.** The mount name goes `fleetd` ->
`fleet`, so a project whose `.claude/settings.json` pre-allows `mcp__fleetd__fleet_whoami` stops
matching. Only this fleet has it installed today, so the cost is small now and grows. Decide once.
**This solves propagation of the charter.** `CLAUDE.md` says the bridge block "must stay
byte-identical with the template in the wiki" and that "other projects carrying the block need the
same edit" — a hand-copy the file itself admits is fragile, with a python snippet to check it. A
plugin skill turns that into a version bump.
### Layer 1b: worker skills reach members through the worktree, not the plugin
This is the change that actually fixes "a worker in kb cannot load `implementer`".
Worktree provisioning already writes into the member's tree — the parity overlay, the neutralised
`.mcp.json`, the IDE overlay. Add one more: seed `<worktree>/.claude/skills/` from a fleetd-owned
source directory, so every member gets `implementer`, `reviewer` and `hunter` whatever repo it is
working in. `.claude/agents/` is already required there by the launcher, so this follows the
grain of the design rather than cutting across it.
Open question for implementation: copy or symlink, and where the source lives (a config key such
as `memberSkills:`, or the plugin's own directory read by the daemon). A symlink is one source of
truth but breaks if the member's tree is archived; a copy drifts but is self-contained.
### Layer 2: host-global fleet settings
`~/.fleet/fleetd.yaml` — the things that are true for the **machine**, not the project:
- `broker:` and `coordinator:` (URIs come from env, no secrets in the file)
- `profiles:` — backends, models, credentials, weights
- `memberCredentials:` policy
- `worktreeRoot`, `worktreeGroup`
### Layer 3: per-project settings, committable
`<project>/.fleet/project.yaml` — the things that are true for the **repo**:
```yaml
lead:
tab: "lead: kb"
profile: opus
ide:
projectDir: "" # kb is a Python repo, everything at the root
worktree: true
```
**This split fixes a contradiction that exists today.** `ideProjectDir` is a property of a *repo*
(fleetd's Maven module is a subdirectory; kb's code is at the root) but the config key is
per-*profile*. One profile therefore cannot serve both repos — measured on fleet01 on 2026-09-05,
where the key had to be commented out to make kb work. Moving it to a project file removes the
contradiction rather than working around it.
It is also committable, because it holds no secrets. A project that has been fleet-enabled once
stays fleet-enabled for everyone who clones it.
## 5. The `fleet-setup` skill
What the operator actually asked for: one command that makes any project fleet-compatible.
```mermaid
sequenceDiagram
participant Op as Operator
participant Sk as fleet-setup skill
participant Fs as project files
participant Fd as fleetd
Op->>Sk: /fleet-setup (in any project)
Sk->>Fs: write .fleet/project.yaml
Sk->>Fs: add the bridge block to CLAUDE.md (if absent)
Sk->>Fd: register the lead (tab + cwd)
Fd-->>Sk: tab created, lead launched
Sk-->>Op: report what changed, and what is still manual
```
*The skill writes the project half and asks the daemon for the host half.*
The registration step needs something that does not exist yet: an MCP tool such as
`fleet_workspace_add{path, tab, profile}`, or a `fleetd` config include so a project file is picked
up without hand-editing the host file. **This is the one genuinely new piece of daemon work.**
## 6. Does this reduce fleetd's complexity?
Honestly: **partly**. Claiming more than this would be wrong.
**Yes, in three places.**
1. Skill and agent delivery leaves the daemon and the repos entirely.
2. Worktree config neutralisation (fleetd #134) gets safer. It blanks `.mcp.json` so the primary's
IDE and forge servers do not leak into a worker. Today the fleet mount survives only because
the launcher re-adds it by argv. With a user-scope plugin the fleet mount is outside the file
being neutralised, so the two concerns stop fighting.
3. The `ideProjectDir` per-profile/per-repo contradiction disappears.
**No, in the places that matter most.** fleetd still owns spawn, authorization, worktrees, herdr,
the broker, tickets, and identity. A plugin cannot do any of those. The plugin is a **distribution**
mechanism, not a replacement for the daemon.
**And it adds one new risk.** The mount URL becomes a second source of truth. `fleetd.yaml` has the
port; the plugin has the URL. Mitigation: the plugin reads `${FLEETD_MCP_URL}` only, and the host
env is the single place it is set.
## 7. Rollout stages
| Stage | Content | Ends with |
|---|---|---|
| 0 | **Make it visible.** One `CLAUDE.md` line and one Features entry saying the plugin exists and where | nobody re-plans it a third time |
| 1 | Fix the plugin's drift: mount name `fleet`, `${FLEETD_MCP_URL}`, names, README, stale `primary.terminal` advice | a lead in any project can install one plugin and get the mount |
| 1b | Seed `.claude/skills/` into provisioned worktrees | **a worker in *kb* can load `implementer`** |
| 2 | `.fleet/project.yaml` schema + `FleetConfig` reads it; `ideProjectDir` moves there | kb and fleetd both work off one profile |
| 3 | Fix #359, then config-include for lead registration, wired into the existing `setup` skill | `/fleet:setup` in a fresh project produces a working lead |
| 4 | Roll out to fleet01; retire the hand-copied CLAUDE.md block in favour of the skill | one `git pull` propagates the charter |
Stage 0 is minutes of work and is the one that stops this happening again, so it goes first.
**Stage 1b carries most of the value** and is independent of the plugin — it could ship first if the
plugin rename needs more thought. Stage 1 alone ends with tools a hand-started session mostly
cannot use, because authority still comes from the tab; that is fixed in stage 3, not stage 1.
#359 moves ahead of stage 3 on the architect's advice, because stage 3 is what creates the second
lead.
## 8. Questions for the architect
1. **Is the layer-2 / layer-3 split right?** Specifically: should `profiles:` stay host-global, or
should a project be able to pin which profiles it uses? Cost of getting this wrong is a config
that has to be re-split later.
2. **Config include, or a new MCP tool, for registering a project's lead?** An include is passive
and survives a restart; a tool is live but writes to a gitignored file the daemon owns.
3. **What happens when the plugin is absent?** Should `LeadLauncher` keep its `--mcp-config`
belt-and-braces, or is that the drift risk we should remove? Note opencode members cannot use
Claude plugins at all, so `OpenCodeLauncher` keeps its ephemeral config either way.
4. **Does a user-scope plugin mount leak into members in a way we do not want?** Members already
inherit user-scope MCP servers (`--mcp-config` adds, it does not replace). A worker getting
`fleet_*` is correct and already happens. Confirm nothing else in the plugin should be
worker-invisible.
5. **Two leads on one host both hold a subscription seat.** Is per-project leads the right unit, or
should one lead serve several projects by changing cwd?
6. **Blocking defect to fix first or alongside:** `LeadCoordLoop.resolveLocalLead()` (lines
174-190) routes a peer message to "the sole lead" when no lead is *named* after
`coordinator.selfId`. The moment a host has two leads — exactly what this plan encourages —
cross-host coordination silently stops. Tracked as #359.
+3 -3
View File
@@ -1,7 +1,7 @@
{
"name": "claude-bridge",
"description": "Make a project bridge-ready: mount the fleetd MCP gateway and set up standard Claude Code settings so this session can orchestrate a fleet of delegated workers. Ships no credentials.",
"version": "0.1.0",
"name": "fleet",
"description": "Make a project fleet-ready: mount the fleetd MCP gateway and set up standard Claude Code settings so this session can orchestrate a fleet of delegated workers. Lead-side only — member skills and agents travel in the worktree. Ships no credentials.",
"version": "0.2.0",
"author": {
"name": "LTMS"
},
+2 -2
View File
@@ -1,8 +1,8 @@
{
"mcpServers": {
"fleetd": {
"fleet": {
"type": "http",
"url": "http://127.0.0.1:8765/mcp"
"url": "${FLEETD_MCP_URL}"
}
}
}
+33 -7
View File
@@ -1,11 +1,24 @@
# claude-bridge (Claude Code plugin)
# fleet (Claude Code plugin)
Makes a project **bridge-ready**: mounts the `fleetd` MCP gateway and applies standard Claude Code
Makes a project **fleet-ready**: mounts the `fleetd` MCP gateway and applies standard Claude Code
settings, so the session can orchestrate a fleet of delegated workers.
**This plugin ships no credentials.** Every secret is referenced by environment-variable *name*;
the values stay with the user. Nothing the plugin writes is unsafe to commit.
## Scope — lead-side only
This plugin configures **the session you are sitting in**: a lead, or any human-started Claude Code
session that wants to talk to the daemon. It deliberately does **not** carry the worker playbook
skills or the role agent definitions, and it cannot:
- the launcher adds `--agent` only when `<worktree>/.claude/agents/<role>.md` exists in the
member's own tree (`ClaudeCodeLauncher.java:371,391`), so agent files must live in the repo;
- a member's `CLAUDE_CONFIG_DIR` points at its profile's config directory
(`ClaudeCodeLauncher.java:285`), so it never reads the operator's plugin store.
Member-facing assets travel in the worktree, not in this plugin. See fleetd #362.
## What it is not
The plugin is the **client-side setup**, not the bridge. `fleetd` is a separate daemon and `herdr`
@@ -16,22 +29,35 @@ not try to install system services on your behalf.
## Install
```shell
/plugin marketplace add ltms/claude-bridge
/plugin install claude-bridge@claude-bridge
/plugin marketplace add https://git.ltms.dev/fleet/fleetd
/plugin install fleet@fleetd
```
Export the gateway URL — the plugin mounts `${FLEETD_MCP_URL}`, not a hardcoded address, so one
plugin serves hosts that run the daemon on different ports:
```shell
export FLEETD_MCP_URL=http://127.0.0.1:8765/mcp
```
Then, in the project you want to onboard:
```shell
/claude-bridge:setup
/fleet:setup
```
## What you get
| Component | Effect |
|---|---|
| `.mcp.json` | mounts `fleetd` at `http://127.0.0.1:8765/mcp` for any session with the plugin enabled |
| `skills/setup` | `/claude-bridge:setup` — preflight, project settings, credential guidance, and verification |
| `.mcp.json` | mounts `fleet` at `${FLEETD_MCP_URL}` for any session with the plugin enabled |
| `skills/setup` | `/fleet:setup` — preflight, project settings, credential guidance, and verification |
The server is named **`fleet`** on purpose: that is `PeerLauncher.MCP_MOUNT_NAME` in the daemon and
the name a spawned member's own mount carries. Version 0.1.0 named it `fleetd`, which produced two
mounts of one daemon for anyone who also had a project-level `.mcp.json`. Upgrading from 0.1.0 is a
**breaking change** — a project that pre-allowed `mcp__fleetd__fleet_whoami` in
`.claude/settings.json` must be updated to `mcp__fleet__*`.
Because the plugin carries its own `.mcp.json`, an installed plugin needs no project-level MCP
file at all. The setup skill writes one only when you want the mount to work *without* the plugin —
+43 -18
View File
@@ -36,9 +36,17 @@ a time.
command -v herdr && herdr --version 2>&1 | head -1 || echo "MISSING: herdr"
command -v ccs && ccs version 2>&1 | head -1 || echo "MISSING: ccs (needed for worker profiles)"
command -v codex && codex --version 2>&1 | head -1 || echo "absent: codex (optional)"
curl -s -m 5 http://127.0.0.1:8765/healthz || echo "MISSING: fleetd daemon is not reachable"
curl -s -m 5 "${FLEETD_MCP_URL%/mcp}/healthz" 2>/dev/null \
|| curl -s -m 5 http://127.0.0.1:8765/healthz \
|| echo "MISSING: fleetd daemon is not reachable"
[ -n "$FLEETD_MCP_URL" ] && echo "FLEETD_MCP_URL is set" || echo "MISSING: FLEETD_MCP_URL"
```
**`FLEETD_MCP_URL` is required.** The plugin's own `.mcp.json` mounts `${FLEETD_MCP_URL}` rather
than a hardcoded address, so one plugin can serve hosts that run the daemon on different ports. If
it is unset the mount does not resolve. The usual value is `http://127.0.0.1:8765/mcp`; tell the
user to export it, do not write it into a file for them.
A healthy daemon answers with its status **and the herdr protocol it negotiated**:
```json
@@ -70,7 +78,7 @@ The entry to add, exactly:
```json
{
"mcpServers": {
"fleetd": {
"fleet": {
"type": "http",
"url": "http://127.0.0.1:8765/mcp"
}
@@ -78,8 +86,13 @@ The entry to add, exactly:
}
```
If `.mcp.json` already exists, add only the `fleetd` key and leave every other server untouched.
If a `fleetd` entry is already there with a different URL, **ask** rather than assuming yours is
**The server must be named `fleet`.** That is `PeerLauncher.MCP_MOUNT_NAME` in the daemon, the name
a spawned member's mount carries, and the name the `mcp__fleet__*` role heuristic in `CLAUDE.md`
keys on. An earlier version of this plugin named it `fleetd`, which gave a lead with both a project
file and the plugin **two mounts of the same daemon** and a duplicated `fleet_*` tool set.
If `.mcp.json` already exists, add only the `fleet` key and leave every other server untouched.
If a `fleet` entry is already there with a different URL, **ask** rather than assuming yours is
right — a non-default port usually means a deliberate second daemon.
> **If this plugin is installed, you can skip this step entirely.** The plugin ships its own
@@ -111,11 +124,11 @@ project already set.
"$schema": "https://json.schemastore.org/claude-code-settings.json",
"permissions": {
"allow": [
"mcp__fleetd__fleet_whoami",
"mcp__fleetd__fleet_list",
"mcp__fleetd__fleet_status",
"mcp__fleetd__fleet_profiles",
"mcp__fleetd__fleet_poll"
"mcp__fleet__fleet_whoami",
"mcp__fleet__fleet_list",
"mcp__fleet__fleet_status",
"mcp__fleet__fleet_profiles",
"mcp__fleet__fleet_poll"
]
}
}
@@ -161,19 +174,31 @@ fleet_whoami
```
- `{"role":"primary"}` — correct, you are done with this step.
- `{"role":"worker", …}` — **this is the trap.** If the primary runs inside a herdr pane, the
daemon resolves it to a terminal and classifies it as a worker, refusing `spawn`/`send`/`stop`:
every verb an orchestrator exists to call. It is **self-locking**, because the daemon can only
*learn* the primary's terminal from those same refused calls. The only way out is an
operator-set pin in the daemon's config:
- `{"role":"worker", …}` — **this is the trap.** If the lead runs inside a herdr pane whose tab the
daemon does not recognise, it is classified as a worker and refused on `spawn`/`send`/`stop`:
every verb an orchestrator exists to call. It is **self-locking**, because those are the same
calls that would tell the daemon who you are.
Identity is the **tab label**, matched exactly and case-insensitively:
```yaml
primary:
terminal: term_xxxxxxxxxxxx # the terminalId fleet_whoami just reported
fleet:
leaders:
kb: # name it after coordinator.selfId if this host uses lead-to-lead
profile: opus
tab: "lead: kb" # the exact label of the tab this lead sits in
cwd: /path/to/the/project
```
The daemon reads this **at boot**, so it needs a restart. Re-pin whenever the primary moves
panes — a stale pin fails exactly as silently as no pin.
A tab label is stable across restarts of the agent inside it, which is why CB-579 replaced the
older `primary.terminal:` pin — a herdr `terminal_id` changed on every restart and cost a config
edit each time. `primary.terminal:` still parses, but it is no longer the mechanism; do not
reach for it.
The daemon reads `leaders:` **at boot**, so a new entry needs a restart. Two things to check
afterwards: that `fleet_whoami` now answers `primary`, and that no *stale* tab carries the same
label — duplicate lead tabs are their own failure (#359), and they stall lead-to-lead delivery
until one lead is named after `coordinator.selfId`.
Then prove the fleet actually works, with a real spawn: