Compare commits

...

18 Commits

Author SHA1 Message Date
Dai Ha c50f5b2d61 fleetd #176 stage 2: make effectiveCredentialId() subscription-aware
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 1m44s
Stage 1's lead-seat matcher (leadSeatLookup) was correct but inert on
the live host: the lead runs on profile 'opus', members on 'sonnet',
both subscription:true with no explicit credentialId. Because
effectiveCredentialId() fell back to the profile's own name, opus and
sonnet never matched even though they share one Claude login, so the
matcher charged zero seats.

FleetConfig.Profile.effectiveCredentialId() now falls back to a shared
sentinel (SUBSCRIPTION_CREDENTIAL_ID = "<subscription>") instead of the
profile name when subscription:true and credentialId is unset. An
explicit credentialId still wins, so two separate Claude logins on one
host can still be kept apart.

This is also BackendQuarantine's and BackendOutagePolicy's grouping
key and CompositePeerLauncher's spawn-time enforcement key, so the fix
also links quarantine/cool-off across subscription profiles sharing an
account -- intentional: one usage limit really does take out every
profile on that login, mirroring credentialId: openai-shared already
doing this for off-subscription profiles. Every caller was reviewed;
none wants "this exact profile" over "this account".

Tests added:
- FleetdLeadSeatLookupTest: the live shape itself (lead on a
  DIFFERENT subscription profile than the target, same account,
  neither sets credentialId) -- the case stage 1's suite never covered
- FleetMcpTest: quarantining one subscription profile's shared
  account zeroes free on another sharing it, via the same
  effectiveCredentialId()-driven wiring Fleetd.main uses

Mutation-tested: reverting the subscription branch to the old
fall-back-to-profile-name behavior sends both new tests RED with 0
compile errors; reverting the mutation restores byte-identical
(diff -q) source and green tests.

fleetd.example.yaml's fleetd #176 notes are rewritten for the sentinel
semantics and when to override it with an explicit credentialId.
2026-09-03 13:10:10 +07:00
Dai Ha c796eac09c fleetd #176: subtract the lead's own subscription seat from free
CI / contract (pull_request) Successful in 1m11s
CI / build (pull_request) Successful in 1m16s
maxLoad counted panes, never subscription seats: a subscription:true
profile's lead is itself a live claude session on that same account,
so free overstated capacity by the lead's own seat (measured free:1
with a real ceiling of 0, and free:3 on an idle fleet with a real
ceiling of 2).

Add FleetMcp.LeadSeatSource (same shape as QuarantineSource/
OutageSource) and Fleetd.leadSeatLookup, which derives the seat count
from fleet.leaders.<name>.profile matched against the target profile
by effectiveCredentialId() - no hardcoded "-1", and no new config key:
profile: already exists for this exact "which account does this lead
share" question. maxLoad itself is left untouched; only free (and a
new, additive-only leadSeats field) changes.

Exhaustion quarantine (cause 2 in the ticket) already forced free to 0
via the same BackendQuarantine capacityView already reads - confirmed
by reading the exhaustionSink wiring, no code change needed there.
2026-09-03 12:55:08 +07:00
Dai Ha 9d37f3aa29 fleetd #201: the coverage line must name the key its caller actually means
CI / contract (push) Successful in 43s
CI / build (push) Successful in 1m56s
Found by reading a real boot log after the redeploy, not by a test.

coverage() is shared by two call sites — CB-578's exhaustedPattern line
and Unit 5's errorPattern line — but its 'off' branch hard-coded the
word exhaustedPattern. So this daemon printed:

  backend-exhausted classification (CB-578 stage A): partial
      (configured: [sol, terra]; not configured: [...])
  backend-error classification (fleetd #201 Unit 5): off
      (no profile has an exhaustedPattern configured; profiles: [...])

Two lines, one directly under the other, disagreeing about whether any
profile has an exhaustedPattern. Both were individually defensible and
together they were nonsense. Worse, the message sends an operator to
set the wrong key: the thing that is missing is errorPattern.

coverage now takes the key name. I changed the signature rather than
adding an overload, so the compiler found all three existing callers
instead of leaving them silently on the old path.

Every earlier coverage test passed the exhaustion case only, which is
why none of them could see this. The new test pins the errorPattern
case. Reverting the fix turns it red with 0 compile errors.

1229 tests, 0 failures, BUILD SUCCESS.

This is the second defect in two hours found only by reading the live
startup log — see #115, where the noise of a false warning had been
hiding a correct line saying a whole feature was off.
2026-09-03 12:29:27 +07:00
Dai Ha eaf89abaf6 fleetd #248: make Fleetd's CompletionResolver wiring provable
CI / contract (push) Successful in 42s
CI / build (push) Successful in 1m40s
Before this, dropping either #241's worktree lookup or Unit 5's
backend-error pair at Fleetd.main's new CompletionResolver(...) call
left all 1216 tests green with 0 compile errors. Every existing test
built its own CompletionResolver, so they proved the class and never
the wiring. BackendOutageFlowTest was the sharpest case: it copies
main's sink lambda line-for-line, so it proves the copy and cannot
notice the original being deleted.

The three inline arguments are now package-private static factories on
Fleetd, following the deliverableTo pattern the file already had, each
with its own behaviour test. backendErrorSink is public so a
cross-package test can drive the real production object rather than a
hand-mirrored copy.

The test that was actually missing is a source-text assertion. That is
the honest fallback for a composition root with no seam, and it is
labelled [SOURCE TEXT] in every test name and message so it cannot be
misread as a behaviour check. It is not vacuous: two tests pin that the
variables are assigned from the factories, and two pin that those
variables reach the call site, so renaming a variable while assigning
an inert value does not slip through.

Known cost, accepted: the assertions match exact source substrings, so
reformatting that statement will break them. That is the price of
covering a main method, and a spurious failure here is loud and
obvious, which is the right direction to fail.

Verified by the lead, both mutations re-run against the merged code —
see the merge check.

PR #251
2026-09-03 12:22:32 +07:00
Dai Ha d895f02bc1 fleetd #248: prove main() wires CompletionResolver's arguments, not just the class
CI / contract (pull_request) Successful in 48s
CI / build (pull_request) Failing after 1m40s
Fleetd.main built three of CompletionResolver's 8 constructor arguments inline
(a worktree/branch lookup lambda, and the backend-error pattern lookup + sink
locals). Dropping any of them at the call site compiled clean and left every
existing test green, because every existing test constructs its own
CompletionResolver and only ever proves the class, never main's wiring.

Extract each into a static factory on Fleetd (worktreeBranchLookup,
backendErrorPatternLookup, backendErrorSink — the same static-factory pattern
Fleetd.deliverableTo already uses), test each factory's own behaviour, and add
a source-text assertion (FleetdCompletionResolverWiringTest) proving main's
CompletionResolver call still passes all three. backendErrorSink is public so
BackendOutageFlowTest can exercise the real production sink directly instead
of the hand-mirrored copy its own class doc used to describe.

No production behaviour changes — mechanical extraction only.
2026-09-03 12:20:28 +07:00
Dai Ha 43206cac2f fleetd #148 point 2: drop .envrc from the default parity overlay
CI / contract (push) Successful in 1m9s
CI / build (push) Successful in 1m18s
The default is now [.env], not [.env, .envrc]. .env is data, so copying
it into a worker worktree can only move values. .envrc is executable
shell that direnv runs on every cd, so copying it moves behaviour. Those
are different risks and should not share a default.

The knob is unchanged. An operator who wants .envrc copied writes
parityOverlay: ['.env', '.envrc'] and owns that choice; a new test pins
that escape hatch, because without it this would be a removal rather
than a re-default.

Decision recorded on the ticket, with the evidence it asked for first:
this checkout has no .env and no .envrc, and direnv is not on PATH, so
there was no live exposure. Point 1 (extend the credential scrub to
direnv) is declined and the reason is on the ticket — the scrub is a
one-shot .zlogin and a direnv hook runs on every cd, so no amount of
work on the scrub can cover it. Not copying the executable file is the
smaller change and removes the need.

The worker also fixed WorktreeSessionManagerTest, which hardcoded the
same default at another layer and broke the build. Outside its named
scope, correctly flagged rather than done silently.

Verified by the lead: 1216 tests, 0 failures, 0 compile errors.

PR #250
2026-09-03 12:03:21 +07:00
Dai Ha 8bba3a8184 fleetd #148 (point 2): drop .envrc from the default parityOverlay
CI / contract (pull_request) Successful in 1m20s
CI / build (pull_request) Successful in 1m32s
.env is data; .envrc is executable shell that direnv runs on every cd, so
copying it into a worker moves behaviour, not just values. The default
parityOverlay is now [.env] only. The knob is unchanged: an operator who
wants .envrc copied can still write parityOverlay: [.env, .envrc]
explicitly.

Updates FleetConfig's default and javadoc, fleetd.example.yaml's two
mentions of the default, and the FleetConfigTest coverage: renamed the
default test, added parityOverlayExplicitEnvrcOptInStillWorks to prove
the .envrc opt-in escape hatch still works, and fixed
WorktreeSessionManagerTest#worktreeAcquireRunsParityOverlayWithProfileDefaults
which also hardcoded the old default.
2026-09-03 12:00:13 +07:00
Dai Ha 5cf3ca9a89 fleetd #241: never hand the lead back its own brief as the member's report
CI / build (push) Successful in 1m15s
CI / contract (push) Successful in 1m27s
The completion fallback scrapes a member's pane when a turn ends with no
fleet_reply. If the pane still shows the brief the lead injected, the
scrape returned that brief, and the lead read its own words as the
member's answer. A silent member looked like a member that had reported.

echoesInjectedBrief now recognises that case and refuses it. Round 1 used
plain containment in both directions, which destroyed real reports: a
genuine report that quotes the brief contains it. Round 2 keeps the safe
direction unbounded (the brief contains the scrape) and bounds the other
one at MAX_ECHO_EXCESS_CHARS, so a scrape only counts as an echo when it
adds almost nothing to the brief.

Merge note — the Fleetd.java conflict:

This call site was changed by both #201/#227 Unit 5 (backendErrorPatterns
+ backendErrorSink) and by this ticket (the worktree/branch lookup). I
resolved it onto the full 8-argument constructor so neither feature is
dropped; nowNanos has to be passed explicitly to reach that overload.

Verified by the lead: 1215 tests, 0 failures, 0 compile errors.

I also measured whether the resolution itself is protected, and it is
NOT. Both mutations at this call site stay green:
  - drop the worktree lookup (pass _ -> null): 1215 tests, 0 failures
  - drop Unit 5's patterns/sink (legacy()/none()): 1215 tests, 0 failures
Nothing in the suite covers Fleetd's composition root, so either feature
could be silently unwired here and the build would still be clean. The
tests prove the seams, not the caller. Filed separately rather than
fixed in a merge commit.

PR #245, branch worker/cb241-fallback-echo-1175e9-11
2026-09-03 11:54:06 +07:00
Dai Ha ac474981e4 fleetd #201/#227 Unit 5: wire the backend-error cool-off into config, placement and the MCP surface
A profile's credential that throws two distinct backend errors inside 60
seconds now cools off for 60 seconds. Automatic placement skips it,
an explicit fleet_spawn naming it is refused before the adapter is
called, and fleet_list/fleet_profiles report it as coolingOffForSeconds
next to the separate CB-578 quarantinedForSeconds.

Verified by the lead: see the merge check below. The worker ran 7
mutations, all killed with 0 compile errors; M7 was NOT killed on the
first pass (the assertion only checked .contains("quarantined"), which
is true of both the correct message and the mutated fallback), and the
worker strengthened it to assertEquals on the exact literal and kept
that change. That is the right call and it is reported honestly.

Two deviations, both justified in the PR:
- FixedPlacementPolicy needed the same coolingOff filter because it
  filters candidates inline instead of using PlacementPolicyUtil.
- BackendOutageFlowTest sits in dev.ltms.fleet.inject because
  CompletionResolver.InFlight is package-private there.

The startup coverage log line is a code-reading claim, not a captured
line from a live daemon. The worker said so rather than overclaiming.

PR #246, branch worker/cb201-unit5-wiring-6c12e6-8
2026-09-03 11:48:42 +07:00
Dai Ha ba04b2359b fleetd #201/#227 unit 5: wire errorPattern, cool-off spawn gate, and fleet views
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Successful in 1m52s
Wires the already-merged units into production:
- Per-profile errorPattern config (beside exhaustedPattern), compiled once at
  startup; falls back to the legacy (?i)\bAPI Error\s*: pattern when unset.
  Startup logs configured-vs-legacy coverage, same as exhaustedPattern.
- One production BackendErrorSink in Fleetd.java: mark backend_error on the
  session, resolve the profile's credential fail-loud (never
  Optional.ifPresent), record it in BackendOutagePolicy, and push a lead
  nudge on a new incident.
- CompositePeerLauncher's explicit and automatic spawn paths both refuse a
  cooling-off credential; exhaustion quarantine wins when both are active.
  PlacementContext gets a separate coolingOff set so refusal text says
  "cooling off", never "exhausted".
- fleet_list/fleet_profiles report coolingOffForSeconds as an independent
  fact from quarantinedForSeconds; both can appear together.
- fleetd.example.yaml documents errorPattern and the 2/60/60 cool-off policy;
  CLAUDE.md tells leads how to read the two independent outage states.

Also: FixedPlacementPolicy.java, not in the original file list, needed the
same coolingOff filtering as PlacementPolicyUtil (it does its own inline
candidate filtering rather than delegating).

1190 tests, 0 failures (up from the 1163 baseline); BUILD SUCCESS.
2026-09-03 11:45:11 +07:00
Dai Ha 2e5b63f6f6 fleetd #149: seed the workspace-trust entry before a claude-code spawn
CI / contract (push) Successful in 1m13s
CI / build (push) Successful in 1m29s
Claude Code asks 'is this a project you trust?' the first time it starts in a directory it
has not seen. It is interactive with no timeout, and every member spawned with worktree:true
lands in a brand-new directory. The member never reaches its first turn and never replies,
while herdr reports blocked/interactive_ready — which reads as healthy.

ClaudeCodeLauncher now seeds projects.<cwd>.hasTrustDialogAccepted in the profile's
.claude.json before the process starts. This is not a new grant: the operator already
trusted the repo by configuring the profile against it, and a worktree is a checkout of it.

The write is gated on isProvisionedWorktree(cwd) — a .git that is a regular gitdir-pointer
file, never a real checkout. That gate exists because an earlier revision of this change,
run under mutation testing, wrote to the operator's real ~/.claude.json and truncated it
from 72KB to 919 bytes. Tests using a null configDir fall back to the real user.home, so
an ungated seed reaches real files.

The write is atomic (sibling temp file + ATOMIC_MOVE, never truncate-in-place) and the
whole read-modify-write is under a lock, because .claude.json is large, live, and rewritten
by Claude Code itself while fleetd runs. Two parallel spawns are normal here.

Verified by the lead: 1173 tests, 0 failures. A truncating write turns the torn-read test
red; removing the lock turns concurrentSeedsForDifferentCwdsBothSurvive red. Both with 0
compile errors. copyPosixPermissionsIfPresent is NOT covered by a test — its mutation stays
green — but createTempFile is 0600 on POSIX by default, so the not-world-readable property
holds without it; the line only preserves a non-default mode.
2026-09-03 11:44:03 +07:00
Dai Ha 743377d6cd fleetd #149 review round 2: make the trust-dialog seed atomic and lock-protected
CI / build (pull_request) Successful in 1m21s
CI / contract (pull_request) Successful in 1m25s
Files.writeString truncates the target in place before writing, so there
was a window where .claude.json could be observed empty or half-written
- exactly the shape of the incident this ticket already hit once, but
reachable in production too: a crash/kill mid-write, or two concurrent
claude-code spawns (normal here - several run in parallel routinely)
racing a naive read-modify-write and silently discarding one spawn's
entry.

Two independent fixes, each with its own dedicated test proving it (not
the other):

- ClaudeCodeLauncher.writeAtomically: serialise to a sibling temp file in
  the same directory, then Files.move with ATOMIC_MOVE + REPLACE_EXISTING,
  preserving the target's existing POSIX permissions (.claude.json ships
  0600). A reader now only ever observes the fully-old or fully-new file,
  never a torn one. Package-visible so a test can drive it directly.
- TRUST_JSON_LOCK: a process-wide lock around seedTrustDialog's whole
  read-modify-write, so two concurrent spawns for different cwds both
  keep their entry instead of the second write discarding the first.
  Sufficient because every spawn on this daemon runs in one JVM; it does
  NOT protect against a second daemon process or the operator's own live
  Claude Code writing at the same instant - writeAtomically covers that
  case instead.

Both fail soft, same as before: any I/O failure here must never block a
spawn.

Four new tests: a large (30-project) existing file survives without
collapsing (asserted on the restored key set, not just that the result
parses); two concurrent spawns for different cwds both keep their entry
(CountDownLatch-synchronised, not a sleep); existing 0600 permissions
survive the write; and a direct test of writeAtomically with a busy-poll
reader thread proving a concurrent reader never observes a torn file.

See PR body for the full mutation-testing table, including an
honest note on which of these tests the atomicity mutation actually
caught (not the one implied by the numbering in review) and why.
2026-09-03 11:38:10 +07:00
Dai Ha 5952d559c7 fleetd #134 point 3: tell the lead which files a worktree neutralizes
CI / contract (push) Successful in 57s
CI / build (push) Failing after 1m26s
The daemon log and the worktree git config are both new in 205ad82, but neither helps a
lead who is writing a brief. This is the line that does: never brief a worker to edit
.mcp.json, opencode.json or .autoenv in its worktree, because the edit cannot be
committed and nothing will say so.

Goes in the project addendum, not the canonical block, so the byte-sync with
wiki/7-Use-Cases.md is unaffected — re-checked and still true.
2026-09-03 11:32:03 +07:00
Dai Ha 205ad823b0 fleetd #134: make tool-surface neutralization visible to the daemon and the worker
CI / contract (push) Successful in 1m18s
CI / build (push) Successful in 1m34s
isolateToolSurface replaces .mcp.json, opencode.json and .autoenv with stubs in every
provisioned worktree and marks them --skip-worktree. That neutralisation is correct and
is unchanged here — the committed files would mount the primary's credentials.

The problem was that it was invisible. A worker told to edit opencode.json read a 3-byte
stub and reported, truthfully and wrongly, that the mount key did not exist. A missing
file would have prompted a question; a plausible stub did not.

Two changes, both visibility only. The daemon now logs one info summary per provisioning
with the denominator, the files neutralised, and the consequence. And the list is recorded
in worktree-scoped git config (fleet.neutralizedConfig / fleet.neutralizedConfigNote) so a
worker can discover it from inside its own worktree with
'git config --worktree --get-all fleet.neutralizedConfig'.

Worktree-scoped config was chosen over a file in the working tree because it lives in
.git/worktrees/<nonce>/config.worktree and so can never appear in git status, and because
configureEnvironmentCredentialHelper already uses the same mechanism in the same add() call.

Verified by the lead: baseline 1172 tests, 0 failures. Reverting the summary to log.debug
goes red (2 tests), and recording into --local rather than --worktree — which would leak
the record into the shared repo config — goes red too. Both with 0 compile errors.
2026-09-03 11:31:24 +07:00
Dai Ha d654ccb818 fleetd #134: make tool-surface neutralization visible to the daemon and the worker
CI / contract (pull_request) Successful in 53s
CI / build (pull_request) Successful in 1m35s
isolateToolSurface replaced .mcp.json/opencode.json/.autoenv with neutral stubs and
marked them --skip-worktree, but said nothing anywhere. A real worker read a 3-byte
{} stub for opencode.json, where the repo's real file is 30+ lines, and truthfully
(but wrongly) reported a mount key did not exist.

Two readers, two fixes:
- the daemon operator gets one info log per provisioning, naming the denominator,
  what was neutralized, and why anything was not (same shape as overlayParity's
  fix in #148 point 3).
- the worker gets the same fact recorded in worktree-scoped git config
  (fleet.neutralizedConfig / fleet.neutralizedConfigNote), discoverable with
  `git config --worktree --get-all fleet.neutralizedConfig` from inside its own
  worktree, without asking the lead. Not a working-tree file: this repo already
  uses worktree-scoped config for the credential helper and the SSH->HTTPS
  rewrite, and it lives under .git/worktrees/<nonce>/ so it can never appear in
  `git status` for the worker to trip on or commit.

The neutralization itself (stub content, --skip-worktree marking) is unchanged.
2026-09-03 11:24:41 +07:00
Dai Ha a89dcc9b7e fleetd #149: seed the workspace-trust entry before a claude-code spawn
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Successful in 1m43s
A claude-code member spawned into a fresh worktree hits an interactive,
un-timed workspace-trust prompt on its first start in a directory it has
never seen. It never reaches its first turn and never mounts the bridge.

Fix: ClaudeCodeLauncher.seedTrustDialog writes
projects.<cwd>.hasTrustDialogAccepted / hasCompletedProjectOnboarding into
the profile's configDir/.claude.json (or ~/.claude.json when configDir is
unset) BEFORE the herdr spawn call, additively (existing keys/projects are
preserved). Gated to isProvisionedWorktree(cwd) - a .git that is a regular
gitdir-pointer file, never a real checkout's .git directory - the same
signal writeIdeOverlay already used, now shared between both.

That gate is a fix for a real incident hit while building this: an
earlier ungated version ran against this file's own pre-existing tests
(configDir=null, no cwd -> falls back to the real user.dir and
~/.claude.json) and corrupted the operator's actual ~/.claude.json down
to a single entry during a mutation-testing run. See PR body for the
full incident report.

FakeHerdr gained onAgentStart(Runnable) so a test can assert the seed
is on disk at the exact instant herdr's agent.start call is reached -
i.e. strictly before the peer process itself would start.
2026-09-03 11:23:06 +07:00
Dai Ha 0d7b4fb026 fleetd #134 + #148 point 3: make the parity overlay say what it did
CI / contract (push) Successful in 51s
CI / build (push) Successful in 1m35s
Both defects lived in one method. overlayParity logged every step at debug, so at the
default level the copy was silent and nobody could tell which overlay files a member
actually got. It also marked a copied tracked file --skip-worktree and said nothing, so a
worker editing that file later found git ignoring the change with no error anywhere.

The summary now reports the denominator, not a bare count: 'copied 1 of 2 candidates:
.env (.envrc absent)'. A bare 'copied 1' is the same under-reporting shape as #113.
Neutralised files are named with the consequence in the message itself.

No marker file is written into the worktree: acceptance criterion 1 requires the worktree
to hold exactly the configured overlay set, so a marker would violate the fix it documents.

The copy and mark logic is unchanged — only logging is new.

