Compare commits

...

53 Commits

Author SHA1 Message Date
Dai Ha ab0cc71aa4 fleetd #418: barrier the throw-path push-loop test on state decide() reads
CI / contract (pull_request) Successful in 1m27s
CI / build (pull_request) Successful in 1m50s
anAskThatLeavesByThrowingStillClosesItsQuestion barriered on Phase.ASKING,
which markAsyncQuestion sets in ask()'s FIRST step. The assertion right
after it depends on ask()'s THIRD step (pushLoop.onQuestionOpened), which
is what actually populates ReplyPushLoop's pendingQuestions map. Under
load the asker thread can be descheduled between those two steps, so the
barrier released before decide() had anything to see, and it correctly
returned STOP instead of the expected INJECT.

Add ReplyPushLoop#pendingQuestionTurnIdsForTest, a package-private test
seam (modeled on MessageService#isCompletionStampedForTest) exposing the
private pendingQuestionTurnIdsFor. The test now waits for its own turnId
to appear there before asserting on decide() — not for decide() itself to
return INJECT, which would make the barrier assert nothing.

Checked every other awaitTicketPhaseOn(..., Phase.ASKING) in the file
(two, in the CB-582 nudge tests): both are followed by a real awaitNudge()
that waits for an actual agent.prompt push-loop call before any assertion
depends on push-loop state, so they are not exposed to this race.

No production behaviour changed.
2026-09-10 11:24:41 +07:00
ltms 4cd9046353 Merge pull request 'fleetd #409: deterministic test for the #399 completion-stamp ordering race' (#414) from worker/deterministic-stamp-race-409-3cb7b6-10 into main
CI / contract (push) Successful in 59s
CI / build (push) Successful in 1m36s
2026-09-10 04:48:44 +02:00
Dai Ha 4e98a74047 fleetd #409: deterministic test for the #399 completion-stamp ordering race
CI / contract (pull_request) Successful in 45s
CI / build (pull_request) Successful in 2m5s
Widens the completed-hook's real race window (normally instructions-wide, needing
~2x-core host load to hit by chance per #399) by injecting a bounded sleep into the
test clock's completion-stamp read. This makes the ordering invariant — a test must
wait for isCompletionStampedForTest, not just DONE, before advancing the clock past
the TTL — fail deterministically on the first run when the barrier is removed, and
pass deterministically with it present. No production code changed.
2026-09-10 09:39:08 +07:00
ltms a502ba53e0 Merge #404: the armed field reads the startup pattern map, both directions pinned
CI / contract (push) Successful in 49s
CI / build (push) Successful in 2m0s
The armed lookup now uses the same compiled map detection uses. Second test added during review after a mutation proved the true direction was unpinned.
2026-09-10 04:20:54 +02:00
Dai Ha 446cc11d1d fleetd #404: pin the armed field's TRUE direction too
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Successful in 1m51s
Mutation-tested during review: replacing the armed lambda with
`profile -> false` left the whole suite green at 1475 tests, because the
only existing test passes an EMPTY startup map. That mutation would make
#395's visibility feature silently dead.

With this test the same mutation fails, and it is the only test that
fails, so nothing else covers this direction.
2026-09-10 09:20:43 +07:00
ltms ddd81fe174 Merge #391: fleet_reply refuses a lead, and the role argument is now mandatory
CI / contract (push) Successful in 46s
CI / build (push) Successful in 2m1s
The permissive 3-arg overload is deleted, so a handler that forgets the role no longer compiles.
2026-09-10 04:18:36 +02:00
ltms ae74cc081f Merge #399: wait on a real completion stamp, not on DONE
CI / contract (push) Successful in 45s
CI / build (push) Successful in 1m29s
Adds a volatile completedNanos stamp and a bounded test barrier. Mutation-tested: removing the production stamp fails 3 tests.
2026-09-10 04:14:29 +02:00
Dai Ha cf8da1d5fa fleetd #391: require reply caller role
CI / contract (pull_request) Successful in 1m21s
CI / build (pull_request) Successful in 1m45s
2026-09-10 09:13:16 +07:00
Dai Ha b8aedeafcb fleetd #404: test armed startup map behavior
CI / contract (pull_request) Successful in 48s
CI / build (pull_request) Successful in 1m30s
2026-09-10 09:13:07 +07:00
Dai Ha 3d61af6f6f Merge #400: the scrub receipt now measures the blank, not the attempt
CI / contract (push) Successful in 1m8s
CI / build (push) Successful in 1m34s
The blanking loop classified a name as blanked from eval's exit status.
zsh coerces a bare NAME= assignment on an integer special parameter
(SECONDS, RANDOM, SHLVL, HISTSIZE, COLUMNS, LINES, USERNAME) to a number
instead of failing, so eval returned 0 with the value untouched. Measured
7 false receipts in 10 names. This is a security receipt, so a count that
overstates the scrub is worse than no count.

The loop now runs the eval unconditionally and decides from the observed
value, read back with the (P) indirection flag. One check covers all
three shapes a name can take: a real blank, a fatal error eval merely
contained, and this silent no-op. Exit status plays no part.

Lead review: mutation removing the '!' unblankable report line is CAUGHT
(2 failures in EnvAllowListScrubTest, both asserting the name is reported
rather than silently dropped). The eval-site identifier guard is
untouched.

Unplanted evidence the fix works: the #394 test
unblankableNameInTheMiddleDoesNotAbortNamesAfterIt began failing under
the fix, because zsh auto-exports SHLVL and the old exit-status bug had
been miscounting it as blanked all along. Its exact-count assertion had
only ever passed because of the bug beside it; it is now a presence
check, since which names a zsh version auto-exports is not this test's
to pin.

Still not pinned, tracked in #394's follow-up: the eval-site identifier
guard has no test behind it.
2026-09-10 09:05:30 +07:00
Dai Ha ed99c209ac Merge #398: a central allow-list of usable models, with the startup validators pinned
CI / contract (push) Successful in 50s
CI / build (push) Successful in 1m25s
models: allow: is a single place that names every model the fleet may
use. Absent or empty keeps today's behaviour, so this ships inert until
configured. Once set, a profile naming a model outside the list refuses
to start, and refuses a reload, rather than reaching a backend adapter
as a free-form string.

The allow-list cannot be checked against a provider catalogue: for an
opencode profile fleetd SYNTHESIZES the provider from provider/model
plus baseUrl (OpenCodeLauncher:562-583), so a valid fleetd model id
appears in no published catalogue. An operator-owned list is therefore
the only workable gate.

Includes the #398 follow-up: FleetConfig.validateAll() reflectively
sweeps this class's validateXxx() methods, and both real call sites
(Fleetd.main and ConfigRef.reload) call that one method. Before this,
deleting a validateXxx() call from either caller left the whole suite
green. FleetdStartupValidationTest now drives the real Fleetd.main.

Lead review: mutation on the reload call site is caught (ConfigRefTest,
2 failures). Mutation replacing the reflective sweep with a hardcoded
list is NOT caught (1491 green) — so the sweep is a convenience and the
denominator test is the real guarantee; two false statements in the test
javadoc were corrected to say so (af4c88d).

Recovered work: the worker's agent died mid-turn with the follow-up
uncommitted and the startup call left disabled as
'// MUTATION-TEST-TEMP: cfg.validateAll();'. I restored it before
committing.

NOT covered: the five log-only reportXxx(cfg) calls in main are still
unpinned — filed as #407.
2026-09-10 09:02:18 +07:00
Dai Ha af4c88d54b t398: correct two false statements in FleetConfigValidateAllTest
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Successful in 1m58s
Measured at review: reverting validateAll() to a hardcoded list of
today's six calls leaves the suite green (1491 tests, 0 failures). The
class javadoc claimed that mutation fails a test. It does not — claim 1
pins the generic helper on an unrelated class, claim 2 pins today's six,
and a hardcoded list satisfies both.

The interaction was the real hazard. The denominator assertion IS a
tripwire (declaring a seventh validator fails it), but its failure
message said the sweep reaches new validators 'by construction' and told
the author to just update the expected set. If the sweep were ever
replaced by a name list, the one assertion that fires would hand back a
false all-clear at the moment it fired.

Javadoc now states the measurement, and names the denominator test as
the actual guarantee. The assertion message now says to confirm
validateAll() still delegates to invokeAllValidators(this) BEFORE
updating the expected set.
2026-09-10 09:01:49 +07:00
Dai Ha 44c735f6f5 fleetd #404: report armed detection from startup map
CI / build (pull_request) Successful in 1m25s
CI / contract (pull_request) Successful in 1m20s
2026-09-10 08:56:39 +07:00
Dai Ha b540a1744b fleetd #398 follow-up: pin the startup validators with a reflective validateAll()
CI / contract (pull_request) Successful in 1m9s
CI / build (pull_request) Successful in 1m28s
Mutation testing found that deleting a cfg.validateXxx() call from
Fleetd.main left the full suite green: every test called a validator
directly and none exercised main as the caller.

FleetConfig.validateAll() sweeps this class's own public no-arg void
validateXxx() methods by reflection and invokes each in alphabetical
order, so a newly written validator is wired into both callers
(Fleetd.main and ConfigRef.reload) with no second step to forget.
FleetdStartupValidationTest calls the real Fleetd.main with six configs,
each failing exactly one validator.

Recovered by the lead: the worker's agent died mid-turn with this work
uncommitted, and had left the startup call commented out as
'// MUTATION-TEST-TEMP: cfg.validateAll();' from its own mutation run.
I restored the call before committing. Build after restoring:
Tests run: 1491, Failures: 0, BUILD SUCCESS.

NOT covered, and not claimed to be: the five log-only reporters in
main (reportRequiredSecrets, reportGitHostShape, reportMemberTrustModel,
reportMemberCredentialsGap, and reportExhaustedPatternGap on current
main) are not validateXxx() methods, so the sweep does not reach them
and their call sites stay unpinned.
2026-09-10 08:56:20 +07:00
Dai Ha 0f51d53098 fleetd #399: fix TTL test race by waiting on a real completion stamp, not DONE
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Successful in 2m2s
poll() can report Phase.DONE for a ticket before the whenComplete hook that
stamps Task.completedNanos has run — CompletableFuture.complete() publishes
its result and only then runs dependents. The TTL tests advanced an injected
clock right after observing DONE, so on a host where the hook runs late it
stamps the ADVANCED time and the eviction never happens (fails on Linux,
passes on macOS).

Add a package-private test seam, MessageService.isCompletionStampedForTest,
that reports whether completedNanos is stamped. Both TTL tests now wait on
that (a real volatile read/write happens-before edge) before advancing the
clock, instead of on Phase.DONE. The prune condition in pruneTerminalTickets
is untouched.
2026-09-10 08:54:38 +07:00
Dai Ha c670792ffe #400: classify the blanking loop's result on the observed value, not eval's exit status
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Successful in 2m1s
eval "export NAME=" can return success even when zsh coerces the bare
assignment on an integer special parameter (SECONDS, RANDOM, SHLVL,
HISTSIZE, COLUMNS, LINES, USERNAME) instead of failing, leaving the
value unchanged. The old exit-status check then reported the name as
blanked when it was not -- a false receipt.

Classify on the observed effect instead: attempt the export, then read
the name's value back with the (P) indirection flag and decide from
whether it is now empty. One check now covers all three shapes a name
can take here -- a genuine blank, a fatal read-only error eval merely
contains, and this silent no-op -- with the exit status playing no
part in the decision.

Adds a test driving all three shapes through the real scrubScript in
one run (a normal name, LINENO for the fatal case, SECONDS for the
silent no-op), with the parent environment explicitly carrying those
names since a cleared ProcessBuilder parent does not expose them on
its own. Also corrects the previously-merged
unblankableNameInTheMiddleDoesNotAbortNamesAfterIt test, whose "exactly
one failed name" assertion turned out to only pass by accident: zsh
itself auto-exports SHLVL on every shell start, and the old exit-status
bug was silently miscounting it as blanked. The fixed classification
now correctly reports it unblankable too, so the test asserts presence
rather than an exact count.
2026-09-10 08:46:27 +07:00
Dai Ha 3982ace544 fleetd #391: refuse lead fleet replies
CI / contract (pull_request) Successful in 1m11s
CI / build (pull_request) Successful in 1m31s
2026-09-10 08:43:54 +07:00
Dai Ha 7180b1aad0 Merge #395: warn when exhaustedPattern usage-limit detection is off
CI / contract (push) Successful in 1m22s
CI / build (push) Successful in 1m27s
A profile with no exhaustedPattern has usage-limit detection silently
disabled. Startup now reports every unarmed profile (a louder warning for
subscription profiles, which are the ones a limit actually stops), and
fleet_profiles / GET /profiles carry exhaustionDetectionArmed per profile.

Reviewed by the lead: all three requested mutations fail a test, and the
back-compat QuarantineSource ctor defaults to 'not armed' when the source
is unknown. KNOWN GAP, not fixed here: deleting the
reportExhaustedPatternGap(cfg) call at Fleetd.java:140 leaves the suite
green (Tests run: 1472, Failures: 0). That is the same unpinned-startup-
call shape as PR #398's six validators, and fleetd #398's ticket owns it.
2026-09-10 08:43:09 +07:00
Dai Ha c26f695402 fleetd #395: warn when exhaustedPattern detection is silently off
CI / contract (pull_request) Successful in 1m10s
CI / build (pull_request) Successful in 1m25s
exhaustedPattern is opt-in per profile: unset means a usage-limit
refusal on that profile is never classified BACKEND_EXHAUSTED and
never quarantines its credential, with nothing telling the operator.
Add a startup WARN naming every unarmed profile (a louder, separate
WARN for a subscription: true profile, since that is the operator's
own metered plan). Surface the same fact per profile in fleet_profiles
as exhaustionDetectionArmed, so an operator can tell "healthy" from
"can never be caught" without reading fleetd.yaml.
2026-09-10 08:34:33 +07:00
Dai Ha 48877315ca Merge #394: contain the fatal export that aborted the credential scrub
CI / contract (push) Successful in 50s
CI / build (push) Failing after 1m31s
The CB-633 allow-list scrub has been dying mid-loop on every fleet01 member
pane and saying nothing. `export UID=` in zsh is not a failed command — it is
a fatal parameter error that terminates the whole sourced file. The blanking
loop is wrapped in `{ ... } 2>/dev/null`, so the message was swallowed and the
report block after the loop never ran.

Root cause found by the fleet01 lead, with xtrace on a live pane's own ZDOTDIR:

    +scrub.zsh:28> _cb633_n=UID
    +scrub.zsh:28> export 'UID='
    +zsh:1> rc=126        <- file aborted

The severity is the INVERSION, and this is their finding, quoted:

  "env lists inherited names first and the names a startup file exports last.
   So the loop blanks the harmless inherited half and dies immediately before
   the operator's own exports — exactly the credentials the policy exists to
   remove. The selection is inverted, not merely partial."

Measured there: UID is name 42 of 57, and a ~/.zshrc decoy at 58 survived on
8 of 8 spawns. "Partial scrub" reads as "we got most of it"; it got precisely
the wrong half.

Fixed with `eval "export ${n}=" 2>/dev/null` rather than a skip-list of the
known-fatal names (UID EUID GID EGID PPID LINENO). A skip-list has to be
complete forever and this is a security control; eval needs no list. Measured:
plain export dies at UID and every later name keeps its value, while the eval
form completes the loop and blanks all of them. PR #396 proposed the skip-list
and is closed in favour of this; its claim that the abort happens "however the
assignment is wrapped" holds for a direct `if ! export` but not for eval,
which reparses in a nested context.

The report now carries `failed N` and `!`-prefixed unblankable names, and
HerdrPeerLauncher WARNs when any name could not be blanked. The old "no report"
WARN no longer claims the daemon knows what the member saw.

Why the suite stayed green: EnvAllowListScrubTest starts zsh from
pb.environment().clear(), and under a cleared parent UID is not an exported
name at all, so the abort could not reproduce in that harness.

Reviewed by mutation, which found a second gap now also closed: the eval is
only safe because names are filtered to ^[A-Za-z_][A-Za-z0-9_]*$. Replacing
that pattern with .* left the class green, so the line the security property
rests on was unpinned. The guard is now re-asserted at the eval site and pinned
by a test. The reachable vector is a VALUE with an embedded newline, not a
hostile name — measured: zsh strips non-identifier env names outright, while
MULTI=$'keep\njunk.fragment' forges 'junk.fragment' as a candidate name out of
its own value.

Closes #394. Refs #396, #388.
2026-09-10 08:30:15 +07:00
Dai Ha 5c08054533 #394 follow-up: re-assert the identifier guard at the eval call site
CI / contract (pull_request) Successful in 1m29s
CI / build (pull_request) Successful in 1m48s
EnvAllowListScrub's blanking loop splices each name into a string
handed to eval ("export ${n}="). That is only safe because every name
reaching _cb633_blank already passed an identifier check in the
enumeration loop -- 20 lines away, in a different loop. Before eval
was introduced a non-conforming name reaching plain `export "$n="`
was inert either way (the quoting neutralized it); eval removed that
safety net, so the enumeration loop's guard became the ONLY thing
standing between a non-identifier string and code execution in the
member's pane, with nothing at the eval site itself defending that
property.

Re-assert the same [A-Za-z_][A-Za-z0-9_]* check immediately before
the eval call, independent of the enumeration loop's own guard (left
untouched, not moved). A name that fails it is counted unblankable
rather than silently dropped, so a bypass of the upstream guard would
leave real evidence in the report.

New test exploits the "junk from multi-line values" gap the
enumeration loop's own comment already documents: a value with an
embedded newline makes `command env`'s text output split into a
spurious extra "name" line that was never a real variable. Runs the
real generated scrubScript() end-to-end under zsh and asserts the
non-conforming fragment is neither blanked nor counted unblankable.
The fragment used is merely non-conforming (contains a dot) --
never command-shaped.

Mutation-verified both guards. Weakening the enumeration guard alone
DOES break the new test (the fragment then reaches the new eval-site
guard and gets counted unblankable, failing the "not unblankable"
assertion). Removing the new eval-site guard alone, with the
enumeration guard intact, does NOT break it: _cb633_blank has exactly
one producer (the enumeration loop), so nothing can reach the eval
site without already having passed the identical check there. That is
expected given the single-source architecture, and it is exactly why
the eval-site guard is defense-in-depth against a future change that
adds a second path into _cb633_blank or decouples the two loops --
not a currently independently-observable divergence.
2026-09-10 08:24:46 +07:00
Dai Ha b6db9c31f5 charter: a peer lead is answered with fleet_send, not fleet_reply
CI / contract (push) Successful in 46s
CI / build (push) Successful in 2m0s
`fleet_reply` has no route to a peer lead. `AmqpReplyInbox` publishes to
`agent.<target>.inbox`, mandatory, and a lead's own terminal has no such
queue, so the publish is refused. `MessageService.reply()` has no peer
branch at all — `grep -c 'coord\|LeadMailbox'` on it returns 0. The charter
told every lead to use a tool that cannot work, and both leads here hit it.

Three edits to the canonical block, byte-identical with the wiki template
(pushed as 803726a; the in-sync check in this file reports True):

- the intent->tool row now says `fleet_send{coordId}`, or `{sessionId}` for
  a peer on the same host, and says plainly that `fleet_reply` is refused
- the prose says WHY: `fleet_reply` resolves a member's blocked `fleet_send`,
  while a peer's coord-id message is durable and non-blocking, so there is
  nothing for it to resolve
- lead<->lead item 3 gains the data-point rule: N observations are N data
  points only if they differ in the axis you are trusting

Wording for all three drafted by the fleet01 lead, who verified the missing
queue namespace independently in its own tree. The data-point rule has now
caught three separate errors in a day, in both directions: one cause blamed
for N failures, and N agreeing measurements that shared a single instrument.

The refusal message itself is still wrong — it says "queue not declared or
owned", which sends the reader to the broker instead of to this file. That
half stays open on #391.

Tracked as fleetd #391.
2026-09-10 08:22:35 +07:00
Dai Ha e7b33fe3a0 fleetd: central allow-list of usable models (models: + validateModels())
CI / contract (pull_request) Successful in 48s
CI / build (pull_request) Successful in 1m34s
Add an optional top-level `models:` block (Models{allow: List<ModelEntry>})
naming the models any profiles: entry may use. Absent/empty allow: keeps
today's behaviour exactly (no check, no warning). When configured,
FleetConfig.validateModels() fails config load (and reload, via ConfigRef)
naming both the model and the profile, if any profile's model: is outside
the list. The check is one-way: editing profiles: alone can never widen
what is permitted, only models.allow: can.

Wired into Fleetd.main() alongside the other validateXxx() calls, and into
ConfigRef.reload()/DEFERRED_KEYS so a bad edit can't slip in through a
reload either. Each ModelEntry is its own record (not a bare string) so a
later unit can add per-model on/off or load-limit state without changing
the YAML shape. One flat string namespace covers both a bare Claude id and
an opencode provider-prefixed id.
2026-09-10 08:09:27 +07:00
Dai Ha e3e403e5c8 #394: contain a fatal export error instead of letting it abort the scrub
CI / contract (pull_request) Successful in 1m23s
CI / build (pull_request) Successful in 1m52s
EnvAllowListScrub's blanking loop used a plain `export "$n="` on every
name not on the allow-list. For a zsh read-only/special parameter (e.g.
UID) that is a FATAL parameter error, and since the loop runs inside the
sourced startup file, the error aborts the whole file: every name still
to come is never blanked, and scrub-report.txt is never written at all
-- silently, because 2>/dev/null on the group swallows it.

Route each blanking attempt through `eval` instead, which contains the
error to that one iteration. The loop always finishes; a name it could
not blank is now counted separately ("failed" on the report's first
line) and listed !-prefixed rather than disappearing. No skip-list of
known-bad names is added -- every enumerated name is still attempted,
so a name nobody has thought of is still tried and, if it fails, still
counted.

HerdrPeerLauncher: log a WARN when a pane's report carries a nonzero
failed count, and reword the "no report at all" WARN so it no longer
claims the daemon knows the member "saw the full host environment" --
a partial vs. a missing scrub are different situations and only the
first is now distinguishable from the report alone.
2026-09-10 08:08:41 +07:00
Dai Ha 799014e99d Merge #388: scrub a pane shell that is neither login nor interactive
CI / contract (push) Successful in 1m8s
CI / build (push) Successful in 1m38s
2026-09-10 07:21:48 +07:00
Dai Ha 69e09b10fa #388: scrub a pane shell that is neither login nor interactive
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Successful in 1m32s
EnvAllowListScrub generated four zsh startup files but only .zshrc and
.zlogin sourced the scrub — .zshenv (the one file zsh always reads) did
not. A pane shell that is neither login nor interactive reads only
.zshenv and stops, so it was never scrubbed at all (measured on fleet01,
issue #388).

Adding an unguarded scrub to .zshenv (the ticket's own suggested fix) is
wrong: .zshenv is read by every zsh, including a short-lived `zsh -c`
a member's own tooling forks for a single command. Those children are
also neither login nor interactive, so they would scrub the environment
their parent deliberately set for them (GIT_DIR, VIRTUAL_ENV, ...), and
the rewritten scrub-report.txt would describe the last child to exit
instead of the pane.

Fix (per comment 15387, measured): keep .zshrc/.zlogin unconditional,
and add to .zshenv a pass guarded on the exact condition that defines
the gap (neither login nor interactive), plus a per-pane sentinel
(_CB633_SCRUBBED) so it runs once per pane, not once per process. The
sentinel is exported only after the scrub runs, and is folded into the
scrub's own allow-list so a later pass in the same pane cannot blank it
back to empty.

Also corrects the class javadoc's wrong premise (a bare argv[0] proves
NOT login, not "therefore interactive") and its now-stale two-file
walkthrough.

Tests: two new real-zsh tests in EnvAllowListScrubTest run actual
non-login/non-interactive zsh processes (never string-match the
generated files) to prove: a neither-shell pane is scrubbed; a child
that pane forks keeps variables the pane deliberately set for it; the
child does not re-scrub; and scrub-report.txt still describes the pane
after the child exits. Both fail without the production fix (verified
by reverting it and re-running: AssertionFailedError on the sentinel
and on the decoy secret surviving).
2026-09-10 07:21:08 +07:00
Dai Ha b9d09e044e t386: pin the per-member drift baseline the fix's own tests left open
CI / contract (push) Successful in 37s
CI / build (push) Successful in 1m57s
The two tests merged with #386 both start with the member already BUSY, so a
single global drift baseline passes them. This one sleeps the host while nothing
is busy and only then starts a turn, which fails without the per-member map.
2026-09-10 06:59:42 +07:00
Dai Ha fd8650cda4 Merge #386: correct the stall check for a monotonic clock frozen by host sleep 2026-09-10 06:56:37 +07:00
Dai Ha 11050e24ed t384: fix javadoc indentation on the merged shareWithGroup lines
CI / contract (push) Successful in 1m7s
CI / build (push) Successful in 1m29s
2026-09-10 06:54:51 +07:00
Dai Ha 769f282408 #386: give the stall detector a real-time clock, log the divergence
CI / contract (pull_request) Successful in 1m18s
CI / build (pull_request) Successful in 1m26s
FleetHealthMonitor.tick's stalled check compared two monotonic-clock
readings (System.nanoTime(), which macOS freezes across a host sleep),
so a member BUSY for 101 real minutes was never flagged.

The monitor now also takes a wall-clock LongSupplier (realtimeClock),
used only inside the stall check. Each tick measures how far the two
clocks moved apart since the previous tick and folds any positive
divergence into a running total; when a single tick's divergence
exceeds one tick interval (the signature of a sleep, since a tick
cannot run while the process itself is suspended) it logs one WARN
naming how long the detector could not see. The correction is applied
per member, keyed to when that member's current lastActivityAtNanos
was first observed BUSY - not since monitor start - so a sleep that
happened before a member went busy is never charged to it.

Every other use of the monitor's clock (readiness grace, snapshot
timestamp) is unchanged. Backend quarantine/cool-off, the lead tab
scan, the completion resolver, the session reaper and the message
service TTLs are untouched, per the ticket's decision.

Existing FleetHealthMonitor/FleetHealth tests pass unmodified (none of
them ticks a BUSY session more than once, so the drift path never
engages for them). Two new tests: a frozen monotonic clock past the
real-time threshold produces STALL_SUSPECTED, and a single sleep gap
logs the divergence exactly once, not once per tick.
2026-09-10 06:53:30 +07:00
Dai Ha 6f71f40047 Merge #384: pre-create the scrub receipt and give it group write 2026-09-10 06:49:46 +07:00
Dai Ha fb36c5238f Merge #382: SpawnRequest.withProfile() replaces the six-accessor rebuild 2026-09-10 06:49:46 +07:00
Dai Ha cc9cdc938b #384: write shared scrub receipts 2026-09-10 06:45:30 +07:00
Dai Ha 5fede82468 #382: give SpawnRequest a withProfile wither, guard it against the arity trap
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 1m49s
CompositePeerLauncher:372 rebuilt a routed SpawnRequest from a literal
new SpawnRequest(...) call listing six of the original request's own
accessors. That call is only correct because it happens to match the
canonical 6-arg constructor today; add a 7th component plus the
established back-compat constructor at the old (now-shorter) arity and
this call would silently rebind to it, dropping the new field on every
profile-routed spawn with no compile error — the same defect shape
already guarded on FleetConfig.withDefaults() (#357) and MemberSession
(#358).

Add SpawnRequest.withProfile(String), modeled on
MemberSession.withState/withActivity, and use it at the call site
instead. Add a guard test that resolves the true canonical constructor
by exact component types (never by argument count), gives every
component a distinctive value, and asserts every component but
profileName survives withProfile() unchanged.

Proved the guard against the real mechanism: temporarily dropped the
last (role) argument from withProfile()'s constructor call so it bound
to the 5-arg back-compat constructor — it still compiled, and the new
test failed, catching the silently-defaulted role. Restored the fix
and reconfirmed green.
2026-09-10 06:44:15 +07:00
Dai Ha 1515025804 charter: a profile differs in liveness, not just model and cost
CI / contract (push) Successful in 1m12s
CI / build (push) Successful in 1m31s
The old line said profiles differ in model and cost. That is true and it is
not the reason the default hurts. Measured on two hosts: a default sitting on
an exhausted or withdrawn credential either fails the spawn loudly or, worse,
produces a member that starts fine and then returns nothing.

Wording proposed by the fleet01 lead; merged with the existing 'not in tier'
clause, which is still right. Applied byte-identically to the wiki template.
2026-09-10 06:28:51 +07:00
Dai Ha 7754f53662 Merge t385 follow-up: pin the mid-turn redelivery ack 2026-09-10 06:28:08 +07:00
Dai Ha a507f7b31b t385: pin that a redelivery is acked even while the lead is mid-turn
Verification found the fix's dedup check could be moved behind the
injectable-status gate and every test still passed. That placement matters:
this lead is mid-turn most of the time, so gating the ack on an idle pane
leaves the redelivered message held, and the next recovery delivers it again.

The new test fails on that mutant and passes on the fix.
2026-09-10 06:28:03 +07:00
Dai Ha 71c322f104 Merge #385: a redelivered lead message must not write the pane twice 2026-09-10 06:24:40 +07:00
Dai Ha fde2c15627 t385: a redelivered lead message must not write the pane twice
An AMQP recovery clears the held delivery tags, the broker redelivers with
fresh ones, and the coordination loop wrote the same peer message into the
lead pane again. Measured on the live daemon: one msgId reached the pane 12
times in 9 hours, across 19 recovery events.

LeadCoordLoop now remembers the msgIds it has written to a pane (bounded at
1024) and acks a redelivery without a second write. LeadMailbox.ack no longer
returns quietly for an unknown msgId: it throws, so an ack that never reached
the broker is reported instead of hidden. A repeat ack that this connection
already completed stays quiet, tracked in a bounded set.
2026-09-10 06:24:35 +07:00
Dai Ha 2830735644 t377: raise the logger level in the test, or it asserts on an empty list
CI / contract (push) Successful in 46s
CI / build (push) Failing after 1m32s
The 9 tests the member wrote failed 7 of 9 on first build. The production
code was correct; the test harness was not. logback-test.xml sets
dev.ltms.fleet to WARN, so the INFO shape lines were dropped by the level
check before any appender saw them.

The member copied attach()/detach() from MemberTrustModelReportTest but not
the setLevel(INFO) those siblings do at each call site. Doing it inside
attach()/detach() covers all nine at once, and restores the original level
(null, meaning inherit) rather than a concrete one.

Proven by mutation: making the code log the value fails 4 tests, including
theValueNeverAppearsInLogOutput — 'the full GITEA_HOST value must never
reach the log'. Full build 1459 tests green.
2026-09-09 07:52:43 +07:00
Dai Ha 4e3ac91a22 Merge worker/t377-7f587b-7: startup reports the git host value's shape, never its value 2026-09-09 07:49:09 +07:00
Dai Ha 29d3f0b41f t377: report the git host value's SHAPE at startup, never its value
A member receives the git host as GITEA_HOST and often gets a full URL
(scheme, trailing slash) where it expects a bare host, so it builds
https://https://... and the request never leaves the machine.

Logs set/unset, length, startsWithScheme and trailingSlash next to the
existing secret report. The value itself is never logged, and the value
is still passed to members unchanged — rewriting it here would change
what works on one host and breaks on another.

Written by the gx member on branch worker/t377-7f587b-7; it could not
build, commit or push because the host command classifier refused every
shell command (see #381). Build and verification are mine.
2026-09-09 07:49:02 +07:00
Dai Ha cbc732444f fleetd #365: the push-nudge metric outcome is 'sent', not 'delivered'
CI / contract (push) Successful in 1m22s
CI / build (push) Successful in 1m31s
The label was renamed in the #365 merge. This table still named the old one.
'sent' counts the herdr paste-and-submit call returning, never a confirmation
that the pane read it.
2026-09-09 07:40:57 +07:00
Dai Ha 6f828b8c38 Merge worker/t365-3920c5-3: t365: fleet_reply/REST reply distinguish resolved vs queued; rename nudge metric outcome
CI / contract (push) Successful in 46s
CI / build (push) Successful in 1m48s
2026-09-09 07:35:00 +07:00
Dai Ha 145a8c8862 Merge worker/t373-336973-2: t373: pin the production XDG-excludes seam GitWorktreesTest.seedingGitWorktrees builds 2026-09-09 07:35:00 +07:00
Dai Ha 7057291739 Merge worker/t358-6e989b-1: t358: guard MemberSession's 5 rebuild sites and Profile.withProfile() against the back-compat-arity trap 2026-09-09 07:35:00 +07:00
Dai Ha 2302b3bc11 #376: the too-fast failure stops asserting a cause it cannot know
CI / contract (push) Successful in 50s
CI / build (push) Successful in 1m31s
A turn that settles inside MIN_TURN_NANOS is still reported FAILED. Only the claim
about WHY is withdrawn, and the pane is still carried.

Measured on fleet01 on 2026-09-08 UTC: an opencode member on mimo-v2.5-free answered
a real question in 1575ms, below the 2000ms floor. The daemon reported the turn as
failed with 'most likely a backend error before any work started'. The answer was
right there in the scrape. With no errorPattern configured — the live state on both
hosts, which both log as 'backend-error classification: off' — that sentence is a
guess, and a reader who believes it stops looking at the pane.

WHAT I REJECTED, because the next person will try it. A worker implemented the
ticket's first suggested direction: inspect the pane inside the floor and resolve a
COMPLETION when the text looks like a real reply. Its test for 'looks like a real
reply' was non-blank plus a '.', '!' or '?' anywhere in the text. That is unsafe
twice over. lastAssistantBlock falls back to the WHOLE pane when it finds no U+23FA
marker, so on a crash the candidate reply is the entire screen; and a crash pane
almost always contains a full stop, in a file path, a version or a hostname. I ran
that implementation against the new guard test and it resolved

    Error: connection reset while loading src/main/java/Foo.java v1.2.3

as a COMPLETION — expected: <FAILED> but was: <COMPLETION>. A loud wrong answer
became a silent one, which is the trade the ticket brief forbade.

The obvious repair does not work either. Requiring the U+23FA marker as positive
evidence would be safe, but that marker is Claude Code chrome and an opencode pane
never carries it — and an opencode member is what raised this ticket. There is no
reliable cross-backend marker for 'this is a real reply', so this path must not try
to judge one. That reasoning is now in the failTooFast javadoc.

Two guard tests. aPlausibleLookingReplyInsideTheFloorStillFails pins the safety
property against exactly the rejected approach; it passes today and fails against
that implementation, which is how it was verified rather than assumed.
theTooFastFailureDoesNotAssertACauseItCannotKnow pins the wording.

One existing test changed. aNonMatchInsideTheFloorStaysGenericAndNeverNotifiesTheSink
asserted the phrase 'too fast to be real work', which carried the withdrawn claim. It
now asserts what it was really guarding: the floor alone fails the turn, the reason
stays generic, the pane is carried, and the typed sink is never notified.

MIN_TURN_NANOS is unchanged at 2000ms. mvn clean install: 1441 tests, 0 failures.
2026-09-09 07:28:04 +07:00
Dai Ha 3a004dc1b3 t373: pin the production XDG-excludes seam GitWorktreesTest.seedingGitWorktrees builds
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Successful in 1m32s
fleetd #362 review finding 2 protects GitWorktrees#previouslyEffectiveExcludesFileContent's
Java-side XDG_CONFIG_HOME/HOME read (it never goes through a git subprocess, so no
GIT_CONFIG_GLOBAL/GIT_CONFIG_SYSTEM isolation reaches it) with a gitEnv constructor seam. A
mutation run during the #372/#369 merge found that seam unpinned: stripping hermeticGitEnv(tmp)
from seedingGitWorktrees left every test green, poisoned XDG_CONFIG_HOME or not.

Adds seedingGitWorktreesResolvesTheExcludesFileFallbackInsideItsThrowawayDirectory, which asserts
the property directly (a GitWorktrees built for seeding resolves the fallback inside its own
throwaway directory) using a self-contained marker instead of relying on an externally poisoned
env var. Refactors hermeticGitEnv/seedingGitWorktrees into two-argument overloads (one taking an
explicit XDG_CONFIG_HOME / gitEnv) so the new test can pre-populate the marker before construction
while still going through the same production construction every other seeding test uses; no
behavior change for the 4 existing call sites.
2026-09-09 07:24:22 +07:00
Dai Ha a9a3c12232 t365: fleet_reply/REST reply distinguish resolved vs queued; rename nudge metric outcome
CI / contract (pull_request) Successful in 1m17s
CI / build (pull_request) Successful in 1m51s
fleetd #365. fleet_reply always returned the literal "delivered" and
POST /sessions/{id}/reply always returned {"delivered": true}, whether
the reply resolved a live waiting send/ticket or was merely queued in
the inbox for a later drain (CB-307) — both are successes, but not the
same fact.

MessageService.reply() now returns a ReplyOutcome (RESOLVED_SEND,
RESOLVED_ASYNC_TICKET, or QUEUED) instead of an always-true boolean.
FleetMcp.reply and FleetApp.replyMessage both read it: the MCP tool
result names which happened, and the REST body's "delivered" field is
now accurate, with an added "outcome" field.

Also renames the heartbeat/push-loop nudge metric's "delivered" outcome
to "sent" (LeadHeartbeatLoop, ReplyPushLoop, FleetMetrics): it only
records that the herdr agent.prompt paste-and-submit call succeeded,
never that the lead's pane actually read it — there is no read-receipt
concept at that layer, so "delivered" overclaimed there too.

Tests: MessageServiceTest/FleetMcpTest/FleetAppTest strengthened to
assert the specific outcome per case (a resolved send, a resolved async
ticket, and a queued reply); ReplyPushLoopTest updated for the outcome
rename.
2026-09-09 07:23:41 +07:00
Dai Ha 94ec77a1bc t358: guard MemberSession's 5 rebuild sites and Profile.withProfile() against the back-compat-arity trap
CI / contract (pull_request) Successful in 1m23s
CI / build (pull_request) Successful in 1m53s
Follow-up to #357 (FleetConfig.withDefaults()). Reflectively enumerate each
record's own components, resolve the canonical constructor by exact
component types, build a real non-null value per component, run each
rebuild site, and assert every component survives (except the one it is
documented to change). Exclusion lists pinned at 0 for both.

Re-counted the Profile back-compat ladder directly against the source:
8 constructors (arities 25, 24, 22, 20, 18, 15, 14, 12) against a
canonical arity of 26 — the ticket's own number was explicitly untrusted.
2026-09-09 07:19:12 +07:00
Dai Ha 49f285cfda charter: a measured fact in an addendum must carry its own deletion trigger
CI / contract (push) Successful in 49s
CI / build (push) Successful in 1m36s
The operator chose this and its size; the reasoning below is the fleet01 lead's.

Placed in the canonical block's boundary paragraph rather than the orchestration
body. That paragraph already talks about the addendum layer instead of protocol,
every project that mounts the bridge inherits it, and it sits about 3800 characters
before the primary's step list, so it does not dilute the steps a lead reads while
working. Perishability is structurally an addendum problem: the block is
byte-identical across projects by construction, so a dated local measurement in the
block body would already be a layering violation.

What happened. The fleet01 lead's kb addendum held a dated merge-refusal section
that carried an instruction to delete itself once it stopped reproducing. On
2026-09-08 UTC the operator granted merge rights on akb/kb, the lead re-ran the
probe, got 409 'head out of date' where the identical request had returned 405
'User not allowed to merge PR' on 2026-09-06, and deleted the section as instructed.

Why four parts and not one. The lead's finding is that the banner did not work
because it was emphatic. It worked because the falsification condition was
executable: it carried the exact probe, the reason for the all-zeroes
head_commit_id, and what each response code meant. The lead did not have to
reconstruct the experiment or decide what would count as refutation, and just ran
it. A banner saying 'this may be out of date, verify before relying on it' costs
the same space and does nothing, because deciding what would falsify a claim is the
expensive step and a reader in the middle of another task will not pay it. So: the
date, the command, what each outcome means, and the instruction to delete. The
fourth without the second is decoration.

The closing clause is the justification for the machinery. Most stale notes are
merely wrong. This one went stale in the dangerous direction: it would have told a
future lead it could not merge at the exact moment merging became its job, silently
and with confidence. A note that goes harmlessly stale does not need this.

Note what is NOT centralized here. The banner text itself cannot be. What fired for
the lead was a specific instruction sitting on top of the specific stale fact, which
it could not read past on its way to acting. A rule elsewhere saying 'date your
measurements' would not have fired, because nobody reads that rule at the moment
they re-measure. This sentence sets the convention; the trigger still has to live
next to the fact it governs.

Propagated to the wiki template in the same turn, wiki 8c4f152 on main; the sync
check in this file's addendum reports 'in sync: True'.
2026-09-09 04:24:17 +07:00
Dai Ha 5a12ae7930 charter: test a refusal, and do not count a transport failure as one
CI / contract (push) Successful in 1m11s
CI / build (push) Successful in 1m38s
Step 8 gained a refusal paragraph in 2f71a30, which said what a lead does once the
forge refuses a merge. It did not say how a lead establishes that it was refused.
Both halves of this amendment come from the fleet01 lead, measured on akb/kb on
2026-09-08 UTC, and both are ways to be wrong about a permission you never tested.

Do not read a refusal off a permissions field. After the operator granted merge
rights, the lead re-ran its probe: POST .../pulls/53/merge with an all-zeroes
head_commit_id, chosen so the request cannot succeed on its merits and a rejection
can only mean the refusal. It returned 409 'head out of date' where the identical
request returned 405 'User not allowed to merge PR' on 2026-09-06. A 409 is payload
validation and sits after the permission gate, so the grant took. The lead reports
the repository permissions object did not change across that flip -- still
admin:false, push:true, pull:true. I did not read that object myself; my forge token
is a different identity and would return a different one, so this stays the lead's
measurement and not mine. Merge rights on a protected branch live in branch
protection, so a permissions field can be wrong in both directions.

Do not count a transport failure as a refusal. The lead's first attempt returned
HTTP 000, because GITEA_HOST already carries a scheme and a trailing slash and the
URL came out as https://https://git.ltms.dev//api/... Under a 'not 200' test that is
indistinguishable from being refused. A probe exists to separate a refusal from
everything else, so an error that never reached the gate has to be a third answer
that concludes nothing.

Propagated to the wiki template in the same turn, wiki 8c2ef96 on main; the sync
check in this file's addendum reports 'in sync: True'.
2026-09-09 04:09:08 +07:00
Dai Ha 127e6832a9 Merge #374: fleetd holds off idle sleep while any member is live
CI / contract (push) Successful in 1m18s
CI / build (push) Successful in 2m3s
Lands PR #355 (fleetd #354's sibling), rebased onto current main by a worker
after 39 commits of drift left it unmergeable.

The problem, measured on the original branch: a fleetd host idle-slept after as
little as one minute (pmset -g custom reported 'sleep 1' on battery). Overnight
the daemon's AMQP link dropped 13 times, and every drop minute had a sleep or
wake event in pmset -g log in the same minute or the one before. The AMQP churn
is the visible symptom; the real cost is a member mid-turn freezing with the
host, and a long turn with nobody typing is exactly the case that goes idle.

IdleSleepGuard holds an OS-level assertion for as long as at least one member is
live. It is driven by SessionManager's existing onAcquire/onRelease hooks rather
than a second member count kept in parallel, so it reads the same registry
fleet_list's numbers come from, and only a real 0->1 or 1->0 crossing touches the
OS. It fails safe: a mechanism that cannot acquire means nothing is ever held,
and it never throws, never blocks a spawn, a release, or shutdown.

Conflict resolution was the whole job, and all three were in config plumbing:
ConfigRef, FleetConfig and ConfigRefTopLevelReportingCoverageTest. The power
package is byte-identical to the original branch commit.

Verified on this merge, not taken from the worker's report:
  mvn clean install -> Tests run: 1439, Failures: 0, Errors: 0, BUILD SUCCESS
  (1425 on main + 14 new: 4 caffeinate, 5 guard, 1 wiring, 4 config)

The denominator recount, which the worker flagged as its own weakest number
because this file's count has drifted three times before (#330/#333/#337). I
counted it mechanically rather than reading it: FleetConfig has 24 canonical
record components; COLD_KEYS 5, DEFERRED_KEYS 13, SPLIT_KEYS 3, plus the 3 the
javadoc names as hot-excluded (placement, memberCredentials, memberLoginShell).
5+13+3+3 = 24. The javadoc's '24 components: 5 cold, 13 deferred, 3 split, 3
hot-excluded' is correct. The worker's prose called idleSleepGuard the 25th
constructor argument; it is the 24th. The code is right, the report was off by
one.

Mutation run on merge, on the half the worker verified by READING rather than by
proving -- it said it had checked that withDefaults()'s final call binds the true
canonical constructor. I dropped the trailing idleSleepGuard argument so the call
silently binds the 23-arg back-compat overload. It compiles, which is the whole
hazard. Caught: 1 failure, 3 errors, BUILD FAILURE, and
FleetConfigWithDefaultsPreservesEveryComponentTest names the dropped component
and prints its own denominator -- '24 components, 24 checked, 0 excluded, 23
survived'. That test was added on main after this exact defect happened live when
idleSleepGuard was added on a sibling branch; the worker had to add the missing
entry to it, and doing so is what makes the guard cover this component at all.
2026-09-07 20:31:41 +07:00
43 changed files with 3783 additions and 180 deletions
+29 -6
View File
@@ -7,6 +7,14 @@
> wiki ([Use Cases](https://git.ltms.dev/fleet/fleetd/wiki/7-Use-Cases) → *The portable
> CLAUDE.md block*); improvements go to the template first, then out to each project. Anything
> specific to *this* repo lives under §Project addendum below, never inline above it.
>
> **Anything you measure in an addendum is perishable.** Date it, give the command that
> re-measures it and what each outcome means, and tell the reader to delete the section once
> it stops reproducing. The four parts work together: deciding what would falsify a claim is
> the expensive step, and a reader in the middle of another task will not pay it, so a bare
> "verify before relying on this" costs the same space and does nothing. The case this is for
> is a note that goes stale as a live restriction — it will tell a future session it cannot do
> the thing at the moment doing it becomes the job.
If no `fleet_*` MCP tools are mounted in this session, this section does not apply — skip it.
@@ -71,8 +79,11 @@ below are the procedure — run them in order, every task, not only the big ones
the final judgment call, verification, merges, and anything that depends on context only you
hold. Nothing else is yours by default.
3. **Spawn every delegated unit first** — `fleet_spawn{profile, worktree:true, ticket}`, one per
unit, *before* sending any. Pass `profile` explicitly: profiles differ in model and cost, not in
tier, so the default is rarely what you want.
unit, *before* sending any. Pass `profile` explicitly: profiles differ in model, cost and
LIVENESS, not in tier, so the default is rarely what you want. The default is whatever the
daemon reports, and on a host where it sits on an exhausted or withdrawn credential every
unqualified spawn fails — sometimes loudly, sometimes as a member that spawns fine and then
produces nothing. `fleet_profiles` reports the default; check it once per session.
4. **Then send them all** — `fleet_send{sessionId, content, wait:false}`. Line 1 of every brief is
`Load the <name> skill.` naming the worker's playbook; those skills are opt-in and that line is
what makes them reliable. Where the project ships no such skill, spell the procedure out in the
@@ -100,6 +111,12 @@ below are the procedure — run them in order, every task, not only the big ones
without having read the diff yourself. A refusal is exactly when that shortcut is tempting,
because no action is left that forces you to look, and taking it turns this step into
forwarding a reviewer's verdict — which is delegating the merge by proxy, two lines above.
**Test a refusal; do not read it off a permissions field.** A protected branch holds its merge
rights separately from the repository permissions, so that field can say yes while the merge is
refused, and still say no after a grant makes it work. Probe instead, with a request that cannot
succeed on its merits, so a rejection can only mean the refusal. Treat a transport failure as a
third answer that proves nothing: a timeout, a DNS error or a bad URL is not a refusal, and
counting it as one makes you sure of something you never measured.
**Steps 3 and 4 are separate on purpose** — spawning and sending in one loop is how parallel work
silently becomes serial, and it is the most common way this layer is wasted. For the same reason,
@@ -120,7 +137,7 @@ the merge — and merging on a reviewer's word is delegating it by proxy.
| Answer a member's `fleet_ask` | `fleet_send{turnId, content}` — **not** `sessionId` |
| Message a **peer lead** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task |
| Message a **peer lead** on another daemon or host | `fleet_send{coordId: <their coord-id>, content}` — needs a `coordinator:` block; your own coord-id is in `fleet_list`. Coordination only, **never** a task |
| Answer a peer lead that messaged you | `fleet_reply{content}` — the one case a lead replies |
| Answer a peer lead that messaged you | `fleet_send{coordId}` — or `{sessionId}` if they are on this host. **Not** `fleet_reply`: it has no peer route and the publish is refused |
| Collect a held reply | `fleet_poll{target}` · then `fleet_ack{target, msgId}` |
| Tear down a member | `fleet_stop{paneId}` |
@@ -146,15 +163,21 @@ The traffic between leads is coordination and nothing else:
3. **Verify a peer exactly as you verify yourself.** Peer status buys nothing: check the claim
against the code, and re-run the build. A peer's correction gets the same treatment — right or
wrong on the evidence, not on who said it. Neither of you merges the other's work unreviewed.
**N observations are N data points only if they differ in the axis you are trusting.** This cuts
both ways. N *failures* blamed on one cause are one data point when the cases share what you are
not varying. N *agreeing measurements* are also one data point when they share an instrument —
two hosts, two operators and the same formula is one formula, not two confirmations.
4. **Ask a peer to read your project addendum.** Your addendum is instruction surface: every future
session on your host obeys it, and a wrong one is obeyed just as faithfully as a right one. The
author is the worst reader of their own qualifier placement — measured here, one addendum carried
two defects and a non-author found both. If you have no peer, at least re-read it asking "which
sentence goes false first, and would a reader reach the caveat before acting?"
Being messaged by a peer does not make you its worker: answer with `fleet_reply`, and push back on
the substance if it is wrong. A peer that simply complies has thrown away the reason there are two of
you.
Being messaged by a peer does not make you its worker: answer the way you would open —
`fleet_send{coordId}` for another daemon, `fleet_send{sessionId}` on this host — and push back on
the substance if it is wrong. `fleet_reply` resolves a member's blocked `fleet_send`; a peer's
coord-id message is durable and non-blocking, so there is nothing for it to resolve. A peer that
simply complies has thrown away the reason there are two of you.
### Member (worker or architect) — the turn contract
+1 -1
View File
@@ -192,7 +192,7 @@ Deliberately small; every one maps to a failure mode we have actually hit.
| `fleet_send_duration_seconds` | histogram | delegated turn latency |
| `fleet_replies_total{path}` | counter | path ∈ rendezvous\|inbox — how often a reply strands (CB-307's whole reason to exist) |
| `fleet_inbox_depth{target}` | gauge | undrained replies; steady-state should be 0 |
| `fleet_push_nudges_total{outcome}` | counter | outcome ∈ delivered\|exhausted — a rising `exhausted` means the primary is not draining |
| `fleet_push_nudges_total{outcome}` | counter | outcome ∈ sent\|exhausted — a rising `exhausted` means the primary is not draining. `sent` was called `delivered` until fleetd #365; it counts the herdr paste-and-submit call returning, never a confirmation the pane read it |
| `fleet_spawns_total{kind,outcome}` | counter | outcome ∈ ready\|timeout\|guard_rejected; per peer kind (CB-402) |
| `fleet_sessions{state}` | gauge | SPAWNING/READY/BUSY/DONE census |
| `fleet_herdr_calls_total{method,outcome}` | counter | socket health — the dependency everything rests on |
+29
View File
@@ -877,3 +877,32 @@ guard:
# terminal: term_65619bd6174568
# pushReminders: 5
# pushBackoffMs: 15000
# Central allow-list of models any profiles: entry may name. Nothing checked a profile's model:
# value before this block existed — it was a free-form string handed straight to the backend
# adapter, and a withdrawn or misspelled name failed silently instead of at config load (opencode
# falls back to a default model rather than erroring on an unknown -m).
#
# Absent, or present with an empty allow:, is OFF: no profile's model: is checked, exactly like
# before this block existed. fleetd.yaml is gitignored on every host, so an upgrade must not force
# every operator to enumerate their models before the daemon will start.
#
# The list is the authority; profiles: is checked against it, never the reverse — adding or
# editing a profiles: entry cannot, by itself, widen what is permitted here.
#
# Enforcement is at CONFIG LOAD only (a bad model: fails the daemon at startup, naming both the
# model and the profile). There is no spawn-time enforcement, no runtime on/off switch, and no
# interaction with BackendQuarantine — those are separate, later units.
#
# allow → the permitted models. Each entry is its own block (not a bare string) so a later unit
# can add an on/off state or a load limit per model without changing this shape.
# model → the model id exactly as a profiles: entry's model: field would write it. One flat,
# opaque-string namespace: a bare Claude id (claude-sonnet-5) and an opencode
# provider-prefixed id (openai/gpt-5.6-terra) both fit here unchanged — the check is a
# plain string match, never a parse of the provider prefix or a branch on kind:.
# models:
# allow:
# - model: claude-sonnet-5
# - model: claude-opus-5
# - model: openai/gpt-5.6-terra
# - model: amazon.nova-pro-v1:0
+155 -18
View File
@@ -123,12 +123,21 @@ public final class Fleetd {
// else can fail on a silently-empty one. A daemon started without a login shell (launchd)
// boots fine either way — this is the only thing that says so out loud.
reportRequiredSecrets(cfg);
// fleetd #377: the git host value a member receives as GITEA_HOST is often a full URL
// (scheme and trailing slash), not a host name. A member that assumes a bare host then
// builds https://https://... and the request never leaves the machine. Report the shape
// next to the secret report — shape only, never the value.
reportGitHostShape(cfg);
reportMemberTrustModel(cfg);
// CB-596: an absent (or empty) memberCredentials: block blocks NOTHING — no credential
// name is hardcoded any more to fall back on. Say so loudly, the same way a missing
// secret is reported above, so upgrading past this commit never silently drops CB-592's
// protection.
reportMemberCredentialsGap(cfg);
// fleetd #395: an unset exhaustedPattern is a silent opt-out of usage-limit detection for
// that profile — say so loudly, the same way the two reports above do, rather than let an
// operator discover it only when a limit goes undetected.
reportExhaustedPatternGap(cfg);
// CB-559: `cfg` stays the startup snapshot — every validation and every piece of one-time
// wiring below reads it, and must, because those decisions cannot be unmade. `config` is the
// live reference the hot paths read per use. Which keys can actually move is ConfigRef's
@@ -139,19 +148,16 @@ public final class Fleetd {
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
guard.assertPrimaryClean(System.getenv());
// CB-501: refuse to start if the bind is wider than the auth mode can defend. Under
// loopback-trust, "not a known worker" means "the primary" — sound only because the OS
// refuses remote connections to a loopback socket. This throws rather than warns so the
// dangerous configuration cannot be reached by ignoring a log line.
cfg.validateAuthExposure();
cfg.validateLeadTabPrefixes();
// CB-542: a subscription:true profile whose env: reseats ANTHROPIC_BASE_URL/AUTH_TOKEN would
// reach an unguarded endpoint (the launcher skips SubscriptionGuard for it). Refuse at load.
cfg.validateSubscriptionProfiles();
cfg.validateCharters();
// CB-548: every architect slot must name a configured workers: profile — the strong-model
// backend the future spawn lifecycle would read. A stale reference dies here, not later.
cfg.validateMembers();
// Every FleetConfig.validateXxx() the operator's config can fail — CB-501's auth-exposure
// check, CB-531's lead-tab-prefix check, CB-542's subscription-profile check, the charter
// and member-slot checks, and the "central allow-list of usable models" check — must run
// here, at load, before anything below opens a socket or spawns a member. fleetd ticket
// "central allow-list of usable models" follow-up: mutation testing found six individual
// calls here with nothing proving any of them still ran (deleting one left the full suite
// green). validateAll() replaces them with the one call that FleetConfigValidateAllTest
// and the Fleetd-startup tests actually pin — see FleetConfig#validateAll's javadoc for
// why a name-by-name list here would have the same defect it replaces.
cfg.validateAll();
Path socket = cfg.herdrSocket() != null && !cfg.herdrSocket().isBlank()
? Path.of(cfg.herdrSocket())
@@ -579,8 +585,12 @@ public final class Fleetd {
if (cfg.health() != null && cfg.health().isEnabled()) {
// CB-580: a member found GONE/NEVER_READY must fail whatever ticket is waiting on it,
// through the same idempotent target-wide operation CB-516 already uses on release.
// fleetd #386: System::nanoTime freezes across a macOS sleep, so the stall check also
// gets a wall-clock source to detect and correct for that freeze. Every other decision
// in FleetHealthMonitor stays on the monotonic clock, unchanged.
healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler,
System::nanoTime, cfg.health().intervalOrDefault(),
System::nanoTime, () -> TimeUnit.MILLISECONDS.toNanos(System.currentTimeMillis()),
cfg.health().intervalOrDefault(),
cfg.health().workingSuspectAfterOrDefault(), messages::abandon);
String coverage = FleetHealthMonitor.coverage(true,
cfg.health().notifications() != null && cfg.health().notifications().configured());
@@ -652,10 +662,8 @@ public final class Fleetd {
// built a second time — two independently-constructed sources reading the SAME BackendQuarantine
// / BackendOutagePolicy would still be able to drift (e.g. a future edit to the credentialIdFor
// closure in only one of the two places), exactly the shape #284 was.
FleetMcp.QuarantineSource quarantineSource = new FleetMcp.QuarantineSource(profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.effectiveCredentialId();
}, quarantine);
FleetMcp.QuarantineSource quarantineSource = quarantineSource(config, quarantine,
exhaustedPatternsByProfile);
FleetMcp.OutageSource outageSource = new FleetMcp.OutageSource(profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.effectiveCredentialId();
@@ -789,6 +797,19 @@ public final class Fleetd {
return target -> presence.isPresent(target) || leads.get().containsKey(target);
}
/**
* fleetd #404: production source for quarantine reporting. Credential IDs are hot, but
* exhausted patterns are compiled once at startup for {@link CompletionResolver}, so the armed
* field must use that same compiled map until restart.
*/
static FleetMcp.QuarantineSource quarantineSource(ConfigRef config, BackendQuarantine quarantine,
Map<String, Pattern> startupExhaustedPatterns) {
return new FleetMcp.QuarantineSource(profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.effectiveCredentialId();
}, quarantine, profile -> startupExhaustedPatterns.containsKey(profile));
}
/**
* fleetd #248: package-private factory for the member worktree/branch lookup {@link
* CompletionResolver} uses to name a fallback report's worktree and branch (fleetd#241).
@@ -1195,6 +1216,76 @@ public final class Fleetd {
});
}
/**
* fleetd #377: the env var names holding the git host value that members receive as
* {@code GITEA_HOST}. {@code HerdrPeerLauncher.applyGitToken} injects {@code GITEA_HOST}
* only for profiles that opted in via {@code gitTokenEnv} (CB-302), reading the value from
* that profile's {@code gitHostEnv} (default {@code GITEA_HOST}), so the shape only matters
* where a git token is opted in. A var used by more than one profile is one entry naming
* every profile that reads it, the same shape as {@link #requiredSecretEnvVars}.
*
* <p>Package-private and pure (no I/O, no logging) so the derivation is unit-testable
* without capturing log output; {@link #reportGitHostShape(FleetConfig)} is the logging caller.
*/
static Map<String, List<String>> gitHostEnvVars(FleetConfig cfg) {
Map<String, List<String>> hostsBy = new LinkedHashMap<>();
cfg.profiles().forEach((name, profile) -> {
if (profile.hasGitToken()) {
hostsBy.computeIfAbsent(profile.gitHostEnv(), _ -> new ArrayList<>())
.add("profile '" + name + "' gitHostEnv");
}
});
return hostsBy;
}
/**
* True when the value already starts with a URI scheme ({@code https://...}, {@code
* http://...}). A bare host name and a host:port must both report {@code false} — the shape
* this ticket exists for is a value that <em>looks like</em> a host but is a full URL, and
* confusing those in the report would move the failure to the log instead of the network.
*/
static boolean startsWithScheme(String value) {
return value.matches("[A-Za-z][A-Za-z0-9+.-]*://.*");
}
/**
* fleetd #377: log, on the same startup path as {@link #reportRequiredSecrets}, the SHAPE of
* each git host value that members receive as {@code GITEA_HOST}: set or unset, its length,
* whether it starts with a scheme, whether it ends with a slash. Never the value itself — the
* same discipline as {@link #reportRequiredSecrets}, which logs by name only. Either shape is
* legitimate: the value is passed to members unchanged, and a line that quietly rewrites it
* would change what works on one host and breaks on another. The shape line only tells the
* operator which URL form to expect from a member that builds on {@code GITEA_HOST}. An
* unset variable is logged at INFO — useful information, not an error — and the daemon
* starts on either way.
*
* <p>The env read lives in the overload below so a test can drive the line with a known
* value and prove that value never reaches the log.
*/
static void reportGitHostShape(FleetConfig cfg) {
reportGitHostShape(cfg, System.getenv());
}
static void reportGitHostShape(FleetConfig cfg, Map<String, String> env) {
Map<String, List<String>> hostsBy = gitHostEnvVars(cfg);
if (hostsBy.isEmpty()) {
log.info("startup git host: no profile sets a gitTokenEnv — nothing to check");
return;
}
hostsBy.forEach((varName, sources) -> {
String value = env.get(varName);
if (value == null || value.isBlank()) {
log.info("startup git host {}: unset ({}) — a member gets GITEA_TOKEN but no "
+ "GITEA_HOST value", varName, String.join(", ", sources));
} else {
log.info("startup git host {}: set ({}) — length={}, startsWithScheme={}, "
+ "trailingSlash={}",
varName, String.join(", ", sources),
value.length(), startsWithScheme(value), value.endsWith("/"));
}
});
}
/**
* fleetd #184: state the member trust model at startup. Environment controls and worktrees do
* not make a sandbox when fleetd and its members use the same OS user. A separate herdr may
@@ -1244,6 +1335,52 @@ public final class Fleetd {
+ "fleetd.yaml — see fleetd.example.yaml — and restart.");
}
/**
* fleetd #395: {@code exhaustedPattern} (see {@link FleetConfig.Profile#exhaustedPattern}) is
* deliberately opt-in — {@code null}/blank means a backend refusal on that profile is never
* classified as {@code BACKEND_EXHAUSTED}, so its credential is never quarantined. That is a
* legitimate choice (guessing the vendor's wording would be worse), but an operator who never
* opted a profile in should not discover the gap only when a usage limit silently goes
* undetected. Warn once at startup, naming every unarmed profile, exactly like {@link
* #reportMemberCredentialsGap} — never refuse to start over it.
*
* <p>A {@code subscription: true} profile that is unarmed gets a SECOND, louder WARN of its
* own: it bills the operator's metered Claude plan, the case where an undetected usage limit
* costs the most.
*
* <p>Package-private so a test can capture the real log via a {@link
* ch.qos.logback.core.read.ListAppender}, the same pattern {@link
* #reportMemberCredentialsGap}'s own test uses.
*/
static void reportExhaustedPatternGap(FleetConfig cfg) {
List<String> unarmedSubscription = new ArrayList<>();
List<String> unarmedOther = new ArrayList<>();
cfg.profiles().forEach((name, profile) -> {
if (!profile.hasExhaustedPattern()) {
(profile.isSubscription() ? unarmedSubscription : unarmedOther).add(name);
}
});
if (unarmedSubscription.isEmpty() && unarmedOther.isEmpty()) {
log.info("exhaustedPattern: every configured profile has usage-limit detection armed");
return;
}
List<String> allUnarmed = new ArrayList<>(unarmedSubscription);
allUnarmed.addAll(unarmedOther);
allUnarmed = allUnarmed.stream().sorted().toList();
log.warn("exhaustedPattern: profile(s) {} have no exhaustedPattern configured — a "
+ "usage-limit refusal on any of them is never detected and never "
+ "quarantines its credential. Set exhaustedPattern (see "
+ "fleetd.example.yaml) to arm detection for a profile.",
allUnarmed);
if (!unarmedSubscription.isEmpty()) {
List<String> sortedSubscription = unarmedSubscription.stream().sorted().toList();
log.warn("exhaustedPattern: subscription profile(s) {} run on the operator's metered "
+ "Claude plan and have NO usage-limit detection armed — this is the "
+ "case where a missed usage limit costs the most.",
sortedSubscription);
}
}
/**
* Poll herdr's {@code ping} until it answers or {@link #HERDR_WAIT_SECONDS} elapses (CB-504).
*
@@ -43,6 +43,12 @@ import java.util.function.Supplier;
* running daemon keeps whatever this was at startup regardless of a later edit),
* {@code spawnReadyTimeoutMs} / {@code spawnReadyPollMs}, {@code quarantineCooldownSeconds}
* (CB-578 stage B — baked once into the {@code BackendQuarantine} built at startup),
* {@code models:} (fleetd ticket "central allow-list of usable models" — {@link
* FleetConfig#validateModels()} re-runs against the fresh config in {@link #reload()}
* (via {@link FleetConfig#validateAll()}), so a
* models.allow: edit that would refuse to boot still refuses the reload; a change that
* passes has nothing built at startup to rebuild, so it is reported deferred rather than
* silently accepted with no report at all),
* {@code guard:}, {@code worktreeRoot:}, {@code worktreeGroup:} and {@code memberSkills:}
* (all three of the latter baked once into the {@code GitWorktrees} built at
* {@code Fleetd.java:251} and never rebuilt — fleetd #323 instance 2 found
@@ -135,8 +141,9 @@ import java.util.function.Supplier;
* </ul>
*
* <p><strong>The denominator, measured on 2026-09-04 (fleetd #330; recounted for fleetd #333);
* recounted again for fleetd #362, and again after {@code idleSleepGuard:} was added.</strong>
* {@code FleetConfig} has 24 top-level record components: 5 cold, 13 deferred, 3 split, 3
* recounted again for fleetd #362, again after {@code idleSleepGuard:} was added, and again after
* {@code models:} was added.</strong>
* {@code FleetConfig} has 25 top-level record components: 5 cold, 14 deferred, 3 split, 3
* hot-excluded. Three of them are named nowhere in this file, and the reason is the same for all
* three: {@code placement}, {@code memberCredentials} and {@code memberLoginShell} are
* <strong>hot</strong> and correctly absent — all three are read live off {@code config.get()}
@@ -219,7 +226,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
static final Set<String> DEFERRED_KEYS = Set.of(
"guard", "worktreeRoot", "worktreeGroup", "memberSkills", "primary", "configReload",
"leadHeartbeat", "lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs",
"quarantineCooldownSeconds", "profiles", "idleSleepGuard");
"quarantineCooldownSeconds", "profiles", "idleSleepGuard", "models");
private final Path path;
private final AtomicReference<FleetConfig> current;
@@ -322,12 +329,12 @@ public final class ConfigRef implements Supplier<FleetConfig> {
fresh = FleetConfig.load(path);
// The same gate startup runs. A config that would have refused to boot must not be able
// to slip in through a reload — that is how a daemon ends up in a state it could never
// have started in, which is the hardest kind to debug.
fresh.validateAuthExposure();
fresh.validateLeadTabPrefixes();
fresh.validateSubscriptionProfiles();
fresh.validateCharters();
fresh.validateMembers();
// have started in, which is the hardest kind to debug. fleetd ticket "central allow-list
// of usable models" follow-up: this used to be six individual validateXxx() calls, and
// mutation testing found two of the six unpinned here even though startup pinned nothing
// at all — see FleetConfig#validateAll's javadoc for why the fix is one reflective call,
// not a longer hand-maintained list.
fresh.validateAll();
} catch (RuntimeException e) {
String msg = e.getMessage() == null ? e.toString() : e.getMessage();
log.warn("config reload from {} refused, keeping the running config: {}", path, msg);
@@ -439,6 +446,15 @@ public final class ConfigRef implements Supplier<FleetConfig> {
if (!Objects.equals(old.idleSleepGuard(), fresh.idleSleepGuard())) {
changed.add("idleSleepGuard");
}
// fleetd ticket "central allow-list of usable models": validateModels() runs again in
// reload() above (via validateAll()), so a bad edit is already refused as cold-adjacent
// (the whole reload is refused via the catch block, never partially applied). A GOOD edit
// to the allow-list
// itself has nothing built at startup to rebuild — it only ever mattered to the validation
// call that already ran — so report it deferred rather than silently swallowing the change.
if (!Objects.equals(old.models(), fresh.models())) {
changed.add("models");
}
if (!Objects.equals(old.spawnReadyTimeoutMs(), fresh.spawnReadyTimeoutMs())
|| !Objects.equals(old.spawnReadyPollMs(), fresh.spawnReadyPollMs())) {
changed.add("spawnReady*");
@@ -16,11 +16,15 @@ import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.io.UncheckedIOException;
import java.lang.reflect.InvocationTargetException;
import java.lang.reflect.Method;
import java.lang.reflect.Modifier;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.Comparator;
import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.List;
@@ -124,6 +128,13 @@ import java.util.regex.PatternSyntaxException;
* under a member's long turn. {@code null} (the block omitted) behaves the
* same as an explicit {@code enabled: true}; set {@code enabled: false} to
* turn it off. See {@link dev.ltms.fleet.power.IdleSleepGuard}.
* @param models central allow-list of models any {@code profiles:} entry may name. {@code
* null} or an empty {@code allow:} ⇒ off: {@link #validateModels()} checks
* nothing and every existing config keeps working exactly as it does today.
* When non-empty, a profile whose {@code model:} is not one of {@link
* Models#ids()} fails config load, naming both the model and the profile —
* see {@link #validateModels()}. This block only decides what may be
* CONFIGURED; nothing here enforces it at spawn time. See {@link Models}.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record FleetConfig(
@@ -150,7 +161,22 @@ public record FleetConfig(
String worktreeGroup,
String memberLoginShell,
String memberSkills,
IdleSleepGuard idleSleepGuard) {
IdleSleepGuard idleSleepGuard,
Models models) {
/** Back-compat form before the {@code models:} block was added. */
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
Guard guard, String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
ConfigReload configReload, Integer quarantineCooldownSeconds,
MemberCredentials memberCredentials, Coordinator coordinator, String worktreeGroup,
String memberLoginShell, String memberSkills, IdleSleepGuard idleSleepGuard) {
this(bind, herdrSocket, memberHerdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, worktreeGroup,
memberLoginShell, memberSkills, idleSleepGuard, null);
}
/** Back-compat form before the {@code idleSleepGuard:} block was added. */
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
@@ -1318,6 +1344,72 @@ public record FleetConfig(
}
}
/**
* Central allow-list of models any {@code profiles:} entry may name (fleetd ticket: "a central
* allow-list of usable models"). Nothing before this block checked a profile's {@code model:}
* against anything — it was a free-form string handed straight to the backend adapter, and a
* withdrawn or misspelled name failed silently (opencode falls back to a default model rather
* than erroring on an unknown {@code -m}) rather than at config load, where a mistake is cheap.
*
* <p><b>Absent or empty {@code allow:} is "off"</b>, on purpose: this is a large deployment
* with gitignored {@code fleetd.yaml} on more than one host, and a change that forced every
* operator to enumerate their models before the daemon would start would break every one of
* them on upgrade. See {@link FleetConfig#validateModels()}, which is where the allow-list is
* actually enforced, at config load.
*
* <p><b>The list is the authority; a {@code profiles:} entry is checked against it, never the
* other way around.</b> Adding or editing a {@code profiles:} entry cannot, by itself, widen
* the set of permitted models — only editing {@code models.allow:} itself can. This is the
* invariant the ticket asked for: the two blocks are validated in one direction only.
*
* <p><b>Out of scope here, deliberately:</b> nothing in this block is read at spawn time —
* enforcing it against a live spawn, an on/off runtime switch, and any interaction with {@code
* BackendQuarantine} are separate units. This block is config-load validation only.
*
* @param allow the permitted models, each its own {@link ModelEntry} rather than a bare
* string — see that record's javadoc for why. {@code null}/empty ⇒ the block is
* treated as absent: {@link FleetConfig#validateModels()} checks nothing.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Models(List<ModelEntry> allow) {
public Models {
allow = (allow == null) ? List.of() : List.copyOf(allow);
}
/**
* One permitted model, named as a record rather than a bare string on purpose: a later unit
* needs to hang an on/off state and a load-limit state off each entry, and a bare {@code
* List<String>} cannot grow those fields without changing the YAML shape underneath every
* operator who already wrote one. {@link #model()} is intentionally a single flat,
* opaque-string namespace — a bare Claude id ({@code claude-sonnet-5}) and an opencode
* provider-prefixed id ({@code openai/gpt-5.6-terra}) both fit it unchanged, because
* {@link FleetConfig#validateModels()} only ever compares a profile's {@code model:} value
* against this string for exact equality; it never parses a provider prefix or branches on
* a profile's {@code kind:}.
*
* @param model the model id exactly as a {@code profiles:} entry's {@code model:} field
* would name it
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record ModelEntry(String model) {
public ModelEntry {
model = (model == null || model.isBlank()) ? null : model.trim();
}
}
/** {@link #allow}'s model ids, as a set for membership checks. Blank/null entries are dropped. */
public Set<String> ids() {
Set<String> ids = new java.util.LinkedHashSet<>();
for (ModelEntry e : allow) {
if (e != null && e.model() != null) {
ids.add(e.model());
}
}
return Collections.unmodifiableSet(ids);
}
}
/**
* The terminal → lead-name map seeded from the legacy singular {@code primary:} pin (CB-530).
*
@@ -1569,7 +1661,7 @@ public record FleetConfig(
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds",
"memberCredentials", "coordinator", "worktreeGroup", "memberLoginShell", "memberSkills",
"idleSleepGuard");
"idleSleepGuard", "models");
/** Load and validate config from {@code path}. */
public static FleetConfig load(Path path) {
@@ -2254,10 +2346,15 @@ public record FleetConfig(
// enabled: true — see its javadoc), so defaulting the block here would change nothing a
// reader observes and would only obscure that "block omitted" and "block present and
// enabled" are deliberately the same outcome.
// models is left as-is, like broker/primary/coordinator above: null/empty is "off", and an
// absent block must validate nothing (see Models's javadoc) — defaulting it here to an
// empty Models would be a no-op for validateModels() either way, since an empty allow-list
// already means "check nothing", so there is nothing to gain and one more null check to
// avoid by leaving it exactly as configured.
return new FleetConfig(b, herdrSocket, memberHerdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
quarantineCooldown, mc, coordinator, worktreeGroup, memberLoginShell, memberSkills,
idleSleepGuard);
idleSleepGuard, models);
}
/**
@@ -2477,6 +2574,118 @@ public record FleetConfig(
}
}
/**
* Reject a {@code profiles:} entry whose {@code model:} is not on the configured {@link
* #models} allow-list.
*
* <p>Absent or empty {@code models.allow:} validates nothing — see {@link Models}'s javadoc:
* an existing config with no such block must keep working exactly as it does today. Once the
* operator declares at least one entry, every profile's {@code model:} (when set — a profile
* may legitimately leave it {@code null}, e.g. a {@code subscription: true} profile relying on
* the account's own default) must equal one of {@link Models#ids()} exactly. The comparison is
* a flat string match: a bare Claude id and an opencode {@code provider/model} id are both
* just opaque strings here, so nothing here needs to know which {@code kind:} a profile runs.
*
* <p>The check runs one direction only, by construction: it reads {@link #profiles} and
* {@link #models}, and only ever adds to {@code bad} when a profile's model is missing from
* the allow-list. Nothing here can be satisfied by widening a {@code profiles:} entry — only
* editing {@code models.allow:} itself changes what passes. That is the invariant the ticket
* asked for: the list is the authority, profiles are checked against it.
*
* @throws IllegalStateException when any profile names a model outside the configured
* allow-list, naming both the model and the profile that wanted it
*/
public void validateModels() {
if (models == null || models.allow().isEmpty()) {
return;
}
Set<String> allowed = models.ids();
List<String> bad = new ArrayList<>();
profiles.forEach((name, p) -> {
String model = p.model();
if (model != null && !model.isBlank() && !allowed.contains(model)) {
bad.add("profile '" + name + "' names model '" + model + "', which is not in "
+ "models.allow: (have: " + allowed + ").");
}
});
if (!bad.isEmpty()) {
throw new IllegalStateException("refusing to start: " + String.join(" ", bad));
}
}
/**
* Runs every validator this class declares — found by reflection, not by name.
*
* <p>fleetd ticket "central allow-list of usable models", follow-up: mutation testing found
* that although each of the six validators above was well pinned on its own, nothing proved
* either real caller ({@code Fleetd.main} and {@link ConfigRef#reload()}) still
* invoked it — deleting a call site left the full suite green. The fix is not a seventh test
* per caller; a hand-maintained list of six names here would have the exact same defect its
* own javadoc would warn against: the seventh validator someone adds next month has no reason
* to be added to it. So this method does not name any validator. It sweeps {@link
* #getClass()}'s own public, no-argument, {@code void} methods whose name starts with {@code
* "validate"} (excluding itself) and invokes every one it finds, via {@link
* #invokeAllValidators}. A new {@code validateXxx()} method is therefore wired into both
* callers the moment it is written — there is no second step to forget, and so no state in
* which it silently never runs.
*
* <p>{@code Fleetd.main} and {@link ConfigRef#reload()} each call this one method instead of
* the six individually — see the comments at those two call sites for why
* each must run it.
*
* <p>Methods run in a fixed (alphabetical) order, so a config with more than one violation
* always names the same one first, on every run.
*
* @throws IllegalStateException (or whatever unchecked exception a validator itself throws),
* propagated unchanged from the first validator, in that order,
* that finds a problem
*/
public void validateAll() {
invokeAllValidators(this);
}
/**
* The reflective sweep behind {@link #validateAll()}, kept as its own method — taking any
* {@code target}, not just {@code this} — so a test can prove the MECHANISM is generic (it
* would sweep a seventh {@code validateXxx()} method added to any class, not just something
* special-cased to today's six on {@link FleetConfig}) without needing to add a real, unwanted
* seventh validator to this class just to exercise that claim. See {@code
* FleetConfigValidateAllTest} for that proof.
*
* @param target an object whose public, no-argument, {@code void} methods named {@code
* validateXxx} (any name starting with {@code "validate"}, excluding {@code
* validateAll} itself) should all run, in alphabetical-by-name order
*/
static void invokeAllValidators(Object target) {
List<Method> methods = new ArrayList<>();
for (Method m : target.getClass().getMethods()) {
if (Modifier.isPublic(m.getModifiers())
&& m.getParameterCount() == 0
&& m.getReturnType() == void.class
&& m.getName().startsWith("validate")
&& !m.getName().equals("validateAll")) {
methods.add(m);
}
}
methods.sort(Comparator.comparing(Method::getName));
for (Method m : methods) {
try {
m.invoke(target);
} catch (InvocationTargetException e) {
Throwable cause = e.getCause();
if (cause instanceof RuntimeException re) {
throw re;
}
if (cause instanceof Error err) {
throw err;
}
throw new IllegalStateException("validator " + m.getName() + " failed", cause);
} catch (IllegalAccessException e) {
throw new IllegalStateException("cannot invoke validator " + m.getName(), e);
}
}
}
/** True for the loopback addresses and the unspecified-but-local forms we treat as same-host. */
private static boolean isLoopbackBind(String host) {
if (host == null || host.isBlank()) {
@@ -38,15 +38,46 @@ public final class FleetHealthMonitor {
*/
static final long ASK_LAPSE_RECHECK_DELAY_SECONDS = 120;
/**
* fleetd #386: {@code System.nanoTime()} (or whatever {@link #clock} is) does not advance while
* macOS sleeps, so a raw {@code nowNanos - lastActivityAtNanos} comparison freezes with the
* host and can never cross {@link #workingSuspectAfterNanos}. This is a second, wall-clock
* source used ONLY inside the stall check ({@link #stallElapsedNanos}) to detect and correct
* for that freeze. Nothing else in this class reads it — every other decision (readiness grace,
* the fault classification itself) stays exactly on {@link #clock}, as the ticket requires.
*/
private static final LongSupplier DEFAULT_REALTIME_CLOCK =
() -> TimeUnit.MILLISECONDS.toNanos(System.currentTimeMillis());
private final AgentControl agents;
private final Supplier<List<MemberSession>> roster;
private final MessageService messages;
private final ScheduledExecutorService scheduler;
private final LongSupplier clock;
private final LongSupplier realtimeClock;
private final long intervalSeconds;
private final long tickIntervalNanos;
private final long workingSuspectAfterNanos;
private final BiConsumer<String, String> failTarget;
private final Map<String, HealthPrior> priors = new HashMap<>();
/**
* fleetd #386 clock-drift bookkeeping. {@code haveClockBaseline}/{@code lastTickMonoNanos}/
* {@code lastTickRealNanos} track the previous tick's pair of readings so each new tick can
* measure how far the two clocks moved apart since then. {@code accumulatedDriftNanos} is the
* running total of every such divergence observed since this monitor started (never decreases —
* the monotonic clock can only lag real time, never lead it). {@code busyDriftBaselineNanos}/
* {@code busyBaselineActivityNanos} record, per target, the value of {@code accumulatedDriftNanos}
* at the moment this monitor first saw that target's CURRENT {@code lastActivityAtNanos} while
* BUSY — so {@link #stallElapsedNanos} adds back only the drift observed DURING this BUSY span,
* never drift from a sleep that happened before the member went busy. All five fields are touched
* only from {@code tick()}, like {@link #priors}.
*/
private boolean haveClockBaseline = false;
private long lastTickMonoNanos;
private long lastTickRealNanos;
private long accumulatedDriftNanos = 0;
private final Map<String, Long> busyDriftBaselineNanos = new HashMap<>();
private final Map<String, Long> busyBaselineActivityNanos = new HashMap<>();
/**
* The live classification per member, and the only one of this class's three maps that more
* than one scheduler task touches. {@code tick} writes it (and prunes it to the roster);
@@ -88,12 +119,29 @@ public final class FleetHealthMonitor {
public FleetHealthMonitor(AgentControl agents, Supplier<List<MemberSession>> roster, MessageService messages,
ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds,
long workingSuspectAfterSeconds, BiConsumer<String, String> failTarget) {
this(agents, roster, messages, scheduler, clock, DEFAULT_REALTIME_CLOCK, intervalSeconds,
workingSuspectAfterSeconds, failTarget);
}
/**
* @param realtimeClock fleetd #386: a wall-clock nanosecond source (e.g.
* {@code System.currentTimeMillis()} converted to nanos) that keeps
* advancing while {@code clock} is frozen by a host sleep. Used only to
* correct the stall check — see the class-level javadoc on the
* clock-drift fields.
*/
public FleetHealthMonitor(AgentControl agents, Supplier<List<MemberSession>> roster, MessageService messages,
ScheduledExecutorService scheduler, LongSupplier clock, LongSupplier realtimeClock,
long intervalSeconds, long workingSuspectAfterSeconds,
BiConsumer<String, String> failTarget) {
this.agents = agents;
this.roster = roster;
this.messages = messages;
this.scheduler = scheduler;
this.clock = clock;
this.realtimeClock = Objects.requireNonNull(realtimeClock, "realtimeClock");
this.intervalSeconds = intervalSeconds;
this.tickIntervalNanos = TimeUnit.SECONDS.toNanos(intervalSeconds);
this.workingSuspectAfterNanos = TimeUnit.SECONDS.toNanos(workingSuspectAfterSeconds);
this.failTarget = Objects.requireNonNull(failTarget, "failTarget");
}
@@ -123,6 +171,7 @@ public final class FleetHealthMonitor {
for (Agent agent : agentsNow) live.put(agent.terminalId(), agent);
HashSet<String> current = new HashSet<>();
long nowNanos = clock.getAsLong();
long driftBeforeThisTick = observeClockDrift(nowNanos);
for (MemberSession session : rosterNow) {
current.add(session.terminalId());
Agent agent = live.get(session.terminalId());
@@ -133,7 +182,7 @@ public final class FleetHealthMonitor {
&& session.state() != MemberSession.State.SPAWNING;
boolean readinessGraceElapsed = nowNanos - session.spawnedAtNanos() >= READINESS_GRACE_NANOS;
boolean stalled = session.state() == MemberSession.State.BUSY
&& nowNanos - session.lastActivityAtNanos() >= workingSuspectAfterNanos;
&& stallElapsedNanos(session, nowNanos, driftBeforeThisTick) >= workingSuspectAfterNanos;
// CB-643: the three message-layer facts CB-640 published. Read them here rather than
// leaving them false — that constant is what made 8 of the 9 fault states dead.
boolean queuedDelivery = messages.hasQueuedDelivery(session.terminalId());
@@ -151,6 +200,8 @@ public final class FleetHealthMonitor {
priors.keySet().retainAll(current);
states.keySet().retainAll(current);
orphanStreaks.keySet().retainAll(current);
busyDriftBaselineNanos.keySet().retainAll(current);
busyBaselineActivityNanos.keySet().retainAll(current);
} catch (Throwable error) {
// Any unclassified collection failure must never kill the monitor's only scheduler task.
log.warn("fleet health collection failed; will retry next tick", error);
@@ -175,6 +226,65 @@ public final class FleetHealthMonitor {
return streak >= ORPHAN_CONFIRM_TICKS;
}
/**
* fleetd #386: compare this tick's monotonic and real-time readings against the previous
* tick's, and fold any positive divergence into {@link #accumulatedDriftNanos} (a ratchet — it
* never decreases, since the monotonic clock can only fall behind real time, never ahead of
* it). Logs once, at WARN, when that single tick's divergence exceeds one full tick interval —
* the signature of a host that slept between the two ticks (a tick literally cannot run while
* the process itself is suspended, so the whole sleep duration lands inside one tick's gap).
*
* @return {@link #accumulatedDriftNanos} as it stood BEFORE this tick's divergence was folded
* in — the baseline {@link #stallElapsedNanos} needs when a target is observed BUSY
* for the first time this tick, so a sleep that happened before this member went busy
* is not attributed to it.
*/
private long observeClockDrift(long nowNanos) {
long nowRealNanos = realtimeClock.getAsLong();
long driftBeforeThisTick = accumulatedDriftNanos;
if (haveClockBaseline) {
long monoDelta = nowNanos - lastTickMonoNanos;
long realDelta = nowRealNanos - lastTickRealNanos;
long tickDrift = realDelta - monoDelta;
if (tickDrift > tickIntervalNanos) {
log.warn("fleet health: the monotonic clock did not advance for about {}s that the "
+ "real clock did since the last tick (host likely slept); the stall "
+ "detector could not see that time", TimeUnit.NANOSECONDS.toSeconds(tickDrift));
}
if (tickDrift > 0) {
accumulatedDriftNanos = driftBeforeThisTick + tickDrift;
}
}
lastTickMonoNanos = nowNanos;
lastTickRealNanos = nowRealNanos;
haveClockBaseline = true;
return driftBeforeThisTick;
}
/**
* fleetd #386: {@code nowNanos - lastActivityAtNanos} alone freezes across a host sleep, since
* both come from the monotonic {@link #clock}. This adds back the real-time drift observed
* since this BUSY span started — not the monitor's whole lifetime, so a sleep that happened
* before this member went busy never leaks into its stall reading (see the class-level javadoc
* on the drift fields). The baseline resets whenever {@code lastActivityAtNanos} changes (a new
* turn) or the member is not currently BUSY.
*/
private long stallElapsedNanos(MemberSession session, long nowNanos, long driftBeforeThisTick) {
String target = session.terminalId();
if (session.state() != MemberSession.State.BUSY) {
busyDriftBaselineNanos.remove(target);
busyBaselineActivityNanos.remove(target);
return nowNanos - session.lastActivityAtNanos();
}
Long baselineActivity = busyBaselineActivityNanos.get(target);
if (baselineActivity == null || baselineActivity != session.lastActivityAtNanos()) {
busyBaselineActivityNanos.put(target, session.lastActivityAtNanos());
busyDriftBaselineNanos.put(target, driftBeforeThisTick);
}
long driftSinceBusyStart = accumulatedDriftNanos - busyDriftBaselineNanos.get(target);
return (nowNanos - session.lastActivityAtNanos()) + driftSinceBusyStart;
}
void reportTransition(String target, HealthState next) {
HealthState previous = states.put(target, next);
if (previous == next) return;
@@ -548,8 +548,19 @@ public final class CompletionResolver implements TurnListener {
* fast completion. Runs the same backend-error classification the normal and raw-scrape paths
* apply, against whatever is on screen right now: a match together with the too-fast crash
* signature notifies {@link #backendErrorSink} (only on the resolution that wins the race). A
* non-match stays the original generic too-fast failure, naming the member and both timings,
* with whatever the pane shows appended so the caller sees the cause, not just "it failed".
* non-match stays the generic too-fast failure, naming the member and both timings, with
* whatever the pane shows appended so the caller sees the evidence, not just "it failed".
*
* <p>fleetd#376: <strong>this path must never resolve a completion.</strong> A fast backend can
* genuinely answer inside the floor, so the failure is sometimes wrong — but it is wrong in the
* loud direction, and the fix for that is honest wording, not a guess at the pane's meaning.
* Reclassifying from the scrape was tried and rejected: there is no reliable positive marker for
* "this is a real reply" across backends. {@link #lastAssistantBlock} falls back to the entire
* pane when it finds no {@code ⏺} marker, so on a crash the candidate "reply" is the whole
* screen; and {@code ⏺} itself is a Claude Code marker that an opencode pane never carries — the
* very backend whose speed raised this ticket. Any weaker test (non-blank, or "contains sentence
* punctuation") passes on almost every crash, because a pane holding a file path, a version
* number or a hostname contains a full stop. That trades a loud wrong answer for a silent one.
*/
private void failTooFast(String target, InFlight turn, CompletableFuture<Rendezvous.Resolution> waiter,
long elapsedNanos) {
@@ -560,13 +571,13 @@ public final class CompletionResolver implements TurnListener {
scrape = "";
}
String clippedScrape = clip(scrape);
String baseReason = String.format(
"member %s went BUSY -> DONE in %dms (floor %dms) — too fast to be real work, most "
+ "likely a backend error before any work started",
String timing = String.format(
"member %s went BUSY -> DONE in %dms (floor %dms)",
target, elapsedNanos / 1_000_000, MIN_TURN_NANOS / 1_000_000);
String backendError = firstMatchingLine(scrape, backendErrorPatternOrFallback(target));
if (backendError != null) {
String reason = baseReason + ": " + clippedScrape;
String reason = timing + " — too fast to be real work, and the pane carries a backend "
+ "error: " + clippedScrape;
if (rendezvous.resolveFailure(waiter, reason)) {
inFlight.remove(target, turn);
log.warn("failing send to {} via turn-stall fallback: {}", target, reason);
@@ -575,7 +586,14 @@ public final class CompletionResolver implements TurnListener {
}
return;
}
fail(target, turn, clippedScrape.isBlank() ? baseReason : baseReason + ": " + clippedScrape);
// fleetd#376: no pattern matched, so the cause is genuinely unknown. Say that, rather than
// asserting a backend error the way this message used to — a fast backend really can finish
// inside the floor, and a reader who trusts a wrong cause stops looking at the pane.
String reason = timing + " — inside the floor. That is usually a backend error before any "
+ "work started, but a fast backend can answer inside it too, and nothing here can "
+ "tell those apart, so the turn is reported failed rather than guessed. Read the "
+ "pane below before deciding which it was";
fail(target, turn, clippedScrape.isBlank() ? reason : reason + ": " + clippedScrape);
}
/**
@@ -119,9 +119,28 @@ public final class FleetMcp {
/**
* CB-578 stage B quarantine facts used by {@code fleet_profiles}: a profile → credential id
* lookup, plus the shared {@link BackendQuarantine} to read remaining cooldowns off.
*
* @param exhaustedPatternArmed fleetd #395: profile → whether that profile's {@code
* exhaustedPattern} is configured (see {@code
* FleetConfig.Profile#hasExhaustedPattern}), i.e. whether a backend refusal
* on it can EVER be classified {@code BACKEND_EXHAUSTED} and quarantine its
* credential. Bundled here, not a separate Source, because it answers the
* exact question {@code fleet_profiles}'s quarantine facts already answer
* for a QUARANTINED profile — "can this profile's usage limit ever be
* caught?" — just for every profile, not only one currently caught.
*/
public record QuarantineSource(Function<String, String> credentialIdFor, BackendQuarantine quarantine) {
/** Inert source — no profile is ever reported quarantined. Explicit stand-in, not a default. */
public record QuarantineSource(Function<String, String> credentialIdFor, BackendQuarantine quarantine,
Function<String, Boolean> exhaustedPatternArmed) {
/**
* Backward-compatible 2-arg form, before fleetd #395 added {@code exhaustedPatternArmed} —
* reports every profile unarmed. Keeps every pre-existing call site (production and test)
* compiling and behaving identically for the quarantine facts they actually asked for.
*/
public QuarantineSource(Function<String, String> credentialIdFor, BackendQuarantine quarantine) {
this(credentialIdFor, quarantine, _ -> false);
}
/** Inert source — no profile is ever reported quarantined or armed. Explicit stand-in, not a default. */
public static QuarantineSource none() { return new QuarantineSource(_ -> null, BackendQuarantine.none()); }
}
@@ -346,7 +365,7 @@ public final class FleetMcp {
String self = callerTerminal(exchange);
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_reply", req.arguments()), self);
if (denied != null) return denied;
return reply(messages, self, str(req.arguments(), "content"));
return reply(messages, self, principal(exchange).role(), str(req.arguments(), "content"));
};
// fleet_ask (CB-205): a worker's mid-turn question — identity from the CONNECTION.
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> askHandler =
@@ -868,14 +887,26 @@ public final class FleetMcp {
/**
* {@code fleet_reply}: the worker returns its structured answer, resolving the awaiting send
* or — when no send is open — queueing the reply in the inbox for later drain (CB-307).
* {@code callerTerminal} is resolved from the connection (never an argument); a {@code null}
* means the caller is not a known worker (e.g. the primary called it by mistake).
* {@code callerTerminal} and {@code callerRole} are resolved from the connection (never an
* argument). A {@code null} terminal means the caller is not a known worker. A PRIMARY with a
* terminal is a lead and must use {@code fleet_send}, because reply has no peer-lead route.
*
* <p>fleetd #365: the result text names which of those actually happened
* ({@link MessageService.ReplyOutcome#description()}) instead of the single word "delivered"
* for both — a queued reply is a real success, but it is not the same fact as one that resolved
* a live waiter, and the caller could not previously tell them apart.
*/
static McpSchema.CallToolResult reply(MessageService messages, String callerTerminal, String content) {
static McpSchema.CallToolResult reply(MessageService messages, String callerTerminal, Role callerRole, String content) {
if (callerTerminal == null) {
return error("fleet_reply is for workers only — could not identify the calling worker "
+ "from the connection");
}
if (callerRole == Role.PRIMARY) {
return error("fleet_reply has no route to a peer lead. Use fleet_send{coordId: ...} for a peer on another "
+ "daemon or fleet_send{sessionId: ...} for a peer on this host. fleet_reply resolves a member's "
+ "blocked fleet_send, and a peer's coord-id message is durable and non-blocking, so there is "
+ "nothing for it to resolve.");
}
// fleetd #302: isBlank, not == null, to match fleet_send's own guard above. MessageService
// .reply now REJECTS blank content, and this handler is a bare BiFunction with no try/catch
// around it — so a whitespace-only fleet_reply would leave here as an uncaught
@@ -884,8 +915,8 @@ public final class FleetMcp {
if (isBlank(content)) {
return error("content is required");
}
messages.reply(callerTerminal, content);
return text("delivered");
MessageService.ReplyOutcome outcome = messages.reply(callerTerminal, content);
return text(outcome.description());
}
/** {@code fleet_ack}: acknowledge (remove) a specific reply from the inbox. */
@@ -1104,6 +1135,14 @@ public final class FleetMcp {
* not the other, and the two then disagree about a live outage. That is exactly what fleetd
* #284 was, where one rule computed in two places was widened in only one and a single response
* contradicted itself. Shared inputs do not make duplicated computation safe.
*
* <p>fleetd #395: also reports {@code exhaustionDetectionArmed}, one boolean per configured
* profile — {@code true} when that profile's {@code exhaustedPattern} is set, {@code false}
* when it is not, so an operator can tell "this profile is healthy" from "nothing can ever
* quarantine this profile" without reading {@code fleetd.yaml}. Unlike {@code quarantined}/
* {@code coolingOff}, this map always names every profile: an unarmed profile never enters a
* transient state to be absent from, so silence here would read as "healthy" rather than "not
* being watched at all".
*/
public static Map<String, Object> profilesView(PeerLauncher workers, QuarantineSource quarantine, OutageSource outage) {
Map<String, Object> result = new LinkedHashMap<>();
@@ -1111,7 +1150,9 @@ public final class FleetMcp {
result.put("default", workers.defaultProfile() == null ? "" : workers.defaultProfile());
Map<String, Object> quarantined = new LinkedHashMap<>();
Map<String, Object> coolingOff = new LinkedHashMap<>();
Map<String, Object> exhaustionDetectionArmed = new LinkedHashMap<>();
for (String profile : workers.profiles()) {
exhaustionDetectionArmed.put(profile, quarantine.exhaustedPatternArmed().apply(profile));
String credentialId = quarantine.credentialIdFor().apply(profile);
if (credentialId != null) {
quarantine.quarantine().remainingSeconds(credentialId).ifPresent(remaining -> {
@@ -1131,6 +1172,7 @@ public final class FleetMcp {
});
}
}
result.put("exhaustionDetectionArmed", exhaustionDetectionArmed);
if (!quarantined.isEmpty()) {
result.put("quarantined", quarantined);
}
@@ -1658,7 +1700,11 @@ public final class FleetMcp {
"List the configured worker profiles (backends) and which one fleet_spawn uses by "
+ "default. A 'quarantined' map is present when a backend-exhausted refusal put "
+ "a profile's credential on cooldown — fleet_spawn onto it is refused until "
+ "quarantinedForSeconds elapses; a profile sharing that credential is listed too.",
+ "quarantinedForSeconds elapses; a profile sharing that credential is listed too. "
+ "'exhaustionDetectionArmed' reports, per profile, whether a usage-limit refusal "
+ "on it can EVER be classified and quarantined (its exhaustedPattern is "
+ "configured) — false means that profile's credential can never be quarantined "
+ "by this mechanism, however many usage-limit refusals it sees.",
objectSchema(Map.of(), List.of()));
}
@@ -369,8 +369,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
// CB-547a: route the chosen profile but keep the caller's session identity — dropping it
// here would silently sever the resume handle on every policy-routed spawn. CB-557: the
// role rides along for the same reason, or a routed spawn would be labelled as a dev.
SpawnRequest routedReq = new SpawnRequest(chosen.profile(), req.requestedCwd(), req.callerCwd(),
req.sessionName(), req.resumeSessionId(), req.role());
SpawnRequest routedReq = req.withProfile(chosen.profile());
try {
PeerHandle handle = d.spawn(routedReq);
spawnedBy.put(handle.id(), d);
@@ -13,6 +13,7 @@ import java.nio.file.attribute.PosixFilePermissions;
import java.time.Duration;
import java.time.Instant;
import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
import java.util.stream.Stream;
@@ -28,33 +29,85 @@ import java.util.stream.Stream;
* control lacked: herdr applies that overlay BEFORE the shell starts, so any sourced file can undo
* it — and did.
*
* <p><b>Which file is last depends on the platform, so the scrub runs from two of them.</b> zsh
* <p><b>zsh reads its four startup files under three different conditions, so no single file is
* guaranteed to run — the scrub has to cover the gap between them, not just the platforms.</b> zsh
* reads {@code .zshenv} always, {@code .zprofile} and {@code .zlogin} only for a LOGIN shell, and
* {@code .zshrc} only for an INTERACTIVE one. herdr does not open the same kind of shell
* everywhere — measured on herdr 0.8.0: macOS panes run {@code -zsh} (login, so {@code .zlogin}
* runs), Linux panes run a plain {@code /usr/bin/zsh} (interactive but NOT login, so
* {@code .zlogin} never runs at all). A scrub in {@code .zlogin} alone is therefore a control that
* silently does nothing on Linux — the exact failure this class exists to remove, one platform
* over.
* {@code .zshrc} only for an INTERACTIVE one. A pane shell that is at least one of login or
* interactive is covered by sourcing the scrub from {@code .zshrc} and {@code .zlogin} (below), but
* a pane shell that is NEITHER reads only {@code .zshenv} and stops — fleetd #388, measured: a herdr
* pane can be neither login nor interactive, and such a pane read {@code .zshenv}, never reached
* {@code scrub.zsh}, and left no report at all. A bare {@code argv[0]} of {@code /usr/bin/zsh}
* proves the shell is NOT a login shell; it says nothing about whether it is interactive, so it
* must never be read as "therefore interactive" — that wrong inference is what let #388 ship.
*
* <p>So both {@code .zshrc} and {@code .zlogin} source the same generated {@code scrub.zsh} after
* sourcing their {@code $HOME} counterpart. On Linux only the first fires; on macOS both do, and
* the second pass is deliberate rather than merely harmless — it re-scrubs anything the operator's
* own {@code ~/.zlogin} exported after {@code .zshrc} had finished. Re-running is idempotent: a
* name already blank is blanked again, and the report is rewritten with the same counts.
* <p>So {@code .zshenv} carries a THIRD pass, guarded by the exact condition that defines the gap:
* {@code [[ ! -o login && ! -o interactive ]]}. That guard is why this pass cannot double-scrub a
* pane that {@code .zshrc} or {@code .zlogin} will also cover — one of {@code -o login}/
* {@code -o interactive} is always true there, so the {@code .zshenv} pass never fires for them, and
* their own unconditional sourcing is untouched. The guard also carries a sentinel
* ({@value #SCRUB_SENTINEL}) so it fires once per PANE and not once per PROCESS: {@code .zshenv} is
* read by every zsh a member's own tooling forks (a plain {@code zsh -c '...'} for a single
* command is itself neither login nor interactive), and those children inherit variables their
* parent deliberately set for them (git hooks get {@code GIT_DIR}, a venv gets
* {@code VIRTUAL_ENV}, a build tool gets {@code NODE_OPTIONS} or {@code JAVA_TOOL_OPTIONS}).
* Re-scrubbing every such child would blank all of that, and would also make the pane's own
* {@code scrub-report.txt} — rewritten on every pass — describe whichever child exited last
* instead of the pane. The sentinel is exported only AFTER {@code scrub.zsh} runs, so the pass
* that sets it never sees it and cannot blank it; it must also be on the scrub's own allow-list
* (see {@link #generate(Path, Set)}) so a later pass, in the same pane, cannot blank it back to
* empty — an exported-but-empty sentinel reads as unset to the {@code -z} guard and would silently
* re-enable scrubbing for every subsequent child of that pane.
*
* <p>So all three of {@code .zshenv} (gap only, guarded), {@code .zshrc}, and {@code .zlogin}
* source the same generated {@code scrub.zsh} after sourcing their {@code $HOME} counterpart. A
* login-and-interactive pane runs the {@code .zshrc} and {@code .zlogin} passes, and the second is
* deliberate rather than merely harmless — it re-scrubs anything the operator's own
* {@code ~/.zlogin} exported after {@code .zshrc} had finished. A pane that is neither runs only the
* {@code .zshenv} pass. Re-running is idempotent: a name already blank is blanked again, and the
* report is rewritten with the same counts.
*
* <p>Each generated file sources its {@code $HOME} counterpart FIRST, so {@code PATH} and every
* toolchain binary still resolve exactly as the operator configured them; only afterwards does
* {@code .zlogin} run the scrub: every EXPORTED variable not on the derived allow-list is re-exported
* toolchain binary still resolve exactly as the operator configured them; only afterwards does the
* scrub run: every EXPORTED variable not on the derived allow-list is re-exported
* blank. Blank, not credential-shaped-pattern-filtered: a pattern list ({@code *TOKEN*}, …) is an
* enumeration and misses what it did not think of — a username is the other half of a credential and
* is shaped like none. Credential-SHAPED names among the blanked set go to the WARN log only,
* never to the control.
*
* <p>The scrub also writes {@code scrub-report.txt} into its own directory: one {@code allowed N of
* M} line (N = exports left untouched, M = exports present when the scrub ran), then the blanked
* NAMES — never values. The launcher reads this back at teardown and logs it, because a blocked
* count next to an unknown denominator is not a finding.
* M failed F} line (N = exports left untouched, M = exports present when the scrub ran, F = names
* the scrub attempted to blank but could not), then the NAMES — blanked ones bare, unblankable ones
* {@code !}-prefixed — never values. The launcher reads this back at teardown and logs it, because a
* blocked count next to an unknown denominator is not a finding.
*
* <p><b>fleetd #394:</b> plain {@code export "$n="} is a FATAL error for a zsh read-only or special
* parameter (for example {@code UID}) — it aborts the whole sourced file, so every name still to
* come is never blanked and the report above is never written at all. The blanking loop instead
* routes each attempt through {@code eval}, which contains that error to the single iteration: the
* loop always finishes, and a name that could not be blanked is counted as {@code failed} and
* listed {@code !}-prefixed rather than silently disappearing. This is deliberately not a skip-list
* of known-bad names — every enumerated name is still attempted, so a name nobody has thought of
* yet still gets tried and, if it fails, still gets counted.
*
* <p>The blanking loop also re-asserts, on its own, the same {@code [A-Za-z_][A-Za-z0-9_]*} shape
* check the enumeration loop already applied. Before {@code eval} was introduced a non-conforming
* name reaching {@code export "$n="} was harmless either way — the quoting made it inert. With
* {@code eval}, the name is spliced into a string and interpreted as shell syntax, so the enumeration
* loop's check is no longer sufficient on its own to keep that call site safe — it is a guard on a
* different loop, and the two must not silently drift apart. Re-checking right before the
* {@code eval} keeps that call site safe by its own reading, independent of whatever the enumeration
* loop does or stops doing in a later change.
*
* <p><b>fleetd #400:</b> {@code eval}'s exit status is not proof that the blank actually happened.
* zsh coerces a bare {@code NAME=} assignment on an integer special parameter (measured on macOS zsh
* 5.9: {@code SECONDS}, {@code RANDOM}, {@code SHLVL}, {@code HISTSIZE}, {@code COLUMNS},
* {@code LINES}, {@code USERNAME}) to a number instead of failing — {@code eval} returns success,
* the value is untouched, and a status-based classification reports it as blanked when it was not.
* The fix classifies on the observed effect instead: after the attempt, the name's value is read
* back with the {@code (P)} indirection flag and the decision is made from whether that is now
* empty. This one check covers all three shapes a name can take at this point — a genuine blank, a
* fatal read-only error {@code eval} merely contained, and this silent no-op — and the exit status
* plays no part in the decision at all.
*/
public final class EnvAllowListScrub {
@@ -63,7 +116,11 @@ public final class EnvAllowListScrub {
/** Name of the report file written into the generated directory by the scrub itself. */
static final String REPORT_FILE = "scrub-report.txt";
/** The scrub body, generated once and sourced from both {@code .zshrc} and {@code .zlogin}. */
/**
* The scrub body, generated once and sourced from {@code .zshrc} and {@code .zlogin}
* unconditionally, and from {@code .zshenv} when the pane shell is neither login nor
* interactive (fleetd #388) — see the class javadoc.
*/
static final String SCRUB_FILE = "scrub.zsh";
/** Prefix of every generated directory — also what {@link #reapOrphans} matches on. */
@@ -79,14 +136,42 @@ public final class EnvAllowListScrub {
private static final String SOURCE_SCRUB =
"source \"$ZDOTDIR/" + SCRUB_FILE + "\"\n";
/**
* fleetd #388: marks a pane, not a process, as already scrubbed. Set only by the guarded
* {@code .zshenv} pass (see {@link #NEITHER_LOGIN_NOR_INTERACTIVE_SCRUB}) after
* {@code scrub.zsh} has run, so it must also be folded into that pass's own allow-list — see
* the class javadoc's "must also be on the scrub's own allow-list" paragraph.
*/
static final String SCRUB_SENTINEL = "_CB633_SCRUBBED";
/**
* Appended to {@code .zshenv}, after its {@code $HOME} source: the third pass, guarded on the
* exact condition that defines the gap {@code .zshrc}/{@code .zlogin} do not cover — a shell
* that is neither login nor interactive. The sentinel export happens only once the scrub has
* already run, and only for as long as the current pane's environment has not been rebuilt from
* scratch (a fresh {@code env -i} child would not inherit it — that is out of scope here, since
* such a child is no longer running under the pane's own environment at all).
*/
private static final String NEITHER_LOGIN_NOR_INTERACTIVE_SCRUB =
"if [[ ! -o login && ! -o interactive && -z \"${" + SCRUB_SENTINEL + ":-}\" ]]; then\n"
+ " " + SOURCE_SCRUB
+ " export " + SCRUB_SENTINEL + "=1\n"
+ "fi\n";
private EnvAllowListScrub() {
}
/**
* A parsed {@code scrub-report.txt}: how many exported variables existed when the scrub ran,
* how many were left untouched (allowed), and the NAMES that were blanked. Values never appear.
* A parsed {@code scrub-report.txt}: how many exported variables existed when the scrub ran
* ({@code total}), how many were left untouched ({@code allowed}), how many the scrub attempted
* to blank but could not ({@code failed} — fleetd #394: a zsh read-only/special parameter such
* as {@code UID} fatally errors on plain {@code export NAME=}, so those attempts go through
* {@code eval} instead so the loop keeps going and the failure is counted rather than left
* invisible), and the NAMES in each of the latter two categories. {@code allowed +
* blanked.size() + unblankable.size() == total}, and {@code unblankable.size() == failed}.
* Values never appear.
*/
record ScrubReport(int allowed, int total, List<String> blanked) {
record ScrubReport(int allowed, int total, int failed, List<String> blanked, List<String> unblankable) {
}
/**
@@ -105,11 +190,16 @@ public final class EnvAllowListScrub {
reapOrphans(parentDir);
Path dir = Files.createTempDirectory(parentDir, DIR_PREFIX);
dir.toFile().deleteOnExit();
// The report is written by zsh, after these hooks are registered, so register its path
// too — otherwise the directory is non-empty at JVM exit and cannot be removed at all.
dir.resolve(REPORT_FILE).toFile().deleteOnExit();
write(dir, SCRUB_FILE, scrubScript(allowedNames));
write(dir, ".zshenv", homeSourcingFile(".zshenv"));
// zsh truncates this pre-created receipt after these hooks are registered. Register its
// path too — otherwise the directory is non-empty at JVM exit and cannot be removed.
Files.createFile(dir.resolve(REPORT_FILE)).toFile().deleteOnExit();
// fleetd #388: scrub.zsh's OWN allow-list must also keep SCRUB_SENTINEL, or a later
// pass in the same pane blanks it back to empty and the .zshenv guard below thinks it
// was never scrubbed — see the class javadoc.
Set<String> namesForScrubScript = new HashSet<>(allowedNames);
namesForScrubScript.add(SCRUB_SENTINEL);
write(dir, SCRUB_FILE, scrubScript(namesForScrubScript));
write(dir, ".zshenv", homeSourcingFile(".zshenv") + NEITHER_LOGIN_NOR_INTERACTIVE_SCRUB);
write(dir, ".zprofile", homeSourcingFile(".zprofile"));
write(dir, ".zshrc", homeSourcingFile(".zshrc") + SOURCE_SCRUB);
write(dir, ".zlogin", homeSourcingFile(".zlogin") + SOURCE_SCRUB);
@@ -147,12 +237,11 @@ public final class EnvAllowListScrub {
* into it: owner keeps full access, {@code group} gets traverse+read on the directory ({@code
* rwxr-x---}, so a member process — a login shell reading it via {@code ZDOTDIR}, or another
* process simply opening a file under it — running under that group can find and read the
* files) and read-only on each file ({@code rw-r-----}) — deliberately no group WRITE anywhere,
* since a member never needs to add or change what fleetd generated. (For the ZDOTDIR scrub
* specifically, this also means the scrub script's own report write inside the pane fails
* closed rather than open — see {@code scrub.zsh}'s trailing {@code 2>/dev/null} — which {@link
* dev.ltms.fleet.member.HerdrPeerLauncher#releaseZdotdir} already treats as "cannot be
* confirmed to have run" rather than success.)
* files) and read-only on each file ({@code rw-r-----}), except the pre-created ZDOTDIR
* {@code scrub-report.txt}. That receipt gets group write ({@code rw-rw----}), so
* {@code scrub.zsh} can truncate and write it without granting group write on the directory.
* If its optional permission change fails, the member cannot write a receipt and the launcher
* keeps its existing WARN rather than failing the spawn.
*
* <p>Package-private and named generically on purpose: fleetd #213 built this for the ZDOTDIR
* scrub directory, and fleetd #219 reuses it verbatim for {@link
@@ -167,7 +256,15 @@ public final class EnvAllowListScrub {
setGroupAndPermissions(dir, principal, "rwxr-x---");
try (Stream<Path> entries = Files.list(dir)) {
for (Path file : entries.toList()) {
setGroupAndPermissions(file, principal, "rw-r-----");
if (REPORT_FILE.equals(file.getFileName().toString())) {
try {
setGroupAndPermissions(file, principal, "rw-rw----");
} catch (IOException | UnsupportedOperationException ignored) {
// The receipt is optional. Its absence keeps the existing WARN path.
}
} else {
setGroupAndPermissions(file, principal, "rw-r-----");
}
}
}
} catch (IOException e) {
@@ -245,11 +342,13 @@ public final class EnvAllowListScrub {
}
return """
# generated by fleetd (CB-633 memberCredentials policy=allow-list) — do not edit.
# Sourced from .zshrc and again from .zlogin, each time AFTER that file has sourced
# its $HOME counterpart — so this runs after everything the operator sourced, on a
# login shell (macOS panes) and on a plain interactive one (Linux panes) alike.
# Running twice is idempotent and deliberate: the second pass catches anything
# ~/.zlogin exported after ~/.zshrc had finished.
# Sourced unconditionally from .zshrc and again from .zlogin, each time AFTER that
# file has sourced its $HOME counterpart — so this runs after everything the
# operator sourced, on any pane that is login and/or interactive. Running twice is
# idempotent and deliberate: the second pass catches anything ~/.zlogin exported
# after ~/.zshrc had finished. Also sourced, once, from a guarded pass in .zshenv
# (fleetd #388) when the pane shell is NEITHER login nor interactive — the one gap
# those two files do not cover.
typeset -A _cb633_allowed
for _cb633_n in %s; do _cb633_allowed[$_cb633_n]=1; done
@@ -271,15 +370,57 @@ public final class EnvAllowListScrub {
_cb633_blank+=("$_cb633_n")
done
{ for _cb633_n in "${_cb633_blank[@]}"; do export "$_cb633_n="; done; } 2>/dev/null
# fleetd #394: plain `export "$n="` is FATAL for a zsh read-only/special parameter
# (e.g. UID) and aborts this whole sourced file — every name still to come is never
# blanked, and the report below is never written, silently. `eval` contains that
# error to the single iteration instead: it still fails for that one name, but the
# loop continues and we can tell allowed / blanked / unblankable apart afterwards.
# This is not a skip-list of known-bad names (that would miss the next one nobody
# thought of) — every name in _cb633_blank is still attempted, unconditionally.
# Every name reaching this loop already passed the identical identifier check in the
# enumeration loop above — but that guard is 20 lines away in a different loop, and
# this line is about to splice the name into a string handed to `eval`. Before this
# fix the name only ever reached `export` quoted ("$n="), which is inert on a
# non-identifier string either way; `eval` makes THIS line the only thing standing
# between such a string and code execution in the member's pane, so it re-asserts the
# same check on its own rather than trusting a guard it does not own. Under normal
# operation this can never fire (the enumeration guard already filtered everything
# reaching _cb633_blank), so a name caught here is counted as unblankable rather than
# silently dropped — it is real evidence that the upstream guard was bypassed.
#
# fleetd #400: the attempt's own exit status is NOT proof of its effect. zsh coerces
# a bare `NAME=` assignment on an integer special parameter (SECONDS, RANDOM, SHLVL,
# HISTSIZE, COLUMNS, LINES, USERNAME on this host) to a number instead of failing —
# `eval` returns 0, the value is untouched, and the old exit-status check reported it
# as blanked when it was not. Classify on the observed effect instead: attempt the
# export, then read the name's value back with the `(P)` indirection flag and decide
# from whether it is now empty. One check then covers all three shapes a name can
# take here — a genuine blank, a fatal read-only error `eval` merely contained, and
# this silent no-op — without the exit status entering the decision at all.
typeset -a _cb633_ok _cb633_unblankable
_cb633_ok=()
_cb633_unblankable=()
for _cb633_n in "${_cb633_blank[@]}"; do
if [[ ! "$_cb633_n" =~ ^[A-Za-z_][A-Za-z0-9_]*$ ]]; then
_cb633_unblankable+=("$_cb633_n")
continue
fi
eval "export ${_cb633_n}=" 2>/dev/null
if [[ -z "${(P)_cb633_n}" ]]; then
_cb633_ok+=("$_cb633_n")
else
_cb633_unblankable+=("$_cb633_n")
fi
done
integer _cb633_kept=$(( _cb633_total - ${#_cb633_blank} ))
{
print -r -- "allowed $_cb633_kept of $_cb633_total"
for _cb633_n in "${_cb633_blank[@]}"; do print -r -- "$_cb633_n"; done
print -r -- "allowed $_cb633_kept of $_cb633_total failed ${#_cb633_unblankable}"
for _cb633_n in "${_cb633_ok[@]}"; do print -r -- "$_cb633_n"; done
for _cb633_n in "${_cb633_unblankable[@]}"; do print -r -- "!$_cb633_n"; done
} > "$ZDOTDIR/%s" 2>/dev/null
unset _cb633_allowed _cb633_names _cb633_blank _cb633_n _cb633_total _cb633_kept
unset _cb633_allowed _cb633_names _cb633_blank _cb633_ok _cb633_unblankable _cb633_n _cb633_total _cb633_kept
""".formatted(names, MemberEnvAllowList.zshCasePattern(), REPORT_FILE);
}
@@ -293,6 +434,11 @@ public final class EnvAllowListScrub {
* Read and parse {@link #REPORT_FILE} out of a generated ZDOTDIR directory. Returns {@code null}
* when absent or unreadable (the pane may have been torn down before its login shell ever got to
* the scrub) — callers treat that as "no measurement available", never as success.
*
* <p>First line is {@code "allowed <N> of <M> failed <F>"} (fleetd #394 added the trailing
* {@code failed <F>} — a count of names the scrub attempted to blank but could not, e.g. a zsh
* read-only/special parameter). Every following non-blank line is a name: a bare name was
* blanked, a {@code !}-prefixed name was attempted and failed. Values never appear on either.
*/
static ScrubReport readReport(Path zdotdir) {
Path report = zdotdir.resolve(REPORT_FILE);
@@ -305,17 +451,24 @@ public final class EnvAllowListScrub {
return null;
}
String[] parts = lines.getFirst().substring("allowed ".length()).trim().split("\\s+");
if (parts.length != 3 || !"of".equals(parts[1])) {
if (parts.length != 5 || !"of".equals(parts[1]) || !"failed".equals(parts[3])) {
return null;
}
List<String> blanked = new ArrayList<>();
List<String> unblankable = new ArrayList<>();
for (int i = 1; i < lines.size(); i++) {
if (!lines.get(i).isBlank()) {
blanked.add(lines.get(i));
String line = lines.get(i);
if (line.isBlank()) {
continue;
}
if (line.startsWith("!")) {
unblankable.add(line.substring(1));
} else {
blanked.add(line);
}
}
return new ScrubReport(Integer.parseInt(parts[0]), Integer.parseInt(parts[2]),
List.copyOf(blanked));
Integer.parseInt(parts[4]), List.copyOf(blanked), List.copyOf(unblankable));
} catch (IOException | NumberFormatException e) {
return null;
}
@@ -1643,11 +1643,19 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
log.warn("memberCredentials allow-list: pane {} left no scrub report in {} — the "
+ "environment scrub cannot be confirmed to have run. Either the pane ended "
+ "before its shell finished starting, or its shell never read our generated "
+ "startup files, in which case that member saw the full host environment.",
+ "startup files. Either way, we cannot tell from here whether the scrub ran, "
+ "so we do not know what that member's environment contained.",
paneId, dir);
} else {
log.info("memberCredentials allow-list: pane {} allowed {} of {} environment variables",
paneId, report.allowed(), report.total());
if (report.failed() > 0) {
log.warn("memberCredentials allow-list: pane {} could not blank {} environment "
+ "variable(s) — {} (likely a zsh read-only/special parameter) — those "
+ "names were left in the member's environment. Confirm none of them is a "
+ "credential.",
paneId, report.failed(), report.unblankable());
}
List<String> shaped = report.blanked().stream()
.filter(name -> CREDENTIAL_SHAPED_NAME.matcher(name).matches())
.toList();
@@ -56,10 +56,13 @@ public final class FleetMetrics {
m.describe(REPLIES, "counter",
"Worker replies by delivery path (rendezvous=resolved an open send, inbox=stranded and held).");
m.describe(PUSH_NUDGES, "counter",
"CB-307 push-loop nudges to the primary (delivered|exhausted).");
"CB-307 push-loop nudges to the primary (sent|exhausted). fleetd #365: \"sent\" means "
+ "the herdr paste-and-submit call succeeded, not that the pane read it — this "
+ "layer has no read-receipt concept.");
m.describe(HEARTBEAT_NUDGES, "counter",
"CB-551 idle-lead heartbeat nudges (delivered|failed|exhausted). Quiet-cap exhaustion "
+ "means the lead idled with nothing pending and was told to stand down.");
"CB-551 idle-lead heartbeat nudges (sent|failed|exhausted). Quiet-cap exhaustion "
+ "means the lead idled with nothing pending and was told to stand down. "
+ "fleetd #365: \"sent\" means the herdr call succeeded, not that the lead read it.");
m.describe(SPAWNS, "counter",
"Worker spawn attempts by peer kind and outcome (ready|timeout|guard_rejected).");
m.describe(HERDR_CALLS, "counter",
@@ -32,7 +32,11 @@ public interface LeadChannel {
/** Non-destructive FIFO snapshot of the messages held for this daemon's own coord-id. */
List<LeadMessage> peek();
/** Drop {@code msgId} from the held set and ack it on the broker. A no-op if it is not held. */
/**
* Drop {@code msgId} from the held set and ack it on the broker. A repeated ack that this
* connection already completed may be a no-op. Any other unknown msgId must throw rather than
* report an ack that did not reach the broker.
*/
void ack(String msgId);
/** This daemon's own lead coordination id — the mailbox it owns, and the {@code from} it sends as. */
@@ -5,6 +5,7 @@ import dev.ltms.fleet.herdr.AgentStatus;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ScheduledExecutorService;
@@ -43,6 +44,9 @@ public final class LeadCoordLoop {
private static final Logger log = LoggerFactory.getLogger(LeadCoordLoop.class);
/** A bounded window is enough: redelivery can only follow a recent failed ack or recovery. */
private static final int RECENT_DELIVERY_LIMIT = 1_024;
/** How an arriving peer message is rendered into the lead's pane — the sender's coord-id, then its text. */
static final String DELIVERY_FORMAT = "[lead %s] %s";
@@ -51,6 +55,8 @@ public final class LeadCoordLoop {
private final Supplier<Map<String, String>> leads;
private final ScheduledExecutorService scheduler;
private final long intervalMs;
/** msgIds already written to the pane, so recovery redelivery is acked without another pane write. */
private final LinkedHashMap<String, Boolean> delivered = new LinkedHashMap<>();
private volatile boolean running;
@@ -117,6 +123,13 @@ public final class LeadCoordLoop {
if (held.isEmpty()) {
return;
}
LeadMessage msg = held.getFirst();
if (wasDelivered(msg.msgId())) {
// This lives here, rather than in LeadMailbox, because only this loop knows a pane write
// happened. The mailbox only knows broker delivery tags and must still redeliver after a crash.
ackDelivered(msg);
return;
}
String lead = resolveLocalLead();
if (lead == null) {
// Left unacked on purpose: the broker keeps holding it until a lead pane exists.
@@ -136,7 +149,6 @@ public final class LeadCoordLoop {
lead, status, held.size());
return;
}
LeadMessage msg = held.getFirst();
try {
agents.send(lead, DELIVERY_FORMAT.formatted(msg.from(), msg.content()));
} catch (RuntimeException e) {
@@ -145,6 +157,13 @@ public final class LeadCoordLoop {
msg.msgId(), msg.from(), lead, e.toString());
return;
}
rememberDelivered(msg.msgId());
if (ackDelivered(msg)) {
log.debug("lead coordination: delivered message {} from {} to lead {}", msg.msgId(), msg.from(), lead);
}
}
private boolean ackDelivered(LeadMessage msg) {
try {
channel.ack(msg.msgId());
} catch (RuntimeException e) {
@@ -152,9 +171,24 @@ public final class LeadCoordLoop {
// deliberate direction of this trade.
log.warn("lead coordination: delivered message {} but could not ack it: {}",
msg.msgId(), e.toString());
return;
return false;
}
return true;
}
private boolean wasDelivered(String msgId) {
synchronized (delivered) {
return delivered.containsKey(msgId);
}
}
private void rememberDelivered(String msgId) {
synchronized (delivered) {
delivered.put(msgId, Boolean.TRUE);
if (delivered.size() > RECENT_DELIVERY_LIMIT) {
delivered.remove(delivered.keySet().iterator().next());
}
}
log.debug("lead coordination: delivered message {} from {} to lead {}", msg.msgId(), msg.from(), lead);
}
/**
@@ -96,7 +96,14 @@ public final class LeadHeartbeatLoop {
this.metrics = metrics;
}
/** Count one nudge outcome when a registry is wired; a no-op in unit tests. */
/**
* Count one nudge outcome when a registry is wired; a no-op in unit tests.
*
* <p>fleetd #365: the {@code "sent"} outcome (renamed from {@code "delivered"}) records only
* that {@link #injectNudge} — a one-way herdr {@code agent.prompt} paste-and-submit — returned
* without throwing, not that the lead's pane actually read or acted on the text. This layer has
* no read-receipt concept, so "sent" is the honest word for what this call can ever establish.
*/
private void countNudge(String outcome) {
if (metrics != null) {
metrics.inc(FleetMetrics.HEARTBEAT_NUDGES, "outcome", outcome);
@@ -245,7 +252,7 @@ public final class LeadHeartbeatLoop {
agents.send(leadTerminal, fleet.nudgeText());
log.debug("idle-heartbeat: nudge sent to lead {} (quiet nudges so far in this stretch: {})",
leadTerminal, quietCount);
countNudge("delivered");
countNudge("sent");
} catch (RuntimeException e) {
log.warn("idle-heartbeat: failed to nudge lead {}: {}", leadTerminal, e.toString());
countNudge("failed");
@@ -59,7 +59,8 @@ import java.util.concurrent.TimeoutException;
* <p><strong>Recovery.</strong> The connection is opened with automatic + topology recovery
* enabled, mirroring {@code AmqpReplyInbox}: on reconnect the broker hands out fresh delivery tags,
* so the held snapshot is cleared (dedup by {@code msgId} still prevents any double-queue on
* redelivery) and any publish still awaiting its confirm is failed rather than left to idle out
* redelivery). {@link LeadCoordLoop} separately deduplicates pane writes, since it alone knows
* which messages reached a lead. Any publish still awaiting its confirm is failed rather than left to idle out
* the confirm timeout against a sequence number that means nothing on the new channel.
*/
public final class LeadMailbox implements LeadChannel, AutoCloseable {
@@ -85,6 +86,10 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
private final Object channelLock = new Object();
/** msgId → held delivery, for this mailbox's own queue only (there is exactly one). */
private final LinkedHashMap<String, Held> held = new LinkedHashMap<>();
/** Successful broker acks on this connection, retained only to make a repeated caller ack quiet. */
private final LinkedHashMap<String, Boolean> recentlyAcked = new LinkedHashMap<>();
/** Bounds {@link #recentlyAcked}: it is only an idempotency aid, never delivery state. */
private static final int RECENT_ACK_LIMIT = 1_024;
/**
* A dedicated channel for {@link #publish}, kept separate from {@link #channel} (consume + ack)
@@ -162,16 +167,15 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
}
// On automatic recovery the broker redelivers unacked messages with FRESH delivery-tags; the
// tags we were holding are now stale. Drop the held snapshot so the re-attached consumer
// repopulates it with valid tags (dedup by msgId still prevents any double-queue). Any publish
// repopulates it with valid tags (dedup by msgId still prevents any double-queue). LeadCoordLoop
// remembers successful pane writes separately, so that redelivery cannot write a pane twice. Any publish
// confirm still in flight when the connection dropped is equally stale — fail it now rather
// than let it silently ride out CONFIRM_TIMEOUT_MS.
if (connection instanceof Recoverable recoverable) {
recoverable.addRecoveryListener(new RecoveryListener() {
@Override
public void handleRecovery(Recoverable recoverable) {
synchronized (held) {
held.clear();
}
clearHeldForRecovery();
failPendingPublishesOnRecovery();
log.info("AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery");
}
@@ -350,15 +354,26 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
return snapshot;
}
/** Remove the held message {@code msgId} and ack it on the broker. No-op if not held. */
/**
* Remove the held message {@code msgId} and ack it on the broker.
*
* <p>A repeated ack that this connection already completed is a no-op, tracked in the bounded
* {@link #recentlyAcked} set. Any other missing entry throws: recovery clears {@link #held} while
* the broker still owns the unacked delivery, and quiet success there would hide a required retry.
* The set is bounded because it only distinguishes a recent duplicate caller ack from an unknown
* delivery; it is not a substitute for broker state across a reconnect.
*/
@Override
public void ack(String msgId) {
Held h;
synchronized (held) {
h = held.remove(msgId);
if (h == null && recentlyAcked.containsKey(msgId)) {
return;
}
}
if (h == null) {
return; // never held (or already acked) — no-op
throw new IllegalStateException("cannot ack lead message " + msgId + ": it is not held");
}
try {
synchronized (channelLock) {
@@ -372,6 +387,19 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
}
throw new IllegalStateException("cannot ack lead message " + msgId, e);
}
synchronized (held) {
recentlyAcked.put(msgId, Boolean.TRUE);
if (recentlyAcked.size() > RECENT_ACK_LIMIT) {
recentlyAcked.remove(recentlyAcked.keySet().iterator().next());
}
}
}
/** Clear stale delivery tags after recovery; package-private so the recovery contract test drives this exact path. */
void clearHeldForRecovery() {
synchronized (held) {
held.clear();
}
}
private DeliverCallback deliverCallback() {
@@ -138,6 +138,53 @@ public final class MessageService {
public record AskResult(AskOutcome outcome, String answer) {
}
/**
* How a worker's {@code fleet_reply} ({@link #reply(String, String)}) actually landed
* (fleetd #365) — the two doors that expose it, {@code fleet_reply} and {@code POST
* /sessions/{id}/reply}, both used to report the single word "delivered" whichever of these
* happened, so a caller could not tell an active handoff from a reply merely held for later
* drain. Both are successes; they are not the same fact.
*/
public enum ReplyOutcome {
/** Resolved a {@code fleet_send}/{@code fleet_ask} that was actively waiting on this reply. */
RESOLVED_SEND("resolved_send", true,
"delivered — resolved the fleet_send that was waiting for it"),
/**
* No live waiter was open, but the reply completed a parked async ticket directly
* ({@link #askAnsweredAsyncTasks}) — a {@code fleet_poll} caller sees it immediately.
*/
RESOLVED_ASYNC_TICKET("resolved_async_ticket", true,
"delivered — resolved a pending async ticket (visible to fleet_poll)"),
/** Nothing was waiting; the reply was queued in the inbox for a later drain (CB-307). */
QUEUED("queued", false,
"queued — no send or ticket was waiting; held in the inbox for a later drain");
private final String wireName;
private final boolean delivered;
private final String description;
ReplyOutcome(String wireName, boolean delivered, String description) {
this.wireName = wireName;
this.delivered = delivered;
this.description = description;
}
/** Stable machine-readable name for a JSON/metrics label (REST's {@code outcome} field). */
public String wireName() {
return wireName;
}
/** Whether something was actively waiting and received this reply right now. */
public boolean delivered() {
return delivered;
}
/** Shared human-readable text — the one place both {@code fleet_reply} and REST word this. */
public String description() {
return description;
}
}
/** Lifecycle phase of an async delegation ticket. */
public enum Phase {
/** Delegated and in flight — queued for the worker or being worked. */
@@ -439,16 +486,15 @@ public final class MessageService {
* @throws IllegalArgumentException if {@code content} is {@code null} or blank — the caller must
* report this as a client error (REST: 400 {@code bad_request}) rather than resolve
* anything
* @return always {@code true} — the reply resolved a live send, completed a parked ticket, or
* was queued
* @return which of the three ways (fleetd #365) the reply actually landed — never {@code null}
*/
public boolean reply(String session, String content) {
public ReplyOutcome reply(String session, String content) {
if (content == null || content.isBlank()) {
throw new IllegalArgumentException("content is required");
}
if (rendezvous.resolve(session, content)) {
count(FleetMetrics.REPLIES, "path", "rendezvous");
return true; // a live send took it — unchanged fast path
return ReplyOutcome.RESOLVED_SEND; // a live send took it — unchanged fast path
}
// #137/fleetd #307: no live rendezvous waiter, but this may be the worker's real fleet_reply resuming
// a turn that either answer() (#137) or ask() (fleetd #307) already gave up waiting on:
@@ -480,7 +526,7 @@ public final class MessageService {
asyncTasksByTurn.remove(turnId, orphan);
}
count(FleetMetrics.REPLIES, "path", "async-recovered");
return true; // the ticket itself took it — no inbox stranding at all
return ReplyOutcome.RESOLVED_ASYNC_TICKET; // the ticket itself took it — no inbox stranding
}
} else if (candidates.size() > 1) {
List<String> tickets = candidates.stream().map(t -> t.ticket).toList();
@@ -498,7 +544,7 @@ public final class MessageService {
if (pushLoop != null) {
pushLoop.onReplyQueued(session);
}
return true; // held, not lost
return ReplyOutcome.QUEUED; // held, not lost — but not delivered either
}
/**
@@ -1271,6 +1317,27 @@ public final class MessageService {
return new TaskView(ticket, Phase.FAILED, null, null, detail, null);
}
/**
* Test seam only — carries no production behaviour, and nothing in this class calls it;
* {@link #pruneTerminalTickets} still reads {@link Task#completedNanos} directly.
*
* <p>Reports whether {@code ticket}'s completion hook (the {@code whenComplete} registered in
* {@link Task}'s constructor) has actually run yet. Exists because {@link #poll} can report
* {@link Phase#DONE} for a ticket before that hook fires: {@code CompletableFuture.complete()}
* publishes its result and only afterwards runs dependent actions such as {@code whenComplete}
* (fleetd #399), so a caller that observes the future done via {@link #poll} is not thereby
* guaranteed to also observe {@link Task#completedNanos} stamped. A test that must order both
* events — e.g. before advancing an injected clock past the TTL, to avoid stamping the
* *advanced* time and masking a real eviction bug — waits on this instead of on
* {@link Phase#DONE}.
*
* @return {@code false} for an unknown ticket or one whose completion hook has not run yet
*/
boolean isCompletionStampedForTest(String ticket) {
Task task = tasks.get(ticket);
return task != null && task.completedNanos != null;
}
/** Best-effort live worker status for a pending poll; never throws (a lookup error is just noise). */
private String liveStatus(String target) {
try {
@@ -132,7 +132,15 @@ public final class ReplyPushLoop {
this.metrics = metrics;
}
/** Count one nudge outcome when a registry is wired; a no-op in unit tests. */
/**
* Count one nudge outcome when a registry is wired; a no-op in unit tests.
*
* <p>fleetd #365: the {@code "sent"} outcome (renamed from {@code "delivered"}) records only
* that {@code agents.send} — a one-way herdr {@code agent.prompt} paste-and-submit — returned
* without throwing. Nothing in this loop, or anywhere downstream of it, confirms the pane
* actually read or acted on the text; there is no read-receipt concept at this layer. "Sent"
* says exactly that; "delivered" claimed more than this call can ever establish.
*/
private void countNudge(String outcome) {
if (metrics != null) {
metrics.inc(FleetMetrics.PUSH_NUDGES, "outcome", outcome);
@@ -222,6 +230,26 @@ public final class ReplyPushLoop {
.collect(Collectors.toUnmodifiableSet());
}
/**
* Test seam only (fleetd #418): carries no production behaviour, and nothing in this class
* calls it. Exposes {@link #pendingQuestionTurnIdsFor} — the exact state {@link #decide} reads
* to decide whether a question keeps a lead's schedule alive.
*
* <p>{@code MessageService.ask()} does three things in order before a question is fully open to
* this loop: it flips the ticket's {@code poll()} phase to {@code Phase.ASKING}, then resolves
* the reverse-rendezvous waiter, then calls {@link #onQuestionOpened}, which is what actually
* populates {@link #pendingQuestions}. A test that barriers on {@code Phase.ASKING} observes only
* the first of those three steps — under load the asker thread can be descheduled between steps
* one and three, so the barrier releases before this method's underlying map is populated, and
* {@link #decide} correctly reports nothing pending yet. A test that must order itself after the
* state {@link #decide} actually reads waits on this instead of on the phase.
*
* @return an unmodifiable snapshot; empty for a lead with no open questions
*/
Set<String> pendingQuestionTurnIdsForTest(String lead) {
return pendingQuestionTurnIdsFor(lead);
}
private List<PendingIncident> pendingIncidentsFor(String lead) {
return pendingIncidents.values().stream().filter(i -> lead.equals(i.key().lead())).toList();
}
@@ -730,7 +758,7 @@ public final class ReplyPushLoop {
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
questionReminderCount + 1, maxReminders,
replyTargets.size(), tickets.size(), questions.size());
countNudge("delivered");
countNudge("sent");
for (PendingIncident incident : incidents) {
if (pendingIncidents.remove(incident.key(), incident)) {
deliveredIncidents.add(incident.key());
@@ -29,4 +29,9 @@ public record SpawnRequest(String profileName, String requestedCwd, String calle
String sessionName, String resumeSessionId) {
this(profileName, requestedCwd, callerCwd, sessionName, resumeSessionId, null);
}
/** Return a copy of this request with {@code profileName} replaced by {@code profile}. */
public SpawnRequest withProfile(String profile) {
return new SpawnRequest(profile, requestedCwd, callerCwd, sessionName, resumeSessionId, role);
}
}
@@ -696,6 +696,10 @@ public final class FleetApp {
/**
* The worker's structured reply ({@code fleet_reply}) — resolves the blocking send awaiting
* on this session, or queues the reply in the inbox when no send is open (CB-307).
*
* <p>fleetd #365: the response body's {@code delivered} field used to be unconditionally
* {@code true} for either case; it now reports whether a send/ticket was actually resolved,
* with {@code outcome} naming which (see {@link MessageService.ReplyOutcome}).
*/
private void replyMessage(Context ctx) {
String id = ctx.pathParam("id");
@@ -719,13 +723,19 @@ public final class FleetApp {
// a WRONG value instead of failing loudly. The check lives in MessageService.reply so both
// this door and FleetMcp.reply inherit the same rule; this catch only translates it into the
// {error, detail} envelope this file uses everywhere else.
MessageService.ReplyOutcome outcome;
try {
messages.reply(id, content);
outcome = messages.reply(id, content);
} catch (IllegalArgumentException e) {
ctx.status(400).json(Map.of("error", "bad_request", "detail", e.getMessage()));
return;
}
ctx.status(200).json(Map.of("sessionId", id, "delivered", true));
// fleetd #365: "delivered": true used to be unconditional here, whether the reply resolved
// a waiting send or was merely queued in the inbox for a later drain — the same gap
// FleetMcp.reply had over MCP. `delivered` now reflects which actually happened, and
// `outcome` names the specific case (see MessageService.ReplyOutcome).
ctx.status(200).json(Map.of("sessionId", id, "delivered", outcome.delivered(),
"outcome", outcome.wireName()));
}
/**
@@ -0,0 +1,150 @@
package dev.ltms.fleet;
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 dev.ltms.fleet.config.FleetConfig;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.slf4j.LoggerFactory;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #395: {@code exhaustedPattern} (see {@link FleetConfig.Profile#exhaustedPattern}) is
* deliberately opt-in — an unset one leaves usage-limit detection silently OFF for that profile,
* and nothing quarantines its credential. {@link Fleetd#reportExhaustedPatternGap} must say so at
* startup, naming every unarmed profile, and must never fire when every profile is armed. Mirrors
* {@link MemberCredentialsGapReportTest}'s pattern, capturing the real log via a
* {@link ListAppender}.
*
* <p>The 8-profile shape in {@link #theLiveEightProfileShapeWarnsExactlyTheSixUnarmedProfiles} is
* the live {@code fleetd.yaml} shape measured 2026-09-10 (fleetd #395's own ticket): 6 of 8
* profiles unarmed, 2 of those 6 ({@code opus}, {@code sonnet}) running on the operator's Claude
* subscription. {@code fleetd.yaml} itself is gitignored and unavailable to this test, so the
* shape is reproduced as a throwaway config in a {@code @TempDir} rather than read off disk.
*/
class ExhaustedPatternGapReportTest {
private static FleetConfig load(Path dir, String yaml) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml);
return FleetConfig.load(f);
}
private static ListAppender<ILoggingEvent> attach() {
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
return appender;
}
private static void detach(ListAppender<ILoggingEvent> appender) {
((Logger) LoggerFactory.getLogger(Fleetd.class)).detachAppender(appender);
}
/** Every profile name mentioned by a WARN-level log line, across every WARN this call produced. */
private static List<String> warnMessages(ListAppender<ILoggingEvent> appender) {
return appender.list.stream()
.filter(e -> e.getLevel() == Level.WARN)
.map(ILoggingEvent::getFormattedMessage)
.toList();
}
@Test
void theLiveEightProfileShapeWarnsExactlyTheSixUnarmedProfiles(@TempDir Path dir) throws Exception {
// Reproduces the live shape measured 2026-09-10: 8 profiles, 2 armed (sol, terra), 6
// unarmed (local, local-direct, gx, opus, sonnet, xf) — 2 of the unarmed 6 (opus, sonnet)
// are subscription: true.
FleetConfig cfg = load(dir, """
profiles:
local:
baseUrl: http://gx00.gw:8000
local-direct:
baseUrl: http://gx01.gw:8000
gx:
kind: opencode
baseUrl: https://llm.ltms.dev/v1
opus:
subscription: true
model: claude-opus-5
sonnet:
subscription: true
model: claude-sonnet-5
sol:
baseUrl: https://llm.ltms.dev/v1
exhaustedPattern: "The usage limit has been reached"
terra:
baseUrl: https://llm.ltms.dev/v1
exhaustedPattern: "The usage limit has been reached"
xf:
baseUrl: https://llm.ltms.dev/v1
""");
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportExhaustedPatternGap(cfg);
} finally {
detach(appender);
}
List<String> warns = warnMessages(appender);
assertFalse(warns.isEmpty(), "6 of 8 profiles are unarmed — at least one WARN must fire");
// Exactly one WARN aggregates every unarmed profile, naming all 6 and none of the 2 armed.
String aggregate = warns.stream()
.filter(m -> m.contains("no exhaustedPattern configured"))
.findFirst()
.orElseThrow(() -> new AssertionError("expected an aggregate unarmed-profiles WARN: " + warns));
for (String unarmed : List.of("local", "local-direct", "gx", "opus", "sonnet", "xf")) {
assertTrue(aggregate.contains(unarmed), "aggregate WARN must name '" + unarmed + "': " + aggregate);
}
for (String armed : List.of("sol", "terra")) {
assertFalse(aggregate.contains(armed), "aggregate WARN must NOT name armed profile '" + armed + "': " + aggregate);
}
// A second, louder WARN calls out the subscription profiles specifically.
String subscriptionWarn = warns.stream()
.filter(m -> m.contains("metered Claude plan"))
.findFirst()
.orElseThrow(() -> new AssertionError("expected a subscription-specific WARN: " + warns));
assertTrue(subscriptionWarn.contains("opus"), subscriptionWarn);
assertTrue(subscriptionWarn.contains("sonnet"), subscriptionWarn);
assertFalse(subscriptionWarn.contains("local-direct"),
"the subscription WARN must not name a non-subscription profile: " + subscriptionWarn);
}
@Test
void everyProfileArmedProducesNoWarningAtAll(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, """
profiles:
sol:
baseUrl: https://llm.ltms.dev/v1
exhaustedPattern: "The usage limit has been reached"
terra:
baseUrl: https://llm.ltms.dev/v1
exhaustedPattern: "The usage limit has been reached"
opus:
subscription: true
model: claude-opus-5
exhaustedPattern: "5-hour limit reached"
""");
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportExhaustedPatternGap(cfg);
} finally {
detach(appender);
}
assertTrue(warnMessages(appender).isEmpty(),
"every profile is armed — a checker that warns anyway always fires: " + warnMessages(appender));
}
}
@@ -0,0 +1,94 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.mcp.FleetMcp;
import dev.ltms.fleet.placement.BackendQuarantine;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Map;
import java.util.regex.Pattern;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #404: {@code exhaustionDetectionArmed} must describe the startup pattern map, not the
* reloaded config snapshot.
*
* <p>This test needs a reload. At startup the two snapshots agree, so a test of only a newly
* started daemon would not detect a live {@code config.get()} lookup in the report field.
*/
class FleetdExhaustionDetectionArmedWiringTest {
private static final String NO_PATTERN = """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
terra:
baseUrl: http://gx00.gw:8000
model: terra
guard:
offSubscriptionHosts:
- gx00.gw
""";
private static final String WITH_PATTERN = """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
terra:
baseUrl: http://gx00.gw:8000
model: terra
exhaustedPattern: "usage limit"
guard:
offSubscriptionHosts:
- gx00.gw
""";
@Test
@DisplayName("reloading an exhaustedPattern does not arm the startup detection source")
void reloadedPatternDoesNotChangeTheArmedFieldUntilRestart(@TempDir Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, NO_PATTERN);
ConfigRef config = new ConfigRef(file, FleetConfig.load(file));
Files.writeString(file, WITH_PATTERN);
assertTrue(config.reload().applied());
assertTrue(config.get().profiles().get("terra").hasExhaustedPattern());
FleetMcp.QuarantineSource source = Fleetd.quarantineSource(config, BackendQuarantine.none(),
Map.of());
assertFalse(source.exhaustedPatternArmed().apply("terra"),
"exhaustionDetectionArmed must use the startup pattern map, not config.get()");
}
@Test
@DisplayName("a profile in the startup pattern map is reported as armed")
void aProfileInTheStartupMapIsArmed(@TempDir Path dir) throws Exception {
// fleetd #404, second direction. The test above only ever passes an EMPTY startup map, so
// it cannot tell a correct lookup from one that is permanently off. Measured: replacing the
// armed lambda with `profile -> false` left the whole suite green at 1475 tests. That
// mutation would make #395's visibility feature dead — an operator fixing a detection gap
// would be told the gap is still open after fixing it, forever. Both directions are needed:
// this test is the only thing that fails when the field stops reporting armed at all.
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, WITH_PATTERN);
ConfigRef config = new ConfigRef(file, FleetConfig.load(file));
FleetMcp.QuarantineSource source = Fleetd.quarantineSource(config, BackendQuarantine.none(),
Map.of("terra", Pattern.compile("usage limit")));
assertTrue(source.exhaustedPatternArmed().apply("terra"),
"a profile whose pattern was compiled at startup must report armed");
assertFalse(source.exhaustedPatternArmed().apply("sonnet"),
"a profile absent from the startup map must not report armed");
}
}
@@ -0,0 +1,132 @@
package dev.ltms.fleet;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Files;
import java.nio.file.Path;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* Proves the one thing no test proved before this ticket's follow-up: that {@code Fleetd.main}
* ITSELF — not a copy of its logic, not the validator called directly — still refuses to start on
* a bad config. Mutation testing found that deleting {@code cfg.validateAll();} (née six
* individual {@code cfg.validateXxx();} calls) from {@code Fleetd.main} left the full 1478-test
* suite green; every existing test called a validator directly and none exercised {@code
* Fleetd.main} as the caller. See {@code FleetConfigValidateAllTest} for why the fix collapses
* those six calls into one reflective {@link FleetConfig#validateAll()}, and {@code
* ConfigRefTest} for the equivalent proof on the {@link
* dev.ltms.fleet.config.ConfigRef#reload()} path.
*
* <p>This is deliberately the real {@code static void main(String[] args)} — package-private, so
* only a test in this package can call it, which is exactly what makes this proof strong: it is
* not a helper extracted for testability, it is the literal method {@code java -jar fleetd.jar}
* invokes. Every fixture below is otherwise-valid and fails exactly one validator, and — because
* {@link FleetConfig#validateAll()} runs immediately after {@code SubscriptionGuard.
* assertPrimaryClean}, before {@code main} opens the herdr socket, binds Javalin, or touches
* anything else with a real side effect — calling {@code Fleetd.main} with one of these configs is
* safe: it is guaranteed to throw before reaching any of that, precisely because the config is
* deliberately invalid.
*/
class FleetdStartupValidationTest {
private static void assertMainRefuses(Path dir, String fileName, String yaml,
String mustContain) throws Exception {
Path f = dir.resolve(fileName);
Files.writeString(f, yaml);
IllegalStateException e = assertThrows(IllegalStateException.class,
() -> Fleetd.main(new String[]{f.toString()}),
fileName + ": Fleetd.main must refuse this config before doing anything else");
assertTrue(e.getMessage().contains(mustContain),
fileName + ": expected message to contain \"" + mustContain + "\" but was: "
+ e.getMessage());
}
@Test
void mainRefusesANonLoopbackBindWithoutTokenMode(@TempDir Path dir) throws Exception {
assertMainRefuses(dir, "auth-exposure.yaml", """
bind:
host: 0.0.0.0
port: 8765
""", "auth.mode: token");
}
@Test
void mainRefusesALeadTabPrefixCollision(@TempDir Path dir) throws Exception {
assertMainRefuses(dir, "lead-tab-prefixes.yaml", """
bind:
host: 127.0.0.1
port: 8765
fleet:
tabLabel: "lead: {role} {profile}"
leaders:
opus:
tab: "lead: opus"
""", "fleet.tabLabel");
}
@Test
void mainRefusesASubscriptionProfileThatReseatsAnthropicBaseUrl(@TempDir Path dir)
throws Exception {
assertMainRefuses(dir, "subscription-profiles.yaml", """
bind:
host: 127.0.0.1
port: 8765
profiles:
sonnet:
subscription: true
argv: ["ccs", "sonnet"]
env:
ANTHROPIC_BASE_URL: http://anything-not-on-the-allowlist
""", "ANTHROPIC_BASE_URL");
}
@Test
void mainRefusesAnUnknownCharterKey(@TempDir Path dir) throws Exception {
assertMainRefuses(dir, "charters.yaml", """
bind:
host: 127.0.0.1
port: 8765
fleet:
charters:
architetc: text
""", "architetc");
}
@Test
void mainRefusesAnArchitectSlotNamingAnUnconfiguredProfile(@TempDir Path dir) throws Exception {
assertMainRefuses(dir, "members.yaml", """
bind:
host: 127.0.0.1
port: 8765
profiles:
gx10:
baseUrl: http://gx10.gw:8000
fleet:
architects:
lead-designer:
profile: sonnet
""", "lead-designer");
}
@Test
void mainRefusesAProfileNamingAModelOutsideTheAllowList(@TempDir Path dir) throws Exception {
assertMainRefuses(dir, "models.yaml", """
bind:
host: 127.0.0.1
port: 8765
profiles:
sonnet:
baseUrl: http://gx10.gw:8000
model: claude-sonnet-5
rogue:
baseUrl: http://gx11.gw:8000
model: claude-opus-9000
models:
allow:
- model: claude-sonnet-5
""", "rogue");
}
}
@@ -0,0 +1,241 @@
package dev.ltms.fleet;
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 dev.ltms.fleet.config.FleetConfig;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.slf4j.LoggerFactory;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import java.util.Map;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #377: the startup git host line must report the <em>shape</em> of the value a
* member receives as GITEA_HOST — set or unset, length, scheme, trailing slash — and the
* value itself must never reach the log.
*/
class GitHostShapeReportTest {
/** A profile that opts in to the git-forge token, so GITEA_HOST is the var that matters. */
private static final String GIT_TOKEN_CONFIG = """
profiles:
local:
baseUrl: http://gx00.gw:8000
gitTokenEnv: WORKER_GITEA_TOKEN
""";
private static FleetConfig load(Path dir, String yaml) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml);
return FleetConfig.load(f);
}
/**
* The level this logger had before {@link #attach()} raised it, so {@link #detach} can put it
* back. {@code null} is a real value here — it means "inherit from the parent" — and that is
* exactly the state this logger starts in, so it must be restored as {@code null} rather than
* as some concrete level.
*/
private static Level originalLevel;
/**
* fleetd #377: the shape lines are logged at INFO, and {@code logback-test.xml} sets
* {@code dev.ltms.fleet} to WARN — so INFO events are dropped by the level check BEFORE any
* appender sees them. Attaching an appender is therefore not enough: without raising the level
* the list stays empty and every assertion below fails against correct production code. The
* sibling report tests do the same thing at each call site (see
* {@code MemberTrustModelReportTest}); doing it here keeps it in one place.
*/
private static ListAppender<ILoggingEvent> attach() {
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
originalLevel = logger.getLevel();
logger.setLevel(Level.INFO);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
return appender;
}
private static void detach(ListAppender<ILoggingEvent> appender) {
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
logger.detachAppender(appender);
logger.setLevel(originalLevel);
}
private static List<String> messages(ListAppender<ILoggingEvent> appender) {
return appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
}
@Test
void setLineReportsLengthSchemeAndTrailingSlash(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, GIT_TOKEN_CONFIG);
String value = "https://git.example.test/";
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportGitHostShape(cfg, Map.of("GITEA_HOST", value));
} finally {
detach(appender);
}
String expected = "startup git host GITEA_HOST: set (profile 'local' gitHostEnv) — "
+ "length=" + value.length() + ", startsWithScheme=true, trailingSlash=true";
assertTrue(messages(appender).contains(expected),
"a value with a scheme and a trailing slash must be reported by shape only: "
+ "its length, startsWithScheme=true, trailingSlash=true");
}
@Test
void bareHostLineReportsNoSchemeNoTrailingSlash(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, GIT_TOKEN_CONFIG);
String value = "git.example.test";
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportGitHostShape(cfg, Map.of("GITEA_HOST", value));
} finally {
detach(appender);
}
String expected = "startup git host GITEA_HOST: set (profile 'local' gitHostEnv) — "
+ "length=" + value.length() + ", startsWithScheme=false, trailingSlash=false";
assertTrue(messages(appender).contains(expected),
"a bare host name must report startsWithScheme=false, trailingSlash=false");
}
@Test
void unsetVariableIsLoggedAtInfoAndDoesNotThrow(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, GIT_TOKEN_CONFIG);
ListAppender<ILoggingEvent> appender = attach();
try {
// GITEA_HOST absent from the map entirely — the unset case must not throw.
Fleetd.reportGitHostShape(cfg, Map.of());
} finally {
detach(appender);
}
assertTrue(appender.list.stream().anyMatch(e ->
e.getLevel() == Level.INFO
&& e.getFormattedMessage()
.startsWith("startup git host GITEA_HOST: unset (profile 'local' gitHostEnv)")),
"an unset git host is useful information, not an error — say it at INFO level");
}
/**
* The important test: the value must NEVER appear in the log output. This test fails if the
* line is ever changed to include the value, because it drives the real logging path with a
* value that carries a marker no shape field could contain.
*/
@Test
void theValueNeverAppearsInLogOutput(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, GIT_TOKEN_CONFIG);
String marker = "never-logged-host-shape-377";
String value = "https://" + marker + "/";
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportGitHostShape(cfg, Map.of("GITEA_HOST", value));
} finally {
detach(appender);
}
List<String> msgs = messages(appender);
assertFalse(msgs.stream().anyMatch(m -> m.contains(value)),
"the full GITEA_HOST value must never reach the log");
assertFalse(msgs.stream().anyMatch(m -> m.contains(marker)),
"no fragment of the value may reach the log — a line that embeds the value"
+ " would leak at least this marker");
assertTrue(msgs.stream().anyMatch(m -> m.contains("startsWithScheme=true")
&& m.contains("trailingSlash=true")),
"the shape must still be reported although the value is not");
}
@Test
void unsetReportAlsoCarriesNoValue(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, GIT_TOKEN_CONFIG);
// A blank value is not injected either (putIfPresent skips it), so it must report unset
// without echoing anything of it.
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportGitHostShape(cfg, Map.of("GITEA_HOST", " "));
} finally {
detach(appender);
}
assertTrue(messages(appender).stream().anyMatch(m ->
m.startsWith("startup git host GITEA_HOST: unset")),
"a blank value is skipped by the launcher, so the line reports unset");
}
@Test
void defaultGitHostEnvIsGITEA_HOST(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, GIT_TOKEN_CONFIG);
assertEquals(Map.of("GITEA_HOST", List.of("profile 'local' gitHostEnv")),
Fleetd.gitHostEnvVars(cfg));
}
@Test
void explicitGitHostEnvReportsUnderItsOwnName(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, """
profiles:
local:
baseUrl: http://gx00.gw:8000
gitTokenEnv: WORKER_GITEA_TOKEN
gitHostEnv: MY_FORGE_HOST
""");
String value = "https://forge.example.test/";
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportGitHostShape(cfg, Map.of("MY_FORGE_HOST", value));
} finally {
detach(appender);
}
assertTrue(messages(appender).contains(
"startup git host MY_FORGE_HOST: set (profile 'local' gitHostEnv) — "
+ "length=" + value.length()
+ ", startsWithScheme=true, trailingSlash=true"),
"an explicitly named gitHostEnv is reported under that name");
}
@Test
void noGitTokenMeansNoHostLineIsNeeded(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, """
profiles:
local:
baseUrl: http://gx00.gw:8000
""");
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportGitHostShape(cfg, Map.of());
} finally {
detach(appender);
}
assertEquals(List.of("startup git host: no profile sets a gitTokenEnv — nothing to check"),
messages(appender));
}
@Test
void aHostWithAPortIsNotTreatedAsAScheme() {
assertFalse(Fleetd.startsWithScheme("git.example.test"));
assertFalse(Fleetd.startsWithScheme("git.example.test:3000"));
assertFalse(Fleetd.startsWithScheme("git.example.test:3000/"));
assertTrue(Fleetd.startsWithScheme("https://git.example.test"));
assertTrue(Fleetd.startsWithScheme("https://git.example.test:3000/"));
assertTrue(Fleetd.startsWithScheme("ssh://git.example.test"));
}
}
@@ -109,6 +109,8 @@ class ConfigRefTopLevelReportingCoverageTest {
v.put("memberLoginShell", null);
v.put("memberSkills", "/skills/a");
v.put("idleSleepGuard", new FleetConfig.IdleSleepGuard(true));
v.put("models", new FleetConfig.Models(
List.of(new FleetConfig.Models.ModelEntry("model-a"))));
assertNamesMatchComponents(v);
return v;
}
@@ -151,6 +153,8 @@ class ConfigRefTopLevelReportingCoverageTest {
v.put("memberLoginShell", null);
v.put("memberSkills", "/skills/b");
v.put("idleSleepGuard", new FleetConfig.IdleSleepGuard(false));
v.put("models", new FleetConfig.Models(
List.of(new FleetConfig.Models.ModelEntry("model-b"))));
assertNamesMatchComponents(v);
return v;
}
@@ -0,0 +1,172 @@
package dev.ltms.fleet.config;
import org.junit.jupiter.api.Test;
import java.lang.reflect.Constructor;
import java.lang.reflect.RecordComponent;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.TreeSet;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* Fleetd #358, the same "defect factory" #357 guarded on {@code FleetConfig.withDefaults()}
* (see {@code FleetConfigWithDefaultsPreservesEveryComponentTest}), reproduced here on
* {@link FleetConfig.Profile}. {@code Profile} carries a long back-compat constructor ladder — 8
* constructors, re-counted directly against the source rather than trusted from the ticket, at
* arities 25, 24, 22, 20, 18, 15, 14 and 12, against a canonical arity of 26 — and exactly ONE
* rebuild site, {@link FleetConfig.Profile#withProfile(String)}, whose own
* {@code return new Profile(...)} call is written at a literal 26-arg count. Add a 27th component
* and its established back-compat constructor at the old (26-arg) arity, and {@code withProfile}'s
* own call becomes a legal match for that new overload — silently dropping the new component every
* time a profile's name is defaulted from its {@code workers:} key.
*
* <p>Builds one {@link FleetConfig.Profile} through the TRUE canonical constructor — resolved by
* the record's own component types via {@code getDeclaredConstructor}, never by argument count —
* with a real, distinctive, non-null value in every component, calls {@link
* FleetConfig.Profile#withProfile(String)}, and asserts every component except {@code profile}
* itself survives unchanged, while {@code profile} comes back as the new name it was given.
*
* <p>Every value here is chosen so {@code Profile}'s own compact constructor (which normalizes
* several components — defaults {@code argv}/{@code kind}/{@code placement}/{@code workspace}/
* {@code gitHostEnv}, nulls a handful of blank-checked strings, clamps {@code weight}, coerces
* {@code subscription}) leaves it unchanged: every String is non-blank and already in the shape the
* compact constructor would otherwise coerce it to (e.g. {@code placement} is already lowercase),
* and every collection is non-empty. That is what makes "must survive unchanged" a valid assertion
* for every component below, the same reasoning {@code FleetConfigWithDefaultsPreservesEveryComponentTest}
* documents for {@code withDefaults()}.
*
* <p>{@link #EXCLUDED_FROM_SURVIVAL_CHECK} is kept deliberately empty and size-pinned by
* {@link #exclusionListSizeIsPinned()} — a checker whose escape hatch can grow to silence a failure
* is not a checker. Every one of {@code Profile}'s 26 current components has a real, non-null,
* non-blank value here and none is excluded.
*/
class FleetConfigProfileWithProfilePreservesEveryComponentTest {
private static final RecordComponent[] COMPONENTS = FleetConfig.Profile.class.getRecordComponents();
/** Deliberately empty today; grow it only with a matching justification, and re-pin the size. */
private static final Set<String> EXCLUDED_FROM_SURVIVAL_CHECK = Set.of();
/** One real, distinctive, non-null value per component, chosen to survive the compact ctor. */
private static Map<String, Object> baseValues() {
Map<String, Object> v = new LinkedHashMap<>();
v.put("profile", "profile-guard");
v.put("baseUrl", "https://guard.example/base");
v.put("model", "model-guard");
v.put("configDir", "/config/guard");
v.put("tokenEnv", "GUARD_TOKEN");
v.put("argv", List.of("guard-cmd"));
v.put("placement", "guard-placement");
v.put("workspace", "workspace-guard");
v.put("tabLabel", "tab-guard");
v.put("mcpUrl", "https://mcp.guard/");
v.put("cwd", "/cwd/guard");
v.put("parityOverlay", List.of(".guardrc"));
v.put("gitTokenEnv", "GUARD_GIT_TOKEN");
v.put("gitHostEnv", "GUARD_GIT_HOST");
v.put("kind", "claude-code");
v.put("env", Map.of("GUARD_ENV", "1"));
v.put("weight", 2.5f);
v.put("maxLoad", 4);
v.put("subscription", Boolean.TRUE);
v.put("exhaustedPattern", "pattern-guard");
v.put("credentialId", "cred-guard");
v.put("ideMcpUrl", "https://ide.guard/");
v.put("ideProjectDir", "ide-project-guard");
v.put("ideOpenCommand", "open-guard {dir}");
v.put("autoCompactWindow", 150_000);
v.put("errorPattern", "error-pattern-guard");
assertNamesMatchComponents(v);
return v;
}
/**
* Guards {@link #baseValues()} itself against drifting from the record's real shape — forgetting
* to add a new component here fails this assertion by name, rather than silently checking one
* component fewer than the record has.
*/
private static void assertNamesMatchComponents(Map<String, Object> values) {
Set<String> names = new TreeSet<>();
for (RecordComponent rc : COMPONENTS) {
names.add(rc.getName());
}
assertEquals(names, new TreeSet<>(values.keySet()),
"this test's value map has drifted from FleetConfig.Profile's actual components — "
+ "update baseValues() alongside the record");
}
/**
* Builds a {@link FleetConfig.Profile} through the TRUE canonical constructor — resolved by the
* record's own component types, not by argument count — so this never accidentally exercises a
* back-compat overload the way a literal {@code new Profile(...)} call risks doing.
*/
private static FleetConfig.Profile profileOf(Map<String, Object> values) throws ReflectiveOperationException {
Class<?>[] types = Arrays.stream(COMPONENTS).map(RecordComponent::getType).toArray(Class<?>[]::new);
Object[] args = Arrays.stream(COMPONENTS).map(rc -> values.get(rc.getName())).toArray();
Constructor<FleetConfig.Profile> ctor = FleetConfig.Profile.class.getDeclaredConstructor(types);
return ctor.newInstance(args);
}
@Test
void exclusionListSizeIsPinned() {
assertEquals(0, EXCLUDED_FROM_SURVIVAL_CHECK.size(),
"EXCLUDED_FROM_SURVIVAL_CHECK grew from 0 — every entry needs a justification in "
+ "this test class's javadoc AND this assertion re-pinned to the new size; a "
+ "growing exclusion list that silences failures on its own is not a guard");
}
/**
* The mutation this is built to catch: make {@code withProfile(String)}'s final constructor call
* literal at some arg count, add one more component to the record with a new back-compat
* constructor at the old arity, and the stale call silently rebinds. Every component here is real
* and non-null/non-blank, so none of it should be replaced by {@code withProfile}, except
* {@code profile} itself, which the method is documented to replace.
*/
@Test
void withProfilePreservesEveryOtherComponent() throws ReflectiveOperationException {
Map<String, Object> base = baseValues();
FleetConfig.Profile profile = profileOf(base);
FleetConfig.Profile renamed = profile.withProfile("renamed-profile-guard");
List<String> dropped = new ArrayList<>();
int checked = 0;
for (RecordComponent rc : COMPONENTS) {
String name = rc.getName();
if (EXCLUDED_FROM_SURVIVAL_CHECK.contains(name)) {
continue;
}
checked++;
Object expected = "profile".equals(name) ? "renamed-profile-guard" : base.get(name);
Object actual;
try {
actual = rc.getAccessor().invoke(renamed);
} catch (ReflectiveOperationException e) {
throw new RuntimeException("failed to read FleetConfig.Profile." + name + "()", e);
}
if (!Objects.equals(expected, actual)) {
dropped.add(String.format(Locale.ROOT,
"%s: withProfile() was expected to carry (%s) for '%s' but returned %s — a "
+ "component silently dropped by withProfile(), the shape of the "
+ "defect this test exists to catch (its final \"return new "
+ "Profile(...)\" call binding to a back-compat constructor instead "
+ "of the true canonical one)",
name, expected, name, actual));
}
}
System.out.printf(Locale.ROOT,
"FleetConfig.Profile.withProfile() component-survival coverage — %d components, %d "
+ "checked, %d excluded, %d survived%n",
COMPONENTS.length, checked, EXCLUDED_FROM_SURVIVAL_CHECK.size(), checked - dropped.size());
assertEquals(List.of(), dropped,
"withProfile() silently dropped these components: " + dropped);
}
}
@@ -2683,4 +2683,222 @@ class FleetConfigTest {
assertTrue(cfg.idleSleepGuard().isEnabled(),
"unlike ConfigReload/Health, this block defaults to ON even when present but empty");
}
// --- models: central allow-list of usable models -----------------------------------------
/**
* An absent {@code models:} block is today's behaviour exactly: no profile's {@code model:} is
* checked against anything, whatever it says. This is the "existing config keeps working"
* invariant — an operator on a gitignored {@code fleetd.yaml} that predates this feature must
* not be broken by upgrading the daemon.
*/
@Test
void absentModelsBlockValidatesNothing(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
gx10:
baseUrl: http://gx10.gw:8000
model: totally-unheard-of-model-xyz
""");
FleetConfig cfg = FleetConfig.load(f);
assertDoesNotThrow(cfg::validateModels);
}
/**
* {@code models:} present but with an empty (or absent) {@code allow:} must behave exactly like
* the block being absent — an operator adding the block for the first time with nothing in it
* yet must not be surprised by every profile suddenly refusing to start.
*/
@Test
void modelsBlockPresentButEmptyValidatesNothing(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
gx10:
baseUrl: http://gx10.gw:8000
model: whatever-the-operator-typed
models:
allow: []
""");
FleetConfig cfg = FleetConfig.load(f);
assertDoesNotThrow(cfg::validateModels);
}
/**
* The core of the ticket: a profile naming a model outside the configured allow-list fails
* config load, naming both the model and the profile that wanted it.
*/
@Test
void aProfileNamingAModelOutsideTheAllowListRefusesToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
sonnet:
baseUrl: http://gx10.gw:8000
model: claude-sonnet-5
rogue:
baseUrl: http://gx11.gw:8000
model: claude-opus-9000
models:
allow:
- model: claude-sonnet-5
- model: claude-haiku-5
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e = assertThrows(IllegalStateException.class, cfg::validateModels);
assertTrue(e.getMessage().contains("rogue"), "the refusal names the offending profile");
assertTrue(e.getMessage().contains("claude-opus-9000"), "the refusal names the offending model");
assertFalse(e.getMessage().contains("'sonnet'"),
"the profile whose model IS allowed must not be reported");
}
/** A profile whose {@code model:} is on the allow-list loads fine. */
@Test
void aProfileNamingAModelOnTheAllowListLoadsFine(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
sonnet:
baseUrl: http://gx10.gw:8000
model: claude-sonnet-5
models:
allow:
- model: claude-sonnet-5
""");
FleetConfig cfg = FleetConfig.load(f);
assertDoesNotThrow(cfg::validateModels);
}
/**
* A profile that never sets {@code model:} (a {@code subscription: true} profile relying on the
* account's own default is the live example) must not be refused just because an allow-list is
* active — there is nothing to check it against.
*/
@Test
void aProfileWithNoModelPassesEvenWithAnActiveAllowList(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
opus:
subscription: true
models:
allow:
- model: claude-sonnet-5
""");
FleetConfig cfg = FleetConfig.load(f);
assertDoesNotThrow(cfg::validateModels);
}
/**
* Reproduces the live config shape this ticket measured: a mix of {@code amazon-bedrock},
* {@code opencode} and {@code openai}-backed profiles, some naming a bare Claude id and some a
* provider-prefixed opencode id, in ONE allow-list. Both forms are just opaque strings compared
* for exact equality — proves the shape decision holds against the shape actually seen live,
* not just against a synthetic single-provider example.
*/
@Test
void aBareClaudeIdAndAnOpencodeProviderPrefixedIdBothFitOneAllowList(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
sonnet:
baseUrl: http://gx10.gw:8000
model: claude-sonnet-5
terra:
kind: opencode
model: openai/gpt-5.6-terra
nova:
kind: opencode
model: amazon-bedrock/amazon.nova-pro-v1:0
models:
allow:
- model: claude-sonnet-5
- model: openai/gpt-5.6-terra
- model: amazon-bedrock/amazon.nova-pro-v1:0
""");
FleetConfig cfg = FleetConfig.load(f);
assertDoesNotThrow(cfg::validateModels);
}
/** The same live shape, but one opencode profile's model is missing from the allow-list. */
@Test
void anUnlistedOpencodeProviderPrefixedModelRefusesToStart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
sonnet:
baseUrl: http://gx10.gw:8000
model: claude-sonnet-5
terra:
kind: opencode
model: openai/gpt-5.6-terra-withdrawn
models:
allow:
- model: claude-sonnet-5
- model: openai/gpt-5.6-terra
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e = assertThrows(IllegalStateException.class, cfg::validateModels);
assertTrue(e.getMessage().contains("terra"), "the refusal names the offending profile");
assertTrue(e.getMessage().contains("openai/gpt-5.6-terra-withdrawn"),
"the refusal names the offending model, with its provider prefix intact");
}
/**
* Editing a {@code profiles:} entry alone must never be able to widen what is permitted — only
* editing {@code models.allow:} itself can. This is the invariant the ticket states explicitly;
* this test pins it by giving a profile a plausible-looking model that was never added to the
* allow-list and confirming it is still refused.
*/
@Test
void addingAProfileCannotWidenTheAllowListByItself(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
sonnet:
baseUrl: http://gx10.gw:8000
model: claude-sonnet-5
brand-new:
baseUrl: http://gx12.gw:8000
model: claude-sonnet-6-preview
models:
allow:
- model: claude-sonnet-5
""");
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e = assertThrows(IllegalStateException.class, cfg::validateModels);
assertTrue(e.getMessage().contains("brand-new"));
assertTrue(e.getMessage().contains("claude-sonnet-6-preview"));
}
/** {@code models} is a brand-new top-level key and must be recognized, not WARN-ed as unknown. */
@Test
void modelsIsAKnownTopLevelKey() {
assertTrue(FleetConfig.KNOWN_TOP_LEVEL_KEYS.contains("models"));
}
}
@@ -0,0 +1,358 @@
package dev.ltms.fleet.config;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.lang.reflect.Method;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.List;
import java.util.Set;
import java.util.TreeSet;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* The gap this class exists to close: mutation testing on the fleetd ticket "central allow-list
* of usable models" found that although {@link FleetConfig#validateModels()}'s own logic was well
* pinned, nothing proved either real caller ({@code Fleetd.main} and {@link ConfigRef#reload()})
* still invoked it — deleting the call site left the full suite green (1478/0/0/0). A follow-up
* measurement (same technique — remove one call site, run the suite, not read the code) found the
* SAME gap for all five of {@link FleetConfig}'s other validators at startup, and for four of the
* six inside {@link ConfigRef#reload()}. This is a class of gap, not one line's mistake: every one
* of those thirteen tests called the validator itself directly, never the real caller that was
* supposed to.
*
* <p>The fix replaces the six individual {@code cfg.validateXxx()} calls at each of the two real
* call sites with one {@link FleetConfig#validateAll()}, which reaches every validator by
* reflection rather than by a hand-maintained list of names. A hand-maintained list of six names
* would have exactly the defect it replaces: the seventh validator someone adds next month has no
* reason to be added to it, and nothing would say so. This class proves TWO separate claims, and
* keeps them separate on purpose:
*
* <ol>
* <li>{@link #theSweepMechanismIsGenericNotHardcodedToFleetConfigsSixNames()} and its neighbours
* prove the reflective sweep itself ({@link FleetConfig#invokeAllValidators}) is a general
* mechanism — it runs whatever public, no-arg, void {@code validateXxx()} methods a class
* happens to declare today, including a class with more of them than {@link FleetConfig}
* has right now. This is the proof that a future, real seventh validator on {@link
* FleetConfig} would be swept automatically, without needing to add a real (unwanted)
* seventh validator just to exercise the claim.</li>
* <li>{@link #validateAllReachesEveryOneOfTodaysSixValidators()} proves {@link
* FleetConfig#validateAll()} itself is wired to that same generic mechanism and genuinely
* reaches each of today's six real validators — reusing the exact minimal failing
* configurations {@code FleetConfigTest} already established for each one directly, so a
* single call to {@code validateAll()} is shown to reproduce every one of those six
* failures.</li>
* </ol>
*
* <p>Together with the direct-{@code Fleetd.main}-invocation tests in {@code
* FleetdStartupValidationTest} (which prove the real startup call site still calls {@code
* validateAll()}) and the {@code ConfigRefTest} reload tests (which prove the same for {@link
* ConfigRef#reload()}), removing {@code cfg.validateAll();} from either real call site now fails
* a test in this module.
*
* <p><b>What is NOT pinned, measured rather than assumed.</b> Reverting {@link
* FleetConfig#validateAll()} to a hardcoded list of today's six method calls leaves the whole
* suite green (measured at review: 1491 tests, 0 failures). Nothing ties {@code validateAll()} to
* the generic sweep — claim 1 proves {@link FleetConfig#invokeAllValidators} is generic, and claim
* 2 proves {@code validateAll()} reaches today's six, and a hardcoded list satisfies both. So the
* reflective sweep is a convenience, not the guarantee. The guarantee is {@link
* #fleetConfigDeclaresExactlyTheseSixValidatorsToday()}: it fails the moment a seventh validator
* is declared, which forces whoever adds it to look at this file.
*/
class FleetConfigValidateAllTest {
// ── Claim 1: the reflective sweep is a general mechanism, not six names in disguise ──────────
/**
* A throwaway fixture class, unrelated to {@link FleetConfig} in every way except shape: three
* public, no-arg, void methods named {@code validateXxx}. Proves the sweep works on ANY class
* with this shape, not on something special-cased to {@link FleetConfig}.
*/
static class ThreeValidators {
final List<String> ran = new ArrayList<>();
public void validateAlpha() {
ran.add("validateAlpha");
}
public void validateBeta() {
ran.add("validateBeta");
}
public void validateGamma() {
ran.add("validateGamma");
}
}
@Test
void theSweepMechanismIsGenericNotHardcodedToFleetConfigsSixNames() {
ThreeValidators target = new ThreeValidators();
FleetConfig.invokeAllValidators(target);
assertEquals(List.of("validateAlpha", "validateBeta", "validateGamma"), target.ran,
"every validateXxx() method on this unrelated class must run, in alphabetical "
+ "order — the sweep reads the class's own shape, not a name FleetConfig "
+ "happens to know about");
}
/**
* The core of the "self-maintaining" requirement: the exact same class shape as {@link
* ThreeValidators}, plus one more method — standing in for "a developer adds a validator next
* month". Nothing about the sweep changes to pick it up; the new method is invoked purely
* because it exists and matches the shape. This is what makes adding a seventh real validator
* to {@link FleetConfig} safe without touching {@link FleetConfig#validateAll()} or either
* call site — there is no "wire it in" step left to forget.
*/
static class FourValidators {
final List<String> ran = new ArrayList<>();
public void validateAlpha() {
ran.add("validateAlpha");
}
public void validateBeta() {
ran.add("validateBeta");
}
public void validateGamma() {
ran.add("validateGamma");
}
public void validateDelta() {
ran.add("validateDelta");
}
}
@Test
void addingAFourthValidatorMethodGetsSweptWithNoOtherChange() {
FourValidators target = new FourValidators();
FleetConfig.invokeAllValidators(target);
assertEquals(List.of("validateAlpha", "validateBeta", "validateDelta", "validateGamma"),
sorted(target.ran),
"the fourth method must be reached automatically — proving a class can grow the "
+ "set of things it validates with no change to the sweep itself");
}
private static List<String> sorted(List<String> in) {
List<String> copy = new ArrayList<>(in);
copy.sort(String::compareTo);
return copy;
}
@Test
void aFailingValidatorStopsTheSweepAndPropagatesUnchanged() {
class OneFails {
@SuppressWarnings("unused")
public void validateOk() {
// passes
}
public void validateBoom() {
throw new IllegalStateException("refusing to start: boom");
}
}
IllegalStateException e = assertThrows(IllegalStateException.class,
() -> FleetConfig.invokeAllValidators(new OneFails()));
assertEquals("refusing to start: boom", e.getMessage(),
"the real exception must propagate unchanged, not be wrapped or swallowed");
}
/**
* Every rule the sweep's filter applies, proven independently: only public, no-arg, void
* methods whose name starts with {@code "validate"} run, {@code validateAll} itself is
* excluded (so a class that declares one of its own — as {@link FleetConfig} does — cannot
* recurse into itself), and a same-shaped-but-wrongly-named or wrongly-shaped method never
* runs. A reader who "simplifies" the filter in {@link FleetConfig#invokeAllValidators} in a
* way that widens or narrows it breaks one of these.
*/
static class FilterEdgeCases {
final List<String> ran = new ArrayList<>();
public void validateReal() {
ran.add("validateReal");
}
/** Wrong name — must not run. */
public void checkSomething() {
ran.add("checkSomething");
}
/** Wrong shape — takes an argument. */
public void validateWithArg(String ignored) {
ran.add("validateWithArg");
}
/** Wrong shape — returns something. */
public boolean validateReturnsBoolean() {
ran.add("validateReturnsBoolean");
return true;
}
/** Excluded by name on purpose, so the sweep cannot call itself. */
public void validateAll() {
ran.add("validateAll");
}
}
@Test
void onlyPublicNoArgVoidValidateNamedMethodsRun() {
FilterEdgeCases target = new FilterEdgeCases();
FleetConfig.invokeAllValidators(target);
assertEquals(List.of("validateReal"), target.ran,
"checkSomething (wrong name), validateWithArg (wrong shape), "
+ "validateReturnsBoolean (wrong shape), and validateAll (excluded by "
+ "name) must all be skipped");
}
// ── Claim 2: FleetConfig.validateAll() is wired to that mechanism and reaches all six today ──
/**
* Reflectively enumerates {@link FleetConfig}'s own public, no-arg, void {@code validateXxx()}
* methods (excluding {@code validateAll} itself) — the exact same filter {@link
* FleetConfig#invokeAllValidators} applies. This is not the mechanism proof (that is claim 1,
* above, on an unrelated class) — it is a visible denominator: today there are six, named
* here, so a reader adding a seventh sees this assertion name the new count rather than a
* silent pass at the old one.
*/
@Test
void fleetConfigDeclaresExactlyTheseSixValidatorsToday() {
Set<String> names = new TreeSet<>();
for (Method m : FleetConfig.class.getMethods()) {
if (java.lang.reflect.Modifier.isPublic(m.getModifiers())
&& m.getParameterCount() == 0
&& m.getReturnType() == void.class
&& m.getName().startsWith("validate")
&& !m.getName().equals("validateAll")) {
names.add(m.getName());
}
}
assertEquals(new TreeSet<>(Set.of("validateAuthExposure", "validateLeadTabPrefixes",
"validateSubscriptionProfiles", "validateCharters", "validateMembers",
"validateModels")), names,
"FleetConfig's public validate*() methods changed. Do TWO things, in this "
+ "order. First confirm validateAll() still delegates to "
+ "invokeAllValidators(this) — a hardcoded list there passes every other "
+ "test in this class, so this assertion is the only place that will ever "
+ "make you check. Only then update the expected set to match.");
}
/** A minimal, otherwise-valid file — same shape FleetConfigTest and ConfigRefTest use. */
private static String minimalValidYaml() {
return """
bind:
host: 127.0.0.1
port: 8765
profiles:
sonnet:
baseUrl: http://gx00.gw:8000
model: sonnet
""";
}
@Test
void aFullyValidConfigPassesValidateAll(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, minimalValidYaml());
assertDoesNotThrow(() -> FleetConfig.load(f).validateAll());
}
/**
* The heart of claim 2: for each of today's six real validators, a minimal file that fails
* ONLY that one — the exact fixtures {@code FleetConfigTest} uses to test each validator
* directly — must also fail through {@link FleetConfig#validateAll()}. If a future edit to
* {@code validateAll()} silently dropped one validator from the sweep (e.g. a typo'd name
* filter), exactly one of these six would start passing when it must not.
*/
@Test
void validateAllReachesEveryOneOfTodaysSixValidators(@TempDir Path dir) throws Exception {
// validateAuthExposure: a non-loopback bind without token mode.
assertValidateAllRefuses(dir, "auth-exposure.yaml", """
bind:
host: 0.0.0.0
port: 8765
""", "auth.mode: token");
// validateLeadTabPrefixes: a fleet-wide tabLabel that starts with a lead's own tabPrefix.
assertValidateAllRefuses(dir, "lead-tab-prefixes.yaml", """
bind:
host: 127.0.0.1
port: 8765
fleet:
tabLabel: "lead: {role} {profile}"
leaders:
opus:
tab: "lead: opus"
""", "fleet.tabLabel");
// validateSubscriptionProfiles: subscription: true with env: reseating ANTHROPIC_BASE_URL.
assertValidateAllRefuses(dir, "subscription-profiles.yaml", """
bind:
host: 127.0.0.1
port: 8765
profiles:
sonnet:
subscription: true
argv: ["ccs", "sonnet"]
env:
ANTHROPIC_BASE_URL: http://anything-not-on-the-allowlist
""", "ANTHROPIC_BASE_URL");
// validateCharters: a charter key that is not a role wire name.
assertValidateAllRefuses(dir, "charters.yaml", """
bind:
host: 127.0.0.1
port: 8765
fleet:
charters:
architetc: text
""", "architetc");
// validateMembers: an architect slot naming an unconfigured profile.
assertValidateAllRefuses(dir, "members.yaml", """
bind:
host: 127.0.0.1
port: 8765
profiles:
gx10:
baseUrl: http://gx10.gw:8000
fleet:
architects:
lead-designer:
profile: sonnet
""", "lead-designer");
// validateModels: a profile naming a model outside the configured allow-list.
assertValidateAllRefuses(dir, "models.yaml", """
bind:
host: 127.0.0.1
port: 8765
profiles:
sonnet:
baseUrl: http://gx10.gw:8000
model: claude-sonnet-5
rogue:
baseUrl: http://gx11.gw:8000
model: claude-opus-9000
models:
allow:
- model: claude-sonnet-5
""", "rogue");
}
private static void assertValidateAllRefuses(Path dir, String fileName, String yaml,
String mustContain) throws Exception {
Path f = dir.resolve(fileName);
Files.writeString(f, yaml);
FleetConfig cfg = FleetConfig.load(f);
IllegalStateException e = assertThrows(IllegalStateException.class, cfg::validateAll,
fileName + ": validateAll() must refuse this config");
assertTrue(e.getMessage().contains(mustContain),
fileName + ": expected message to contain \"" + mustContain + "\" but was: "
+ e.getMessage());
}
}
@@ -97,6 +97,8 @@ class FleetConfigWithDefaultsPreservesEveryComponentTest {
v.put("memberLoginShell", "/bin/zsh");
v.put("memberSkills", "/skills/guard");
v.put("idleSleepGuard", new FleetConfig.IdleSleepGuard(true));
v.put("models", new FleetConfig.Models(
List.of(new FleetConfig.Models.ModelEntry("model-guard"))));
assertNamesMatchComponents(v);
return v;
}
@@ -30,6 +30,7 @@ import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.ScheduledThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.BiConsumer;
import java.util.function.LongSupplier;
@@ -206,6 +207,139 @@ class FleetHealthMonitorTest {
.workingSuspectAfterOrDefault());
}
// --- fleetd #386: a stall detector whose only clock freezes with a sleeping host is worse
// than a silent one — it reports "quiet" for a member that was genuinely busy for hours.
@Test void monotonicClockFrozenPastThresholdOnRealClockStillReportsStallSuspected() {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
AtomicLong mono = new AtomicLong(0);
AtomicLong real = new AtomicLong(0);
FakeHerdr herdr = new FakeHerdr().withAgent("busy", "term_busy", "pane_busy", "tab_busy");
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = monitorWithClocks(herdr,
List.of(member("term_busy", MemberSession.State.BUSY, 0, 0)), scheduler,
mono::get, real::get, 60, 600, (_, _) -> { });
monitor.tick(); // establishes the clock baseline; nothing has diverged yet
assertEquals(0, appender.list.stream().filter(event -> event.getFormattedMessage()
.contains("state=STALL_SUSPECTED")).count());
// The host "sleeps": the monotonic clock stands completely still while the real clock
// keeps moving, past the 600s stall threshold.
real.set(TimeUnit.SECONDS.toNanos(700));
monitor.tick();
monitor.stop();
assertTrue(appender.list.stream().anyMatch(event -> event.getFormattedMessage()
.contains("member=term_busy state=STALL_SUSPECTED")),
"the real clock crossed the stall threshold even though the monotonic clock never moved");
} finally {
logger.detachAppender(appender);
}
}
@Test void clockDivergenceIsLoggedOnceNotOncePerTick() {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
AtomicLong mono = new AtomicLong(0);
AtomicLong real = new AtomicLong(0);
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = monitorWithClocks(new FakeHerdr(), List.of(), scheduler,
mono::get, real::get, 60, 600, (_, _) -> { });
monitor.tick(); // baseline: no divergence possible yet
// One sleep gap: the monotonic clock is frozen while the real clock jumps far past one
// tick interval (60s).
real.set(TimeUnit.SECONDS.toNanos(700));
monitor.tick();
// The host is awake again: both clocks advance together from here, so no more divergence.
mono.set(TimeUnit.SECONDS.toNanos(10));
real.set(TimeUnit.SECONDS.toNanos(710));
monitor.tick();
mono.set(TimeUnit.SECONDS.toNanos(20));
real.set(TimeUnit.SECONDS.toNanos(720));
monitor.tick();
monitor.stop();
assertEquals(1, appender.list.stream().filter(event -> event.getFormattedMessage()
.contains("the monotonic clock did not advance"))
.count(), "one sleep gap must produce exactly one divergence line, not one per tick");
} finally {
logger.detachAppender(appender);
}
}
/**
* fleetd #386 follow-up, added on merge. The fix carries a PER-MEMBER drift baseline, so drift
* from a sleep that happened BEFORE a member went busy is never charged to that member. The
* two tests shipped with the fix both start with the member already BUSY, so a single global
* baseline passes them — this one fails without the per-member map.
*
* <p>Order matters: the host sleeps while nothing is busy, and only then does a member take a
* turn. Its stall clock must start at zero.
*/
@Test void driftFromASleepBeforeAMemberWentBusyIsNotChargedToThatMember() {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
AtomicLong mono = new AtomicLong(0);
AtomicLong real = new AtomicLong(0);
AtomicReference<List<MemberSession>> roster = new AtomicReference<>(List.of());
FakeHerdr herdr = new FakeHerdr().withAgent("busy", "term_busy", "pane_busy", "tab_busy");
var scheduler = Executors.newSingleThreadScheduledExecutor();
AgentControl agents = new AgentControl(herdr);
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, roster::get,
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
scheduler, mono::get, real::get, 60, 600, (_, _) -> { });
monitor.tick(); // baseline, no members yet
// The host sleeps for 700s with nobody busy: the monotonic clock stands still.
real.set(TimeUnit.SECONDS.toNanos(700));
monitor.tick();
// Awake again. Only NOW does a member start a turn, with a fresh activity stamp taken
// from the monotonic clock. Both clocks advance together from here.
mono.set(TimeUnit.SECONDS.toNanos(10));
real.set(TimeUnit.SECONDS.toNanos(710));
roster.set(List.of(member("term_busy", MemberSession.State.BUSY,
0, TimeUnit.SECONDS.toNanos(10))));
monitor.tick();
mono.set(TimeUnit.SECONDS.toNanos(20));
real.set(TimeUnit.SECONDS.toNanos(720));
monitor.tick();
monitor.stop();
assertEquals(0, appender.list.stream().filter(event -> event.getFormattedMessage()
.contains("member=term_busy state=STALL_SUSPECTED")).count(),
"the member has been busy for 10s, not 710s — the earlier sleep is not its stall");
} finally {
logger.detachAppender(appender);
}
}
private static FleetHealthMonitor monitorWithClocks(FakeHerdr herdr, List<MemberSession> roster,
java.util.concurrent.ScheduledExecutorService scheduler, LongSupplier clock,
LongSupplier realtimeClock, long intervalSeconds, long workingSuspectAfterSeconds,
BiConsumer<String, String> failTarget) {
AgentControl agents = new AgentControl(herdr);
return new FleetHealthMonitor(agents, () -> roster,
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
scheduler, clock, realtimeClock, intervalSeconds, workingSuspectAfterSeconds, failTarget);
}
@Test void goneMemberRecoveryLogsOnceWithoutRefiringTargetFailure() {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
Level previousLevel = logger.getLevel();
@@ -357,6 +357,54 @@ class CompletionResolverTest {
"the failure carries whatever was on screen: " + waiter.getNow(null).text());
}
@Test
void aPlausibleLookingReplyInsideTheFloorStillFails() {
// fleetd#376 guard. A fix was attempted that inspected the pane inside the floor and resolved
// a COMPLETION when the text "looked like a real reply". Every cheap test for that is unsafe:
// lastAssistantBlock falls back to the WHOLE pane when there is no ⏺ marker, and a crash pane
// almost always contains sentence punctuation — in a file path, a version, or a hostname.
// This pane is the trap: it reads like a finished answer and it is a backend failure.
FakeHerdr herdr = new FakeHerdr().readText(
"Error: connection reset while loading src/main/java/Foo.java v1.2.3\n❯ ");
Rendezvous rendezvous = new Rendezvous();
long[] clock = {10_000_000_000L};
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
var waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, null, clock[0]);
clock[0] += CompletionResolver.MIN_TURN_NANOS - 1; // inside the floor
resolver.resolve("term_a", turn);
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
"inside the floor the verdict is always FAILED — never guess a completion from pane text");
}
@Test
void theTooFastFailureDoesNotAssertACauseItCannotKnow() {
// fleetd#376: the message used to say "most likely a backend error before any work started".
// When no error pattern matches, that cause is a guess, and a reader who believes it stops
// looking at the pane. The verdict stays FAILED; only the claim about WHY is withdrawn.
FakeHerdr herdr = new FakeHerdr().readText("I am running on opencode/mimo-v2.5-free.\n❯ ");
Rendezvous rendezvous = new Rendezvous();
long[] clock = {10_000_000_000L};
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
var waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, null, clock[0]);
clock[0] += CompletionResolver.MIN_TURN_NANOS - 1; // inside the floor
resolver.resolve("term_a", turn);
String text = waiter.getNow(null).text();
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(), "still fails, still loud");
assertFalse(text.contains("most likely a backend error"),
"an unmatched fast turn must not assert a backend error: " + text);
assertTrue(text.contains("mimo-v2.5-free"), "the pane is still carried: " + text);
}
@Test
void aBusyToDoneTransitionJustOutsideTheFloorResolvesNormally() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ a real, if quick, answer\n❯ ");
@@ -1117,8 +1165,16 @@ class CompletionResolverTest {
assertTrue(waiter.isDone());
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
"still a failure — the floor itself, not the pattern, is why");
assertTrue(waiter.getNow(null).text().contains("too fast to be real work"),
"a non-match inside the floor stays the generic too-fast reason: " + waiter.getNow(null).text());
// fleetd#376: this used to assert the phrase "too fast to be real work", which carried the
// claim "most likely a backend error before any work started". With no pattern matched that
// cause is a guess, so the wording was withdrawn. What this test really guards is unchanged:
// the floor alone still fails the turn, it stays generic, and it never notifies the sink.
String reason = waiter.getNow(null).text();
assertTrue(reason.contains("inside the floor"),
"a non-match inside the floor stays the generic floor reason: " + reason);
assertFalse(reason.contains("most likely a backend error"),
"a non-match must not assert a cause it did not establish: " + reason);
assertTrue(reason.contains("still starting up"), "the pane is still carried: " + reason);
assertTrue(notified.isEmpty(), "a non-match must never notify the typed sink");
}
}
@@ -3,6 +3,7 @@ package dev.ltms.fleet.mcp;
import dev.ltms.fleet.auth.CallerResolver;
import dev.ltms.fleet.auth.MemberRegistry;
import dev.ltms.fleet.auth.Principal;
import dev.ltms.fleet.auth.Role;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
@@ -28,6 +29,7 @@ import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementPolicies;
import io.modelcontextprotocol.spec.McpSchema;
import dev.ltms.fleet.msg.InMemoryReplyInbox;
import dev.ltms.fleet.msg.ReplyInbox;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
@@ -77,7 +79,7 @@ class FleetMcpTest {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting(target), "send should be accepted for " + target);
FleetMcp.reply(messages, target, "received");
FleetMcp.reply(messages, target, Role.WORKER, "received");
assertEquals("received", textOf(send.get(6, TimeUnit.SECONDS)));
}
@@ -108,8 +110,10 @@ class FleetMcpTest {
}
assertTrue(rendezvous.isWaiting("term_a"), "send should have opened its waiter");
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "LGTM");
assertEquals("delivered", textOf(reply));
// fleetd #365: a resolved live send must read distinctly from a merely-queued reply —
// see replyWithNoPendingSendIsQueuedNotError below for the other case.
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", Role.WORKER, "LGTM");
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND.description(), textOf(reply));
McpSchema.CallToolResult res = send.get(6, TimeUnit.SECONDS);
assertNotEquals(Boolean.TRUE, res.isError());
@@ -134,8 +138,8 @@ class FleetMcpTest {
}
assertTrue(rendezvous.isWaiting("term_a"), "send should have opened its waiter");
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "async LGTM");
assertEquals("delivered", textOf(reply));
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", Role.WORKER, "async LGTM");
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND.description(), textOf(reply));
// Poll until the async send completes and reports the reply.
McpSchema.CallToolResult polled = FleetMcp.poll(messages, ticket, null);
@@ -181,7 +185,7 @@ class FleetMcpTest {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting("term_a"));
FleetMcp.reply(messages, "term_a", "done");
FleetMcp.reply(messages, "term_a", Role.WORKER, "done");
assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS)));
McpSchema.CallToolResult done = FleetMcp.poll(messages, ticket, null);
@@ -213,7 +217,7 @@ class FleetMcpTest {
// fleet_poll{ticket} stuck PENDING forever and later force-failed with a false "session
// released before it replied" reason. This used to land in the inbox instead (see the old
// assertion this replaced: messages.drainReplies("term_a").getFirst()...) — that was the bug.
FleetMcp.reply(messages, "term_a", "finished after timeout");
FleetMcp.reply(messages, "term_a", Role.WORKER, "finished after timeout");
assertEquals("finished after timeout", textOf(FleetMcp.poll(messages, ticket, null)));
assertTrue(messages.drainReplies("term_a").isEmpty(),
"the reply completed its own ticket directly and never touched the inbox");
@@ -267,7 +271,7 @@ class FleetMcpTest {
}
assertEquals(MessageService.Phase.FAILED, second.phase());
FleetMcp.reply(messages, "term_a", "late reply");
FleetMcp.reply(messages, "term_a", Role.WORKER, "late reply");
assertEquals("late reply", messages.drainReplies("term_a").getFirst().content());
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
@@ -276,7 +280,7 @@ class FleetMcpTest {
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
FleetMcp.reply(messages, "term_a", "done");
FleetMcp.reply(messages, "term_a", Role.WORKER, "done");
assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS)));
}
@@ -328,9 +332,10 @@ class FleetMcpTest {
@Test
void replyWithNoPendingSendIsQueuedNotError() {
// CB-307: a reply with no open send is now queued in the inbox, not an error.
McpSchema.CallToolResult res = FleetMcp.reply(messages, "term_a", "orphan");
// fleetd #365: it must also no longer claim "delivered" — nothing was waiting for it.
McpSchema.CallToolResult res = FleetMcp.reply(messages, "term_a", Role.WORKER, "orphan");
assertNotEquals(Boolean.TRUE, res.isError(), "a queued reply is not an error");
assertEquals("delivered", textOf(res));
assertEquals(MessageService.ReplyOutcome.QUEUED.description(), textOf(res));
// The reply is drainable by target.
var drained = messages.drainReplies("term_a");
@@ -338,6 +343,40 @@ class FleetMcpTest {
assertEquals("orphan", drained.getFirst().content());
}
@Test
void replyFromLeadIsRefusedBeforeItCanPublishToTheWorkerInbox() {
ReplyInbox inboxThatRejectsPublishes = new ReplyInbox() {
@Override public void own(String target) { }
@Override public void release(String target) { }
@Override public void publish(String target, String msgId, String content) {
fail("a lead fleet_reply must not publish to the worker inbox");
}
@Override public List<InboxMessage> peek(String target) { return List.of(); }
@Override public void ack(String target, String msgId) { }
};
MessageService leadMessages = new MessageService(agents, new Injector(agents), new Rendezvous(),
inboxThatRejectsPublishes);
McpSchema.CallToolResult res = assertDoesNotThrow(
() -> FleetMcp.reply(leadMessages, "term_lead", Role.PRIMARY, "peer reply"));
assertTrue(res.isError());
assertEquals("fleet_reply has no route to a peer lead. Use fleet_send{coordId: ...} for a peer on another "
+ "daemon or fleet_send{sessionId: ...} for a peer on this host. fleet_reply resolves a member's "
+ "blocked fleet_send, and a peer's coord-id message is durable and non-blocking, so there is "
+ "nothing for it to resolve.",
textOf(res));
}
@Test
void replyFromUnidentifiedCallerKeepsItsOwnError() {
McpSchema.CallToolResult res = FleetMcp.reply(messages, null, Role.PRIMARY, "reply");
assertTrue(res.isError());
assertEquals("fleet_reply is for workers only — could not identify the calling worker from the connection",
textOf(res));
}
@Test
void replyWithBlankContentIsACleanToolErrorNotAnUncaughtException() {
// fleetd #302: MessageService.reply now REJECTS blank content by throwing. fleet_reply's
@@ -348,7 +387,7 @@ class FleetMcpTest {
// has always used isBlank for exactly this reason.
for (String blank : new String[] {null, "", " ", "\n\t"}) {
McpSchema.CallToolResult res = assertDoesNotThrow(
() -> FleetMcp.reply(messages, "term_a", blank),
() -> FleetMcp.reply(messages, "term_a", Role.WORKER, blank),
"blank content must be refused as a tool error, never thrown out of the handler");
assertEquals(Boolean.TRUE, res.isError(), "blank content is an error result");
assertTrue(textOf(res).contains("content is required"),
@@ -361,7 +400,7 @@ class FleetMcpTest {
@Test
void bridgePollWithTargetDrainsReplies() {
// A reply with no open send queues it in the inbox.
FleetMcp.reply(messages, "term_a", "queued-msg");
FleetMcp.reply(messages, "term_a", Role.WORKER, "queued-msg");
// fleet_poll with target drains the inbox.
McpSchema.CallToolResult res = FleetMcp.poll(messages, null, "term_a");
@@ -412,8 +451,8 @@ class FleetMcpTest {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting("term_a"), "the answer should have reopened a waiter");
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "done");
assertEquals("delivered", textOf(reply));
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", Role.WORKER, "done");
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND.description(), textOf(reply));
assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS)));
}
@@ -1251,13 +1290,13 @@ class FleetMcpTest {
@Test
void bridgeAckRemovesSpecificReply() {
// Queue a reply and capture its msgId.
FleetMcp.reply(messages, "term_a", "orphan");
FleetMcp.reply(messages, "term_a", Role.WORKER, "orphan");
var before = messages.drainReplies("term_a");
assertEquals(1, before.size(), "one reply in the inbox");
String msgId = before.getFirst().msgId();
// Publish the same reply again and ack it via fleet_ack surface.
FleetMcp.reply(messages, "term_a", "orphan-again");
FleetMcp.reply(messages, "term_a", Role.WORKER, "orphan-again");
var peeked = messages.drainReplies("term_a");
assertEquals(1, peeked.size(), "one fresh reply in the inbox");
@@ -0,0 +1,62 @@
package dev.ltms.fleet.mcp;
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.WorkspaceControl;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.peer.PeerLauncher;
import io.modelcontextprotocol.spec.McpSchema;
import org.junit.jupiter.api.Test;
import java.util.LinkedHashMap;
import java.util.Map;
import java.util.Set;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #395: {@code fleet_profiles} must let an operator tell "this profile's usage-limit
* detection is armed" from "nothing can ever quarantine this profile" — see {@link
* FleetMcp.QuarantineSource#exhaustedPatternArmed()} and {@link
* FleetMcp#profilesView(PeerLauncher, FleetMcp.QuarantineSource, FleetMcp.OutageSource)}.
*/
class FleetProfilesArmedFieldTest {
private static PeerLauncher twoProfileLauncher(FakeHerdr h) {
FleetConfig.Profile armed = new FleetConfig.Profile(
"armed-profile", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
FleetConfig.Profile unarmed = new FleetConfig.Profile(
"unarmed-profile", "http://gx01.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
Map<String, FleetConfig.Profile> profiles = new LinkedHashMap<>();
profiles.put(armed.profile(), armed);
profiles.put(unarmed.profile(), unarmed);
return new ClaudeCodeLauncher(new AgentControl(h), new WorkspaceControl(h),
new SubscriptionGuard(Set.of("gx00.gw", "gx01.gw")), profiles, armed.profile(), _ -> "tok");
}
private static String textOf(McpSchema.CallToolResult r) {
return ((McpSchema.TextContent) r.content().getFirst()).text();
}
@Test
void armedProfileReportsArmedAndUnarmedReportsUnarmed() {
FakeHerdr h = new FakeHerdr();
PeerLauncher workers = twoProfileLauncher(h);
FleetMcp.QuarantineSource source = new FleetMcp.QuarantineSource(
_ -> null, dev.ltms.fleet.placement.BackendQuarantine.none(),
profile -> "armed-profile".equals(profile));
McpSchema.CallToolResult res = FleetMcp.profiles(workers, source);
assertNotEquals(Boolean.TRUE, res.isError());
String out = textOf(res);
assertTrue(out.contains("\"exhaustionDetectionArmed\""), out);
assertTrue(out.contains("\"armed-profile\":true"), out);
assertTrue(out.contains("\"unarmed-profile\":false"), out);
}
}
@@ -8,6 +8,7 @@ import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
@@ -17,6 +18,7 @@ import java.util.regex.Matcher;
import java.util.regex.Pattern;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -45,6 +47,15 @@ class EnvAllowListScrubTest {
/** Env var names appearing in command output; anything else (prompts, wrapped lines) is noise. */
private static final Pattern ENV_NAME = Pattern.compile("^([A-Za-z_][A-Za-z0-9_]*)$");
/**
* fleetd #388: the sentinel name the generated {@code .zshenv} guard uses, kept here as a
* literal rather than referencing {@link EnvAllowListScrub#SCRUB_SENTINEL} — the two tests that
* use it must still compile and run against the pre-fix production class (which has no such
* constant), so the revert-and-prove-it-fails step exercises a real assertion instead of a
* compilation error.
*/
private static final String SCRUB_SENTINEL_NAME = "_CB633_SCRUBBED";
/**
* The equality test. Expected survivors = baseline exports ∩ allowed — i.e. every survivor is
* allowed AND every allowed name that existed survives. The operator's own secret-store exports
@@ -108,6 +119,230 @@ class EnvAllowListScrubTest {
"allowed N of M with N <= M — the denominator is always reported");
}
/**
* fleetd #394: the actual defect. Plain {@code export "$n="} is FATAL for a zsh read-only or
* special parameter (e.g. {@code UID}) and aborts the whole sourced file — every name still to
* come is never blanked, and the {@code scrub-report.txt} below is never written at all,
* silently ({@code 2>/dev/null} swallows the error). This plants an unblankable, exported,
* read-only variable in the MIDDLE of the names the scrub attempts to blank, with two more
* names after it, and asserts that both of those later names are STILL blanked and the report
* is STILL written with the failure counted — a test that only checked names BEFORE the failure
* point would pass today and prove nothing.
*
* <p>The planted name is a made-up one ({@code FLEETD_TEST_UNBLANKABLE}), not {@code UID} or
* any other name a skip-list might already know about — invariant 1 is that the loop survives
* ANY unblankable name, not a known one, so the test must not lean on one either.
*
* <p>Exercises the real artefact: {@link EnvAllowListScrub#scrubScript} is run verbatim under a
* real {@code /bin/zsh}, not just asserted on as a Java string. The four planted names are
* exported one at a time via {@code typeset -x}/{@code typeset -rx} immediately before the
* script runs, in a fixed order — zsh's {@code export}/{@code typeset -x} appends to the
* process's environment table in call order (verified empirically: a freshly-exported name
* always sorts after every inherited one and after every earlier freshly-exported name in
* {@code command env}'s own output), which is what makes the "middle" position deterministic
* here, unlike relying on the OS's own inherited-environment order.
*/
@Test
void unblankableNameInTheMiddleDoesNotAbortNamesAfterIt(@TempDir Path tmp) throws Exception {
assumeTrue(Files.isExecutable(ZSH), "/bin/zsh not present — nothing to prove here");
// Only ZDOTDIR is allowed — it must survive the scrub itself, since the report is written
// to "$ZDOTDIR/..." AFTER the blanking loop runs; if ZDOTDIR were blanked as a side effect,
// the report write would silently go to the wrong place instead of testing anything.
String script = EnvAllowListScrub.scrubScript(Set.of("ZDOTDIR"));
String setup = """
typeset -x FLEETD_TEST_BEFORE=1
typeset -rx FLEETD_TEST_UNBLANKABLE=1
typeset -x FLEETD_TEST_AFTER_A=1
typeset -x FLEETD_TEST_AFTER_B=1
""";
ProcessBuilder pb = new ProcessBuilder("/bin/zsh");
pb.environment().clear();
pb.environment().put("PATH", "/usr/bin:/bin");
pb.environment().put("ZDOTDIR", tmp.toAbsolutePath().toString());
pb.redirectError(ProcessBuilder.Redirect.DISCARD);
Process zsh = pb.start();
zsh.getOutputStream().write((setup + script).getBytes(StandardCharsets.UTF_8));
zsh.getOutputStream().flush();
zsh.getOutputStream().close();
assertTrue(zsh.waitFor(60, java.util.concurrent.TimeUnit.SECONDS),
"the scrub script did not exit within 60s");
assertEquals(0, zsh.exitValue(),
"the scrub script itself must never abort — an unblankable name must not kill the "
+ "sourced file");
EnvAllowListScrub.ScrubReport report = EnvAllowListScrub.readReport(tmp);
assertNotNull(report, "the report must still be written even though one name could not be "
+ "blanked — a report that silently never appears is the #394 bug");
assertTrue(report.blanked().contains("FLEETD_TEST_BEFORE"),
"sanity: the name before the unblankable one must be blanked");
assertTrue(report.blanked().contains("FLEETD_TEST_AFTER_A"),
"the FIRST name AFTER the unblankable one must still be blanked — before the fix, "
+ "the whole loop aborted at the unblankable name and every later name was "
+ "silently left untouched");
assertTrue(report.blanked().contains("FLEETD_TEST_AFTER_B"),
"the SECOND name after the unblankable one must also still be blanked");
assertTrue(report.unblankable().contains("FLEETD_TEST_UNBLANKABLE"),
"the unblankable name is reported by name, not silently dropped");
// fleetd #400 note: zsh itself auto-exports SHLVL on every shell start (measured: it appears
// in `command env` even from a fully cleared parent), and a bare assignment to it is coerced
// rather than failing — exactly the shape #400 fixes. Before that fix, eval's exit status
// alone silently misclassified SHLVL as blanked, so this test's old "exactly one" assertion
// passed by accident: it never actually proved SHLVL was absent from the candidates, only
// that the old bug hid it. Now that classification reads the value back, SHLVL and
// FLEETD_TEST_UNBLANKABLE both correctly land in unblankable() — real, unplanted evidence
// the #400 fix works, not just the synthetic case in the dedicated #400 test above.
assertTrue(report.unblankable().contains("SHLVL"),
"fleetd #400: zsh's own auto-exported SHLVL must also be reported unblankable, not "
+ "silently miscounted as blanked");
assertEquals(report.unblankable().size(), report.failed(),
"the failed count must equal the number of names actually reported unblankable");
}
/**
* fleetd #400: {@code eval}'s exit status is not proof that a name was actually blanked. zsh
* coerces a bare {@code NAME=} assignment on an integer special parameter to a number instead of
* failing, so {@code eval} reports success while the value stays non-empty — a status-based
* classification calls that "blanked" when it was not. This drives all three shapes a name can
* take through the real {@code scrubScript} in ONE run: a normal, genuinely blankable name; a
* fatal one ({@code LINENO} — deliberately not {@code UID}, so this test does not depend on the
* harness's uid); and the silent-no-op one the ticket is about ({@code SECONDS}, rc 0 but
* unchanged). Under the pre-#400 exit-status check, {@code SECONDS} would land in
* {@code blanked()} — that is the exact false receipt this fix removes.
*
* <p>Criterion 3: a cleared {@code ProcessBuilder} parent does not, by itself, give the child
* zsh any of these names — {@code SECONDS}/{@code LINENO} are zsh's own built-in parameters and
* only become CANDIDATES the enumeration loop can see (i.e. show up in {@code command env}) when
* they arrive via the process's own environment table, not merely by existing as zsh parameters
* inside the shell. So each is put into {@code pb.environment()} explicitly, after
* {@code clear()} — confirmed empirically first (a throwaway probe piping
* {@code env -i PATH=... SECONDS=999 LINENO=999 FLEETD_TEST_NORMAL=1 zsh -c 'command env | cut
* -d= -f1'}) that all three names really appear in {@code command env}'s output under exactly
* this construction, not relying on whatever the test-runner's own ambient environment happens
* to contain.
*/
@Test
void classifiesByObservedValueNotExitStatusAcrossAllThreeShapes(@TempDir Path tmp) throws Exception {
assumeTrue(Files.isExecutable(ZSH), "/bin/zsh not present — nothing to prove here");
String script = EnvAllowListScrub.scrubScript(Set.of("ZDOTDIR"));
ProcessBuilder pb = new ProcessBuilder("/bin/zsh");
pb.environment().clear();
pb.environment().put("PATH", "/usr/bin:/bin");
pb.environment().put("ZDOTDIR", tmp.toAbsolutePath().toString());
// Explicitly placed in the child's environment table — see the javadoc above on why a
// cleared parent alone does not put these on the enumeration loop's candidate list.
pb.environment().put("SECONDS", "999"); // rc 0, value coerced/unchanged — the #400 bug
pb.environment().put("LINENO", "999"); // fatal on assignment, eval rc != 0, contained
pb.environment().put("FLEETD_TEST_NORMAL", "1"); // genuinely blankable, the control case
pb.redirectError(ProcessBuilder.Redirect.DISCARD);
Process zsh = pb.start();
zsh.getOutputStream().write(script.getBytes(StandardCharsets.UTF_8));
zsh.getOutputStream().flush();
zsh.getOutputStream().close();
assertTrue(zsh.waitFor(60, java.util.concurrent.TimeUnit.SECONDS),
"the scrub script did not exit within 60s");
assertEquals(0, zsh.exitValue(),
"the scrub script must still reach its end with both a fatal name and a silent "
+ "no-op name among the candidates");
EnvAllowListScrub.ScrubReport report = EnvAllowListScrub.readReport(tmp);
assertNotNull(report, "the report must still be written");
assertTrue(report.blanked().contains("FLEETD_TEST_NORMAL"),
"the control case: an ordinary name is genuinely blankable and must be reported so");
assertTrue(report.unblankable().contains("LINENO"),
"a name fatal to assign to must be reported unblankable — sanity check that "
+ "containment still works under the new classification");
assertTrue(report.unblankable().contains("SECONDS"),
"the #400 defect: eval returns rc 0 for SECONDS (zsh coerces the assignment instead "
+ "of failing) but the value is left non-empty — classifying on the observed "
+ "value catches this; classifying on eval's exit status would have called "
+ "this \"blanked\" and produced a false receipt");
assertFalse(report.blanked().contains("SECONDS"),
"SECONDS must never appear as blanked — it was never actually emptied");
}
/**
* fleetd #394 follow-up: the blanking loop's {@code eval "export ${n}="} splices {@code n} into
* a string that zsh then interprets as shell syntax. That is only safe because every name
* reaching {@code _cb633_blank} already passed an identifier check in the ENUMERATION loop
* (20 lines away, in a different loop) — so the fix re-asserts the identical check immediately
* before the {@code eval} call, rather than trusting that distant guard to keep holding.
*
* <p>This test plants a value with an embedded newline, exploiting the exact "junk from
* multi-line values" gap the enumeration loop's own comment already documents: {@code command
* env}'s text output is read line-by-line, so a value's second line becomes a spurious extra
* "name" that was never a real exported variable. The fragment used here ({@code
* junk.fragment}) is merely non-conforming (it contains a dot) — never command-shaped; this
* test must never demonstrate command execution and plants no command-shaped payload.
*
* <p>Exercises the real artefact end-to-end: {@link EnvAllowListScrub#scrubScript} runs
* verbatim under a real {@code /bin/zsh}, exactly as {@code generate()} would produce it — this
* is not a synthetic call into just the blanking loop.
*/
@Test
void nonIdentifierJunkFromAMultilineValueIsSkippedNotBlankedOrUnblankable(@TempDir Path tmp)
throws Exception {
assumeTrue(Files.isExecutable(ZSH), "/bin/zsh not present — nothing to prove here");
String script = EnvAllowListScrub.scrubScript(Set.of("ZDOTDIR"));
ProcessBuilder pb = new ProcessBuilder("/bin/zsh");
pb.environment().clear();
pb.environment().put("PATH", "/usr/bin:/bin");
pb.environment().put("ZDOTDIR", tmp.toAbsolutePath().toString());
// Embedded newline: `command env`'s own text output splits this into two lines, and the
// second ("junk.fragment") has no "=" at all, so `cut -d= -f1` returns it unchanged as a
// spurious candidate "name" — it was never an actual exported variable by that name.
pb.environment().put("FLEETD_TEST_MULTILINE", "keep\njunk.fragment");
pb.redirectError(ProcessBuilder.Redirect.DISCARD);
Process zsh = pb.start();
zsh.getOutputStream().write(script.getBytes(StandardCharsets.UTF_8));
zsh.getOutputStream().flush();
zsh.getOutputStream().close();
assertTrue(zsh.waitFor(60, java.util.concurrent.TimeUnit.SECONDS),
"the scrub script did not exit within 60s");
assertEquals(0, zsh.exitValue(), "the scrub script must reach its end");
EnvAllowListScrub.ScrubReport report = EnvAllowListScrub.readReport(tmp);
assertNotNull(report, "the report must still be written");
assertTrue(report.blanked().contains("FLEETD_TEST_MULTILINE"),
"sanity: the real, identifier-shaped variable must still be blanked normally");
assertFalse(report.blanked().contains("junk.fragment"),
"a non-identifier fragment is not a real variable and must never be blanked");
assertFalse(report.unblankable().contains("junk.fragment"),
"a non-identifier fragment must never even become a candidate the blanking loop "
+ "attempts — it must be filtered before either guard has to catch it, so "
+ "it is neither blanked nor counted as a failed attempt");
}
/** A group-shared ZDOTDIR still lets the member truncate and write its pre-created receipt. */
@Test
void groupSharedScrubWritesAndReadsItsReport(@TempDir Path tmp) throws Exception {
assumeTrue(Files.isExecutable(ZSH), "/bin/zsh not present — nothing to prove here");
Set<String> allowed = MemberEnvAllowList.derive(List.of());
Path zdotdir = EnvAllowListScrub.generate(tmp, allowed, currentUserGroup());
Map<String, String> cleanParent = Map.of(
"HOME", System.getProperty("user.home"),
"PATH", "/usr/bin:/bin",
"SHELL", "/bin/zsh");
exportedNamesFromCleanParent(cleanParent, zdotdir);
EnvAllowListScrub.ScrubReport report = EnvAllowListScrub.readReport(zdotdir);
assertNotNull(report, "a group-shared completed login shell must leave a report behind");
assertTrue(report.allowed() >= 0 && report.total() >= report.allowed(),
"allowed N of M with N <= M — the denominator is always reported");
assertEquals("rw-rw----", java.nio.file.attribute.PosixFilePermissions.toString(
Files.getPosixFilePermissions(zdotdir.resolve(EnvAllowListScrub.REPORT_FILE))),
"the pre-created receipt must be group-writable");
assertEquals("rwxr-x---", java.nio.file.attribute.PosixFilePermissions.toString(
Files.getPosixFilePermissions(zdotdir)),
"group sharing must not make the ZDOTDIR directory group-writable");
}
/** Report parsing is lenient: absent file → null (no measurement), not an exception. */
@Test
void readReportReturnsNullForADirectoryWithoutOne(@TempDir Path dir) {
@@ -224,6 +459,140 @@ class EnvAllowListScrubTest {
+ "shell. A difference here means the scrub is dead on Linux.");
}
/**
* fleetd #388: the actual gap. zsh reads {@code .zshenv} always, {@code .zprofile}/
* {@code .zlogin} only for a LOGIN shell, and {@code .zshrc} only for an INTERACTIVE one — so a
* shell that is NEITHER (a bare {@code /bin/zsh} reading a script off a non-tty stdin, no
* {@code -l}, no {@code -i}) reads only {@code .zshenv} and stops. Before this fix, that shell
* never reached {@code scrub.zsh} at all: the decoy secret below would survive untouched. This
* test injects that decoy directly into the process environment (not via a sourced dotfile,
* since the whole point of the gap is that {@code .zshenv} is normally close to empty) so the
* test does not depend on any real {@code ~/.zshrc} content existing on the host.
*/
@Test
void scrubRunsInAShellThatIsNeitherLoginNorInteractive(@TempDir Path tmp) throws Exception {
assumeTrue(Files.isExecutable(ZSH), "/bin/zsh not present — nothing to prove here");
Set<String> allowed = MemberEnvAllowList.derive(List.of());
Path zdotdir = EnvAllowListScrub.generate(tmp, allowed);
Map<String, String> cleanParent = new HashMap<>(Map.of(
"HOME", System.getProperty("user.home"),
"PATH", "/usr/bin:/bin",
"SHELL", "/bin/zsh",
"USER", System.getProperty("user.name", "nobody"),
"TMPDIR", tmp.toString()));
cleanParent.put("FLEETD_TEST_DECOY_SECRET", "x"); // not on any allow-list; must be blanked
List<String> neither = List.of(); // no -l, no -i; stdin is a pipe (never a tty) either way
Set<String> baseline = exportedNamesFromCleanParent(cleanParent, null, neither);
Set<String> scrubbed = exportedNamesFromCleanParent(cleanParent, zdotdir, neither);
Set<String> expected = new TreeSet<>();
for (String name : baseline) {
if (MemberEnvAllowList.keeps(allowed, name)) {
expected.add(name);
}
}
assertTrue(baseline.contains("FLEETD_TEST_DECOY_SECRET"),
"sanity: the decoy must actually reach the un-scrubbed baseline, or this test proves "
+ "nothing");
expected.add("ZDOTDIR"); // the harness set it and it is infrastructure, so it must survive
expected.add(SCRUB_SENTINEL_NAME); // set by the new .zshenv guard once scrubbed
assertEquals(expected, scrubbed,
"a pane shell that is NEITHER login nor interactive must still be scrubbed — its "
+ "surviving exported names must EQUAL baseline ∩ allow-list, plus the "
+ "sentinel the guard sets once it has run. FLEETD_TEST_DECOY_SECRET surviving "
+ "here means the gap is still open.");
}
/**
* fleetd #388 invariants 3 and 4, which a name-set equality cannot show: a member's own tooling
* forks plain, non-login, non-interactive zsh processes for a single command (the same shape as
* the pane shell itself), and such a child must (a) keep whatever its parent deliberately set
* for it, never (b) re-run the scrub and blank it, and never (c) overwrite the pane's own
* {@code scrub-report.txt} with a description of itself instead of the pane. All three can only
* be shown by actually running a child process from within the scrubbed pane shell.
*
* <p>The pane and the child both report presence via {@code ${NAME:+present}} — empty when a
* name is unset OR blanked (exported empty), {@code present} when it is set and non-empty. No
* value is ever printed, only these two shapes and the literal word {@code set}/{@code unset}
* for the sentinel.
*/
@Test
void neitherShellChildKeepsParentVariablesAndReceiptStillDescribesThePane(@TempDir Path tmp) throws Exception {
assumeTrue(Files.isExecutable(ZSH), "/bin/zsh not present — nothing to prove here");
Set<String> allowed = MemberEnvAllowList.derive(List.of());
Path zdotdir = EnvAllowListScrub.generate(tmp, allowed);
Map<String, String> paneEnv = new HashMap<>(Map.of(
"HOME", System.getProperty("user.home"),
"PATH", "/usr/bin:/bin",
"SHELL", "/bin/zsh",
"USER", System.getProperty("user.name", "nobody"),
"TMPDIR", tmp.toString()));
paneEnv.put("FLEETD_TEST_DECOY_SECRET", "x"); // not allow-listed; the pane must blank it
paneEnv.put("ZDOTDIR", zdotdir.toAbsolutePath().toString());
// The pane's own script reports what IT sees, then forks a plain non-login, non-interactive
// child — the shape a member's own tooling uses — carrying a variable the "parent" (this
// pane) deliberately set for it, the way git sets GIT_DIR for a hook.
String outerScript = """
print -r -- "PANE_SENTINEL=${%1$s:+set}"
print -r -- "PANE_DECOY=${FLEETD_TEST_DECOY_SECRET:+present}"
FLEETD_TEST_TOOL_VAR=keep /bin/zsh <<'CHILD'
print -r -- "CHILD_LOGIN=$([[ -o login ]] && echo yes || echo no)"
print -r -- "CHILD_INTERACTIVE=$([[ -o interactive ]] && echo yes || echo no)"
print -r -- "CHILD_TOOL_VAR=${FLEETD_TEST_TOOL_VAR:+present}"
print -r -- "CHILD_DECOY=${FLEETD_TEST_DECOY_SECRET:+present}"
print -r -- "CHILD_SENTINEL=${%1$s:+set}"
CHILD
exit
""".formatted(SCRUB_SENTINEL_NAME);
ProcessBuilder pb = new ProcessBuilder("/bin/zsh"); // no -l, no -i: the pane's own shape
pb.environment().clear();
pb.environment().putAll(paneEnv);
pb.redirectError(ProcessBuilder.Redirect.DISCARD);
Process zsh = pb.start();
zsh.getOutputStream().write(outerScript.getBytes(StandardCharsets.UTF_8));
zsh.getOutputStream().flush();
zsh.getOutputStream().close();
String stdout = new String(zsh.getInputStream().readAllBytes(), StandardCharsets.UTF_8);
assertTrue(zsh.waitFor(60, java.util.concurrent.TimeUnit.SECONDS),
"the pane+child probe did not exit within 60s");
assertTrue(zsh.exitValue() == 0, "probe zsh exited non-zero: " + stdout);
Map<String, String> reported = new HashMap<>();
for (String line : stdout.split("\n")) {
int eq = line.indexOf('=');
if (eq > 0) {
reported.put(line.substring(0, eq).trim(), line.substring(eq + 1).trim());
}
}
assertEquals("set", reported.get("PANE_SENTINEL"),
"the pane itself is neither login nor interactive, so the .zshenv guard must have "
+ "run the scrub and exported the sentinel");
assertEquals("", reported.get("PANE_DECOY"),
"the pane must blank a non-allow-listed name — invariant 1");
assertEquals("no", reported.get("CHILD_LOGIN"), "sanity: the child must also be non-login");
assertEquals("no", reported.get("CHILD_INTERACTIVE"), "sanity: the child must also be non-interactive");
assertEquals("present", reported.get("CHILD_TOOL_VAR"),
"invariant 3: a variable the pane deliberately set for its child must survive — a "
+ "child that re-ran the scrub would have blanked it");
assertEquals("", reported.get("CHILD_DECOY"),
"a name already blanked by the pane must stay blanked in the child, never resurrected");
assertEquals("set", reported.get("CHILD_SENTINEL"),
"the child must inherit the sentinel from the pane's environment, or it would re-scrub");
EnvAllowListScrub.ScrubReport report = EnvAllowListScrub.readReport(zdotdir);
assertNotNull(report, "the pane's own scrub pass must leave a report behind");
assertTrue(report.blanked().stream().noneMatch(n -> n.startsWith("FLEETD_TEST_TOOL_VAR")),
"invariant 4: the receipt must still describe the PANE, not the child — a child that "
+ "re-ran the scrub would have rewritten this file to list its own "
+ "FLEETD_TEST_TOOL_VAR as blanked");
}
/**
* Run {@code /bin/zsh -l -i} from a clean parent and return the NAMES it has exported by prompt
* time. With {@code zdotdir} non-null, {@code ZDOTDIR} points at a generated scrub directory, so
@@ -54,6 +54,40 @@ class LeadCoordLoopTest {
assertTrue(channel.peek().isEmpty(), "and is no longer held");
}
@Test
void redeliveryOfAMessageAlreadyWrittenToThePaneIsAckedWithoutAnotherPaneWrite() {
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "recover me"));
var herdr = new FakeHerdr().agentStatus("idle");
var loop = loop(channel, herdr, Map.of(LEAD_TERM, SELF));
loop.tick();
channel.hold(new LeadMessage("m1", PEER, SELF, "recover me"));
loop.tick();
assertEquals(1, prompts(herdr).size(), "a redelivery must not consume the lead pane twice");
assertEquals(List.of("m1", "m1"), channel.acked(), "the redelivery still needs a fresh broker ack");
}
@Test
void aRedeliveryIsAckedEvenWhileTheLeadIsMidTurn() {
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "recover me"));
var herdr = new FakeHerdr().agentStatus("idle");
var loop = loop(channel, herdr, Map.of(LEAD_TERM, SELF));
loop.tick();
channel.hold(new LeadMessage("m1", PEER, SELF, "recover me"));
herdr.agentStatus("working");
loop.tick();
assertEquals(1, prompts(herdr).size(), "the pane is still written exactly once");
assertEquals(List.of("m1", "m1"), channel.acked(),
"a message already written to the pane must be acked even mid-turn: the mid-turn "
+ "gate exists to protect the pane, and this message needs no pane. Gating the ack "
+ "on it leaves the message held on a lead that is busy most of the time, and every "
+ "recovery redelivers it again — which is the loop this fix exists to stop");
assertTrue(channel.peek().isEmpty(), "so it is no longer held");
}
@Test
void leavesTheMessageUnackedWhenTheLeadIsMidTurn() {
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello"));
@@ -92,6 +92,33 @@ class LeadMailboxTest {
}
}
@Test
void ackThrowsWhenRecoveryClearedTheHeldMessage() throws Exception {
String to = coordId("lead-recovery-ack");
try (LeadMailbox inbox = LeadMailbox.open(uri(), to)) {
inbox.publish(to, new LeadMessage("recovery-ack", "lead-from", to, "in flight"));
assertEquals(1, awaitPeek(inbox).size(), "the broker delivery must be held before recovery clears it");
inbox.clearHeldForRecovery();
assertThrows(IllegalStateException.class, () -> inbox.ack("recovery-ack"),
"a cleared delivery has no valid tag, so ack must report that it did not reach the broker");
}
}
@Test
void ackOfAMessageAlreadyAckedOnThisConnectionStaysQuiet() throws Exception {
String to = coordId("lead-repeat-ack");
try (LeadMailbox inbox = LeadMailbox.open(uri(), to)) {
inbox.publish(to, new LeadMessage("repeat-ack", "lead-from", to, "once"));
awaitPeek(inbox);
inbox.ack("repeat-ack");
inbox.ack("repeat-ack");
}
}
@Test
void duplicateMsgIdIsNotDoubleQueued() throws Exception {
String to = coordId("lead-dedup");
@@ -479,7 +479,8 @@ class MessageServiceTest {
// The worker resumes on its own (per the ask() contract) and eventually sends its real
// fleet_reply; the async ticket must still resolve with it, not strand at PENDING.
assertTrue(messages.reply(T, "real result"), "the worker's real reply must still be accepted");
assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET, messages.reply(T, "real result"),
"the worker's real reply must still be accepted, resolving the parked async ticket");
} finally {
messages.setAskTimeoutRaceHookForTest(null);
}
@@ -767,7 +768,9 @@ class MessageServiceTest {
@Test
void replyQueuesInInboxWhenNoSendIsOpen() {
// No send is open for this session — reply should queue in the inbox.
assertTrue(messages.reply(T, "queued-text"), "reply should succeed (queued)");
// fleetd #365: this is the case that must read as QUEUED, not "delivered".
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "queued-text"),
"reply should succeed but only as queued — nothing was waiting for it");
var drained = messages.drainReplies(T);
assertEquals(1, drained.size());
@@ -780,7 +783,9 @@ class MessageServiceTest {
awaitUninterruptibly(T);
// An explicit reply resolves the open send.
assertTrue(messages.reply(T, "send-resolved"), "reply should succeed (resolved live send)");
// fleetd #365: this is the other case — RESOLVED_SEND, distinct from QUEUED above.
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND, messages.reply(T, "send-resolved"),
"reply should succeed by resolving the live waiting send");
// The inbox should be empty — the reply went to the send, not the inbox.
assertTrue(messages.drainReplies(T).isEmpty(), "no reply in the inbox");
@@ -1094,8 +1099,11 @@ class MessageServiceTest {
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(),
"the primary's own bounded wait gives up before the worker finishes resuming");
// The worker keeps working past that window and only now calls fleet_reply.
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
// The worker keeps working past that window and only now calls fleet_reply. The forward
// waiter answer() opened already timed out, so this resolves via the parked async ticket,
// not a live send (fleetd #365).
assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET,
messages.reply(T, "PR opened: https://example/pulls/42"));
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", done.reply(),
@@ -1119,7 +1127,10 @@ class MessageServiceTest {
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome());
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
// No live waiter (answer()'s own forward wait already timed out) — resolves the parked
// async ticket instead (fleetd #365).
assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET,
messages.reply(T, "PR opened: https://example/pulls/42"));
// fleet_stop tears the worker's session down right after the reply landed — this must never
// report the misleading "the worker session was released before it replied": a reply is
@@ -1170,7 +1181,9 @@ class MessageServiceTest {
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting(); // answer() opened its own forward waiter for the resumed worker turn
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
// A live waiter is open (the forward wait above) — this resolves it directly (fleetd #365).
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND,
messages.reply(T, "PR opened: https://example/pulls/42"));
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome(),
"the lead's own answer() call must not throw because ask()'s timeout cleanup raced it");
@@ -1221,8 +1234,9 @@ class MessageServiceTest {
// The worker keeps working past the timeout and only now calls fleet_reply — with no live
// rendezvous waiter open (ask()'s timeout already closed it) and no new send() having
// reopened one for this target.
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
// reopened one for this target. So it resolves the parked async ticket (fleetd #365).
assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET,
messages.reply(T, "PR opened: https://example/pulls/42"));
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", done.reply(),
@@ -1279,7 +1293,10 @@ class MessageServiceTest {
// real reply — reproduce that interleaving directly instead of trying to win a real race.
messages.forgetTurnForTest(turnId);
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
// answer() is still waiting on its own forward waiter for the resumed turn — a live send —
// so this resolves it directly, not the async ticket (fleetd #365).
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND,
messages.reply(T, "PR opened: https://example/pulls/42"));
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome(),
"the primary's own answer() call must still see the worker's real reply");
@@ -1321,7 +1338,9 @@ class MessageServiceTest {
messages.setReplyOrphanTurnIdRaceHookForTest(() -> messages.forgetTurnForTest(turnId));
try {
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
// No live waiter — resolves the parked async ticket (fleetd #365).
assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET,
messages.reply(T, "PR opened: https://example/pulls/42"));
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", done.reply(),
@@ -1353,7 +1372,8 @@ class MessageServiceTest {
injectDelivery();
assertEquals(MessageService.AskOutcome.TIMED_OUT, messages.ask(T, "Q2?", 200).outcome());
assertTrue(messages.reply(T, "which task does this answer?"));
// Ambiguous — two candidates, so it must fall back to the inbox rather than guess (fleetd #365).
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "which task does this answer?"));
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket1).phase(),
"an ambiguous reply must not guess ticket1");
@@ -1782,7 +1802,11 @@ class MessageServiceTest {
Thread asker = new Thread(() -> assertThrows(IllegalStateException.class,
() -> service.ask(T, "which config file?", 30_000)));
asker.start();
awaitTicketPhaseOn(service, ticket, MessageService.Phase.ASKING);
// Phase.ASKING (markAsyncQuestion) is only the FIRST of ask()'s three steps; the assertion
// below depends on the THIRD (pushLoop.onQuestionOpened). Barrier on the push loop's own
// pending-question state instead of the phase — see ReplyPushLoop#pendingQuestionTurnIdsForTest.
MessageService.TaskView asking = awaitTicketPhaseOn(service, ticket, MessageService.Phase.ASKING);
awaitQuestionPendingOn(pushLoop, LEAD, asking.turnId());
assertEquals(ReplyPushLoop.Action.INJECT, pushLoop.decide(LEAD, 0, 0, 0),
"the open question should be the one thing keeping this lead's schedule alive");
@@ -1897,6 +1921,12 @@ class MessageServiceTest {
// Only now does it reply. Under the old clock this reply was born already expired.
assertTrue(rendezvous.resolve(T, "the long report"));
awaitTicketPhaseOn(wiring.service(), slow, MessageService.Phase.DONE);
// #399: same barrier fix as the sibling eviction test, for consistency — this test's
// own assertion happens to survive a late stamp today (the clock is not advanced any
// further between here and the sweep below, so cutoff cannot move past whichever
// value completedNanos ends up stamped with), but DONE is still the wrong thing to
// order on before a sweep that matters for the TTL.
awaitCompletionStamped(wiring.service(), slow);
// A second delegation runs pruneTerminalTickets before it returns.
wiring.service().sendAsync(T, "an unrelated second task");
@@ -1926,6 +1956,11 @@ class MessageServiceTest {
injectDelivery();
assertTrue(rendezvous.resolve(T, "quick result"));
awaitTicketPhaseOn(wiring.service(), done, MessageService.Phase.DONE);
// #399: DONE can be observed before the completion hook stamps completedNanos. Wait
// for the real stamp before advancing the clock, or the hook can run late and stamp
// the ADVANCED time — making cutoff = advanced - TTL unreachable and hiding the very
// eviction this test exists to pin.
awaitCompletionStamped(wiring.service(), done);
// Nobody collected it, and the TTL has now passed since it FINISHED.
clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1));
@@ -1936,6 +1971,80 @@ class MessageServiceTest {
}
}
/**
* fleetd #409: pins the #399 ordering invariant deterministically, on the first run, without
* relying on host load.
*
* <p>fleetd #399 was itself only reproducible probabilistically: the real race window between
* {@code CompletableFuture.complete()} making {@link MessageService.Phase#DONE} observable and
* the constructor's {@code whenComplete} hook actually stamping {@code completedNanos} (see
* {@link MessageService#isCompletionStampedForTest}) is normally a handful of instructions wide,
* and needed heavy background load on the host to show up in a run at all. This test does not
* try to hit that narrow window by chance — it widens it on purpose: {@link #STAMP_DELAY_MILLIS}
* is injected into the clock itself, so the completion hook's one read of {@code nowNanos} for
* this ticket sleeps before returning, and everything the test does in the meantime (advance the
* clock, run a sweep, assert) happens for certain inside that window, on any host.
*
* <p>Removing the {@link #awaitCompletionStamped} call below reproduces the pre-#399 ordering:
* the test then advances the clock and runs its sweep while the hook is still asleep, so the
* hook wakes up and stamps {@code completedNanos} with the clock's ALREADY-ADVANCED value
* instead of the real completion time. {@code cutoff = advanced - TICKET_TTL_NANOS} can then
* never exceed that stamp (the gap between them is fixed at exactly {@code TICKET_TTL_NANOS}),
* so the sweep that already ran never evicts the ticket and the very next assertion — expecting
* eviction — fails immediately. That is a deterministic, first-run RED failure, not a flaky one
* and not a false pass: I ran the test with the barrier call removed and confirmed
* {@code assertNull} fails because {@code poll} still returns the ticket's DONE view, which is
* exactly this masked-eviction mechanism and not some unrelated defect in the test's own wiring.
*/
@Test
void aTicketOrderedOnDoneInsteadOfTheCompletionStampSurvivesAnEvictionItMustNotSurvive() throws Exception {
java.util.concurrent.atomic.AtomicLong clock = new java.util.concurrent.atomic.AtomicLong(1_000_000_000L);
// Armed for exactly one call: MessageService's only three nowNanos() call sites are the Task
// constructor's createdNanos, this constructor's whenComplete hook's completedNanos, and
// pruneTerminalTickets' cutoff — none of which run between arming this (right before
// triggering the reply below) and the completion hook firing, so the delayed call is
// unambiguously that hook's stamp for `ticket`, never a createdNanos or cutoff read.
java.util.concurrent.atomic.AtomicBoolean delayArmed = new java.util.concurrent.atomic.AtomicBoolean(false);
java.util.function.LongSupplier gatedClock = () -> {
if (delayArmed.compareAndSet(true, false)) {
try {
Thread.sleep(STAMP_DELAY_MILLIS);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
return clock.get();
};
MessageService service = new MessageService(agents, injector, rendezvous, inbox, null, null, gatedClock);
try {
String ticket = service.sendAsync(T, "a quick task");
awaitWaiting();
injectDelivery();
delayArmed.set(true);
assertTrue(rendezvous.resolve(T, "quick result"));
awaitTicketPhaseOn(service, ticket, MessageService.Phase.DONE);
// The barrier under test (fleetd #409, same fix as fleetd #399): remove this one call to
// reproduce the pre-#399 ordering — see the class-level note above for what happens then.
awaitCompletionStamped(service, ticket);
// A real clock would need the whole TTL to pass; the injected one does it instantly, and
// by now completedNanos already holds the REAL (small, unadvanced) completion time.
clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1));
service.sendAsync(T, "an unrelated second task"); // runs pruneTerminalTickets before returning
assertNull(service.poll(ticket),
"a ticket whose real completion time is long past the advanced cutoff must be "
+ "evicted, whatever the completion hook's clock read was delayed by");
} finally {
service.close();
}
}
/** Bounded delay the injected clock sleeps for in {@link #aTicketOrderedOnDoneInsteadOfTheCompletionStampSurvivesAnEvictionItMustNotSurvive}. */
private static final long STAMP_DELAY_MILLIS = 300;
private MessageService.TaskView awaitTicketPhaseOn(MessageService svc, String ticket,
MessageService.Phase phase) throws Exception {
long deadline = System.currentTimeMillis() + 3000;
@@ -1970,6 +2079,45 @@ class MessageServiceTest {
return view;
}
/**
* fleetd #399: waits until {@code ticket}'s completion hook has actually stamped
* {@code completedNanos}, not just until {@link MessageService#poll} reports
* {@link MessageService.Phase#DONE} for it. {@code poll} can observe {@code DONE} the instant
* the task's future resolves, before the {@code whenComplete} hook that stamps the completion
* time has run — {@code CompletableFuture.complete()} publishes its result and only then runs
* dependents. A test that is about to advance an injected clock past the TTL must order itself
* after the stamp, not after {@code DONE}: winning the race the other way stamps the
* *advanced* clock value and can hide a real eviction bug behind a false pass.
*/
private void awaitCompletionStamped(MessageService svc, String ticket) throws Exception {
long deadline = System.currentTimeMillis() + 3000;
while (!svc.isCompletionStampedForTest(ticket)) {
assertTrue(System.currentTimeMillis() < deadline,
"completedNanos for " + ticket + " was never stamped");
Thread.sleep(5);
}
}
/**
* fleetd #418: waits until {@code turnId} actually appears in {@code pushLoop}'s own open-question
* state for {@code lead} — {@link ReplyPushLoop#pendingQuestionTurnIdsForTest} — rather than until
* {@link MessageService#poll} reports {@link MessageService.Phase#ASKING}. {@code Phase.ASKING} is
* set by {@code markAsyncQuestion}, the FIRST of three steps {@code MessageService.ask()} performs;
* {@code pushLoop.onQuestionOpened} (the one that actually publishes to {@code pendingQuestions})
* is the THIRD. Under load the asker thread can be descheduled between those two steps, so a
* barrier on the phase alone can release before the push loop has anything pending — a test that
* then asserts on {@link ReplyPushLoop#decide} is asserting on state that has not been published
* yet, not on the throwing path it is named for.
*/
private void awaitQuestionPendingOn(ReplyPushLoop pushLoop, String lead, String turnId) throws Exception {
long deadline = System.currentTimeMillis() + 3000;
while (!pushLoop.pendingQuestionTurnIdsForTest(lead).contains(turnId)) {
assertTrue(System.currentTimeMillis() < deadline,
"turnId " + turnId + " never appeared in the push loop's pending questions for " + lead);
Thread.sleep(5);
}
}
// --- CB-640: fleet health evidence accessors --------------------------------------------
@Test
@@ -2032,7 +2180,7 @@ class MessageServiceTest {
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
assertTrue(messages.reply(T, "resolved-live"));
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND, messages.reply(T, "resolved-live"));
assertFalse(messages.hasStrandedReply(T), "a reply that resolved an open send is not stranded");
MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
@@ -2042,14 +2190,14 @@ class MessageServiceTest {
@Test
void hasStrandedReplyIsTrueWhenNoSendWasWaiting() {
// No send is open for T — the reply queues into the inbox and is recorded as stranded.
assertTrue(messages.reply(T, "nobody was waiting"));
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "nobody was waiting"));
assertTrue(messages.hasStrandedReply(T),
"a reply with no open send strands, even though it is safely queued in the inbox");
}
@Test
void hasStrandedReplyClearsOnceTheTargetsNextDeliveryIsAccepted() throws Exception {
assertTrue(messages.reply(T, "stray"));
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "stray"));
assertTrue(messages.hasStrandedReply(T));
// The next accepted delivery for T clears the stale stranding fact — the one case the
@@ -2067,7 +2215,7 @@ class MessageServiceTest {
@Test
void hasStrandedReplyClearsOnAbandon() {
assertTrue(messages.reply(T, "stray"));
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "stray"));
assertTrue(messages.hasStrandedReply(T));
messages.abandon(T, "session released");
@@ -984,7 +984,7 @@ class ReplyPushLoopTest {
// --- metrics (CB-512) ----------------------------------------------------------------------
@Test
void successfulNudgeIncrementsDelivered() throws Exception {
void successfulNudgeIncrementsSent() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
@@ -994,11 +994,13 @@ class ReplyPushLoopTest {
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS),
"one nudge (1 agent.prompt call) should have been sent");
// The delivered count is bumped on the scheduler thread right after the send that releases
// The sent count is bumped on the scheduler thread right after the send that releases
// the latch — settle briefly so the counter is published before we read it.
Thread.sleep(200);
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "delivered"),
"a successfully sent nudge must count as delivered");
// fleetd #365: "sent", not "delivered" — this only proves the herdr call succeeded, not
// that the primary's pane read it.
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "sent"),
"a successfully sent nudge must count as sent");
}
@Test
@@ -1013,11 +1015,11 @@ class ReplyPushLoopTest {
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "exhausted"),
"hitting the reminder cap must count as exhausted");
assertEquals(0, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "delivered"));
assertEquals(0, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "sent"));
}
@Test
void successfulTicketNudgeIncrementsDelivered() throws Exception {
void successfulTicketNudgeIncrementsSent() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
Metrics metrics = new Metrics();
@@ -1026,8 +1028,8 @@ class ReplyPushLoopTest {
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one ticket nudge should have been sent");
Thread.sleep(200);
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "delivered"),
"a successfully sent ticket nudge must count as delivered, same metric as CB-307");
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "sent"),
"a successfully sent ticket nudge must count as sent, same metric as CB-307");
}
@Test
@@ -0,0 +1,144 @@
package dev.ltms.fleet.peer;
import org.junit.jupiter.api.Test;
import java.lang.reflect.Constructor;
import java.lang.reflect.RecordComponent;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.TreeSet;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* Fleetd #382, the same "defect factory" #357 and #358 guarded on {@code FleetConfig.withDefaults()}
* and {@link dev.ltms.fleet.session.MemberSession}'s rebuild sites — reproduced here on
* {@link SpawnRequest#withProfile(String)}.
*
* <p>{@code SpawnRequest} has back-compat constructors at arity 3 and 5 alongside its canonical
* arity-6 constructor. Before this ticket, {@code CompositePeerLauncher} routed a profile by
* building a fresh {@code SpawnRequest} from a literal {@code new SpawnRequest(...)} call listing
* six of the original request's own accessors. That call is only correct because it happens to
* name exactly six arguments today — add a 7th component and the established back-compat pattern
* (a new constructor at the old, now-shorter arity) and a call one argument short of the new
* canonical arity would silently rebind to that back-compat constructor, dropping the new
* component on every profile-routed spawn without any compile error. {@link #withProfile} replaces
* that literal call, so this test guards the ONE rebuild site instead of a call site scattered
* through a launcher.
*
* <p>The check below builds a {@link SpawnRequest} through the TRUE canonical constructor —
* resolved by the record's own component types via {@code getDeclaredConstructor}, never by
* argument count, so it can never itself land on a back-compat overload — with a real, distinctive,
* non-null value in every component, calls {@link SpawnRequest#withProfile(String)}, and asserts
* every component the method is not documented to change survives unchanged, while {@code
* profileName} comes back as the new value it was given. A component that comes back anything else
* was silently dropped — the shape of the defect this test exists to catch.
*
* <p>{@link #EXCLUDED_FROM_SURVIVAL_CHECK} is kept deliberately empty and size-pinned by
* {@link #exclusionListSizeIsPinned()}, for the same reason the other two guards pin theirs at
* zero: a checker whose escape hatch can grow to silence a failure is not a checker. Every one of
* {@link SpawnRequest}'s 6 current components has a real, non-null value here and none is excluded.
*/
class SpawnRequestWithProfilePreservesEveryComponentTest {
private static final RecordComponent[] COMPONENTS = SpawnRequest.class.getRecordComponents();
/** Deliberately empty today; grow it only with a matching justification, and re-pin the size. */
private static final Set<String> EXCLUDED_FROM_SURVIVAL_CHECK = Set.of();
/** One real, distinctive, non-null value per component — none of the 6 is excluded. */
private static Map<String, Object> baseValues() {
Map<String, Object> v = new LinkedHashMap<>();
v.put("profileName", "profile-guard");
v.put("requestedCwd", "/wt/requested-guard");
v.put("callerCwd", "/wt/caller-guard");
v.put("sessionName", "session-guard");
v.put("resumeSessionId", "resume-guard");
v.put("role", MemberRole.REVIEWER);
assertNamesMatchComponents(v);
return v;
}
/**
* Guards {@link #baseValues()} itself against drifting from the record's real shape — forgetting
* to add a new component here fails this assertion by name, rather than silently checking one
* component fewer than the record has.
*/
private static void assertNamesMatchComponents(Map<String, Object> values) {
Set<String> names = new TreeSet<>();
for (RecordComponent rc : COMPONENTS) {
names.add(rc.getName());
}
assertEquals(names, new TreeSet<>(values.keySet()),
"this test's value map has drifted from SpawnRequest's actual components — "
+ "update baseValues() alongside the record");
}
/**
* Builds a {@link SpawnRequest} through the TRUE canonical constructor — resolved by the
* record's own component types, not by argument count — so this never accidentally exercises a
* back-compat overload the way a literal {@code new SpawnRequest(...)} call risks doing.
*/
private static SpawnRequest requestOf(Map<String, Object> values) throws ReflectiveOperationException {
Class<?>[] types = Arrays.stream(COMPONENTS).map(RecordComponent::getType).toArray(Class<?>[]::new);
Object[] args = Arrays.stream(COMPONENTS).map(rc -> values.get(rc.getName())).toArray();
Constructor<SpawnRequest> ctor = SpawnRequest.class.getDeclaredConstructor(types);
return ctor.newInstance(args);
}
@Test
void exclusionListSizeIsPinned() {
assertEquals(0, EXCLUDED_FROM_SURVIVAL_CHECK.size(),
"EXCLUDED_FROM_SURVIVAL_CHECK grew from 0 — every entry needs a justification in "
+ "this test class's javadoc AND this assertion re-pinned to the new size; a "
+ "growing exclusion list that silences failures on its own is not a guard");
}
@Test
void withProfilePreservesEveryOtherComponent() throws ReflectiveOperationException {
Map<String, Object> base = baseValues();
SpawnRequest request = requestOf(base);
SpawnRequest result = request.withProfile("profile-updated");
Map<String, Object> expectedOverrides = Map.of("profileName", "profile-updated");
List<String> dropped = new ArrayList<>();
int checked = 0;
for (RecordComponent rc : COMPONENTS) {
String name = rc.getName();
if (EXCLUDED_FROM_SURVIVAL_CHECK.contains(name)) {
continue;
}
checked++;
Object expected = expectedOverrides.containsKey(name) ? expectedOverrides.get(name) : base.get(name);
Object actual;
try {
actual = rc.getAccessor().invoke(result);
} catch (ReflectiveOperationException e) {
throw new RuntimeException("failed to read SpawnRequest." + name + "()", e);
}
if (!Objects.equals(expected, actual)) {
dropped.add(String.format(Locale.ROOT,
"%s: withProfile() was expected to carry (%s) for '%s' but returned %s — a "
+ "component silently dropped by withProfile(), the shape of the defect "
+ "this test exists to catch (its final \"return new SpawnRequest(...)\" "
+ "call binding to a back-compat constructor instead of the true "
+ "canonical one)",
name, expected, name, actual));
}
}
System.out.printf(Locale.ROOT,
"SpawnRequest.withProfile() component-survival coverage — %d components, %d checked, "
+ "%d excluded, %d survived%n",
COMPONENTS.length, checked, EXCLUDED_FROM_SURVIVAL_CHECK.size(), checked - dropped.size());
assertEquals(List.of(), dropped,
"withProfile() silently dropped these components: " + dropped);
}
}
@@ -476,6 +476,11 @@ class FleetAppTest {
Thread.sleep(200);
HttpResponse<String> reply = postJson(port, "/sessions/term_a/reply", "{\"content\":\"LGTM ship it\"}");
assertEquals(200, reply.statusCode());
// fleetd #365: "delivered" used to be unconditionally true; a live send was actually waiting
// here, so this is the case where it must genuinely read true, with outcome naming why.
JsonNode replyBody = mapper.readTree(reply.body());
assertEquals(true, replyBody.get("delivered").asBoolean());
assertEquals("resolved_send", replyBody.get("outcome").asText());
HttpResponse<String> res = send.get(6, java.util.concurrent.TimeUnit.SECONDS);
assertEquals(200, res.statusCode());
@@ -527,6 +532,11 @@ class FleetAppTest {
int port = startHealthy();
HttpResponse<String> res = postJson(port, "/sessions/term_a/reply", "{\"content\":\"orphan\"}");
assertEquals(200, res.statusCode());
// fleetd #365: nothing was waiting, so "delivered" must now read false, not the old
// unconditional true — outcome names this as queued.
JsonNode resBody = mapper.readTree(res.body());
assertEquals(false, resBody.get("delivered").asBoolean());
assertEquals("queued", resBody.get("outcome").asText());
// The queued reply is drainable.
HttpResponse<String> drain = req(port, "GET", "/sessions/term_a/replies");
@@ -1532,26 +1532,52 @@ class GitWorktreesTest {
* for repo setup.
*
* <p>Scope, measured on the fleetd #369 merge and narrower than an earlier version of this
* comment claimed: this protects the 5 {@link #seedingGitWorktrees} sites plus — through
* comment claimed: this protects the {@link #seedingGitWorktrees} call sites plus — through
* {@link #gitProcessBuilder} — every {@code git} subprocess the TEST itself starts. It does
* NOT cover the other 53 {@code new GitWorktrees(...)} constructions in this file, which pass
* NOT cover the {@code new GitWorktrees(...)} constructions elsewhere in this file that pass
* no env override, so a production instance built that way still inherits the JVM's real
* environment. Stripping this override from {@code seedingGitWorktrees} leaves the class green
* both with and without the poison command above, so that half is currently unpinned.
* environment. (Re-measured for fleetd #373, on this file as it stands here: 4 call sites go
* through {@link #seedingGitWorktrees(Path, String, Path)} — not 5, an earlier count this
* comment and fleetd #373's own ticket text both repeated without re-running it — out of 59
* total {@code new GitWorktrees(...)} occurrences, one of which is the shared construction
* inside {@link #seedingGitWorktrees(Path, String, Map)} itself. This class-wide count moves
* every time a test is added, so treat any number here as a snapshot, not a fact to cite
* without recounting.) Stripping the {@code gitEnv} override from a {@link
* #seedingGitWorktrees} call site leaves the class green both with and without the poison
* command above for that call site's OWN test, so that half was unpinned until fleetd #373
* added {@link #seedingGitWorktreesResolvesTheExcludesFileFallbackInsideItsThrowawayDirectory}
* below, which asserts the property directly instead of relying on a poisoned real machine.
*/
private static Map<String, String> hermeticGitEnv(Path tmp) {
return hermeticGitEnvAt(tmp.resolve("hermetic-xdg-config-home-" + System.nanoTime()));
}
/** Same isolation as {@link #hermeticGitEnv(Path)}, with an explicit {@code XDG_CONFIG_HOME}
* instead of a fresh nanoTime-unique one under {@code tmp} — used by
* {@link #seedingGitWorktreesResolvesTheExcludesFileFallbackInsideItsThrowawayDirectory}
* (fleetd #373) so it can pre-populate that directory with a marker BEFORE the production
* {@link GitWorktrees} instance reads it, something the random per-call name from
* {@link #hermeticGitEnv(Path)} makes impossible to predict from outside. */
private static Map<String, String> hermeticGitEnvAt(Path xdgConfigHome) {
return Map.of(
"GIT_CONFIG_GLOBAL", "/dev/null",
"GIT_CONFIG_SYSTEM", "/dev/null",
"GIT_TERMINAL_PROMPT", "0",
"XDG_CONFIG_HOME", tmp.resolve("hermetic-xdg-config-home-" + System.nanoTime()).toString());
"XDG_CONFIG_HOME", xdgConfigHome.toString());
}
/** {@link GitWorktrees}'s full test seam, with a {@code memberSkillsSource} and no other
* overrides — the shape every seeding test below needs, isolated via {@link #hermeticGitEnv}. */
private static GitWorktrees seedingGitWorktrees(Path root, String memberSkillsSource, Path tmp) {
return new GitWorktrees(root.toString(), null, _ -> {}, null, null, memberSkillsSource,
hermeticGitEnv(tmp));
return seedingGitWorktrees(root, memberSkillsSource, hermeticGitEnv(tmp));
}
/** Same shape as {@link #seedingGitWorktrees(Path, String, Path)}, taking an already-built
* {@code gitEnv} directly rather than computing one via {@link #hermeticGitEnv(Path)} — lets
* fleetd #373's test drive the exact production construction a real member spawn uses, with a
* {@code gitEnv} it has already pre-populated a marker into. */
private static GitWorktrees seedingGitWorktrees(Path root, String memberSkillsSource, Map<String, String> gitEnv) {
return new GitWorktrees(root.toString(), null, _ -> {}, null, null, memberSkillsSource, gitEnv);
}
/** Acceptance criterion 2 (part 1): a worktree with no {@code .claude/} at all gets the skill
@@ -1784,6 +1810,63 @@ class GitWorktreesTest {
+ "after skill seeding ran — got:\n" + porcelain);
}
/**
* fleetd #373. Pins the production seam that fleetd #362 review finding 2 protects: {@link
* GitWorktrees#previouslyEffectiveExcludesFileContent}'s XDG-fallback branch reads {@code
* XDG_CONFIG_HOME}/{@code HOME} straight in Java, not through a {@code git} subprocess, so
* {@code gitEnv} — the constructor seam every {@link #seedingGitWorktrees} instance in this
* class is built with — is the ONLY thing that can isolate it. A mutation run during the
* fleetd #372/#369 merge found this unpinned: replacing {@code hermeticGitEnv(tmp)} with
* {@code null} in {@link #seedingGitWorktrees(Path, String, Path)} left every test in this
* class green — the 56 tests that would fail against a real machine's poisoned {@code
* XDG_CONFIG_HOME} were fixed by fleetd #369's subprocess-level isolation, but none of them
* looks at what THIS Java-side read resolves, so deleting the override stays invisible.
*
* <p>This test asserts the PROPERTY, not the constructor argument: a {@link GitWorktrees}
* built for seeding — through the very same {@link #seedingGitWorktrees(Path, String, Map)}
* construction every other seeding test in this class goes through — must resolve the
* excludes-file fallback inside its own throwaway {@code gitEnv}-supplied directory. It needs
* NO externally-set poisoned environment variable: the marker pattern below is written ONLY
* inside a throwaway {@code XDG_CONFIG_HOME} this test controls directly (bypassing {@link
* #hermeticGitEnv(Path)}'s unpredictable nanoTime-named directory, via {@link
* #hermeticGitEnvAt}, so the marker can be in place before the production instance ever reads
* it), reachable ONLY through the {@code gitEnv} seam. If that seam is stripped, the
* production code instead falls back to resolving the REAL {@code XDG_CONFIG_HOME}/{@code
* HOME} of the machine running the test — which does not carry this marker — so the marker
* file below shows up as untracked and the assertion fails on any machine, with no poison
* command required. See the PR body for the pasted failure from actually running that
* mutation (removing the {@code gitEnv} override from this test's own construction).
*/
@Test
void seedingGitWorktreesResolvesTheExcludesFileFallbackInsideItsThrowawayDirectory(@TempDir Path tmp)
throws Exception {
Path xdgConfigHome = tmp.resolve("cb373-xdg-config-home");
Files.createDirectories(xdgConfigHome.resolve("git"));
Files.writeString(xdgConfigHome.resolve("git").resolve("ignore"), "cb373-xdg-fallback-marker\n");
Map<String, String> gitEnv = hermeticGitEnvAt(xdgConfigHome);
Path repo = initRepo(tmp.resolve("repo"));
Path skillsSource = tmp.resolve("skills-src");
writeSkill(skillsSource, "implementer", "IMPLEMENTER SKILL\n");
GitWorktrees seeding = seedingGitWorktrees(tmp.resolve("wts"), skillsSource.toString(), gitEnv);
String wt = seeding.add(repo.toString(), "cb-373-xdg-seam", "HEAD");
assertEquals("IMPLEMENTER SKILL\n",
Files.readString(Path.of(wt, ".claude", "skills", "implementer", "SKILL.md")),
"fixture check — the skill really was seeded, so previouslyEffectiveExcludesFileContent ran");
Files.writeString(Path.of(wt, "cb373-xdg-fallback-marker"),
"would only be invisible to git status if the fallback resolved THIS throwaway "
+ "XDG_CONFIG_HOME rather than the real machine's\n");
String porcelain = fullStatus(Path.of(wt));
assertEquals("", porcelain,
"the marker pattern lives only in this test's throwaway XDG_CONFIG_HOME; git "
+ "status must still be empty, proving the production seam resolved the "
+ "excludes-file fallback through the gitEnv seam rather than the JVM's "
+ "real environment — got:\n" + porcelain);
}
/**
* fleetd #369, acceptance criterion 4 — make the fix hard to undo by accident. Every git
* subprocess this class starts is required to go through {@link #gitProcessBuilder}, the one
@@ -0,0 +1,190 @@
package dev.ltms.fleet.session;
import dev.ltms.fleet.peer.CharterReceipt;
import dev.ltms.fleet.peer.MemberRole;
import org.junit.jupiter.api.Test;
import java.lang.reflect.Constructor;
import java.lang.reflect.RecordComponent;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.TreeSet;
import java.util.function.Function;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* Fleetd #358, the same "defect factory" #357 guarded on {@code FleetConfig.withDefaults()}
* (see {@code FleetConfigWithDefaultsPreservesEveryComponentTest}), reproduced here on
* {@link MemberSession} — the worse of the two sibling cases named in #358, because this record
* has FIVE independent rebuild sites instead of one: {@link MemberSession#withState},
* {@link MemberSession#withActivity}, {@link MemberSession#bumpTurn},
* {@link MemberSession#withAgentSessionId} and {@link MemberSession#withFailureReason} each end in
* their own literal {@code new MemberSession(...)} call. Add a 16th component and add the
* established back-compat constructor at the old (15-arg) arity, and every one of those five
* literal calls becomes a legal match for that new overload — silently dropping the new component,
* independently, on whichever of the five paths a missed update leaves behind. That is harder to
* spot than #357's single call site: the field would survive through some transitions and vanish
* through others.
*
* <p>Each check below builds one {@link MemberSession} through the TRUE canonical constructor —
* resolved by the record's own component types via {@code getDeclaredConstructor}, never by
* argument count, so it can never itself land on a back-compat overload — with a real, distinctive,
* non-null value in every component, calls the real rebuild method under test, and asserts every
* component the method is not documented to change survives unchanged, while the component(s) it IS
* documented to change come back as the new value it was given. A component that comes back
* anything else was silently dropped or lost — the shape of the defect this test exists to catch.
*
* <p>{@link #EXCLUDED_FROM_SURVIVAL_CHECK} is kept deliberately empty and size-pinned by
* {@link #exclusionListSizeIsPinned()}, for the same reason {@code FleetConfig}'s guard pins its own
* exclusion list at zero: a checker whose escape hatch can grow to silence a failure is not a
* checker. Every one of {@link MemberSession}'s 15 current components has a real, non-null,
* non-blank value here and none is excluded.
*/
class MemberSessionRebuildPreservesEveryComponentTest {
private static final RecordComponent[] COMPONENTS = MemberSession.class.getRecordComponents();
/** Deliberately empty today; grow it only with a matching justification, and re-pin the size. */
private static final Set<String> EXCLUDED_FROM_SURVIVAL_CHECK = Set.of();
/** One real, distinctive, non-null value per component — none of the 15 is excluded. */
private static Map<String, Object> baseValues() {
Map<String, Object> v = new LinkedHashMap<>();
v.put("paneId", "pane-guard");
v.put("terminalId", "term-guard");
v.put("profile", "profile-guard");
v.put("role", MemberRole.REVIEWER);
v.put("cwd", "/wt/guard");
v.put("ownerTerminal", "owner-guard");
v.put("spawnedAtNanos", 111_111L);
v.put("lastActivityAtNanos", 222_222L);
v.put("turnCount", 7);
v.put("state", MemberSession.State.BUSY);
v.put("worktree", "/wt/guard-tree");
v.put("branch", "worker/guard-branch");
v.put("charterReceipt", new CharterReceipt(
MemberRole.DEV, "profile-guard", "fleet.charters.dev", "deadbeefguard", 42));
v.put("agentSessionId", "agent-guard");
v.put("failureReason", "reason-guard");
assertNamesMatchComponents(v);
return v;
}
/**
* Guards {@link #baseValues()} itself against drifting from the record's real shape — forgetting
* to add a new component here fails this assertion by name, rather than silently checking one
* component fewer than the record has.
*/
private static void assertNamesMatchComponents(Map<String, Object> values) {
Set<String> names = new TreeSet<>();
for (RecordComponent rc : COMPONENTS) {
names.add(rc.getName());
}
assertEquals(names, new TreeSet<>(values.keySet()),
"this test's value map has drifted from MemberSession's actual components — "
+ "update baseValues() alongside the record");
}
/**
* Builds a {@link MemberSession} through the TRUE canonical constructor — resolved by the
* record's own component types, not by argument count — so this never accidentally exercises a
* back-compat overload the way a literal {@code new MemberSession(...)} call risks doing.
*/
private static MemberSession sessionOf(Map<String, Object> values) throws ReflectiveOperationException {
Class<?>[] types = Arrays.stream(COMPONENTS).map(RecordComponent::getType).toArray(Class<?>[]::new);
Object[] args = Arrays.stream(COMPONENTS).map(rc -> values.get(rc.getName())).toArray();
Constructor<MemberSession> ctor = MemberSession.class.getDeclaredConstructor(types);
return ctor.newInstance(args);
}
@Test
void exclusionListSizeIsPinned() {
assertEquals(0, EXCLUDED_FROM_SURVIVAL_CHECK.size(),
"EXCLUDED_FROM_SURVIVAL_CHECK grew from 0 — every entry needs a justification in "
+ "this test class's javadoc AND this assertion re-pinned to the new size; a "
+ "growing exclusion list that silences failures on its own is not a guard");
}
/**
* Shared check for one rebuild site: build a base session with a real value in every component,
* call {@code rebuild}, and assert every component comes back equal to {@code expectedOverrides}
* when named there, or equal to the base value otherwise. Prints the same denominator style as
* {@code FleetConfigWithDefaultsPreservesEveryComponentTest}.
*/
private void checkRebuildSite(String siteName, Function<MemberSession, MemberSession> rebuild,
Map<String, Object> expectedOverrides) throws ReflectiveOperationException {
Map<String, Object> base = baseValues();
MemberSession session = sessionOf(base);
MemberSession result = rebuild.apply(session);
List<String> dropped = new ArrayList<>();
int checked = 0;
for (RecordComponent rc : COMPONENTS) {
String name = rc.getName();
if (EXCLUDED_FROM_SURVIVAL_CHECK.contains(name)) {
continue;
}
checked++;
Object expected = expectedOverrides.containsKey(name) ? expectedOverrides.get(name) : base.get(name);
Object actual;
try {
actual = rc.getAccessor().invoke(result);
} catch (ReflectiveOperationException e) {
throw new RuntimeException("failed to read MemberSession." + name + "()", e);
}
if (!Objects.equals(expected, actual)) {
dropped.add(String.format(Locale.ROOT,
"%s: %s() was expected to carry (%s) for '%s' but returned %s — a component "
+ "silently dropped by %s(), the shape of the defect this test exists "
+ "to catch (its final \"return new MemberSession(...)\" call binding "
+ "to a back-compat constructor instead of the true canonical one)",
name, siteName, expected, name, actual, siteName));
}
}
System.out.printf(Locale.ROOT,
"MemberSession.%s() component-survival coverage — %d components, %d checked, %d "
+ "excluded, %d survived%n",
siteName, COMPONENTS.length, checked, EXCLUDED_FROM_SURVIVAL_CHECK.size(),
checked - dropped.size());
assertEquals(List.of(), dropped,
siteName + "() silently dropped these components: " + dropped);
}
@Test
void withStatePreservesEveryOtherComponent() throws ReflectiveOperationException {
checkRebuildSite("withState", s -> s.withState(MemberSession.State.FAILED),
Map.of("state", MemberSession.State.FAILED));
}
@Test
void withActivityPreservesEveryOtherComponent() throws ReflectiveOperationException {
checkRebuildSite("withActivity", s -> s.withActivity(999_999L),
Map.of("lastActivityAtNanos", 999_999L));
}
@Test
void bumpTurnPreservesEveryOtherComponent() throws ReflectiveOperationException {
checkRebuildSite("bumpTurn", s -> s.bumpTurn(999_999L),
Map.of("lastActivityAtNanos", 999_999L, "turnCount", 8));
}
@Test
void withAgentSessionIdPreservesEveryOtherComponent() throws ReflectiveOperationException {
checkRebuildSite("withAgentSessionId", s -> s.withAgentSessionId("agent-updated"),
Map.of("agentSessionId", "agent-updated"));
}
@Test
void withFailureReasonPreservesEveryOtherComponent() throws ReflectiveOperationException {
checkRebuildSite("withFailureReason", s -> s.withFailureReason("reason-updated"),
Map.of("failureReason", "reason-updated"));
}
}