Compare commits

..

12 Commits

Author SHA1 Message Date
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
22 changed files with 2306 additions and 94 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
+50 -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.
#
@@ -273,9 +310,10 @@ profiles:
# 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 +407,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 +439,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
@@ -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
@@ -342,6 +352,26 @@ public final class Fleetd {
.orElse(null);
log.info("backend-exhausted classification (CB-578 stage A): {}",
CompletionResolver.coverage(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()));
}
});
BackendErrorPatternLookup backendErrorPatterns = target -> sessions.roster().stream()
.filter(session -> target.equals(session.terminalId()))
.findFirst()
.map(session -> errorPatternsByProfile.get(session.profile()))
.orElse(null);
log.info("backend-error classification (fleetd #201 Unit 5): {}",
CompletionResolver.coverage(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,8 +419,59 @@ 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<>();
// Order: (1) mark the member BACKEND_ERROR; (2) resolve profile/credential through the
// roster — fail loud (never Optional.ifPresent, the fleetd #234 lesson applied to this new
// sink) and notify the lead via onBackendTargetUnmapped when it cannot be resolved; (3)
// record the error in BackendOutagePolicy; (4) on a NEW incident (the record() call that
// actually crosses the threshold), tell the lead via onBackendIncident.
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 : config.get().profiles().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 = 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 = config.get().profiles().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 = pushLoopRef.get();
if (loop != null) {
loop.onBackendIncident(inc.id(), inc.targets(), credentialId, affectedProfiles,
(int) inc.remainingCoolOffSeconds());
}
});
};
AgentControl agents = router.memberAgents();
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns, exhaustionSink,
// 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.
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns,
exhaustionSink, backendErrorPatterns, backendErrorSink, System::nanoTime,
target -> sessions.roster().stream()
.filter(session -> target.equals(session.terminalId()))
.findFirst()
@@ -477,6 +558,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 +668,11 @@ 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));
// 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
@@ -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);
}
/**
@@ -589,6 +634,15 @@ public record FleetConfig(
return exhaustedPattern != null;
}
/**
* 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 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
@@ -1389,6 +1443,7 @@ public record FleetConfig(
rejectDuplicateMemberSlots(yaml);
rejectNegativeMaxLoad(yaml);
rejectAutoCompactWindowOutOfRange(yaml);
rejectMalformedErrorPattern(yaml);
rejectUnknownKind(yaml);
rejectUnknownAuthMode(yaml);
rejectUnknownPlacement(yaml);
@@ -1737,6 +1792,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);
@@ -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,8 @@ 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;
/** CB-637: this daemon's lead-to-lead channel; {@code null} when no coordinator is configured. */
private final LeadChannel leadChannel;
@@ -116,6 +119,22 @@ 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));
}
}
/**
* @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 +149,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());
}
/**
@@ -143,9 +162,25 @@ 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());
}
/**
* As above, with fleetd #201 Unit 5 cool-off facts for {@code fleet_list}/{@code fleet_profiles}
* (see {@link OutageSource}). This is what {@code Fleetd.main} actually wires up.
*
* @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.leadChannel = leadChannel;
this.capacity = capacity;
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
this.outage = Objects.requireNonNull(outage, "outage");
this.healthCoverage = healthCoverage;
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
this.transport = HttpServletStreamableServerTransportProvider.builder()
@@ -274,7 +309,7 @@ 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,
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
callers == null ? Map.of() : callers.leads(),
callerTerminal(exchange),
leadChannel == null ? null : leadChannel.selfCoordId());
@@ -290,7 +325,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 +910,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 +989,14 @@ 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, 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 +1010,15 @@ 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(),
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) {
try {
Map<String, Agent> live = workers.list().stream()
.map(Agent.class::cast)
@@ -971,7 +1045,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)).toList());
return text(json(result));
} catch (HerdrException e) {
return error("herdr error listing the fleet: " + e.getMessage());
@@ -1004,10 +1078,17 @@ 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.
*/
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) {
Integer cap = maxLoad.apply(profile);
int live = liveCount.apply(profile);
int reclaimable = (int) roster.stream().filter(s -> profile.equals(s.profile()))
@@ -1025,6 +1106,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
@@ -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");
}
}
@@ -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,81 @@ class FleetMcpTest {
assertEquals(2, out.split("\"free\":0", -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();
@@ -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");
}