Verified by the lead: baseline 1168 tests, 0 failures. Reverting either log.info to
log.debug goes red (2 reds and 1 red, 0 compile errors each), which is the regression that
matters since the whole fix is the log level.
2026-09-03 11:14:21 +07:00
Dai Ha ef8c97871e fleetd #134/#148 point 3: make overlayParity's copy and skip-worktree visible
CI / contract (pull_request) Successful in 1m15s
CI / build (pull_request) Successful in 1m53s
overlayParity logged everything at debug, so at the default level nobody could
tell which overlay files a spawn actually received (#148 pt 3), and a tracked
file marked --skip-worktree gave no warning that it can no longer be edited
from that worktree (#134).

Report the outcome at info: a per-spawn summary naming the denominator (every
configured candidate), what was copied, and why anything was not — plus a
separate line naming every file marked --skip-worktree, stating plainly that
it cannot be committed from this worktree. No worktree-local marker file: the
worktree must hold exactly the configured overlay set and nothing else, so an
extra file would violate that invariant.
2026-09-03 11:09:50 +07:00
30 changed files with 3383 additions and 114 deletions
+20
View File
@@ -188,6 +188,18 @@ must obey belongs in the charter, not here.
- **This repo is the bridge.** The daemon is `fleetd`, its MCP mount is `http://127.0.0.1:8765/mcp`,
and the code behind the rules above is `mcp/FleetMcp` (tools), `auth/Authz` (the role table),
`mcp/ConnectionIdentity` (connection→role), and `worker/*Launcher` (`REPLY_CHARTER`).
- **`fleet_profiles`/`fleet_list` report two separate outage states, and they are not the same
thing.** *Quarantined* (CB-578) means the backend told us it is out of capacity — a long,
1800s-default cooldown. *Cooling off* (fleetd #201/#227) means a profile's credential threw two
distinct backend errors (a non-exhaustion failure such as an HTTP 5xx) within 60 seconds — a
short, fixed 60s cooldown, not configurable per profile. Each check runs independently, so a
profile can show both at once. In the JSON: a cooling profile carries `credentialId` and
`coolingOffForSeconds`; a quarantined profile carries `quarantinedForSeconds`; a profile hit by
both carries all three fields, and either state alone already sets that profile's `free` to `0`.
A `fleet_spawn` naming a cooling-off profile is refused before it ever reaches the backend
adapter, with a message naming the credential and the remaining seconds ("cooling off after
repeated backend errors") — distinct wording from a quarantine refusal, so don't conflate the
two when reading a spawn failure.
- **Skills available to delegate:** `implementer` (worktree → commit → push → own PR) and
`reviewer` (scoped review → one structured finding). Name one in every delegation.
- **Primary-side skills** (not delegation playbooks — a worker cannot use them):
@@ -195,6 +207,14 @@ must obey belongs in the charter, not here.
`fleets-status` (report every fleet that shares one LavinMQ instance).
- **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
committed copies would otherwise mount the primary's IDE and forge servers (fleetd #134). The
worktree's copy of each is a stub, **not** the repo's real file, so a worker that reads one and
reports what it found is reporting on the stub. The daemon logs a per-spawn summary, but the
worker cannot see that log. From inside its own worktree a worker — or a lead debugging one —
reads the list with `git config --worktree --get-all fleet.neutralizedConfig`, and the
consequence with `git config --worktree --get fleet.neutralizedConfigNote`. Never brief a worker
to edit one of these files: the edit cannot be committed, and it will not tell you so.
- **Flows and the error model** — rendezvous, `fleet_ask`, detached delivery, turn-done fallback —
are diagrammed in `docs/MCP-Contract.md` **§6 only**. The rest of that page is a pre-build design
doc whose tool names, parameter names and REST paths never caught up with the code, so do not use
+95 -4
View File
@@ -172,7 +172,11 @@ herdrSocket: ~/.config/herdr/herdr.sock
# skills/MCP/hooks. Omit to leave the worker on the host default.
# parityOverlay → repo-relative paths copied primary→worktree so a worker in a provisioned
# worktree sees the same local config (CB-301-ext). Omit for the default set:
# [.env, .envrc]. (.claude/settings.local.json is NOT in the default — it
# [.env] only (CB-148). .envrc is left out of the default on purpose: it is
# executable shell that direnv runs on every cd, so copying it carries
# behaviour into the worker, not just values, unlike .env. An operator who
# wants it copied can still write parityOverlay: [.env, .envrc] explicitly.
# (.claude/settings.local.json is NOT in the default — it
# pre-approves IDE/tool grants a member must not hold ambiently; CB-525/CB-634.)
#
# Do NOT add .mcp.json (CB-525). A worker's tools are whatever its launcher
@@ -206,6 +210,39 @@ herdrSocket: ~/.config/herdr/herdr.sock
# profile quarantines alone, under its own name, exactly as if the field did
# not exist. Cooldown length is the top-level quarantineCooldownSeconds below.
# HOT: read live at every spawn/exhaustion check — no restart needed.
# errorPattern → fleetd #201 / #227: regex matched against a completion-fallback scrape to
# classify a turn that ended with no fleet_reply as a BACKEND ERROR — a
# credential outage or a provider 5xx — rather than a real answer or a
# usage-limit exhaustion (exhaustedPattern above always wins when a line
# matches both). Opt-in. Omit it and this profile falls back to fleetd's
# built-in legacy pattern `(?i)\bAPI Error\s*:` — classification still
# happens, just without a profile-specific match; every backend words its
# failure differently, so a hardcoded sentence would only ever match one
# of them.
# DEFERRED: compiled once into a startup pattern map, same as exhaustedPattern
# — editing it needs a daemon restart.
# # errorPattern: "503 Service Unavailable" # opt-in: classify a backend outage
#
# What happens once a match fires (BackendOutagePolicy, credentialId-keyed,
# SEPARATE from the CB-578 stage B quarantine above and never merged with it):
# - threshold 2 — TWO DISTINCT TARGETS (never raw events) on the same
# effective credential inside a 60-second window start an "incident" and a
# 60-second cool-off for that credential. One member repeating the same
# classified line twice never cools anything off — a real outage hits
# every target on that credential, so requiring a second, independent
# target loses nothing against the case this guards against, while
# protecting against a heuristic misfire on one flaky member.
# - a fresh error while a credential is already cooling off is ignored
# outright: it neither extends the 60s deadline nor starts a new incident.
# - `fleet_list`/`fleet_profiles` report a cooling credential with
# `coolingOffForSeconds` (never `quarantinedForSeconds`, unless CB-578
# exhaustion quarantine is ALSO independently active for the same
# credential — the two checks can both fire at once). A spawn onto a
# cooling profile is refused with a message naming the credential and
# remaining seconds — "cooling off", never "exhausted", so an operator can
# tell a short transient fault from a spent subscription at a glance.
# - the lead gets ONE nudge per incident (not one per affected target), via
# the same push loop that already delivers ticket/question reminders.
# env → extra environment for this profile's workers, as a literal key/value map
# (CB-511). Use it to give workers a toolchain.
#
@@ -269,13 +306,43 @@ profiles:
# GOTCHA 2 — `maxLoad` is the ONLY throttle you have here. There is no metering, no budget
# and no refusal on cost; the cap on live members is the single thing standing between a
# fan-out and your monthly limit. Set it deliberately and keep it small.
#
# GOTCHA 3 (fleetd #176) — `maxLoad` counts members, never the lead itself. The lead is a live
# `claude` session on this SAME account (a lead is never moved off-subscription, whatever its
# own profile says), so it already holds one seat before any member spawns. If a lead's
# `fleet.leaders.<name>.profile` names THIS profile — or ANY OTHER `subscription: true`
# profile that shares this one's account (see THE SENTINEL, just below, next to
# `credentialId:`) — `fleet_list`'s `free` for this profile subtracts that lead's live
# seat(s) automatically; see `profile:` under THE FLEET below. If no lead entry names a
# profile sharing this account, fleetd has no way to know a lead holds a seat here, and `free`
# will overstate what a fresh `fleet_spawn` actually gets by exactly the seats the lead is
# quietly holding.
#
# THE SENTINEL (fleetd #176 stage 2, correcting an inert stage 1 fix): every `subscription:
# true` profile that leaves `credentialId` unset shares ONE implicit account-wide credential
# id with every other such profile on this host — because a subscription profile doesn't
# authenticate with a credential of its own, it authenticates as the operator's own Claude
# login, and there is exactly one of those. So on a typical host, `opus` (the lead's profile)
# and `sonnet` (the members' profile) are linked automatically, with NOTHING to set here — that
# is what makes GOTCHA 3 above work without also writing matching `credentialId:` values on
# both. This linkage is not just cosmetic: it is the same key `BackendQuarantine`/cool-off use,
# so a usage-limit hit on `opus` now quarantines `sonnet` too (and vice versa) — correct, since
# they are one Claude account, but worth knowing before you wonder why an unrelated-looking
# profile went quarantined.
#
# WHEN TO OVERRIDE — set explicit, DIFFERENT `credentialId:` values on two `subscription: true`
# profiles only when they are genuinely two separate Claude logins on the same host (a real,
# if unusual, setup). An explicit `credentialId` always wins over the sentinel, so this is the
# one way to keep two subscription profiles from being treated as one account for lead-seat
# counting AND for quarantine/cool-off grouping alike.
# gitTokenEnv: GITEA_TOKEN # opt-in: let this profile's workers open their own PR (CB-302)
# gitHostEnv: GITEA_HOST # defaults to GITEA_HOST; injected only with gitTokenEnv
# exhaustedPattern: "usage limit has been reached" # opt-in: classify a usage-limit refusal (CB-578)
# credentialId: shared-openai # opt-in: quarantine together with every other profile sharing this id (CB-578)
# errorPattern: "503 Service Unavailable" # opt-in: classify a backend outage (fleetd #201/#227) — see the key doc above
# configDir: /Users/me/.ccs/instances/gx10 # CLAUDE_CONFIG_DIR — inherit that profile's skills/MCP
# cwd: /Users/me/src/myrepo # pin the working dir; omit to inherit the primary's
# parityOverlay: [".env", ".envrc"] # the default; never add .mcp.json or .claude/settings.local.json — see above
# parityOverlay: [".env"] # the default; add ".envrc" explicitly if you want it copied too (CB-148) — never add .mcp.json or .claude/settings.local.json — see above
# ideMcpUrl: http://127.0.0.1:29170/index-mcp/streamable-http # opt-in (CB-634): IDE code intelligence, pinned to the worktree
# ideProjectDir: fleetd # CB-634: module dir the IDE opens + the overlay pins (this repo's pom is in fleetd/)
# ideOpenCommand: env DISPLAY=:10.0 idea {dir} # CB-634 auto-open: opens {dir} in the IDE at spawn; omit to open by hand
@@ -369,6 +436,12 @@ placement: weighted
# DEFERRED: baked once into the BackendQuarantine built at startup — a running quarantine keeps
# its original cooldown regardless; a new value only applies to a quarantine that starts after a
# restart. Editing this needs a daemon restart to take effect.
#
# This does NOT govern the fleetd #201 / #227 backend-error cool-off documented under errorPattern
# above — that mechanism is a separate, shorter-lived, NOT-configurable policy (threshold 2 distinct
# targets, 60-second window, 60-second cool-off), on purpose: it exists to survive a brief transient
# fault, not to replace this 30-minute exhaustion quarantine. Do not conflate the two when reading
# fleet_list/fleet_profiles — coolingOffForSeconds and quarantinedForSeconds are independent facts.
# quarantineCooldownSeconds: 1800
# Re-read this file without restarting the daemon (CB-559). Off unless you add this block, so an
@@ -395,8 +468,10 @@ placement: weighted
# stage B — baked once into the quarantine tracker built at startup), ADDING or
# REMOVING a profile (a new backend needs its own launcher, and launchers are built
# once), AND an existing profile's launch settings — model, baseUrl, argv, env,
# configDir, mcpUrl, tabLabel, exhaustedPattern. The launcher takes a copy of
# `profiles:` at startup and resolves every spawn out of that copy, so those never
# configDir, mcpUrl, tabLabel, exhaustedPattern, errorPattern (fleetd #201 / #227 —
# compiled once into a startup pattern map the same way exhaustedPattern is). The
# launcher takes a copy of `profiles:` at startup and resolves every spawn out of
# that copy, so those never
# reach a launch until you restart. The reload logs them by name rather than
# pretending they applied.
# COLD → cannot change at all: `bind:`, `herdrSocket:`, `broker:` and `auth:`. The socket is
@@ -457,6 +532,22 @@ fleet:
# recognised: give it a `profile:` and the daemon launches the shortfall when fewer than
# `instances` are live. Omit `profile:` and it is recognise-only, as before.
#
# `profile:` has a SECOND job as of fleetd #176, even for a recognise-only lead you never want
# auto-launched: it is also how fleetd learns which account this lead's own session shares. A
# `subscription: true` profile bills the operator's Claude account, and the lead itself is always
# a live `claude` session on that same account — `maxLoad` never counted that seat. If a lead
# entry here names a profile that shares a worker profile's account, `fleet_list`'s `free` for
# that worker profile subtracts the lead's live seat(s) automatically. "Shares the account" is
# decided by matching `effectiveCredentialId()`, which (fleetd #176 stage 2 — see THE SENTINEL,
# next to `credentialId:`, in THE WORKERS above) means: an explicit, matching `credentialId:` on
# both, OR — the common case, needing NO extra config — both being `subscription: true` with
# `credentialId` left unset, since those all share one implicit account-wide id. A lead on `opus`
# and workers on `sonnet` link automatically this way; they do NOT need the same profile name.
# Setting `profile:` on an already-running, recognise-only lead is safe — the daemon only launches
# the SHORTFALL below `instances`, so naming a profile here does not, by itself, start anything.
# Omit it and fleetd has no way to derive the sharing — there is no other reliable signal on the
# daemon's side — so that lead's seat goes uncounted, exactly as before this ticket.
#
# `tab:` (CB-579) is REQUIRED and is the only field identity depends on — the exact label of the
# tab hosting the lead, matched case-insensitively. Label the tab yourself and put that same
# string here, and the pane is recognised on the next rescan. Reopen the tab later, or the session
+245 -9
View File
@@ -13,6 +13,8 @@ import dev.ltms.fleet.lead.LeadLauncher;
import dev.ltms.fleet.herdr.PaneLocator;
import dev.ltms.fleet.herdr.UnixSocketHerdrClient;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.BackendErrorPatternLookup;
import dev.ltms.fleet.inject.BackendErrorSink;
import dev.ltms.fleet.inject.CompletionResolver;
import dev.ltms.fleet.inject.ExhaustedPatternLookup;
import dev.ltms.fleet.inject.ExhaustionSink;
@@ -50,6 +52,7 @@ import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.member.CompositePeerLauncher;
import dev.ltms.fleet.member.HerdrPeerLauncher;
import dev.ltms.fleet.member.OpenCodeLauncher;
import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
import io.javalin.Javalin;
import org.slf4j.Logger;
@@ -61,6 +64,7 @@ import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
@@ -205,12 +209,18 @@ public final class Fleetd {
// here, at startup, and a config reload only changes it for a daemon restart.
BackendQuarantine quarantine = new BackendQuarantine(System::nanoTime,
TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));
// fleetd #201 Unit 5: one outage-cool-off tracker for the whole daemon, shared between the
// launcher (checked at spawn, like `quarantine` above) and the backend-error sink wired in
// below (written on a classified backend error). A SEPARATE, shorter-lived mechanism from
// `quarantine` — see BackendOutagePolicy's class doc — never merged with it.
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(System::nanoTime);
PeerLauncher workers = new CompositePeerLauncher(
adapters,
cfg.effectiveDefaultProfile(),
config,
profileName -> liveCountRef.get().apply(profileName),
quarantine);
quarantine,
outagePolicy);
// CB-504: under supervision (launchd/systemd) fleetd can start before herdr's socket
// exists. The client itself is lazy — it connects per call — but the orphan reap below is
// the first thing that actually talks to herdr, so without this wait a boot-order race
@@ -341,7 +351,29 @@ public final class Fleetd {
.map(session -> exhaustedPatternsByProfile.get(session.profile()))
.orElse(null);
log.info("backend-exhausted classification (CB-578 stage A): {}",
CompletionResolver.coverage(cfg.profiles().keySet(), exhaustedPatternsByProfile.keySet()));
CompletionResolver.coverage("exhaustedPattern", cfg.profiles().keySet(),
exhaustedPatternsByProfile.keySet()));
// fleetd #201 Unit 5: classify a completion-fallback scrape that matches a profile's
// configured backend-error refusal (a credential outage, a provider 5xx) as a backend error
// rather than handing it back as a real answer. Compiled once at startup, keyed by profile
// name, mirroring exhaustedPatternsByProfile above — a profile with no configured
// errorPattern is simply absent here, so CompletionResolver falls back to its built-in
// narrow {@code (?i)\bAPI Error\s*:} compatibility pattern for that profile's targets
// (BackendErrorPatternLookup#legacy's contract — see backendErrorPatterns below).
Map<String, Pattern> errorPatternsByProfile = new LinkedHashMap<>();
cfg.profiles().forEach((name, profile) -> {
if (profile.hasErrorPattern()) {
errorPatternsByProfile.put(name, Pattern.compile(profile.errorPattern()));
}
});
// fleetd #248: extracted to a static factory (see backendErrorPatternLookup below) so a
// test can prove main() actually PASSES this into CompletionResolver, not only that the
// lookup itself behaves correctly — the exact gap fleetd #248 exists to close.
BackendErrorPatternLookup backendErrorPatterns = backendErrorPatternLookup(sessions::roster,
errorPatternsByProfile);
log.info("backend-error classification (fleetd #201 Unit 5): {}",
CompletionResolver.coverage("errorPattern", cfg.profiles().keySet(),
errorPatternsByProfile.keySet()));
// CB-578 stage B: on a classification that actually wins, quarantine the exhausted profile's
// CREDENTIAL — not the profile name — so a profile sharing that credential (e.g. two models
// on one OpenAI account) is refused too, not just the one that happened to report it. Reads
@@ -389,13 +421,30 @@ public final class Fleetd {
// fleetd #175: point the forwarding sink handed to OpenCodeLauncher above at the real one,
// now that `sessions` exists to resolve target -> session -> profile.
exhaustionSinkRef.set(exhaustionSink);
// fleetd #201 Unit 5: the production BackendErrorSink needs `pushLoop` (built further below,
// after `sessions`) to tell a lead about an incident or an unmapped target — the same
// construction-order cycle `exhaustionSinkRef` breaks above, broken the same way: a mutable
// holder set once `pushLoop` exists, read lazily from inside the lambda built here.
AtomicReference<ReplyPushLoop> pushLoopRef = new AtomicReference<>();
// fleetd #248: extracted to a static factory (see backendErrorSink below), public rather
// than package-private like the other two factories here, so
// dev.ltms.fleet.inject.BackendOutageFlowTest can exercise the REAL production sink
// directly instead of a hand-mirrored copy of this lambda — that copy was precisely the
// gap fleetd #248 exists to close (see that test's class doc for the history).
BackendErrorSink backendErrorSink = backendErrorSink(sessions, () -> config.get().profiles(),
outagePolicy, pushLoopRef::get);
AgentControl agents = router.memberAgents();
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns, exhaustionSink,
target -> sessions.roster().stream()
.filter(session -> target.equals(session.terminalId()))
.findFirst()
.map(session -> new CompletionResolver.WorktreeBranch(session.worktree(), session.branch()))
.orElse(null));
// Both fleetd#201 Unit 5 (backend-error patterns + sink) and fleetd#241 (the worktree/branch
// lookup the fallback report names) land on this one call. The full constructor takes both,
// so neither feature is dropped; nowNanos must be passed explicitly to reach it.
//
// fleetd #248: every argument built specifically for this call (backendErrorPatterns and
// backendErrorSink above, and the worktree/branch lookup right here) now comes from a
// static factory tested on its own; FleetdCompletionResolverWiringTest source-asserts that
// THIS call actually passes them, which is the coverage that was missing before.
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns,
exhaustionSink, backendErrorPatterns, backendErrorSink, System::nanoTime,
worktreeBranchLookup(sessions::roster));
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
MemberPresence presence = sessions.asPresence();
@@ -477,6 +526,9 @@ public final class Fleetd {
Metrics metrics = FleetMetrics.create(sessions, replyInbox);
var pushLoop = new ReplyPushLoop(primaryRegistry, router.leadAgents(), replyInbox,
pushScheduler, maxReminders, backoffMs, metrics);
// fleetd #201 Unit 5: point the forwarding holder captured by the backendErrorSink lambda
// above at the real push loop, now that it exists.
pushLoopRef.set(pushLoop);
// CB-551: idle-lead heartbeat. Opt-in; absent `leadHeartbeat:` this is never constructed, so
// an upgraded daemon cannot silently start spending subscription on nudging an idle lead.
// It has its own single-thread scheduler and holds its own scheduler shutdown via close().
@@ -584,7 +636,12 @@ public final class Fleetd {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.effectiveCredentialId();
}, quarantine),
leadMailbox);
leadMailbox,
new FleetMcp.OutageSource(profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.effectiveCredentialId();
}, outagePolicy),
new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)));
// 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
@@ -682,6 +739,185 @@ public final class Fleetd {
return target -> presence.isPresent(target) || leads.get().containsKey(target);
}
/**
* fleetd #248: package-private factory for the member worktree/branch lookup {@link
* CompletionResolver} uses to name a fallback report's worktree and branch (fleetd#241).
*
* <p>Before this ticket the lookup was an anonymous lambda built inline inside {@code main}'s
* {@code CompletionResolver} constructor call — provably untested wiring, the whole reason
* fleetd #248 exists: dropping that one argument (passing {@code _ -> null} instead) compiled
* clean and left every test green. Extracted here, {@code main} now calls this factory instead
* of building the lambda inline, and a source assertion on that call site
* ({@code FleetdCompletionResolverWiringTest}) proves the argument is still actually passed.
*
* <p>Takes the roster as a plain {@link Supplier} — not a {@link SessionManager} — so this is
* directly testable with a hand-built session list; no real {@code SessionManager} (launcher,
* worktrees, …) needs constructing. Follows the same {@code static} factory pattern as
* {@link #deliverableTo} above.
*
* @param roster the live member roster, normally {@code sessions::roster}
*/
static Function<String, CompletionResolver.WorktreeBranch> worktreeBranchLookup(
Supplier<List<MemberSession>> roster) {
return target -> roster.get().stream()
.filter(session -> target.equals(session.terminalId()))
.findFirst()
.map(session -> new CompletionResolver.WorktreeBranch(session.worktree(), session.branch()))
.orElse(null);
}
/**
* fleetd #176: per-profile factory for {@link FleetMcp.LeadSeatSource} — how many seats a
* profile's own live LEAD session(s) hold on the same Claude subscription.
*
* <p>{@code maxLoad} counts only members; the lead itself is a live {@code claude} session that
* is never moved off-subscription ({@code LeadLauncher} strips {@code ANTHROPIC_BASE_URL}/
* {@code AUTH_TOKEN} from a lead's env whatever its profile says), so a {@code subscription:
* true} profile's real ceiling is lower than its configured {@code maxLoad} by exactly the
* number of lead seats sharing that same account.
*
* <p><b>The derivation, and why this route was chosen over a new config key.</b> The link is
* {@code fleet.leaders.<name>.profile} — the field the operator already sets to name which
* {@code profiles:} entry a lead runs on (CB-557; see {@code fleetd.example.yaml}) — matched
* against the profile passed in here via {@link FleetConfig.Profile#effectiveCredentialId()},
* the same grouping key {@link dev.ltms.fleet.placement.BackendQuarantine} already uses to say
* two profiles share one account. Nothing new is added to the config schema: this reuses a field
* that already exists and already means "the profile this lead's own session runs on". A lead
* entry that names no {@code profile:} (recognise-only, CB-558) says nothing about which account
* it shares, and there is no other reliable signal on the daemon's side to derive that from — so
* such a lead contributes no seats, exactly as before this ticket. Making that lead's seat count
* requires the operator to add one line (`profile: sonnet` under its {@code fleet.leaders} entry)
* — a config statement, not a code change, and the smallest one available given the field
* already exists for a closely related purpose.
*
* <p>Only counts leads {@code liveLeadTerminals} currently reports — CB-531's live tab scan (or
* the legacy {@code primary.terminal} pin) — never every configured lead: an entry whose
* {@code instances} nobody has actually started is not really competing for a seat, and must not
* shrink capacity for one that was never live.
*
* @param profiles the live profile map, normally {@code () -> config.get().profiles()}
* in {@code main} — hot, like every other {@code maxLoad}/
* {@code credentialId} read {@link FleetMcp.CapacitySource} already does
* @param leaders {@code fleet.leaders}, read once at startup like the rest of that
* block ({@code fleetd.example.yaml} notes it is not hot) — passed as a
* plain map, never re-read from {@code config.get()}
* @param liveLeadTerminals terminal_id → lead name for every CURRENTLY recognised lead, normally
* the same supplier {@link dev.ltms.fleet.auth.CallerResolver#leads()}
* and {@code LeadCoordLoop} already consult
*/
static Function<String, Integer> leadSeatLookup(Supplier<Map<String, FleetConfig.Profile>> profiles,
Map<String, FleetConfig.Leader> leaders, Supplier<Map<String, String>> liveLeadTerminals) {
return profileName -> {
FleetConfig.Profile target = profiles.get().get(profileName);
if (target == null || !target.isSubscription()) {
return 0;
}
String targetCredential = target.effectiveCredentialId();
int seats = 0;
for (String leadName : liveLeadTerminals.get().values()) {
FleetConfig.Leader lead = leaders.get(leadName);
if (lead == null || lead.profile() == null || lead.profile().isBlank()) {
continue;
}
FleetConfig.Profile leadProfile = profiles.get().get(lead.profile());
if (leadProfile != null && targetCredential.equals(leadProfile.effectiveCredentialId())) {
seats++;
}
}
return seats;
};
}
/**
* fleetd #248 / fleetd#201 Unit 5: package-private factory for the per-target backend-error
* pattern lookup {@link CompletionResolver} classifies a pane scrape against. Closes over the
* live roster (to resolve a target to a profile) and {@code errorPatternsByProfile} (each
* profile's configured {@code errorPattern}, already compiled by the caller — the same map also
* feeds the coverage log next to where this is called) — nothing else, so it is directly
* testable. See {@link #worktreeBranchLookup} above for why this ticket exists and why the
* factory takes a roster {@link Supplier} rather than a {@link SessionManager}.
*
* @param roster the live member roster, normally {@code sessions::roster}
* @param errorPatternsByProfile every profile that has an {@code errorPattern} configured,
* keyed by profile name
*/
static BackendErrorPatternLookup backendErrorPatternLookup(Supplier<List<MemberSession>> roster,
Map<String, Pattern> errorPatternsByProfile) {
return target -> roster.get().stream()
.filter(session -> target.equals(session.terminalId()))
.findFirst()
.map(session -> errorPatternsByProfile.get(session.profile()))
.orElse(null);
}
/**
* fleetd #248 / fleetd#201 Unit 5: factory for the production {@link BackendErrorSink} — the
* collaborator {@link CompletionResolver} notifies when a pane-scrape classification actually
* resolves a waiter as a backend error. Order: (1) mark the member BACKEND_ERROR; (2) resolve
* profile/credential through the roster — fail loud (never {@code Optional.ifPresent}, the
* fleetd #234 lesson applied to this sink) and notify the lead via {@code
* onBackendTargetUnmapped} when it cannot be resolved; (3) record the error in {@code
* outagePolicy}; (4) on a NEW incident (the {@code record()} call that actually crosses the
* threshold), tell the lead via {@code onBackendIncident}.
*
* <p>{@code public}, unlike {@link #worktreeBranchLookup} and {@link #backendErrorPatternLookup}
* above: {@code dev.ltms.fleet.inject.BackendOutageFlowTest} exercises this exact object as the
* real, wired production path, replacing what its own class doc used to call out as a
* hand-mirrored copy of this lambda ("mirrors {@code Fleetd.main}'s {@code backendErrorSink}
* lambda line-for-line") — that copy proved only itself, never that {@code main} still wires the
* real thing. That was precisely the gap fleetd #248 exists to close.
*
* @param sessions the session registry; both read (roster) and written (onBackendError)
* @param profiles the live profile map, normally {@code () -> config.get().profiles()} in
* {@code main}, or a fixed test map via {@code () -> profiles}
* @param pushLoop the lead-nudge loop, read lazily: {@code main} builds this sink before the
* real {@link ReplyPushLoop} exists (a genuine construction-order cycle, broken
* the same way {@code exhaustionSinkRef} is a few lines above it), so a
* {@link Supplier} reads whatever {@code main} has filled in by the time a real
* backend error fires
*/
public static BackendErrorSink backendErrorSink(SessionManager sessions,
Supplier<Map<String, FleetConfig.Profile>> profiles, BackendOutagePolicy outagePolicy,
Supplier<ReplyPushLoop> pushLoop) {
return (target, matchedLine, reason) -> {
sessions.onBackendError(target, reason);
String profileName = sessions.roster().stream()
.filter(session -> target.equals(session.terminalId()))
.findFirst()
.map(MemberSession::profile)
.orElse(null);
FleetConfig.Profile profile = profileName == null ? null : profiles.get().get(profileName);
if (profile == null) {
log.error("backend error on target '{}' ({}) but no profile could be resolved — the "
+ "target is not (yet) in the roster — no cool-off applied (fleetd #201 Unit 5)",
target, reason);
ReplyPushLoop loop = pushLoop.get();
if (loop != null) {
loop.onBackendTargetUnmapped(target, reason);
}
return;
}
String credentialId = profile.effectiveCredentialId();
Optional<BackendOutagePolicy.Incident> incident = outagePolicy.record(credentialId, target, reason);
incident.ifPresent(inc -> {
List<String> affectedProfiles = profiles.get().values().stream()
.filter(p -> credentialId.equals(p.effectiveCredentialId()))
.map(FleetConfig.Profile::profile)
.sorted()
.toList();
log.warn("credential '{}' cooling off for {}s after backend errors on {} distinct "
+ "target(s) (profile '{}'): {}", credentialId,
inc.remainingCoolOffSeconds(), inc.evidenceCount(), profile.profile(), reason);
ReplyPushLoop loop = pushLoop.get();
if (loop != null) {
loop.onBackendIncident(inc.id(), inc.targets(), credentialId, affectedProfiles,
(int) inc.remainingCoolOffSeconds());
}
});
};
}
/** Injection seam for {@link #selectReplyInbox}: production binds {@link AmqpReplyInbox#open}. */
@FunctionalInterface
interface AmqpOpener {
@@ -43,7 +43,9 @@ import java.util.function.Supplier;
* which is constructed once), <em>and an existing profile's launch settings</em> —
* {@code model}, {@code baseUrl}, {@code argv}, {@code env}, {@code mcpUrl},
* {@code exhaustedPattern} (CB-578 stage A — compiled once into {@code Fleetd.main}'s
* pattern map at startup), and the rest. {@code credentialId} (CB-578 stage B) is NOT on
* pattern map at startup), {@code errorPattern} (fleetd #201 Unit 5 — compiled once into
* {@code Fleetd.main}'s backend-error pattern map at startup, the same way), and the rest.
* {@code credentialId} (CB-578 stage B) is NOT on
* this list — it is read live off the config supplier at every quarantine check and
* exhaustion event, exactly like {@code weight} / {@code maxLoad}, so it is hot instead.
* {@code HerdrPeerLauncher} takes {@code Map.copyOf(profiles)} at construction and resolves
@@ -292,6 +294,10 @@ public final class ConfigRef implements Supplier<FleetConfig> {
// CB-578 stage B: exhaustedPattern is compiled once into Fleetd.main's pattern map
// at startup (see ExhaustedPatternLookup wiring) — a reload never re-reads it, so a
// changed pattern must be reported as deferred, exactly like model/baseUrl/argv.
&& Objects.equals(a.exhaustedPattern(), b.exhaustedPattern());
&& Objects.equals(a.exhaustedPattern(), b.exhaustedPattern())
// fleetd #201 Unit 5: errorPattern is compiled once into Fleetd.main's backend-error
// pattern map at startup (see BackendErrorPatternLookup wiring), the same way
// exhaustedPattern is — a reload never re-reads it either.
&& Objects.equals(a.errorPattern(), b.errorPattern());
}
}
@@ -24,6 +24,8 @@ import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.regex.Pattern;
import java.util.regex.PatternSyntaxException;
/**
* {@code fleetd} configuration, loaded from a YAML file (see
@@ -257,7 +259,13 @@ public record FleetConfig(
* @param cwd fixed working directory for this profile's workers (CB-112 "told otherwise");
* {@code null}/blank → inherit the primary's cwd, else the daemon's
* @param parityOverlay repo-relative paths copied primary→worktree for config parity; null/empty
* defaults to a sensible set of local config files.
* defaults to {@code [.env]} only (CB-148 point 2). {@code .env} is data — a
* copy of it can only carry values. {@code .envrc} is executable shell that
* {@code direnv} evaluates on every {@code cd}, so copying it moves behaviour
* into the worker, not just values, and that is a different risk from copying
* data. It is deliberately left out of the default for that reason. The knob
* still supports it: an operator who wants it copied writes
* {@code parityOverlay: [.env, .envrc]} explicitly and owns that choice.
* <p><strong>Every overlay path must be gitignored or tracked-and-skipped.</strong>
* CB-576 made {@code release()} preserve a worktree that {@code git status
* --porcelain} reports as dirty, and it deliberately counts untracked files —
@@ -267,11 +275,10 @@ public record FleetConfig(
* worktrees then accumulate with no error anywhere.
* <p>Checked on 2026-08-15 (CB-581): inert as configured. Tracked overlay
* files carry {@code --skip-worktree} so {@code --porcelain} cannot see them,
* {@code fleetd.yaml} is gitignored, and the default pair {@code .env} /
* {@code .envrc} does not exist in this repo. Note the default applies to
* <em>every</em> profile, so creating either file at the repo root is enough
* to make it live. Add a new overlay path to {@code .gitignore} in the same
* change that adds it here.
* {@code fleetd.yaml} is gitignored, and {@code .env} does not exist in this
* repo. Note the default applies to <em>every</em> profile, so creating
* {@code .env} at the repo root is enough to make it live. Add a new overlay
* path to {@code .gitignore} in the same change that adds it here.
* <p>CB-578 stage C raised the stakes: a preserved dirty worktree is now also
* committed to {@code refs/wip/<branch>} via {@code git add -A}. The overlay
* carries the primary's own environment files into the worktree, so a
@@ -356,6 +363,16 @@ public record FleetConfig(
* {@code limit.context} instead, which bounds the window opencode compacts
* <em>within</em>, and only when {@code model} resolves to a
* {@code provider/model} pair.
* @param errorPattern regex matched against a completion-fallback scrape (fleetd #201 Unit 5)
* to classify a turn that ended with no {@code fleet_reply} as a backend
* error (credential outage, provider 5xx) rather than a real answer or an
* exhausted usage limit. {@code null}/blank ⇒ this profile relies on
* {@link dev.ltms.fleet.inject.CompletionResolver}'s built-in
* {@code (?i)\bAPI Error\s*:} compatibility pattern instead — classification
* still happens, just without a profile-specific match (reported separately
* at startup as legacy-default coverage, never full coverage). Compiled once
* at startup ({@code Fleetd.main}), like {@code exhaustedPattern}, so a
* reload of this key is DEFERRED (see {@code ConfigRef}).
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Profile(String profile, String baseUrl, String model,
@@ -374,7 +391,8 @@ public record FleetConfig(
String ideMcpUrl,
String ideProjectDir,
String ideOpenCommand,
Integer autoCompactWindow) {
Integer autoCompactWindow,
String errorPattern) {
/** Peer kind spawned by {@link dev.ltms.fleet.member.ClaudeCodeLauncher} (the default). */
public static final String KIND_CLAUDE_CODE = "claude-code";
@@ -408,7 +426,7 @@ public record FleetConfig(
// servers do not exist there to be enabled. This is defence in depth, not a live fix: the
// worker stays isolated only because this separate mechanism already removes the servers.
parityOverlay = (parityOverlay == null || parityOverlay.isEmpty())
? List.of(".env", ".envrc")
? List.of(".env")
: List.copyOf(parityOverlay);
// gitTokenEnv stays null when unset (opt-in). gitHostEnv defaults so operators enabling
// checkpoints need only set gitTokenEnv; it is injected only alongside a resolved token.
@@ -448,6 +466,32 @@ public record FleetConfig(
// there is no blank-string form to normalize (it's an Integer), and the [100000, 1000000]
// range is enforced eagerly at config load (rejectAutoCompactWindowOutOfRange), naming the
// profile, rather than silently clamped here. A profile that never sets it keeps null.
// errorPattern stays null when unset/blank (opt-in) — same rule as exhaustedPattern: no
// defaulting, no vendor wording. A profile that never sets it relies on
// CompletionResolver's built-in compatibility pattern instead (never "off").
errorPattern = (errorPattern == null || errorPattern.isBlank()) ? null : errorPattern;
}
/**
* Backward-compatible constructor without the fleetd #201 {@code errorPattern} field — the
* profile relies on {@code CompletionResolver}'s built-in compatibility pattern instead
* (legacy-default coverage). This is the shape the canonical constructor had before the
* field was added; every pre-existing Java call site (and any YAML that omits the key) keeps
* compiling and behaving identically. Jackson binds the canonical (longest) constructor, so
* YAML omitting {@code errorPattern:} still lands here as {@code null} via that path too.
*/
public Profile(String profile, String baseUrl, String model,
String configDir, String tokenEnv, List<String> argv,
String placement, String workspace, String tabLabel, String mcpUrl,
String cwd, List<String> parityOverlay, String gitTokenEnv, String gitHostEnv,
String kind, Map<String, String> env, Float weight, Integer maxLoad,
Boolean subscription, String exhaustedPattern, String credentialId,
String ideMcpUrl, String ideProjectDir, String ideOpenCommand,
Integer autoCompactWindow) {
this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad,
subscription, exhaustedPattern, credentialId, ideMcpUrl, ideProjectDir, ideOpenCommand,
autoCompactWindow, null);
}
/**
@@ -492,7 +536,8 @@ public record FleetConfig(
public Profile withProfile(String p) {
return new Profile(p, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, subscription,
exhaustedPattern, credentialId, ideMcpUrl, ideProjectDir, ideOpenCommand, autoCompactWindow);
exhaustedPattern, credentialId, ideMcpUrl, ideProjectDir, ideOpenCommand, autoCompactWindow,
errorPattern);
}
/**
@@ -590,14 +635,55 @@ public record FleetConfig(
}
/**
* The credential group this profile quarantines with (CB-578 stage B): the configured
* {@link #credentialId} when set, else this profile's own name — so an unconfigured profile
* quarantines alone, exactly as it did before this field existed. Two profiles that set the
* same non-blank {@code credentialId} share one quarantine: a {@code BACKEND_EXHAUSTED}
* classification on either one quarantines both.
* True when this profile configures its own fleetd #201 Unit 5 backend-error classification
* pattern. {@code false} means this profile relies on {@code CompletionResolver}'s built-in
* {@code (?i)\bAPI Error\s*:} compatibility pattern instead (legacy-default coverage).
*/
public boolean hasErrorPattern() {
return errorPattern != null;
}
/**
* The shared credential id every {@code subscription: true} profile falls back to when it
* sets no explicit {@link #credentialId} (fleetd #176 stage 2, correcting an inert first cut
* of that ticket). A subscription profile has no credential of its own to fall back to its
* name for: it authenticates as the operator's own Claude login, and a host has exactly one
* of those, whatever names the operator gives the profiles running on it. Falling back to the
* profile's own name (the way an ordinary off-subscription profile does) would keep two
* subscription profiles on one login apart from each other, which is the opposite of what
* "one account" means.
*
* <p>Measured live and what it broke: a lead on profile {@code opus}, members on profile
* {@code sonnet}, same Claude login, neither setting {@code credentialId}. Before this
* sentinel, {@code opus.effectiveCredentialId()} was {@code "opus"} and {@code sonnet
* .effectiveCredentialId()} was {@code "sonnet"} — so fleetd #176's lead-seat matcher (and,
* this sentinel now also fixes, {@code CompositePeerLauncher}'s quarantine/cool-off spawn
* refusal and {@code BackendOutagePolicy}'s incident grouping) silently never linked them: the
* fix shipped, and stayed inert on the one host it was written for.
*/
public static final String SUBSCRIPTION_CREDENTIAL_ID = "<subscription>";
/**
* The credential group this profile quarantines with (CB-578 stage B; extended fleetd #176
* stage 2 — see {@link #SUBSCRIPTION_CREDENTIAL_ID}): the configured {@link #credentialId}
* when set — that always wins, so an operator with two separate Claude logins on one host can
* still keep them apart. Otherwise, a {@code subscription: true} profile falls back to
* {@link #SUBSCRIPTION_CREDENTIAL_ID} rather than its own name; an ordinary off-subscription
* profile falls back to its own name, exactly as it did before this field existed, so an
* unconfigured off-subscription profile still quarantines alone.
*
* <p>A {@code BACKEND_EXHAUSTED} (or repeated backend-error) classification on any profile
* sharing the result quarantines/cools off every profile that shares it — including, now,
* every {@code subscription: true} profile with no explicit {@code credentialId}. That is
* intended, not incidental: one Claude subscription hitting a usage limit really does take out
* every profile running on it, the same way {@code credentialId: openai-shared} already lets
* two OpenAI-backed profiles share one quarantine.
*/
public String effectiveCredentialId() {
return (credentialId == null || credentialId.isBlank()) ? profile : credentialId;
if (credentialId != null && !credentialId.isBlank()) {
return credentialId;
}
return isSubscription() ? SUBSCRIPTION_CREDENTIAL_ID : profile;
}
/** True when this profile's workers are granted a forge token to open their own PR (CB-302). */
@@ -892,7 +978,15 @@ public record FleetConfig(
* survives restarts of the agent inside it — so identity is now the tab label alone.
*
* @param profile the {@code profiles:} entry to launch this lead on when one must
* be created; {@code null} ⇒ recognise-only, never create
* be created; {@code null} ⇒ recognise-only, never create.
* <p>fleetd #176: also the field {@code Fleetd.leadSeatLookup} reads
* to learn which account this lead's own live session shares — set it
* (safely, even on an already-running recognise-only lead: naming a
* profile here never starts anything beyond {@code instances}) so a
* {@code subscription: true} worker profile sharing its
* {@code effectiveCredentialId()} has this lead's seat subtracted from
* {@code fleet_list}'s {@code free}. {@code null} here also means this
* lead's seat cannot be derived and is not counted.
* @param tab the exact tab label hosting this lead, matched case-insensitively;
* the only field identity depends on. Required — a lead with no
* {@code tab} can never be discovered, launched or not
@@ -1389,6 +1483,7 @@ public record FleetConfig(
rejectDuplicateMemberSlots(yaml);
rejectNegativeMaxLoad(yaml);
rejectAutoCompactWindowOutOfRange(yaml);
rejectMalformedErrorPattern(yaml);
rejectUnknownKind(yaml);
rejectUnknownAuthMode(yaml);
rejectUnknownPlacement(yaml);
@@ -1737,6 +1832,51 @@ public record FleetConfig(
}
}
/**
* Reject a profile whose {@code errorPattern} (fleetd #201 Unit 5) is not a valid Java regex,
* naming the profile, the key, and the parser's own message.
*
* <p>Unset/{@code null} means "use {@code CompletionResolver}'s built-in {@code (?i)\bAPI
* Error\s*:} compatibility pattern" and passes silently. A profile that DOES set the key gets it
* compiled once at daemon startup ({@code Fleetd.main}, mirroring {@code exhaustedPattern}) — an
* uncaught {@link java.util.regex.PatternSyntaxException} there crashes startup without naming
* which profile or key is at fault. Validate eagerly here instead, at config load, the same
* "fail loud at load, not lazily later" reasoning as {@link #rejectAutoCompactWindowOutOfRange}.
*
* @param yaml the raw config text
* @throws IllegalStateException when any profile's {@code errorPattern} fails to compile
*/
static void rejectMalformedErrorPattern(String yaml) {
Map<?, ?> raw;
try {
raw = YAML.readValue(yaml, Map.class);
} catch (IOException | IllegalArgumentException e) {
return; // a malformed file is reported by the real parse, not here
}
if (raw == null || !(raw.get("profiles") instanceof Map<?, ?> profiles)) {
return;
}
List<String> bad = new ArrayList<>();
for (Map.Entry<?, ?> e : profiles.entrySet()) {
if (!(e.getValue() instanceof Map<?, ?> p)) {
continue;
}
if (!(p.get("errorPattern") instanceof String pattern) || pattern.isBlank()) {
continue;
}
try {
Pattern.compile(pattern);
} catch (PatternSyntaxException ex) {
bad.add("profiles." + e.getKey() + ".errorPattern (\"" + pattern + "\"): " + ex.getMessage());
}
}
bad.sort(String::compareTo);
if (!bad.isEmpty()) {
throw new IllegalStateException("refusing to start: malformed errorPattern — "
+ String.join("; ", bad));
}
}
/** The peer kinds this build has an adapter for — {@link Profile#kind()}'s only valid values. */
private static final Set<String> KNOWN_KINDS = Set.of(Profile.KIND_CLAUDE_CODE, Profile.KIND_OPENCODE);
@@ -620,9 +620,9 @@ public final class CompletionResolver implements TurnListener {
* @param allProfiles every configured profile name
* @param configuredProfiles the subset of {@code allProfiles} that carry an exhausted pattern
*/
public static String coverage(Set<String> allProfiles, Set<String> configuredProfiles) {
public static String coverage(String patternKey, Set<String> allProfiles, Set<String> configuredProfiles) {
if (configuredProfiles.isEmpty()) {
return "off (no profile has an exhaustedPattern configured; profiles: " + sorted(allProfiles) + ")";
return "off (no profile has an " + patternKey + " configured; profiles: " + sorted(allProfiles) + ")";
}
Set<String> unconfigured = new TreeSet<>(allProfiles);
unconfigured.removeAll(configuredProfiles);
@@ -17,6 +17,7 @@ import dev.ltms.fleet.msg.LeadMessage;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.peer.PeerUnreachableException;
import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementException;
import dev.ltms.fleet.session.SessionManager;
@@ -93,6 +94,10 @@ public final class FleetMcp {
private final CapacitySource capacity;
private final HealthCoverageSource healthCoverage;
private final QuarantineSource quarantine;
/** fleetd #201 Unit 5: SEPARATE from {@link #quarantine} — see {@link OutageSource}'s doc. */
private final OutageSource outage;
/** fleetd #176: SEPARATE from both of the above — see {@link LeadSeatSource}'s doc. */
private final LeadSeatSource leadSeats;
/** CB-637: this daemon's lead-to-lead channel; {@code null} when no coordinator is configured. */
private final LeadChannel leadChannel;
@@ -116,6 +121,40 @@ public final class FleetMcp {
public static QuarantineSource none() { return new QuarantineSource(_ -> null, BackendQuarantine.none()); }
}
/**
* fleetd #201 Unit 5 cool-off facts used by {@code fleet_list}/{@code fleet_profiles}: a
* profile → credential id lookup, plus the shared {@link BackendOutagePolicy} to read remaining
* cool-offs off. A SEPARATE source from {@link QuarantineSource} — never merged into it — so a
* credential that is cooling off after repeated backend errors is reported independently of
* whether it is ALSO under CB-578 stage B exhaustion quarantine; the two checks can both fire
* for the same profile, and when they do, {@code fleet_list}/{@code fleet_profiles} report both
* (see {@link #capacityView}/{@link #profiles}).
*/
public record OutageSource(Function<String, String> credentialIdFor, BackendOutagePolicy outagePolicy) {
/** Inert source — no profile is ever reported cooling off. Explicit stand-in, not a default. */
public static OutageSource none() {
return new OutageSource(_ -> null, new BackendOutagePolicy(System::nanoTime));
}
}
/**
* fleetd #176: the seats a profile's own live LEAD session(s) hold on the same Claude
* subscription — the third reason (alongside {@link QuarantineSource} and {@link OutageSource})
* {@code free} can overstate what a fresh {@code fleet_spawn} would actually get.
*
* <p>{@code maxLoad} counts only <em>members</em>, never the lead itself. But a
* {@code subscription: true} profile bills the operator's own Claude account, and the lead is
* always a live {@code claude} session on that same account (it is never moved off-subscription
* — see {@code LeadLauncher}). So a fan-out that fills every member slot still leaves the lead's
* own seat unaccounted for, and the daemon reports a slot that was never really free. See
* {@code Fleetd.leadSeatLookup} for how the count is derived — from {@code fleet.leaders.<name>
* .profile} and each profile's {@code effectiveCredentialId()}, never a hardcoded constant.
*/
public record LeadSeatSource(Function<String, Integer> seatsFor) {
/** Inert source — no profile is ever reported as sharing a seat with a lead. */
public static LeadSeatSource none() { return new LeadSeatSource(_ -> 0); }
}
/**
* @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
@@ -130,7 +169,7 @@ public final class FleetMcp {
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine) {
this(messages, workers, sessions, identity, presence, primaryRegistry, callers, metrics, capacity,
healthCoverage, quarantine, null);
healthCoverage, quarantine, null, OutageSource.none(), LeadSeatSource.none());
}
/**
@@ -143,9 +182,43 @@ public final class FleetMcp {
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, LeadChannel leadChannel) {
this(messages, workers, sessions, identity, presence, primaryRegistry, callers, metrics, capacity,
healthCoverage, quarantine, leadChannel, OutageSource.none(), LeadSeatSource.none());
}
/**
* As above, with fleetd #201 Unit 5 cool-off facts for {@code fleet_list}/{@code fleet_profiles}
* (see {@link OutageSource}).
*
* @param outage required — pass {@link OutageSource#none()} for a caller that does not want the
* feature, never a defaulting overload (the same rule {@code quarantine} follows).
*/
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) {
this(messages, workers, sessions, identity, presence, primaryRegistry, callers, metrics, capacity,
healthCoverage, quarantine, leadChannel, outage, LeadSeatSource.none());
}
/**
* 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).
*/
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.leadChannel = leadChannel;
this.capacity = capacity;
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
this.outage = Objects.requireNonNull(outage, "outage");
this.leadSeats = Objects.requireNonNull(leadSeats, "leadSeats");
this.healthCoverage = healthCoverage;
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
this.transport = HttpServletStreamableServerTransportProvider.builder()
@@ -274,8 +347,8 @@ public final class FleetMcp {
(exchange, _) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied;
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine,
callers == null ? Map.of() : callers.leads(),
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
leadSeats, callers == null ? Map.of() : callers.leads(),
callerTerminal(exchange),
leadChannel == null ? null : leadChannel.selfCoordId());
};
@@ -290,7 +363,7 @@ public final class FleetMcp {
(exchange, _) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied;
return profiles(workers, quarantine);
return profiles(workers, quarantine, outage);
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> whoamiHandler =
(exchange, _) -> {
@@ -875,25 +948,47 @@ public final class FleetMcp {
* has ever been quarantined gets exactly the pre-stage-B shape.
*/
static McpSchema.CallToolResult profiles(PeerLauncher workers, QuarantineSource quarantine) {
return profiles(workers, quarantine, OutageSource.none());
}
/**
* As above, plus a SEPARATE {@code coolingOff} map (fleetd #201 Unit 5) — never merged into
* {@code quarantined} — for a profile whose credential is cooling off after repeated backend
* errors ({@link BackendOutagePolicy}). The two checks are independent: a profile can appear in
* both maps at once when it is both exhaustion-quarantined AND cooling off.
*/
static McpSchema.CallToolResult profiles(PeerLauncher workers, QuarantineSource quarantine, OutageSource outage) {
Map<String, Object> result = new LinkedHashMap<>();
result.put("profiles", workers.profiles());
result.put("default", workers.defaultProfile() == null ? "" : workers.defaultProfile());
Map<String, Object> quarantined = new LinkedHashMap<>();
Map<String, Object> coolingOff = new LinkedHashMap<>();
for (String profile : workers.profiles()) {
String credentialId = quarantine.credentialIdFor().apply(profile);
if (credentialId == null) {
continue;
if (credentialId != null) {
quarantine.quarantine().remainingSeconds(credentialId).ifPresent(remaining -> {
Map<String, Object> row = new LinkedHashMap<>();
row.put("credentialId", credentialId);
row.put("quarantinedForSeconds", remaining);
quarantined.put(profile, row);
});
}
String outageCredentialId = outage.credentialIdFor().apply(profile);
if (outageCredentialId != null) {
outage.outagePolicy().remainingCoolOffSeconds(outageCredentialId).ifPresent(remaining -> {
Map<String, Object> row = new LinkedHashMap<>();
row.put("credentialId", outageCredentialId);
row.put("coolingOffForSeconds", remaining);
coolingOff.put(profile, row);
});
}
quarantine.quarantine().remainingSeconds(credentialId).ifPresent(remaining -> {
Map<String, Object> row = new LinkedHashMap<>();
row.put("credentialId", credentialId);
row.put("quarantinedForSeconds", remaining);
quarantined.put(profile, row);
});
}
if (!quarantined.isEmpty()) {
result.put("quarantined", quarantined);
}
if (!coolingOff.isEmpty()) {
result.put("coolingOff", coolingOff);
}
return text(json(result));
}
@@ -932,6 +1027,15 @@ public final class FleetMcp {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, leads, selfTerm, null);
}
/** 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) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
LeadSeatSource.none(), leads, selfTerm, null);
}
/**
* 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
@@ -945,6 +1049,25 @@ public final class FleetMcp {
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, Map<String, String> leads, String selfTerm,
String selfCoordId) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, OutageSource.none(),
LeadSeatSource.none(), leads, selfTerm, selfCoordId);
}
/** 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) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
LeadSeatSource.none(), leads, selfTerm, selfCoordId);
}
/** As above, plus fleetd #176 lead-seat facts (see {@link LeadSeatSource}). */
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, OutageSource outage,
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
String selfCoordId) {
try {
Map<String, Agent> live = workers.list().stream()
.map(Agent.class::cast)
@@ -971,7 +1094,7 @@ public final class FleetMcp {
}
if (capacity.available()) result.put("capacity", profiles.stream()
.map(profile -> capacityView(profile, capacity.liveCount(), capacity.maxLoad(), roster, messages,
capacity.clock().getAsLong(), quarantine)).toList());
capacity.clock().getAsLong(), quarantine, outage, leadSeats)).toList());
return text(json(result));
} catch (HerdrException e) {
return error("herdr error listing the fleet: " + e.getMessage());
@@ -1004,19 +1127,39 @@ public final class FleetMcp {
* reports, reusing {@link QuarantineSource} rather than a second lookup. Both new keys are
* added only when the profile is actually quarantined, so an ordinary fleet's rows are
* byte-identical to before this change.
*
* <p>fleetd #201 Unit 5: a profile whose credential is cooling off (a SEPARATE, independent
* check from quarantine — see {@link OutageSource}) also forces {@code free} to 0 and carries
* {@code credentialId}/{@code coolingOffForSeconds}, but never {@code quarantinedForSeconds} —
* that key is added only when exhaustion quarantine is ALSO active for this profile, since the
* two checks are independent and either, both, or neither can be true.
*
* <p>fleetd #176: {@code maxLoad} counts panes, not subscription seats — it never counted the
* lead's own seat on a {@code subscription: true} profile's account. {@link LeadSeatSource}
* reports that count (0 for a non-subscription profile, or when no live lead shares its
* credential), and it is subtracted from {@code free} the same way {@code live} already is —
* {@code maxLoad} itself is left untouched, so the row still reports the configured cap. The
* {@code leadSeats} key is added only when the count is positive, for the same
* byte-identical-when-unused reason as the quarantine/cool-off keys above.
*/
private static Map<String, Object> capacityView(String profile, Function<String, Integer> liveCount,
Function<String, Integer> maxLoad, List<MemberSession> roster,
MessageService messages, long nowNanos, QuarantineSource quarantine) {
MessageService messages, long nowNanos, QuarantineSource quarantine,
OutageSource outage, LeadSeatSource leadSeats) {
Integer cap = maxLoad.apply(profile);
int live = liveCount.apply(profile);
int leadSeatCount = leadSeats.seatsFor().apply(profile);
int reclaimable = (int) roster.stream().filter(s -> profile.equals(s.profile()))
.filter(s -> (s.state() == MemberSession.State.READY || s.state() == MemberSession.State.DONE))
.filter(s -> messages == null || (!messages.hasAcceptedDelivery(s.terminalId()) && !messages.hasInboxMessage(s.terminalId())))
.count();
Map<String, Object> row = new LinkedHashMap<>();
row.put("profile", profile); row.put("maxLoad", cap); row.put("live", live);
row.put("free", cap == null ? null : Math.max(0, cap - live)); row.put("reclaimable", reclaimable);
row.put("free", cap == null ? null : Math.max(0, cap - live - leadSeatCount));
row.put("reclaimable", reclaimable);
if (leadSeatCount > 0) {
row.put("leadSeats", leadSeatCount);
}
String credentialId = quarantine.credentialIdFor().apply(profile);
if (credentialId != null) {
quarantine.quarantine().remainingSeconds(credentialId).ifPresent(remaining -> {
@@ -1025,6 +1168,14 @@ public final class FleetMcp {
row.put("quarantinedForSeconds", remaining);
});
}
String outageCredentialId = outage.credentialIdFor().apply(profile);
if (outageCredentialId != null) {
outage.outagePolicy().remainingCoolOffSeconds(outageCredentialId).ifPresent(remaining -> {
row.put("free", 0);
row.put("credentialId", outageCredentialId);
row.put("coolingOffForSeconds", remaining);
});
}
return row;
}
@@ -1,5 +1,8 @@
package dev.ltms.fleet.member;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.Agent;
@@ -14,6 +17,8 @@ import java.io.IOException;
import java.io.UncheckedIOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.StandardCopyOption;
import java.nio.file.attribute.PosixFileAttributeView;
import java.util.EnumSet;
import java.util.List;
import java.util.Map;
@@ -47,6 +52,21 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
private static final Logger log = LoggerFactory.getLogger(ClaudeCodeLauncher.class);
/** JSON codec for the additive workspace-trust seed (fleetd #149) — Jackson's default settings. */
private static final ObjectMapper TRUST_JSON = new ObjectMapper();
/**
* Serialises every {@link #seedTrustDialog} read-modify-write for the whole daemon process.
* Two claude-code spawns starting at once are normal (fleetd runs several members in parallel
* routinely) and both would otherwise read the same {@code .claude.json}, add their own entry
* to their own in-memory copy, and write — the second write wins and the first spawn's trust
* entry silently disappears. A single process-wide lock is enough because every spawn on this
* daemon runs in this one JVM; it does not protect against a second daemon process or the
* operator's own Claude Code process writing at the same instant, which {@link #writeAtomically}
* covers instead (each writer only ever sees a fully-old or fully-new file, never a torn one).
*/
private static final Object TRUST_JSON_LOCK = new Object();
private final SubscriptionGuard guard;
/**
@@ -263,6 +283,10 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
putIfPresent(workerEnv, "ANTHROPIC_MODEL", cfg.model());
putIfPresent(workerEnv, "CLAUDE_CONFIG_DIR", cfg.configDir());
applyGitToken(workerEnv, cfg);
// fleetd #149: seed the workspace-trust entry BEFORE this spawn ever reaches herdr — see
// seedTrustDialog for why, and isProvisionedWorktree for why this is gated to a worktree
// fleetd itself provisioned (never a real checkout, never an un-configured fallback cwd).
seedTrustDialog(cfg.configDir(), spec.cwd());
// CB-547a: Claude Code can MINT its own session id, so fleetd chooses it — a fresh spawn
// gets a UUID we pass as --session-id and return from agentSessionId(), so the resume
@@ -412,12 +436,12 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
*/
private static void writeIdeOverlay(String cwd, String projectPath) {
try {
Path dotGit = Path.of(cwd, ".git");
if (!Files.isRegularFile(dotGit)) {
if (!isProvisionedWorktree(cwd)) {
// Not a provisioned worktree (primary's real checkout has a .git directory, or the
// cwd is not a repo at all). Never write into it.
return;
}
Path dotGit = Path.of(cwd, ".git");
// The overlay FILE lives at the worktree root (claude-code's cwd), but its CONTENT pins
// project_path to the module dir the IDE opened (projectPath), not the worktree root.
Files.writeString(Path.of(cwd, "CLAUDE.local.md"), PeerLauncher.ideOverlayText(projectPath));
@@ -447,6 +471,189 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
}
}
/**
* fleetd #149: pre-seed the workspace-trust entry for {@code cwd} in Claude Code's own config
* file, BEFORE this launch ever reaches herdr (called from {@link #buildLaunch}, which always
* runs before the base class starts the process). Claude Code asks an interactive, un-timed
* "Is this a project you created or one you trust?" the first time it starts in a directory it
* has not seen before, and every {@code worktree: true} spawn lands in a brand-new directory —
* so without this seed the member sits on that dialog forever, never mounts the bridge MCP, and
* never calls {@code fleet_reply}. herdr reports it as healthy the whole time ({@code
* agent_status: blocked}, {@code interactive_ready: true}), so nothing else catches it. Measured
* live on fleet01 2026-08-23: an unseeded fresh cwd sat on the dialog indefinitely; a seeded one
* reached {@code idle} clean.
*
* <p>This is not a new grant — the operator already trusted this repo by configuring the
* profile against it, and a worktree is a checkout of that same repo.
*
* <p>The key Claude Code reads is per-project, in {@code .claude.json}: {@code
* projects.<cwd>.hasTrustDialogAccepted}. The file lives at {@code <configDir>/.claude.json}
* when the profile sets {@code CLAUDE_CONFIG_DIR} (mirrors this launcher's own env var above),
* else the default {@code ~/.claude.json} — the same file Claude Code itself would read either
* way, so this seeds exactly what the spawned peer is about to open.
*
* <p><b>Additive, not a rewrite.</b> {@code .claude.json} is large (tens of KB, dozens of
* projects) and Claude Code itself rewrites it while running, so this reads the file as a JSON
* tree (missing or unreadable → treated as an empty object) and changes only
* {@code projects.<cwd>.hasTrustDialogAccepted} / {@code .hasCompletedProjectOnboarding} —
* every other top-level key and every other project entry is written back untouched. Only the
* one project entry for {@code cwd} is replaced/created; an existing entry for a DIFFERENT cwd
* (or the operator's own project history) is never touched.
*
* <p>Best-effort, like {@link #writeIdeOverlay}: a failure here (unwritable configDir, a
* corrupt existing file, …) must never fail the spawn — it is logged at debug and swallowed. A
* peer that starts without the seed still starts; it just may hit the dialog fleetd #149
* describes.
*
* <p><b>Gated to a provisioned worktree</b> ({@link #isProvisionedWorktree}) — see that
* method's javadoc for the incident that made this gate mandatory, not optional: this must
* never run against a real checkout or an un-configured fallback cwd, only the exact
* always-fresh-directory population fleetd #149 describes.
*
* <p><b>Atomic and lock-protected.</b> {@code .claude.json} is a live file — Claude Code itself
* rewrites it while running, and this daemon routinely spawns several members at once, each
* calling this method for its own cwd. Every write goes through {@link #writeAtomically} (a
* sibling-temp-file + {@code ATOMIC_MOVE}, never a truncate-in-place) so a crash mid-write or a
* concurrent reader never observes a half-written file, and through {@link #TRUST_JSON_LOCK} so
* two concurrent spawns' entries both survive instead of the second write silently discarding
* the first. Both exist because of a real incident: see {@link #isProvisionedWorktree}'s javadoc
* and {@link #writeAtomically}'s javadoc.
*
* @param configDir the profile's {@code CLAUDE_CONFIG_DIR} ({@code cfg.configDir()}), or
* {@code null}/blank to target the default {@code ~/.claude.json}
* @param cwd the spawn's resolved working directory — the exact key Claude Code will look
* up for itself once it starts there
*/
private static void seedTrustDialog(String configDir, String cwd) {
if (!isProvisionedWorktree(cwd)) {
return;
}
Path target = (configDir == null || configDir.isBlank())
? Path.of(System.getProperty("user.home"), ".claude.json")
: Path.of(configDir, ".claude.json");
synchronized (TRUST_JSON_LOCK) {
try {
if (target.getParent() != null) {
Files.createDirectories(target.getParent());
}
ObjectNode root = null;
if (Files.isRegularFile(target)) {
JsonNode existing = TRUST_JSON.readTree(target.toFile());
if (existing instanceof ObjectNode existingObject) {
root = existingObject;
}
}
if (root == null) {
root = TRUST_JSON.createObjectNode();
}
JsonNode projectsNode = root.get("projects");
ObjectNode projects = projectsNode instanceof ObjectNode projectsObject
? projectsObject : TRUST_JSON.createObjectNode();
if (!(projectsNode instanceof ObjectNode)) {
root.set("projects", projects);
}
JsonNode projectNode = projects.get(cwd);
ObjectNode project = projectNode instanceof ObjectNode projectObject
? projectObject : TRUST_JSON.createObjectNode();
if (!(projectNode instanceof ObjectNode)) {
projects.set(cwd, project);
}
project.put("hasTrustDialogAccepted", true);
project.put("hasCompletedProjectOnboarding", true);
writeAtomically(target, TRUST_JSON.writerWithDefaultPrettyPrinter().writeValueAsString(root));
} catch (Exception e) {
log.debug("cannot seed workspace-trust entry for cwd '{}' into '{}'", cwd, target, e);
}
}
}
/**
* Write {@code content} to {@code target} atomically: serialise to a sibling temp file in the
* <strong>same directory</strong> as {@code target} (an atomic move is only guaranteed within
* one filesystem — a different directory could mean a different filesystem), then
* {@link StandardCopyOption#ATOMIC_MOVE} it into place. A reader — Claude Code itself, or
* another {@code seedTrustDialog} call — only ever observes the fully-old file or the
* fully-new one, never a truncated or half-written one.
*
* <p><b>fleetd #149 incident.</b> The original implementation used
* {@code Files.writeString(target, content)} directly, which truncates {@code target} in place
* before writing the replacement bytes. Combined with an ungated {@code cwd} (see
* {@link #isProvisionedWorktree}'s javadoc), a mutation-testing run hit that truncation window
* against the operator's real {@code ~/.claude.json} and left it at 178 bytes. The gate closes
* <em>which file</em> this can ever target; this closes <em>how</em> the target is written, so
* that even a legitimate write against a real, live, concurrently-read {@code .claude.json}
* cannot leave it observably empty or partial.
*
* <p>Preserves {@code target}'s existing POSIX permissions (Claude Code ships {@code
* .claude.json} as {@code 0600}) when the filesystem reports them; a freshly created temp file
* already defaults to owner-only permissions on a POSIX filesystem, so a first-ever write (no
* existing {@code target}) is no less private without this. On a non-POSIX filesystem (e.g.
* Windows) the permission copy is a silent no-op rather than a failure.
*
* <p>Package-visible (not {@code private}) so a test can drive it directly with a concurrent
* reader thread and prove the torn-file property this method exists for — the LOCK in
* {@link #seedTrustDialog} already fully serialises every call this launcher itself makes, so a
* test that only ever goes through {@code seedTrustDialog}/{@code spawn()} could never observe
* a torn file regardless of whether this method is atomic; it would be proving the lock, not
* this method. Atomicity's actual job is protecting against a writer the lock cannot reach at
* all — a second daemon process, or the operator's own live Claude Code — so the test for it
* has to reach this method on its own.
*/
static void writeAtomically(Path target, String content) throws IOException {
Path parent = target.getParent();
Path tmp = Files.createTempFile(parent, target.getFileName() + ".", ".tmp");
try {
Files.writeString(tmp, content);
copyPosixPermissionsIfPresent(target, tmp);
Files.move(tmp, target, StandardCopyOption.ATOMIC_MOVE, StandardCopyOption.REPLACE_EXISTING);
} catch (IOException e) {
Files.deleteIfExists(tmp);
throw e;
}
}
/** Copy {@code target}'s POSIX permissions onto {@code tmp}, or no-op where either is unsupported. */
private static void copyPosixPermissionsIfPresent(Path target, Path tmp) {
try {
if (!Files.isRegularFile(target)) {
return; // nothing to inherit from — first-ever write, temp file's own default stands
}
PosixFileAttributeView view = Files.getFileAttributeView(target, PosixFileAttributeView.class);
if (view == null) {
return; // non-POSIX filesystem — nothing this JVM can read/set here
}
Files.setPosixFilePermissions(tmp, Files.getPosixFilePermissions(target));
} catch (IOException e) {
log.debug("cannot preserve permissions of '{}' onto its replacement", target, e);
}
}
/**
* Whether {@code cwd} is a fleetd-provisioned git worktree — signalled the same way
* {@link #writeIdeOverlay} already gates on: a {@code .git} that is a <strong>regular
* file</strong> holding a {@code gitdir:} pointer, as opposed to a real checkout's {@code .git}
* <strong>directory</strong>. {@code null}/blank never qualifies.
*
* <p>Shared by every write that must land only in a worktree fleetd itself created for a
* member — never in a real checkout, an arbitrary configured directory, or (see the incident
* below) the daemon's own fallback cwd.
*
* <p><b>fleetd #149 incident.</b> {@link #seedTrustDialog} originally ran unconditionally on
* any non-blank {@code cwd}. Most of this launcher's OWN tests spawn a profile with no
* {@code cwd} configured, so the base class's {@code resolveCwd} falls through to the real
* {@code user.dir} — and with no {@code configDir} either (also the common case in this
* file's fixtures), the seed's target falls through the same way to the real
* {@code ~/.claude.json}. Running this repo's own test suite corrupted the operator's actual
* config file (it shrank from ~72 KB to a single seeded entry) the first time a mutation
* happened to make the write non-additive. Gating both cwd-targeted writes on "this is a
* worktree fleetd provisioned" — exactly the population fleetd #149 describes
* ({@code worktree: true} always lands in a brand-new directory) — makes that class of write
* impossible against a real checkout or an untouched fallback cwd, in production or in tests.
*/
private static boolean isProvisionedWorktree(String cwd) {
return cwd != null && !cwd.isBlank() && Files.isRegularFile(Path.of(cwd, ".git"));
}
/** {@code s}, or {@code null} when {@code s} is null/blank — the charter-presence test used above. */
private static String nonBlank(String s) {
return (s == null || s.isBlank()) ? null : s;
@@ -10,6 +10,7 @@ import dev.ltms.fleet.peer.PeerHandle;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.peer.PeerUnreachableException;
import dev.ltms.fleet.peer.SpawnRequest;
import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementCandidate;
import dev.ltms.fleet.placement.PlacementContext;
@@ -92,6 +93,18 @@ public final class CompositePeerLauncher implements PeerLauncher {
/** CB-578 stage B: credential cooldown, checked before an explicit spawn and filtered into placement. */
private final BackendQuarantine quarantine;
/**
* fleetd #201 Unit 5: credential cool-off after repeated backend errors, checked before an
* explicit spawn and filtered into placement — a SEPARATE, shorter-lived source from
* {@link #quarantine}. A never-{@code record}-called instance is naturally inert (its
* {@code remainingCoolOffSeconds} always returns empty), so back-compat constructors that predate
* this feature share one fixed instance rather than needing a {@code none()} sentinel.
*/
private final BackendOutagePolicy outagePolicy;
/** The shared inert instance back-compat constructors wire in — never {@code record}-called. */
private static final BackendOutagePolicy NO_OUTAGE_POLICY = new BackendOutagePolicy(System::nanoTime);
/**
* CB-557: the role pools an unqualified spawn draws its candidates from. A supplier that yields
* {@code null}, and an empty pool for a role, both fall back to every configured profile — the
@@ -130,7 +143,8 @@ public final class CompositePeerLauncher implements PeerLauncher {
Map<String, FleetConfig.Profile> profileConfigs,
PlacementPolicy placementPolicy,
Function<String, Integer> liveCount) {
this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, null, BackendQuarantine.none());
this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, null,
BackendQuarantine.none(), NO_OUTAGE_POLICY);
}
/**
@@ -148,13 +162,14 @@ public final class CompositePeerLauncher implements PeerLauncher {
PlacementPolicy placementPolicy,
Function<String, Integer> liveCount,
FleetConfig.Fleet fleet) {
this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, fleet, BackendQuarantine.none());
this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, fleet,
BackendQuarantine.none(), NO_OUTAGE_POLICY);
}
/**
* Production constructor with role pools and quarantine (CB-578 stage B). The full-featured
* non-reloading form; {@link #CompositePeerLauncher(List, String, Supplier, Function, BackendQuarantine)}
* is what {@code Fleetd.main} actually wires up.
* Production constructor with role pools and quarantine (CB-578 stage B), no cool-off (fleetd
* #201 Unit 5). Kept for callers that predate the cool-off feature; use the 8-arg overload below
* to wire a real {@link BackendOutagePolicy}.
*
* @param quarantine required — pass {@link BackendQuarantine#none()} for a caller that does not
* want the feature, never a defaulting overload (CB-578 stage B's own rule).
@@ -166,17 +181,42 @@ public final class CompositePeerLauncher implements PeerLauncher {
Function<String, Integer> liveCount,
FleetConfig.Fleet fleet,
BackendQuarantine quarantine) {
this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, fleet, quarantine,
NO_OUTAGE_POLICY);
}
/**
* Production constructor with role pools, quarantine (CB-578 stage B), and cool-off (fleetd
* #201 Unit 5). The full-featured non-reloading form;
* {@link #CompositePeerLauncher(List, String, Supplier, Function, BackendQuarantine, BackendOutagePolicy)}
* is what {@code Fleetd.main} actually wires up.
*
* @param quarantine required — pass {@link BackendQuarantine#none()} for a caller that does not
* want the feature, never a defaulting overload (CB-578 stage B's own rule).
* @param outagePolicy required — pass a fresh, never-{@code record}-called {@link
* BackendOutagePolicy} for a caller that does not want the feature; same
* "explicit opt-out, never a silent default" rule as {@code quarantine}.
*/
public CompositePeerLauncher(List<HerdrPeerLauncher> delegates,
String defaultProfile,
Map<String, FleetConfig.Profile> profileConfigs,
PlacementPolicy placementPolicy,
Function<String, Integer> liveCount,
FleetConfig.Fleet fleet,
BackendQuarantine quarantine,
BackendOutagePolicy outagePolicy) {
// LinkedHashMap, not Map.copyOf: candidates() promises definition order and the weighted
// policy breaks exact-weight ties on it, so a salted iteration order would make placement
// differ from one JVM run to the next.
this(delegates, defaultProfile,
constant(Collections.unmodifiableMap(new LinkedHashMap<>(profileConfigs))),
constant(placementPolicy), liveCount, constant(fleet), quarantine);
constant(placementPolicy), liveCount, constant(fleet), quarantine, outagePolicy);
}
/**
* Production constructor that re-reads its placement inputs per spawn (CB-559), so a config
* reload retargets the next member without a restart.
* reload retargets the next member without a restart. No cool-off (fleetd #201 Unit 5); use the
* 6-arg overload below to wire a real {@link BackendOutagePolicy}.
*
* @param config the live configuration — read at every spawn, never captured
* @param quarantine required — CB-578 stage B; pass {@link BackendQuarantine#none()} to opt out
@@ -186,12 +226,31 @@ public final class CompositePeerLauncher implements PeerLauncher {
Supplier<FleetConfig> config,
Function<String, Integer> liveCount,
BackendQuarantine quarantine) {
this(delegates, defaultProfile, config, liveCount, quarantine, NO_OUTAGE_POLICY);
}
/**
* Production constructor that re-reads its placement inputs per spawn (CB-559), with cool-off
* (fleetd #201 Unit 5). This is what {@code Fleetd.main} actually wires up.
*
* @param config the live configuration — read at every spawn, never captured
* @param quarantine required — CB-578 stage B; pass {@link BackendQuarantine#none()} to opt out
* @param outagePolicy required — fleetd #201 Unit 5; pass a fresh, never-{@code record}-called
* {@link BackendOutagePolicy} to opt out
*/
public CompositePeerLauncher(List<HerdrPeerLauncher> delegates,
String defaultProfile,
Supplier<FleetConfig> config,
Function<String, Integer> liveCount,
BackendQuarantine quarantine,
BackendOutagePolicy outagePolicy) {
this(delegates, defaultProfile,
() -> config.get().profiles(),
() -> PlacementPolicies.fromName(config.get().placement()),
liveCount,
() -> config.get().fleet(),
quarantine);
quarantine,
outagePolicy);
}
/** The all-suppliers form every other constructor funnels into. */
@@ -201,9 +260,11 @@ public final class CompositePeerLauncher implements PeerLauncher {
Supplier<PlacementPolicy> placementPolicy,
Function<String, Integer> liveCount,
Supplier<FleetConfig.Fleet> fleet,
BackendQuarantine quarantine) {
BackendQuarantine quarantine,
BackendOutagePolicy outagePolicy) {
this.fleet = fleet;
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
this.outagePolicy = Objects.requireNonNull(outagePolicy, "outagePolicy");
if (delegates.isEmpty()) {
throw new IllegalArgumentException("at least one peer adapter must be configured");
}
@@ -266,7 +327,10 @@ public final class CompositePeerLauncher implements PeerLauncher {
// the charter makes explicit-profile spawns the normal path — so skipping the check
// here would leave the cap dead config in real operation.
HerdrPeerLauncher d = route(requestedProfile);
// Checked in this order so exhaustion quarantine wins when both are active: quarantine
// throws first and short-circuits before the cool-off check ever runs (fleetd #201 Unit 5).
enforceNotQuarantined(requestedProfile);
enforceNotCoolingOff(requestedProfile);
enforceMaxLoad(requestedProfile);
PeerHandle handle = d.spawn(req);
spawnedBy.put(handle.id(), d);
@@ -283,7 +347,10 @@ public final class CompositePeerLauncher implements PeerLauncher {
// CB-578 stage B: computed once up front — a quarantine's expiry cannot pass within one spawn
// call, so re-deriving it per retry would only cost work, never change the answer.
Set<String> quarantined = quarantinedProfiles(candidates);
PlacementContext ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable, quarantined);
// fleetd #201 Unit 5: a distinct set from quarantined — see PlacementContext.coolingOff.
Set<String> coolingOff = coolingOffProfiles(candidates);
PlacementContext ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable,
quarantined, coolingOff);
int maxAttempts = candidates.isEmpty() ? 1 : candidates.size();
for (int attempt = 0; attempt < maxAttempts; attempt++) {
@@ -313,7 +380,8 @@ public final class CompositePeerLauncher implements PeerLauncher {
chosen.profile(), e.getMessage());
unreachable.add(chosen.profile());
// Update the context for the next selection so the policy excludes this profile.
ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable, quarantined);
ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable,
quarantined, coolingOff);
}
}
@@ -379,6 +447,31 @@ public final class CompositePeerLauncher implements PeerLauncher {
.collect(Collectors.toSet());
}
/**
* Refuse an explicit-profile spawn whose credential is cooling off after repeated backend errors
* (fleetd #201 Unit 5 — {@link BackendOutagePolicy}): a SEPARATE, shorter-lived source from
* {@link #enforceNotQuarantined}'s exhaustion quarantine. Checked after quarantine so exhaustion
* wins when both are active — see the call site in {@link #spawn}.
*
* @throws PlacementException naming the profile, its credential, and the remaining cool-off
*/
private void enforceNotCoolingOff(String profile) {
String credentialId = credentialIdFor(profile);
outagePolicy.remainingCoolOffSeconds(credentialId).ifPresent(remaining -> {
throw new PlacementException("worker profile '" + profile + "' credential '" + credentialId
+ "' is cooling off after repeated backend errors; ~" + remaining
+ "s remaining — refusing spawn");
});
}
/** The subset of {@code candidates} whose credential is currently cooling off (fleetd #201 Unit 5). */
private Set<String> coolingOffProfiles(List<PlacementCandidate> candidates) {
return candidates.stream()
.map(PlacementCandidate::profile)
.filter(p -> outagePolicy.remainingCoolOffSeconds(credentialIdFor(p)).isPresent())
.collect(Collectors.toSet());
}
private void enforceMaxLoad(String profile) {
// Absent config, or a config whose maxLoad normalized to null (ABSENT ⇒ unlimited at load),
// means no cap — never cap what wasn't configured. Note "non-positive ⇒ unlimited" was true
@@ -1,55 +1,68 @@
package dev.ltms.fleet.placement;
import java.util.ArrayList;
import java.util.List;
/**
* Backward-compatible placement: an unqualified spawn always resolves to the configured default
* profile, exactly as {@code CompositePeerLauncher} did before CB-518. This ignores caps and
* reachability so that a pre-existing config behaves identically after upgrade.
*
* <p>Two exceptions walk past the default instead of returning it unconditionally:
* <p>Three exceptions walk past the default instead of returning it unconditionally:
* <ul>
* <li>Quarantine (CB-578 stage B): a quarantined default is a credential that just refused on
* a usage limit, not a transient capacity or reachability concern.
* <li>Cooling off (fleetd #201 Unit 5): a credential cooling off after repeated backend errors
* ({@code BackendOutagePolicy}) — a separate, shorter-lived source from quarantine. When a
* profile is both quarantined and cooling off, only the quarantine reason is reported
* (exhaustion takes priority), matching {@code CompositePeerLauncher}'s explicit-spawn order.
* <li>Weight 0 (CB-554): {@code fixed} is still automatic selection, so a profile the operator
* marked "never auto-select me" ({@code weight <= 0}) must be skipped here exactly as
* {@code weighted}/{@code round-robin} skip it — an explicit {@code fleet_spawn} naming
* the profile is unaffected, only this automatic fallback walk.
* </ul>
* A fleet where nothing is ever quarantined or weight-0 never exercises either path, so today's
* behaviour is unchanged.
* A fleet where nothing is ever quarantined, cooling off, or weight-0 never exercises any of these
* paths, so today's behaviour is unchanged.
*/
final class FixedPlacementPolicy implements PlacementPolicy {
@Override
public PlacementCandidate select(PlacementContext ctx) {
String d = ctx.defaultProfile();
if (d != null && !d.isBlank() && !ctx.quarantined().contains(d) && !weightExcluded(ctx, d)) {
if (d != null && !d.isBlank() && !ctx.quarantined().contains(d) && !ctx.coolingOff().contains(d)
&& !weightExcluded(ctx, d)) {
return new PlacementCandidate(d, null, 1.0f, null);
}
for (PlacementCandidate c : ctx.candidates()) {
if (!ctx.quarantined().contains(c.profile()) && !c.excluded()) {
if (!ctx.quarantined().contains(c.profile()) && !ctx.coolingOff().contains(c.profile())
&& !c.excluded()) {
return new PlacementCandidate(c.profile(), null, c.weight(), c.maxLoad());
}
}
if (d != null && !d.isBlank()) {
boolean dQuarantined = ctx.quarantined().contains(d);
// Exhaustion quarantine takes priority: reported only when quarantine is absent, so the
// message never claims "cooling off" for a profile that is really backend-exhausted.
boolean dCoolingOff = !dQuarantined && ctx.coolingOff().contains(d);
boolean dWeightExcluded = weightExcluded(ctx, d);
if (dQuarantined && dWeightExcluded) {
throw new PlacementException("worker profile '" + d + "' is quarantined (backend "
+ "exhausted) and has weight 0 (excluded from automatic selection), and no "
+ "available candidate remains");
}
if (dWeightExcluded) {
throw new PlacementException("worker profile '" + d + "' has weight 0 (excluded "
+ "from automatic selection) and no available candidate remains");
}
if (dQuarantined) {
throw new PlacementException("worker profile '" + d + "' is quarantined (backend "
+ "exhausted) and no un-quarantined candidate is available");
if (dQuarantined || dCoolingOff || dWeightExcluded) {
List<String> reasons = new ArrayList<>();
if (dQuarantined) {
reasons.add("is quarantined (backend exhausted)");
}
if (dCoolingOff) {
reasons.add("is cooling off after repeated backend errors");
}
if (dWeightExcluded) {
reasons.add("has weight 0 (excluded from automatic selection)");
}
throw new PlacementException("worker profile '" + d + "' "
+ String.join(" and ", reasons) + ", and no available candidate remains");
}
}
if (!ctx.candidates().isEmpty()) {
throw new PlacementException(
"all worker profiles are excluded from automatic selection (quarantined or weight-0)");
throw new PlacementException("all worker profiles are excluded from automatic "
+ "selection (quarantined, cooling off, or weight-0)");
}
throw new PlacementException("no worker profiles configured");
}
@@ -14,10 +14,19 @@ import java.util.function.Function;
* @param quarantined profiles whose credential is currently quarantined (CB-578 stage B) — a
* {@code BACKEND_EXHAUSTED} classification put it, or a profile it shares a
* credential with, on cooldown. Filtered the same way as {@code unreachable}.
* @param coolingOff profiles whose credential is currently cooling off after repeated backend
* errors (fleetd #201 Unit 5 — {@code BackendOutagePolicy}), a SEPARATE,
* shorter-lived source from {@code quarantined}: a credential outage cools off
* even when no member was ever exhausted. Deliberately its own set rather than
* merged into {@code quarantined} — {@link PlacementPolicyUtil} needs to tell
* the two apart so its refusal message says "cooling off", not "exhausted",
* when only this one is active. A profile can be in both sets at once; when it
* is, exhaustion quarantine is reported (it takes priority).
*/
public record PlacementContext(String defaultProfile,
List<PlacementCandidate> candidates,
Function<String, Integer> liveCount,
Set<String> unreachable,
Set<String> quarantined) {
Set<String> quarantined,
Set<String> coolingOff) {
}
@@ -14,14 +14,16 @@ final class PlacementPolicyUtil {
/**
* Candidates that are not weight-excluded (CB-554: explicit {@code weight <= 0}, checked
* first because it is a static config choice rather than transient state), not
* known-unreachable, not quarantined (CB-578 stage B), and have not reached their maxLoad.
* A {@code null} maxLoad means unlimited.
* known-unreachable, not quarantined (CB-578 stage B), not cooling off after repeated backend
* errors (fleetd #201 Unit 5 — a separate, shorter-lived source from quarantine), and have not
* reached their maxLoad. A {@code null} maxLoad means unlimited.
*/
static List<PlacementCandidate> available(PlacementContext ctx) {
List<PlacementCandidate> out = new ArrayList<>();
for (PlacementCandidate c : ctx.candidates()) {
if (c.excluded() || ctx.unreachable().contains(c.profile())
|| ctx.quarantined().contains(c.profile())) {
|| ctx.quarantined().contains(c.profile())
|| ctx.coolingOff().contains(c.profile())) {
continue;
}
Integer cap = c.maxLoad();
@@ -38,21 +40,27 @@ final class PlacementPolicyUtil {
/**
* Build a clear exception describing why every candidate was dropped: all weight-0, all
* quarantined, all at capacity, all unreachable, or a mix. Each candidate is counted into
* exactly one bucket (weight-excluded takes priority) so a candidate excluded for more than
* one reason is never double-counted.
* quarantined, all cooling off, all at capacity, all unreachable, or a mix. Each candidate is
* counted into exactly one bucket (weight-excluded first, then quarantined, then cooling off)
* so a candidate excluded for more than one reason is never double-counted — a candidate that is
* both quarantined (CB-578 stage B, backend exhausted) and cooling off (fleetd #201 Unit 5,
* repeated backend errors) counts only as quarantined, matching {@code CompositePeerLauncher}'s
* explicit-spawn ordering: exhaustion quarantine takes priority when both are active.
*/
static PlacementException emptyException(PlacementContext ctx) {
int weightExcluded = 0;
int atCap = 0;
int unreachable = 0;
int quarantined = 0;
int coolingOff = 0;
for (PlacementCandidate c : ctx.candidates()) {
Integer cap = c.maxLoad();
if (c.excluded()) {
weightExcluded++;
} else if (ctx.quarantined().contains(c.profile())) {
quarantined++;
} else if (ctx.coolingOff().contains(c.profile())) {
coolingOff++;
} else if (ctx.unreachable().contains(c.profile())) {
unreachable++;
} else if (cap != null && ctx.liveCount().apply(c.profile()) >= cap) {
@@ -71,6 +79,10 @@ final class PlacementPolicyUtil {
if (quarantined == total) {
return new PlacementException("all worker profiles are quarantined (backend exhausted)");
}
if (coolingOff == total) {
return new PlacementException(
"all worker profiles are cooling off after repeated backend errors");
}
if (atCap == total) {
return new PlacementException("all worker profiles are at maxLoad");
}
@@ -79,7 +91,9 @@ final class PlacementPolicyUtil {
}
return new PlacementException("no worker profile available: " + atCap + " at maxLoad, "
+ unreachable + " unreachable, " + quarantined + " quarantined, "
+ coolingOff + " cooling off, "
+ weightExcluded + " weight-0, "
+ (total - atCap - unreachable - quarantined - weightExcluded) + " remaining");
+ (total - atCap - unreachable - quarantined - coolingOff - weightExcluded)
+ " remaining");
}
}
@@ -426,19 +426,75 @@ public final class GitWorktrees implements Worktrees {
* neutralized copy from ever showing up as a local modification the worker might commit. A config
* the repo does not carry is skipped silently — no stub is invented for a file the repo does not
* have, and one missing file must never fail provisioning.
*
* <p>fleetd #134. The neutralization above is correct and stays unconditional — the defect was
* that it was invisible on both sides. Neither the daemon's own log nor the worker sitting in the
* worktree could tell a stub from the repo's real file: a real worker read a 3-byte {@code {}}
* where the repo's {@code opencode.json} is 30+ lines, and reported — truthfully from what it
* could see, and wrongly — that a mount key did not exist. Two fixes, aimed at two different
* readers:
* <ul>
* <li>the daemon operator reads {@link #log}, so the summary below names the denominator, what
* was neutralized, and why anything was not — the same shape {@code overlayParity} reports
* its own copy in;
* <li>the worker reads its own worktree, not the daemon's log, so the same fact is recorded a
* second time in worktree-scoped git config ({@code fleet.neutralizedConfig} /
* {@code fleet.neutralizedConfigNote}, readable with {@code git config --worktree --get-all
* fleet.neutralizedConfig}) rather than as a file in the working tree. A working-tree file
* would show up in {@code git status} for the worker to trip on or commit; worktree-scoped
* config lives in {@code .git/worktrees/<nonce>/config.worktree} and can never appear there.
* This reuses the exact mechanism {@link #configureEnvironmentCredentialHelper} and
* {@link #configureHttpsUrlRewriteForSshOrigin} already use for other worktree-local state.
* </ul>
*/
private void isolateToolSurface(String worktreePath) {
Path root = Path.of(worktreePath).toAbsolutePath().normalize();
List<String> neutralized = new ArrayList<>();
List<String> skipped = new ArrayList<>();
for (WorktreeHostileConfig cfg : WORKTREE_HOSTILE_CONFIGS) {
neutralize(root, worktreePath, cfg);
if (neutralize(root, worktreePath, cfg)) {
neutralized.add(cfg.file());
} else {
skipped.add(cfg.file() + " absent");
}
}
String detail = neutralized.isEmpty() ? String.join(", ", skipped)
: skipped.isEmpty() ? String.join(", ", neutralized)
: String.join(", ", neutralized) + " (" + String.join(", ", skipped) + ")";
log.info("tool-surface isolation: neutralized {} of {} configs: {} — the worktree copy is a "
+ "stub, not the repo's file; edit the real file in the primary checkout instead",
neutralized.size(), WORKTREE_HOSTILE_CONFIGS.size(), detail);
recordNeutralizedConfigForWorker(worktreePath, neutralized);
}
private void neutralize(Path root, String worktreePath, WorktreeHostileConfig cfg) {
/**
* The worker-readable half of fleetd #134: record which files were neutralized where the worker
* itself can read it, without a working-tree file that would show up in {@code git status}.
* Worktree-scoped git config is per-worktree, lives under {@code .git/worktrees/<nonce>/} rather
* than the working tree, and this repo already relies on the same mechanism (and the same
* {@code extensions.worktreeConfig} enablement) for the credential helper and the SSH→HTTPS
* rewrite — see {@link #configureEnvironmentCredentialHelper}.
*/
private void recordNeutralizedConfigForWorker(String worktreePath, List<String> neutralized) {
if (neutralized.isEmpty()) {
return;
}
exec("git", "-C", worktreePath, "config", "extensions.worktreeConfig", "true");
for (String file : neutralized) {
exec("git", "-C", worktreePath, "config", "--worktree", "--add", "fleet.neutralizedConfig", file);
}
exec("git", "-C", worktreePath, "config", "--worktree", "fleet.neutralizedConfigNote",
"the worktree copy of each fleet.neutralizedConfig path is a stub, not the repo's "
+ "committed file; edit the real file from the primary checkout instead");
}
/** @return true if {@code cfg} was neutralized (present, or created because {@link
* WorktreeHostileConfig#createIfAbsent()}); false if the repo does not carry it and it was
* left alone. */
private boolean neutralize(Path root, String worktreePath, WorktreeHostileConfig cfg) {
Path target = root.resolve(cfg.file());
if (!Files.exists(target) && !cfg.createIfAbsent()) {
log.debug("{} absent in the worktree — skipping (repo does not carry it)", cfg.file());
return;
return false;
}
try {
Files.writeString(target, cfg.stub());
@@ -449,7 +505,7 @@ public final class GitWorktrees implements Worktrees {
if (isTracked(root, cfg.file())) {
exec("git", "-C", worktreePath, "update-index", "--skip-worktree", cfg.file());
}
log.debug("neutralized {} — worker tool surface is launcher-mounted only", cfg.file());
return true;
}
@Override
@@ -486,25 +542,46 @@ public final class GitWorktrees implements Worktrees {
}
Path srcRoot = Path.of(repoRoot).toAbsolutePath().normalize();
Path dstRoot = Path.of(worktreePath).toAbsolutePath().normalize();
List<String> copied = new ArrayList<>();
List<String> skipped = new ArrayList<>();
List<String> neutralized = new ArrayList<>();
for (String rel : overlay) {
Path src = srcRoot.resolve(rel).normalize();
if (!Files.exists(src)) {
log.debug("parity overlay source missing — skipping {}", rel);
skipped.add(rel + " absent");
continue;
}
Path dst = dstRoot.resolve(rel).normalize();
try {
Files.createDirectories(dst.getParent());
Files.copy(src, dst, StandardCopyOption.REPLACE_EXISTING, StandardCopyOption.COPY_ATTRIBUTES);
log.debug("copied parity overlay {}", rel);
copied.add(rel);
} catch (IOException e) {
throw new WorktreeException("cannot copy overlay " + rel + ": " + e.getMessage(), e);
}
if (isTracked(dstRoot, rel)) {
exec("git", "-C", worktreePath, "update-index", "--skip-worktree", rel);
log.debug("marked overlay --skip-worktree {}", rel);
neutralized.add(rel);
}
}
// CB-148 point 3: a bare count ("copied 1") hides which candidates were even considered — the
// same shape of under-reporting this repo has been bitten by before. Name the denominator
// (every configured candidate), what was actually copied, and — for anything not copied —
// why, so a spawn's overlay outcome is legible from the log alone, no filesystem dig required.
String detail = copied.isEmpty() ? String.join(", ", skipped)
: skipped.isEmpty() ? String.join(", ", copied)
: String.join(", ", copied) + " (" + String.join(", ", skipped) + ")";
log.info("parity overlay: copied {} of {} candidates: {}", copied.size(), overlay.size(), detail);
// CB-134: a neutralized tracked file is otherwise a silent trap — a worker edits it, git
// ignores the change with no error, and nothing anywhere said the file could not be
// committed from this worktree. Name every file marked --skip-worktree here, with the
// consequence stated in the message itself, rather than adding a worktree-local marker
// file: acceptance criterion 1 requires the worktree hold exactly the configured overlay
// set and nothing else, so an extra marker would itself violate the fix.
if (!neutralized.isEmpty()) {
log.info("parity overlay marked --skip-worktree (cannot be committed from this worktree): {}",
String.join(", ", neutralized));
}
}
@Override
@@ -0,0 +1,59 @@
package dev.ltms.fleet;
import dev.ltms.fleet.inject.BackendErrorPatternLookup;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.session.MemberSession;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import java.util.List;
import java.util.Map;
import java.util.regex.Pattern;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
/**
* fleetd #248 / fleetd#201 Unit 5: {@link Fleetd#backendErrorPatternLookup} is the factory that
* replaced the local lambda {@code Fleetd.main} used to build {@code backendErrorPatterns} — one
* of the two arguments {@code CompletionResolver} lost cleanly (0 compile errors, every test still
* green) when this ticket's measurement dropped it alongside {@code backendErrorSink}. This class
* proves the factory's own behaviour; {@code FleetdCompletionResolverWiringTest} proves {@code
* main} still passes its result into {@code CompletionResolver}.
*/
class FleetdBackendErrorPatternLookupTest {
private static MemberSession session(String terminal, String profile) {
return new MemberSession("pane-" + terminal, terminal, profile, MemberRole.DEV,
"/cwd", null, 0L, 0L, 0, MemberSession.State.READY, null, null);
}
@Test
@DisplayName("a target on a profile with a configured pattern resolves to that pattern")
void configuredProfileResolves() {
Map<String, Pattern> byProfile = Map.of("terra", Pattern.compile("(?i)503"));
BackendErrorPatternLookup lookup =
Fleetd.backendErrorPatternLookup(() -> List.of(session("term1", "terra")), byProfile);
assertEquals("(?i)503", lookup.patternFor("term1").pattern());
}
@Test
@DisplayName("a target on a profile with no configured pattern resolves to null")
void unconfiguredProfileResolvesToNull() {
Map<String, Pattern> byProfile = Map.of("terra", Pattern.compile("x"));
BackendErrorPatternLookup lookup =
Fleetd.backendErrorPatternLookup(() -> List.of(session("term1", "sol")), byProfile);
assertNull(lookup.patternFor("term1"));
}
@Test
@DisplayName("an unknown target resolves to null")
void unknownTargetResolvesToNull() {
BackendErrorPatternLookup lookup =
Fleetd.backendErrorPatternLookup(List::of, Map.of("terra", Pattern.compile("x")));
assertNull(lookup.patternFor("term_stranger"));
}
}
@@ -0,0 +1,229 @@
package dev.ltms.fleet;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.BackendErrorSink;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.member.CompositePeerLauncher;
import dev.ltms.fleet.mcp.PrimaryRegistry;
import dev.ltms.fleet.msg.InMemoryReplyInbox;
import dev.ltms.fleet.msg.ReplyPushLoop;
import dev.ltms.fleet.peer.Capability;
import dev.ltms.fleet.peer.PeerHandle;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.peer.SpawnRequest;
import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementPolicies;
import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.SessionManager;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #248 / fleetd#201 Unit 5: {@link Fleetd#backendErrorSink} is the factory that replaced
* the local lambda {@code Fleetd.main} used to build {@code backendErrorSink} — the other half of
* the pair this ticket's measurement dropped cleanly (0 compile errors, every test still green).
*
* <p>Before this ticket, the closest thing to coverage was {@code
* dev.ltms.fleet.inject.BackendOutageFlowTest}, whose own class doc said it "mirrors {@code
* Fleetd.main}'s {@code backendErrorSink} lambda line-for-line" — a hand-copy that proves itself,
* never that {@code main} still wires the real thing. This class exercises the actual production
* factory instead. {@code FleetdCompletionResolverWiringTest} proves {@code main} still passes its
* result into {@code CompletionResolver}.
*/
class FleetdBackendErrorSinkTest {
private final List<ScheduledExecutorService> schedulers = new ArrayList<>();
@AfterEach
void tearDown() {
schedulers.forEach(ScheduledExecutorService::shutdownNow);
}
private static FleetConfig.Profile stubWorker(String profile, String credentialId) {
return new FleetConfig.Profile(profile, "http://gx00.gw:8000", "coder",
null, "FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers",
"w #{n}", null, null, null, null, null, null, null, null, null,
null, null, credentialId, null);
}
private static Map<String, FleetConfig.Profile> orderedProfiles() {
Map<String, FleetConfig.Profile> m = new LinkedHashMap<>();
m.put("terra", stubWorker("terra", "shared-openai"));
m.put("sol", stubWorker("sol", "shared-openai"));
return m;
}
/** Minimal recording {@code HerdrClient} for the LEAD pane — mirrors ReplyPushLoopTest's own. */
private static final class RecordingLeadClient implements HerdrClient {
private static final ObjectMapper MAPPER = new ObjectMapper();
private final List<Object> prompts = new CopyOnWriteArrayList<>();
volatile CountDownLatch sendLatch = new CountDownLatch(1);
@Override
public JsonNode call(String method, Object params) {
if ("agent.get".equals(method)) {
return MAPPER.createObjectNode().set("agent", MAPPER.createObjectNode()
.put("terminal_id", "term_primary").put("agent_status", "idle"));
}
if ("agent.prompt".equals(method)) {
prompts.add(params);
sendLatch.countDown();
}
return MAPPER.createObjectNode();
}
@Override
public void close() {
}
int sendCount() {
return prompts.size();
}
}
/** A {@link PeerLauncher} that never actually spawns — enough to construct a bare {@link SessionManager}. */
private static final class NeverSpawnsLauncher implements PeerLauncher {
@Override
public Set<Capability> capabilities() {
return Set.of();
}
@Override
public Set<Capability> capabilitiesFor(String profileName) {
return Set.of();
}
@Override
public PeerHandle spawn(SpawnRequest req) {
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
}
@Override
public Set<String> profiles() {
return Set.of();
}
@Override
public String defaultProfile() {
return null;
}
@Override
public String effectiveCwd(SpawnRequest req) {
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
}
@Override
public List<String> parityOverlay(String profileName) {
return List.of();
}
@Override
public List<?> list() {
return List.of();
}
@Override
public int reapOrphanWorkers() {
return 0;
}
@Override
public void stop(String id) {
}
@Override
public boolean clearContext(String id) {
return false;
}
}
@Test
@DisplayName("a target with no resolvable profile logs and returns without recording an incident (never throws)")
void unresolvableProfileDoesNotRecordOrThrow() {
SessionManager sessions = new SessionManager(new NeverSpawnsLauncher());
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L);
RecordingLeadClient leadClient = new RecordingLeadClient();
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
PrimaryRegistry registry = new PrimaryRegistry(null);
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
schedulers.add(scheduler);
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, new AgentControl(leadClient), inbox, scheduler, 3, 50);
BackendErrorSink sink = Fleetd.backendErrorSink(sessions, Map::of, outagePolicy, () -> pushLoop);
sink.onBackendError("term_unmapped", "matched line", "503 Service Unavailable");
assertTrue(outagePolicy.remainingCoolOffSeconds("shared-openai").isEmpty(),
"no credential is ever resolvable here, so nothing must be recorded");
}
@Test
@DisplayName("two distinct targets classified through the real factory start an incident and cool the credential")
void twoDistinctTargetsStartAnIncident() throws Exception {
FakeHerdr herdr = new FakeHerdr()
.readText("⏺ 503 Service Unavailable: upstream credential rejected\n❯ ");
Map<String, FleetConfig.Profile> profiles = orderedProfiles();
AtomicLong clockNanos = new AtomicLong(0L);
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(clockNanos::get);
ClaudeCodeLauncher adapter = new ClaudeCodeLauncher(new AgentControl(herdr),
new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
profiles, "terra", _ -> "tok");
CompositePeerLauncher workers = new CompositePeerLauncher(List.of(adapter), "terra", profiles,
PlacementPolicies.weighted(), _ -> 0, null, BackendQuarantine.none(), outagePolicy);
SessionManager sessions = new SessionManager(workers);
MemberSession s1 = sessions.acquire("terra", null, null, null);
MemberSession s2 = sessions.acquire("terra", null, null, null);
PrimaryRegistry registry = new PrimaryRegistry(null);
registry.recordDelegation(s1.terminalId(), "term_primary");
registry.recordDelegation(s2.terminalId(), "term_primary");
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
inbox.own(s1.terminalId());
inbox.own(s2.terminalId());
RecordingLeadClient leadClient = new RecordingLeadClient();
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
schedulers.add(scheduler);
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, new AgentControl(leadClient), inbox, scheduler, 3, 50);
AtomicReference<ReplyPushLoop> pushLoopRef = new AtomicReference<>(pushLoop);
// The exact object under test: Fleetd's real production factory, not a hand copy.
BackendErrorSink sink = Fleetd.backendErrorSink(sessions, () -> profiles, outagePolicy, pushLoopRef::get);
sink.onBackendError(s1.terminalId(), "matched line", "503 Service Unavailable");
assertTrue(outagePolicy.remainingCoolOffSeconds("shared-openai").isEmpty(),
"one distinct target must not start a cool-off");
sink.onBackendError(s2.terminalId(), "matched line", "503 Service Unavailable");
assertTrue(leadClient.sendLatch.await(3, TimeUnit.SECONDS),
"the second distinct target must cross the threshold and nudge the lead");
var remaining = outagePolicy.remainingCoolOffSeconds("shared-openai");
assertTrue(remaining.isPresent(), "two distinct targets must start a cool-off");
assertEquals(1, leadClient.sendCount());
}
}
@@ -0,0 +1,89 @@
package dev.ltms.fleet;
import java.nio.file.Files;
import java.nio.file.Path;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #248: this is the test that was actually missing. {@code Fleetd.main} builds its {@code
* CompletionResolver} from an 8-argument constructor, and the ticket's own measurement proved two
* ways to silently unwire it — both compiled with 0 errors and left every existing test green:
*
* <ul>
* <li>replacing the worktree/branch argument (the 8th) with {@code _ -> null} — drops
* fleetd#241's fallback-report location entirely;</li>
* <li>replacing {@code backendErrorPatterns, backendErrorSink} (5th/6th) with {@code
* BackendErrorPatternLookup.legacy(), BackendErrorSink.none()} — drops fleetd#201 Unit 5's
* backend-error classification and cool-off entirely.</li>
* </ul>
*
* <p>Neither mutation could be caught by any test that constructs its own {@code
* CompletionResolver} (every test before this one did exactly that) or by a test of {@link
* Fleetd#worktreeBranchLookup}, {@link Fleetd#backendErrorPatternLookup}, or {@link
* Fleetd#backendErrorSink} in isolation (see {@code FleetdWorktreeBranchLookupTest}, {@code
* FleetdBackendErrorPatternLookupTest}, {@code FleetdBackendErrorSinkTest}) — those prove the
* factories work, never that {@code main} still calls them. This class is a plain source-text
* assertion on {@code Fleetd.java} — crude, but honest about what it checks, and it turns red the
* instant the wiring is dropped, mirroring the same fallback shape {@link
* FleetdFleetAppConstructionTest} already uses for a different constructor argument.
*
* <p><b>This test checks source text, not runtime behaviour.</b> It never constructs a {@code
* CompletionResolver} and never runs {@code main}.
*/
class FleetdCompletionResolverWiringTest {
private static String fleetdSource() throws Exception {
return Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
}
@Test
@DisplayName("[SOURCE TEXT] CompletionResolver's construction call still names backendErrorPatterns and backendErrorSink")
void backendErrorArgumentsAreStillNamedAtTheCallSite() throws Exception {
String source = fleetdSource();
assertTrue(source.contains(
"exhaustionSink, backendErrorPatterns, backendErrorSink, System::nanoTime,"),
"CompletionResolver's construction call must still pass backendErrorPatterns and "
+ "backendErrorSink as its 5th/6th arguments. Replacing them with "
+ "BackendErrorPatternLookup.legacy()/BackendErrorSink.none() (fleetd #248's measured "
+ "mutation) compiles with 0 errors and leaves every behavioural test green — this "
+ "source check is what must go red instead.");
}
@Test
@DisplayName("[SOURCE TEXT] CompletionResolver's construction call still passes worktreeBranchLookup(sessions::roster)")
void worktreeBranchLookupIsStillPassedAtTheCallSite() throws Exception {
String source = fleetdSource();
assertTrue(source.contains("worktreeBranchLookup(sessions::roster)"),
"CompletionResolver's construction call must still pass worktreeBranchLookup(sessions::roster) "
+ "as its 8th (last) argument. Replacing it with the inert `_ -> null` (fleetd #248's "
+ "other measured mutation) compiles with 0 errors and leaves every behavioural test "
+ "green — this source check is what must go red instead.");
assertFalse(source.contains("System::nanoTime,\n _ -> null"),
"the worktree/branch argument must never regress to the inert `_ -> null` literal");
}
@Test
@DisplayName("[SOURCE TEXT] backendErrorPatterns is assigned from the extracted backendErrorPatternLookup(...) factory")
void backendErrorPatternsComesFromTheFactory() throws Exception {
String source = fleetdSource();
assertTrue(source.contains(
"BackendErrorPatternLookup backendErrorPatterns = backendErrorPatternLookup(sessions::roster,"),
"backendErrorPatterns must be assigned from Fleetd.backendErrorPatternLookup(...), not an "
+ "inline lambda that a source check on the CompletionResolver call alone cannot see "
+ "through");
}
@Test
@DisplayName("[SOURCE TEXT] backendErrorSink is assigned from the extracted backendErrorSink(...) factory")
void backendErrorSinkComesFromTheFactory() throws Exception {
String source = fleetdSource();
assertTrue(source.contains(
"BackendErrorSink backendErrorSink = backendErrorSink(sessions, () -> config.get().profiles(),"),
"backendErrorSink must be assigned from Fleetd.backendErrorSink(...), not an inline lambda "
+ "that a source check on the CompletionResolver call alone cannot see through");
}
}
@@ -0,0 +1,156 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.FleetConfig;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import java.util.Map;
import java.util.function.Function;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* fleetd #176: {@link Fleetd#leadSeatLookup} is the factory {@code Fleetd.main} wires into {@code
* FleetMcp.LeadSeatSource} so {@code fleet_list}'s {@code free} can subtract the seat(s) a
* {@code subscription: true} profile's own live LEAD session holds on that same account —
* {@code maxLoad} never counted the lead, only members. {@code FleetdLeadSeatWiringTest} proves
* {@code main} still passes this factory's result in; this class proves the factory's own matching
* logic: subscription-only, credential-matched, and counting only CURRENTLY LIVE leads.
*/
class FleetdLeadSeatLookupTest {
private static FleetConfig.Profile subscriptionProfile(String name, String credentialId) {
return new FleetConfig.Profile(name, null, "claude-sonnet-5", null, null, null,
"tab", "fleet", "w #{n}", null, null, null, null, null, null, null,
null, 3, true, null, credentialId, null);
}
private static FleetConfig.Profile offSubscriptionProfile(String name, String credentialId) {
return new FleetConfig.Profile(name, "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
null, "tab", "fleet", "w #{n}", null, null, null, null, null, null, null,
null, 2, false, null, credentialId, null);
}
private static FleetConfig.Leader leadOnProfile(String profile) {
return new FleetConfig.Leader(profile, "lead: primary", 1, "lead:", 10, "claude", "claude-sonnet-5");
}
@Test
@DisplayName("a live lead sharing the target profile's credential counts as one seat")
void liveLeadSharingCredentialCountsAsOneSeat() {
Map<String, FleetConfig.Profile> profiles = Map.of("sonnet", subscriptionProfile("sonnet", null));
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("sonnet"));
Function<String, Integer> lookup = Fleetd.leadSeatLookup(() -> profiles, leaders,
() -> Map.of("term_primary", "primary"));
assertEquals(1, lookup.apply("sonnet"));
}
@Test
@DisplayName("no live lead names this profile ⇒ zero seats, exactly as before this ticket")
void noLiveLeadOnTheProfileCountsAsZero() {
Map<String, FleetConfig.Profile> profiles = Map.of("sonnet", subscriptionProfile("sonnet", null));
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("sonnet"));
Function<String, Integer> lookup = Fleetd.leadSeatLookup(() -> profiles, leaders, Map::of);
assertEquals(0, lookup.apply("sonnet"));
}
@Test
@DisplayName("a lead entry with no `profile:` (recognise-only) contributes no seats — cannot be derived")
void recogniseOnlyLeadWithNoProfileContributesNothing() {
Map<String, FleetConfig.Profile> profiles = Map.of("sonnet", subscriptionProfile("sonnet", null));
FleetConfig.Leader recogniseOnly = new FleetConfig.Leader(null, "lead: primary", 1, "lead:", 10,
"claude", "claude-sonnet-5");
Map<String, FleetConfig.Leader> leaders = Map.of("primary", recogniseOnly);
Function<String, Integer> lookup = Fleetd.leadSeatLookup(() -> profiles, leaders,
() -> Map.of("term_primary", "primary"));
assertEquals(0, lookup.apply("sonnet"));
}
@Test
@DisplayName("a non-subscription profile never has a lead seat subtracted, whatever the credential match")
void nonSubscriptionProfileIsNeverAdjusted() {
Map<String, FleetConfig.Profile> profiles = Map.of(
"terra", offSubscriptionProfile("terra", "shared-openai"),
"sonnet", subscriptionProfile("sonnet", "shared-openai"));
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("sonnet"));
Function<String, Integer> lookup = Fleetd.leadSeatLookup(() -> profiles, leaders,
() -> Map.of("term_primary", "primary"));
assertEquals(0, lookup.apply("terra"), "terra is not subscription:true, so it must never be adjusted");
}
@Test
@DisplayName("explicit, different credentialIds still separate two subscription profiles (post fleetd #176 "
+ "stage 2 sentinel)")
void differentCredentialIsNotCounted() {
Map<String, FleetConfig.Profile> profiles = Map.of(
"sonnet", subscriptionProfile("sonnet", "claude-account-a"),
"opus", subscriptionProfile("opus", "claude-account-b"));
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
Function<String, Integer> lookup = Fleetd.leadSeatLookup(() -> profiles, leaders,
() -> Map.of("term_primary", "primary"));
assertEquals(0, lookup.apply("sonnet"), "different accounts must never be conflated into one seat count "
+ "— an explicit credentialId on both sides must still win over the subscription sentinel, so an "
+ "operator with two separate Claude logins on one host can keep them apart");
}
/**
* fleetd #176 stage 2 — the exact live shape that shipped inert: a lead on subscription profile
* {@code opus}, members on a DIFFERENTLY NAMED subscription profile {@code sonnet}, same Claude
* login, and NEITHER profile sets {@code credentialId}. Every other test in this class puts the
* lead on the SAME profile name as the target, which happened to keep working even with the old
* fall-back-to-profile-name {@code effectiveCredentialId()} — this is the one that did not, and
* its absence is what let the stage-1 fix ship without ever catching the bug it was filed for.
*/
@Test
@DisplayName("[LIVE SHAPE] lead on a DIFFERENT subscription profile, same account, neither sets "
+ "credentialId ⇒ still counts as a seat")
void leadOnADifferentSubscriptionProfileSameAccountStillCountsAsASeat() {
Map<String, FleetConfig.Profile> profiles = Map.of(
"opus", subscriptionProfile("opus", null),
"sonnet", subscriptionProfile("sonnet", null));
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
Function<String, Integer> lookup = Fleetd.leadSeatLookup(() -> profiles, leaders,
() -> Map.of("term_primary", "primary"));
assertEquals(1, lookup.apply("sonnet"), "opus and sonnet are both subscription:true with no explicit "
+ "credentialId, so they share one Claude login and the lead's live seat on opus must be charged "
+ "against sonnet too — this is the live host's actual shape (fleetd #176 stage 2)");
}
@Test
@DisplayName("two live instances of the same lead count as two seats")
void twoLiveInstancesOfTheSameLeadCountAsTwoSeats() {
Map<String, FleetConfig.Profile> profiles = Map.of("sonnet", subscriptionProfile("sonnet", null));
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("sonnet"));
Function<String, Integer> lookup = Fleetd.leadSeatLookup(() -> profiles, leaders,
() -> Map.of("term_a", "primary", "term_b", "primary"));
assertEquals(2, lookup.apply("sonnet"));
}
@Test
@DisplayName("an unconfigured target profile resolves to zero, not a thrown exception")
void unconfiguredTargetProfileIsZero() {
Function<String, Integer> lookup = Fleetd.leadSeatLookup(Map::of, Map.of(), Map::of);
assertEquals(0, lookup.apply("ghost"));
}
@Test
@DisplayName("live leads are read through the supplier on every call, not snapshotted")
void liveLeadsAreReadThroughOnEveryCall() {
Map<String, FleetConfig.Profile> profiles = Map.of("sonnet", subscriptionProfile("sonnet", null));
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("sonnet"));
java.util.Map<String, String> live = new java.util.HashMap<>();
Function<String, Integer> lookup = Fleetd.leadSeatLookup(() -> profiles, leaders, () -> live);
assertEquals(0, lookup.apply("sonnet"));
live.put("term_primary", "primary");
assertEquals(1, lookup.apply("sonnet"));
}
}
@@ -0,0 +1,43 @@
package dev.ltms.fleet;
import java.nio.file.Files;
import java.nio.file.Path;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #176: {@code Fleetd.main} builds its {@code FleetMcp} from a 14-argument constructor whose
* last argument is a {@code FleetMcp.LeadSeatSource} wrapping {@link Fleetd#leadSeatLookup}. That
* argument is exactly the kind of wiring fleetd #248 warned about: dropping it (or swapping it for
* the inert {@code FleetMcp.LeadSeatSource.none()}) compiles with 0 errors and leaves every test
* that builds its own {@code FleetMcp}/{@code CapacitySource} directly — every test that predates
* this ticket — green, because none of them go through {@code main} at all.
*
* <p>{@link FleetdLeadSeatLookupTest} proves the factory's own matching logic; this class is the
* plain source-text assertion that proves {@code main} still passes its result in, mirroring
* {@code FleetdCompletionResolverWiringTest}'s approach for the same class of gap.
*
* <p><b>This test checks source text, not runtime behaviour.</b> It never constructs a
* {@code FleetMcp} and never runs {@code main}.
*/
class FleetdLeadSeatWiringTest {
private static String fleetdSource() throws Exception {
return Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
}
@Test
@DisplayName("[SOURCE TEXT] FleetMcp's construction call still passes a LeadSeatSource built from leadSeatLookup(...)")
void fleetMcpConstructionStillWiresLeadSeatLookup() throws Exception {
String source = fleetdSource();
assertTrue(source.contains("new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), "
+ "leaders, leads))"),
"FleetMcp's construction call must still pass a LeadSeatSource built from "
+ "Fleetd.leadSeatLookup(...). Dropping it or swapping in "
+ "FleetMcp.LeadSeatSource.none() (fleetd #176's would-be silent regression, the same "
+ "shape as fleetd #248's measured mutations) compiles with 0 errors and leaves every "
+ "existing behavioural test green — this source check is what must go red instead.");
}
}
@@ -0,0 +1,64 @@
package dev.ltms.fleet;
import dev.ltms.fleet.inject.CompletionResolver;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.session.MemberSession;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import java.util.ArrayList;
import java.util.List;
import java.util.function.Function;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
/**
* fleetd #248: {@link Fleetd#worktreeBranchLookup} is the factory that replaced the anonymous
* lambda {@code Fleetd.main} used to build inline, as the 8th (last) argument to {@code
* CompletionResolver}'s constructor. Before this ticket that argument was untestable wiring:
* replacing it with {@code _ -> null} compiled clean and every existing test stayed green, because
* every existing test builds its own {@code CompletionResolver} directly rather than going through
* {@code main}. This class proves the factory's own behaviour; {@code
* FleetdCompletionResolverWiringTest} proves {@code main} still passes it in.
*/
class FleetdWorktreeBranchLookupTest {
private static MemberSession session(String terminal, String worktree, String branch) {
return new MemberSession("pane-" + terminal, terminal, "terra", MemberRole.DEV,
"/cwd", null, 0L, 0L, 0, MemberSession.State.READY, worktree, branch);
}
@Test
@DisplayName("a known target resolves to its session's worktree and branch")
void knownTargetResolves() {
Function<String, CompletionResolver.WorktreeBranch> lookup =
Fleetd.worktreeBranchLookup(() -> List.of(session("term1", "/wt/worker_x", "worker/x")));
CompletionResolver.WorktreeBranch resolved = lookup.apply("term1");
assertEquals("/wt/worker_x", resolved.worktree());
assertEquals("worker/x", resolved.branch());
}
@Test
@DisplayName("an unknown target resolves to null, not a thrown exception")
void unknownTargetResolvesToNull() {
Function<String, CompletionResolver.WorktreeBranch> lookup =
Fleetd.worktreeBranchLookup(() -> List.of(session("term1", "/wt/worker_x", "worker/x")));
assertNull(lookup.apply("term_stranger"));
}
@Test
@DisplayName("the roster is read through the supplier on every call, not snapshotted")
void rosterIsReadThroughOnEveryCall() {
List<MemberSession> roster = new ArrayList<>();
Function<String, CompletionResolver.WorktreeBranch> lookup = Fleetd.worktreeBranchLookup(() -> roster);
assertNull(lookup.apply("term_late"));
roster.add(session("term_late", "/wt/late", "worker/late"));
assertEquals("/wt/late", lookup.apply("term_late").worktree());
}
}
@@ -386,6 +386,54 @@ class ConfigRefTest {
assertEquals("rate limit exceeded", ref.get().profiles().get("sonnet").exhaustedPattern());
}
/**
* fleetd #201 Unit 5: errorPattern is compiled once into Fleetd.main's backend-error pattern
* map at startup (see BackendErrorPatternLookup), the same way exhaustedPattern is above — a
* reload never re-reads it either, so a changed value must be reported deferred.
*/
@Test
void changingAProfilesErrorPatternIsReportedAsDeferred(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
sonnet:
baseUrl: http://gx00.gw:8000
model: sonnet
errorPattern: "credential outage"
guard:
offSubscriptionHosts:
- gx00.gw
""");
ConfigRef ref = refFor(f);
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
sonnet:
baseUrl: http://gx00.gw:8000
model: sonnet
errorPattern: "provider 5xx"
guard:
offSubscriptionHosts:
- gx00.gw
""");
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertEquals(1, out.deferred().size(), out.deferred().toString());
assertTrue(out.deferred().getFirst().contains("sonnet"), out.deferred().toString());
assertTrue(out.deferred().getFirst().contains("launch settings"), out.deferred().toString());
// The snapshot still carries the new value — a restart is what makes it take effect.
assertEquals("provider 5xx", ref.get().profiles().get("sonnet").errorPattern());
}
@Test
void aFixedRefHasNoFileAndRefusesToReload() {
FleetConfig cfg = new FleetConfig(null, null, null, null, null, null,
@@ -86,6 +86,85 @@ class FleetConfigTest {
"unset means off — today's behaviour, unchanged");
}
// ── fleetd #201 Unit 5: errorPattern ────────────────────────────────────────────────────────
@Test
void aProfileWithAnErrorPatternBindsAndHasErrorPatternIsTrue(@TempDir Path dir) throws Exception {
Path f = dir.resolve("error-pattern.yaml");
Files.writeString(f, """
profiles:
ltms-local:
baseUrl: http://gx00.gw:8000
errorPattern: "(?i)\\\\bcredential outage\\\\b"
""");
FleetConfig.Profile w = FleetConfig.load(f).profiles().get("ltms-local");
assertTrue(w.hasErrorPattern(), "errorPattern binds and enables fleetd #201 Unit 5 classification");
assertEquals("(?i)\\bcredential outage\\b", w.errorPattern());
}
@Test
void aProfileWithNoErrorPatternLeavesItNullAndHasErrorPatternIsFalse(@TempDir Path dir) throws Exception {
Path f = dir.resolve("no-error-pattern.yaml");
Files.writeString(f, """
profiles:
ltms-local:
baseUrl: http://gx00.gw:8000
""");
FleetConfig.Profile w = FleetConfig.load(f).profiles().get("ltms-local");
assertNull(w.errorPattern(), "unset means 'use the built-in fallback' — never off");
assertFalse(w.hasErrorPattern());
}
@Test
void aBlankErrorPatternNormalizesToNullJustLikeUnset(@TempDir Path dir) throws Exception {
Path f = dir.resolve("blank-error-pattern.yaml");
Files.writeString(f, """
profiles:
ltms-local:
baseUrl: http://gx00.gw:8000
errorPattern: " "
""");
FleetConfig.Profile w = FleetConfig.load(f).profiles().get("ltms-local");
assertNull(w.errorPattern());
assertFalse(w.hasErrorPattern());
}
@Test
void aProfileWithAMalformedErrorPatternIsRejectedAtLoadNamingTheProfileAndKey(@TempDir Path dir)
throws Exception {
Path f = dir.resolve("malformed-error-pattern.yaml");
Files.writeString(f, """
profiles:
ltms-local:
baseUrl: http://gx00.gw:8000
errorPattern: "(unterminated["
""");
IllegalStateException e = assertThrows(IllegalStateException.class, () -> FleetConfig.load(f));
assertTrue(e.getMessage().contains("ltms-local"), "the offending profile is named: " + e.getMessage());
assertTrue(e.getMessage().contains("errorPattern"), "the offending key is named: " + e.getMessage());
}
@Test
void withProfileCarriesErrorPatternThrough(@TempDir Path dir) throws Exception {
Path f = dir.resolve("with-profile-error-pattern.yaml");
Files.writeString(f, """
profiles:
ltms-local:
baseUrl: http://gx00.gw:8000
errorPattern: "credential outage"
""");
// The compact constructor defaults `profile` to the map key via withProfile(name) — see
// loadsFullConfig above. errorPattern must survive that same rebuild.
FleetConfig.Profile w = FleetConfig.load(f).profiles().get("ltms-local");
assertEquals("ltms-local", w.profile());
assertEquals("credential outage", w.errorPattern());
}
@Test
void appliesDefaultsForMissingSections(@TempDir Path dir) throws Exception {
Path f = dir.resolve("minimal.yaml");
@@ -1313,6 +1392,7 @@ class FleetConfigTest {
weight: 0.5
maxLoad: 2
exhaustedPattern: "usage limit has been reached"
errorPattern: "credential outage"
placement: weighted
lifecycle:
idleTtlSeconds: 300
@@ -1350,6 +1430,8 @@ class FleetConfigTest {
assertEquals(2, w.maxLoad(), "maxLoad binds as an integer");
assertTrue(w.hasExhaustedPattern(), "exhaustedPattern binds and enables the CB-578 stage A classification");
assertEquals("usage limit has been reached", w.exhaustedPattern());
assertTrue(w.hasErrorPattern(), "errorPattern binds and enables fleetd #201 Unit 5 classification");
assertEquals("credential outage", w.errorPattern());
assertEquals("weighted", cfg.placement(), "placement binds at the top level");
assertEquals(300, cfg.lifecycle().idleTtlSeconds());
@@ -1781,7 +1863,7 @@ class FleetConfigTest {
}
@Test
void parityOverlayDefaultsToEnvFilesOnly(@TempDir Path dir) throws Exception {
void parityOverlayDefaultsToDotEnvOnly(@TempDir Path dir) throws Exception {
Path f = dir.resolve("no-overlay.yaml");
Files.writeString(f, """
bind:
@@ -1792,9 +1874,11 @@ class FleetConfigTest {
""");
FleetConfig cfg = FleetConfig.load(f);
assertEquals(List.of(".env", ".envrc"),
assertEquals(List.of(".env"),
cfg.profiles().get("gx10").parityOverlay(),
"the default parity overlay is the env files; settings.local.json is no longer copied by default");
"the default parity overlay is .env only (CB-148): .envrc is executable shell that "
+ "direnv runs on every cd, so it is no longer copied by default; settings.local.json "
+ "is also not copied by default");
}
@Test
@@ -1815,6 +1899,25 @@ class FleetConfigTest {
"an operator's explicit list survives verbatim — the default only changes when unset");
}
@Test
void parityOverlayExplicitEnvrcOptInStillWorks(@TempDir Path dir) throws Exception {
Path f = dir.resolve("explicit-envrc-overlay.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
gx10:
baseUrl: http://gx10.gw:8000
parityOverlay: [".env", ".envrc"]
""");
FleetConfig cfg = FleetConfig.load(f);
assertEquals(List.of(".env", ".envrc"),
cfg.profiles().get("gx10").parityOverlay(),
"an operator can still opt into copying .envrc explicitly (CB-148); it is only "
+ "dropped from the unset default, not removed as a capability");
}
// ── CB-542: subscription:true must not smuggle an unguarded endpoint via env: ───────────────
@Test
@@ -49,6 +49,7 @@ public final class FakeHerdr implements HerdrClient {
private int pinnedStarts = 0; // how many upcoming agent.start calls report a fixed pane
private String pinnedStartTerminal;
private String pinnedStartPane;
private Runnable onAgentStart; // fires the instant agent.start is called — see onAgentStart(Runnable)
private volatile int agentGetOkCalls = Integer.MAX_VALUE; // how many agent.get calls succeed first
private volatile String agentGetFailCode = null; // error code every agent.get call after that reports
@@ -152,6 +153,18 @@ public final class FakeHerdr implements HerdrClient {
return this;
}
/**
* Run {@code hook} synchronously the instant an {@code agent.start} call reaches this fake —
* i.e. the instant the peer PROCESS would start against a real herdr daemon. A test uses this
* to assert something is already true at that exact point (rather than merely true once
* {@code spawn()} returns), e.g. fleetd #149's trust-dialog seed having already been written to
* disk before the process herdr would launch ever starts.
*/
public FakeHerdr onAgentStart(Runnable hook) {
this.onAgentStart = hook;
return this;
}
/**
* Seed a named agent into {@code agent.list} (e.g. an orphaned worker for CB-117 reaper tests).
@@ -250,6 +263,9 @@ public final class FakeHerdr implements HerdrClient {
case "agent.read" -> mapper.readTree(mapper.writeValueAsString(
java.util.Map.of("type", "agent_read", "read", java.util.Map.of("text", readText))));
case "agent.start" -> {
if (onAgentStart != null) {
onAgentStart.run();
}
// Protocol 19: kind and pane_id are required — reject like the real daemon.
java.util.Map<?, ?> p = params instanceof java.util.Map<?, ?> m ? m : java.util.Map.of();
for (String required : new String[]{"kind", "pane_id"}) {
@@ -0,0 +1,295 @@
package dev.ltms.fleet.inject;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.mcp.PrimaryRegistry;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.member.CompositePeerLauncher;
import dev.ltms.fleet.msg.InMemoryReplyInbox;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.msg.ReplyPushLoop;
import dev.ltms.fleet.peer.SpawnRequest;
import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementException;
import dev.ltms.fleet.placement.PlacementPolicies;
import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.SessionManager;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
import java.util.regex.Pattern;
import static org.junit.jupiter.api.Assertions.*;
/**
* fleetd #201 / #227 Unit 5: proves the WIRED production path — a real {@link CompletionResolver}
* pane-scrape classification, through the exact {@code BackendErrorSink} ordering {@code Fleetd.main}
* builds (mirrored here line for line since the lambda itself lives inline in {@code Fleetd.main}),
* into a real {@link BackendOutagePolicy}, a real {@link CompositePeerLauncher} spawn gate, and a
* real {@link ReplyPushLoop} lead nudge — never by calling {@link BackendOutagePolicy#record} or the
* sink directly, which is what the finer-grained unit tests in
* {@code dev.ltms.fleet.member.CompositePeerLauncherTest} and {@code ReplyPushLoopTest} do for their
* own concerns. This class exists to catch a wiring mistake those unit tests cannot see: each of them
* hands the collaborator a value it assumes was computed correctly one layer up.
*
* <p>{@code fleet_list}/{@code fleet_profiles} rendering of the resulting cool-off (criteria 11/12)
* is proven separately in {@code dev.ltms.fleet.mcp.FleetMcpTest} — {@code FleetMcp}'s view methods
* are package-private to {@code dev.ltms.fleet.mcp} and this class needs {@code CompletionResolver}'s
* package-private {@code InFlight}, so the two cannot share one test class. What matters for THIS
* test is only that the same {@link BackendOutagePolicy} object the real scrape path populates is the
* one {@code FleetMcp} reads from — proven by asserting directly against
* {@link BackendOutagePolicy#remainingCoolOffSeconds} after the real classification below.
*/
class BackendOutageFlowTest {
private final List<ScheduledExecutorService> schedulers = new ArrayList<>();
@AfterEach
void tearDown() {
schedulers.forEach(ScheduledExecutorService::shutdownNow);
}
private static FleetConfig.Profile stubWorker(String profile, String credentialId) {
return new FleetConfig.Profile(profile, "http://gx00.gw:8000", "coder",
null, "FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers",
"w #{n}", null, null, null, null, null, null, null, null, null,
null, null, credentialId, null);
}
/** Minimal recording {@code HerdrClient} for the LEAD pane — mirrors ReplyPushLoopTest's own. */
private static final class RecordingLeadClient implements HerdrClient {
private static final ObjectMapper MAPPER = new ObjectMapper();
private final List<Object> prompts = new CopyOnWriteArrayList<>();
volatile CountDownLatch sendLatch = new CountDownLatch(1);
@Override
public JsonNode call(String method, Object params) {
if ("agent.get".equals(method)) {
return MAPPER.createObjectNode().set("agent", MAPPER.createObjectNode()
.put("terminal_id", "term_primary").put("agent_status", "idle"));
}
if ("agent.prompt".equals(method)) {
prompts.add(params);
sendLatch.countDown();
}
return MAPPER.createObjectNode();
}
@Override
public void close() {
}
int sendCount() {
return prompts.size();
}
String lastText() {
return ((Map<?, ?>) prompts.get(prompts.size() - 1)).get("text").toString();
}
}
/** Every collaborator the flow wires together, built exactly once per test. */
private final class Flow {
final FakeHerdr herdr = new FakeHerdr()
.readText("⏺ 503 Service Unavailable: upstream credential rejected\n❯ ");
final Map<String, FleetConfig.Profile> profiles = orderedProfiles();
final AtomicLong clockNanos = new AtomicLong(0L);
final BackendOutagePolicy outagePolicy = new BackendOutagePolicy(clockNanos::get);
final CompositePeerLauncher workers;
final SessionManager sessions;
final MemberSession s1;
final MemberSession s2;
final PrimaryRegistry registry = new PrimaryRegistry(null);
final InMemoryReplyInbox inbox = new InMemoryReplyInbox();
final RecordingLeadClient leadClient = new RecordingLeadClient();
final ReplyPushLoop pushLoop;
final CompletionResolver resolver;
final Rendezvous rendezvous = new Rendezvous();
Flow() {
ClaudeCodeLauncher adapter = new ClaudeCodeLauncher(new AgentControl(herdr),
new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
profiles, "terra", _ -> "tok");
workers = new CompositePeerLauncher(List.of(adapter), "terra", profiles,
PlacementPolicies.weighted(), _ -> 0, null, BackendQuarantine.none(), outagePolicy);
sessions = new SessionManager(workers);
// Two REAL member sessions on "terra" — the roster lookup the production sink does.
s1 = sessions.acquire("terra", null, null, null);
s2 = sessions.acquire("terra", null, null, null);
registry.recordDelegation(s1.terminalId(), "term_primary");
registry.recordDelegation(s2.terminalId(), "term_primary");
inbox.own(s1.terminalId());
inbox.own(s2.terminalId());
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
schedulers.add(scheduler);
pushLoop = new ReplyPushLoop(registry, new AgentControl(leadClient), inbox, scheduler, 3, 50);
AtomicReference<ReplyPushLoop> pushLoopRef = new AtomicReference<>(pushLoop);
// --- mirrors Fleetd.main's backendErrorSink lambda EXACTLY: (1) mark BACKEND_ERROR,
// (2) resolve profile/credential via the roster, fail-loud + notify unmapped-target,
// (3) record in BackendOutagePolicy, (4) on a NEW incident, notify the lead. ------------
BackendErrorSink backendErrorSink = (target, matchedLine, reason) -> {
sessions.onBackendError(target, reason);
String profileName = sessions.roster().stream()
.filter(session -> target.equals(session.terminalId()))
.findFirst()
.map(MemberSession::profile)
.orElse(null);
FleetConfig.Profile profile = profileName == null ? null : profiles.get(profileName);
if (profile == null) {
ReplyPushLoop loop = pushLoopRef.get();
if (loop != null) {
loop.onBackendTargetUnmapped(target, reason);
}
return;
}
String credentialId = profile.effectiveCredentialId();
Optional<BackendOutagePolicy.Incident> incident = outagePolicy.record(credentialId, target, reason);
incident.ifPresent(inc -> {
List<String> affectedProfiles = profiles.values().stream()
.filter(p -> credentialId.equals(p.effectiveCredentialId()))
.map(FleetConfig.Profile::profile)
.sorted()
.toList();
ReplyPushLoop loop = pushLoopRef.get();
if (loop != null) {
loop.onBackendIncident(inc.id(), inc.targets(), credentialId, affectedProfiles,
(int) inc.remainingCoolOffSeconds());
}
});
};
BackendErrorPatternLookup patterns = target -> Pattern.compile("(?i)503 Service Unavailable");
resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(),
ExhaustionSink.none(), patterns, backendErrorSink);
}
/** Drives one real pane-scrape classification for {@code target} through {@code resolver}. */
void classifyBackendErrorOn(String target) {
var waiter = rendezvous.open(target);
resolver.resolve(target, new CompletionResolver.InFlight(waiter, null));
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
"a classified backend error still resolves the blocked send as a failure");
}
}
private static Map<String, FleetConfig.Profile> orderedProfiles() {
Map<String, FleetConfig.Profile> m = new LinkedHashMap<>();
m.put("terra", stubWorker("terra", "shared-openai"));
m.put("sol", stubWorker("sol", "shared-openai"));
return m;
}
private long agentStartCalls(FakeHerdr herdr) {
return herdr.calls.stream().filter(c -> c.method().equals("agent.start")).count();
}
// --- criteria 4/5: two DISTINCT targets start an incident; one target never does ------------
@Test
void aSingleClassifiedBackendErrorOnOneTargetNeverCoolsTheCredential() {
Flow flow = new Flow();
flow.classifyBackendErrorOn(flow.s1.terminalId());
assertTrue(flow.outagePolicy.remainingCoolOffSeconds("shared-openai").isEmpty(),
"one distinct target's classified error must not start a cool-off");
assertEquals(0, flow.leadClient.sendCount(), "no incident yet, so no lead nudge");
assertDoesNotThrow(() -> flow.workers.spawn(new SpawnRequest("terra", null, null)),
"terra must still be spawnable after only one target's error");
}
@Test
void twoDistinctTargetsClassifiedThroughTheRealScrapeStartAnIncidentAndCoolTheCredential() throws Exception {
Flow flow = new Flow();
flow.classifyBackendErrorOn(flow.s1.terminalId());
assertTrue(flow.outagePolicy.remainingCoolOffSeconds("shared-openai").isEmpty());
flow.classifyBackendErrorOn(flow.s2.terminalId());
assertTrue(flow.leadClient.sendLatch.await(3, TimeUnit.SECONDS),
"the second distinct target must cross the threshold and nudge the lead");
Thread.sleep(150);
var remaining = flow.outagePolicy.remainingCoolOffSeconds("shared-openai");
assertTrue(remaining.isPresent(), "two distinct targets must start a cool-off");
assertEquals(60L, remaining.getAsLong(),
"this is the exact fact FleetMcp.capacityView/profiles reads as coolingOffForSeconds "
+ "(see FleetMcpTest for the rendering side)");
}
// --- criterion 13: the incident reaches the lead through ReplyPushLoop exactly once ----------
@Test
void theIncidentSendsExactlyOneNudgeToTheLeadNamingBothCredentialAndTargets() throws Exception {
Flow flow = new Flow();
flow.classifyBackendErrorOn(flow.s1.terminalId());
flow.classifyBackendErrorOn(flow.s2.terminalId());
assertTrue(flow.leadClient.sendLatch.await(3, TimeUnit.SECONDS));
Thread.sleep(200);
assertEquals(1, flow.leadClient.sendCount(), "one incident must send exactly one nudge");
String nudge = flow.leadClient.lastText();
assertTrue(nudge.contains("shared-openai"), nudge);
assertTrue(nudge.contains(flow.s1.terminalId()) && nudge.contains(flow.s2.terminalId()), nudge);
assertTrue(nudge.contains("60"), nudge);
}
// --- criterion 6: explicit spawn is refused BEFORE the adapter is ever called -----------------
@Test
void explicitSpawnOnACooledCredentialIsRefusedBeforeReachingTheAdapter() throws Exception {
Flow flow = new Flow();
flow.classifyBackendErrorOn(flow.s1.terminalId());
flow.classifyBackendErrorOn(flow.s2.terminalId());
assertTrue(flow.leadClient.sendLatch.await(3, TimeUnit.SECONDS));
long startsBefore = agentStartCalls(flow.herdr);
PlacementException e = assertThrows(PlacementException.class,
() -> flow.workers.spawn(new SpawnRequest("terra", null, null)));
assertTrue(e.getMessage().contains("cooling off"), e.getMessage());
assertTrue(e.getMessage().contains("shared-openai"), e.getMessage());
assertEquals(startsBefore, agentStartCalls(flow.herdr),
"the refusal must happen before the adapter is ever called — no new agent.start");
}
// --- criteria 7/9: automatic placement skips every profile sharing the cooled credential -----
@Test
void automaticPlacementRefusesAndNamesCoolingOffWhenEveryCandidateSharesTheCredential() throws Exception {
Flow flow = new Flow();
flow.classifyBackendErrorOn(flow.s1.terminalId());
flow.classifyBackendErrorOn(flow.s2.terminalId());
assertTrue(flow.leadClient.sendLatch.await(3, TimeUnit.SECONDS));
PlacementException e = assertThrows(PlacementException.class,
() -> flow.workers.spawn(new SpawnRequest(null, null, null)),
"terra and sol both share the cooled-off credential — nothing is available");
assertTrue(e.getMessage().contains("cooling off"), e.getMessage());
assertThrows(PlacementException.class, () -> flow.workers.spawn(new SpawnRequest("sol", null, null)),
"sol shares terra's credential, so it must be locked out too");
}
}
@@ -768,19 +768,37 @@ class CompletionResolverTest {
@Test
void coverageIsOffWhenNoProfileHasAPatternConfigured() {
assertEquals("off (no profile has an exhaustedPattern configured; profiles: [terra])",
CompletionResolver.coverage(Set.of("terra"), Set.of()));
CompletionResolver.coverage("exhaustedPattern", Set.of("terra"), Set.of()));
}
@Test
void coverageIsFullWhenEveryProfileHasAPatternConfigured() {
assertEquals("full (all profiles configured: [gx10, terra])",
CompletionResolver.coverage(Set.of("terra", "gx10"), Set.of("terra", "gx10")));
CompletionResolver.coverage("exhaustedPattern", Set.of("terra", "gx10"), Set.of("terra", "gx10")));
}
@Test
void coverageIsPartialAndNamesWhichProfilesAreConfigured() {
assertEquals("partial (configured: [terra]; not configured: [gx10])",
CompletionResolver.coverage(Set.of("terra", "gx10"), Set.of("terra")));
CompletionResolver.coverage("exhaustedPattern", Set.of("terra", "gx10"), Set.of("terra")));
}
/**
* Found live on 2026-09-03, reading a real boot log rather than a test. {@code coverage} is
* shared by two call sites — CB-578's {@code exhaustedPattern} line and fleetd#201 Unit 5's
* {@code errorPattern} line — but its "off" branch hard-coded the word {@code exhaustedPattern}.
* So a daemon with no {@code errorPattern} anywhere printed "no profile has an exhaustedPattern
* configured" directly beneath a line reporting that two profiles DO have one. Both lines were
* individually defensible and together they were nonsense, and the message sent an operator to
* set the wrong key.
*
* <p>Every earlier test here passed the exhaustion case only, so none of them could see it. This
* one pins that the message names the key the caller actually meant.
*/
@Test
void coverageNamesTheConfigKeyItsCallerMeansRatherThanAlwaysSayingExhaustedPattern() {
assertEquals("off (no profile has an errorPattern configured; profiles: [gx10, terra])",
CompletionResolver.coverage("errorPattern", Set.of("terra", "gx10"), Set.of()));
}
// --- fleetd#201 Unit 1: target-keyed backend-error pattern + typed sink ----------------------
@@ -20,6 +20,7 @@ import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.WorktreeRequest;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.member.CompositePeerLauncher;
import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementPolicies;
import io.modelcontextprotocol.spec.McpSchema;
@@ -500,6 +501,25 @@ class FleetMcpTest {
assertTrue(out.contains("\"quarantinedForSeconds\":1800"), out);
}
/** fleetd #201 Unit 5: {@code coolingOff} is a SEPARATE map from {@code quarantined}. */
@Test
void profilesReportsACoolingOffCredentialInASeparateMap() {
FakeHerdr h = new FakeHerdr();
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L);
outagePolicy.record("shared-openai", "t1", "API Error: rate limited");
outagePolicy.record("shared-openai", "t2", "API Error: rate limited"); // 2nd distinct target starts the incident
FleetMcp.OutageSource outage = new FleetMcp.OutageSource(
profile -> "ltms-local".equals(profile) ? "shared-openai" : null, outagePolicy);
McpSchema.CallToolResult res = FleetMcp.profiles(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), FleetMcp.QuarantineSource.none(), outage);
assertNotEquals(Boolean.TRUE, res.isError());
String out = textOf(res);
assertTrue(out.contains("\"coolingOff\""), out);
assertTrue(out.contains("shared-openai"), out);
assertTrue(out.contains("\"coolingOffForSeconds\":60"), out);
assertFalse(out.contains("\"quarantined\""), "nothing is exhaustion-quarantined here: " + out);
}
@Test
void listReportsTrackedWorkers() {
FakeHerdr h = new FakeHerdr();
@@ -631,6 +651,123 @@ class FleetMcpTest {
assertEquals(2, out.split("\"free\":0", -1).length - 1, out);
}
/**
* fleetd #176 stage 2 (correcting the inert stage 1): two {@code subscription: true} profiles,
* {@code opus} and {@code sonnet}, neither setting an explicit {@code credentialId} — the exact
* shape measured on the live Mac fleet. This is INTENDED, not a regression: a real Claude usage
* limit on the one login behind both profiles really does take out every profile running on it,
* the same way {@code credentialId: openai-shared} already lets two OpenAI-backed profiles share
* one quarantine (see {@code everyProfileSharingTheQuarantinedCredentialReportsZeroFree} above).
* The {@code credentialIdFor} function here is built the same way {@code Fleetd.main} wires it —
* {@code profile -> profiles.get(profile).effectiveCredentialId()} — so this proves the actual
* config-driven behaviour, not just {@code capacityView}'s arithmetic with a hand-picked string.
*/
@Test
void quarantiningOneSubscriptionProfileZeroesFreeOnTheOtherSharingTheSameAccount() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
Map<String, FleetConfig.Profile> profiles = Map.of(
"opus", new FleetConfig.Profile("opus", null, "claude-opus-4", null, null, null,
"tab", "fleet", "w #{n}", null, null, null, null, null, null, null,
null, 3, true, null, null, null),
"sonnet", new FleetConfig.Profile("sonnet", null, "claude-sonnet-5", null, null, null,
"tab", "fleet", "w #{n}", null, null, null, null, null, null, null,
null, 3, true, null, null, null));
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(20));
quarantine.quarantine(FleetConfig.Profile.SUBSCRIPTION_CREDENTIAL_ID);
FleetMcp.QuarantineSource source = new FleetMcp.QuarantineSource(profile -> {
FleetConfig.Profile configured = profiles.get(profile);
return configured == null ? null : configured.effectiveCredentialId();
}, quarantine);
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("opus", "sonnet"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"),
source, Map.of(), ""));
assertEquals(2, out.split("\"free\":0", -1).length - 1,
"opus and sonnet share one Claude login with neither setting credentialId, so quarantining "
+ "opus's account must also zero sonnet's free — this is intended, not a side effect: "
+ out);
assertEquals(2, out.split("\"credentialId\":\"" + FleetConfig.Profile.SUBSCRIPTION_CREDENTIAL_ID + "\"", -1)
.length - 1, out);
}
/** fleetd #201 Unit 5: cool-off forces {@code free:0} but never adds {@code quarantinedForSeconds}. */
@Test
void coolingOffProfileReportsZeroFreeButNeverQuarantinedForSeconds() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L);
outagePolicy.record("shared-openai", "t1", "API Error: rate limited");
outagePolicy.record("shared-openai", "t2", "API Error: rate limited");
FleetMcp.OutageSource outage = new FleetMcp.OutageSource(
profile -> "terra".equals(profile) ? "shared-openai" : null, outagePolicy);
String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
sessions, null, new FleetMcp.CapacitySource(profile -> 0, profile -> 2,
() -> Set.of("terra"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), outage, Map.of(), ""));
assertTrue(out.contains("\"profile\":\"terra\""), out);
assertTrue(out.contains("\"free\":0"), out);
assertTrue(out.contains("\"credentialId\":\"shared-openai\""), out);
assertTrue(out.contains("\"coolingOffForSeconds\":60"), out);
assertFalse(out.contains("quarantinedForSeconds"),
"exhaustion quarantine was never active — this key must not appear: " + out);
}
/**
* The two checks are independent — a profile can be BOTH exhaustion-quarantined AND cooling off
* for the same credential at once, and both maps/fields report it simultaneously.
*/
@Test
void bothQuarantineAndCoolingOffCanFireForTheSameProfileAtOnce() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(20));
quarantine.quarantine("shared-openai");
FleetMcp.QuarantineSource quarantineSource = new FleetMcp.QuarantineSource(
profile -> "terra".equals(profile) ? "shared-openai" : null, quarantine);
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L);
outagePolicy.record("shared-openai", "t1", "API Error: rate limited");
outagePolicy.record("shared-openai", "t2", "API Error: rate limited");
FleetMcp.OutageSource outage = new FleetMcp.OutageSource(
profile -> "terra".equals(profile) ? "shared-openai" : null, outagePolicy);
String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
sessions, null, new FleetMcp.CapacitySource(profile -> 0, profile -> 2,
() -> Set.of("terra"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"),
quarantineSource, outage, Map.of(), ""));
assertTrue(out.contains("\"free\":0"), out);
assertTrue(out.contains("\"quarantinedForSeconds\":1200"), out);
assertTrue(out.contains("\"coolingOffForSeconds\":60"), out);
}
@Test
void listAndProfilesAgreeOnWhatIsCoolingOff() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L);
outagePolicy.record("shared-openai", "t1", "API Error: rate limited");
outagePolicy.record("shared-openai", "t2", "API Error: rate limited");
FleetMcp.OutageSource outage = new FleetMcp.OutageSource(
profile -> "ltms-local".equals(profile) ? "shared-openai" : null, outagePolicy);
var workers = workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"));
String listOut = textOf(FleetMcp.listFleet(workers, sessions, null,
new FleetMcp.CapacitySource(profile -> 0, profile -> 2, () -> Set.of("ltms-local"), () -> 0),
new FleetMcp.HealthCoverageSource(() -> "off"), FleetMcp.QuarantineSource.none(), outage,
Map.of(), ""));
String profilesOut = textOf(FleetMcp.profiles(workers, FleetMcp.QuarantineSource.none(), outage));
assertTrue(listOut.contains("\"free\":0"), listOut);
assertTrue(listOut.contains("\"credentialId\":\"shared-openai\""), listOut);
assertTrue(profilesOut.contains("\"coolingOff\""), profilesOut);
assertTrue(profilesOut.contains("shared-openai"), profilesOut);
}
@Test
void listAndProfilesAgreeOnWhatIsQuarantined() {
FakeHerdr h = new FakeHerdr();
@@ -672,6 +809,68 @@ class FleetMcpTest {
assertFalse(out.contains("quarantinedForSeconds"), out);
}
/**
* fleetd #176: this is the exact shape measured on the Mac fleet — {@code maxLoad:3, live:2},
* where one of the "free" three is really the lead's own seat on the same subscription. The old
* formula ({@code max(0, cap - live)}) reported {@code free:1}; the real ceiling is {@code 0}
* (two members plus the lead's own seat already fill all three).
*/
@Test
void leadSeatSubtractsFromFreeTheSameWayLiveDoes() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
FleetMcp.LeadSeatSource leadSeats = new FleetMcp.LeadSeatSource(
profile -> "sonnet".equals(profile) ? 1 : 0);
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));
assertTrue(out.contains("\"maxLoad\":3"), "maxLoad itself must be left untouched: " + out);
assertTrue(out.contains("\"live\":2"), out);
assertTrue(out.contains("\"free\":0"), "2 live + 1 lead seat fills all 3: " + out);
assertTrue(out.contains("\"leadSeats\":1"), out);
}
/**
* fleetd #176: the OTHER measurement in the issue — a completely idle fleet still overstates
* {@code free} by the lead's own seat. {@code maxLoad:3, live:0} must report {@code free:2}, the
* real fan-out ceiling, not {@code 3}.
*/
@Test
void leadSeatLowersFreeOnAnOtherwiseIdleSubscriptionProfile() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
FleetMcp.LeadSeatSource leadSeats = new FleetMcp.LeadSeatSource(
profile -> "sonnet".equals(profile) ? 1 : 0);
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));
assertTrue(out.contains("\"live\":0"), out);
assertTrue(out.contains("\"free\":2"), "an idle fleet's real ceiling is 3 minus the lead's own seat: " + out);
assertTrue(out.contains("\"leadSeats\":1"), out);
}
/** A profile with no lead seats reported must be byte-identical to before this ticket. */
@Test
void zeroLeadSeatsOmitsTheKeyAndLeavesFreeUnchanged() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
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));
assertTrue(out.contains("\"free\":2"), out);
assertFalse(out.contains("leadSeats"), "no lead shares this profile's credential: " + out);
}
@Test
void listReportsLeadsAndFlagsTheCallersOwnRow() {
FakeHerdr h = new FakeHerdr();
@@ -4,6 +4,9 @@ import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.GuardException;
import dev.ltms.fleet.guard.SubscriptionGuard;
@@ -23,11 +26,16 @@ import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.attribute.PosixFileAttributeView;
import java.nio.file.attribute.PosixFilePermissions;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Function;
import java.util.function.Supplier;
@@ -2098,4 +2106,433 @@ class ClaudeCodeLauncherTest {
assertTrue(charterFile.getFileName().toString().startsWith("fleetd-role-charter-"),
"same file-naming scheme as before this fix (no wrapping directory): " + charterFile);
}
// --- fleetd #149: workspace-trust dialog seed ------------------------------------------------
//
// Claude Code asks an interactive, un-timed "Is this a project you created or one you trust?"
// the first time it starts in a directory it has not seen. Every worktree: true spawn lands in
// a brand-new directory, so without a seed the member sits on that dialog forever — herdr still
// reports it healthy (agent_status: blocked, interactive_ready: true) — and never mounts the
// bridge MCP or calls fleet_reply. These tests start the REAL launcher (only the herdr transport
// is faked) so the seed is proven to run inside buildLaunch, before agent.start (the
// process-starting call) ever fires — a test that only checked the JSON writer in isolation
// would prove nothing about whether the launcher actually calls it at the right time.
//
// INCIDENT: the seed originally ran on ANY non-blank cwd. Running this file's own test suite —
// most of whose fixtures spawn with no cwd/configDir set, so both fall back to the real
// user.dir / ~/.claude.json — corrupted the operator's actual ~/.claude.json (it shrank from
// ~72 KB to a single seeded entry) the first time a manual mutation run made the write
// non-additive. The fix restricts the seed to isProvisionedWorktree(cwd) — a real .git FILE
// (not directory) — exactly the writeIdeOverlay gate already used for the same category of
// risk. Every test below marks its own worktree fixture with that .git file, and
// seedTrustDialogNeverWritesWhenCwdIsNotAProvisionedWorktree is the regression test for the
// incident itself.
private FleetConfig.Profile trustProfile(String configDir, String cwd) {
return new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", configDir, "FLEETD_WORKER_TOKEN",
List.of("claude"), "tab", "fleetd-workers", "w #{n}", "http://127.0.0.1:8765/mcp",
cwd, null);
}
/**
* Give {@code dir} the exact signature {@link ClaudeCodeLauncher#isProvisionedWorktree} (and
* {@code writeIdeOverlay} before it) checks for: a {@code .git} REGULAR FILE, never a
* directory. The content is never parsed by the trust seed, so any {@code gitdir:} pointer is
* fine.
*/
private static void markAsProvisionedWorktree(Path dir) throws IOException {
Files.writeString(dir.resolve(".git"), "gitdir: /tmp/not-a-real-gitdir");
}
@Test
void seedsWorkspaceTrustForTheCwdBeforeTheProcessStarts(
@TempDir Path configDir, @TempDir Path worktree) throws Exception {
markAsProvisionedWorktree(worktree);
FakeHerdr herdr = new FakeHerdr();
Path claudeJson = configDir.resolve(".claude.json");
AtomicReference<Boolean> seededBeforeStart = new AtomicReference<>(false);
herdr.onAgentStart(() -> {
try {
if (!Files.exists(claudeJson)) {
return;
}
JsonNode root = new ObjectMapper().readTree(claudeJson.toFile());
JsonNode project = root.path("projects").path(worktree.toString());
seededBeforeStart.set(project.path("hasTrustDialogAccepted").asBoolean(false)
&& project.path("hasCompletedProjectOnboarding").asBoolean(false));
} catch (IOException e) {
seededBeforeStart.set(false);
}
});
FleetConfig.Profile cfg = trustProfile(configDir.toString(), worktree.toString());
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
.spawn();
assertTrue(herdr.called("agent.start"), "the hook must actually have fired during the spawn");
assertEquals(Boolean.TRUE, seededBeforeStart.get(),
"the trust entry for the cwd must already exist at the instant agent.start (the "
+ "process-starting herdr call) fires — not merely once spawn() returns");
}
@Test
void seedTrustDialogWritesBothTrustFlagsForTheResolvedCwd(
@TempDir Path configDir, @TempDir Path worktree) throws Exception {
markAsProvisionedWorktree(worktree);
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = trustProfile(configDir.toString(), worktree.toString());
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
.spawn();
Path claudeJson = configDir.resolve(".claude.json");
assertTrue(Files.exists(claudeJson), "seeded into <configDir>/.claude.json");
JsonNode project = new ObjectMapper().readTree(claudeJson.toFile())
.path("projects").path(worktree.toString());
assertTrue(project.path("hasTrustDialogAccepted").asBoolean(false));
assertTrue(project.path("hasCompletedProjectOnboarding").asBoolean(false));
}
/**
* Criterion 3: seeding is additive. An existing {@code .claude.json} carries the operator's own
* project history and unrelated top-level settings — the seed must change only
* {@code projects.<cwd>} for THIS cwd and leave everything else, including a different project's
* own unrelated data, exactly as it was.
*/
@Test
void seedTrustDialogIsAdditiveAndPreservesUnknownKeysAndOtherProjects(
@TempDir Path configDir, @TempDir Path worktree) throws Exception {
markAsProvisionedWorktree(worktree);
Path claudeJson = configDir.resolve(".claude.json");
Files.writeString(claudeJson, """
{
"numStartups": 42,
"oauthAccount": {"emailAddress": "operator@example.com"},
"projects": {
"/some/other/project": {
"hasTrustDialogAccepted": true,
"mcpServers": {"foo": {"type": "stdio", "command": "foo-mcp"}}
}
}
}
""");
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = trustProfile(configDir.toString(), worktree.toString());
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
.spawn();
JsonNode root = new ObjectMapper().readTree(claudeJson.toFile());
assertEquals(42, root.path("numStartups").asInt(), "unrelated top-level key survives untouched");
assertEquals("operator@example.com", root.path("oauthAccount").path("emailAddress").asText(),
"an unrelated nested top-level key survives untouched");
JsonNode other = root.path("projects").path("/some/other/project");
assertTrue(other.path("hasTrustDialogAccepted").asBoolean(false),
"a different project's own trust entry survives");
assertEquals("foo-mcp", other.path("mcpServers").path("foo").path("command").asText(),
"a different project's own unrelated nested data survives");
JsonNode mine = root.path("projects").path(worktree.toString());
assertTrue(mine.path("hasTrustDialogAccepted").asBoolean(false));
assertTrue(mine.path("hasCompletedProjectOnboarding").asBoolean(false));
}
/**
* Criterion 2: where the profile sets no {@code configDir}, the seed goes to the default
* {@code ~/.claude.json}. {@code user.home} is redirected to a {@code @TempDir} for the
* duration of this test and restored in a {@code finally} — the real operator {@code
* ~/.claude.json} must never be touched by a test.
*/
@Test
void seedTrustDialogTargetsDefaultClaudeJsonWhenConfigDirIsUnset(
@TempDir Path fakeHome, @TempDir Path worktree) throws Exception {
markAsProvisionedWorktree(worktree);
String originalHome = System.getProperty("user.home");
System.setProperty("user.home", fakeHome.toString());
try {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = trustProfile(null, worktree.toString());
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
.spawn();
Path claudeJson = fakeHome.resolve(".claude.json");
assertTrue(Files.exists(claudeJson),
"no configDir set — the default target is ~/.claude.json, here the redirected fake home");
JsonNode project = new ObjectMapper().readTree(claudeJson.toFile())
.path("projects").path(worktree.toString());
assertTrue(project.path("hasTrustDialogAccepted").asBoolean(false));
} finally {
System.setProperty("user.home", originalHome);
}
}
/**
* Regression test for the fleetd #149 incident itself: a cwd that is NOT a fleetd-provisioned
* worktree (no {@code .git} FILE — the exact shape a real checkout, or an un-configured
* fallback cwd, has) must never be written to, however {@code configDir} is set. This is the
* fix for the exact defect that corrupted the operator's real {@code ~/.claude.json}.
*/
@Test
void seedTrustDialogNeverWritesWhenCwdIsNotAProvisionedWorktree(
@TempDir Path configDir, @TempDir Path plainCwd) {
// plainCwd deliberately carries NO .git file — the same shape a real checkout's cwd
// fallback has.
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = trustProfile(configDir.toString(), plainCwd.toString());
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
.spawn();
assertFalse(Files.exists(configDir.resolve(".claude.json")),
"a non-worktree cwd must never get a .claude.json written for it — this is the fix "
+ "for the incident where writing unconditionally corrupted the operator's own "
+ "real ~/.claude.json via this file's own no-cwd/no-configDir test fixtures");
}
/**
* Same regression, for the default (no {@code configDir}) path — the exact combination (no
* {@code configDir}, no worktree-shaped {@code cwd}) that hit the operator's real
* {@code ~/.claude.json} during the incident. {@code user.home} is still redirected to a
* {@code @TempDir} out of caution, so even a reintroduced bug here cannot touch the real file.
*/
@Test
void seedTrustDialogNeverWritesToDefaultHomeWhenCwdIsNotAProvisionedWorktree(
@TempDir Path fakeHome, @TempDir Path plainCwd) {
String originalHome = System.getProperty("user.home");
System.setProperty("user.home", fakeHome.toString());
try {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = trustProfile(null, plainCwd.toString());
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
.spawn();
assertFalse(Files.exists(fakeHome.resolve(".claude.json")),
"the exact incident combination — no configDir, non-worktree cwd — must never "
+ "write, even to the (redirected) default ~/.claude.json");
} finally {
System.setProperty("user.home", originalHome);
}
}
// --- fleetd #149 review round 2: the write must be atomic and lock-protected -----------------
//
// The first round fixed WHICH cwd this can ever target. This round fixes HOW the target is
// written: .claude.json is large (tens of KB, dozens of projects on a real host), Claude Code
// itself rewrites it while running, and this daemon spawns several members in parallel as a
// matter of routine. The pre-fix Files.writeString(target, content) truncates target in place
// before writing the replacement bytes — a crash mid-write, or another writer's read landing in
// that window, loses data. Two more failure modes follow directly: (a) a crash/kill mid-write
// leaves target truncated, and (b) two concurrent spawns racing a naive read-modify-write let
// the second writer's write silently discard the first spawn's entry. The fix is
// ClaudeCodeLauncher.writeAtomically (sibling temp file + ATOMIC_MOVE, package-visible for the
// test below that proves it directly) plus TRUST_JSON_LOCK (a process-wide lock serialising
// every seedTrustDialog call this launcher itself makes).
/**
* Requested test 1: an existing, large-ish (not just a two-key fixture) .claude.json must never
* collapse. Asserts on the restored KEY SET — not merely that the result still parses as JSON,
* which the incident's 178-byte file also did.
*/
@Test
void seedTrustDialogPreservesALargeExistingFileWithoutCollapsing(
@TempDir Path configDir, @TempDir Path worktree) throws Exception {
markAsProvisionedWorktree(worktree);
Path claudeJson = configDir.resolve(".claude.json");
ObjectMapper mapper = new ObjectMapper();
ObjectNode root = mapper.createObjectNode();
root.put("numStartups", 4200);
root.put("firstStartTime", "2025-01-01T00:00:00.000Z");
root.putObject("oauthAccount").put("emailAddress", "operator@example.com");
Set<String> otherProjectPaths = new HashSet<>();
ObjectNode projects = root.putObject("projects");
for (int i = 0; i < 30; i++) {
String path = "/Users/operator/code/project-" + i;
otherProjectPaths.add(path);
ObjectNode project = projects.putObject(path);
project.put("hasTrustDialogAccepted", true);
project.putObject("mcpServers").putObject("server-" + i).put("command", "server-" + i + "-mcp");
}
String before = mapper.writerWithDefaultPrettyPrinter().writeValueAsString(root);
Files.writeString(claudeJson, before);
long sizeBefore = Files.size(claudeJson);
assertTrue(sizeBefore > 4096,
"fixture must actually be large-ish to be a meaningful proof: " + sizeBefore + " bytes");
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = trustProfile(configDir.toString(), worktree.toString());
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
.spawn();
long sizeAfter = Files.size(claudeJson);
assertTrue(sizeAfter >= sizeBefore,
"the file must never collapse below its pre-seed size — before=" + sizeBefore
+ " after=" + sizeAfter + " bytes (the incident shrank ~72 KB to 178 bytes)");
JsonNode after = mapper.readTree(claudeJson.toFile());
assertEquals(4200, after.path("numStartups").asInt(), "unrelated top-level key survives");
assertEquals("operator@example.com", after.path("oauthAccount").path("emailAddress").asText());
JsonNode afterProjects = after.path("projects");
for (String path : otherProjectPaths) {
assertTrue(afterProjects.path(path).path("hasTrustDialogAccepted").asBoolean(false),
"pre-existing project entry " + path + " must survive");
}
assertEquals(otherProjectPaths.size() + 1, afterProjects.size(),
"exactly one NEW project entry (this worktree's) must be added, none dropped");
assertTrue(afterProjects.path(worktree.toString()).path("hasTrustDialogAccepted").asBoolean(false));
}
/**
* Requested test 2: two concurrent spawns for DIFFERENT cwd values sharing one configDir must
* both end up present in the final file — a naive concurrent read-modify-write would let the
* second writer's read (taken before the first writer's write lands) silently discard the
* first. A {@link CountDownLatch} — not a sleep — lines both threads up at the starting line so
* this does not depend on scheduling luck to be meaningful.
*
* <p>This is a test of {@code TRUST_JSON_LOCK}, not of {@code writeAtomically}: the lock fully
* serialises every {@code seedTrustDialog} call this launcher itself makes, so this test would
* pass even without atomicity. See {@link #writeAtomicallyNeverExposesATornFileToAConcurrentReader}
* for the test that exercises atomicity specifically.
*/
@Test
void concurrentSeedsForDifferentCwdsBothSurvive(
@TempDir Path configDir, @TempDir Path worktreeA, @TempDir Path worktreeB) throws Exception {
markAsProvisionedWorktree(worktreeA);
markAsProvisionedWorktree(worktreeB);
CountDownLatch ready = new CountDownLatch(2);
CountDownLatch go = new CountDownLatch(1);
AtomicReference<Exception> failureA = new AtomicReference<>();
AtomicReference<Exception> failureB = new AtomicReference<>();
Thread ta = new Thread(spawnTask(configDir, worktreeA, ready, go, failureA), "spawn-a");
Thread tb = new Thread(spawnTask(configDir, worktreeB, ready, go, failureB), "spawn-b");
ta.start();
tb.start();
assertTrue(ready.await(5, TimeUnit.SECONDS), "both threads must reach the starting line");
go.countDown();
ta.join(5000);
tb.join(5000);
assertFalse(ta.isAlive(), "spawn A must finish within the timeout");
assertFalse(tb.isAlive(), "spawn B must finish within the timeout");
assertNull(failureA.get(), "spawn A must not throw: " + failureA.get());
assertNull(failureB.get(), "spawn B must not throw: " + failureB.get());
JsonNode root = new ObjectMapper().readTree(configDir.resolve(".claude.json").toFile());
assertTrue(root.path("projects").path(worktreeA.toString())
.path("hasTrustDialogAccepted").asBoolean(false),
"worktree A's entry must survive the race");
assertTrue(root.path("projects").path(worktreeB.toString())
.path("hasTrustDialogAccepted").asBoolean(false),
"worktree B's entry must survive the race");
}
private Runnable spawnTask(Path configDir, Path worktree, CountDownLatch ready, CountDownLatch go,
AtomicReference<Exception> failure) {
return () -> {
try {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = trustProfile(configDir.toString(), worktree.toString());
ClaudeCodeLauncher launcher = new ClaudeCodeLauncher(new AgentControl(herdr),
new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
ready.countDown();
go.await();
launcher.spawn();
} catch (Exception e) {
failure.set(e);
}
};
}
/**
* Requested test 3: the atomic write must preserve {@code .claude.json}'s existing {@code 0600}
* permissions, not silently widen them via a fresh temp file's own defaults landing on top of a
* file that had different (e.g. group-readable) permissions. Skips rather than fails on a
* filesystem with no POSIX permissions (e.g. Windows) — the same {@code assumeTrue} pattern this
* file already uses for {@code memberHerdrSocketWithWorktreeRootAndGroupPutsCharterUnderWorktreeRootAndSharesIt}.
*/
@Test
void seedTrustDialogPreservesExisting0600Permissions(
@TempDir Path configDir, @TempDir Path worktree) throws Exception {
markAsProvisionedWorktree(worktree);
Path claudeJson = configDir.resolve(".claude.json");
Files.writeString(claudeJson, "{}");
assumeTrue(Files.getFileAttributeView(claudeJson, PosixFileAttributeView.class) != null,
"no POSIX permissions on this filesystem — skipping rather than failing");
Files.setPosixFilePermissions(claudeJson, PosixFilePermissions.fromString("rw-------"));
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = trustProfile(configDir.toString(), worktree.toString());
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
.spawn();
assertEquals("rw-------", PosixFilePermissions.toString(Files.getPosixFilePermissions(claudeJson)),
"the atomic write must preserve .claude.json's existing 0600 permissions, not widen "
+ "them via a fresh temp file's own defaults");
}
/**
* Direct proof that {@link ClaudeCodeLauncher#writeAtomically} — not {@code TRUST_JSON_LOCK} —
* is what keeps a concurrent reader of {@code target} from ever observing a truncated or
* partially-written file. Calls {@code writeAtomically} directly (package-visible for exactly
* this test) rather than going through {@code seedTrustDialog}/{@code spawn()}, because {@link
* #concurrentSeedsForDifferentCwdsBothSurvive} above is protected by the lock and would pass
* even without atomicity — it proves the lock, not the atomic move. This test proves the atomic
* move specifically: a reader racing a writer OUTSIDE that lock (a second daemon process, or the
* operator's own live Claude Code — exactly what the lock cannot reach) must still never see a
* torn file.
*
* <p>The new content is made large (tens of MB) so a naive truncate-then-write has a real,
* non-instantaneous window for the busy-poll reader thread to land in — this is inherently a
* race, not a guaranteed-deterministic assertion, but it uses no {@code Thread.sleep} and
* reliably reproduced the torn read when run against the pre-fix
* {@code Files.writeString(target, content)} implementation (see the PR's mutation table).
*/
@Test
void writeAtomicallyNeverExposesATornFileToAConcurrentReader(@TempDir Path dir) throws Exception {
Path target = dir.resolve(".claude.json");
String oldContent = "{\"marker\":\"OLD\"}";
Files.writeString(target, oldContent);
String newContent = "{\"marker\":\"NEW\",\"pad\":\"" + "x".repeat(20_000_000) + "\"}";
AtomicReference<String> tornRead = new AtomicReference<>();
AtomicBoolean stop = new AtomicBoolean(false);
Thread reader = new Thread(() -> {
while (!stop.get()) {
try {
String seen = Files.readString(target);
if (!seen.equals(oldContent) && !seen.equals(newContent)) {
tornRead.compareAndSet(null, "torn read of length " + seen.length() + ": "
+ seen.substring(0, Math.min(seen.length(), 80)));
stop.set(true);
}
} catch (IOException ignored) {
// ATOMIC_MOVE guarantees the path always resolves to old-or-new content once
// readable at all — a transient "briefly missing during the rename" is fine to
// ignore and keep sampling.
}
}
});
reader.start();
try {
ClaudeCodeLauncher.writeAtomically(target, newContent);
} finally {
stop.set(true);
reader.join(5000);
}
assertNull(tornRead.get(), "a concurrent reader must never observe a partially-written file: "
+ tornRead.get());
assertEquals(newContent, Files.readString(target), "the final content must be the new content");
}
}
@@ -17,6 +17,7 @@ import dev.ltms.fleet.peer.PeerHandle;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.peer.PeerUnreachableException;
import dev.ltms.fleet.peer.SpawnRequest;
import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementException;
import dev.ltms.fleet.placement.PlacementPolicies;
@@ -911,4 +912,185 @@ class CompositePeerLauncherTest {
assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("claude", null, null)));
assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("gemini", null, null)));
}
// ── fleetd #201 Unit 5: repeated backend errors cool a credential off — a SEPARATE, shorter- ──
// ── lived mechanism from CB-578 stage B quarantine above, never merged with it ───────────────
@Test
void explicitSpawnOntoACoolingOffProfileIsRefusedNamingTheCredentialAndRemainingTime() {
FakeHerdr herdr = new FakeHerdr();
Map<String, FleetConfig.Profile> profiles = ordered(
"sol", stubWorker("sol", "shared-openai"),
"terra", stubWorker("terra", "shared-openai"));
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "sol", Set.of());
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L);
// Two DISTINCT targets on the same credential — evidenceCount() counts distinct targets,
// never raw events, so one target repeating a classified error can never mint an incident.
outagePolicy.record("shared-openai", "term-1", "API Error: 500");
assertTrue(outagePolicy.record("shared-openai", "term-2", "API Error: 500").isPresent(),
"the second distinct target crosses the threshold and starts the cool-off");
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(), outagePolicy);
PlacementException e = assertThrows(PlacementException.class,
() -> composite.spawn(new SpawnRequest("sol", null, null)));
assertTrue(e.getMessage().contains("sol"), "message names the profile: " + e.getMessage());
assertTrue(e.getMessage().contains("shared-openai"), "message names the credential: " + e.getMessage());
assertTrue(e.getMessage().contains("cooling off"), "message says cooling off, not exhausted: " + e.getMessage());
assertTrue(e.getMessage().contains("60"), "message names roughly when it lifts: " + e.getMessage());
assertEquals(0, adapter.spawnCount("sol"), "the cooling-off profile is never delegated to");
}
/**
* A single error from one target must never cool a credential off — the threshold is on
* distinct targets. This proves the spawn gate really reads {@link BackendOutagePolicy}'s
* result rather than reacting to any recorded event.
*/
@Test
void explicitSpawnIsNotRefusedAfterOnlyOneBackendErrorOnOneTarget() {
FakeHerdr herdr = new FakeHerdr();
Map<String, FleetConfig.Profile> profiles = ordered(
"sol", stubWorker("sol", "shared-openai"),
"terra", stubWorker("terra", "shared-openai"));
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "sol", Set.of());
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L);
assertTrue(outagePolicy.record("shared-openai", "term-1", "API Error: 500").isEmpty(),
"one distinct target is below THRESHOLD=2 — no incident yet");
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(), outagePolicy);
assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("sol", null, null)),
"one target's error never cools the credential off");
assertEquals(1, adapter.spawnCount("sol"));
}
/**
* sol and terra share one OpenAI credential; cooling off because of an outage classified via
* ONE of them must lock out the other too, exactly like CB-578 stage B quarantine does.
*/
@Test
void twoProfilesSharingACredentialAreBothCoolingOffByOneOutageIncident() {
FakeHerdr herdr = new FakeHerdr();
Map<String, FleetConfig.Profile> profiles = ordered(
"sol", stubWorker("sol", "shared-openai"),
"terra", stubWorker("terra", "shared-openai"));
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "sol", Set.of());
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L);
outagePolicy.record("shared-openai", "term-1", "API Error: 500");
outagePolicy.record("shared-openai", "term-2", "API Error: 500");
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(), outagePolicy);
assertThrows(PlacementException.class, () -> composite.spawn(new SpawnRequest("sol", null, null)));
assertThrows(PlacementException.class, () -> composite.spawn(new SpawnRequest("terra", null, null)),
"terra shares sol's credential, so it must be locked out too");
assertEquals(0, adapter.spawnCount("sol"));
assertEquals(0, adapter.spawnCount("terra"));
}
@Test
void placementSkipsACoolingOffProfileAndRoutesToAnotherOne() {
FakeHerdr herdr = new FakeHerdr();
Map<String, FleetConfig.Profile> profiles = ordered(
"sol", stubWorker("sol", "shared-openai"),
"b", stubWorker("b"));
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "sol", Set.of());
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L);
outagePolicy.record("shared-openai", "term-1", "API Error: 500");
outagePolicy.record("shared-openai", "term-2", "API Error: 500");
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
PlacementPolicies.weighted(), _ -> 0, null, BackendQuarantine.none(), outagePolicy);
PeerHandle h = composite.spawn(new SpawnRequest(null, null, null));
assertEquals("b", h.profile(), "sol is cooling off, so an unqualified spawn must land on b");
assertEquals(0, adapter.spawnCount("sol"));
}
@Test
void automaticPlacementNamesCoolingOffWhenEveryCandidateIsCoolingOff() {
FakeHerdr herdr = new FakeHerdr();
Map<String, FleetConfig.Profile> profiles = ordered(
"sol", stubWorker("sol", "shared-openai"),
"terra", stubWorker("terra", "shared-openai"));
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "sol", Set.of());
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L);
outagePolicy.record("shared-openai", "term-1", "API Error: 500");
outagePolicy.record("shared-openai", "term-2", "API Error: 500");
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
PlacementPolicies.weighted(), _ -> 0, null, BackendQuarantine.none(), outagePolicy);
// PlacementPolicy.select throws uncaught out of spawn() when available() is empty — never
// wrapped in PeerUnreachableException, which is only for a delegate that actually refused.
PlacementException e = assertThrows(PlacementException.class,
() -> composite.spawn(new SpawnRequest(null, null, null)),
"both candidates share the cooling-off credential — nothing is available");
assertTrue(e.getMessage().contains("cooling off"),
"message says cooling off, not exhausted, since nothing here is quarantined: " + e.getMessage());
}
@Test
void aCoolOffLiftsOnTheInjectedClockAndTheProfileBecomesSpawnableAgain() {
FakeHerdr herdr = new FakeHerdr();
Map<String, FleetConfig.Profile> profiles = ordered(
"sol", stubWorker("sol", "shared-openai"),
"b", stubWorker("b"));
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "sol", Set.of());
AtomicLong nowNanos = new AtomicLong(0L);
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(nowNanos::get);
outagePolicy.record("shared-openai", "term-1", "API Error: 500");
outagePolicy.record("shared-openai", "term-2", "API Error: 500");
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(), outagePolicy);
assertThrows(PlacementException.class, () -> composite.spawn(new SpawnRequest("sol", null, null)),
"still inside the 60s cool-off");
nowNanos.set(TimeUnit.SECONDS.toNanos(61));
PeerHandle h = composite.spawn(new SpawnRequest("sol", null, null));
assertEquals("sol", h.profile(), "the cool-off expired on the injected clock — sol is spawnable again");
assertEquals(1, adapter.spawnCount("sol"));
}
/**
* A credential that is BOTH exhaustion-quarantined (CB-578 stage B) and cooling off (fleetd
* #201 Unit 5) at once — the two checks are independent and can both be true for one profile.
* Exhaustion quarantine must win: the refusal names quarantine, never cooling off, because
* {@code enforceNotQuarantined} is checked (and throws) before {@code enforceNotCoolingOff} ever
* runs.
*/
@Test
void quarantineWinsOverCoolingOffWhenBothAreActiveOnTheSameCredential() {
FakeHerdr herdr = new FakeHerdr();
Map<String, FleetConfig.Profile> profiles = ordered(
"sol", stubWorker("sol", "shared-openai"),
"terra", stubWorker("terra", "shared-openai"));
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "sol", Set.of());
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
quarantine.quarantine("shared-openai");
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L);
outagePolicy.record("shared-openai", "term-1", "API Error: 500");
outagePolicy.record("shared-openai", "term-2", "API Error: 500");
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
PlacementPolicies.fixed(), _ -> 0, null, quarantine, outagePolicy);
PlacementException e = assertThrows(PlacementException.class,
() -> composite.spawn(new SpawnRequest("sol", null, null)));
assertTrue(e.getMessage().contains("quarantined"),
"exhaustion quarantine takes priority in the message: " + e.getMessage());
assertFalse(e.getMessage().contains("cooling off"),
"cooling off is never mentioned when quarantine already refused the spawn: " + e.getMessage());
}
@Test
void aFleetWithNoBackendErrorsEverRecordedBehavesExactlyAsBeforeCoolOffExisted() {
// A never-record()-called BackendOutagePolicy is naturally inert — the 2/5/6-arg
// constructors used throughout this file all wire this in implicitly (NO_OUTAGE_POLICY).
// This test pins that an explicit fresh instance also never refuses a spawn, for any profile.
FakeHerdr herdr = new FakeHerdr();
PeerLauncher composite = composite(herdr);
assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("claude", null, null)));
assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("gemini", null, null)));
}
}
@@ -30,7 +30,14 @@ class PlacementPolicyTest {
private static PlacementContext ctx(List<PlacementCandidate> candidates,
Function<String, Integer> liveCount,
Set<String> unreachable, Set<String> quarantined) {
return new PlacementContext("b", candidates, liveCount, unreachable, quarantined);
return ctx(candidates, liveCount, unreachable, quarantined, Set.of());
}
private static PlacementContext ctx(List<PlacementCandidate> candidates,
Function<String, Integer> liveCount,
Set<String> unreachable, Set<String> quarantined,
Set<String> coolingOff) {
return new PlacementContext("b", candidates, liveCount, unreachable, quarantined, coolingOff);
}
private static PlacementContext ctx(List<PlacementCandidate> candidates,
@@ -52,14 +59,14 @@ class PlacementPolicyTest {
PlacementPolicy policy = PlacementPolicies.fixed();
PlacementContext ctx = new PlacementContext(null,
List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")),
noSessions(), Set.of(), Set.of());
noSessions(), Set.of(), Set.of(), Set.of());
assertEquals("a", policy.select(ctx).profile());
}
@Test
void fixedThrowsWhenNoProfilesAndNoDefault() {
PlacementPolicy policy = PlacementPolicies.fixed();
PlacementContext ctx = new PlacementContext(null, List.of(), noSessions(), Set.of(), Set.of());
PlacementContext ctx = new PlacementContext(null, List.of(), noSessions(), Set.of(), Set.of(), Set.of());
assertThrows(PlacementException.class, () -> policy.select(ctx));
}
@@ -68,7 +75,7 @@ class PlacementPolicyTest {
PlacementPolicy policy = PlacementPolicies.fixed();
PlacementContext ctx = new PlacementContext("b",
List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")),
noSessions(), Set.of(), Set.of("b"));
noSessions(), Set.of(), Set.of("b"), Set.of());
assertEquals("a", policy.select(ctx).profile(),
"the default 'b' is quarantined, so fixed falls through to the first un-quarantined candidate");
}
@@ -78,7 +85,7 @@ class PlacementPolicyTest {
PlacementPolicy policy = PlacementPolicies.fixed();
PlacementContext ctx = new PlacementContext("b",
List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")),
noSessions(), Set.of(), Set.of("a", "b"));
noSessions(), Set.of(), Set.of("a", "b"), Set.of());
PlacementException e = assertThrows(PlacementException.class, () -> policy.select(ctx));
assertTrue(e.getMessage().contains("quarantined"), e.getMessage());
}
@@ -211,6 +218,67 @@ class PlacementPolicyTest {
assertTrue(e.getMessage().contains("quarantined"), e.getMessage());
}
// --- fleetd #201 Unit 5: coolingOff excludes a candidate, SEPARATE from quarantined --------
@Test
void roundRobinSkipsCoolingOffProfiles() {
PlacementPolicy policy = PlacementPolicies.roundRobin();
List<PlacementCandidate> candidates = List.of(
PlacementCandidate.profile("a"),
PlacementCandidate.profile("b"));
PlacementContext ctx = ctx(candidates, noSessions(), Set.of(), Set.of(), Set.of("a"));
for (int i = 0; i < 5; i++) {
assertEquals("b", policy.select(ctx).profile(), "a is cooling off, so every pick lands on b");
}
}
@Test
void weightedSkipsCoolingOffProfile() {
PlacementPolicy policy = PlacementPolicies.weighted();
List<PlacementCandidate> candidates = List.of(
PlacementCandidate.profile("a", 1.0f, null),
PlacementCandidate.profile("b", 1.0f, null));
PlacementContext ctx = ctx(candidates, noSessions(), Set.of(), Set.of(), Set.of("a"));
for (int i = 0; i < 5; i++) {
assertEquals("b", policy.select(ctx).profile(), "a is cooling off, so every pick lands on b");
}
}
@Test
void weightedThrowsWhenAllCoolingOff() {
PlacementPolicy policy = PlacementPolicies.weighted();
List<PlacementCandidate> candidates = List.of(
PlacementCandidate.profile("a"),
PlacementCandidate.profile("b"));
PlacementException e = assertThrows(PlacementException.class,
() -> policy.select(ctx(candidates, noSessions(), Set.of(), Set.of(), Set.of("a", "b"))));
assertTrue(e.getMessage().contains("cooling off"), e.getMessage());
assertFalse(e.getMessage().contains("quarantined"),
"nothing here is quarantined — the message must say cooling off, not exhausted: " + e.getMessage());
}
/**
* A candidate that is BOTH quarantined (CB-578 stage B) and cooling off (fleetd #201 Unit 5)
* counts only as quarantined — exhaustion takes priority, matching
* {@code CompositePeerLauncher}'s explicit-spawn check order.
*/
@Test
void quarantinedAndCoolingOffIsNotDoubleCounted() {
PlacementPolicy policy = PlacementPolicies.weighted();
List<PlacementCandidate> candidates = List.of(
PlacementCandidate.profile("a"),
PlacementCandidate.profile("b"));
// a is BOTH quarantined and cooling off; b is quarantined only. If a were double-bucketed
// as cooling-off instead of quarantined, quarantined would undercount to 1 of 2 candidates
// and this would fall through to the generic "no worker profile available: N quarantined, M
// cooling off, ..." message instead — which also happens to contain the substring
// "quarantined", so a loose contains() check here would pass either way. Pinning the exact
// "all ... quarantined" message is what actually proves the two are not double-counted.
PlacementContext ctx = ctx(candidates, noSessions(), Set.of(), Set.of("a", "b"), Set.of("a"));
PlacementException e = assertThrows(PlacementException.class, () -> policy.select(ctx));
assertEquals("all worker profiles are quarantined (backend exhausted)", e.getMessage());
}
@Test
void weightedThrowsWhenAllUnreachable() {
PlacementPolicy policy = PlacementPolicies.weighted();
@@ -320,7 +388,7 @@ class PlacementPolicyTest {
PlacementContext ctx = new PlacementContext("b",
List.of(PlacementCandidate.profile("a", 1.0f, null),
PlacementCandidate.profile("b", 0.0f, null)),
noSessions(), Set.of(), Set.of());
noSessions(), Set.of(), Set.of(), Set.of());
assertEquals("a", policy.select(ctx).profile(),
"the default 'b' has weight 0, so fixed falls through to the first available candidate");
}
@@ -331,7 +399,7 @@ class PlacementPolicyTest {
PlacementContext ctx = new PlacementContext("b",
List.of(PlacementCandidate.profile("a", 0.0f, null),
PlacementCandidate.profile("b", 0.0f, null)),
noSessions(), Set.of(), Set.of());
noSessions(), Set.of(), Set.of(), Set.of());
PlacementException e = assertThrows(PlacementException.class, () -> policy.select(ctx));
assertTrue(e.getMessage().contains("weight 0"), e.getMessage());
}
@@ -666,6 +666,103 @@ class GitWorktreesTest {
assertEquals("", status(Path.of(wt), ".autoenv"), ".autoenv still shows as modified");
}
// ---- fleetd #134: isolateToolSurface must report what it neutralized (to the daemon operator's
// log) and record it where the worker itself can read it (worktree-scoped git config), without
// ever showing up in the worker's own `git status`. All drive the real provisioning path,
// GitWorktrees#add, per criterion 5 — the whole provisioned worktree is what's under test here. ----
private static Path initRepoWithAllThreeConfigs(Path dir) throws Exception {
Files.createDirectories(dir);
git(dir, "init", "-q", "-b", "main");
git(dir, "config", "user.email", "test@example.invalid");
git(dir, "config", "user.name", "Test");
Files.writeString(dir.resolve(".mcp.json"), WITH_SERVERS);
Files.writeString(dir.resolve("opencode.json"), OPENCODE_WITH_FILE_REF);
Files.writeString(dir.resolve(".autoenv"), AUTOENV_WITH_DIRECTIVE);
Files.writeString(dir.resolve("README.md"), "seed\n");
git(dir, "add", ".mcp.json", "opencode.json", ".autoenv", "README.md");
git(dir, "commit", "-q", "-m", "seed");
return dir;
}
/** Criterion 1, all three present: the summary names the denominator and every neutralized file. */
@Test
void isolateToolSurfaceLogsAllThreeConfigsNeutralized(@TempDir Path tmp) throws Exception {
reportingLogger.setLevel(Level.INFO);
Path repo = initRepoWithAllThreeConfigs(tmp.resolve("repo"));
new GitWorktrees(tmp.resolve("wts").toString()).add(repo.toString(), "cb-134-log-all", "HEAD");
assertTrue(capturedMessages().contains(
"tool-surface isolation: neutralized 3 of 3 configs: .mcp.json, opencode.json, "
+ ".autoenv — the worktree copy is a stub, not the repo's file; edit the "
+ "real file in the primary checkout instead"),
"expected the all-neutralized summary line, got:\n" + capturedMessages());
}
/** Criterion 1, two absent: the summary must still name the denominator and say why. */
@Test
void isolateToolSurfaceLogsAbsentConfigsWithReason(@TempDir Path tmp) throws Exception {
reportingLogger.setLevel(Level.INFO);
Path repo = initRepo(tmp.resolve("repo")); // only .mcp.json + README committed
new GitWorktrees(tmp.resolve("wts").toString()).add(repo.toString(), "cb-134-log-partial", "HEAD");
assertTrue(capturedMessages().contains(
"tool-surface isolation: neutralized 1 of 3 configs: .mcp.json (opencode.json "
+ "absent, .autoenv absent) — the worktree copy is a stub, not the repo's "
+ "file; edit the real file in the primary checkout instead"),
"expected the partial summary line, got:\n" + capturedMessages());
}
/**
* Criterion 2. The daemon's own log is invisible to the worker process — it never reads fleetd's
* stdout. This is the mechanism the worker itself can query, from inside its own worktree, to
* learn "this file is neutralized here, the repo's real file differs" instead of trusting what it
* just read on disk.
*/
@Test
void aWorkerCanDiscoverNeutralizedConfigsFromWorktreeScopedGitConfig(@TempDir Path tmp) throws Exception {
Path repo = initRepoWithAllThreeConfigs(tmp.resolve("repo"));
String wt = new GitWorktrees(tmp.resolve("wts").toString())
.add(repo.toString(), "cb-134-discover", "HEAD");
String recorded = gitOutput(Path.of(wt), "config", "--worktree", "--get-all", "fleet.neutralizedConfig");
Set<String> files = new HashSet<>();
for (String line : recorded.split("\\R")) {
if (!line.isBlank()) {
files.add(line.trim());
}
}
assertEquals(Set.of(".mcp.json", "opencode.json", ".autoenv"), files,
"the worker-readable record must name every neutralized file: " + recorded);
String note = gitOutput(Path.of(wt), "config", "--worktree", "--get", "fleet.neutralizedConfigNote").trim();
assertTrue(note.contains("stub"), "note must say the worktree copy is a stub: " + note);
assertTrue(note.contains("primary checkout"),
"note must state the consequence — where to edit the real file instead: " + note);
}
/**
* Criterion 3. Whatever fleetd#134's worker-discovery mechanism writes must never appear as
* untracked or modified in the worker's own `git status` — a worker that sees a stray file either
* commits it by mistake or burns a turn asking about it. This runs the full porcelain status, not
* a single-file check, so any leftover file anywhere in the worktree would fail it.
*/
@Test
void aProvisionedWorktreeHasCleanGitStatusDespiteNeutralizedConfigRecordkeeping(@TempDir Path tmp)
throws Exception {
Path repo = initRepoWithAllThreeConfigs(tmp.resolve("repo"));
String wt = new GitWorktrees(tmp.resolve("wts").toString())
.add(repo.toString(), "cb-134-clean-status", "HEAD");
assertEquals("", fullStatus(Path.of(wt)),
"a freshly provisioned worktree must show a clean `git status --porcelain`, including "
+ "after the worker-readable neutralized-config record was written");
}
/**
* CB-578 stage C, acceptance criterion 1. A dirty worktree — a tracked edit plus a brand-new
* untracked file, exactly the shape lost in CB-576 — must land in {@code refs/wip/<branch>}'s
@@ -1177,4 +1274,112 @@ class GitWorktreesTest {
+ "the refusal must happen before `git worktree add` ever runs");
}
}
// ---- fleetd #134 / #148 point 3: overlayParity must report what it did, and marking a tracked
// file --skip-worktree must say the file can no longer be committed from this worktree. Drives
// overlayParity directly against a real worktree (git worktree add, no GitWorktrees#add) so these
// tests are independent of origin/credential-helper provisioning, which is not under test here. ----
/** A bare worktree, sibling to {@code repo}, created with plain git — the target overlayParity
* copies into. Deliberately not {@link GitWorktrees#add}: that method does unrelated
* provisioning (origin rewrite, .mcp.json neutralization, credential helper) that would only
* add noise to the log assertions below. */
private static Path bareWorktree(Path repo, Path wtDir, String branch) throws Exception {
git(repo, "worktree", "add", "-q", wtDir.toString(), "-b", branch, "HEAD");
return wtDir;
}
/** Acceptance criterion 1: only the configured candidates land in the worktree — nothing else
* from the source tree leaks in alongside them. */
@Test
void overlayParityCopiesExactlyTheConfiguredFilesAndNothingElse(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
Files.writeString(repo.resolve(".env"), "A=1\n");
Files.writeString(repo.resolve(".envrc"), "export A=1\n");
Files.writeString(repo.resolve("not-overlaid.txt"), "must not be copied\n");
Path wt = bareWorktree(repo, tmp.resolve("wt"), "cb134-exact");
new GitWorktrees(tmp.resolve("wts").toString())
.overlayParity(repo.toString(), wt.toString(), List.of(".env", ".envrc"));
assertEquals("A=1\n", Files.readString(wt.resolve(".env")));
assertEquals("export A=1\n", Files.readString(wt.resolve(".envrc")));
assertFalse(Files.exists(wt.resolve("not-overlaid.txt")),
"overlayParity must copy only the configured candidates, not the whole source tree");
}
/** Criterion 2: both candidates present — the summary line names both and the denominator. */
@Test
void overlayParityLogsBothCopiedWhenBothCandidatesArePresent(@TempDir Path tmp) throws Exception {
reportingLogger.setLevel(Level.INFO);
Path repo = initRepo(tmp.resolve("repo"));
Files.writeString(repo.resolve(".env"), "A=1\n");
Files.writeString(repo.resolve(".envrc"), "export A=1\n");
Path wt = bareWorktree(repo, tmp.resolve("wt"), "cb134-both");
new GitWorktrees(tmp.resolve("wts").toString())
.overlayParity(repo.toString(), wt.toString(), List.of(".env", ".envrc"));
assertTrue(capturedMessages().contains("parity overlay: copied 2 of 2 candidates: .env, .envrc"),
"expected the both-copied summary line, got:\n" + capturedMessages());
}
/** Criterion 2: one candidate present, one absent — the summary must name the copied file, the
* denominator, and why the other candidate was not copied. */
@Test
void overlayParityLogsOneCopiedOneAbsent(@TempDir Path tmp) throws Exception {
reportingLogger.setLevel(Level.INFO);
Path repo = initRepo(tmp.resolve("repo"));
Files.writeString(repo.resolve(".env"), "A=1\n");
// .envrc deliberately not created — the absent candidate.
Path wt = bareWorktree(repo, tmp.resolve("wt"), "cb134-partial");
new GitWorktrees(tmp.resolve("wts").toString())
.overlayParity(repo.toString(), wt.toString(), List.of(".env", ".envrc"));
assertTrue(Files.exists(wt.resolve(".env")));
assertFalse(Files.exists(wt.resolve(".envrc")));
assertTrue(capturedMessages().contains("parity overlay: copied 1 of 2 candidates: .env (.envrc absent)"),
"expected the copied/absent summary line, got:\n" + capturedMessages());
}
/** Criterion 3: a tracked candidate is marked --skip-worktree, and that must be named in the log
* with the consequence spelled out — a worker editing it afterward finds git ignoring the
* change, silently, unless this line told it so beforehand. */
@Test
void overlayParityLogsSkipWorktreeConsequenceForATrackedFile(@TempDir Path tmp) throws Exception {
reportingLogger.setLevel(Level.INFO);
Path repo = initRepo(tmp.resolve("repo"));
Files.writeString(repo.resolve(".env"), "A=1\n");
git(repo, "add", ".env");
git(repo, "commit", "-q", "-m", "track env");
Path wt = bareWorktree(repo, tmp.resolve("wt"), "cb134-tracked");
// Change the source after the worktree checkout, so the overlay copy actually overwrites it.
Files.writeString(repo.resolve(".env"), "A=2\n");
new GitWorktrees(tmp.resolve("wts").toString())
.overlayParity(repo.toString(), wt.toString(), List.of(".env"));
assertEquals("A=2\n", Files.readString(wt.resolve(".env")));
assertEquals("", status(wt, ".env"),
"the skip-worktree'd file must not show as modified even though its content changed");
assertTrue(capturedMessages().contains(
"parity overlay marked --skip-worktree (cannot be committed from this worktree): .env"),
"expected the skip-worktree consequence line, got:\n" + capturedMessages());
}
/** Criterion 5: null and empty overlay lists return quietly — no exception, no log noise. */
@Test
void overlayParityWithNoCandidatesLogsNothing(@TempDir Path tmp) throws Exception {
reportingLogger.setLevel(Level.INFO);
Path repo = initRepo(tmp.resolve("repo"));
Path wt = bareWorktree(repo, tmp.resolve("wt"), "cb134-empty");
GitWorktrees worktrees = new GitWorktrees(tmp.resolve("wts").toString());
worktrees.overlayParity(repo.toString(), wt.toString(), null);
worktrees.overlayParity(repo.toString(), wt.toString(), List.of());
assertTrue(reportingAppender.list.isEmpty(),
"a null/empty overlay must log nothing, got:\n" + capturedMessages());
}
}
@@ -140,7 +140,7 @@ class WorktreeSessionManagerTest {
void worktreeAcquireRunsParityOverlayWithProfileDefaults() {
FakeHerdr herdr = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt")
.track(".envrc")
.track(".env")
.exists(".claude/settings.local.json");
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
@@ -151,14 +151,15 @@ class WorktreeSessionManagerTest {
FakeWorktrees.OverlayCall overlay = worktrees.lastOverlay();
assertNotNull(overlay);
assertEquals("/repo", overlay.repoRoot());
assertEquals(List.of(".env", ".envrc"),
overlay.requested(), "default parity overlay is used when unset");
assertEquals(List.of(".env"),
overlay.requested(), "default parity overlay is used when unset (CB-148: .env only, "
+ ".envrc is no longer defaulted because it is executable shell direnv runs on cd)");
assertFalse(overlay.requested().contains(".mcp.json"),
"CB-525: replicating the primary's MCP config gives a worker the primary's IDE "
+ "servers, which navigate its edits out of its own worktree");
assertEquals(List.of(".envrc"), overlay.copied(),
assertEquals(List.of(".env"), overlay.copied(),
"existing paths are copied; missing paths are skipped");
assertEquals(List.of(".envrc"), overlay.skipWorktree(),
assertEquals(List.of(".env"), overlay.skipWorktree(),
"tracked copied paths are --skip-worktree'd");
}