Compare commits

...

38 Commits

Author SHA1 Message Date
ltms 3762aca307 Merge #606: wire the lead context gauge to the real configDir, and stop flapping on a torn line
CI / shell-tests (pull_request) Failing after 6s
CI / build (pull_request) Failing after 1m21s
CI / contract (pull_request) Successful in 1m50s
Fixes the two defects recorded in #602's own body.

1. FleetMcp.contextView passed configDir=null, so the gauge read <user.home>/
   .claude while this host's lead profile sets an override. Measured before the
   fix: 19 transcripts under the real directory, 0 under the fallback. The gauge
   would have deployed green and reported UNKNOWN forever, for every lead.
   Now threaded via FleetMcp.LeadConfigDirSource, built by Fleetd
   .leadConfigDirSource, following fleet.leaders.<name>.profile to that
   profile's configDir and reading config.get() live inside the lambda.

2. A torn final line no longer means UNKNOWN. fleet_list reads a transcript
   Claude Code may be mid-write on, so the last line can be cut. The old code
   treated that as fatal, which would make the gauge flap at random. The stated
   reason ("a format change should show as UNKNOWN") does not hold: a real
   format change makes EVERY line unparseable, and that case is still caught.

Verified by the lead on the combined state with current main merged in, not on
the branch alone: mvn -o clean install exit 0, 147 reports, 1841 tests, 0
failures, 0 errors, 0 skipped.

Mutation-checked independently by the lead:
  Fleetd.leadConfigDirSource body -> none()   -> KILLED (1 failure)
  FleetMcp.contextView configDir -> null      -> KILLED (worker-measured)

KNOWN RESIDUAL, documented rather than overclaimed. main's own one-line call to
leadConfigDirSource could be swapped for none() and the suite stays green. Every
member of this wiring-test family (loopHealthSource, capacitySource,
healthCoverageSource) has the identical gap — no test runs Fleetd.main far enough
to observe which factory it called. Filed separately as a class-wide problem
rather than patched here. The empirical close is the dogfood check after redeploy.
2026-09-20 11:38:22 +02:00
Dai Ha 4e27bde2d7 fleetd #602 gauge-wiring follow-up: pin Fleetd.main's LeadConfigDirSource wiring
Extract the inline new FleetMcp.LeadConfigDirSource(leadConfigDirLookup(...))
construction in Fleetd.main into a package-private factory,
Fleetd.leadConfigDirSource, mirroring loopHealthSource/capacitySource/
healthCoverageSource. Add FleetdLeadConfigDirSourceWiringTest, which calls the
factory directly with real Profile/Leader fixtures and asserts the returned
source resolves a real configDir -- a property that is false if the factory's
body is mutated to return LeadConfigDirSource.none().

Neither FleetMcpLeadContextGaugeWiringTest nor FleetdLeadConfigDirLookupTest
could catch main losing this wiring: each builds its own instance instead of
calling what main calls. This closes that gap at the factory level, matching
the standard already accepted for loopHealthSource's own wiring test.
2026-09-20 16:34:49 +07:00
Dai Ha d345b14e43 fleetd #602 gauge-wiring: thread a lead's configured configDir into the context gauge
FleetMcp.contextView hardcoded LeadContextGauge.read(null, ...), so a lead whose
profile sets its own CLAUDE_CONFIG_DIR always read the wrong transcript directory
and reported UNKNOWN forever, with no error anywhere.

- Add FleetMcp.LeadConfigDirSource (same idiom as LeadSeatSource) and thread it
  through the constructor / listFleet overload chain / leadView / contextView.
- Add Fleetd.leadConfigDirLookup, wired at construction, following the same
  fleet.leaders.<name>.profile link leadSeatLookup already uses, one step
  further to that profile's own configDir.
- LeadContextGauge.parse: a single unparseable line (typically the final one,
  torn by a write this read raced) is now skipped rather than forcing UNKNOWN;
  only when every line in the read window fails to parse does it report
  UNKNOWN, which is the real format-change signal.
- Tests: FleetdLeadConfigDirLookupTest (lookup logic), FleetMcpLeadContextGaugeWiringTest
  (end-to-end: config naming directory A vs B decides which is read; a lead with
  no configured dir degrades without throwing), and two replacement properties in
  LeadContextGaugeTest for the torn-line fix plus its all-unparseable control.
2026-09-20 16:19:37 +07:00
Dai Ha 3ed7bfca67 fleetd: report a lead's live context usage in fleet_list
CI / shell-tests (pull_request) Failing after 7s
CI / build (pull_request) Failing after 1m37s
CI / contract (pull_request) Successful in 2m7s
fleetd had no way to see how full a lead's context window is. On this host a
lead auto-compacted 30 times in one session, discarding roughly 250,000 tokens
and costing 46s to 3m16s each time, and nothing could see it coming.

LeadContextGauge reads the transcript Claude Code itself writes, never the
lead's pane. It finds <sessionId>.jsonl by NAME under <configDir>/projects/
rather than deriving the project slug, which is an undocumented internal.

Three states, not two: OK, HIGH, UNKNOWN. Every path that cannot positively
establish a token count reports UNKNOWN with no number, so a lead is never
told it is fine when the honest answer is "I could not look".

Bounded two ways: TAIL_BYTES caps bytes read off disk, and a 5s cache TTL caps
how often that read happens, because fleet_list is polled constantly.

Recovered by the lead: the authoring member ended on a backend error (DNS
ENOTFOUND) with this work uncommitted and unpushed in its worktree. Verified
before committing: mvn -o clean install exit 0, 1813 tests, 0 failures,
0 errors, 0 skipped, 144 reports; LeadContextGaugeTest 8/8.

KNOWN INCOMPLETE - see the PR. The fleet_list call site passes configDir=null,
which falls back to ~/.claude, but this host's lead profile sets configDir to
an override. Measured: 19 transcripts under the real configDir, 0 under the
fallback. The gauge is therefore INERT on this fleet until that is wired.
2026-09-20 16:05:15 +07:00
ltms 6eb34a654f Merge #599: fleetd #589 groups 1+2 — wiring-test 6 sites in Fleetd.main()
CI / shell-tests (push) Failing after 11s
CI / contract (push) Successful in 1m18s
CI / build (push) Successful in 1m43s
Extracts 6 inline constructions in Fleetd.main() (:214-:468) to
package-private factories, each pinned by a behavioural wiring test.
Production behaviour unchanged.

Gate: test-merged onto current main (already carrying #598, which rewrote
117 lines of the same file). No conflict — the two units insert at different
anchors, #599 after capacitySource and #598 after loopHealthSource, as their
briefs specified. mvn exit=0, 1805 tests / 0 failures / 0 errors from 143
surefire reports; all 6 new classes confirmed to have run with their
expected counts. Diff confirmed extraction-only.

No source-text assertions in any of the 6 new files; that zero carries a
positive control (the same pattern finds 18 such files elsewhere in the
repo). Worker self-reported fixing a mutation that threw NullPointerException
rather than failing an assertion — the assertNotNull guard is present ahead
of the matcher call, confirmed.
2026-09-19 10:39:47 +02:00
ltms 7084d99b89 Merge #597: fleetd #593 — running_pid() counts only the daemon
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 1m39s
CI / build (push) Successful in 1m50s
running_pid() filters pgrep -f hits by `comm = java` (allowlist) instead of
denying a fixed list of shell names. Closes two false-positive holes: an
exited pid (empty comm matched no denied name) and any non-shell wrapper
(ssh, perl, python3, ruby) carrying the pattern in its own argv.

Gate: merged onto current main, bash suite exit 0 and structurally identical
to the baseline run on main (the one "Unattributable mutation" line is
pre-existing, confirmed by running the suite on origin/main). All four
running_pid tests confirmed defined AND invoked. Mutation check run by me:
neutering the allowlist makes the suite exit 1 with a named failure; restore
is byte-identical to baseline by git hash-object and green again.

Also checked and cleared: the new die-message advice `ps -eo pid,comm,args`
does NOT expose process environments on macOS — `-e` with `-o` selects all
processes, it does not imply `-E`. Verified with an isolated two-phase probe
and a positive control, after three earlier probes gave false positives by
self-matching (the grep's own argv, and the probe script's own text).
2026-09-19 10:37:50 +02:00
Dai Ha ae7845c375 fleetd #589 (Groups 1 & 2): pin 6 main() wiring sites with named factories
CI / shell-tests (pull_request) Successful in 9s
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Successful in 1m52s
Extracts 6 inline wiring expressions from Fleetd.main() into named,
directly-testable package-private static factories, following the
FleetdLoopHealthSourceWiringTest (#584) shape, and adds one wiring test
per factory:

Group 1 (exhaustion/quarantine):
- forwardingExhaustionSink(exhaustionSinkRef) — was inline
  ExhaustionSink.forwardingTo(exhaustionSinkRef::get)
- publishExhaustionSink(...) — was two untested statements building the
  real sink and .set()-ing it into exhaustionSinkRef
- liveExhaustedPatterns(config) — was inline
  new LiveExhaustedPatterns(() -> config.get().profiles())
- exhaustedPatternLookup(roster, liveExhaustedPatterns) — was an inline
  lambda resolving a herdr target to its profile's live pattern; silently
  losing this is the worst regression in the sweep, since a real
  usage-limit refusal would stop being classified as BACKEND_EXHAUSTED

Group 2 (CB-596 credential policy):
- claudeCodeLauncher(...) — was an inline `new ClaudeCodeLauncher(...)`
  whose memberCredentials supplier argument was untestable wiring
- openCodeLauncher(...) — same, for OpenCodeLauncher

Each new test pins its factory behaviorally (never via source-text
assertions): built and confirmed RED by name against the named inert
mutation, then confirmed GREEN again after restoring, and separately
confirmed GREEN after a behavior-preserving reformat/local-variable
extraction of the same call, to rule out a disguised source-text test.

Suite: 1789 -> 1799 tests (+10, matching the 10 tests added), 0
failures, mvn -o clean install BUILD SUCCESS.

Scope strictly limited to main()'s :214-:468 range per the ticket split
with the concurrent worker handling Group 3 at line 500+.
2026-09-19 15:35:32 +07:00
ltms 61115f6f61 Merge #598: fleetd #589 group 3 — wiring-test 5 sites in Fleetd.main()
CI / shell-tests (push) Successful in 6s
CI / contract (push) Successful in 50s
CI / build (push) Successful in 2m9s
Extracts 5 inline lambdas/method-refs in Fleetd.main() to package-private
factories and pins each with a wiring test. Production behaviour unchanged;
the releaseCleanup body moved verbatim.

Gate: test-merged onto current main in a scratch worktree, mvn exit=0,
1795 tests / 0 failures / 0 errors from 137 surefire reports. Diff read in
full. Worker's self-disclosed bare `git stash push` verified as recovered —
all 3 surviving stash entries predate today, so no other worktree lost work.
2026-09-19 10:33:07 +02:00
Dai Ha 4b9ebda1b3 fleetd #593 CORRECTION 1: allowlist comm=java, not a denylist of shells
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 1m27s
CI / build (pull_request) Successful in 1m48s
The round-1 fix excluded known shell names (sh/bash/zsh/dash/ksh) from
running_pid()'s pgrep candidates. Two holes remained, both the same
false-positive shape the ticket exists to remove:

1. A pid pgrep lists can exit before the following `ps -o comm=` lookup
   runs. On a gone pid, ps prints nothing, comm is empty, and an empty
   string matches no denied shell name -- so a dead pid was still counted.
2. The denylist only knows the shells someone thought to name. ssh, perl,
   python3, ruby, tail -- anything else carrying the pattern in its own
   argv -- was still counted alongside the real daemon. The ticket names
   ssh as a live route.

Both close with one change: allowlist comm=java instead of denying shells.
The daemon is always `java -jar target/fleetd.jar`, so its comm is always
`java`; an empty comm (hole 1) is not `java` either, closing that hole for
free.

Answers the objection in the code comment: an allowlist can under-count if
fleetd ever stops being launched by `java` (a native image, a renamed
launcher). That's a false negative, the worse direction for a guard -- but
it is not a new assumption: PATTERN='target/fleetd.jar' already assumes a
jar run by java, and that pattern breaks before this allowlist would.

Replaces the round-1 "real second process" test (which gave its exec -a
standin an argv[0] holding the pattern, but not comm=java) with one that
forces comm=java via `exec -a java sh -c '...'`. Adds two stubbed
pgrep/ps tests pinning the two holes directly (a non-java, non-shell comm
such as perl; an empty comm from an already-exited pid) -- deterministic
on every platform, unlike a live-process fixture, and immune to the BSD
vs Linux difference in how `comm` is derived from a fabricated process.
Adds a stubbed positive backstop (comm=java is counted).

Confirmed the regression is caught: reverted to the round-1 denylist,
reran the suite, watched the new non-shell-comm test fail at
`set -e`'s first failure, then isolated the exited-pid test separately
and confirmed it also fails against the same broken code. Restored the
fix and reran green.

Branch merged with origin/main (3 commits: hunter role + CLAUDE.md
addendum) before this commit; unrelated, no conflicts.
2026-09-19 15:31:15 +07:00
Dai Ha 1e68d7ee39 Merge origin/main into worker/593-1a8025-5 2026-09-19 15:27:46 +07:00
Dai Ha 6cccd458d4 #589 Group 3: wiring-test the 5 sites below line 500 in Fleetd.main()
CI / shell-tests (pull_request) Successful in 11s
CI / contract (pull_request) Successful in 1m26s
CI / build (pull_request) Successful in 2m11s
Extracts the inline lambdas/method references at the 5 assigned wiring
sites into named package-private factories on Fleetd, following the
FleetdLoopHealthSourceWiringTest pattern from #584:

- turnRegistrar(CompletionResolver) — was completion::register (Injector)
- healthFailTarget(MessageService) — was messages::abandon (FleetHealthMonitor)
- releaseCleanup(MessageService, ReplyInbox, PrimaryRegistry) — was the
  inline sessions.onRelease(detail -> {...}) cleanup lambda
- replyInboxOpener() — was AmqpReplyInbox::open passed to selectReplyInbox
- leadMailboxOpener() — was LeadMailbox::open passed to openLeadMailbox

Each factory has a new runtime test (not source-text) that drives real
collaborators through public APIs: MessageService.poll(ticket).phase(),
InMemoryReplyInbox.peek(), PrimaryRegistry.nudgeTargetFor(), and the
opener tests connect to a guaranteed-closed local port to prove a real
network attempt vs. an inert stub.

releaseCleanup was done first per the brief: MessageService.abandon's
javadoc documents that losing this cleanup leaves a torn-down worker's
rendezvous waiter open forever.

Tests: 1789 -> 1794 (+5), 0 failures, 0 errors. mvn -q -o test exit 0,
no BUILD FAILURE, no piped exit status. Each new test verified RED on
the inert form named in the ticket, and GREEN after reformatting the
call across lines and extracting the argument into a local/factory.
2026-09-19 15:24:11 +07:00
Dai Ha 42820fbe75 fleetd #593 (pid-count half): running_pid() no longer matches the caller
CI / shell-tests (pull_request) Successful in 9s
CI / contract (pull_request) Successful in 1m18s
CI / build (pull_request) Successful in 1m45s
running_pid() was a bare `pgrep -f "$PATTERN"`, which matches ANY process whose
full command line contains the pattern text -- including a shell that merely
embeds it as literal text (a hand-typed investigation, an ssh-shaped
`sh -c '...; ...'`, or a pipeline) rather than being the daemon. That self-match
turns a working redeploy into a reported "racing supervisor" failure via
assert_single_daemon.

pgrep -c does not exist on BSD/macOS, so this can't be fixed by switching flags.
running_pid() now keeps pgrep to find candidates (portable), then drops any
candidate whose process name (comm) names a shell -- the daemon is always
`java`, so a self-matching wrapper of this shape is always excluded while a
genuine second daemon-shaped process still counts.

assert_single_daemon's die message no longer hands the operator a bare
`pgrep -f "$PATTERN"` as remediation -- that was exactly the self-matching
invocation -- and now says in words that a pattern can match the caller.

Adds three tests: a self-matching wrapper shell must be excluded, a real
second daemon-shaped process must still be found, and the die message must
not recommend the self-matching command. Verified the first test fails
against the pre-fix implementation (confirmed the regression is caught).

Leaves instance 1 (the fleetd.out log source, systemd-only) for a Linux host,
per the ticket's scope split.
2026-09-19 15:18:32 +07:00
Dai Ha d91ff886da #568 follow-up: fix the text defects the hunter-role merge introduced
CI / shell-tests (push) Successful in 6s
CI / contract (push) Successful in 1m17s
CI / build (push) Successful in 1m48s
Found by reading the diff at the merge gate, not reported by the worker.

1. FleetConfig.java: the operator-facing "unknown key" hint read
   "'fleet.architects', 'fleet.developers' or 'fleet.hunters' or
   'fleet.reviewers'" — a double "or". This is text an operator reads at the
   moment their config is already wrong, so it should not itself be wrong.
2. MemberLifecycle.java: javadoc continuation asterisk indented 6 spaces, not 5.
3. MemberRegistry.java: javadoc asterisks moved from column 2 to column 4.
4. CallerResolver.java: a // comment indented one space past its block.

2-4 are the worker mangling alignment while widening enum lists to include
HUNTER. No behaviour changes.

CORRECTION to the #596 merge commit message. It claimed a fifth defect, "two
javadoc lines pushed past the 100-column convention". There is no such
convention in this repo: no checkstyle, no spotless, no .editorconfig, and 2975
of 32728 lines under fleetd/src/main/java already exceed 100 characters. I
asserted the rule before measuring it. Those two lines are untouched.

Verified: built in a scratch worktree, 1790 tests, 0 failures, 0 errors,
0 skipped, counted from the surefire XML.
2026-09-19 15:15:54 +07:00
ltms 386e760a5c Merge #596: fleetd #568 — add the hunter member role
CI / shell-tests (push) Successful in 5s
CI / contract (push) Successful in 49s
CI / build (push) Successful in 2m37s
Verified by the lead before merge, not taken on the worker's report:

- branch contains a639969; 1 ahead, 0 behind — clean fast-forward
- CLAUDE.md change is +2 lines in the Project addendum, NOT the canonical block
- canonical block sync check prints True on main and on this branch
- built in a scratch worktree (never mvn clean in the main clone): 1790 tests,
  0 failures, 0 errors, 0 skipped, counted from the surefire XML. Baseline 1789.
- live fleetd.yaml still loads: the hunters pool is optional

Five text defects found by reading the diff, not reported by the worker. They are
fixed in a follow-up commit on main rather than a round trip:
- FleetConfig.java operator-facing message reads "... 'fleet.developers' or
  'fleet.hunters' or 'fleet.reviewers'" — a double "or"
- misaligned javadoc continuation asterisks in MemberLifecycle and MemberRegistry
- a misaligned // comment in CallerResolver
- two javadoc lines pushed past the 100-column convention

Known gap, tracked separately: fleet.hunters is absent from the live config, so a
hunter cannot spawn on this host until the pool is added after the redeploy. The
role ships correct and inert.
2026-09-19 10:13:26 +02:00
Dai Ha 2e349139e9 #568: add hunter member role
CI / shell-tests (pull_request) Successful in 6s
CI / contract (pull_request) Successful in 1m21s
CI / build (pull_request) Successful in 1m44s
2026-09-19 15:09:07 +07:00
Dai Ha a639969a9a CLAUDE.md: a blocked lead consults architects, not the operator
CI / shell-tests (push) Successful in 6s
CI / contract (push) Successful in 1m32s
CI / build (push) Successful in 1m42s
The operator set this rule on 2026-09-19: when a decision blocks a lead, it
consults one or more architect members, who are authorized to agree on one
decision and unblock. The operator is not asked. Escalation stays open only
for things outside the fleet's authority -- money, credentials, or a promise
made to someone else.

The paragraph also carries the reason the ticket record is mandatory rather
than optional. The operator's old notification channel was the block itself:
work stopped, so they found out. Taking the operator out of the loop removes
that signal with it, so the decision goes on the ticket, which reaches them
whether or not they are at a terminal when it is made.

The rule has a second half aimed at architects, which lives in
fleet.charters.architect and is applied per daemon -- filed as #591, because
a charter can never reach a lead and charters do not travel between hosts.
The missing notification event is #592.

Block verified byte-identical with wiki 7-Use-Cases.md at 1d9bd1b.
2026-09-19 14:57:04 +07:00
Dai Ha 17c3a69c57 docs: CB-591 gateway page had two claims that went stale
CI / shell-tests (push) Successful in 5s
CI / contract (push) Successful in 1m23s
CI / build (push) Failing after 1m37s
The gateway's chat model is served under the stable alias `acoder`, and the
model behind that alias changed on 2026-08-28 — it is Qwen3.8-27B now, not
DeepSeek-V4-Flash. The old name is still served, so nothing broke, but it
names a model this is not.

Two claims on the page were wrong as a result, and both were written as
current facts rather than dated measurements:

- `/v1/models` returns exactly `["deepseek-v4-flash"]` — it returns 6 ids
  now. This sat under a heading saying it needs no re-testing.
- the status banner said `local` and `gx` are both at `weight: 100` —
  `local` is at 0.

Measured today against the live gateway: /v1/models returns acoder,
qwen3.8-27b-nvfp4, deepseek-v4-flash and three embedding names; a completion
sent as `deepseek-v4-flash` comes back reporting `"model": "acoder"`, which
is the alias in plain sight. /v1/deployment reports generation
2026-08-28-qwen3.8-27b-nvfp4.

§2 and §3 are left alone. They are the August plan, and rewriting them would
destroy the record of the migration.

fleetd.yaml moved to `acoder` in the same change. It is not tracked here.
2026-09-13 07:01:15 +07:00
ltms 49a5875586 Merge #583: fleetd #582 — assert pending message-id cleanup at every publish cleanup site
CI / shell-tests (push) Successful in 9s
CI / contract (push) Successful in 48s
CI / build (push) Failing after 2m1s
All ten assertions proven live: six by the implementer, the last four by the lead.
One contract build with four deleted removal lines produced exactly four named failures,
one per site, with the total unchanged at 1825.
2026-09-12 15:52:25 +02:00
Dai Ha 634d33b50b Merge worker/562-loop-health-wiring-test-99611c-5
CI / shell-tests (push) Successful in 9s
CI / contract (push) Successful in 1m17s
CI / build (push) Successful in 1m45s
2026-09-12 20:28:14 +07:00
Dai Ha 1db79bcaa9 Merge worker/581-completionresolver-cas-sites-0542b7-6 2026-09-12 20:28:14 +07:00
Dai Ha 4507bc5a70 Merge worker/571-attempted-outcome-5739f7-2 2026-09-12 20:28:14 +07:00
Dai Ha d7239ed23b fleetd #571: pin FleetMcp.formatReply's TIMED_OUT_UNCONFIRMED wording
CI / shell-tests (pull_request) Successful in 5s
CI / contract (pull_request) Successful in 1m14s
CI / build (pull_request) Successful in 1m42s
CORRECTION 5 on the ticket: mutating the new arm's message text to the
queued/working arm's text survived every existing test, because nothing
asserted the specific wording. This adds one test that asserts the
unconfirmed-delivery message and asserts it does NOT carry the
queued/working arm's retry invitation — the distinction #571 exists for.

No production code changes; formatReply's TIMED_OUT_UNCONFIRMED arm was
already correct.
2026-09-12 20:19:37 +07:00
Dai Ha dfeb9340b4 fleetd #581: cover completion CAS removals
CI / shell-tests (pull_request) Successful in 5s
CI / build (pull_request) Failing after 1m31s
CI / contract (pull_request) Successful in 1m35s
2026-09-12 20:19:20 +07:00
Dai Ha 1513d4f260 fleetd #562 follow-up: extract loopHealthSource factory, pin its wiring
CI / shell-tests (pull_request) Successful in 6s
CI / contract (pull_request) Successful in 59s
CI / build (pull_request) Successful in 1m43s
PR #579's inline `new FleetMcp.LoopHealthSource(poller::health, ...)` in
Fleetd.main had nothing a test could call directly. Measured: replacing
poller::health with a constant () -> RUNNING compiled clean and left all
1771 tests green (see issue #562 comment "HOLD on PR #579").

Extracts the inline construction to a package-private Fleetd.loopHealthSource
factory, the same style as the sibling capacitySource/healthCoverageSource
factories, and adds FleetdLoopHealthSourceWiringTest with three separate
assertions: the statusPoller half, the sessionReaper half, and the
reaper == null branch (still STOPPED).
2026-09-12 20:16:23 +07:00
Dai Ha 4ca7d72303 #582: assert pending message-id cleanup
CI / shell-tests (pull_request) Successful in 5s
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Successful in 2m32s
2026-09-12 20:15:48 +07:00
Dai Ha c0545d003d fleetd #571: make FleetApp.writeReply's inner Outcome switch exhaustive, no default
CI / shell-tests (pull_request) Successful in 8s
CI / contract (pull_request) Successful in 1m16s
CI / build (pull_request) Successful in 1m48s
Ticket comments (17126, 17127) corrected the original acceptance criterion after this
unit was already in flight: a hand-listed grep for the enum's constant names goes stale
silently the moment a new constant lands, so the compiler must be the enumeration instead.
sendOutcomeLabel (MessageService.java) and formatReply (FleetMcp.java) were already
default-free switch expressions. The one gap was writeReply's inner "status" switch, which
had `default -> "done"` — the exact value that would have lied about TIMED_OUT_UNCONFIRMED.
Remove the default and list every Outcome constant explicitly; REPLIED, COMPLETED_UNREPLIED,
QUESTION and STALE_TURN get an arm too even though the outer switch always dispatches them
first, so the inner switch stays exhaustive on its own. The outer switch (a statement, not
an expression) keeps its own default — Java does not require exhaustiveness there regardless,
and "everything not terminal is a 202" is an intentional catch-all.

Verified with the proof the ticket asked for: added a scratch 11th Outcome constant after
deleting all default arms and confirmed all three switch-expression sites (and no test file)
fail to compile without an arm for it, one at a time, then removed the scratch constant.
2026-09-12 19:28:59 +07:00
Dai Ha 275ac0d251 fleetd #562: surface loop health
CI / shell-tests (pull_request) Successful in 12s
CI / contract (pull_request) Successful in 1m31s
CI / build (pull_request) Successful in 2m11s
2026-09-12 19:26:37 +07:00
Dai Ha c1e06c9e12 fleetd #571: add TIMED_OUT_UNCONFIRMED so an ATTEMPTED delivery is not reported as never-arriving
MessageService.send's TimeoutException branch collapsed Injector.Cancellation.ATTEMPTED
(fleetd #551 — the send call was made but its outcome is unknown) into
Outcome.TIMED_OUT_QUEUED, which promises the caller the message will never arrive. On this
route agent.prompt may already have pasted and submitted the text, so a caller's natural
recovery (resend) risks a double delivery.

Add Outcome.TIMED_OUT_UNCONFIRMED and route ATTEMPTED to it. Update the three readers found
by searching for the enum's constant names (not `Outcome.`, which misses FleetMcp's
unqualified `case REPLIED ->` switches and would false-positive on ConfigRef's unrelated
Outcome record):
 - MessageService.sendOutcomeLabel: add it to the "timeout" metric label group.
 - FleetMcp.formatReply: its own case, warning against a blind retry (distinct from the
   generic "retry or poll status" message the other timeouts get).
 - FleetApp.writeReply: its own "unconfirmed" status and detail text, so it no longer falls
   through the switch's default -> "done" arm, which would have reported "the delegation
   completed" for the one case where delivery is unconfirmed.
2026-09-12 19:21:10 +07:00
Dai Ha 204da67d66 Merge #576: fleetd #575 — one finally covers answer()'s STALE_TURN exit
CI / shell-tests (push) Successful in 16s
CI / build (push) Successful in 1m42s
CI / contract (push) Successful in 1m57s
Third instance of the #572 shape on this file: one invariant kept at N sites,
asserted at fewer than N. Here answer()'s inner try opened AFTER the Task
lookup/registration and the STALE_TURN early return, so that return was covered
only by a hand-rolled copy of the finally's cleanup pair. The fix widens the try
upward and deletes the copy, so every exit runs the one finally exactly once.
rendezvous.open(workerSession) correctly stays outside it — nothing to clean if
it never opened.

This fixed no live leak, and the code comment says so: neither
rendezvous.answerAsk nor clearAsyncQuestion(turnId, false) can throw, so nothing
ever left through the old gap uncovered. It is a structure fix, the ticket's own
fallback case.

Verified here, not taken from the worker's report:
- baseline on the branch: 1766 tests, 0 failures, from Maven and from an
  independent sum over target/surefire-reports/*.txt.
- my own mutation, located fresh: move the inner `try {` back down below the
  STALE_TURN return — the exact pre-fix structure, minus the hand-rolled pair.
  Result: 1766 run, exactly 1 failure, and it is the new test —
  MessageServiceTest.answerLosingTheRaceToAnAlreadyAnsweredAskStillReturnsStaleTurnAndCleansUpOnce:545
  "the forward waiter this answer() call opened must be closed after a STALE_TURN
  return". One failure, not a crowd: the new test is the only thing holding this
  path.
- the worker's own mutation removed the whole finally and took 8 other tests with
  it. That proves the finally runs; it does not prove the STALE_TURN path reaches
  it. Mine does.

The new test hook answerAskLapseRaceHookForTest follows the file's existing
askTimeoutRaceHookForTest convention.

Closes #575.
2026-09-12 19:05:00 +07:00
Dai Ha b091c51eee fleetd #575: widen answer()'s try so one finally covers its STALE_TURN exit
CI / shell-tests (pull_request) Successful in 5s
CI / contract (pull_request) Successful in 58s
CI / build (pull_request) Successful in 2m30s
The waiter cleanup pair (asyncTasksByWaiter.remove + rendezvous.close) was
duplicated: two sites sit in a finally, the third was hand-rolled inline
before answer()'s early STALE_TURN return, structurally outside any finally.

Both rendezvous.answerAsk and clearAsyncQuestion(turnId, false) are total
(cannot throw), so the gap never leaked in practice. But the duplicate was
untested: mutating it away left all 1765 tests green, while the two
finally-protected sites are each killed by 8-36 tests. Same shape as #572.

Fix: widen the try to wrap the Task registration and the STALE_TURN check,
so the single finally covers every exit and the hand-rolled copy is gone.
Added a race hook + regression test that deterministically reproduces the
'ask lapsed between the lookup and the unblock' case and proves the fix
still returns STALE_TURN and cleans up exactly once.
2026-09-12 18:56:02 +07:00
ltms 84d631b030 Merge #572: answer()'s session-lock release is pinned on all four exits
CI / shell-tests (push) Successful in 9s
CI / contract (push) Successful in 1m27s
CI / build (push) Successful in 1m44s
fleetd #572. MessageService releases its per-session lock in a finally at two sites. Removing the
send() one failed 22 tests and errored 1. Removing the answer() one left the whole suite green.
Skip that unlock and the thread holds the session lock forever, so every later send or answer to
that session blocks permanently — no exception, no log line.

Tests only. Against current main the diff is one file, MessageServiceTest.java, +178 lines.

Four new tests, one per exit of answer(): normal REPLIED, TIMED_OUT_WORKING, the ExecutionException
rethrow, and the InterruptedException rethrow. Each proves REACQUISITION rather than the return
value — a bounded follow-up send on the SAME session must not come back BUSY, and send reports BUSY
only when tryLock itself timed out, so a non-BUSY probe is specifically evidence the lock was free.

Lead verification, re-running rather than accepting the worker's numbers, on the branch merged with
main at 3f8c38f:
  baseline   -> exit 0, 1765 tests from Maven and an independent sum over 131 reports (1761 + 4)
  mutation L -> BUILD FAILURE, exactly 4 failures, all four new tests by name

THE TRAP THE WORKER CAUGHT ITSELF, which is the most valuable thing in this PR. Their first draft of
the TIMED_OUT_WORKING test called answer() inline on the test thread and then probed with send() on
that same thread. It FALSE-PASSED under the mutation: lock is a ReentrantLock, so the same thread
reenters for free whether or not unlock() ran. A same-thread probe proves reentrancy, not release.
They found it only because they actually ran the mutation instead of trusting a green baseline, moved
answer() onto a background thread, and re-proved the kill. The pitfall is now documented in the test's
own javadoc, with the measurement that exposed it.

NOTE FOR ANYONE RE-RUNNING THE TICKET'S RECIPE: the ticket records `sed '1218s|...'`, measured before
#551 merged. #551 edited MessageService.java, and the line is now 1231. Locate it fresh — a
line-anchored sed against a stale number mutates the wrong line and reports a meaningless green. The
anchor count is the check that catches this: `grep -Fxc '            lock.unlock();'` must read 2
pristine and 1 after.

Shape reported by the worker, not investigated and not fixed: the same two-line cleanup
`asyncTasksByWaiter.remove(reply); rendezvous.close(...)` appears at three sites — send()'s finally,
answer()'s inner finally, and answer()'s inline STALE_TURN early return, which is NOT inside a
finally. Same maintained-at-N-sites shape as this finding. Filed separately.
2026-09-12 13:30:22 +02:00
ltms 3f8c38fc54 Merge #567: pin LeadMailbox.inspect's probe-channel close
CI / shell-tests (push) Successful in 8s
CI / contract (push) Successful in 1m1s
CI / build (push) Successful in 1m49s
fleetd #567. LeadMailbox.inspect opens a probe channel and closes it in a finally. The production
code was already correct; nothing asserted it, so a future refactor could drop the close and leak an
AMQP channel per inspect() call with the suite green.

Test only. LeadMailbox.java is untouched — sha256 a2cd99be77345b7e... before and after.

The test asserts the CONSEQUENCE rather than the return value: it caps the connection at three
channels (LeadMailbox uses two, consume and publish), runs a successful inspect, then requires a
replacement channel. If the probe is left open, the broker has no channel number left and
createChannel() returns null. A test that only checked inspect()'s MailboxState would pass under the
mutation, which is the whole reason this hole existed.

Lead verification, re-running rather than accepting the worker's numbers, on the branch merged with
main at ed2fd66:
  mvn -o clean install -Pcontract  ->  exit 0, 1792 tests from Maven and from an independent sum
                                       over 137 surefire reports.
  1792 = 1761 (main) + 30 (contract-only) + 1 (new), which also confirms the default-profile count
  is untouched: the new test is in an @Tag("contract") class.

My own mutation, located fresh rather than assuming the reported line number: delete
LeadMailbox.java:338 `probe.close();`, anchor by grep -Fxc 1 -> 0. Result under -Pcontract: exactly
1 failure, inspectClosesItsSuccessfulProbeChannel:226, "the replacement channel was null". Restored
to a2cd99be77345b7e..., git status --short empty.

THE PROFILE TRAP THIS TICKET EXISTS BECAUSE OF. An earlier sweep worker instrumented this line, saw
zero hits under the DEFAULT profile, and concluded "no test executes this line". That was false:
fleetd/pom.xml:264 sets excludedGroups=contract, so the covering tests were excluded from the run,
not absent. A surviving mutation has THREE causes — never executes, executes with nothing asserted,
or the covering tests were excluded from the profile — and only naming the profile tells them apart.
The true finding was covered-but-unasserted.

Known fragility, recorded rather than fixed: the test's channel cap of three assumes LeadMailbox
holds exactly two channels. If it ever holds more, this test fails loudly, which is fine. If it ever
holds fewer, a leak would no longer exhaust the cap and the test would go vacuous silently. Worth
re-checking if LeadMailbox's channel usage changes.

Not covered, and stated rather than faked: the defensive catch (RuntimeException) around the passive
declare, which needs a connection dying between createChannel() and the declare landing. The
method's own javadoc already admits that branch is unproven; the worker did not invent a test for it.
2026-09-12 13:26:30 +02:00
Dai Ha a4dbc8f8b7 fleetd #572: pin answer()'s session-lock release across all four exits
CI / shell-tests (pull_request) Successful in 7s
CI / contract (pull_request) Successful in 47s
CI / build (pull_request) Successful in 2m10s
MessageService.answer() releases its per-session lock in an outer finally
(MessageService.java:1218) that mutation testing showed was covered but
unasserted: removing that line left all 1750 existing tests green, because
every existing test on this path checks answer()'s return value, never that
the lock it took is actually reacquirable afterward. If it leaked, a session
would be wedged forever with no exception and no log line.

Adds four tests, one per exit of answer() (normal REPLIED reply,
TIMED_OUT_WORKING, ExecutionException rethrow, InterruptedException
rethrow), each proving the lock is reacquirable via a bounded (300ms)
follow-up send on the same session rather than merely checking answer()'s
own outcome. No production change.
2026-09-12 18:23:54 +07:00
ltms ed2fd6646a Merge #561: the completion/session listener fan-out survives either half throwing, and both sites are pinned
CI / shell-tests (push) Successful in 6s
CI / contract (push) Successful in 50s
CI / build (push) Successful in 2m7s
fleetd #561. Fleetd composed two TurnListener halves as `completion.X(); sessions.X();`, so a throw
from the first half skipped the second. The anonymous class is now a package-private factory,
Fleetd.turnListener(completion, sessions), built on two helpers that always attempt both halves and
rethrow whatever escaped — a second failure attached with addSuppressed rather than dropped, so it
still reaches StatusPoller's catch (Throwable).

The wiring at Fleetd.java:498 calls that factory, so the seam under test is the real caller.

Lead verification, on a merged tree, re-running the checks rather than accepting the worker's:
exit 0, 1761 tests from Maven and from an independent sum over 131 surefire reports.

The interesting part is what the first round MISSED. Two helpers maintain one invariant — "the
second half always runs" — and the first round's five tests asserted it at only one site. Measured:

  bothMustRun                      reverted to the pre-fix bug -> 1 named failure   (pinned)
  bothMustRunKeepingSecondResult   the SAME bug                -> 1755/1755 GREEN   (unpinned)
  failure.addSuppressed(t) deleted                             -> 1755/1755 GREEN   (unpinned)

Both survivors are now killed by new tests, re-verified by the lead after the fix:
sessionHalfStillRunsWhenTheCompletionHalfThrowsSynchronouslyForPostAction and
bothFailuresEscapeWhenBothHalvesThrowDistinctExceptions, each failing alone under its own mutation.

The rule this cost us, worth carrying: COUNT ASSERTIONS PER SITE, NOT PER INVARIANT. The total being
non-zero is what hides a zero at one site, and extracting a shared helper makes it worse rather than
better — it does not reduce the number of sites, only how many are visible. Credit to the fleet01
lead, who predicted this shape before an instance was found.

onDelivered stays deliberately unguarded. Its comment now gives the real reason —
CompletionResolver.captureBaseline already catches RuntimeException around its scrape and fails open,
so that half does not realistically throw — instead of the previous reason, which was true but about
registration rather than about this pair. A correct conclusion resting on a wrong premise reads
exactly like a verified one.
2026-09-12 13:20:56 +02:00
Dai Ha 5441a2b321 fleetd #567: assert inspect closes probe channel
CI / shell-tests (pull_request) Successful in 9s
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Successful in 2m23s
2026-09-12 18:19:43 +07:00
ltms 384867dfa3 Merge #551: ATTEMPTED is its own cancellation answer, and the javadoc stops claiming a timed-out send never arrived
CI / shell-tests (push) Successful in 6s
CI / contract (push) Successful in 1m11s
CI / build (push) Successful in 1m44s
fleetd #551. The injector polls a queued entry off the queue and marks it ATTEMPTED BEFORE the
irreversible AgentControl#send call, not after. So a Throwable escaping that call can never leave
the entry QUEUED at the head (the #546 re-send hazard) and can never be recorded as a confident
NOT_DELIVERED for text that may already be in the pane.

Cancellation.ATTEMPTED is added as a third answer. NOT_DELIVERED stays reserved for confirmed
absence: the readiness grace expiring, drop(), or a herdr *_not_found error, which the rest of this
codebase already reads as definitely-absent rather than inconclusive.

Also fixes three javadoc/comment sites that claimed a timed-out send definitely did not arrive:
TIMED_OUT_QUEUED, hasQueuedDelivery, the queuedDeliveries field, and the comment in send()'s timeout
branch. Two claims were wrong, not merely stale: "on every route it will not arrive later" is false
on the ATTEMPTED route, and "may already hold a partial paste" understates it — agent.prompt pastes
AND SUBMITS in one call, so the target may hold a complete, running turn.

Verified by the lead on a merged tree: mvn -o clean install from fleetd/, exit 0, 1754 tests from
Maven and from an independent sum over 130 surefire reports. Mutation: folding ATTEMPTED back into
NOT_DELIVERED in cancellationOf gives 3 red, each naming the property
(anErrorFromSendRemovesTheMessageAndMarksItAttempted:740,
aHerdrExceptionFromSendStillSurfacesButNowReportsAttempted:781,
aHerdrExceptionAfterThePasteIsNeverRecordedAsConfidentlyNotDelivered:806).

The final round is comment-only, proven mechanically rather than by reading: stripping every comment
from MessageService.java before and after and collapsing whitespace gives byte-identical code.
2026-09-12 13:19:20 +02:00
Dai Ha 034e17bb32 fleetd #561 follow-up: pin the session half of bothMustRunKeepingSecondResult
CI / shell-tests (pull_request) Successful in 7s
CI / contract (pull_request) Successful in 1m14s
CI / build (pull_request) Successful in 1m41s
Two helpers maintain one invariant (the second callback half always runs,
even when the first throws): bothMustRun and bothMustRunKeepingSecondResult.
Only bothMustRun's "session half still runs" direction was asserted
(sessionHalfStillRunsWhenTheCompletionHalfThrowsSynchronously, via
onTurnComplete). bothMustRunKeepingSecondResult — the helper
onTurnCompleteWithPostAction uses — could be reverted to the pre-#561
broken shape and the suite stayed green.

Adds two tests to FleetdTurnListenerCompositionTest:
- sessionHalfStillRunsWhenTheCompletionHalfThrowsSynchronouslyForPostAction:
  mirrors the existing onTurnComplete case for onTurnCompleteWithPostAction/
  bothMustRunKeepingSecondResult.
- bothFailuresEscapeWhenBothHalvesThrowDistinctExceptions: proves a second,
  distinct failure from the session half is preserved via addSuppressed
  rather than silently dropped when both halves of bothMustRun throw.

Also rewords the onDelivered comment in Fleetd.turnListener: it previously
said this pair is safe because registration survives a throw via #556's
Injector wiring, which is true but is not why THIS pair is unguarded.
CompletionResolver.captureBaseline already catches RuntimeException around
its scrape read and fails open, so completion.onDelivered does not
realistically throw. Comment text only, no logic change.
2026-09-12 18:10:31 +07:00
Dai Ha e20ccab1eb fleetd #561: harden the completion/session TurnListener fan-out
CI / shell-tests (pull_request) Successful in 8s
CI / contract (pull_request) Successful in 1m4s
CI / build (pull_request) Successful in 1m56s
Fleetd's turnListener composition had four callbacks (onTurnComplete,
onTurnCompleteWithPostAction, and both onTurnFailed overloads) built from two
bare, unguarded statements each. onDelivered's registration was already fixed
structurally by #556; these four had the identical fragility and were still
untested: nothing enforced that the completion resolver's half ran before the
session half beyond call order in the source, so a future reorder (or a
throwing session listener sequenced first) could silently skip the
completion resolver's effect and strand a caller for its full timeout.

Extracted the composition to a package-private static factory,
Fleetd.turnListener(completion, sessions), and hardened it with
bothMustRun/bothMustRunKeepingSecondResult: both callback halves are always
attempted regardless of whether the other throws, and whatever escapes is
rethrown afterward (never swallowed) so it still reaches StatusPoller's
catch (Throwable) and logs at ERROR.

FleetdTurnListenerCompositionTest builds this real composition from a real
CompletionResolver and a throwing fake sessions half, and asserts the
completion resolver's effect (the waiter resolving) survives the session
half throwing, for all four callbacks, plus a mirror case showing the
session half still runs when the completion half throws first.

onTurnCompleteWithPostAction keeps completion-before-session as a functional
requirement (resolveBeforePostAction must run before the context-reset
housekeeping can erase the pane), not just fault tolerance, so it is not
reorder-symmetric like the other three — documented in Fleetd.turnListener's
javadoc.
2026-09-12 17:47:53 +07:00
43 changed files with 4181 additions and 230 deletions
+23
View File
@@ -0,0 +1,23 @@
---
name: hunter
description: Sweep one assigned scope for defects and report ranked findings without changes.
---
<!-- CB-617: The model comes from fleetd.yaml because the launch flag overrides model here on both backends. -->
You sweep the assigned package or scope for real defects. Read the full assigned scope before you
judge it. Report several ranked findings when the evidence supports them. Change nothing: do not
edit code, commit, push, or open a pull request.
You may run the build or tests to check a finding. Read the complete output and report the real
result. Do not hide failures with a pipe. State only checks you actually ran. The primary's IDE
tools are not yours. A mounted forge tool may use a blocked credential and fail by design.
Do only the assigned scope. Note anything outside it in one line and do not investigate it further.
Use `fleet_ask{question}` only when a decision belongs to the lead, such as an unclear requirement
or two defensible fixes. Do not ask about something you can decide by reading more code.
Your handoff must name the files you read, each ranked finding or `NO FINDINGS`, the checks you ran,
and any caveat for review.
The launcher provides the required bridge reply instructions for every member.
+12 -1
View File
@@ -139,12 +139,21 @@ prefer `wait:false` + `fleet_poll` for anything non-trivial: a blocking `fleet_s
**Delegating does not delegate responsibility.** Workers open PRs; you are the gate. Never delegate
the merge — and merging on a reviewer's word is delegating it by proxy.
**When a decision blocks you, consult architects — not the operator.** Spawn one or more architect
members, give them the question and the evidence you have, and act on what they agree. They are
authorized to settle it, not only to advise. If two of them still disagree after two rounds, they
return both positions and you decide. Go to the operator only for something outside the fleet's
authority: money, credentials, or a promise made to someone else. **Then write the decision on the
ticket.** Taking the operator out of the loop also removes the signal they used to get, because
that signal was the block itself — work stopped, so they found out. A ticket comment replaces it,
and it reaches them whether or not they are at a terminal when you decide.
| Intent | Tool |
|---|---|
| Confirm your own role | `fleet_whoami` |
| See backends available | `fleet_profiles` |
| Start a member | `fleet_spawn{role?, profile?, cwd?, worktree?, ticket?, sessionName?, resumeSessionId?}` → `sessionId` + `paneId` |
| See the fleet | `fleet_list` → `leads` (your peers) + `members` (each carries `agentSessionId` when its backend knows one) · one peer's state: `fleet_status{sessionId}` |
| See the fleet | `fleet_list` → `leads` (your peers) + `members` (each carries `agentSessionId` when its backend knows one) + `loopHealth` (`RUNNING`, `STALLED`, or `STOPPED` for `statusPoller` and `sessionReaper`) · one peer's state: `fleet_status{sessionId}` |
| Delegate (blocking) | `fleet_send{sessionId, content}` |
| Delegate (long task) | `fleet_send{sessionId, content, wait:false}` → ticket → `fleet_poll{ticket}` |
| Answer a member's `fleet_ask` | `fleet_send{turnId, content}` — **not** `sessionId` |
@@ -275,6 +284,8 @@ must obey belongs in the charter, not here.
it for a multi-finding sweep hands the worker two contradictory output contracts. That has
already cost three workers' turns: each wrote a good report to its terminal and ended the turn
with no `fleet_reply`, and the scrape returned the tail of the brief instead.
Spawn `implementer` with role `dev`, `reviewer` with role `reviewer`, and `hunter` with role
`hunter`.
- **Primary-side skills** (not delegation playbooks — a worker cannot use them):
`port-to-opencode` (make an OpenCode session a participant in this workspace),
`fleets-status` (report every fleet that shares one LavinMQ instance),
+19 -1
View File
@@ -5,6 +5,21 @@
here took a revert and two upstream fixes — see §7.1, which is the useful part of this document. One
risk is **accepted rather than solved**: a stream cut by any mid-response timer arrives as HTTP 200
with no terminator, and our third-party members cannot detect it (§7.2).
> **Superseded in part — 2026-09-13.** Two claims on this page are no longer true of the live fleet.
> I measured both on this host today.
>
> 1. **The model is named `acoder` now, not `deepseek-v4-flash`.** `acoder` is a stable alias, and
> the model behind it changed on 2026-08-28: it is Qwen3.8-27B, not DeepSeek. The old name is
> still served, so nothing broke — the gateway answers it and reports `"model": "acoder"` in the
> reply, which is how you can see for yourself that it is an alias. `fleetd.yaml` moved to
> `acoder` on 2026-09-13. Do not guess behaviour from the name; ask the gateway's own manifest,
> `GET https://llm.ltms.dev/v1/deployment`, and read its `generation` field.
> 2. **`local` sits at `weight: 0`, not 100.** Only `gx` is auto-selected today.
>
> §2 and §3 below are the plan as written in August. They are the record of the migration, so they
> stay as they are. If this note stops matching `fleetd.yaml`, re-measure and rewrite the note.
· **Upstream:** [systems/vms wiki → LLM and MCP Gateway](https://git.ltms.dev/systems/vms/wiki/LLM-and-MCP-Gateway)
· **Upstream issue:** [systems/vms#31](https://git.ltms.dev/systems/vms/issues/31)
@@ -380,7 +395,10 @@ one turn; this one costs the whole task and is indistinguishable from a slow wor
- Token accepted on both surfaces. **Unauthenticated → 401**, so the Caddy proxy really does gate —
the wiki's "SecurityPolicy fails open" warning is about the gateway itself, not the edge.
- `/v1/models` returns exactly `["deepseek-v4-flash"]`, so trap 3 is clear.
- `/v1/models` returned exactly `["deepseek-v4-flash"]` **on 2026-08-15**, so trap 3 was clear
then. It returns 6 ids now — `acoder`, `qwen3.8-27b-nvfp4`, `deepseek-v4-flash` and three
embedding names — measured on this host 2026-09-13. The exact-name rule still holds; the
one-item list does not.
- **Reasoning survives both surfaces** — see §3b above.
- The launcher's generated opencode provider block is correct, carrying a real 48-character `llmk-`
key rather than the `fleetd-local-noauth` placeholder.
+12 -6
View File
@@ -560,7 +560,7 @@ placement: weighted
# older keys: `leaders:`, `members:`, `leadScan:` and `defaultProfile:`.
#
# A member is anything a lead spawns, and every member has two INDEPENDENT attributes:
# role — which contract: architect, dev or reviewer. It picks the launch charter, the role
# role — which contract: architect, dev, hunter or reviewer. It picks the launch charter, the role
# file, the playbook skill and the authz row.
# profile — which backend: one of the `profiles:` keys above (model, CLI adapter, cost).
# They vary on their own. A reviewer may run on the same profile as the dev whose diff it reads,
@@ -572,13 +572,13 @@ placement: weighted
#
# Each pool lists the profiles that role MAY run on — these are pools, not identities. That is also
# what replaced `defaultProfile:`: an unqualified spawn names a role, and that role's pool supplies
# the candidates, in definition order. A dev and a reviewer staying anonymous is exactly compatible
# with being listed here; the entry key just names the entry.
# the candidates, in definition order. A dev, hunter and reviewer staying anonymous is exactly
# compatible with being listed here; the entry key just names the entry.
fleet:
# Optional launch-charter text, keyed only by the singular role wire names: architect, dev,
# reviewer. Changes are HOT and reach the next spawn without a daemon restart. Do not put secrets
# here: a later launch step writes this text to a world-readable temp file, and ${ENV} interpolation
# is deliberately not supported.
# hunter, reviewer. Changes are HOT and reach the next spawn without a daemon restart. Do not put
# secrets here: a later launch step writes this text to a world-readable temp file, and ${ENV}
# interpolation is deliberately not supported.
charters:
architect: |-
You are an architect in this fleet. You refine work before anyone builds it:
@@ -589,6 +589,9 @@ fleet:
dev: |-
You implement the one unit you were given, and nothing else. You test it,
commit it, and open your own pull request. You never merge.
hunter: |-
You sweep the assigned scope for real defects. You may run the build or tests
to check a finding. You change nothing, and report several ranked findings.
reviewer: |-
You review the diff you were given. You report bugs, risks and missing tests.
You do not change code.
@@ -666,6 +669,9 @@ fleet:
developers:
gx10:
profile: gx10
# hunters:
# gx10:
# profile: gx10 # a hunt may run checks, but never changes code
# reviewers:
# gx10:
# profile: gx10 # the same backend may serve two roles; that is the point
+494 -82
View File
@@ -21,8 +21,10 @@ import dev.ltms.fleet.inject.ExhaustedPatternLookup;
import dev.ltms.fleet.inject.ExhaustionSink;
import dev.ltms.fleet.inject.LiveExhaustedPatterns;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.inject.LoopWatchdog;
import dev.ltms.fleet.inject.StatusPoller;
import dev.ltms.fleet.inject.TurnListener;
import dev.ltms.fleet.inject.TurnRegistrar;
import dev.ltms.fleet.inject.MemberPresence;
import dev.ltms.fleet.auth.MemberRegistry;
import dev.ltms.fleet.auth.CallerResolver;
@@ -79,6 +81,9 @@ import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.BiConsumer;
import java.util.function.BooleanSupplier;
import java.util.function.Consumer;
import java.util.function.Function;
import java.util.function.LongSupplier;
import java.util.function.Predicate;
@@ -216,23 +221,23 @@ public final class Fleetd {
// there is no 2-arg overload left for any lambda to silently bind to instead), but so a
// test can call the exact same object this line builds, instead of asserting a copy of its
// shape (round 3's lesson).
ExhaustionSink forwardingExhaustionSink = ExhaustionSink.forwardingTo(exhaustionSinkRef::get);
// fleetd #589 Group 1: extracted to forwardingExhaustionSink(...) below (see that method's
// javadoc) so a dedicated test can prove this factory keeps reading the reference live,
// rather than a rebuilt copy of its shape.
ExhaustionSink forwardingExhaustionSink = forwardingExhaustionSink(exhaustionSinkRef);
// The claude-code adapter is the always-present default; keep it even with no profiles (so a
// bridge configured with no workers, or opencode-only, still has a well-defined base adapter)
// unless opencode is the only kind configured.
// fleetd #589 Group 2: extracted to claudeCodeLauncher(...)/openCodeLauncher(...) below (see
// those methods' javadoc) so a dedicated test can prove the CB-596 memberCredentials policy
// supplier is actually wired to each adapter, not silently replaced with `() -> null`.
if (!claudeProfiles.isEmpty() || opencodeProfiles.isEmpty()) {
adapters.add(new ClaudeCodeLauncher(router.memberAgents(), router.memberSpaces(), guard,
claudeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
() -> config.get().fleet(),
() -> config.get().memberCredentials(), null, config::get));
adapters.add(claudeCodeLauncher(router.memberAgents(), router.memberSpaces(), guard,
claudeProfiles, cfg, config));
}
if (!opencodeProfiles.isEmpty()) {
adapters.add(new OpenCodeLauncher(router.memberAgents(), router.memberSpaces(),
opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
() -> config.get().fleet(),
() -> config.get().memberCredentials(), config::get, forwardingExhaustionSink));
adapters.add(openCodeLauncher(router.memberAgents(), router.memberSpaces(),
opencodeProfiles, cfg, config, forwardingExhaustionSink));
}
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
// CB-578 stage B: one quarantine tracker for the whole daemon, shared between the launcher
@@ -406,12 +411,12 @@ public final class Fleetd {
// LiveExhaustedPatterns's class doc for why this replaces the old compiled-once-at-startup
// map. A profile with no exhaustedPattern simply returns null here, so its workers keep
// today's completion-fallback behaviour unchanged.
LiveExhaustedPatterns liveExhaustedPatterns = new LiveExhaustedPatterns(() -> config.get().profiles());
ExhaustedPatternLookup exhaustedPatterns = target -> sessions.roster().stream()
.filter(session -> target.equals(session.terminalId()))
.findFirst()
.map(session -> liveExhaustedPatterns.patternFor(session.profile()))
.orElse(null);
// fleetd #589 Group 1: both extracted to liveExhaustedPatterns(...)/
// exhaustedPatternLookup(...) below (see those methods' javadoc) — this is the worst
// consequence in the whole #589 sweep: silently losing either wiring means a genuine
// usage-limit refusal is handed back as a real completion instead of BACKEND_EXHAUSTED.
LiveExhaustedPatterns liveExhaustedPatterns = liveExhaustedPatterns(config);
ExhaustedPatternLookup exhaustedPatterns = exhaustedPatternLookup(sessions::roster, liveExhaustedPatterns);
// The startup coverage line still reports the boot-time snapshot only — it is printed once,
// here, and a reload no longer needs to change what it said; exhaustionDetectionArmed (via
// liveExhaustedPatterns.armed, wired into quarantineSource below) is what stays live.
@@ -459,11 +464,13 @@ public final class Fleetd {
// method's javadoc for the full fleetd #175/#234/#446 history this used to carry inline —
// so a dedicated test can drive the exact ExhaustionSink main() builds, not a hand-rebuilt
// copy of its shape.
ExhaustionSink exhaustionSink = exhaustionSink(sessions, config, quarantine,
quarantineReasonByCredential, cfg);
// fleetd #175: point the forwarding sink handed to OpenCodeLauncher above at the real one,
// now that `sessions` exists to resolve target -> session -> profile.
exhaustionSinkRef.set(exhaustionSink);
// fleetd #589 Group 1: both statements (build + set) folded into publishExhaustionSink(...)
// below (see that method's javadoc), so a test can prove the reference is actually
// repointed at the real sink, not silently left at ExhaustionSink.none().
ExhaustionSink exhaustionSink = publishExhaustionSink(exhaustionSinkRef, sessions, config,
quarantine, quarantineReasonByCredential, cfg);
// fleetd #201 Unit 5: the production BackendErrorSink needs `pushLoop` (built further below,
// after `sessions`) to tell a lead about an incident or an unmapped target — the same
// construction-order cycle `exhaustionSinkRef` breaks above, broken the same way: a mutable
@@ -491,61 +498,34 @@ public final class Fleetd {
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
MemberPresence presence = sessions.asPresence();
TurnListener turnListener = new TurnListener() {
@Override
public void onTurnComplete(String target) {
completion.onTurnComplete(target);
sessions.onTurnComplete(target);
}
@Override
public boolean hasPostTurnAction(String target) {
return sessions.hasPostTurnAction(target);
}
@Override
public boolean onTurnCompleteWithPostAction(String target) {
completion.resolveBeforePostAction(target);
return sessions.onTurnCompleteWithPostAction(target);
}
@Override
public void onDelivered(String target, dev.ltms.fleet.msg.TurnToken token) {
completion.onDelivered(target, token);
sessions.onDelivered(target, token);
}
@Override
public void onTurnFailed(String target) {
completion.onTurnFailed(target);
sessions.onTurnFailed(target);
}
@Override
public void onTurnFailed(String target, String reason) {
completion.onTurnFailed(target, reason);
sessions.onTurnFailed(target);
}
};
// fleetd #561: extracted to a static factory (see turnListener below) — the anonymous class
// this replaced had two bare, unguarded statements per callback, and nothing enforced that
// the completion half went first beyond call order in the source.
TurnListener turnListener = turnListener(completion, sessions);
Predicate<String> deliverable = deliverableTo(presence, leads);
// fleetd #556: registration is wired directly to `completion`, not folded into the
// `turnListener` fan-out above — so it survives `sessions.onDelivered` (or any future
// listener) throwing, regardless of call order. See TurnRegistrar's javadoc.
// fleetd #589 Group 3 (:505): extracted to turnRegistrar(...) — see FleetdTurnRegistrarWiringTest.
Injector injector = new Injector(router, turnListener, deliverable,
presence::forget, completion::register);
presence::forget, turnRegistrar(completion));
StatusPoller poller = new StatusPoller(router, injector, Injector.POLL_INTERVAL_MILLIS);
poller.start();
// CB-307: reply inbox. A broker: block selects the AMQP-backed durable adapter; absent (or
// unusable), fleetd stays soft-state on the in-memory inbox. The AMQP inbox owns a broker
// connection, so keep the reference to close it in the ordered shutdown hook.
final ReplyInbox replyInbox = selectReplyInbox(cfg.broker(), System.getenv(), AmqpReplyInbox::open);
// fleetd #589 Group 3 (:512): opener extracted to replyInboxOpener() — see
// FleetdReplyInboxOpenerWiringTest.
final ReplyInbox replyInbox = selectReplyInbox(cfg.broker(), System.getenv(), replyInboxOpener());
// CB-637: this daemon's lead-to-lead mailbox on the SHARED coordination vhost — a separate
// broker from the reply inbox by design (see FleetConfig.Coordinator). Absent a coordinator:
// block this is null and every lead path below is simply not wired, which is exactly the
// behaviour before this ticket. It owns a broker connection, so keep the reference for the
// ordered shutdown hook.
final LeadMailbox leadMailbox = openLeadMailbox(cfg.coordinator(), System.getenv(), LeadMailbox::open);
// fleetd #589 Group 3 (:518): opener extracted to leadMailboxOpener() — see
// FleetdLeadMailboxOpenerWiringTest.
final LeadMailbox leadMailbox = openLeadMailbox(cfg.coordinator(), System.getenv(), leadMailboxOpener());
// CB-307: learn the primary's terminal from orchestration tool calls (or pin from config).
// The pin also feeds CallerResolver below: a primary running inside a herdr pane would
// otherwise resolve as a worker and be refused every orchestration tool.
@@ -615,10 +595,12 @@ public final class Fleetd {
// 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.
// fleetd #589 Group 3 (:591): failTarget extracted to healthFailTarget(...) — see
// FleetdHealthFailTargetWiringTest.
healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler,
System::nanoTime, () -> TimeUnit.MILLISECONDS.toNanos(System.currentTimeMillis()),
cfg.health().intervalOrDefault(),
cfg.health().workingSuspectAfterOrDefault(), messages::abandon);
cfg.health().workingSuspectAfterOrDefault(), healthFailTarget(messages));
String coverage = FleetHealthMonitor.coverage(true,
cfg.health().notifications() != null && cfg.health().notifications().configured());
if ("detection-only".equals(coverage)) {
@@ -638,27 +620,11 @@ public final class Fleetd {
// CB-516: releasing a worker must fail whatever send was waiting on it. Without this a
// torn-down delegation kept reporting PENDING until the 30-minute async timeout, and never
// reached /metrics — the delegation was unresolvable and nothing said so.
sessions.onRelease(detail -> {
// CB-578 stage C, acceptance criterion 10: a failed ticket's detail should tell a lead
// where to re-dispatch onto the same tree, not just that the worker vanished.
String reason = "the worker session was released before it replied";
if (detail.worktreePath() != null) {
reason += "; worktree=" + detail.worktreePath() + " branch=" + detail.branch()
+ " snapshot=" + (detail.snapshotRef() != null ? detail.snapshotRef() : "none");
}
// CB-584 (issue #65 criterion 5): also name the agent session, so a lead can resume the
// member's conversation instead of only re-dispatching a fresh one onto the same files.
if (detail.agentSessionId() != null) {
reason += " agentSessionId=" + detail.agentSessionId();
}
// fleetd #275: this is an explicit teardown (fleet_stop, or the idle reaper) — the
// worker's pane is being stopped right now, so an open fleet_ask has no turn left to
// resume into. Sweep it too, unlike FleetHealthMonitor's health-classification call
// (see MessageService.abandon's javadoc for why those two must differ).
messages.abandon(detail.terminalId(), reason, true);
replyInbox.release(detail.terminalId());
primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding
});
// fleetd #589 Group 3 (:611-631): the whole cleanup lambda extracted to releaseCleanup(...)
// — see FleetdReleaseCleanupWiringTest, and MessageService.abandon's javadoc for the
// documented incident (a torn-down worker's rendezvous waiter left open) this lambda exists
// to prevent.
sessions.onRelease(releaseCleanup(messages, replyInbox, primaryRegistry));
// MCP server face (CB-105): fleet_send/fleet_reply/fleet_status, mounted at /mcp.
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
@@ -696,14 +662,20 @@ public final class Fleetd {
return configured == null ? null : configured.effectiveCredentialId();
}, outagePolicy);
FleetMcp.LoopHealthSource loopHealth = loopHealthSource(poller, reaper);
FleetMcp mcp = new FleetMcp(messages, workers, sessions, identity, presence,
primaryRegistry, callers, FleetMcp.AuthorizationMode.ENFORCED, metrics,
capacitySource(config, cfg, profile -> liveCountRef.get().apply(profile)),
healthCoverageSource(config),
loopHealth,
quarantineSource,
leadMailbox,
outageSource,
new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)),
// fleetd #602 gauge-wiring: threads each lead's configured configDir into the
// context gauge — see leadConfigDirSource's own doc for why this, not a hardcoded
// null, is what fleet_list's context row now reads.
leadConfigDirSource(() -> config.get().profiles(), leaders),
// fleetd #361: the operator-declared peers this daemon's fleet_list should try to
// reach. Read from the SAME snapshot leadMailbox itself opened from (cfg.coordinator()),
// not the live config.get() — coordinator wiring is already boot-time-fixed (see
@@ -791,7 +763,7 @@ public final class Fleetd {
Javalin app = new FleetApp(herdr, memberHerdr, workers, sessions, messages, presence, mcp.servlet(),
callers, metrics, deliverable,
() -> MemberCredentialPolicyView.of(config.get().memberCredentials()),
quarantineSource, outageSource).build();
quarantineSource, outageSource, loopHealth).build();
app.start(cfg.bind().host(), cfg.bind().port());
log.info("fleetd listening on {}:{}, herdr socket {}",
cfg.bind().host(), cfg.bind().port(), socket);
@@ -1037,6 +1009,120 @@ public final class Fleetd {
cfg.profiles()::keySet, System::nanoTime);
}
/**
* fleetd #589 Group 1: the forwarding {@link ExhaustionSink} handed to the adapters built
* before {@code sessions} exists (see the {@code exhaustionSinkRef}/{@code
* forwardingExhaustionSink} locals in {@code main}, just above {@link #capacitySource}'s call
* site). Before this ticket, {@code ExhaustionSink.forwardingTo(exhaustionSinkRef::get)} was
* built inline — nothing a test could call directly, so a mutation swapping the supplier for a
* hardcoded {@code () -> ExhaustionSink.none()} compiled clean and left the suite green: the
* forwarder would silently stop reading the reference at all, and {@link
* #publishExhaustionSink} repointing that reference later would have no effect.
*
* <p>Extracted the same way {@link #capacitySource}/{@link #loopHealthSource} were, so {@code
* FleetdExhaustionSinkForwardingWiringTest} can call this factory directly with a real {@link
* AtomicReference}, mutate the reference AFTER the forwarder is built, and prove the forwarder
* still reads it live rather than a fixed target captured at construction time.
*/
static ExhaustionSink forwardingExhaustionSink(AtomicReference<ExhaustionSink> exhaustionSinkRef) {
return ExhaustionSink.forwardingTo(exhaustionSinkRef::get);
}
/**
* fleetd #589 Group 1: publish the real {@link ExhaustionSink} — built the same way {@link
* #exhaustionSink} always was — into the forwarding reference {@link #forwardingExhaustionSink}
* built above, replacing {@code main}'s previously untested two-statement sequence ({@code
* ExhaustionSink exhaustionSink = exhaustionSink(...); exhaustionSinkRef.set(exhaustionSink);}).
* Before this ticket, nothing proved the {@code .set(...)} call actually received the real sink
* rather than a hardcoded {@code ExhaustionSink.none()} — the whole point of {@code
* exhaustionSinkRef} existing (fleetd #175) is that {@link
* dev.ltms.fleet.member.OpenCodeLauncher}'s model-mismatch check, built before {@code sessions}
* exists, keeps working once this line runs; silently keeping the reference at {@code none()}
* would mean that check permanently does nothing, with the full suite still green because no
* existing test drives this exact call site.
*
* <p>Returns the built sink so {@code main} can still pass it to {@link CompletionResolver}'s
* constructor at the same call site it already does, without building it twice.
*/
static ExhaustionSink publishExhaustionSink(AtomicReference<ExhaustionSink> exhaustionSinkRef,
SessionManager sessions, ConfigRef config, BackendQuarantine quarantine,
Map<String, String> quarantineReasonByCredential, FleetConfig cfg) {
ExhaustionSink sink = exhaustionSink(sessions, config, quarantine, quarantineReasonByCredential, cfg);
exhaustionSinkRef.set(sink);
return sink;
}
/**
* fleetd #589 Group 1: the live {@code exhaustedPattern} source (fleetd #446) {@link
* CompletionResolver} enforces on, extracted out of {@code main} for the same reason {@link
* #capacitySource} was. Before this ticket {@code new LiveExhaustedPatterns(() ->
* config.get().profiles())} was built inline; replacing the supplier with a hardcoded {@code ()
* -> Map.of()} compiled clean and left the suite green, meaning every profile's {@code
* exhaustedPattern} would silently stop being recognised and a genuine usage-limit refusal
* would be handed back as a real completion instead of {@code BACKEND_EXHAUSTED}.
*/
static LiveExhaustedPatterns liveExhaustedPatterns(ConfigRef config) {
return new LiveExhaustedPatterns(() -> config.get().profiles());
}
/**
* fleetd #589 Group 1: the {@link ExhaustedPatternLookup} {@link CompletionResolver} enforces
* on, resolving a herdr {@code target} to its session's profile and then to that profile's live
* {@link LiveExhaustedPatterns#patternFor}. Extracted out of {@code main} the same way {@link
* #worktreeBranchLookup} was — same {@code Supplier<List<MemberSession>>} roster shape, same
* reason: before this ticket the lambda was built inline, and replacing it with {@code target ->
* null} (the exact shape of {@link ExhaustedPatternLookup#none()}) compiled clean and left the
* suite green. This is the worst consequence in the whole #589 sweep (see the ticket): a
* genuine usage-limit refusal would stop being classified as {@code BACKEND_EXHAUSTED} and
* would be handed back to a waiting {@code fleet_send} as if it were real completed work.
*
* @param roster the live member roster, normally {@code sessions::roster}
*/
static ExhaustedPatternLookup exhaustedPatternLookup(Supplier<List<MemberSession>> roster,
LiveExhaustedPatterns liveExhaustedPatterns) {
return target -> roster.get().stream()
.filter(session -> target.equals(session.terminalId()))
.findFirst()
.map(session -> liveExhaustedPatterns.patternFor(session.profile()))
.orElse(null);
}
/**
* fleetd #589 Group 2: the production {@link ClaudeCodeLauncher} adapter, extracted out of
* {@code main} the same way {@link #capacitySource} was. Before this ticket the constructor
* call (11 arguments, including the CB-596 {@code memberCredentials} policy supplier) was built
* inline; replacing the {@code () -> config.get().memberCredentials()} argument with {@code ()
* -> null} compiled clean and left the suite green — {@code memberCredentials} is not {@code
* null} itself (a lambda is never {@code null}), so {@link
* dev.ltms.fleet.member.HerdrPeerLauncher#applyMemberCredentialPolicy} sees {@code
* memberCredentials.get() == null} and silently shadows nothing, reopening the exact CB-592
* exposure gap CB-596's policy closed. {@code FleetdClaudeCodeLauncherCredentialWiringTest}
* calls this factory with a real {@link ConfigRef} carrying a {@code memberCredentials:} block
* and proves a known-but-not-allowed name is actually shadowed on {@code spawn()}.
*/
static ClaudeCodeLauncher claudeCodeLauncher(AgentControl agents, WorkspaceControl spaces,
SubscriptionGuard guard, Map<String, FleetConfig.Profile> claudeProfiles, FleetConfig cfg,
ConfigRef config) {
return new ClaudeCodeLauncher(agents, spaces, guard, claudeProfiles, cfg.effectiveDefaultProfile(),
System::getenv, cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
() -> config.get().fleet(), () -> config.get().memberCredentials(), null, config::get);
}
/**
* fleetd #589 Group 2: the production {@link OpenCodeLauncher} adapter, the {@code opencode}
* counterpart to {@link #claudeCodeLauncher} above and extracted for the identical reason: the
* same {@code () -> config.get().memberCredentials()} argument, reopening the same CB-592
* exposure gap if silently replaced with {@code () -> null}.
*/
static OpenCodeLauncher openCodeLauncher(AgentControl agents, WorkspaceControl spaces,
Map<String, FleetConfig.Profile> opencodeProfiles, FleetConfig cfg, ConfigRef config,
ExhaustionSink forwardingExhaustionSink) {
return new OpenCodeLauncher(agents, spaces, opencodeProfiles, cfg.effectiveDefaultProfile(),
System::getenv, cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
() -> config.get().fleet(), () -> config.get().memberCredentials(), config::get,
forwardingExhaustionSink);
}
/**
* fleetd #426: package-private factory for {@code fleet_list}'s {@code healthCoverage} source,
* extracted out of {@code main} for the same reason {@link #capacitySource} and {@link
@@ -1069,6 +1155,129 @@ public final class Fleetd {
});
}
/**
* fleetd #562 follow-up: package-private factory for {@code fleet_list}'s and {@code
* /healthz}'s {@code loopHealth} source, extracted out of {@code main} for the same reason
* {@link #capacitySource} and {@link #healthCoverageSource} were. Before this ticket the
* {@link FleetMcp.LoopHealthSource} was built inline with a bare {@code new}, so there was
* nothing a test could call directly — measured: replacing {@code poller::health} with a
* constant {@code () -> LoopWatchdog.State.RUNNING} at the call site compiled clean and left
* the full suite green, meaning the daemon could report the {@link StatusPoller} as always
* {@code RUNNING} even while it was actually stalled. That is a false negative on the exact
* signal this ticket exists to surface, and is the mirror of a false positive muting a real
* monitoring component — worse, because there is no noise for anyone to notice and then
* silence. {@link FleetdLoopHealthSourceWiringTest} calls this factory directly and pins both
* halves separately, plus the {@code reaper == null} branch below.
*
* <p>{@code reaper} may be {@code null} — a {@link SessionReaper} is only constructed when
* {@code lifecycle.idleTtlSeconds} is configured (see the {@code reaper} local above) — and
* this factory preserves the existing behaviour of reporting {@link LoopWatchdog.State#STOPPED}
* in that case, rather than a {@code NullPointerException} on the first {@code fleet_list} or
* {@code /healthz} call.
*/
static FleetMcp.LoopHealthSource loopHealthSource(StatusPoller poller, SessionReaper reaper) {
return new FleetMcp.LoopHealthSource(poller::health,
() -> reaper == null ? LoopWatchdog.State.STOPPED : reaper.health());
}
/**
* fleetd #589 Group 3 (site {@code :505}): package-private factory for the {@link Injector}'s
* {@link TurnRegistrar}, extracted out of {@code main} for the same reason {@link
* #loopHealthSource} was — before this ticket {@code completion::register} was an inline
* argument to {@code new Injector(...)}, so nothing could pin it directly. Replacing it with
* {@link TurnRegistrar#NOOP} compiles clean and leaves every existing test green: {@code
* onDelivered}'s own {@code captureBaseline} does the identical {@code inFlight} check-and-put a
* moment later on the ordinary path, so the two are indistinguishable once a turn's delivery
* finishes normally. The gap {@link TurnRegistrar}'s own javadoc (fleetd #556) exists to close is
* a {@code turnListener} callback throwing between the two — {@link
* FleetdTurnRegistrarWiringTest} pins that {@code register} itself (not just {@code
* captureBaseline}) makes a delivered turn's waiter resolvable.
*/
static TurnRegistrar turnRegistrar(CompletionResolver completion) {
return completion::register;
}
/**
* fleetd #589 Group 3 (site {@code :591}): package-private factory for {@link
* FleetHealthMonitor}'s {@code failTarget} callback, extracted out of {@code main} for the same
* reason {@link #loopHealthSource} was. Before this ticket {@code messages::abandon} was an
* inline argument to {@code new FleetHealthMonitor(...)}; replacing it with a no-op {@code
* BiConsumer} compiles clean and leaves every existing test green, and in production it means a
* member found {@code GONE}/{@code NEVER_READY} never fails the ticket waiting on it — the
* caller reports {@code PENDING} for the full 30-minute async timeout instead of the immediate,
* accurate failure CB-580 exists to give it. {@link FleetdHealthFailTargetWiringTest} pins that
* the returned callback actually reaches the real {@link MessageService#abandon}.
*/
static BiConsumer<String, String> healthFailTarget(MessageService messages) {
return messages::abandon;
}
/**
* fleetd #589 Group 3 (site {@code :611-631}): package-private factory for the whole {@link
* SessionManager#onRelease} cleanup callback, extracted out of {@code main} for the same reason
* {@link #loopHealthSource} was. Before this ticket this was an inline lambda built directly
* inside {@code main}; replacing its body with a no-op {@code detail -> { }} compiles clean and
* leaves every existing test green, and in production it is the exact incident {@link
* MessageService#abandon}'s own javadoc documents: a torn-down worker's rendezvous waiter is
* left open, so a blocking {@code fleet_send} keeps blocking and an async one reports {@code
* PENDING} for a hardcoded thirty minutes on every {@code fleet_stop} and every idle-reap.
*
* <p>{@link FleetdReleaseCleanupWiringTest} pins all three collaborator calls this lambda makes
* — {@code messages.abandon}, {@code replyInbox.release}, and {@code
* primaryRegistry.forgetDelegation} — each already tested on its own ({@code MessageServiceTest},
* {@code PrimaryRegistryTest}), but never before proven to actually be reached from here.
*/
static Consumer<SessionManager.ReleaseDetail> releaseCleanup(MessageService messages, ReplyInbox replyInbox,
PrimaryRegistry primaryRegistry) {
return detail -> {
// CB-578 stage C, acceptance criterion 10: a failed ticket's detail should tell a lead
// where to re-dispatch onto the same tree, not just that the worker vanished.
String reason = "the worker session was released before it replied";
if (detail.worktreePath() != null) {
reason += "; worktree=" + detail.worktreePath() + " branch=" + detail.branch()
+ " snapshot=" + (detail.snapshotRef() != null ? detail.snapshotRef() : "none");
}
// CB-584 (issue #65 criterion 5): also name the agent session, so a lead can resume the
// member's conversation instead of only re-dispatching a fresh one onto the same files.
if (detail.agentSessionId() != null) {
reason += " agentSessionId=" + detail.agentSessionId();
}
// fleetd #275: this is an explicit teardown (fleet_stop, or the idle reaper) — the
// worker's pane is being stopped right now, so an open fleet_ask has no turn left to
// resume into. Sweep it too, unlike FleetHealthMonitor's health-classification call
// (see MessageService.abandon's javadoc for why those two must differ).
messages.abandon(detail.terminalId(), reason, true);
replyInbox.release(detail.terminalId());
primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding
};
}
/**
* fleetd #589 Group 3 (site {@code :512}): package-private factory for {@code selectReplyInbox}'s
* production {@link AmqpOpener}, extracted out of {@code main} for the same reason {@link
* #loopHealthSource} was. Before this ticket {@code AmqpReplyInbox::open} was an inline argument
* to the {@code selectReplyInbox(...)} call; replacing it with {@code (uri, prefetch) -> new
* InMemoryReplyInbox()} compiles clean and leaves every existing test green — {@link
* FleetdReplyInboxSelectionTest} drives {@code selectReplyInbox} with its own injected opener and
* never sees what {@code main} actually passes. {@link FleetdReplyInboxOpenerWiringTest} pins
* that this returns the real opener by pointing it at a guaranteed-closed local port and
* asserting the real network attempt throws — the inert stub never attempts a connection at all.
*/
static AmqpOpener replyInboxOpener() {
return AmqpReplyInbox::open;
}
/**
* fleetd #589 Group 3 (site {@code :518}): package-private factory for {@code
* openLeadMailbox}'s production {@link LeadMailboxOpener}, extracted out of {@code main} for the
* same reason {@link #replyInboxOpener} was — same gap, same fix, the lead-coordination mailbox
* instead of the reply inbox. {@link FleetdLeadMailboxOpenerWiringTest} pins that this returns
* the real opener the same way.
*/
static LeadMailboxOpener leadMailboxOpener() {
return LeadMailbox::open;
}
/**
* 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).
@@ -1096,6 +1305,157 @@ public final class Fleetd {
.orElse(null);
}
/**
* fleetd #561: compose the production {@link TurnListener} from its two halves — the
* completion resolver (which resolves a blocked {@code fleet_send}'s waiter) and the session
* manager (which drives the member's lifecycle state) — hardened so a throw from either
* half's callback can never suppress the other half's callback for the same event.
*
* <p>Before this, each method below was two bare, unguarded statements: whichever ran first
* throwing meant the second one never ran at all, and nothing beyond call order in the
* source enforced "the completion half goes first". #556 fixed the identical shape for {@code
* onDelivered}'s registration by moving it off the fan-out entirely (see {@code
* TurnRegistrar}); this fixes the four remaining callbacks — {@code onTurnComplete}, {@code
* onTurnCompleteWithPostAction}, and both {@code onTurnFailed} overloads — by hardening the
* fan-out itself instead, since none of them can be pulled out of the listener the way
* registration was.
*
* <p>The invariant this composition guarantees: <b>a throwing session-listener must not
* prevent the completion resolver from being told the turn ended.</b> The completion half is
* always attempted first, and — as a bonus the resolver does not depend on — its own throw
* does not stop the session half from running either. Whatever escapes (from one half or
* both) is rethrown once both have been attempted, with a second failure recorded via {@link
* Throwable#addSuppressed} on the first, so it still reaches {@link StatusPoller}'s {@code
* catch (Throwable)} and logs at ERROR. Nothing here swallows a failure to make the two halves
* look safe.
*
* <p>{@code onTurnCompleteWithPostAction} is the one place order is not just fault-tolerance
* but a functional requirement: {@code completion.resolveBeforePostAction} must resolve the
* scrape before {@code sessions.onTurnCompleteWithPostAction}'s adapter housekeeping can erase
* the pane's rendered output (see {@link CompletionResolver#resolveBeforePostAction}). That is
* why this composition does not treat the pair symmetrically the way {@link #bothMustRun}
* does for the other three callbacks: the mirror case (completion half throws, session half's
* return value still observed) is not preserved here — once the completion half's failure
* escapes, the session half's return value is discarded, matching how the {@link Injector}
* already treats any throw from this callback as "the action did not start" (see the {@code
* started} default at its call site, {@code Injector.java} ~line 616).
*
* <p>Package-private so {@code FleetdTurnListenerCompositionTest} can build this listener
* directly from a real {@link CompletionResolver} and a fake {@link TurnListener} standing in
* for {@code sessions}, without booting the rest of {@code main}.
*/
static TurnListener turnListener(CompletionResolver completion, TurnListener sessions) {
return new TurnListener() {
@Override
public void onTurnComplete(String target) {
bothMustRun(() -> completion.onTurnComplete(target), () -> sessions.onTurnComplete(target));
}
@Override
public boolean hasPostTurnAction(String target) {
return sessions.hasPostTurnAction(target);
}
@Override
public boolean onTurnCompleteWithPostAction(String target) {
return bothMustRunKeepingSecondResult(() -> completion.resolveBeforePostAction(target),
() -> sessions.onTurnCompleteWithPostAction(target));
}
@Override
public void onDelivered(String target, dev.ltms.fleet.msg.TurnToken token) {
// fleetd #561: left unguarded on purpose, not because registration survives a
// throw elsewhere. completion.onDelivered runs CompletionResolver.captureBaseline,
// which already wraps its scrape read in its own catch (RuntimeException) and
// fails open (baseline = null) — so this call does not realistically throw, and
// there is nothing here for bothMustRun to protect.
completion.onDelivered(target, token);
sessions.onDelivered(target, token);
}
@Override
public void onTurnFailed(String target) {
bothMustRun(() -> completion.onTurnFailed(target), () -> sessions.onTurnFailed(target));
}
@Override
public void onTurnFailed(String target, String reason) {
bothMustRun(() -> completion.onTurnFailed(target, reason), () -> sessions.onTurnFailed(target));
}
};
}
/**
* fleetd #561: run two listener-callback halves for one lifecycle event, guaranteeing the
* SECOND always runs even when the FIRST throws. Whatever is thrown is rethrown once both
* halves have been attempted — a second failure is attached to the first via {@link
* Throwable#addSuppressed} rather than dropped. Never swallows.
*/
private static void bothMustRun(Runnable completionHalf, Runnable sessionsHalf) {
Throwable failure = null;
try {
completionHalf.run();
} catch (Throwable t) {
failure = t;
}
try {
sessionsHalf.run();
} catch (Throwable t) {
if (failure == null) {
failure = t;
} else {
failure.addSuppressed(t);
}
}
if (failure != null) {
throwUnchecked(failure);
}
}
/**
* fleetd #561: like {@link #bothMustRun}, but for {@code onTurnCompleteWithPostAction}, whose
* session half returns the value the {@link Injector} needs. The second (session) half's
* result is what this method returns; if the first (completion) half throws, the second half
* still runs and its result is still computed here, but the throw is rethrown afterward
* regardless — so that result is discarded at the {@link Injector} call site exactly as it
* already is today when this callback throws (see {@link #turnListener}'s javadoc).
*/
private static boolean bothMustRunKeepingSecondResult(Runnable completionHalf,
BooleanSupplier sessionsHalf) {
Throwable failure = null;
try {
completionHalf.run();
} catch (Throwable t) {
failure = t;
}
boolean result = false;
try {
result = sessionsHalf.getAsBoolean();
} catch (Throwable t) {
if (failure == null) {
failure = t;
} else {
failure.addSuppressed(t);
}
}
if (failure != null) {
throwUnchecked(failure);
}
return result;
}
/**
* fleetd #561: rethrow a captured {@link Throwable} without a checked-exception wrapper. The
* two callback halves above never declare a checked exception (both existing production
* halves — {@code CompletionResolver} and {@code SessionManager} — only ever throw unchecked),
* so this only ever actually rethrows a {@link RuntimeException} or {@link Error}; the generic
* cast is the standard "sneaky throw" idiom, not a claim that a checked exception is expected.
*/
@SuppressWarnings("unchecked")
private static <T extends Throwable> void throwUnchecked(Throwable t) throws T {
throw (T) t;
}
/**
* fleetd #480: construct the {@link LeadRollover} executor only when {@code leadRollover:} is
* present at startup — the same presence gate {@code leadHeartbeat:} uses just above this
@@ -1220,6 +1580,58 @@ public final class Fleetd {
};
}
/**
* fleetd #602 gauge-wiring: per-lead-name factory for {@link FleetMcp.LeadConfigDirSource} — the
* {@code CLAUDE_CONFIG_DIR} the named lead's own profile runs on, so {@code LeadContextGauge}
* reads the transcript directory that lead's Claude Code process actually writes to, not always
* the built-in {@code <user.home>/.claude} default.
*
* <p>Follows the same {@code fleet.leaders.<name>.profile} link {@link #leadSeatLookup} already
* uses to find a lead's profile, one step further to that profile's own {@code configDir:}. A
* lead entry that names no {@code profile:}, or whose named profile is not configured, or whose
* profile sets no {@code configDir:} override, returns {@code null} — {@code LeadContextGauge}
* then falls back to its own default, exactly as before this ticket.
*
* @param profiles the live profile map, normally {@code () -> config.get().profiles()} in
* {@code main} — read live, like every other {@code configDir} lookup, so a
* config reload takes effect on the very next {@code fleet_list} call
* @param leaders {@code fleet.leaders}, read once at startup like {@link #leadSeatLookup}'s own
* {@code leaders} parameter — passed as a plain map, never re-read from
* {@code config.get()}
*/
static Function<String, String> leadConfigDirLookup(Supplier<Map<String, FleetConfig.Profile>> profiles,
Map<String, FleetConfig.Leader> leaders) {
return leadName -> {
FleetConfig.Leader lead = leaders.get(leadName);
if (lead == null || lead.profile() == null || lead.profile().isBlank()) {
return null;
}
FleetConfig.Profile leadProfile = profiles.get().get(lead.profile());
return leadProfile == null ? null : leadProfile.configDir();
};
}
/**
* fleetd #602 gauge-wiring follow-up (PR #606 review comment 17353): {@code main} used to build
* {@code new FleetMcp.LeadConfigDirSource(leadConfigDirLookup(...))} inline, with nothing a test
* could call directly. Measured on that shape: replacing the whole expression with {@code
* FleetMcp.LeadConfigDirSource.none()} at the call site compiled with 0 errors and left the full
* 1822-test suite green — the daemon could be changed to always report every lead's context as
* {@code UNKNOWN}, forever, while every test stayed green. That is the same hand-built-vs-wired
* shape as fleetd #561/#248/#426/#562 ({@link #loopHealthSource}).
*
* <p>The fix extracts the inline {@code new} into this factory, in the same style as {@link
* #loopHealthSource}/{@link #capacitySource}/{@link #healthCoverageSource} — which is exactly
* what makes it directly callable from {@code FleetdLeadConfigDirSourceWiringTest}. That test
* calls this factory with real {@link FleetConfig.Profile}/{@link FleetConfig.Leader} fixtures and
* asserts the returned source resolves a real {@code configDir} — a property that would be false
* if this method's body were mutated to {@code return FleetMcp.LeadConfigDirSource.none();}.
*/
static FleetMcp.LeadConfigDirSource leadConfigDirSource(Supplier<Map<String, FleetConfig.Profile>> profiles,
Map<String, FleetConfig.Leader> leaders) {
return new FleetMcp.LeadConfigDirSource(leadConfigDirLookup(profiles, leaders));
}
/**
* fleetd #248 / fleetd#201 Unit 5: package-private factory for the per-target backend-error
* pattern lookup {@link CompletionResolver} classifies a pane scrape against. Closes over the
@@ -221,7 +221,8 @@ public final class CallerResolver {
// The config/live binding names this pane as an architect slot's own. Same
// unforgeable pane mapping; the live binding, never a request argument, decides.
// Check the slot role too: this defence in depth prevents a bad lifecycle bind from
// escalating a dev or reviewer into an architect. Checked before the worker fallback.
// escalating a dev, hunter or reviewer into an architect. Checked before
// the worker fallback.
return Principal.architect(memberSlotNames.apply(slot), c.terminal(), c.pid());
}
return Principal.worker(c.terminal(), c.pid()); // unforgeable; never token-gated
@@ -42,7 +42,8 @@ public interface MemberLifecycle {
* Try to bind a newly spawned {@code terminal} into the role it was granted.
*
* @return the role this session actually holds: {@code role} unchanged for a role with no
* live slot-binding semantics (dev, reviewer), or when the bind succeeded; a fallback
* live slot-binding semantics (dev, hunter, reviewer), or when the bind
* succeeded; a fallback
* role — never {@code role} — when a slot-bound role (architect) could not be bound.
* Callers must record THIS value on the session, never the requested {@code role}, so
* a later roster read never reports a role the session does not hold (CB-619). In
@@ -20,7 +20,8 @@ import java.util.function.Supplier;
*
* <p>Two halves, split by who owns each:
* <ul>
* <li><b>slots</b> — read from {@code fleet.architects}/{@code developers}/{@code reviewers}
* <li><b>slots</b> — read from {@code fleet.architects}/{@code developers}/
* {@code hunters}/{@code reviewers}
* (see {@link #slots()}), each carrying the {@code profile} reference the spawn lifecycle
* reads when it stands the slot up. <strong>Live, since fleetd #424</strong>: {@link #live}
* re-reads {@code fleet:} on every call, through a supplier the same shape as
@@ -322,7 +323,7 @@ public final class MemberRegistry implements MemberLifecycle {
* CB-619 / fleetd #123: refuse an architect acquire before anything spawns when no configured
* slot carries {@code profile} — the config-gap case from the original defect report (a spawn
* asked for {@code role=architect, profile=sonnet}, and {@code fleet.architects} carried only
* {@code opus} and {@code sol}). A dev/reviewer acquire is always a no-op: those pools are
* {@code opus} and {@code sol}). A dev/hunter/reviewer acquire is always a no-op: those pools are
* placement candidates only (see {@code CompositePeerLauncher}), never a live identity binding,
* so there is nothing here to refuse — an explicit profile outside the pool for those roles is a
* documented operator override, not a defect.
@@ -31,7 +31,7 @@ import java.util.function.Supplier;
* {@code placement:}, and an existing profile's {@code weight} / {@code maxLoad}. Both are
* read through a supplier on {@code CompositePeerLauncher}, which is what makes them hot —
* not the fact that they are config. Most of {@code fleet:} — every role pool
* ({@code architects}/{@code developers}/{@code reviewers}), {@code charters}, and
* ({@code architects}/{@code developers}/{@code hunters}/{@code reviewers}), {@code charters}, and
* {@code tabLabel} — is read the same live way, through the same supplier
* ({@code () -> config.get().fleet()}). {@code architects} in particular is hot for
* <strong>two independent consumers</strong> (fleetd #424): {@code CompositePeerLauncher}
@@ -605,7 +605,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
+ "opened once and needs a restart; the broker URI env-var name kept out of a "
+ "member's environment is read live on every spawn and already applied");
}
// fleetd #333: unlike health/coordinator above, most of `fleet:` (developers, reviewers,
// fleetd #333: unlike health/coordinator above, most of `fleet:` (developers, hunters, reviewers,
// charters, tabLabel) is genuinely hot — ConfigRefTest.aHotChangeIsAppliedAndRead-
// ThroughGet and aCharterChangeIsHotAndReachesTheLiveConfig prove it reaches the live config
// with no restart note. `architects` is hot too, and — since fleetd #424 — hot for BOTH of
@@ -630,7 +630,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
+ "identity map and to auto-launch leads, and neither is rebuilt on reload, so a "
+ "lead added, removed, or given a new tab: label needs a restart — until then it "
+ "stays unrecognised, and a caller from its new tab resolves as a worker, not a "
+ "lead; the rest of fleet: (developers, reviewers, charters, tabLabel) is read "
+ "lead; the rest of fleet: (developers, hunters, reviewers, charters, tabLabel) is read "
+ "live through the supplier on CompositePeerLauncher, and architects is read "
+ "live through that same supplier for placement AND through a separate supplier "
+ "on MemberRegistry for spawn-time identity — both already applied");
@@ -62,7 +62,7 @@ import java.util.regex.PatternSyntaxException;
* @param fleet who the daemon may run and under which role (CB-557). One block replacing
* the former {@code leaders:}, {@code members:}, {@code leadScan:} and
* {@code defaultProfile:}. Role is the containing key — {@code leaders},
* {@code architects}, {@code developers}, {@code reviewers} — and each entry
* {@code architects}, {@code developers}, {@code hunters}, {@code reviewers} — and each entry
* names the {@code profiles:} backend it runs on. See {@link Fleet}
* @param leadHeartbeat opt-in idle-lead heartbeat (CB-551); {@code null} ⇒ off, and an upgraded
* daemon never nudges an idle lead on its own initiative
@@ -1200,6 +1200,7 @@ public record FleetConfig(
* @param leaders panes that orchestrate rather than are orchestrated, keyed by lead name
* @param architects profiles the {@code architect} role may run on
* @param developers profiles the {@code dev} role may run on
* @param hunters profiles the {@code hunter} role may run on
* @param reviewers profiles the {@code reviewer} role may run on
* @param charters optional launch-charter text keyed by singular role wire name
* @param tabLabel template for a member tab's label; {@code {role}}, {@code {profile}},
@@ -1210,6 +1211,7 @@ public record FleetConfig(
public record Fleet(Map<String, Leader> leaders,
Map<String, Slot> architects,
Map<String, Slot> developers,
Map<String, Slot> hunters,
Map<String, Slot> reviewers,
Map<String, String> charters,
String tabLabel) {
@@ -1226,6 +1228,7 @@ public record FleetConfig(
leaders = unmodifiableOrEmpty(leaders);
architects = unmodifiableOrEmpty(architects);
developers = unmodifiableOrEmpty(developers);
hunters = unmodifiableOrEmpty(hunters);
reviewers = unmodifiableOrEmpty(reviewers);
charters = unmodifiableOrEmpty(charters);
tabLabel = (tabLabel == null || tabLabel.isBlank()) ? DEFAULT_TAB_LABEL : tabLabel;
@@ -1240,9 +1243,15 @@ public record FleetConfig(
* constructor: the launcher reads {@code fleet.charters()} from the live config. Jackson
* binds the canonical constructor, so this one cannot swallow an operator's YAML.
*/
public Fleet(Map<String, Leader> leaders, Map<String, Slot> architects,
Map<String, Slot> developers, Map<String, Slot> reviewers,
Map<String, String> charters, String tabLabel) {
this(leaders, architects, developers, null, reviewers, charters, tabLabel);
}
public Fleet(Map<String, Leader> leaders, Map<String, Slot> architects,
Map<String, Slot> developers, Map<String, Slot> reviewers, String tabLabel) {
this(leaders, architects, developers, reviewers, null, tabLabel);
this(leaders, architects, developers, null, reviewers, null, tabLabel);
}
/**
@@ -1264,6 +1273,7 @@ public record FleetConfig(
return switch (role) {
case ARCHITECT -> architects;
case DEV -> developers;
case HUNTER -> hunters;
case REVIEWER -> reviewers;
};
}
@@ -1872,7 +1882,7 @@ public record FleetConfig(
/** The {@code fleet:} child blocks whose direct children are slot names. */
private static final Set<String> FLEET_POOL_KEYS =
Set.of("leaders", "architects", "developers", "reviewers");
Set.of("leaders", "architects", "developers", "hunters", "reviewers");
/**
* Reject a {@code fleet:} role pool whose slot names repeat (CB-548, re-homed by CB-557).
@@ -1882,7 +1892,7 @@ public record FleetConfig(
* daemon would never know. Jackson's YAML parser does not fail on duplicate mapping keys by
* default, so duplicates are caught here, at parse time, before the map is built.
*
* <p>Only the four pools <em>directly under the top-level {@code fleet:}</em> are considered,
* <p>Only the five pools <em>directly under the top-level {@code fleet:}</em> are considered,
* and only their direct child keys (the slot names). A nested field elsewhere, even one also
* named {@code developers:}, is ignored, so parsing of the rest of the config is unaffected.
*
@@ -2049,8 +2059,9 @@ public record FleetConfig(
"defaultProfile", "a role pool under 'fleet:' — an unqualified spawn now names a role,"
+ " and that role's pool supplies the candidate profiles",
"architects", "'fleet.architects'",
"members", "a role pool under 'fleet:' — 'fleet.architects', 'fleet.developers' or"
+ " 'fleet.reviewers'; the role is the containing key, not a 'role:' field",
"members", "a role pool under 'fleet:' — 'fleet.architects', 'fleet.developers',"
+ " 'fleet.hunters' or 'fleet.reviewers'; the role is the containing key, not"
+ " a 'role:' field",
"leaders", "'fleet.leaders'",
"leadScan", "'fleet.leaders.<name>.tabPrefix' and '.scanIntervalSeconds' — lead"
+ " discovery is now configured on the lead it discovers");
@@ -0,0 +1,307 @@
package dev.ltms.fleet.lead;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import java.io.IOException;
import java.io.RandomAccessFile;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.LongSupplier;
/**
* Reads how full a lead's own Claude Code context window is, from the transcript Claude Code
* itself writes — never from the lead's pane (fleetd has {@code AgentControl.read} for that, and
* this must not use it: a pane holds terminal text, not the structured usage numbers a transcript
* carries, and scraping it would also race the lead's own rendering).
*
* <p><strong>Why this exists.</strong> A lead auto-compacts when its context fills — on the host
* this was built for, that happened 30 times in one session, discarding roughly 250,000 tokens
* and costing 46 seconds to 3 minutes each time, and fleetd had no way to see it coming. This
* class is the first thing that looks.
*
* <p><strong>The route.</strong> Claude Code appends one JSON object per line to
* {@code <configDir>/projects/<slug>/<sessionId>.jsonl}. {@code <slug>} is an undocumented,
* internal encoding of the working directory — this class never derives it. Instead it lists the
* one-level-deep subdirectories of {@code <configDir>/projects/} and looks for
* {@code <sessionId>.jsonl} by name, so the slug rule can change without breaking this reader.
*
* <ul>
* <li><strong>Live context</strong> is read off the last record in the read window that carries
* a {@code message.usage} object: {@code input_tokens + cache_read_input_tokens +
* cache_creation_input_tokens}. This is what actually fills the window — a plain
* {@code input_tokens} count alone understates it once the conversation has any cached
* prefix, which on a long-lived lead is always.</li>
* <li><strong>Compaction history</strong> is a count of {@code subtype: "compact_boundary"}
* records seen in the same read window — see {@link Reading#compactions()}. It is a count
* within the window this reader actually looked at, not a lifetime total: a session with
* more compactions than fit in {@link #TAIL_BYTES} of transcript will undercount. That
* trade-off is deliberate — see {@link #TAIL_BYTES}.</li>
* </ul>
*
* <p><strong>Three states, not two (OK / HIGH / UNKNOWN).</strong> Every path that cannot
* positively establish the live token count — a missing file, an unreadable one, a peer that
* is not a Claude backend, or every line in the read window failing to parse as JSON — returns
* {@link State#UNKNOWN} with no token number, never a default "0" or "ok" that would read as
* "this lead is fine" when the honest answer is "I could not look".
*
* <p><strong>A torn final line does not mean UNKNOWN.</strong> {@code fleet_list} reads this
* transcript while Claude Code may be mid-write on it, so the last line in the window can be cut
* off mid-flush — that is an ordinary, expected race, not a sign the format has changed. Earlier
* this class treated ANY unparseable last line as UNKNOWN, on the theory that "if the format
* changes, we should see UNKNOWN". That reasoning does not hold: a real format change makes
* <em>every</em> line in the window unparseable, not only the last one written. So a single
* malformed line (most often the final, torn one, but the check is not position-specific) is
* simply skipped rather than treated as fatal, and the reading is built from whatever lines in the
* window did parse. Only when <em>none</em> of them parse — the real format-change signal — does
* this return {@link State#UNKNOWN}, still with no stale number standing in for "I could not
* tell".
*
* <p><strong>Bounded cost.</strong> {@code fleet_list} is polled constantly, so every read is
* capped two ways: {@link #TAIL_BYTES} bounds how much of the transcript is ever read from disk
* (never the whole 52 MB a long-lived transcript reaches on the host this was measured on), and
* {@link #DEFAULT_CACHE_TTL_MILLIS} bounds how often that bounded read actually happens — a burst
* of {@code fleet_list} calls inside one TTL window reads the file once. One instance's cache is
* keyed by {@code (configDir, sessionId)}, so it is safe to share across every lead a single
* {@code fleet_list} call reports on.
*/
public final class LeadContextGauge {
private static final String PROJECTS_DIR = "projects";
private static final ObjectMapper MAPPER = new ObjectMapper();
/**
* How many trailing bytes of a transcript a single read ever pulls off disk. Chosen so one
* read comfortably spans many recent turns — each usage or compact_boundary record is at most
* a few KB — while staying nowhere near the 52 MB a long session's real transcript reaches on
* the host this was built for; reading that whole file on every {@code fleet_list} call is
* exactly the cost this bound exists to avoid. 2 MiB holds on the order of hundreds of recent
* lines even when a turn's tool output is unusually large, which is far more than needed to
* find the most recent usage record and any recent compaction.
*/
static final int TAIL_BYTES = 2 * 1024 * 1024;
/**
* How long a {@link Reading} is served from cache before the file is read again.
* {@code fleet_list} is called constantly (by design — it is the fleet's own status probe), so
* without a TTL a burst of calls would re-read the transcript tail once per call. 5 seconds is
* short enough that a caller watching for a state change never waits long, and long enough that
* a poll loop calling every second or two only touches disk once per window.
*/
static final long DEFAULT_CACHE_TTL_MILLIS = 5_000;
/**
* Live tokens at or above this count report {@link State#HIGH}. On the host this was measured
* on, auto-compaction actually fires around 267,000–270,000 tokens, but the point of a HIGH
* state is to warn before that happens, not at it — 200,000 is the standard Claude context
* window size and a sensible built-in default: no config key is required to pick it, and a
* lead crossing it is already deep enough into its window that a compaction is foreseeable.
*/
static final long HIGH_THRESHOLD_TOKENS = 200_000;
/** The only peer kind this reader understands ({@code Agent.agentType()}'s wire value). */
private static final String CLAUDE_AGENT_TYPE = "claude";
public enum State { OK, HIGH, UNKNOWN }
/**
* @param state {@link State#UNKNOWN} whenever {@code tokens} could not be established
* @param tokens live context tokens, or {@code null} exactly when {@code state} is
* {@link State#UNKNOWN}
* @param compactions {@code compact_boundary} records seen in the read window (see class
* javadoc) — {@code 0} both for "genuinely none seen" and for "unknown",
* since a caller that already sees {@code state: UNKNOWN} has no reason to
* trust this number either way
*/
public record Reading(State state, Long tokens, int compactions) {
static Reading unknown() {
return new Reading(State.UNKNOWN, null, 0);
}
}
private record CacheEntry(Reading reading, long readAtMillis) {
}
private final LongSupplier clock;
private final long ttlMillis;
private final Map<String, CacheEntry> cache = new ConcurrentHashMap<>();
/** Test seam only (package-private) — counts real disk reads, i.e. cache misses. */
private final AtomicInteger diskReads = new AtomicInteger();
public LeadContextGauge() {
this(System::currentTimeMillis, DEFAULT_CACHE_TTL_MILLIS);
}
/** Test seam: an injectable clock and TTL so cache expiry is provable without sleeping. */
LeadContextGauge(LongSupplier clock, long ttlMillis) {
this.clock = clock;
this.ttlMillis = ttlMillis;
}
/** How many times this instance has actually read a transcript off disk — test seam only. */
int diskReadCount() {
return diskReads.get();
}
/**
* @param configDir the lead's {@code CLAUDE_CONFIG_DIR}, or {@code null}/blank to use the
* default {@code <user.home>/.claude} — the right answer for the common case
* where the lead's profile sets no {@code configDir} override
* @param sessionId the lead's own Claude session id ({@code Agent.sessionId()}), or
* {@code null} when herdr has not resolved one yet
* @param agentType the detected peer kind ({@code Agent.agentType()}); anything other than
* {@code "claude"} (including {@code null}, meaning undetected) reports
* {@link State#UNKNOWN} — this reader only understands Claude Code's own
* transcript format
*/
public Reading read(String configDir, String sessionId, String agentType) {
if (sessionId == null || sessionId.isBlank()) {
return Reading.unknown();
}
if (!CLAUDE_AGENT_TYPE.equalsIgnoreCase(agentType)) {
return Reading.unknown();
}
String base = (configDir == null || configDir.isBlank())
? System.getProperty("user.home") + "/.claude"
: configDir;
String cacheKey = base + '\u0000' + sessionId;
long now = clock.getAsLong();
CacheEntry cached = cache.get(cacheKey);
if (cached != null && now - cached.readAtMillis() < ttlMillis) {
return cached.reading();
}
Reading fresh = readUncached(base, sessionId);
cache.put(cacheKey, new CacheEntry(fresh, now));
return fresh;
}
private Reading readUncached(String base, String sessionId) {
diskReads.incrementAndGet();
Path file = findTranscript(base, sessionId);
if (file == null) {
return Reading.unknown();
}
TailRead tail;
try {
tail = tailBytes(file, TAIL_BYTES);
} catch (IOException e) {
return Reading.unknown();
}
return parse(tail);
}
/**
* Finds {@code <sessionId>.jsonl} under {@code <base>/projects/}, one level deep — never by
* deriving the slug directory from a working directory (see class javadoc). Bounded to a
* single {@code list()} of {@code projects/} itself: it never recurses further, so the cost is
* the number of project directories, not the size of any transcript inside them.
*/
private Path findTranscript(String base, String sessionId) {
Path projectsDir = Path.of(base, PROJECTS_DIR);
if (!Files.isDirectory(projectsDir)) {
return null;
}
String filename = sessionId + ".jsonl";
Path direct = projectsDir.resolve(filename);
if (Files.isRegularFile(direct)) {
return direct;
}
try (var children = Files.list(projectsDir)) {
return children.filter(Files::isDirectory)
.map(dir -> dir.resolve(filename))
.filter(Files::isRegularFile)
.findFirst()
.orElse(null);
} catch (IOException e) {
return null;
}
}
/** Package-private (not {@code private}): {@link #tailBytes} is a test seam, see its javadoc. */
record TailRead(byte[] bytes, boolean fromStart) {
}
/**
* Reads at most {@code maxBytes} trailing bytes of {@code file}. Package-private (not
* {@code private}) so a test can assert directly on the returned array's length — "bytes
* actually read", not on any parsed answer — without needing a file anywhere near
* {@link #TAIL_BYTES} in size to prove the cap holds.
*/
static TailRead tailBytes(Path file, int maxBytes) throws IOException {
try (RandomAccessFile raf = new RandomAccessFile(file.toFile(), "r")) {
long length = raf.length();
long start = Math.max(0, length - maxBytes);
raf.seek(start);
byte[] buf = new byte[(int) (length - start)];
raf.readFully(buf);
return new TailRead(buf, start == 0);
}
}
/**
* Parses the tail into a {@link Reading}. The first line is dropped unconditionally whenever
* the tail is not the whole file (it starts mid-line, cut by {@link #TAIL_BYTES} — an expected
* artefact of the bound, not a data problem). Every remaining line is then parsed on a
* best-effort basis: a line that fails to parse (most often the last one, torn by a write this
* read raced — see the "torn final line" section of the class javadoc) is skipped, not fatal.
* Only when none of the remaining lines parse does this report {@link State#UNKNOWN}.
*/
private Reading parse(TailRead tail) {
String text = new String(tail.bytes(), StandardCharsets.UTF_8);
List<String> lines = new ArrayList<>(List.of(text.split("\n", -1)));
if (!lines.isEmpty() && lines.get(lines.size() - 1).isEmpty()) {
lines.remove(lines.size() - 1); // trailing newline leaves a phantom empty last element
}
if (!tail.fromStart() && !lines.isEmpty()) {
lines.remove(0); // first line is a fragment cut by our own tail bound, not real data
}
if (lines.isEmpty()) {
return Reading.unknown();
}
Long tokens = null;
int compactions = 0;
boolean anyLineParsed = false;
for (String line : lines) {
JsonNode node = tryParse(line);
if (node == null) {
// A malformed line — typically the last one, cut mid-flush by a write this read
// raced — is skipped rather than treated as fatal. See the class javadoc's "torn
// final line" section for why: a real format change makes EVERY line unparseable,
// not only this one, and that case is still caught below by anyLineParsed.
continue;
}
anyLineParsed = true;
JsonNode usage = node.path("message").path("usage");
if (usage.isObject()) {
tokens = usage.path("input_tokens").asLong(0)
+ usage.path("cache_read_input_tokens").asLong(0)
+ usage.path("cache_creation_input_tokens").asLong(0);
}
if ("compact_boundary".equals(node.path("subtype").asText(null))) {
compactions++;
}
}
if (!anyLineParsed) {
return Reading.unknown();
}
if (tokens == null) {
return new Reading(State.UNKNOWN, null, compactions);
}
State state = tokens >= HIGH_THRESHOLD_TOKENS ? State.HIGH : State.OK;
return new Reading(state, tokens, compactions);
}
private JsonNode tryParse(String line) {
try {
return MAPPER.readTree(line);
} catch (IOException e) {
return null;
}
}
}
@@ -12,6 +12,7 @@ import dev.ltms.fleet.metrics.Metrics;
import dev.ltms.fleet.inject.MemberPresence;
import dev.ltms.fleet.inject.CompletionResolver;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.lead.LeadContextGauge;
import dev.ltms.fleet.lead.LeadRollover;
import dev.ltms.fleet.msg.LeadChannel;
import dev.ltms.fleet.msg.LeadMessage;
@@ -20,6 +21,7 @@ import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.peer.PeerUnreachableException;
import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.inject.LoopWatchdog;
import dev.ltms.fleet.placement.PlacementException;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.session.MemberSession;
@@ -108,11 +110,14 @@ public final class FleetMcp {
private final Metrics metrics; // CB-502: null → auth failures not counted
private final CapacitySource capacity;
private final HealthCoverageSource healthCoverage;
private final LoopHealthSource loopHealth;
private final QuarantineSource quarantine;
/** fleetd #201 Unit 5: SEPARATE from {@link #quarantine} — see {@link OutageSource}'s doc. */
private final OutageSource outage;
/** fleetd #176: SEPARATE from both of the above — see {@link LeadSeatSource}'s doc. */
private final LeadSeatSource leadSeats;
/** fleetd #602 gauge-wiring: see {@link LeadConfigDirSource}. */
private final LeadConfigDirSource leadConfigDirs;
/** CB-637: this daemon's lead-to-lead channel; {@code null} when no coordinator is configured. */
private final LeadChannel leadChannel;
/** fleetd #361: {@code coordinator.peers} — see {@link CoordinationSource}. Empty when unset. */
@@ -124,6 +129,17 @@ public final class FleetMcp {
* clean {@code NOT_CONFIGURED} refusal rather than throwing. See {@link #handover}.
*/
private final LeadRollover leadRollover;
/**
* "Lead context gauge": how full each lead's own Claude Code context window is, reported on
* {@code fleet_list}'s {@code leads} rows (see {@link #contextView}). Built unconditionally, in
* the field initializer rather than a constructor parameter — this is not a togglable feature
* with an on/off config knob the way {@link OutageSource}/{@link LeadSeatSource} are: it needs
* no config at all (see {@link LeadContextGauge}'s own javadoc for the built-in defaults), so
* there is no "off" value to thread through every existing constructor call site. One instance
* per daemon so its read cache (keyed by session, TTL'd) is actually shared across
* {@code fleet_list} calls rather than rebuilt — and therefore useless — on every call.
*/
private final LeadContextGauge leadContextGauge = new LeadContextGauge();
/** Capacity facts used by {@code fleet_list}; production must supply the placement live count. */
public record CapacitySource(Function<String, Integer> liveCount, Function<String, Integer> maxLoad,
@@ -136,6 +152,15 @@ public final class FleetMcp {
/** Coverage is supplied by the health wiring, not inferred from a missing dependency. */
public record HealthCoverageSource(Supplier<String> value) { }
/** Progress states for fleetd's singleton background loops, read by {@code fleet_list} and {@code /healthz}. */
public record LoopHealthSource(Supplier<LoopWatchdog.State> statusPoller,
Supplier<LoopWatchdog.State> sessionReaper) {
/** Inert source for callers that do not wire the background loops. */
public static LoopHealthSource none() {
return new LoopHealthSource(() -> LoopWatchdog.State.STOPPED, () -> LoopWatchdog.State.STOPPED);
}
}
/**
* 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.
@@ -235,6 +260,29 @@ public final class FleetMcp {
public static LeadSeatSource none() { return new LeadSeatSource(_ -> 0); }
}
/**
* fleetd #602 gauge-wiring: a lead name's configured {@code CLAUDE_CONFIG_DIR} override, fed to
* {@link LeadContextGauge#read} so {@code fleet_list}'s {@code context} row reads the transcript
* directory the lead's own profile actually writes to, not always the built-in
* {@code <user.home>/.claude} default.
*
* <p>The same idiom as {@link LeadSeatSource} — a {@code FleetMcp} constructor field, not a
* lookup {@code contextView} performs itself, because {@code FleetMcp} holds no
* {@link dev.ltms.fleet.config.FleetConfig} and {@code contextView} is {@code static}. See
* {@code Fleetd.leadConfigDirLookup} for the derivation: the same
* {@code fleet.leaders.<name>.profile} link {@link LeadSeatSource} already follows, resolved to
* that profile's own {@code configDir:}.
*
* @param configDirFor lead name → {@code configDir}, or {@code null} when the lead's entry names
* no profile, or that profile sets no {@code configDir} override — either
* way {@link LeadContextGauge} then falls back to its own built-in default,
* exactly as before this ticket
*/
public record LeadConfigDirSource(Function<String, String> configDirFor) {
/** Inert source — every lead reads {@link LeadContextGauge}'s built-in default {@code configDir}. */
public static LeadConfigDirSource none() { return new LeadConfigDirSource(_ -> null); }
}
/**
* fleetd #361: peer-visibility facts for {@code fleet_list}'s {@code coordinator} row — this
* daemon's own {@link LeadChannel} (for its self mailbox state and held messages) plus the
@@ -328,11 +376,23 @@ public final class FleetMcp {
* instead of throwing. See {@link #handover}.
*/
public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
CallerResolver callers, AuthorizationMode authorizationMode, Metrics metrics,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage,
LeadSeatSource leadSeats, List<String> peers, LeadRollover leadRollover) {
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
CallerResolver callers, AuthorizationMode authorizationMode, Metrics metrics,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage,
LeadSeatSource leadSeats, List<String> peers, LeadRollover leadRollover) {
this(messages, workers, sessions, identity, presence, primaryRegistry, callers, authorizationMode, metrics,
capacity, healthCoverage, LoopHealthSource.none(), quarantine, leadChannel, outage, leadSeats,
LeadConfigDirSource.none(), peers, leadRollover);
}
public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
CallerResolver callers, AuthorizationMode authorizationMode, Metrics metrics,
CapacitySource capacity, HealthCoverageSource healthCoverage, LoopHealthSource loopHealth,
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage,
LeadSeatSource leadSeats, LeadConfigDirSource leadConfigDirs, List<String> peers,
LeadRollover leadRollover) {
Objects.requireNonNull(callers, "callers");
this.authorizationEnforced = Objects.requireNonNull(authorizationMode, "authorizationMode")
== AuthorizationMode.ENFORCED;
@@ -342,7 +402,9 @@ public final class FleetMcp {
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
this.outage = Objects.requireNonNull(outage, "outage");
this.leadSeats = Objects.requireNonNull(leadSeats, "leadSeats");
this.leadConfigDirs = Objects.requireNonNull(leadConfigDirs, "leadConfigDirs");
this.healthCoverage = healthCoverage;
this.loopHealth = Objects.requireNonNull(loopHealth, "loopHealth");
this.leadRollover = leadRollover;
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
this.transport = HttpServletStreamableServerTransportProvider.builder()
@@ -474,8 +536,8 @@ public final class FleetMcp {
(exchange, _) -> {
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_list", Map.of()), null);
if (denied != null) return denied;
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
leadSeats, callers.leads(),
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage,
leadSeats, leadContextGauge, leadConfigDirs, callers.leads(),
callerTerminal(exchange),
new CoordinationSource(leadChannel, peers),
coordinatorVisibleTo(principal(exchange)));
@@ -805,6 +867,12 @@ public final class FleetMcp {
+ "answered (turnId stale)");
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> text("[no reply within " + timeout + "ms — worker "
+ r.outcome().name().toLowerCase().replace("timed_out_", "") + "; retry or poll status]");
// fleetd #571: delivery is unknown here — agent.prompt pastes and submits in one call,
// so the message may already be sitting in the pane. Do not invite a blind retry the way
// the case above does; a resend on this route can double-deliver the same brief.
case TIMED_OUT_UNCONFIRMED -> text("[no reply within " + timeout + "ms — delivery unconfirmed; "
+ "the message may already have reached the worker, so a retry risks sending it "
+ "twice — poll status before resending]");
};
}
@@ -1546,25 +1614,42 @@ public final class FleetMcp {
* @param selfTerm the calling pane's terminal id, or blank for a caller with no pane
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions,
Map<String, String> leads, String selfTerm) {
Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, null, CapacitySource.none(), new HealthCoverageSource(() -> "off"),
QuarantineSource.none(), leads, selfTerm);
LoopHealthSource.none(), QuarantineSource.none(), leads, selfTerm);
}
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, leads, selfTerm,
CoordinationSource.none());
}
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, leads, selfTerm,
LoopHealthSource loopHealth, QuarantineSource quarantine,
Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, leads, selfTerm,
CoordinationSource.none());
}
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
LoopHealthSource loopHealth, QuarantineSource quarantine,
Map<String, String> leads, String selfTerm,
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine,
OutageSource.none(), LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, OutageSource outage,
Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
LeadSeatSource.none(), leads, selfTerm, CoordinationSource.none());
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, CoordinationSource.none(), false);
}
/**
@@ -1580,8 +1665,8 @@ public final class FleetMcp {
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, Map<String, String> leads, String selfTerm,
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, OutageSource.none(),
LeadSeatSource.none(), leads, selfTerm, coordination);
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, OutageSource.none(),
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
@@ -1589,8 +1674,8 @@ public final class FleetMcp {
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, OutageSource outage,
Map<String, String> leads, String selfTerm, CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
LeadSeatSource.none(), leads, selfTerm, coordination);
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
}
/**
@@ -1611,8 +1696,8 @@ public final class FleetMcp {
QuarantineSource quarantine, OutageSource outage,
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
leadSeats, leads, selfTerm, coordination, false);
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
}
/**
@@ -1631,9 +1716,29 @@ public final class FleetMcp {
* an explicit {@code true}
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, OutageSource outage,
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, OutageSource outage,
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
CoordinationSource coordination, boolean callerIsPrimary) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine,
outage, leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, callerIsPrimary);
}
/**
* The canonical implementation. {@code contextGauge} is the "lead context gauge" (see
* {@link LeadContextGauge}) — every wrapper overload above passes a freshly constructed one,
* which is correct for them (none of them exercise repeated calls where a shared cache would
* matter); the one caller that matters for caching, {@code fleet_list}'s MCP handler, passes
* its own single long-lived instance instead (see {@code FleetMcp}'s {@code leadContextGauge}
* field).
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
LoopHealthSource loopHealth,
QuarantineSource quarantine, OutageSource outage,
LeadSeatSource leadSeats, LeadContextGauge contextGauge,
LeadConfigDirSource leadConfigDirs,
Map<String, String> leads, String selfTerm,
CoordinationSource coordination, boolean callerIsPrimary) {
try {
Map<String, Agent> live = workers.list().stream()
@@ -1642,7 +1747,8 @@ public final class FleetMcp {
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
List<Map<String, Object>> leadRows = leads.entrySet().stream()
.sorted(Map.Entry.comparingByValue())
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm))
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm, contextGauge,
leadConfigDirs))
.toList();
// fleetd #209: this is the caller-driven fleet_list read that actually reports
// agentSessionId (via memberCapacityView -> SessionManager.rosterView), so it uses the
@@ -1656,6 +1762,9 @@ public final class FleetMcp {
Map<String, Object> result = new LinkedHashMap<>();
result.put("leads", leadRows); result.put("members", out);
result.put("healthCoverage", healthCoverage.value().get());
result.put("loopHealth", Map.of(
"statusPoller", loopHealth.statusPoller().get().name(),
"sessionReaper", loopHealth.sessionReaper().get().name()));
// fleetd #439: coordinator/coordinatorView is lead-to-lead coordination state and must
// never reach a worker or an architect -- gate BEFORE assembling it, not after, so the
// key is absent rather than present-and-empty.
@@ -1934,15 +2043,25 @@ public final class FleetMcp {
}
/**
* One lead's row: its address, its name, and whether it can be reached right now.
* One lead's row: its address, its name, whether it can be reached right now, and how full its
* own Claude Code context window is.
*
* <p>{@code status} is herdr's live view, and {@code unknown} when herdr is not tracking that
* pane as an agent — the honest answer, and the one that matters: a lead whose pane herdr cannot
* see is a lead a {@code fleet_send} cannot be typed into. It is reported rather than hidden,
* because a peer that has gone unreachable is exactly what the sender needs to know.
*
* <p>{@code context} is the "lead context gauge" (fleetd's context-usage visibility ticket):
* {@code {state: "ok"|"high"|"unknown", tokens?: number, compactions: number}}, read from the
* lead's transcript — see {@link LeadContextGauge}. Kept small on purpose (a token count, a
* state, and a compaction count) rather than echoing the whole reading history: this is a
* roster row a caller glances at, not a diagnostics dump. {@code tokens} is present only when
* {@code state} is not {@code "unknown"} — never a stale or default number standing in for "I
* could not tell".
*/
private static Map<String, Object> leadView(String terminal, String name, Agent live,
String selfTerm) {
String selfTerm, LeadContextGauge contextGauge,
LeadConfigDirSource leadConfigDirs) {
Map<String, Object> m = new LinkedHashMap<>();
m.put("sessionId", terminal);
m.put("name", name);
@@ -1951,9 +2070,32 @@ public final class FleetMcp {
if (terminal.equals(selfTerm)) {
m.put("self", true);
}
String configDir = leadConfigDirs.configDirFor().apply(name);
m.put("context", contextView(contextGauge, live, configDir));
return m;
}
/**
* Reads the lead context gauge for one lead. {@code configDir} is this lead's configured
* {@code CLAUDE_CONFIG_DIR} override (see {@link LeadConfigDirSource}), derived from
* {@code fleet.leaders.<name>.profile} → that profile's own {@code configDir:} — or {@code null}
* when the lead's entry names no profile, or that profile sets no override, in which case
* {@link LeadContextGauge#read} falls back to its own built-in default
* ({@code <user.home>/.claude}).
*/
private static Map<String, Object> contextView(LeadContextGauge contextGauge, Agent live, String configDir) {
String sessionId = live == null ? null : live.sessionId();
String agentType = live == null ? null : live.agentType();
LeadContextGauge.Reading reading = contextGauge.read(configDir, sessionId, agentType);
Map<String, Object> c = new LinkedHashMap<>();
c.put("state", reading.state().name().toLowerCase());
if (reading.tokens() != null) {
c.put("tokens", reading.tokens());
}
c.put("compactions", reading.compactions());
return c;
}
/** {@code fleet_stop}: tear a worker down by its pane id. */
static McpSchema.CallToolResult stop(SessionManager sessions, String paneId) {
if (isBlank(paneId)) {
@@ -2075,8 +2217,9 @@ public final class FleetMcp {
return tool(FleetTool.SPAWN.wireName(),
"Spawn a new off-subscription member session. A member has two independent attributes: "
+ "role (what it is for) and profile (which backend it runs on). Pass role to pick "
+ "the contract — 'dev' implements a unit and opens its own PR, 'reviewer' reviews a "
+ "diff it did not write, 'architect' refines a ticket before anyone builds it; omit "
+ "the contract — 'dev' implements a unit and opens its own PR, 'hunter' sweeps a "
+ "scope without changing it, 'reviewer' reviews a diff it did not write, 'architect' "
+ "refines a ticket before anyone builds it; omit "
+ "it for 'dev'. Pass profile (from fleet_profiles) to pick the backend, or omit it "
+ "for the default. The two are independent: a reviewer may run on the same profile "
+ "as the dev it reviews. The member opens your current directory by default; pass "
@@ -2093,7 +2236,7 @@ public final class FleetMcp {
+ "one. Returns the member's sessionId (use with fleet_send) and paneId (use with "
+ "fleet_stop).",
objectSchema(Map.of(
"role", stringProp("What the member is for: architect, dev or reviewer (default dev)"),
"role", stringProp("What the member is for: architect, dev, hunter, or reviewer (default dev)"),
"profile", stringProp("Which backend to run it on (omit for the default profile)"),
"cwd", stringProp("Working directory for the member (omit to inherit yours)"),
"worktree", Map.of("type", "string", "description", "'true' or a ticket slug — requests an isolated git worktree"),
@@ -2135,9 +2278,10 @@ public final class FleetMcp {
+ "cannot reliably re-identify: some backends (e.g. opencode) resolve it from the "
+ "member's working directory, which only uniquely identifies a member when it "
+ "was spawned into its own fleetd-provisioned worktree (worktree:true/<slug>); a "
+ "member spawned without one shares its directory with others and never reports "
+ "an id, however long it runs (fleetd #249). An empty 'members' "
+ "means no members are spawned; it says nothing about peers. When capacity "
+ "member spawned without one shares its directory with others and never reports "
+ "an id, however long it runs (fleetd #249). An empty 'members' "
+ "means no members are spawned; it says nothing about peers. 'loopHealth' reports "
+ "the RUNNING, STALLED, or STOPPED state of statusPoller and sessionReaper. When capacity "
+ "facts are configured, a 'capacity' row per profile reports 'free' — the "
+ "slots a fresh fleet_spawn on that profile will actually be granted right "
+ "now (max(0, maxLoad - live)), the same check the spawn gate itself runs. A "
@@ -93,25 +93,30 @@ public final class MessageService {
/** Timed out after the message was delivered — the worker is still working. */
TIMED_OUT_WORKING,
/**
* Timed out with no confirmed delivery. Despite the name, this does not mean the message
* is sitting in a queue. {@link #send} reaches this outcome through {@link Injector#cancel},
* whose result tells three routes apart:
* {@link Injector.Cancellation#CANCELLED} means the message was still queued and this call
* removed it, so the target saw nothing and it will not arrive later;
* {@link Injector.Cancellation#NOT_DELIVERED} means nothing was ever sent — the queue was
* cleared because the target never became ready or was abandoned, or the injector's call to
* the target's terminal ({@link dev.ltms.fleet.herdr.AgentControl#send}) failed with a herdr
* error that this codebase already treats as a confirmed absence — so this route too
* establishes that the target saw nothing and it will not arrive later; but
* {@link Injector.Cancellation#ATTEMPTED} (fleetd #551) means that call was made and its
* outcome is unknown. {@code agent.prompt} pastes <em>and submits</em> in one call, so on
* this route the target may hold a complete, already-submitted turn and be working on it
* right now — {@link Outcome#TIMED_OUT_WORKING}'s meaning, reported here as
* {@code TIMED_OUT_QUEUED} only because this caller never observed the pickup. Only
* {@code CANCELLED} and {@code NOT_DELIVERED} establish that the target saw nothing;
* {@code ATTEMPTED} does not.
* Timed out with no confirmed delivery, and the target saw nothing — the message will not
* arrive later, so a caller may resend. {@link #send} reaches this outcome through {@link
* Injector#cancel} reporting one of two routes: {@link Injector.Cancellation#CANCELLED}
* means the message was still queued and this call removed it; {@link
* Injector.Cancellation#NOT_DELIVERED} means nothing was ever sent — the queue was cleared
* because the target never became ready or was abandoned, or the injector's call to the
* target's terminal ({@link dev.ltms.fleet.herdr.AgentControl#send}) failed with a herdr
* error that this codebase already treats as a confirmed absence. A third route,
* {@link Injector.Cancellation#ATTEMPTED}, used to be folded into this same outcome
* (fleetd #571) — it no longer is; see {@link #TIMED_OUT_UNCONFIRMED}.
*/
TIMED_OUT_QUEUED,
/**
* Timed out with delivery unknown. {@link #send} reaches this outcome when {@link
* Injector#cancel} reports {@link Injector.Cancellation#ATTEMPTED} (fleetd #551): the call
* to the target's terminal ({@link dev.ltms.fleet.herdr.AgentControl#send}) was made, but
* this caller never observed whether it reached the pane. {@code agent.prompt} pastes
* <em>and submits</em> in one call, so the target may already hold a complete, submitted
* turn and be working on it right now — the same reality as {@link #TIMED_OUT_WORKING},
* just not confirmed. The message may or may not have arrived. Treat this as neither a
* confirmed delivery nor a confirmed absence: a caller that resends on this outcome risks a
* double delivery — the same brief typed into the pane twice (fleetd #571).
*/
TIMED_OUT_UNCONFIRMED,
/** Another send to this session was in flight for the whole window. */
BUSY,
/**
@@ -317,20 +322,21 @@ public final class MessageService {
*/
private final ConcurrentHashMap<String, Boolean> strandedReplies = new ConcurrentHashMap<>();
/**
* Targets whose last send timed out with {@link Outcome#TIMED_OUT_QUEUED} (CB-640) — {@link
* #send} called {@link Injector#cancel} and got back something other than {@code DELIVERED}.
* That covers three histories, not one: {@link Injector.Cancellation#CANCELLED} — the message
* was still queued and {@code cancel} removed it right there; {@link
* Injector.Cancellation#NOT_DELIVERED} — nothing was ever sent, because the target never became
* ready, was torn down, or the call to its terminal failed with a herdr error this codebase
* already treats as a confirmed absence; or {@link Injector.Cancellation#ATTEMPTED} (fleetd
* #551) — the call to the target's terminal was made and its outcome is unknown, so the target
* may already hold a complete, submitted turn. Only the first two mean the message will not
* arrive later and the target saw nothing; on the third it may already have arrived in full.
* Set where {@link #send} already computes {@code wasDelivered} for that outcome; no queue is
* kept here, only the fact that the send ended with no confirmed delivery. Cleared the same way
* as {@link #strandedReplies}: the next accepted delivery for the target ({@link #send} opening
* a fresh waiter) or a teardown ({@link #abandon}).
* Targets whose last send timed out with no confirmed delivery (CB-640) — {@link #send} called
* {@link Injector#cancel} and got back something other than {@code DELIVERED}. That covers
* three histories, not one: {@link Injector.Cancellation#CANCELLED} — the message was still
* queued and {@code cancel} removed it right there; {@link Injector.Cancellation#NOT_DELIVERED}
* — nothing was ever sent, because the target never became ready, was torn down, or the call to
* its terminal failed with a herdr error this codebase already treats as a confirmed absence; or
* {@link Injector.Cancellation#ATTEMPTED} (fleetd #551) — the call to the target's terminal was
* made and its outcome is unknown, so the target may already hold a complete, submitted turn.
* Only the first two mean the message will not arrive later and the target saw nothing; on the
* third it may already have arrived in full — and the caller sees a different outcome for it
* ({@link Outcome#TIMED_OUT_UNCONFIRMED}, fleetd #571) than for the first two ({@link
* Outcome#TIMED_OUT_QUEUED}). Set where {@link #send} already computes {@code wasDelivered} for
* that outcome; no queue is kept here, only the fact that the send ended with no confirmed
* delivery. Cleared the same way as {@link #strandedReplies}: the next accepted delivery for the
* target ({@link #send} opening a fresh waiter) or a teardown ({@link #abandon}).
*/
private final ConcurrentHashMap<String, Boolean> queuedDeliveries = new ConcurrentHashMap<>();
private final AtomicLong ticketSeq = new AtomicLong();
@@ -417,21 +423,22 @@ public final class MessageService {
/**
* Read-only delegation fact for fleet views (CB-640): {@code target}'s last send timed out
* with no confirmed delivery — the caller saw {@link Outcome#TIMED_OUT_QUEUED} (see the
* {@code TimeoutException} branch of {@link #send}). Despite the method's name, this is not
* proof that a message is sitting in a queue: {@link Injector#cancel} reports this outcome
* through three routes. {@link Injector.Cancellation#CANCELLED} means the message was still
* queued and got removed right there. {@link Injector.Cancellation#NOT_DELIVERED} means
* nothing was ever sent — the target never became ready, was torn down, or the call to its
* terminal failed with a herdr error this codebase already treats as a confirmed absence.
* Only these two routes mean the message will not arrive later. {@link
* with no confirmed delivery — the caller saw {@link Outcome#TIMED_OUT_QUEUED} or {@link
* Outcome#TIMED_OUT_UNCONFIRMED} (fleetd #571; see the {@code TimeoutException} branch of
* {@link #send}). Despite the method's name, this is not proof that a message is sitting in a
* queue: {@link Injector#cancel} reports this outcome through three routes. {@link
* Injector.Cancellation#CANCELLED} means the message was still queued and got removed right
* there. {@link Injector.Cancellation#NOT_DELIVERED} means nothing was ever sent — the target
* never became ready, was torn down, or the call to its terminal failed with a herdr error this
* codebase already treats as a confirmed absence. Only these two routes mean the message will
* not arrive later, and both report {@code TIMED_OUT_QUEUED}. {@link
* Injector.Cancellation#ATTEMPTED} (fleetd #551) means the call to the target's terminal was
* made and its outcome is unknown: {@code agent.prompt} pastes <em>and submits</em> in one
* call, so on this route the target may already hold a complete, submitted turn and be
* working on it right now — it does NOT follow that the target saw nothing. Distinct from
* {@link Outcome#TIMED_OUT_WORKING}, where delivery already happened and only the reply is
* outstanding. Cleared the next time this target's delivery is accepted or the target is
* abandoned — see {@link #queuedDeliveries}.
* working on it right now — it does NOT follow that the target saw nothing, and this route
* reports {@code TIMED_OUT_UNCONFIRMED} instead. Distinct from {@link Outcome#TIMED_OUT_WORKING},
* where delivery already happened and only the reply is outstanding. Cleared the next time this
* target's delivery is accepted or the target is abandoned — see {@link #queuedDeliveries}.
*/
public boolean hasQueuedDelivery(String target) {
return target != null && queuedDeliveries.containsKey(target);
@@ -646,7 +653,7 @@ public final class MessageService {
return switch (o) {
case REPLIED -> "replied";
case COMPLETED_UNREPLIED -> "completion_fallback";
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> "timeout";
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, TIMED_OUT_UNCONFIRMED, BUSY -> "timeout";
case WORKER_FAILED -> "failed";
case BACKEND_EXHAUSTED -> "backend_exhausted";
case STALE_TURN, QUESTION -> null; // not a completed delegation
@@ -975,6 +982,7 @@ public final class MessageService {
} catch (TimeoutException e) {
boolean wasDelivered = delivery.completion().isDone()
&& !delivery.completion().isCompletedExceptionally();
Injector.Cancellation cancellation = null;
if (!wasDelivered) {
if (timeoutCancellationRaceHookForTest != null) {
// Test-only (fleetd #345): see the field's own javadoc.
@@ -982,20 +990,27 @@ public final class MessageService {
}
// The target monitor makes cancellation atomic with onStatus picking this
// Pending up. If pickup won, report TIMED_OUT_WORKING because the text landed.
wasDelivered = injector.cancel(delivery) == Injector.Cancellation.DELIVERED;
cancellation = injector.cancel(delivery);
wasDelivered = cancellation == Injector.Cancellation.DELIVERED;
}
log.debug("send to {} timed out (delivered={})", target, wasDelivered);
if (!wasDelivered) {
Outcome outcome;
if (wasDelivered) {
outcome = Outcome.TIMED_OUT_WORKING;
} else if (cancellation == Injector.Cancellation.ATTEMPTED) {
// fleetd #571: the call to the target's terminal was made and its outcome is
// unknown — the message may already have arrived in full, so this must not
// be reported as TIMED_OUT_QUEUED, which promises it never will.
outcome = Outcome.TIMED_OUT_UNCONFIRMED;
} else {
// CB-640: record that delivery is not confirmed, for fleet health (see
// queuedDeliveries). Whatever injector.cancel() reported above — this call
// removed a still-queued Pending (CANCELLED), an earlier attempt already
// failed with a confirmed absence (NOT_DELIVERED), or an earlier attempt was
// made and its outcome is unknown (ATTEMPTED, fleetd #551 — the message may
// already have arrived in full) — the send ends with no confirmed delivery.
// queuedDeliveries). cancellation is CANCELLED (this call removed a
// still-queued Pending) or NOT_DELIVERED (an earlier attempt already failed
// with a confirmed absence) — both mean the target saw nothing.
queuedDeliveries.put(target, Boolean.TRUE);
outcome = Outcome.TIMED_OUT_QUEUED;
}
return recorded(new Reply(
wasDelivered ? Outcome.TIMED_OUT_WORKING : Outcome.TIMED_OUT_QUEUED, null));
return recorded(new Reply(outcome, null));
} catch (ExecutionException e) {
Throwable cause = e.getCause();
throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause);
@@ -1135,21 +1150,33 @@ public final class MessageService {
}
try {
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(workerSession);
// #282: mirror send()'s registration (:802) so a SECOND fleet_ask inside this same
// resumed turn can re-associate the async ticket with its new turnId via
// markAsyncQuestion — without this, that second ask has no Task to attach to, and
// markAsyncQuestion silently returns null.
Task task = asyncTasksByTurn.get(turnId);
if (task != null) {
asyncTasksByWaiter.put(reply, task);
}
if (!rendezvous.answerAsk(turnId, content)) {
asyncTasksByWaiter.remove(reply);
rendezvous.close(workerSession, reply);
return new Reply(Outcome.STALE_TURN, null); // lapsed between the lookup and the unblock
}
clearAsyncQuestion(turnId, false);
// fleetd #575: this try used to open below, AFTER the Task lookup/registration and the
// STALE_TURN early return that follows it — so that return was covered only by a
// hand-rolled copy of the finally's own cleanup pair, not the finally itself. Widening the
// try up to wrap the registration closes the gap structurally: every exit from here on,
// STALE_TURN included, now runs through the one finally below exactly once, and the
// duplicated pair is gone. This did not fix a live leak — see the ticket: neither
// rendezvous.answerAsk nor clearAsyncQuestion(turnId, false) can throw, so nothing ever
// actually left through the old gap uncovered — but #572 found this exact drift (one
// finally asserted, an identical sibling not) on this same file, and the hand-rolled copy
// was the wrong shape to keep regardless of whether it was ever exercised.
try {
// #282: mirror send()'s registration (:802) so a SECOND fleet_ask inside this same
// resumed turn can re-associate the async ticket with its new turnId via
// markAsyncQuestion — without this, that second ask has no Task to attach to, and
// markAsyncQuestion silently returns null.
Task task = asyncTasksByTurn.get(turnId);
if (task != null) {
asyncTasksByWaiter.put(reply, task);
}
if (answerAskLapseRaceHookForTest != null) {
// Test-only (fleetd #575): see the field's own javadoc.
answerAskLapseRaceHookForTest.run();
}
if (!rendezvous.answerAsk(turnId, content)) {
return new Reply(Outcome.STALE_TURN, null); // lapsed between the lookup and the unblock
}
clearAsyncQuestion(turnId, false);
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
Reply result = new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
// #282: this waiter can resolve with a FRESH question rather than a terminal reply —
@@ -1649,6 +1676,30 @@ public final class MessageService {
this.askTimeoutRaceHookForTest = hook;
}
/**
* Null in production; test seam for fleetd #575 — invoked from {@link #answer}, right after this
* call's own {@code Task} registration and right before its {@code rendezvous.answerAsk(turnId,
* content)} call. A test installs this to complete the SAME turnId's ask directly via {@link
* Rendezvous#answerAsk} from inside that exact window, deterministically reproducing what a
* second, concurrent {@code answer()} call racing to unblock the same ask can otherwise only
* win by timing luck: this call's own {@code askSession(turnId)} lookup at the top already saw
* the ask as open, but by the time it reaches {@code rendezvous.answerAsk} here, the other call
* already completed it (or the worker's own {@code ask()} teardown already closed it) — so this
* call must see {@code false} and return {@link Outcome#STALE_TURN}, exactly the "lapsed between
* the lookup and the unblock" case named at that call site. Proves the #575 fix (widening this
* method's try so a single finally covers this exit) does not change that outcome and still
* cleans this call's own {@code reply} up exactly once.
*/
private volatile Runnable answerAskLapseRaceHookForTest;
/**
* Test-only (fleetd #575): install {@link #answerAskLapseRaceHookForTest}. Package-private so the
* test, in the same package, can reach it without widening any production API.
*/
void setAnswerAskLapseRaceHookForTest(Runnable hook) {
this.answerAskLapseRaceHookForTest = hook;
}
/** A new send must not open a waiter while an async ticket owns this worker's paused turn. */
private boolean hasAsyncQuestion(String target) {
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
@@ -29,8 +29,8 @@ public enum MemberRole {
* <p>Reads the repo and writes analysis. Never commits code and never opens a pull request —
* an architect that starts implementing has stopped doing the job that makes it useful.
*
* <p>Architects are the one member kind declared in config, because a lead addresses the same
* slots across many tickets and needs a stable name for them.
* <p>Architects are the one member kind with live slot binding, because a lead addresses the
* same slots across many tickets and needs a stable name for them.
*/
ARCHITECT,
@@ -43,6 +43,14 @@ public enum MemberRole {
*/
DEV,
/**
* Sweeps an assigned package for defects and reports several ranked findings.
*
* <p>Never changes code, commits, or opens a pull request. A hunt gathers evidence, which can
* include running the build, but leaves every fix to a later implementation unit.
*/
HUNTER,
/**
* Reviews a diff it did not write and reports one structured finding.
*
@@ -59,7 +67,7 @@ public enum MemberRole {
/**
* The {@code fleet:} block that holds this role's pool — {@code architects},
* {@code developers}, {@code reviewers}.
* {@code developers}, {@code hunters}, {@code reviewers}.
*
* <p>Plural, and not always the wire name: the pool of things a {@code dev} may run on reads
* naturally as {@code developers:}. The wire name stays the singular {@code dev}, because that
@@ -69,6 +77,7 @@ public enum MemberRole {
return switch (this) {
case ARCHITECT -> "architects";
case DEV -> "developers";
case HUNTER -> "hunters";
case REVIEWER -> "reviewers";
};
}
@@ -94,6 +94,7 @@ public final class FleetApp {
// .none() (the honest "feature not wired" view) for every constructor that does not pass one.
private final FleetMcp.QuarantineSource quarantine;
private final FleetMcp.OutageSource outage;
private final FleetMcp.LoopHealthSource loopHealth;
private final ObjectMapper mapper = new ObjectMapper();
/**
@@ -155,7 +156,8 @@ public final class FleetApp {
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials) {
this(herdr, memberHerdr, workers, sessions, messages, presence, mcpServlet, auth, metrics,
deliverable, memberCredentials, FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none());
deliverable, memberCredentials, FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
FleetMcp.LoopHealthSource.none());
}
/**
@@ -169,8 +171,18 @@ public final class FleetApp {
public FleetApp(HerdrClient herdr, HerdrClient memberHerdr, PeerLauncher workers, SessionManager sessions,
MessageService messages, MemberPresence presence,
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials,
FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage) {
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials,
FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage) {
this(herdr, memberHerdr, workers, sessions, messages, presence, mcpServlet, auth, metrics, deliverable,
memberCredentials, quarantine, outage, FleetMcp.LoopHealthSource.none());
}
public FleetApp(HerdrClient herdr, HerdrClient memberHerdr, PeerLauncher workers, SessionManager sessions,
MessageService messages, MemberPresence presence,
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials,
FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage,
FleetMcp.LoopHealthSource loopHealth) {
this.herdr = herdr;
this.memberHerdr = memberHerdr != null ? memberHerdr : herdr;
this.workers = workers;
@@ -183,6 +195,7 @@ public final class FleetApp {
this.memberCredentials = memberCredentials != null ? memberCredentials : MemberCredentialPolicyView::absent;
this.quarantine = quarantine != null ? quarantine : FleetMcp.QuarantineSource.none();
this.outage = outage != null ? outage : FleetMcp.OutageSource.none();
this.loopHealth = loopHealth != null ? loopHealth : FleetMcp.LoopHealthSource.none();
}
/** Wire routes onto a fresh, unstarted Javalin instance. Caller starts it. */
@@ -291,15 +304,19 @@ public final class FleetApp {
* spotted by comparing two numbers by eye.
*/
private void healthz(Context ctx) {
HealthzResponse response = healthzResponse(herdr, memberHerdr, loopHealth);
ctx.status(response.status()).json(response.body());
}
record HealthzResponse(int status, Map<String, Object> body) { }
static HealthzResponse healthzResponse(HerdrClient herdr, HerdrClient memberHerdr,
FleetMcp.LoopHealthSource loopHealth) {
JsonNode pong;
try {
pong = herdr.call("ping");
} catch (HerdrException e) {
ctx.status(503).json(Map.of(
"status", "degraded",
"herdr", "unreachable",
"detail", e.getMessage()));
return;
return degradedResponse("unreachable", e.getMessage(), loopHealth);
}
Map<String, Object> body = new LinkedHashMap<>();
body.put("status", "ok");
@@ -311,11 +328,11 @@ public final class FleetApp {
try {
memberPong = memberHerdr.call("ping");
} catch (HerdrException e) {
ctx.status(503).json(Map.of(
return new HealthzResponse(503, Map.of(
"status", "degraded",
"herdr", "member unreachable",
"detail", e.getMessage()));
return;
"detail", e.getMessage(),
"loopHealth", loopHealthView(loopHealth)));
}
int leadProtocol = pong.path("protocol").asInt();
int memberProtocol = memberPong.path("protocol").asInt();
@@ -326,7 +343,20 @@ public final class FleetApp {
body.put("protocolMismatch", true);
}
}
ctx.status(200).json(body);
body.put("loopHealth", loopHealthView(loopHealth));
return new HealthzResponse(200, body);
}
private static Map<String, String> loopHealthView(FleetMcp.LoopHealthSource loopHealth) {
return Map.of("statusPoller", loopHealth.statusPoller().get().name(),
"sessionReaper", loopHealth.sessionReaper().get().name());
}
private static HealthzResponse degradedResponse(String herdr, String detail,
FleetMcp.LoopHealthSource loopHealth) {
return new HealthzResponse(503, Map.of(
"status", "degraded", "herdr", herdr, "detail", detail,
"loopHealth", loopHealthView(loopHealth)));
}
/**
@@ -640,15 +670,31 @@ public final class FleetApp {
}
default -> ctx.status(202).json(Map.of(
"sessionId", id,
// fleetd #571 (ticket comment 17126): no `default` here on purpose. This switch
// is an expression, so the compiler already demands every Outcome constant have
// an arm — adding an 11th constant to Outcome is a compile error here, not a
// silent fall-through. That is exactly the bug this ticket exists to fix:
// `default -> "done"` used to sit here and would have told a REST caller the
// delegation completed for TIMED_OUT_UNCONFIRMED, the one outcome where delivery
// is unknown. REPLIED, COMPLETED_UNREPLIED, QUESTION and STALE_TURN can never
// actually reach this inner switch — the outer switch above always dispatches
// them first — but they still need an arm to keep this switch exhaustive.
"status", switch (reply.outcome()) {
case TIMED_OUT_WORKING -> "working";
case TIMED_OUT_QUEUED -> "queued";
// Delivery here is unknown, not merely still queued — see
// Outcome#TIMED_OUT_UNCONFIRMED's own javadoc.
case TIMED_OUT_UNCONFIRMED -> "unconfirmed";
case BUSY -> "busy";
case WORKER_FAILED -> "failed";
case BACKEND_EXHAUSTED -> "backend_exhausted";
default -> "done"; // unreachable (terminal outcomes handled above)
case REPLIED, COMPLETED_UNREPLIED, QUESTION, STALE_TURN -> "done"; // unreachable
},
"detail", (reply.outcome() == MessageService.Outcome.WORKER_FAILED
"detail", reply.outcome() == MessageService.Outcome.TIMED_OUT_UNCONFIRMED
? "no reply within " + timeout + "ms; delivery is unconfirmed — the "
+ "message may already have reached the worker, so a resend "
+ "risks sending it twice; poll status first"
: (reply.outcome() == MessageService.Outcome.WORKER_FAILED
|| reply.outcome() == MessageService.Outcome.BACKEND_EXHAUSTED)
&& reply.text() != null
? reply.text()
@@ -0,0 +1,84 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
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.Set;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
/**
* fleetd #589 Group 2: {@link Fleetd#claudeCodeLauncher} is the factory that replaced {@code
* main}'s inline {@code new ClaudeCodeLauncher(...)} call, whose 10th argument is the CB-596 {@code
* memberCredentials} policy supplier ({@code () -> config.get().memberCredentials()}). Before this
* ticket that argument was untestable wiring: replacing it with {@code () -> null} compiled with 0
* errors and left every existing test green, since no existing test builds the exact object {@code
* main} wires and then spawns it. {@code memberCredentials} being a lambda is never itself {@code
* null}, so {@link dev.ltms.fleet.member.HerdrPeerLauncher#applyMemberCredentialPolicy} sees {@code
* memberCredentials.get() == null} and silently shadows nothing — reopening the exact CB-592
* exposure gap CB-596's policy closed (gitea issue #82).
*
* <p>This test drives the factory with a real {@link ConfigRef} carrying a {@code
* memberCredentials:} block, spawns through the resulting launcher, and inspects what {@code
* tab.create} actually carried — the same observable surface {@code ClaudeCodeLauncherTest}'s
* {@code everyKnownNameNotAllowedIsShadowedWithTheSentinel} uses for the launcher's own credential
* policy, applied here to prove {@code main}'s wiring reaches it.
*/
class FleetdClaudeCodeLauncherCredentialWiringTest {
private static final String YAML = """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
ltms-local:
baseUrl: http://gx00.gw:8000
model: coder
memberCredentials:
policy: deny-by-default
known:
- GITEA_ACCESS_TOKEN
""";
@SuppressWarnings("unchecked")
private static Map<String, String> startEnv(FakeHerdr herdr) {
return (Map<String, String>) ((Map<String, Object>) herdr.lastCall("tab.create").params()).get("env");
}
@Test
@DisplayName("main's memberCredentials wiring reaches ClaudeCodeLauncher: a known-but-not-allowed "
+ "name is shadowed on spawn")
void memberCredentialsWiringReachesClaudeCodeLauncher(@TempDir Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, YAML);
FleetConfig cfg = FleetConfig.load(file);
ConfigRef config = ConfigRef.fixed(cfg);
FakeHerdr herdr = new FakeHerdr();
ClaudeCodeLauncher launcher = Fleetd.claudeCodeLauncher(new AgentControl(herdr),
new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
cfg.profiles(), cfg, config);
launcher.spawn();
String shadowed = startEnv(herdr).get("GITEA_ACCESS_TOKEN");
assertNotNull(shadowed,
"GITEA_ACCESS_TOKEN is 'known' but not 'allow'-ed in the loaded config — it must be "
+ "explicitly shadowed on spawn; replacing the memberCredentials supplier with "
+ "() -> null at the Fleetd.claudeCodeLauncher call site must fail this "
+ "assertion, since a null policy shadows nothing");
assertFalse(shadowed.isBlank(), "the overlay value must be non-blank");
}
}
@@ -0,0 +1,86 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.inject.ExhaustedPatternLookup;
import dev.ltms.fleet.inject.LiveExhaustedPatterns;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.session.MemberSession;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import java.util.regex.Pattern;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #589 Group 1: {@link Fleetd#exhaustedPatternLookup} is the factory that replaced {@code
* main}'s inline lambda — resolve a herdr {@code target} to its session's profile, then to that
* profile's live {@link LiveExhaustedPatterns#patternFor}. Same shape as {@link
* Fleetd#worktreeBranchLookup} (which {@code FleetdWorktreeBranchLookupTest} pins the same way).
*
* <p>Before this ticket the lambda was built inline in {@code main} and untestable: replacing it
* with {@code target -> null} — the exact shape of {@link ExhaustedPatternLookup#none()} — compiled
* with 0 errors and left every existing test green. Per the ticket, this is the worst consequence
* in the whole #589 sweep: a genuine usage-limit refusal would stop being classified as {@code
* BACKEND_EXHAUSTED} and would be handed back to a waiting {@code fleet_send} as if it were real
* completed work.
*/
class FleetdExhaustedPatternLookupWiringTest {
private static final String YAML = """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
terra:
baseUrl: http://gx00.gw:8000
model: claude-opus-5
exhaustedPattern: "usage limit"
""";
private static MemberSession session(String terminal, String profile) {
return new MemberSession("pane-" + terminal, terminal, profile, MemberRole.DEV,
"/cwd", null, 0L, 0L, 0, MemberSession.State.READY, null, null);
}
private static LiveExhaustedPatterns liveExhaustedPatterns(Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, YAML);
FleetConfig cfg = FleetConfig.load(file);
return new LiveExhaustedPatterns(() -> ConfigRef.fixed(cfg).get().profiles());
}
@Test
@DisplayName("a known target resolves through its session's profile to that profile's live pattern")
void knownTargetResolvesThroughItsProfile(@TempDir Path dir) throws Exception {
LiveExhaustedPatterns patterns = liveExhaustedPatterns(dir);
ExhaustedPatternLookup lookup = Fleetd.exhaustedPatternLookup(
() -> List.of(session("term1", "terra")), patterns);
Pattern resolved = lookup.patternFor("term1");
assertNotNull(resolved,
"the lookup must resolve term1 -> profile 'terra' -> LiveExhaustedPatterns.patternFor("
+ "'terra') — replacing the lambda body with 'target -> null' at the "
+ "Fleetd.exhaustedPatternLookup call site must fail this assertion");
assertTrue(resolved.matcher("the usage limit has been reached").find());
}
@Test
@DisplayName("an unknown target resolves to null, not a thrown exception")
void unknownTargetResolvesToNull(@TempDir Path dir) throws Exception {
LiveExhaustedPatterns patterns = liveExhaustedPatterns(dir);
ExhaustedPatternLookup lookup = Fleetd.exhaustedPatternLookup(
() -> List.of(session("term1", "terra")), patterns);
assertNull(lookup.patternFor("term_stranger"));
}
}
@@ -0,0 +1,62 @@
package dev.ltms.fleet;
import dev.ltms.fleet.inject.ExhaustionSink;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #589 Group 1: {@link Fleetd#forwardingExhaustionSink} is the factory that replaced
* {@code main}'s inline {@code ExhaustionSink.forwardingTo(exhaustionSinkRef::get)} (fleetd #175's
* construction-order break: the adapters need a sink before {@code sessions} exists to build the
* real one). Before this ticket that call site was untestable wiring: replacing the supplier
* argument with a hardcoded {@code () -> ExhaustionSink.none()} compiled with 0 errors and left
* every existing test green, because no test builds the object {@code main} actually wires and
* then mutates the reference afterward — every existing {@code ExhaustionSink.forwardingTo} caller
* in this codebase reads and writes the SAME reference within one test, so a hardcoded-none supplier
* and a correctly-forwarding one are indistinguishable to them.
*
* <p>This test builds the reference, builds the forwarder from it, and only THEN repoints the
* reference at a spy sink — the discriminating order fleetd #175's whole design depends on
* ({@code exhaustionSinkRef} starts at {@code none()} and is repointed once {@code sessions}
* exists). A forwarder that captured a fixed target at construction time (the inert form) can never
* see that later repoint.
*/
class FleetdExhaustionSinkForwardingWiringTest {
@Test
@DisplayName("the forwarder reads the reference live: repointing it AFTER construction is honoured")
void forwarderReadsTheReferenceLiveNotAFixedTargetCapturedAtConstruction() {
AtomicReference<ExhaustionSink> exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none());
ExhaustionSink forwarder = Fleetd.forwardingExhaustionSink(exhaustionSinkRef);
AtomicBoolean spyCalled = new AtomicBoolean(false);
exhaustionSinkRef.set((target, reason, profile) -> spyCalled.set(true));
forwarder.onExhausted("term_x", "usage limit reached", "terra");
assertTrue(spyCalled.get(),
"forwardingExhaustionSink must delegate to whatever exhaustionSinkRef currently "
+ "holds — hardcoding the supplier to () -> ExhaustionSink.none() at the "
+ "Fleetd.forwardingExhaustionSink call site must fail this assertion, "
+ "since the spy set into the reference after construction would never run");
}
@Test
@DisplayName("before any repoint, the forwarder is inert — it starts at none(), not a crash")
void beforeAnyRepointTheForwarderIsInert() {
AtomicReference<ExhaustionSink> exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none());
ExhaustionSink forwarder = Fleetd.forwardingExhaustionSink(exhaustionSinkRef);
AtomicBoolean spyCalled = new AtomicBoolean(false);
forwarder.onExhausted("term_x", "usage limit reached", "terra");
assertFalse(spyCalled.get(), "nothing was ever wired to be called here — this only pins "
+ "that the factory does not throw before a real sink is published");
}
}
@@ -0,0 +1,110 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.ExhaustionSink;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.session.SessionManager;
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.HashMap;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #589 Group 1: {@link Fleetd#publishExhaustionSink} is the factory that replaced {@code
* main}'s previously untested two-statement sequence — build the real {@link
* Fleetd#exhaustionSink}, then {@code exhaustionSinkRef.set(exhaustionSink)}. {@link
* Fleetd#exhaustionSink} itself is already pinned by {@code FleetdExhaustionSinkWarningTest} (its
* log text) — what was NEVER pinned is the {@code .set(...)} call: {@code main} could replace it
* with {@code exhaustionSinkRef.set(ExhaustionSink.none())} and compile with 0 errors, leaving
* every existing test green, because {@link Fleetd#exhaustionSink}'s own tests build and call the
* sink directly, never through the reference {@code main} publishes it into.
*
* <p>This test proves the PUBLISHED reference — not a freshly rebuilt sink — is the one that
* actually quarantines a credential, by reading {@link BackendQuarantine#isQuarantined} after
* calling {@code exhaustionSinkRef.get().onExhausted(...)}, the same object {@link
* Fleetd#forwardingExhaustionSink} forwards to in production.
*/
class FleetdExhaustionSinkPublishWiringTest {
private static final String YAML = """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
terra:
baseUrl: http://gx00.gw:8000
model: claude-opus-5
guard:
offSubscriptionHosts:
- gx00.gw
""";
private static SessionManager emptyRosterSessions() {
FakeHerdr h = new FakeHerdr();
FleetConfig.Profile dummy = new FleetConfig.Profile(
"dummy", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
ClaudeCodeLauncher launcher = new ClaudeCodeLauncher(new AgentControl(h), new WorkspaceControl(h),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(dummy.profile(), dummy), dummy.profile(), _ -> "tok");
// Never acquires a session — publishExhaustionSink's built sink resolves target -> profile
// via the profileHint fallback (fleetd #234), exactly like OpenCodeLauncher's real call
// site does, so this never needs a populated roster.
return new SessionManager(launcher);
}
@Test
@DisplayName("the published reference actually quarantines — not a rebuilt-but-never-set sink")
void publishedReferenceActuallyQuarantines(@TempDir Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, YAML);
FleetConfig cfg = FleetConfig.load(file);
ConfigRef config = ConfigRef.fixed(cfg);
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
Map<String, String> reasonByCredential = new HashMap<>();
AtomicReference<ExhaustionSink> exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none());
Fleetd.publishExhaustionSink(exhaustionSinkRef, emptyRosterSessions(), config, quarantine,
reasonByCredential, cfg);
exhaustionSinkRef.get().onExhausted("term_x", "The usage limit has been reached", "terra");
assertTrue(quarantine.isQuarantined("terra"),
"publishExhaustionSink must repoint exhaustionSinkRef at the REAL sink — "
+ "replacing the .set(...) call with exhaustionSinkRef.set(ExhaustionSink.none()) "
+ "at the Fleetd.publishExhaustionSink call site must fail this assertion, "
+ "since none()'s onExhausted does nothing");
}
@Test
@DisplayName("before publishing, the reference is still inert — no quarantine, no crash")
void beforePublishingTheReferenceIsStillInert(@TempDir Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, YAML);
FleetConfig cfg = FleetConfig.load(file);
ConfigRef config = ConfigRef.fixed(cfg);
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
AtomicReference<ExhaustionSink> exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none());
exhaustionSinkRef.get().onExhausted("term_x", "The usage limit has been reached", "terra");
assertFalse(quarantine.isQuarantined("terra"),
"nothing was published yet — this only pins the starting state the other test's "
+ "assertion actually distinguishes from");
}
}
@@ -0,0 +1,80 @@
package dev.ltms.fleet;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.msg.Rendezvous;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import java.util.function.BiConsumer;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #589 Group 3, site {@code :591}. {@code Fleetd.main} wires {@link
* dev.ltms.fleet.health.FleetHealthMonitor}'s {@code failTarget} callback with {@code
* messages::abandon} — before this ticket that was an inline argument to {@code new
* FleetHealthMonitor(...)}. Measured: replacing it with a no-op {@code BiConsumer} at the call
* site compiles with 0 errors and leaves the full suite green, because nothing else in the tree
* ever drives that specific constructor argument. In production it means a member the monitor
* classifies {@code GONE}/{@code NEVER_READY} never has its pending ticket failed — the caller
* keeps reporting {@code PENDING} for the full 30-minute async timeout instead of the immediate,
* accurate failure CB-580 exists to give it.
*
* <p>This test calls {@link Fleetd#healthFailTarget} directly — never {@code FleetHealthMonitor}
* or {@code main} — against a real {@link MessageService}, using the same {@code sendAsync} +
* {@code poll} observable {@link MessageServiceTest} already relies on to pin {@code
* MessageService.abandon} itself.
*/
class FleetdHealthFailTargetWiringTest {
private static final String T = "term_a";
@Test
@DisplayName("Fleetd.healthFailTarget delegates to the real MessageService.abandon, not a no-op")
void healthFailTargetDelegatesToMessagesAbandon() throws Exception {
FakeHerdr herdr = new FakeHerdr().readText("$ prompt");
AgentControl agents = new AgentControl(herdr);
Rendezvous rendezvous = new Rendezvous();
Injector injector = new Injector(agents);
MessageService messages = new MessageService(agents, injector, rendezvous);
BiConsumer<String, String> failTarget = Fleetd.healthFailTarget(messages);
String ticket = messages.sendAsync(T, "long task");
awaitWaiting(rendezvous);
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase());
failTarget.accept(T, "member unreachable (health monitor)");
MessageService.TaskView view = null;
long deadline = System.currentTimeMillis() + 3000;
while (System.currentTimeMillis() < deadline) {
view = messages.poll(ticket);
if (view.phase() != MessageService.Phase.PENDING) {
break;
}
//noinspection BusyWait
Thread.sleep(10);
}
assertEquals(MessageService.Phase.FAILED, view == null ? null : view.phase(),
"Fleetd.healthFailTarget(messages) must return messages::abandon — replacing it "
+ "with a no-op BiConsumer at the Fleetd.healthFailTarget call site means "
+ "this ticket is never failed and keeps polling as PENDING");
assertTrue(view.detail() != null && view.detail().contains("member unreachable"),
"the failure reason passed to failTarget.accept must reach MessageService.abandon "
+ "and end up in the ticket's detail");
}
private static void awaitWaiting(Rendezvous rendezvous) throws InterruptedException {
long deadline = System.currentTimeMillis() + 2000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter");
}
}
@@ -0,0 +1,100 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.FleetConfig;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import java.util.Map;
import java.util.function.Function;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* fleetd #602 gauge-wiring: {@link Fleetd#leadConfigDirLookup} is the factory {@code Fleetd.main}
* wires into {@code FleetMcp.LeadConfigDirSource} so {@code fleet_list}'s {@code context} row reads
* the transcript directory a lead's OWN profile actually writes to, instead of always falling back
* to the built-in {@code <user.home>/.claude} default (see {@code LeadContextGauge}).
*
* <p>{@code FleetMcpLeadContextGaugeWiringTest} proves the directory this factory returns is what
* actually gets read; this class proves the factory's own matching logic — the same
* {@code fleet.leaders.<name>.profile} link {@code Fleetd.leadSeatLookup} already follows (see
* {@code FleetdLeadSeatLookupTest}), one step further to that profile's own {@code configDir:}.
*/
class FleetdLeadConfigDirLookupTest {
private static FleetConfig.Profile profileWithConfigDir(String name, String configDir) {
return new FleetConfig.Profile(name, null, "claude-sonnet-5", configDir, null, null,
"tab", "fleet", "w #{n}", null, null, null, null, null, null, null,
null, 3, true, null, null, null);
}
private static FleetConfig.Leader leadOnProfile(String profile) {
return new FleetConfig.Leader(profile, "lead: primary", 1, "lead:", 10, "claude", "claude-sonnet-5");
}
@Test
@DisplayName("a lead on a profile that sets configDir resolves to that directory")
void leadOnAProfileWithConfigDirResolvesToIt() {
Map<String, FleetConfig.Profile> profiles = Map.of("opus", profileWithConfigDir("opus", "/mnt/opus-claude"));
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
Function<String, String> lookup = Fleetd.leadConfigDirLookup(() -> profiles, leaders);
assertEquals("/mnt/opus-claude", lookup.apply("primary"));
}
@Test
@DisplayName("changing the config from directory A to directory B changes what the lookup reports")
void configChangeFromDirectoryAToDirectoryBChangesTheAnswer() {
java.util.concurrent.atomic.AtomicReference<Map<String, FleetConfig.Profile>> profilesRef =
new java.util.concurrent.atomic.AtomicReference<>(
Map.of("opus", profileWithConfigDir("opus", "/mnt/dir-a")));
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
Function<String, String> lookup = Fleetd.leadConfigDirLookup(profilesRef::get, leaders);
assertEquals("/mnt/dir-a", lookup.apply("primary"), "must read directory A before the config changes");
profilesRef.set(Map.of("opus", profileWithConfigDir("opus", "/mnt/dir-b")));
assertEquals("/mnt/dir-b", lookup.apply("primary"), "must read directory B once the LIVE config changes — "
+ "a lookup that snapshotted the profile map at construction would still answer directory A here");
}
@Test
@DisplayName("a lead entry with no `profile:` (recognise-only) resolves to null, not a thrown exception")
void recogniseOnlyLeadWithNoProfileResolvesToNull() {
Map<String, FleetConfig.Profile> profiles = Map.of("opus", profileWithConfigDir("opus", "/mnt/opus-claude"));
FleetConfig.Leader recogniseOnly = new FleetConfig.Leader(null, "lead: primary", 1, "lead:", 10,
"claude", "claude-sonnet-5");
Map<String, FleetConfig.Leader> leaders = Map.of("primary", recogniseOnly);
Function<String, String> lookup = Fleetd.leadConfigDirLookup(() -> profiles, leaders);
assertNull(lookup.apply("primary"));
}
@Test
@DisplayName("a lead naming a profile that is not configured resolves to null, not a thrown exception")
void leadOnAnUnconfiguredProfileResolvesToNull() {
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("ghost-profile"));
Function<String, String> lookup = Fleetd.leadConfigDirLookup(Map::of, leaders);
assertNull(lookup.apply("primary"));
}
@Test
@DisplayName("a lead on a profile that sets no configDir override resolves to null")
void leadOnAProfileWithNoConfigDirResolvesToNull() {
Map<String, FleetConfig.Profile> profiles = Map.of("opus", profileWithConfigDir("opus", null));
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
Function<String, String> lookup = Fleetd.leadConfigDirLookup(() -> profiles, leaders);
assertNull(lookup.apply("primary"));
}
@Test
@DisplayName("an unrecognised lead name resolves to null, not a thrown exception")
void unrecognisedLeadNameResolvesToNull() {
Function<String, String> lookup = Fleetd.leadConfigDirLookup(Map::of, Map.of());
assertNull(lookup.apply("ghost-lead"));
}
}
@@ -0,0 +1,98 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.mcp.FleetMcp;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import java.util.Map;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
/**
* fleetd #602 gauge-wiring follow-up (PR #606 review comment 17353): {@code Fleetd.main}'s {@code
* LeadConfigDirSource} local used to be a bare {@code new FleetMcp.LeadConfigDirSource(
* leadConfigDirLookup(...))} built inline, with nothing a test could call directly. Measured on
* that shape: replacing the whole expression with {@code FleetMcp.LeadConfigDirSource.none()} at
* the call site compiled with 0 errors and left the full 1822-test suite green — the daemon could
* be changed to always report every lead's context as {@code UNKNOWN}, forever, and no test would
* notice. That is the same hand-built-vs-config-wired shape as fleetd #561/#248/#426/#562
* ({@code FleetdLoopHealthSourceWiringTest}).
*
* <p>{@code FleetMcpLeadContextGaugeWiringTest} and {@code FleetdLeadConfigDirLookupTest} both
* predate this class and are both still correct — but neither can catch the mutation above. One
* builds its own {@code FleetMcp} and hands it its own {@code LeadConfigDirSource}; the other
* builds its own lookup and calls {@link Fleetd#leadConfigDirLookup} directly. Neither one ever
* calls the thing {@code Fleetd.main} actually calls.
*
* <p>The fix extracts the inline {@code new} into {@link Fleetd#leadConfigDirSource}, a
* package-private factory in the same style as {@link Fleetd#loopHealthSource}/{@link
* Fleetd#capacitySource}/{@link Fleetd#healthCoverageSource} — which is exactly what makes it
* directly callable here. This test calls that factory with real {@link FleetConfig.Profile}/
* {@link FleetConfig.Leader} fixtures (the same shapes {@code FleetdLeadConfigDirLookupTest}
* already uses) and asserts the returned source resolves a real {@code configDir} — a property
* that would be false if {@link Fleetd#leadConfigDirSource} were mutated to {@code return
* FleetMcp.LeadConfigDirSource.none();}. Measured: mutating exactly that line makes
* {@link #resolvesTheRealConfiguredConfigDir()} fail ({@code expected: </mnt/opus-claude> but was:
* <null>}); restoring it makes the whole suite green again.
*
* <p><b>What this class does not and cannot cover.</b> {@code main}'s own line —
* {@code leadConfigDirSource(() -> config.get().profiles(), leaders)} — could itself be swapped
* for a bare {@code FleetMcp.LeadConfigDirSource.none()}, bypassing this factory entirely. Measured:
* that exact mutation compiles with 0 errors and leaves every test in this file, and the full
* 1825-test suite, green. {@link Fleetd#loopHealthSource}'s own wiring test has the identical gap
* for its own one-line call in {@code main} — no test in this codebase calls {@code Fleetd.main}
* far enough to observe which factory call it made. This class narrows the gap from "nothing tests
* the wiring" (the pre-extraction state this ticket found) to "the factory's own logic is pinned,
* and main's call to it is a one-line, visually-verifiable delegation" — the same standard already
* accepted for {@code loopHealthSource}/{@code capacitySource}/{@code healthCoverageSource}.
*/
class FleetdLeadConfigDirSourceWiringTest {
private static FleetConfig.Profile profileWithConfigDir(String name, String configDir) {
return new FleetConfig.Profile(name, null, "claude-sonnet-5", configDir, null, null,
"tab", "fleet", "w #{n}", null, null, null, null, null, null, null,
null, 3, true, null, null, null);
}
private static FleetConfig.Leader leadOnProfile(String profile) {
return new FleetConfig.Leader(profile, "lead: primary", 1, "lead:", 10, "claude", "claude-sonnet-5");
}
@Test
@DisplayName("the returned source resolves the lead's REAL configured configDir, not a hardcoded null")
void resolvesTheRealConfiguredConfigDir() {
Map<String, FleetConfig.Profile> profiles = Map.of("opus", profileWithConfigDir("opus", "/mnt/opus-claude"));
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
FleetMcp.LeadConfigDirSource source = Fleetd.leadConfigDirSource(() -> profiles, leaders);
assertEquals("/mnt/opus-claude", source.configDirFor().apply("primary"),
"the configDirFor function must delegate to the real leadConfigDirLookup — mutating "
+ "Fleetd.leadConfigDirSource's own body to `return FleetMcp.LeadConfigDirSource.none();` "
+ "must fail this assertion (measured: it does — see this test's class javadoc for the "
+ "companion measurement on main's one-line call to this factory, which this assertion "
+ "does not and structurally cannot cover)");
}
@Test
@DisplayName("a lead on a profile with no configDir override still resolves to null, not a crash")
void leadWithNoConfigDirOverrideResolvesToNull() {
Map<String, FleetConfig.Profile> profiles = Map.of("opus", profileWithConfigDir("opus", null));
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
FleetMcp.LeadConfigDirSource source = Fleetd.leadConfigDirSource(() -> profiles, leaders);
assertNull(source.configDirFor().apply("primary"),
"no configDir: override configured ⇒ null, so LeadContextGauge falls back to its own default");
}
@Test
@DisplayName("an unrecognised lead name resolves to null, not a thrown exception")
void unrecognisedLeadNameResolvesToNull() {
FleetMcp.LeadConfigDirSource source = Fleetd.leadConfigDirSource(Map::of, Map.of());
assertNull(source.configDirFor().apply("ghost-lead"));
}
}
@@ -0,0 +1,48 @@
package dev.ltms.fleet;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import java.net.ServerSocket;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #589 Group 3, site {@code :518}. {@code Fleetd.main} passes {@code LeadMailbox::open} as
* the {@link Fleetd.LeadMailboxOpener} argument to {@code openLeadMailbox(...)} — before this
* ticket that method reference was inline at the call site. Measured: replacing it with the inert
* {@code (uri, selfCoordId, prefetch) -> null} compiles with 0 errors and leaves the full suite
* green — {@code FleetdLeadMailboxSelectionTest} drives {@code openLeadMailbox} with its own
* injected opener and never observes what {@code main} itself actually passes.
*
* <p>This test calls {@link Fleetd#leadMailboxOpener} directly and proves it is the real,
* network-attempting opener rather than a stub: pointed at a guaranteed-closed local port, it must
* throw, exactly mirroring {@link FleetdReplyInboxOpenerWiringTest} for the reply-inbox opener.
* The inert form never attempts a connection and returns {@code null} without throwing, so it
* fails this assertion silently.
*/
class FleetdLeadMailboxOpenerWiringTest {
@Test
@DisplayName("Fleetd.leadMailboxOpener is the real LeadMailbox::open, not a stub that never connects")
void leadMailboxOpenerAttemptsARealConnection() throws Exception {
int closedPort;
try (ServerSocket socket = new ServerSocket(0)) {
closedPort = socket.getLocalPort();
} // released immediately — connecting to it now is a guaranteed refusal, not a fluke
Fleetd.LeadMailboxOpener opener = Fleetd.leadMailboxOpener();
IllegalStateException thrown = assertThrows(IllegalStateException.class,
() -> opener.open("amqp://user:pw@127.0.0.1:" + closedPort + "/coord", "coord-1", 50),
"Fleetd.leadMailboxOpener() must be LeadMailbox::open — a real network attempt "
+ "against a genuinely unreachable broker must throw. The inert form "
+ "(uri, selfCoordId, prefetch) -> null never attempts a connection and "
+ "returns null instead of throwing, so it would fail this assertion "
+ "silently.");
assertTrue(thrown.getMessage().contains("cannot connect to AMQP coordination broker"),
"must be LeadMailbox.open's own real failure message, not a different exception "
+ "shape standing in for it");
}
}
@@ -0,0 +1,72 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.inject.LiveExhaustedPatterns;
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 static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #589 Group 1: {@link Fleetd#liveExhaustedPatterns} is the factory that replaced {@code
* main}'s inline {@code new LiveExhaustedPatterns(() -> config.get().profiles())}. Before this
* ticket, that supplier argument was untestable wiring: replacing it with a hardcoded {@code () ->
* Map.of()} compiled with 0 errors and left every existing test green, because {@code
* LiveExhaustedPatternsTest} builds its own instance directly with a hand-supplied map and never
* goes through {@code main}'s call site.
*
* <p>Silently losing this wiring means every profile's {@code exhaustedPattern} stops being
* recognised — {@link Fleetd#exhaustedPatternLookup} would never see a match, and a genuine
* usage-limit refusal would be handed back to a waiting {@code fleet_send} as real completed work.
*/
class FleetdLiveExhaustedPatternsWiringTest {
private static final String YAML = """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
terra:
baseUrl: http://gx00.gw:8000
model: claude-opus-5
exhaustedPattern: "usage limit"
gx:
baseUrl: http://gx00.gw:8000
""";
private static ConfigRef loadConfig(Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, YAML);
FleetConfig cfg = FleetConfig.load(file);
return ConfigRef.fixed(cfg);
}
@Test
@DisplayName("a profile with a configured exhaustedPattern is armed, with a compiled matcher")
void configuredProfileIsArmed(@TempDir Path dir) throws Exception {
LiveExhaustedPatterns patterns = Fleetd.liveExhaustedPatterns(loadConfig(dir));
assertTrue(patterns.armed("terra"),
"the config's live profiles() supplier must reach LiveExhaustedPatterns — hardcoding "
+ "the supplier to () -> Map.of() at the Fleetd.liveExhaustedPatterns call "
+ "site must fail this assertion");
assertTrue(patterns.patternFor("terra").matcher("the usage limit has been reached").find());
}
@Test
@DisplayName("a profile with no configured exhaustedPattern is not armed, but is still resolvable")
void unconfiguredProfileIsNotArmed(@TempDir Path dir) throws Exception {
LiveExhaustedPatterns patterns = Fleetd.liveExhaustedPatterns(loadConfig(dir));
assertFalse(patterns.armed("gx"), "'gx' has no exhaustedPattern configured");
assertNull(patterns.patternFor("gx"));
}
}
@@ -0,0 +1,174 @@
package dev.ltms.fleet;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.inject.LoopWatchdog;
import dev.ltms.fleet.inject.StatusPoller;
import dev.ltms.fleet.mcp.FleetMcp;
import dev.ltms.fleet.peer.Capability;
import dev.ltms.fleet.peer.PeerHandle;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.peer.SpawnRequest;
import dev.ltms.fleet.placement.PlacementDecision;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.session.SessionReaper;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import java.util.List;
import java.util.Set;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* fleetd #562 follow-up (issue comment "HOLD on PR #579"): {@code Fleetd.main}'s {@code loopHealth}
* local used to be a bare {@code new FleetMcp.LoopHealthSource(poller::health, ...)} built inline,
* with nothing a test could call directly. Measured on that shape: replacing {@code
* poller::health} with a constant {@code () -> LoopWatchdog.State.RUNNING} at the call site
* compiled with 0 errors and left all 1771 existing tests green — the daemon could be changed to
* always report the {@link StatusPoller} as {@code RUNNING}, so the watchdog could never fire and
* a stalled poller would be invisible, while every test stayed green. That is exactly the false
* negative this ticket exists to prevent.
*
* <p>The five tests PR #579 added ({@code FleetMcpTest}, {@code FleetAppTest}) all build their own
* {@link FleetMcp.LoopHealthSource} directly with fixed lambdas — they prove the seam ({@code
* LoopHealthSource} reports what it is given) and nothing about what {@code Fleetd.main} actually
* gives it. This is the same hand-built-vs-config-wired shape as fleetd #561/#248/#426.
*
* <p>The fix extracts the inline {@code new} into {@link Fleetd#loopHealthSource}, a package-private
* factory in the same style as {@link Fleetd#capacitySource} and {@link Fleetd#healthCoverageSource}
* — which is exactly what makes it directly callable here. This test calls that factory with real
* {@link StatusPoller}/{@link SessionReaper} instances (never started, so no herdr or git I/O
* happens) and pins each half separately, plus the {@code reaper == null} branch: one invariant
* wired at three places needs three assertions, not one combined check whose non-zero total could
* hide a gap at any single place.
*/
class FleetdLoopHealthSourceWiringTest {
@Test
@DisplayName("the statusPoller half reports the real poller's health, not a hardcoded state")
void statusPollerHalfReflectsThePollersRealHealth() {
// Stopped without ever being started — stop() still marks the watchdog STOPPED. A poller
// that has never reported RUNNING is the discriminating case: if Fleetd.loopHealthSource
// ever hardcoded RUNNING (the exact mutation this test exists to catch), this would fail.
StatusPoller stoppedPoller = freshPoller();
stoppedPoller.stop();
SessionReaper unusedReaper = freshReaper(); // present only to satisfy the signature
FleetMcp.LoopHealthSource source = Fleetd.loopHealthSource(stoppedPoller, unusedReaper);
assertEquals(LoopWatchdog.State.STOPPED, source.statusPoller().get(),
"the statusPoller supplier must delegate to the real poller's health() — "
+ "replacing poller::health with a constant () -> RUNNING at the "
+ "Fleetd.loopHealthSource call site must fail this assertion");
}
@Test
@DisplayName("the sessionReaper half reports the real reaper's health, not a hardcoded state")
void sessionReaperHalfReflectsTheReapersRealHealth() {
StatusPoller unusedPoller = freshPoller(); // present only to satisfy the signature
SessionReaper stoppedReaper = freshReaper();
stoppedReaper.stop();
FleetMcp.LoopHealthSource source = Fleetd.loopHealthSource(unusedPoller, stoppedReaper);
assertEquals(LoopWatchdog.State.STOPPED, source.sessionReaper().get(),
"the sessionReaper supplier must delegate to the real reaper's health() — "
+ "replacing reaper.health() with a constant at the "
+ "Fleetd.loopHealthSource call site must fail this assertion");
}
@Test
@DisplayName("a null reaper (idle ttl not configured) still reports STOPPED, not a crash")
void nullReaperStillReportsStopped() {
// SessionReaper is only constructed when lifecycle.idleTtlSeconds is configured (see the
// `reaper` local in Fleetd.main) — a real deployment routinely passes null here. That null
// check is real behaviour, not a simplification to delete: it must keep reporting STOPPED
// rather than throwing a NullPointerException on the first fleet_list/healthz call.
StatusPoller runningPoller = freshPoller();
FleetMcp.LoopHealthSource source = Fleetd.loopHealthSource(runningPoller, null);
assertEquals(LoopWatchdog.State.STOPPED, source.sessionReaper().get(),
"reaper == null must still report STOPPED, exactly like an intentionally-stopped "
+ "reaper would — do not delete this null check to simplify the wiring");
}
/** Never started, so no herdr call is ever made; freshly constructed reports RUNNING. */
private static StatusPoller freshPoller() {
AgentControl agents = new AgentControl(new FakeHerdr());
return new StatusPoller(agents, new Injector(agents), 1000);
}
/** Never started, so no git/session I/O is ever made; freshly constructed reports RUNNING. */
private static SessionReaper freshReaper() {
return new SessionReaper(new SessionManager(new NeverSpawnsLauncher()), 60, 1000);
}
/**
* Same minimal shape as {@code FleetdBackendErrorSinkTest.NeverSpawnsLauncher} — every method
* throws or returns an empty/no-op value, since a {@link SessionReaper} that is only ever
* constructed and then stopped (never started) never calls any of them.
*/
private static final class NeverSpawnsLauncher implements PeerLauncher {
@Override
public Set<Capability> capabilities() {
return Set.of();
}
@Override
public Set<Capability> capabilitiesFor(String profileName) {
return Set.of();
}
@Override
public PeerHandle spawn(SpawnRequest req) {
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
}
@Override
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
}
@Override
public Set<String> profiles() {
return Set.of();
}
@Override
public String defaultProfile() {
return null;
}
@Override
public String effectiveCwd(SpawnRequest req) {
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
}
@Override
public List<String> parityOverlay(String profileName) {
return List.of();
}
@Override
public List<?> list() {
return List.of();
}
@Override
public int reapOrphanWorkers() {
return 0;
}
@Override
public void stop(String id) {
}
@Override
public boolean clearContext(String id) {
return false;
}
}
}
@@ -0,0 +1,77 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.ExhaustionSink;
import dev.ltms.fleet.member.OpenCodeLauncher;
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 static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
/**
* fleetd #589 Group 2: {@link Fleetd#openCodeLauncher} is the factory that replaced {@code main}'s
* inline {@code new OpenCodeLauncher(...)} call — the {@code opencode} counterpart to {@link
* Fleetd#claudeCodeLauncher}, extracted for the identical reason. Its {@code memberCredentials}
* argument is the same {@code () -> config.get().memberCredentials()} supplier; replacing it with
* {@code () -> null} compiled with 0 errors and left every existing test green before this ticket,
* reopening the same CB-592 exposure gap CB-596's policy closed.
*
* <p>Same observable surface as {@code OpenCodeLauncherTest}'s own {@code memberCredentials} tests:
* spawn through the launcher {@code main} actually wires and inspect what {@code tab.create}
* carried.
*/
class FleetdOpenCodeLauncherCredentialWiringTest {
private static final String YAML = """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
gemini:
kind: opencode
model: google/gemini-2.5-pro
memberCredentials:
policy: deny-by-default
known:
- GITEA_ACCESS_TOKEN
""";
@SuppressWarnings("unchecked")
private static Map<String, String> startEnv(FakeHerdr herdr) {
return (Map<String, String>) ((Map<String, Object>) herdr.lastCall("tab.create").params()).get("env");
}
@Test
@DisplayName("main's memberCredentials wiring reaches OpenCodeLauncher: a known-but-not-allowed "
+ "name is shadowed on spawn")
void memberCredentialsWiringReachesOpenCodeLauncher(@TempDir Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, YAML);
FleetConfig cfg = FleetConfig.load(file);
ConfigRef config = ConfigRef.fixed(cfg);
FakeHerdr herdr = new FakeHerdr();
OpenCodeLauncher launcher = Fleetd.openCodeLauncher(new AgentControl(herdr),
new WorkspaceControl(herdr), cfg.profiles(), cfg, config, ExhaustionSink.none());
launcher.spawn();
String shadowed = startEnv(herdr).get("GITEA_ACCESS_TOKEN");
assertNotNull(shadowed,
"GITEA_ACCESS_TOKEN is 'known' but not 'allow'-ed in the loaded config — it must be "
+ "explicitly shadowed on spawn; replacing the memberCredentials supplier with "
+ "() -> null at the Fleetd.openCodeLauncher call site must fail this "
+ "assertion, since a null policy shadows nothing");
assertFalse(shadowed.isBlank(), "the overlay value must be non-blank");
}
}
@@ -0,0 +1,107 @@
package dev.ltms.fleet;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.mcp.PrimaryRegistry;
import dev.ltms.fleet.msg.InMemoryReplyInbox;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.session.SessionManager;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import java.util.function.Consumer;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #589 Group 3, site {@code :611-631}. {@code Fleetd.main} wires {@code
* sessions.onRelease(...)} with a lambda that calls three collaborators — {@code
* messages.abandon}, {@code replyInbox.release}, and {@code primaryRegistry.forgetDelegation} —
* before this ticket built inline inside {@code main}. Measured: replacing the whole lambda body
* with {@code detail -> { }} compiles with 0 errors and leaves the full suite green, because each
* collaborator is separately tested in isolation ({@code MessageServiceTest}, {@code
* PrimaryRegistryTest}) but nothing before this ticket drove the lambda that calls all three from
* {@code main}.
*
* <p>{@code MessageService.abandon}'s own javadoc documents the consequence: without this,
* tearing a worker down leaves its rendezvous waiter open, so a blocking {@code fleet_send} keeps
* blocking and an async one reports {@code PENDING} for a hardcoded thirty minutes on every {@code
* fleet_stop} and every idle-reap.
*
* <p>This test calls {@link Fleetd#releaseCleanup} directly — never {@code SessionManager} or
* {@code main} — against real {@link MessageService}, {@link InMemoryReplyInbox}, and {@link
* PrimaryRegistry} instances, and asserts each collaborator's own observable effect: the pending
* async ticket transitions to {@code FAILED} (abandon), the inbox no longer owns the target's
* queue (release), and the recorded delegation is forgotten (forgetDelegation).
*/
class FleetdReleaseCleanupWiringTest {
private static final String T = "term_a";
@Test
@DisplayName("Fleetd.releaseCleanup reaches messages.abandon, replyInbox.release, and primaryRegistry.forgetDelegation")
void releaseCleanupReachesAllThreeCollaborators() throws Exception {
FakeHerdr herdr = new FakeHerdr().readText("$ prompt");
AgentControl agents = new AgentControl(herdr);
Rendezvous rendezvous = new Rendezvous();
Injector injector = new Injector(agents);
MessageService messages = new MessageService(agents, injector, rendezvous);
InMemoryReplyInbox replyInbox = new InMemoryReplyInbox();
PrimaryRegistry primaryRegistry = new PrimaryRegistry(null);
// Set up the "before" state each collaborator's own effect is measured against.
String ticket = messages.sendAsync(T, "long task");
awaitWaiting(rendezvous);
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase(),
"sanity: the async ticket is pending before cleanup runs");
replyInbox.own(T);
replyInbox.publish(T, "msg-1", "hello");
assertEquals(1, replyInbox.peek(T).size(),
"sanity: the inbox owns T and holds one message before cleanup runs");
primaryRegistry.recordDelegation(T, "lead-1");
assertEquals("lead-1", primaryRegistry.nudgeTargetFor(T).orElse(null),
"sanity: the delegation is recorded before cleanup runs");
Consumer<SessionManager.ReleaseDetail> cleanup =
Fleetd.releaseCleanup(messages, replyInbox, primaryRegistry);
cleanup.accept(new SessionManager.ReleaseDetail(T, null, null, null, null));
MessageService.TaskView view = null;
long deadline = System.currentTimeMillis() + 3000;
while (System.currentTimeMillis() < deadline) {
view = messages.poll(ticket);
if (view.phase() != MessageService.Phase.PENDING) {
break;
}
//noinspection BusyWait
Thread.sleep(10);
}
assertEquals(MessageService.Phase.FAILED, view == null ? null : view.phase(),
"releaseCleanup must call messages.abandon(...) — an inert detail -> { } lambda "
+ "leaves this ticket PENDING forever");
assertTrue(view.detail() != null && view.detail().contains("released"),
"the abandon reason must say the worker session was released");
assertTrue(replyInbox.peek(T).isEmpty(),
"releaseCleanup must call replyInbox.release(...) — an inert lambda leaves the "
+ "inbox still owning T with its message");
assertTrue(primaryRegistry.nudgeTargetFor(T).isEmpty(),
"releaseCleanup must call primaryRegistry.forgetDelegation(...) — an inert lambda "
+ "leaves the stale delegation in place");
}
private static void awaitWaiting(Rendezvous rendezvous) throws InterruptedException {
long deadline = System.currentTimeMillis() + 2000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter");
}
}
@@ -0,0 +1,49 @@
package dev.ltms.fleet;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import java.net.ServerSocket;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #589 Group 3, site {@code :512}. {@code Fleetd.main} passes {@code AmqpReplyInbox::open}
* as the {@link Fleetd.AmqpOpener} argument to {@code selectReplyInbox(...)} — before this ticket
* that method reference was inline at the call site. Measured: replacing it with the inert {@code
* (uri, prefetch) -> new InMemoryReplyInbox()} compiles with 0 errors and leaves the full suite
* green — {@code FleetdReplyInboxSelectionTest} drives {@code selectReplyInbox} with its own
* injected opener (including one test that passes the real {@code AmqpReplyInbox::open}
* explicitly) and never observes what {@code main} itself actually passes.
*
* <p>This test calls {@link Fleetd#replyInboxOpener} directly and proves it is the real, network-
* attempting opener rather than a stub: pointed at a guaranteed-closed local port, it must throw —
* the same shape {@code FleetdReplyInboxSelectionTest.aRealUnreachableBrokerFallsBackViaTheRealOpener}
* already relies on for {@code AmqpReplyInbox.open} itself. The inert form never attempts a
* connection and never throws, so it fails this assertion silently (by returning normally).
*/
class FleetdReplyInboxOpenerWiringTest {
@Test
@DisplayName("Fleetd.replyInboxOpener is the real AmqpReplyInbox::open, not a stub that never connects")
void replyInboxOpenerAttemptsARealConnection() throws Exception {
int closedPort;
try (ServerSocket socket = new ServerSocket(0)) {
closedPort = socket.getLocalPort();
} // released immediately — connecting to it now is a guaranteed refusal, not a fluke
Fleetd.AmqpOpener opener = Fleetd.replyInboxOpener();
IllegalStateException thrown = assertThrows(IllegalStateException.class,
() -> opener.open("amqp://user:pw@127.0.0.1:" + closedPort + "/vh", 50),
"Fleetd.replyInboxOpener() must be AmqpReplyInbox::open — a real network attempt "
+ "against a genuinely unreachable broker must throw. The inert form "
+ "(uri, prefetch) -> new InMemoryReplyInbox() never attempts a connection "
+ "and never throws, so it would return normally here and fail this "
+ "assertion silently.");
assertTrue(thrown.getMessage().contains("cannot connect to AMQP broker"),
"must be AmqpReplyInbox.open's own real failure message, not a different exception "
+ "shape standing in for it");
}
}
@@ -0,0 +1,286 @@
package dev.ltms.fleet;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.inject.CompletionResolver;
import dev.ltms.fleet.inject.ExhaustedPatternLookup;
import dev.ltms.fleet.inject.ExhaustionSink;
import dev.ltms.fleet.inject.TurnListener;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.msg.TurnToken;
import org.junit.jupiter.api.Test;
import java.util.LinkedHashSet;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #561: {@code Fleetd.turnListener} composes the completion resolver and the session
* manager into one {@link TurnListener}. Each of the four callbacks below used to be two bare,
* unguarded statements (completion first, then sessions) — a throw from the session half used to
* skip nothing <em>after</em> it (there was nothing after it), but nothing enforced that the
* completion half had to come first either, beyond call order in the source. {@code onDelivered}
* had the identical shape and was fixed by #556 (moving its registration off the fan-out
* entirely); these four callbacks cannot be fixed that way, so the fan-out itself is hardened
* instead (see {@link Fleetd#turnListener} and {@code Fleetd.bothMustRun}/{@code
* bothMustRunKeepingSecondResult}).
*
* <p>The invariant under test: <b>a throwing session-listener must not prevent the completion
* resolver from being told the turn ended.</b> Each test below builds the exact production
* composition ({@link Fleetd#turnListener}) from a real {@link CompletionResolver} and a fake
* {@code sessions} half that throws, then asserts the completion half's effect (the captured
* {@link Rendezvous} waiter resolving) happened anyway — never by inspecting call order directly.
*/
class FleetdTurnListenerCompositionTest {
/** Records which callbacks ran and can be told to throw from a chosen one. */
private static final class RecordingSessions implements TurnListener {
final Set<String> called = new LinkedHashSet<>();
private final Set<String> throwing;
RecordingSessions(String... throwingMethods) {
this.throwing = Set.of(throwingMethods);
}
private void maybeThrow(String method) {
called.add(method);
if (throwing.contains(method)) {
throw new IllegalStateException("boom: sessions." + method);
}
}
@Override
public void onTurnComplete(String target) {
maybeThrow("onTurnComplete");
}
@Override
public boolean onTurnCompleteWithPostAction(String target) {
maybeThrow("onTurnCompleteWithPostAction");
return true;
}
@Override
public void onTurnFailed(String target) {
maybeThrow("onTurnFailed");
}
@Override
public void onTurnFailed(String target, String reason) {
maybeThrow("onTurnFailedWithReason");
}
}
private static CompletionResolver newResolver(FakeHerdr herdr, Rendezvous rendezvous) {
return new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none());
}
/**
* A resolver with a controllable clock, and the clock itself, so a test can register a turn
* "delivered" at time 0 and then jump the clock past {@link CompletionResolver#MIN_TURN_NANOS}
* before resolving it — otherwise {@code onTurnComplete}/{@code onTurnCompleteWithPostAction}
* resolve within microseconds of registering in-test, well inside the fleetd#164 floor, and
* get classified as a too-fast crash rather than a real completion. Mirrors the {@code
* LongSupplier} clock-injection pattern fleetd#164's own tests use.
*/
private static CompletionResolver newResolverPastTheFloor(FakeHerdr herdr, Rendezvous rendezvous,
AtomicLong clock) {
LongSupplier nowNanos = clock::get;
return new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), nowNanos);
}
@Test
void onTurnCompleteResolvesTheWaiterEvenWhenTheSessionHalfThrows() throws Exception {
FakeHerdr herdr = new FakeHerdr().readText("⏺ real answer\n❯ ");
Rendezvous rendezvous = new Rendezvous();
AtomicLong clock = new AtomicLong(0);
CompletionResolver completion = newResolverPastTheFloor(herdr, rendezvous, clock);
var waiter = rendezvous.open("term_a");
completion.register("term_a", new TurnToken("term_a", waiter)); // delivered at clock=0
clock.set(CompletionResolver.MIN_TURN_NANOS + 1_000_000); // past the too-fast floor
RecordingSessions sessions = new RecordingSessions("onTurnComplete");
TurnListener composed = Fleetd.turnListener(completion, sessions);
IllegalStateException thrown = assertThrows(IllegalStateException.class,
() -> composed.onTurnComplete("term_a"),
"the session half's throw must still escape the composed listener");
assertEquals("boom: sessions.onTurnComplete", thrown.getMessage());
assertTrue(sessions.called.contains("onTurnComplete"), "the session half must have run");
// onTurnComplete resolves off-thread (a virtual thread) — wait on the waiter itself,
// exactly like MessageServiceTest.completionFallbackIsNeverQueued does.
Rendezvous.Resolution resolution = waiter.get(5, TimeUnit.SECONDS);
assertEquals(Rendezvous.Kind.COMPLETION, resolution.kind(),
"the completion resolver must still have resolved the send, despite the session "
+ "half throwing");
assertEquals("real answer", resolution.text());
}
@Test
void onTurnFailedResolvesTheWaiterAsFailedEvenWhenTheSessionHalfThrows() throws Exception {
FakeHerdr herdr = new FakeHerdr().readText("⏺ crash context\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver completion = newResolver(herdr, rendezvous);
var waiter = rendezvous.open("term_b");
completion.register("term_b", new TurnToken("term_b", waiter));
RecordingSessions sessions = new RecordingSessions("onTurnFailed");
TurnListener composed = Fleetd.turnListener(completion, sessions);
IllegalStateException thrown = assertThrows(IllegalStateException.class,
() -> composed.onTurnFailed("term_b"),
"the session half's throw must still escape the composed listener");
assertEquals("boom: sessions.onTurnFailed", thrown.getMessage());
assertTrue(sessions.called.contains("onTurnFailed"), "the session half must have run");
Rendezvous.Resolution resolution = waiter.get(5, TimeUnit.SECONDS);
assertEquals(Rendezvous.Kind.FAILED, resolution.kind(),
"the completion resolver must still have failed the send, despite the session "
+ "half throwing");
}
@Test
void onTurnFailedWithReasonResolvesTheWaiterEvenWhenTheSessionHalfThrows() throws Exception {
FakeHerdr herdr = new FakeHerdr();
Rendezvous rendezvous = new Rendezvous();
CompletionResolver completion = newResolver(herdr, rendezvous);
var waiter = rendezvous.open("term_c");
completion.register("term_c", new TurnToken("term_c", waiter));
// fleetd #561: the composed listener calls sessions.onTurnFailed(target) — the ONE-arg
// overload — for this reason-carrying event too (matching the pre-existing production
// behaviour: SessionManager never overrides the two-arg overload either), so the
// throwing key here is "onTurnFailed", not a distinct "...WithReason" one.
RecordingSessions sessions = new RecordingSessions("onTurnFailed");
TurnListener composed = Fleetd.turnListener(completion, sessions);
IllegalStateException thrown = assertThrows(IllegalStateException.class,
() -> composed.onTurnFailed("term_c", "worker unreachable"),
"the session half's throw must still escape the composed listener");
assertEquals("boom: sessions.onTurnFailed", thrown.getMessage());
assertTrue(sessions.called.contains("onTurnFailed"), "the session half must have run");
Rendezvous.Resolution resolution = waiter.get(5, TimeUnit.SECONDS);
assertEquals(Rendezvous.Kind.FAILED, resolution.kind(),
"the completion resolver must still have failed the send, despite the session "
+ "half throwing");
assertEquals("worker unreachable", resolution.text(),
"the explicit reason must still reach the resolved send");
}
@Test
void onTurnCompleteWithPostActionResolvesTheWaiterEvenWhenTheSessionHalfThrows() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ answer before reset\n❯ ");
Rendezvous rendezvous = new Rendezvous();
AtomicLong clock = new AtomicLong(0);
CompletionResolver completion = newResolverPastTheFloor(herdr, rendezvous, clock);
var waiter = rendezvous.open("term_d");
completion.register("term_d", new TurnToken("term_d", waiter)); // delivered at clock=0
clock.set(CompletionResolver.MIN_TURN_NANOS + 1_000_000); // past the too-fast floor
RecordingSessions sessions = new RecordingSessions("onTurnCompleteWithPostAction");
TurnListener composed = Fleetd.turnListener(completion, sessions);
IllegalStateException thrown = assertThrows(IllegalStateException.class,
() -> composed.onTurnCompleteWithPostAction("term_d"),
"the session half's throw must still escape the composed listener");
assertEquals("boom: sessions.onTurnCompleteWithPostAction", thrown.getMessage());
assertTrue(sessions.called.contains("onTurnCompleteWithPostAction"), "the session half must have run");
// resolveBeforePostAction is synchronous by design (it must run before the context reset
// can erase the pane) — the waiter is already resolved by the time the throw propagates.
assertTrue(waiter.isDone(), "resolveBeforePostAction is synchronous — the send must "
+ "already be resolved once the composed call returns (by throwing)");
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind());
assertEquals("answer before reset", waiter.getNow(null).text());
}
/**
* The mirror case (comment 17041's acceptance item 2): the completion half throws, and the
* session half still ran. The production {@link CompletionResolver} is deliberately defensive
* (its scrape reads are wrapped in {@code catch (RuntimeException)}, by the same fail-open
* design {@link CompletionResolver#captureBaseline} documents) so it essentially never throws
* synchronously in normal operation — forcing it to do so needs a genuinely-real trigger, not
* a fabricated one. {@code ConcurrentHashMap.get(null)} is that trigger: passing a {@code null}
* target makes {@code onTurnComplete}'s {@code inFlight.get(target)} throw a
* {@link NullPointerException} before it ever starts its resolving thread — a real code path,
* not a contrived one.
*/
@Test
void sessionHalfStillRunsWhenTheCompletionHalfThrowsSynchronously() {
FakeHerdr herdr = new FakeHerdr();
Rendezvous rendezvous = new Rendezvous();
CompletionResolver completion = newResolver(herdr, rendezvous);
RecordingSessions sessions = new RecordingSessions(); // throws from nothing
TurnListener composed = Fleetd.turnListener(completion, sessions);
assertThrows(NullPointerException.class, () -> composed.onTurnComplete(null),
"ConcurrentHashMap.get(null) inside CompletionResolver.onTurnComplete must still "
+ "escape the composed listener");
assertTrue(sessions.called.contains("onTurnComplete"),
"the session half must still have run even though the completion half threw first");
}
/**
* The mirror case above only exercises {@code onTurnComplete}/{@code bothMustRun}. {@code
* onTurnCompleteWithPostAction} is composed through the OTHER helper,
* {@code bothMustRunKeepingSecondResult}, and nothing previously asserted that its session half
* still runs when its completion half throws — that helper could be reverted to the pre-#561
* broken shape (run the session half only if the completion half did not throw) and the suite
* would still stay green. Uses the same real, non-fabricated trigger as the test above: a
* {@code null} target makes {@code resolveBeforePostAction}'s {@code inFlight.get(target)}
* throw a {@link NullPointerException} before {@code resolve} is ever entered.
*/
@Test
void sessionHalfStillRunsWhenTheCompletionHalfThrowsSynchronouslyForPostAction() {
FakeHerdr herdr = new FakeHerdr();
Rendezvous rendezvous = new Rendezvous();
CompletionResolver completion = newResolver(herdr, rendezvous);
RecordingSessions sessions = new RecordingSessions(); // throws from nothing
TurnListener composed = Fleetd.turnListener(completion, sessions);
assertThrows(NullPointerException.class, () -> composed.onTurnCompleteWithPostAction(null),
"ConcurrentHashMap.get(null) inside CompletionResolver.resolveBeforePostAction must "
+ "still escape the composed listener");
assertTrue(sessions.called.contains("onTurnCompleteWithPostAction"),
"the session half must still have run even though the completion half threw first");
}
/**
* {@code bothMustRun} must not just run both halves — it must not DROP a second failure when
* both halves throw. Forces the completion half to throw (the same real {@code
* inFlight.get(null)} NullPointerException trigger used above) while the session half throws a
* distinct {@link IllegalStateException}, and asserts the completion half's throwable is what
* escapes while the session half's throwable survives as a suppressed exception rather than
* being silently discarded.
*/
@Test
void bothFailuresEscapeWhenBothHalvesThrowDistinctExceptions() {
FakeHerdr herdr = new FakeHerdr();
Rendezvous rendezvous = new Rendezvous();
CompletionResolver completion = newResolver(herdr, rendezvous);
RecordingSessions sessions = new RecordingSessions("onTurnComplete");
TurnListener composed = Fleetd.turnListener(completion, sessions);
NullPointerException thrown = assertThrows(NullPointerException.class,
() -> composed.onTurnComplete(null),
"the completion half's throw (NPE from inFlight.get(null)) must be what escapes");
assertTrue(sessions.called.contains("onTurnComplete"), "the session half must still have run");
assertEquals(1, thrown.getSuppressed().length,
"the session half's distinct failure must be recorded as suppressed, not dropped");
assertEquals(IllegalStateException.class, thrown.getSuppressed()[0].getClass());
assertEquals("boom: sessions.onTurnComplete", thrown.getSuppressed()[0].getMessage());
}
}
@@ -0,0 +1,72 @@
package dev.ltms.fleet;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.inject.CompletionResolver;
import dev.ltms.fleet.inject.ExhaustedPatternLookup;
import dev.ltms.fleet.inject.ExhaustionSink;
import dev.ltms.fleet.inject.TurnRegistrar;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.msg.TurnToken;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #589 Group 3, site {@code :505}. {@code Fleetd.main} wires the {@link
* dev.ltms.fleet.inject.Injector}'s {@link TurnRegistrar} with {@code completion::register} —
* before this ticket that was an inline argument to {@code new Injector(...)}, so nothing could
* pin it directly. Measured: replacing it with {@link TurnRegistrar#NOOP} at the call site
* compiles with 0 errors and leaves the full suite green, because {@code onDelivered}'s own {@code
* captureBaseline} performs the identical {@code inFlight} check-and-put a moment later on the
* ordinary delivery path — the two are indistinguishable unless something reads the resolver
* between {@code register} and {@code onDelivered}, or {@code onDelivered} never runs at all (the
* gap fleetd #556 introduced {@link TurnRegistrar} to close).
*
* <p>This test calls {@link Fleetd#turnRegistrar} directly — never {@code Injector} or {@code
* main} — and drives {@link CompletionResolver} entirely through its public API: {@link
* TurnRegistrar#register} followed by {@link CompletionResolver#resolveBeforePostAction}, which
* looks up the same {@code inFlight} entry {@code onTurnComplete} would. With the real registrar,
* that entry exists and the waiter opened by {@link Rendezvous#open} resolves; with {@link
* TurnRegistrar#NOOP} nothing was ever registered, {@code resolveBeforePostAction} finds no
* in-flight turn, and the waiter is left exactly as it started — never done.
*/
class FleetdTurnRegistrarWiringTest {
private static final String T = "term_a";
@Test
@DisplayName("Fleetd.turnRegistrar delegates to the real CompletionResolver, not a no-op")
void turnRegistrarDelegatesToCompletionRegister() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ BUILD GREEN: 391 files\n❯ ");
AgentControl agents = new AgentControl(herdr);
Rendezvous rendezvous = new Rendezvous();
// fleetd#164: an ever-advancing fake clock stands in for the real time a turn would take
// between delivery and resolution, so the MIN_TURN_NANOS "too fast" floor never trips here —
// see MessageServiceTest's resolverClock for the same technique.
AtomicLong clock = new AtomicLong();
CompletionResolver completion = new CompletionResolver(agents, rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(),
() -> clock.addAndGet(CompletionResolver.MIN_TURN_NANOS + 1));
TurnRegistrar registrar = Fleetd.turnRegistrar(completion);
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.open(T);
registrar.register(T, new TurnToken(T, waiter));
// Mirrors what Injector.onStatus's confirmed working->idle boundary would trigger via
// CompletionResolver.onTurnComplete — resolveBeforePostAction is the public, synchronous
// twin of that path and reads the exact same inFlight entry register() must have written.
completion.resolveBeforePostAction(T);
assertTrue(waiter.isDone(),
"Fleetd.turnRegistrar(completion) must return completion::register — replacing it "
+ "with TurnRegistrar.NOOP at the Fleetd.turnRegistrar call site means this "
+ "turn is never registered with CompletionResolver, so resolveBeforePostAction "
+ "finds no in-flight turn and this waiter is never resolved");
}
}
@@ -345,7 +345,7 @@ class FleetConfigTest {
IllegalStateException unknownError = assertThrows(IllegalStateException.class,
() -> FleetConfig.load(unknown).validateCharters());
assertTrue(unknownError.getMessage().contains("architetc"));
assertTrue(unknownError.getMessage().contains("[architect, dev, reviewer]"));
assertTrue(unknownError.getMessage().contains("[architect, dev, hunter, reviewer]"));
}
/**
@@ -812,13 +812,17 @@ class FleetConfigTest {
reviewers:
b:
profile: sonnet
hunters:
c:
profile: sonnet
""");
FleetConfig cfg = FleetConfig.load(f);
assertEquals(List.of("sonnet"), cfg.fleet().profilesFor(MemberRole.DEV));
assertEquals(List.of("sonnet"), cfg.fleet().profilesFor(MemberRole.HUNTER));
assertEquals(List.of("sonnet"), cfg.fleet().profilesFor(MemberRole.REVIEWER));
assertTrue(cfg.fleet().profilesFor(MemberRole.ARCHITECT).isEmpty());
assertEquals(List.of(MemberRole.DEV, MemberRole.REVIEWER), cfg.fleet().rolesConfigured());
assertEquals(List.of(MemberRole.DEV, MemberRole.HUNTER, MemberRole.REVIEWER), cfg.fleet().rolesConfigured());
}
/** The case the two axes exist for: one backend, two roles, and neither is a duplicate. */
@@ -625,6 +625,149 @@ class CompletionResolverTest {
assertEquals(Rendezvous.Kind.REPLY, waiterB.getNow(null).kind());
}
@Test
void aSupersededDoneTurnMustNotEvictItsSuccessorsRegistration() {
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(new FakeHerdr()), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiterA = rendezvous.open("term_a");
var turnA = new CompletionResolver.InFlight(waiterA, null);
assertTrue(rendezvous.resolve("term_a", "A replied"));
rendezvous.close("term_a", waiterA);
var waiterB = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
resolver.resolve("term_a", turnA);
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
"a done turn must not evict B from resolve()'s early return");
}
@Test
void aSupersededExhaustedTurnMustNotEvictItsSuccessorsRegistration() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ usage limit has been reached\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
target -> Pattern.compile("usage limit has been reached"), ExhaustionSink.none());
var waiterA = rendezvous.open("term_a");
var turnA = new CompletionResolver.InFlight(waiterA, null);
rendezvous.close("term_a", waiterA);
var waiterB = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
resolver.resolve("term_a", turnA);
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
"an exhausted turn must not evict B from resolve()'s exhausted branch");
}
@Test
void aSupersededBackendErrorTurnMustNotEvictItsSuccessorsRegistration() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ API Error: 400 invalid request body\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiterA = rendezvous.open("term_a");
var turnA = new CompletionResolver.InFlight(waiterA, null);
rendezvous.close("term_a", waiterA);
var waiterB = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
resolver.resolve("term_a", turnA);
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
"a backend-error turn must not evict B from resolve()'s error branch");
}
@Test
void aSupersededRawExhaustedTurnMustNotEvictItsSuccessorsRegistration() {
FakeHerdr herdr = new FakeHerdr().readText("╭────\nusage limit has been reached");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
target -> Pattern.compile("usage limit has been reached"), ExhaustionSink.none());
var waiterA = rendezvous.open("term_a");
var turnA = new CompletionResolver.InFlight(waiterA, null);
rendezvous.close("term_a", waiterA);
var waiterB = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
resolver.resolve("term_a", turnA);
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
"a raw exhausted turn must not evict B from the raw-scrape exhausted branch");
}
@Test
void aSupersededRawBackendErrorTurnMustNotEvictItsSuccessorsRegistration() {
FakeHerdr herdr = new FakeHerdr().readText("╭────\nAPI Error: 400 invalid request body");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiterA = rendezvous.open("term_a");
var turnA = new CompletionResolver.InFlight(waiterA, null);
rendezvous.close("term_a", waiterA);
var waiterB = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
resolver.resolve("term_a", turnA);
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
"a raw backend-error turn must not evict B from the raw-scrape error branch");
}
@Test
void aSupersededDoneFailedTurnMustNotEvictItsSuccessorsRegistration() {
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(new FakeHerdr()), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiterA = rendezvous.open("term_a");
var turnA = new CompletionResolver.InFlight(waiterA, null);
assertTrue(rendezvous.resolve("term_a", "A replied"));
rendezvous.close("term_a", waiterA);
var waiterB = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
resolver.fail("term_a", turnA);
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
"a done failed turn must not evict B from fail()'s early return");
}
@Test
void aSupersededTooFastBackendErrorTurnMustNotEvictItsSuccessorsRegistration() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ API Error: 400 invalid request body\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 waiterA = rendezvous.open("term_a");
var turnA = new CompletionResolver.InFlight(waiterA, null, clock[0]);
rendezvous.close("term_a", waiterA);
var waiterB = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
clock[0] += CompletionResolver.MIN_TURN_NANOS - 1;
resolver.resolve("term_a", turnA);
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
"a too-fast backend-error turn must not evict B from failTooFast()");
}
private static void assertSuccessorRegistrationSurvives(CompletionResolver resolver, Rendezvous rendezvous,
Object waiterB, String message) {
CompletionResolver.InFlight afterA = resolver.inFlight("term_a");
assertNotNull(afterA, message + " — a one-arg remove(target) would remove B");
assertEquals(waiterB, afterA.waiter(), message + " — the surviving record must belong to B");
assertTrue(rendezvous.resolve("term_a", "B replied"), message + " — B must still resolve normally");
}
// --- CB-578 stage A: backend-exhausted classification ---------------------------------
@Test
@@ -0,0 +1,238 @@
package dev.ltms.fleet.lead;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.io.IOException;
import java.io.RandomAccessFile;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* Ticket "lead context gauge" — fleetd could not see how full a lead's own Claude Code context
* window was, so a lead auto-compacting was always a surprise. {@link LeadContextGauge} reads that
* off the lead's own transcript. Five acceptance properties from the ticket, one test class each
* (plus the three separate UNKNOWN cases the ticket calls out by name):
* <ol>
* <li>{@link #tokenCountTracksTheLastUsageRecordAndChangesWithIt()}</li>
* <li>{@link #compactionCountTracksCompactBoundaryRecordsAndChangesWithIt()}</li>
* <li>{@link #missingFileIsUnknown()}, {@link #unreadableFileIsUnknown()}</li>
* <li>{@link #tornFinalLineFallsBackToTheLastGoodReading()},
* {@link #everyLineUnparseableIsUnknown()} (fleetd #602 gauge-wiring, finding 2)</li>
* <li>{@link #readNeverExceedsTheTailBound()}</li>
* <li>{@link #secondReadInsideTtlDoesNotTouchDiskAgain()}</li>
* </ol>
*
* <p>Every fixture lives under {@code @TempDir} — never the operator's real config directory (see
* the ticket's hard constraint on this).
*
* <p><strong>fleetd #602 gauge-wiring, finding 2.</strong> The property above used to be "a
* truncated or invalid LAST line reports UNKNOWN" — on the theory that a bad last line signals a
* format change. That reasoning did not hold: {@code fleet_list} reads this transcript while Claude
* Code may be mid-write on it, so a torn LAST line is an ordinary race, not a format change, and a
* real format change makes EVERY line unparseable, not only the last one written. So the property
* is now split in two: {@link #tornFinalLineFallsBackToTheLastGoodReading} (the torn-line case must
* NOT destroy a good earlier reading) and its control, {@link #everyLineUnparseableIsUnknown} (only
* when NOTHING in the window parses does this report UNKNOWN).
*/
class LeadContextGaugeTest {
private static final String SESSION_ID = "11111111-1111-1111-1111-111111111111";
/** One line Claude Code would write for a turn with the given live-context total. */
private static String usageLine(long inputTokens, long cacheRead, long cacheCreation) {
return "{\"type\":\"assistant\",\"message\":{\"role\":\"assistant\",\"usage\":{"
+ "\"input_tokens\":" + inputTokens + ","
+ "\"cache_read_input_tokens\":" + cacheRead + ","
+ "\"cache_creation_input_tokens\":" + cacheCreation + "}}}";
}
/** One line Claude Code writes when an auto-compaction happens. */
private static String compactionLine() {
return "{\"type\":\"system\",\"subtype\":\"compact_boundary\","
+ "\"compactMetadata\":{\"preTokens\":250000,\"postTokens\":5000,"
+ "\"trigger\":\"auto\",\"durationMs\":54000}}";
}
/** Lays out {@code <configDir>/projects/<anySlug>/<sessionId>.jsonl} and writes {@code lines}. */
private static String writeTranscript(Path configDir, String sessionId, String... lines) throws IOException {
Path projectDir = configDir.resolve("projects").resolve("some-project-slug");
Files.createDirectories(projectDir);
Path file = projectDir.resolve(sessionId + ".jsonl");
StringBuilder sb = new StringBuilder();
for (String line : lines) {
sb.append(line).append('\n');
}
Files.writeString(file, sb.toString(), StandardCharsets.UTF_8);
return configDir.toString();
}
// --- property 1: live token count tracks the last usage record -----------------------------
@Test
@DisplayName("reported tokens equal the LAST usage record's total, and change when it does")
void tokenCountTracksTheLastUsageRecordAndChangesWithIt(@TempDir Path tmp) throws IOException {
String configDir = writeTranscript(tmp, SESSION_ID,
usageLine(1_000, 0, 0),
usageLine(40_000, 5_000, 3_000)); // last record: 48,000
LeadContextGauge gauge = new LeadContextGauge();
LeadContextGauge.Reading first = gauge.read(configDir, SESSION_ID, "claude");
assertEquals(48_000L, first.tokens(), "must total input+cache_read+cache_creation of the LAST usage record");
assertEquals(LeadContextGauge.State.OK, first.state());
// Change N in the fixture (a fresh session id avoids the cache) -- the reported number
// must change with it, not stay pinned to the first fixture's total.
String otherSession = "22222222-2222-2222-2222-222222222222";
writeTranscript(tmp, otherSession, usageLine(100_000, 50_000, 50_000)); // last record: 200,000
LeadContextGauge.Reading second = gauge.read(configDir, otherSession, "claude");
assertEquals(200_000L, second.tokens());
assertTrue(second.tokens() != first.tokens(), "changing N in the fixture must change the reported number");
}
// --- property 2: compaction count tracks compact_boundary records --------------------------
@Test
@DisplayName("reported compaction count equals K compact_boundary records, and changes with K")
void compactionCountTracksCompactBoundaryRecordsAndChangesWithIt(@TempDir Path tmp) throws IOException {
String sessionTwoCompactions = "33333333-3333-3333-3333-333333333333";
writeTranscript(tmp, sessionTwoCompactions,
usageLine(1_000, 0, 0),
compactionLine(),
usageLine(2_000, 0, 0),
compactionLine(),
usageLine(3_000, 0, 0));
LeadContextGauge gauge = new LeadContextGauge();
LeadContextGauge.Reading twoCompactions = gauge.read(tmp.toString(), sessionTwoCompactions, "claude");
assertEquals(2, twoCompactions.compactions());
String sessionZeroCompactions = "44444444-4444-4444-4444-444444444444";
writeTranscript(tmp, sessionZeroCompactions, usageLine(3_000, 0, 0));
LeadContextGauge.Reading zeroCompactions = gauge.read(tmp.toString(), sessionZeroCompactions, "claude");
assertEquals(0, zeroCompactions.compactions(), "changing K in the fixture must change the reported count");
}
// --- property 3: two separate UNKNOWN cases -------------------------------------------------
@Test
@DisplayName("a missing transcript file reports UNKNOWN with no token number")
void missingFileIsUnknown(@TempDir Path tmp) {
LeadContextGauge gauge = new LeadContextGauge();
LeadContextGauge.Reading reading = gauge.read(tmp.toString(), SESSION_ID, "claude");
assertEquals(LeadContextGauge.State.UNKNOWN, reading.state());
assertNull(reading.tokens());
}
@Test
@DisplayName("an unreadable transcript file reports UNKNOWN with no token number")
void unreadableFileIsUnknown(@TempDir Path tmp) throws IOException {
String configDir = writeTranscript(tmp, SESSION_ID, usageLine(1_000, 0, 0));
Path file = tmp.resolve("projects").resolve("some-project-slug").resolve(SESSION_ID + ".jsonl");
assertTrue(file.toFile().setReadable(false), "test setup: must be able to revoke read permission");
try {
LeadContextGauge gauge = new LeadContextGauge();
LeadContextGauge.Reading reading = gauge.read(configDir, SESSION_ID, "claude");
assertEquals(LeadContextGauge.State.UNKNOWN, reading.state());
assertNull(reading.tokens());
} finally {
file.toFile().setReadable(true); // so @TempDir cleanup can delete it
}
}
@Test
@DisplayName("fleetd #602 finding 2: a torn final line does not destroy a good earlier reading")
void tornFinalLineFallsBackToTheLastGoodReading(@TempDir Path tmp) throws IOException {
// Deliberately NOT via writeTranscript: that helper appends a trailing newline after every
// line, including the last one, which would not model what a write caught mid-flush looks
// like on disk. Claude Code appends and flushes one line at a time, so the torn line here has
// no trailing newline at all -- exactly the shape a read racing an in-progress write sees.
Path projectDir = tmp.resolve("projects").resolve("some-project-slug");
Files.createDirectories(projectDir);
Path file = projectDir.resolve(SESSION_ID + ".jsonl");
String lastCompleteLine = usageLine(1_000, 2_000, 3_000); // last COMPLETE record: 6,000
String tornLine = "{\"type\":\"assistant\",\"message\":{\"role\":\"assistant\",\"usage\":{\"input_tok";
Files.writeString(file, lastCompleteLine + "\n" + tornLine, StandardCharsets.UTF_8);
LeadContextGauge gauge = new LeadContextGauge();
LeadContextGauge.Reading reading = gauge.read(tmp.toString(), SESSION_ID, "claude");
assertEquals(LeadContextGauge.State.OK, reading.state(),
"a torn final line must not turn a good earlier reading into UNKNOWN");
assertEquals(6_000L, reading.tokens(),
"must report the last COMPLETE line's total, not fail the whole read");
}
@Test
@DisplayName("fleetd #602 finding 2 control: when EVERY line is unparseable, the reading IS UNKNOWN")
void everyLineUnparseableIsUnknown(@TempDir Path tmp) throws IOException {
// The real signal a format change gives: not one bad line, but ALL of them. Without this
// control, code that never reports UNKNOWN at all would still pass the torn-line property.
String configDir = writeTranscript(tmp, SESSION_ID,
"{this is not json at all",
"neither is this{{{");
LeadContextGauge gauge = new LeadContextGauge();
LeadContextGauge.Reading reading = gauge.read(configDir, SESSION_ID, "claude");
assertEquals(LeadContextGauge.State.UNKNOWN, reading.state(),
"every line unparseable is the real format-change signal and must still report UNKNOWN");
assertNull(reading.tokens());
}
// --- property 4: the tail bound is actually enforced ----------------------------------------
@Test
@DisplayName("reading a transcript far larger than the tail bound reads no more than the bound")
void readNeverExceedsTheTailBound(@TempDir Path tmp) throws IOException {
int smallBound = 4_096; // exercise the mechanism without writing a multi-MB fixture
Path file = tmp.resolve("big.jsonl");
// One line far bigger than smallBound, repeated, so the whole file is many times the bound.
String line = usageLine(42, 0, 0) + " ".repeat(500);
StringBuilder sb = new StringBuilder();
for (int i = 0; i < 50; i++) {
sb.append(line).append('\n');
}
Files.writeString(file, sb.toString(), StandardCharsets.UTF_8);
long fileLength = Files.size(file);
assertTrue(fileLength > (long) smallBound * 5, "test setup: fixture must genuinely dwarf the bound");
var tail = LeadContextGauge.tailBytes(file, smallBound);
assertEquals(smallBound, tail.bytes().length,
"must read exactly the bound, not the whole " + fileLength + "-byte file — assert on bytes actually read");
}
// --- property 5: the cache TTL means a burst of calls reads the file once -------------------
@Test
@DisplayName("the same (configDir, sessionId) read twice inside the TTL touches disk once")
void secondReadInsideTtlDoesNotTouchDiskAgain(@TempDir Path tmp) throws IOException {
String configDir = writeTranscript(tmp, SESSION_ID, usageLine(1_000, 0, 0));
AtomicLong now = new AtomicLong(0);
LeadContextGauge gauge = new LeadContextGauge(now::get, 5_000);
gauge.read(configDir, SESSION_ID, "claude");
gauge.read(configDir, SESSION_ID, "claude"); // still inside the TTL window
assertEquals(1, gauge.diskReadCount(), "two reads inside the TTL must touch disk once");
now.set(6_000); // past the TTL
gauge.read(configDir, SESSION_ID, "claude");
assertEquals(2, gauge.diskReadCount(), "a read past the TTL must touch disk again");
}
// --- non-Claude peer / unresolved session -----------------------------------------------------
@Test
@DisplayName("a non-claude agent type or a null session id is UNKNOWN, never OK")
void nonClaudePeerOrUnresolvedSessionIsUnknown(@TempDir Path tmp) throws IOException {
String configDir = writeTranscript(tmp, SESSION_ID, usageLine(1_000, 0, 0));
LeadContextGauge gauge = new LeadContextGauge();
assertEquals(LeadContextGauge.State.UNKNOWN, gauge.read(configDir, SESSION_ID, "opencode").state());
assertEquals(LeadContextGauge.State.UNKNOWN, gauge.read(configDir, SESSION_ID, null).state());
assertEquals(LeadContextGauge.State.UNKNOWN, gauge.read(configDir, null, "claude").state());
}
}
@@ -0,0 +1,126 @@
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.lead.LeadContextGauge;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.session.SessionManager;
import io.modelcontextprotocol.spec.McpSchema;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.atomic.AtomicReference;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #602 gauge-wiring: {@code FleetMcp}'s {@code fleet_list} handler shipped calling
* {@code LeadContextGauge.read(null, ...)} unconditionally, so every lead whose profile sets a
* {@code configDir:} override read the wrong transcript directory forever, with no error anywhere —
* {@link LeadContextGauge}'s own 8 tests all passed because none of them exercised THIS wiring; each
* one hands the gauge its own directory directly.
*
* <p>This class proves the directory {@code fleet.leaders.<name>.profile}'s own {@code configDir:}
* names is the one {@code fleet_list}'s {@code context} row actually reads — not merely that SOME
* directory got passed. If the wiring in {@code FleetMcp.leadView}/{@code contextView} is ever
* reverted to a hardcoded {@code null}, {@link #configuredDirectoryDecidesWhichTranscriptIsRead()}
* must go red: both directories below hold a REAL, DIFFERENT token count for the SAME session id,
* so a hardcoded {@code null} (always reading the built-in default, where neither transcript lives)
* would report {@code UNKNOWN} both times instead of the two distinct numbers this test asserts.
*/
class FleetMcpLeadContextGaugeWiringTest {
private static final String LEAD_TERMINAL = "term_lead";
private static final String LEAD_NAME = "opus";
/** {@code FakeHerdr.withAgent} always projects {@code agent_session.value} as {@code sess-<terminalId>}. */
private static final String LEAD_SESSION_ID = "sess-" + LEAD_TERMINAL;
private static ClaudeCodeLauncher workerService(FakeHerdr h) {
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
return new ClaudeCodeLauncher(new AgentControl(h), new WorkspaceControl(h),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> "tok");
}
private static String textOf(McpSchema.CallToolResult r) {
return ((McpSchema.TextContent) r.content().getFirst()).text();
}
/** One line Claude Code would write for a turn with the given live-context total. */
private static String usageLine(long tokens) {
return "{\"type\":\"assistant\",\"message\":{\"role\":\"assistant\",\"usage\":{"
+ "\"input_tokens\":" + tokens + ",\"cache_read_input_tokens\":0,\"cache_creation_input_tokens\":0}}}";
}
/** Lays out {@code <configDir>/projects/<anySlug>/<sessionId>.jsonl} carrying one usage record. */
private static String writeTranscript(Path configDir, String sessionId, long tokens) throws IOException {
Path projectDir = configDir.resolve("projects").resolve("some-project-slug");
Files.createDirectories(projectDir);
Files.writeString(projectDir.resolve(sessionId + ".jsonl"), usageLine(tokens) + "\n", StandardCharsets.UTF_8);
return configDir.toString();
}
@Test
@DisplayName("the CONFIGURED directory decides which transcript is read, and follows a config change")
void configuredDirectoryDecidesWhichTranscriptIsRead(@TempDir Path tmp) throws IOException {
Path dirA = Files.createDirectory(tmp.resolve("dir-a"));
Path dirB = Files.createDirectory(tmp.resolve("dir-b"));
writeTranscript(dirA, LEAD_SESSION_ID, 11_000);
writeTranscript(dirB, LEAD_SESSION_ID, 22_000);
FakeHerdr herdr = new FakeHerdr().withAgent("lead-opus", LEAD_TERMINAL, "wL:p1", "wL:t1");
SessionManager sessions = new SessionManager(workerService(herdr));
LeadContextGauge contextGauge = new LeadContextGauge();
AtomicReference<String> configuredDir = new AtomicReference<>(dirA.toString());
FleetMcp.LeadConfigDirSource source = new FleetMcp.LeadConfigDirSource(name ->
LEAD_NAME.equals(name) ? configuredDir.get() : null);
String firstRead = textOf(listFleet(herdr, sessions, contextGauge, source));
assertTrue(firstRead.contains("\"tokens\":11000"),
"config names directory A ⇒ fleet_list must report A's token count: " + firstRead);
// The config changes to name directory B instead -- the very next read must follow it.
configuredDir.set(dirB.toString());
String secondRead = textOf(listFleet(herdr, sessions, contextGauge, source));
assertTrue(secondRead.contains("\"tokens\":22000"),
"config now names directory B ⇒ fleet_list must report B's token count, not A's stale one: "
+ secondRead);
}
@Test
@DisplayName("a lead whose config names no directory still degrades to the built-in default, never throws")
void aLeadWhoseConfigNamesNoDirectoryStillFallsBackWithoutThrowing() {
FakeHerdr herdr = new FakeHerdr().withAgent("lead-opus", LEAD_TERMINAL, "wL:p1", "wL:t1");
SessionManager sessions = new SessionManager(workerService(herdr));
LeadContextGauge contextGauge = new LeadContextGauge();
McpSchema.CallToolResult res = listFleet(herdr, sessions, contextGauge, FleetMcp.LeadConfigDirSource.none());
assertNotEquals(Boolean.TRUE, res.isError(), "a lead with no configured configDir must degrade, never throw");
String out = textOf(res);
assertTrue(out.contains("\"state\":\"unknown\""), "no transcript under the built-in default ⇒ UNKNOWN: " + out);
assertFalse(out.contains("\"tokens\""), "UNKNOWN must never carry a stale/default token number: " + out);
}
private static McpSchema.CallToolResult listFleet(FakeHerdr herdr, SessionManager sessions,
LeadContextGauge contextGauge, FleetMcp.LeadConfigDirSource leadConfigDirs) {
return FleetMcp.listFleet(workerService(herdr), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.LoopHealthSource.none(), FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
FleetMcp.LeadSeatSource.none(), contextGauge, leadConfigDirs,
Map.of(LEAD_TERMINAL, LEAD_NAME), LEAD_TERMINAL, FleetMcp.CoordinationSource.none(), false);
}
}
@@ -7,7 +7,9 @@ import dev.ltms.fleet.auth.Role;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.inject.LoopWatchdog;
import dev.ltms.fleet.herdr.PaneLocator;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.Injector;
@@ -57,9 +59,10 @@ class FleetMcpTest {
private final FakeHerdr herdr = new FakeHerdr();
private final AgentControl agents = new AgentControl(herdr);
private final Injector injector = new Injector(agents);
private final Rendezvous rendezvous = new Rendezvous();
private final InMemoryReplyInbox inbox = new InMemoryReplyInbox();
private final MessageService messages = new MessageService(agents, new Injector(agents), rendezvous, inbox);
private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
@BeforeEach
void setUp() {
@@ -298,6 +301,35 @@ class FleetMcpTest {
assertTrue(textOf(res).contains("no reply"), "got: " + textOf(res));
}
/**
* fleetd #571 (ticket CORRECTION 5): {@code formatReply}'s {@code TIMED_OUT_UNCONFIRMED} arm is
* the one message whose whole job is to stop a caller retrying a delivery that may already have
* arrived. Pin that its wording is actually distinct from the queued/working arm's retry
* invitation — a mutation that swapped this arm's text for that one still passed every other
* test in this suite, because nothing asserted the specific wording.
*/
@Test
void sendTimesOutWithAnUnconfirmedNoteNotARetryInvitation() throws Exception {
herdr.agentSendFailsWith("send_failed");
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of()));
long deadline = System.currentTimeMillis() + 2000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter");
injector.onStatus(T, AgentStatus.IDLE); // triggers the failing delivery attempt -> ATTEMPTED
McpSchema.CallToolResult res = send.get(5, TimeUnit.SECONDS);
String text = textOf(res);
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
assertTrue(text.contains("delivery unconfirmed"), "got: " + text);
assertFalse(text.contains("retry or poll status"),
"an unconfirmed delivery must not carry the queued/working arm's retry invitation — "
+ "a resend here can double-deliver the same brief: got " + text);
}
@Test
void sendRejectsMissingArgs() {
assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of()).isError());
@@ -1053,6 +1085,36 @@ class FleetMcpTest {
assertFalse(out.contains("quarantinedForSeconds"), out);
}
@Test
void loopHealthReportsStalledStatusPoller() {
String out = loopHealth(LoopWatchdog.State.STALLED, LoopWatchdog.State.RUNNING);
assertTrue(out.contains("\"statusPoller\":\"STALLED\""),
"fleet_list must report a stalled StatusPoller: " + out);
}
@Test
void loopHealthReportsStoppedSessionReaperAsStopped() {
String out = loopHealth(LoopWatchdog.State.RUNNING, LoopWatchdog.State.STOPPED);
assertTrue(out.contains("\"sessionReaper\":\"STOPPED\""),
"fleet_list must report a deliberately stopped SessionReaper as STOPPED, not an alarm: " + out);
}
@Test
void loopHealthReportsRunningStatusPoller() {
String out = loopHealth(LoopWatchdog.State.RUNNING, LoopWatchdog.State.STOPPED);
assertTrue(out.contains("\"statusPoller\":\"RUNNING\""),
"fleet_list must report a running StatusPoller: " + out);
}
private static String loopHealth(LoopWatchdog.State statusPoller, LoopWatchdog.State sessionReaper) {
FakeHerdr herdr = new FakeHerdr();
return textOf(FleetMcp.listFleet(workerService(herdr, "http://gx00.gw:8000", Set.of("gx00.gw")),
new SessionManager(workerService(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"))), null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
new FleetMcp.LoopHealthSource(() -> statusPoller, () -> sessionReaper),
FleetMcp.QuarantineSource.none(), Map.of(), ""));
}
@Test
void capacityIncludesConfiguredProfileWithoutMembers() {
FakeHerdr h = new FakeHerdr();
@@ -1875,7 +1937,7 @@ class FleetMcpTest {
null, null, null, null, null, null);
assertEquals(Boolean.TRUE, res.isError());
assertTrue(textOf(res).contains("architect, dev, reviewer"), textOf(res));
assertTrue(textOf(res).contains("architect, dev, hunter, reviewer"), textOf(res));
}
// ── CB-619 / fleetd #123: a spawn asking for a role its profile has no slot for must be
@@ -556,6 +556,31 @@ class ClaudeCodeLauncherTest {
"no --agent flag when the role has no agent-definition file");
}
@Test
void hunterRoleUsesItsAgentFileAndStopsUsingItWhenRemoved(@TempDir Path cwd) throws Exception {
FakeHerdr herdr = new FakeHerdr();
Path agentFile = Files.createDirectories(cwd.resolve(".claude/agents")).resolve("hunter.md");
Files.writeString(agentFile, "---\nname: hunter\n---\nSweep for defects.");
FleetConfig.Profile cfg = new FleetConfig.Profile(
"sonnet", "http://gx00.gw:8000", null, null, "FLEETD_WORKER_TOKEN",
List.of("claude"), "tab", "fleetd-workers", "w #{n}", null, null, null);
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
svc.spawn(new SpawnRequest("sonnet", cwd.toString(), null, null, null, MemberRole.HUNTER));
List<String> args = spawnedArgs(herdr);
int flag = args.indexOf("--agent");
assertTrue(flag >= 0, "the hunter role reaches its agent-definition file: " + args);
assertEquals("hunter", args.get(flag + 1));
Files.delete(agentFile);
svc.spawn(new SpawnRequest("sonnet", cwd.toString(), null, null, null, MemberRole.HUNTER));
assertFalse(spawnedArgs(herdr).contains("--agent"),
"the hunter role no longer gets an agent when its file is removed");
}
private ClaudeCodeLauncher multiProfile(FakeHerdr herdr) {
FleetConfig.Profile gx10 = new FleetConfig.Profile("gx10", "http://gx10.gw:8000", "coder",
null, "FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers", "w #{n}", null, null, null);
@@ -8,7 +8,10 @@ import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;
import java.lang.reflect.InvocationHandler;
import java.lang.reflect.Method;
import java.lang.reflect.Proxy;
import java.io.IOException;
import java.util.Map;
import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.CountDownLatch;
@@ -20,6 +23,7 @@ import java.util.concurrent.atomic.AtomicReference;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* CB-528 follow-up: {@link AmqpReplyInbox#failPendingPublishesOnRecovery()} must not fail a publish
@@ -197,11 +201,78 @@ class AmqpReplyInboxRecoveryRaceTest {
+ elapsedMillis.get() + "ms");
}
@Test
void publishIOExceptionRemovesThePendingMessageId() throws Exception {
AtomicLong seqCounter = new AtomicLong();
Channel failing = fakeChannel(seqCounter, new CopyOnWriteArrayList<>(), new CopyOnWriteArrayList<>(),
new AtomicReference<>(), new AtomicReference<>(), true);
AmqpReplyInbox inbox = new AmqpReplyInbox(fakeConnection(failing, failing), AmqpReplyInbox.DEFAULT_PREFETCH);
try {
org.junit.jupiter.api.Assertions.assertThrows(IllegalStateException.class,
() -> inbox.publish("worker", "catch", "body"));
assertEquals(0, pendingByMsgId(inbox).size(), "publish IOException must remove its msgId entry");
} finally {
inbox.close();
}
}
@Test
void interruptedPublishRemovesThePendingMessageIdInFinally() throws Exception {
InboxFixture fixture = new InboxFixture();
Thread publish = fixture.startPublish("finally");
fixture.awaitPublished("finally");
publish.interrupt();
publish.join(5_000);
assertEquals(0, pendingByMsgId(fixture.inbox).size(), "publish finally must remove its msgId entry");
fixture.inbox.close();
}
@Test
void confirmResolutionRemovesThePendingMessageId() throws Exception {
AmqpReplyInbox inbox = new InboxFixture().inbox;
try {
seedPending(inbox, 1, "confirm");
invoke(inbox, "resolveConfirm", new Class<?>[] {long.class, boolean.class, boolean.class}, 1L, false, true);
assertEquals(0, pendingByMsgId(inbox).size(), "confirm resolution must remove its msgId entry");
} finally {
inbox.close();
}
}
@Test
void recoverySweepRemovesThePendingMessageId() throws Exception {
AmqpReplyInbox inbox = new InboxFixture().inbox;
try {
seedPending(inbox, 1, "recovery");
inbox.failPendingPublishesOnRecovery();
assertEquals(0, pendingByMsgId(inbox).size(), "recovery sweep must remove its msgId entry");
} finally {
inbox.close();
}
}
@Test
void closeRemovesThePendingMessageId() throws Exception {
AmqpReplyInbox inbox = new InboxFixture().inbox;
seedPending(inbox, 1, "close");
inbox.close();
assertEquals(0, pendingByMsgId(inbox).size(), "close must remove its msgId entry");
}
/** A {@link Proxy}-backed {@link Channel}: only the calls {@link AmqpReplyInbox} actually makes
* are meaningfully implemented; everything else returns a harmless default. */
private static Channel fakeChannel(AtomicLong seqCounter, List<Long> seqOrder, List<String> msgIdOrder,
AtomicReference<ConfirmCallback> ackCallback,
AtomicReference<ConfirmCallback> nackCallback) {
AtomicReference<ConfirmCallback> ackCallback,
AtomicReference<ConfirmCallback> nackCallback) {
return fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback, false);
}
private static Channel fakeChannel(AtomicLong seqCounter, List<Long> seqOrder, List<String> msgIdOrder,
AtomicReference<ConfirmCallback> ackCallback,
AtomicReference<ConfirmCallback> nackCallback, boolean failPublish) {
InvocationHandler handler = (proxy, method, args) -> {
String name = method.getName();
if (name.equals("getNextPublishSeqNo")) {
@@ -210,6 +281,9 @@ class AmqpReplyInboxRecoveryRaceTest {
return value;
}
if (name.equals("basicPublish")) {
if (failPublish) {
throw new IOException("test publish failure");
}
AMQP.BasicProperties props = (AMQP.BasicProperties) args[3];
msgIdOrder.add(props.getMessageId());
return null;
@@ -285,4 +359,60 @@ class AmqpReplyInboxRecoveryRaceTest {
}
return 0;
}
@SuppressWarnings("unchecked")
private static Map<String, Object> pendingByMsgId(AmqpReplyInbox inbox) throws Exception {
var field = AmqpReplyInbox.class.getDeclaredField("pendingByMsgId");
field.setAccessible(true);
return (Map<String, Object>) field.get(inbox);
}
@SuppressWarnings("unchecked")
private static void seedPending(AmqpReplyInbox inbox, long seq, String msgId) throws Exception {
Class<?> pendingType = Class.forName(AmqpReplyInbox.class.getName() + "$Pending");
var constructor = pendingType.getDeclaredConstructor(String.class);
constructor.setAccessible(true);
Object pending = constructor.newInstance(msgId);
var seqField = AmqpReplyInbox.class.getDeclaredField("pendingBySeq");
seqField.setAccessible(true);
((Map<Long, Object>) seqField.get(inbox)).put(seq, pending);
pendingByMsgId(inbox).put(msgId, pending);
}
private static void invoke(AmqpReplyInbox inbox, String name, Class<?>[] types, Object... args) throws Exception {
Method method = AmqpReplyInbox.class.getDeclaredMethod(name, types);
method.setAccessible(true);
method.invoke(inbox, args);
}
private static final class InboxFixture {
final AtomicLong seqCounter = new AtomicLong();
final List<Long> seqOrder = new CopyOnWriteArrayList<>();
final List<String> msgIdOrder = new CopyOnWriteArrayList<>();
final AtomicReference<ConfirmCallback> ackCallback = new AtomicReference<>();
final AtomicReference<ConfirmCallback> nackCallback = new AtomicReference<>();
final AmqpReplyInbox inbox = new AmqpReplyInbox(
fakeConnection(fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback),
fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback)),
AmqpReplyInbox.DEFAULT_PREFETCH);
Thread startPublish(String msgId) {
Thread thread = Thread.ofVirtual().start(() -> {
try {
inbox.publish("worker", msgId, "body");
} catch (IllegalStateException ignored) {
// Interrupting the confirm wait is the path under test.
}
});
return thread;
}
void awaitPublished(String msgId) throws InterruptedException {
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5);
while (!msgIdOrder.contains(msgId) && System.nanoTime() < deadline) {
Thread.sleep(10);
}
assertTrue(msgIdOrder.contains(msgId), "publish did not register " + msgId);
}
}
}
@@ -12,13 +12,18 @@ import org.testcontainers.junit.jupiter.Testcontainers;
import org.testcontainers.utility.DockerImageName;
import java.io.IOException;
import java.lang.reflect.InvocationHandler;
import java.lang.reflect.Method;
import java.lang.reflect.Proxy;
import java.util.List;
import java.util.Map;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -207,6 +212,31 @@ class LeadMailboxTest {
}
}
/**
* {@link LeadMailbox#inspect} opens a third channel after the mailbox's consume and publish
* channels. Limit this connection to three channels, then require a replacement channel after
* the successful inspect. If inspect leaves its probe open, the broker refuses that replacement.
*/
@Test
void inspectClosesItsSuccessfulProbeChannel() throws Exception {
var factory = LeadMailbox.connectionFactory(uri());
factory.setRequestedChannelMax(3);
Connection connection = factory.newConnection();
try (LeadMailbox mailbox = new LeadMailbox(connection, coordId("lead-inspect-probe-close"))) {
LeadChannel.MailboxState state = mailbox.inspect(mailbox.selfCoordId());
assertTrue(state.exists(), "the owned mailbox must be found before checking the probe channel");
Channel replacement = connection.createChannel();
assertNotNull(replacement,
"inspect must close its successful probe channel; the replacement channel was null");
try {
assertTrue(replacement.isOpen(), "the replacement channel must be open after inspect returns");
} finally {
replacement.close();
}
}
}
/**
* fleetd #440: {@code heldDurable()} must be derived from what {@link LeadMailbox#own} actually
* did against the real broker — a durable queue declare plus a manual-ack consumer — not a
@@ -347,6 +377,72 @@ class LeadMailboxTest {
() -> "expected AlreadyClosedException, got: " + thrown);
}
@Test
void publishIOExceptionRemovesThePendingMessageId() throws Exception {
LeadMailbox mailbox = newMailbox(true);
try {
assertThrows(IllegalStateException.class,
() -> mailbox.publish("target", new LeadMessage("catch", "from", "target", "body")));
assertEquals(0, pendingByMsgId(mailbox).size(), "publish IOException must remove its msgId entry");
} finally {
mailbox.close();
}
}
@Test
void interruptedPublishRemovesThePendingMessageIdInFinally() throws Exception {
LeadMailbox mailbox = newMailbox(false);
Thread publish = Thread.ofVirtual().start(() -> {
try {
mailbox.publish("target", new LeadMessage("finally", "from", "target", "body"));
} catch (IllegalStateException ignored) {
// Interrupting the confirm wait is the path under test.
}
});
awaitPending(mailbox, "finally");
publish.interrupt();
publish.join(5_000);
try {
assertEquals(0, pendingByMsgId(mailbox).size(), "publish finally must remove its msgId entry");
} finally {
mailbox.close();
}
}
@Test
void confirmResolutionRemovesThePendingMessageId() throws Exception {
LeadMailbox mailbox = newMailbox(false);
try {
seedPending(mailbox, 1, "confirm");
invoke(mailbox, "resolveConfirm", new Class<?>[] {long.class, boolean.class, boolean.class}, 1L, false, true);
assertEquals(0, pendingByMsgId(mailbox).size(), "confirm resolution must remove its msgId entry");
} finally {
mailbox.close();
}
}
@Test
void recoverySweepRemovesThePendingMessageId() throws Exception {
LeadMailbox mailbox = newMailbox(false);
try {
seedPending(mailbox, 1, "recovery");
mailbox.failPendingPublishesOnRecovery();
assertEquals(0, pendingByMsgId(mailbox).size(), "recovery sweep must remove its msgId entry");
} finally {
mailbox.close();
}
}
@Test
void closeRemovesThePendingMessageId() throws Exception {
LeadMailbox mailbox = newMailbox(false);
seedPending(mailbox, 1, "close");
mailbox.close();
assertEquals(0, pendingByMsgId(mailbox).size(), "close must remove its msgId entry");
}
/** Poll peek until at least one message is held, or ~10s elapse (broker delivery is async). */
@SuppressWarnings("BusyWait")
private static List<LeadMessage> awaitPeek(LeadMailbox inbox) throws InterruptedException {
@@ -372,4 +468,108 @@ class LeadMailboxTest {
}
return state;
}
private static LeadMailbox newMailbox(boolean failPublish) {
AtomicLong sequence = new AtomicLong();
Channel consume = fakeChannel(sequence, false);
Channel publish = fakeChannel(sequence, failPublish);
return new LeadMailbox(fakeConnection(consume, publish), "self");
}
private static Channel fakeChannel(AtomicLong sequence, boolean failPublish) {
InvocationHandler handler = (proxy, method, args) -> {
if (method.getName().equals("getNextPublishSeqNo")) {
return sequence.incrementAndGet();
}
if (method.getName().equals("basicPublish") && failPublish) {
throw new IOException("test publish failure");
}
if (method.getName().equals("equals")) {
return proxy == args[0];
}
if (method.getName().equals("hashCode")) {
return System.identityHashCode(proxy);
}
return defaultValue(method.getReturnType());
};
return (Channel) Proxy.newProxyInstance(LeadMailboxTest.class.getClassLoader(), new Class<?>[] {Channel.class}, handler);
}
private static Connection fakeConnection(Channel first, Channel second) {
AtomicLong calls = new AtomicLong();
InvocationHandler handler = (proxy, method, args) -> {
if (method.getName().equals("createChannel") && (args == null || args.length == 0)) {
return calls.getAndIncrement() == 0 ? first : second;
}
if (method.getName().equals("equals")) {
return proxy == args[0];
}
if (method.getName().equals("hashCode")) {
return System.identityHashCode(proxy);
}
return defaultValue(method.getReturnType());
};
return (Connection) Proxy.newProxyInstance(LeadMailboxTest.class.getClassLoader(), new Class<?>[] {Connection.class}, handler);
}
@SuppressWarnings("unchecked")
private static Map<String, Object> pendingByMsgId(LeadMailbox mailbox) throws Exception {
var field = LeadMailbox.class.getDeclaredField("pendingByMsgId");
field.setAccessible(true);
return (Map<String, Object>) field.get(mailbox);
}
@SuppressWarnings("unchecked")
private static void seedPending(LeadMailbox mailbox, long seq, String msgId) throws Exception {
Class<?> pendingType = Class.forName(LeadMailbox.class.getName() + "$Pending");
var constructor = pendingType.getDeclaredConstructor(String.class);
constructor.setAccessible(true);
Object pending = constructor.newInstance(msgId);
var seqField = LeadMailbox.class.getDeclaredField("pendingBySeq");
seqField.setAccessible(true);
((Map<Long, Object>) seqField.get(mailbox)).put(seq, pending);
pendingByMsgId(mailbox).put(msgId, pending);
}
private static void invoke(LeadMailbox mailbox, String name, Class<?>[] types, Object... args) throws Exception {
Method method = LeadMailbox.class.getDeclaredMethod(name, types);
method.setAccessible(true);
method.invoke(mailbox, args);
}
private static void awaitPending(LeadMailbox mailbox, String msgId) throws Exception {
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5);
while (!pendingByMsgId(mailbox).containsKey(msgId) && System.nanoTime() < deadline) {
Thread.sleep(10);
}
assertTrue(pendingByMsgId(mailbox).containsKey(msgId), "publish did not register " + msgId);
}
private static Object defaultValue(Class<?> type) {
if (!type.isPrimitive() || type == void.class) {
return null;
}
if (type == boolean.class) {
return Boolean.FALSE;
}
if (type == long.class) {
return 0L;
}
if (type == short.class) {
return (short) 0;
}
if (type == byte.class) {
return (byte) 0;
}
if (type == char.class) {
return (char) 0;
}
if (type == double.class) {
return 0.0d;
}
if (type == float.class) {
return 0.0f;
}
return 0;
}
}
@@ -22,6 +22,7 @@ import java.util.concurrent.TimeUnit;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
@@ -505,6 +506,52 @@ class MessageServiceTest {
"an answer to a turn that never existed (or already lapsed) is stale, not a hang");
}
/**
* fleetd #575. {@code answer()}'s STALE_TURN return for "the ask lapsed between the lookup and
* the unblock" ({@code rendezvous.answerAsk(turnId, content)} returning {@code false} even though
* this call's own {@code rendezvous.askSession(turnId)} check at the top saw the ask as open) used
* to be covered only by a hand-rolled copy of the cleanup pair its own finally already runs, not
* the finally itself — a structural gap, closed by widening the try up to cover the registration
* above it. This test pins the exact race with {@link
* MessageService#setAnswerAskLapseRaceHookForTest}, which fires right before this call's own
* {@code rendezvous.answerAsk} call and completes the same turnId's ask directly — reproducing
* what a second, concurrent {@code answer()} winning that race could otherwise only do by timing
* luck. Proves the fix changed no behaviour on this path: it still returns {@code STALE_TURN},
* and this call's own forward waiter is still closed exactly once (not left open, and not closed
* twice — there is now only one cleanup site left to run).
*/
@Test
void answerLosingTheRaceToAnAlreadyAnsweredAskStillReturnsStaleTurnAndCleansUpOnce() throws Exception {
String ticket = messages.sendAsync(T, "long task");
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
String turnId = asking.turnId();
assertNotNull(turnId, "an ASKING view carries the turnId to answer on");
messages.setAnswerAskLapseRaceHookForTest(() -> rendezvous.answerAsk(turnId, "raced in first"));
try {
MessageService.Reply r = messages.answer(turnId, "too late", 500);
assertEquals(MessageService.Outcome.STALE_TURN, r.outcome(),
"an ask already answered by the race must be seen as lapsed, not double-delivered");
} finally {
messages.setAnswerAskLapseRaceHookForTest(null);
}
// Cleanup ran exactly once: the forward waiter THIS call opened is closed, not leaked.
assertNull(rendezvous.currentWaiter(T),
"the forward waiter this answer() call opened must be closed after a STALE_TURN return");
// The worker's own ask() call, unblocked by the hook's direct answerAsk, still completes
// normally — the race this test simulates does not strand it.
MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.AskOutcome.ANSWERED, a.outcome());
assertEquals("raced in first", a.answer());
}
// --- timeout, answer, poll, and lock-contention edges ----------------------------------
@Test
@@ -554,6 +601,30 @@ class MessageServiceTest {
}
}
/**
* fleetd #571 (the acceptance test the ticket was filed for). The worker is idle so the injector
* attempts delivery, but the {@code agent.prompt} call itself fails with a herdr error that is
* not a confirmed absence (not a {@code *_not_found} code) — {@link Injector} marks the Pending
* {@code ATTEMPTED} (fleetd #551), meaning the call was made and whether it reached the pane is
* unknown. Before this fix, {@code send}'s {@code TimeoutException} branch collapsed
* {@code ATTEMPTED} into {@code TIMED_OUT_QUEUED} — a promise that the message will never arrive,
* which may already be false: {@code agent.prompt} pastes and submits in one call.
*/
@Test
void sendTimesOutWithAttemptedDeliveryReportsUnconfirmedNotQueued() throws Exception {
herdr.agentSendFailsWith("send_failed");
CompletableFuture<MessageService.Reply> send =
CompletableFuture.supplyAsync(() -> messages.send(T, "brief", 150));
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // triggers the failing delivery attempt → ATTEMPTED
MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.TIMED_OUT_UNCONFIRMED, r.outcome(),
"an ATTEMPTED delivery must not collapse into TIMED_OUT_QUEUED — the message may "
+ "already have arrived in full, and TIMED_OUT_QUEUED promises it never will");
assertNull(r.text());
}
@Test
void answerTimesOutWhenTheResumedWorkerNeverReplies() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
@@ -579,6 +650,183 @@ class MessageServiceTest {
assertEquals("config.yaml", a.answer());
}
// --- fleetd #572: answer() must release the session lock on EVERY exit, not just return it -----
//
// answer()'s session lock is released in an outer `finally` (MessageService.java:1217-1219) that
// wraps its whole body. Removing that one line survives the entire suite: coverage is not the
// gap, assertion is — every existing answer() test checks the RETURN VALUE, never that the lock
// it took is reacquirable afterward. If it is not, the session is wedged forever with no
// exception and no log line. Each test below drives answer() through one of its four exits
// (normal reply, TIMED_OUT_WORKING, ExecutionException, InterruptedException) and then proves
// reacquisition the only way that actually proves it: a bounded follow-up `send` on the SAME
// session must not come back BUSY. A `send` only reports BUSY when `tryLock` itself timed out
// (MessageService.java:926-928) — every other early return in `send` still passes through its own
// lock-acquired `try`, so a non-BUSY probe result is specifically evidence the lock was free.
/** The REPLIED exit (the happy path) — the resumed worker's real {@code fleet_reply} arrives. */
@Test
void answerReleasesTheSessionLockAfterANormalReply() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.Reply q = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
ask.get(5, TimeUnit.SECONDS); // worker resumed with the answer
awaitWaiting(); // the answering call has (re)opened its own forward waiter
assertTrue(rendezvous.resolve(T, "done"), "the worker's final reply resolves the answering send");
MessageService.Reply done = answer.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.REPLIED, done.outcome());
MessageService.Reply probe = messages.send(T, "probe after normal reply", 300);
assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(),
"the session lock must be released after a normal REPLIED answer(), or this bounded "
+ "follow-up send would come back BUSY instead of timing out on its own work");
}
/**
* The {@code TIMED_OUT_WORKING} exit. Same setup as {@link #answerTimesOutWhenTheResumedWorkerNeverReplies}
* (which only checks {@code answer()}'s return value — exactly the assertion the fleetd #572
* mutation survives), plus the reacquisition proof that test does not make.
*
* <p>{@code answer()} MUST run on its own thread here, not the test's main thread: {@code lock}
* is a {@link java.util.concurrent.locks.ReentrantLock}, so a probe {@code send} issued by the
* SAME thread that (under the mutation) still "holds" it would reenter for free and report a
* false pass — reentrancy, not release. Measured while writing this test: with {@code answer()}
* called inline, this test stayed green under the mutation while its three siblings correctly
* went red.
*/
@Test
void answerReleasesTheSessionLockAfterATimedOutWorkingReturn() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.Reply q = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
// The primary answers, unblocking the worker; the worker never sends its follow-up
// fleet_reply, so the answering call rides out its short window as still-working.
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 200));
MessageService.Reply answered = answer.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answered.outcome(),
"an answered worker that never replies times out as still working");
ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer
MessageService.Reply probe = messages.send(T, "probe after timed-out-working", 300);
assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(),
"the session lock must be released after a TIMED_OUT_WORKING answer(), or this bounded "
+ "follow-up send would come back BUSY instead of timing out on its own work");
}
/**
* The {@code ExecutionException} exit. Forces it directly on answer()'s own reopened forward
* waiter — {@code rendezvous.currentWaiter(T)} is the exact {@code CompletableFuture} its
* {@code reply.get(...)} is blocked on — rather than trying to make a real worker fail, since the
* failure mode under test is in answer()'s own wait, not in how it got triggered.
*/
@Test
void answerReleasesTheSessionLockWhenTheReplyFutureFailsExceptionally() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.Reply q = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
String turnId = q.turnId();
java.util.concurrent.atomic.AtomicReference<Throwable> caught =
new java.util.concurrent.atomic.AtomicReference<>();
Thread answerer = new Thread(() -> {
try {
messages.answer(turnId, "config.yaml", 5000);
caught.set(new AssertionError("expected answer() to throw"));
} catch (Throwable t) {
caught.set(t);
}
});
answerer.start();
awaitWaiting(); // answer() re-opened its forward waiter and is about to block on it
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(T);
assertNotNull(waiter, "answer() should have registered a forward waiter for the worker session");
waiter.completeExceptionally(new RuntimeException("boom"));
answerer.join(5000);
assertFalse(answerer.isAlive(), "answer() should have left by throwing once its reply future failed");
Throwable thrown = caught.get();
assertNotNull(thrown, "answer() should have thrown");
assertTrue(thrown instanceof RuntimeException, "a RuntimeException cause is rethrown as-is: " + thrown);
assertEquals("boom", thrown.getMessage());
ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer
MessageService.Reply probe = messages.send(T, "probe after exceptional failure", 300);
assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(),
"the session lock must be released when answer()'s reply future fails exceptionally, "
+ "or this bounded follow-up send would come back BUSY instead of timing out on its own work");
}
/**
* The {@code InterruptedException} exit. Same shape as the {@code ExecutionException} case above,
* but the thread blocked in {@code reply.get(...)} is interrupted instead of the future failing.
*/
@Test
void answerReleasesTheSessionLockWhenInterrupted() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.Reply q = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
String turnId = q.turnId();
java.util.concurrent.atomic.AtomicReference<Throwable> caught =
new java.util.concurrent.atomic.AtomicReference<>();
Thread answerer = new Thread(() -> {
try {
messages.answer(turnId, "config.yaml", 5000);
caught.set(new AssertionError("expected answer() to throw"));
} catch (Throwable t) {
caught.set(t);
}
});
answerer.start();
awaitWaiting(); // answer() re-opened its forward waiter and is about to block on it
answerer.interrupt();
answerer.join(5000);
assertFalse(answerer.isAlive(), "answer() should have left by throwing once its thread was interrupted");
Throwable thrown = caught.get();
assertNotNull(thrown, "answer() should have thrown");
assertTrue(thrown instanceof IllegalStateException,
"an interrupted wait is wrapped in IllegalStateException: " + thrown);
ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer
MessageService.Reply probe = messages.send(T, "probe after interruption", 300);
assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(),
"the session lock must be released when answer()'s wait is interrupted, or this bounded "
+ "follow-up send would come back BUSY instead of timing out on its own work");
}
@Test
void pollReturnsNullForAnUnknownTicket() {
assertNull(messages.poll("task-999999"), "a ticket that was never minted is unknown");
@@ -10,12 +10,13 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
class MemberRoleTest {
@Test
void theThreeRolesAreArchitectDevAndReviewer() {
assertEquals(3, MemberRole.values().length,
void theFourRolesAreArchitectDevHunterAndReviewer() {
assertEquals(4, MemberRole.values().length,
"a new role changes the charter, the role file, the skill and the authz row — "
+ "adding one is a deliberate act, so this count is meant to fail first");
assertEquals("architect", MemberRole.ARCHITECT.wireName());
assertEquals("dev", MemberRole.DEV.wireName());
assertEquals("hunter", MemberRole.HUNTER.wireName());
assertEquals("reviewer", MemberRole.REVIEWER.wireName());
}
@@ -30,6 +31,7 @@ class MemberRoleTest {
void parseIsCaseInsensitiveAndTrimsSurroundingSpace() {
assertSame(MemberRole.ARCHITECT, MemberRole.parse("Architect"));
assertSame(MemberRole.DEV, MemberRole.parse(" DEV "));
assertSame(MemberRole.HUNTER, MemberRole.parse("HuNtEr"));
assertSame(MemberRole.REVIEWER, MemberRole.parse("ReViEwEr"));
}
@@ -38,7 +40,7 @@ class MemberRoleTest {
IllegalArgumentException e =
assertThrows(IllegalArgumentException.class, () -> MemberRole.parse("archtiect"));
assertTrue(e.getMessage().contains("archtiect"), e.getMessage());
assertTrue(e.getMessage().contains("architect, dev, reviewer"),
assertTrue(e.getMessage().contains("architect, dev, hunter, reviewer"),
"a typo in config should be fixable from the message alone: " + e.getMessage());
}
@@ -8,6 +8,7 @@ import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.inject.LoopWatchdog;
import dev.ltms.fleet.inject.StatusPoller;
import dev.ltms.fleet.inject.MemberPresence;
import dev.ltms.fleet.msg.InMemoryReplyInbox;
@@ -165,6 +166,29 @@ class FleetAppTest {
assertEquals("degraded", mapper.readTree(res.body()).get("status").asText());
}
@Test
void healthzKeepsOkStatusAndReportsLoopHealthInItsBody() {
FleetMcp.LoopHealthSource loops = new FleetMcp.LoopHealthSource(
() -> LoopWatchdog.State.RUNNING, () -> LoopWatchdog.State.STOPPED);
FleetApp.HealthzResponse ok = FleetApp.healthzResponse(new FakeHerdr(), new FakeHerdr(), loops);
assertEquals(200, ok.status(), "a healthy herdr must keep /healthz at 200 regardless of loop states");
assertEquals(Map.of("statusPoller", "RUNNING", "sessionReaper", "STOPPED"), ok.body().get("loopHealth"),
"the /healthz body must report each loop state without making STOPPED an alarm");
}
@Test
void healthzKeepsDegradedStatusAndReportsLoopHealthInItsBody() {
FleetMcp.LoopHealthSource loops = new FleetMcp.LoopHealthSource(
() -> LoopWatchdog.State.RUNNING, () -> LoopWatchdog.State.STOPPED);
FleetApp.HealthzResponse degraded = FleetApp.healthzResponse(new FakeHerdr().healthy(false),
new FakeHerdr(), loops);
assertEquals(503, degraded.status(),
"an unreachable herdr must keep /healthz at 503 regardless of loop states");
assertEquals(Map.of("statusPoller", "RUNNING", "sessionReaper", "STOPPED"),
degraded.body().get("loopHealth"), "the degraded /healthz body must retain loop states");
}
@Test
void sessionsMapsWorkspaceList() throws Exception {
int port = startHealthy();
@@ -623,6 +647,28 @@ class FleetAppTest {
assertTrue(herdr.called("agent.prompt"), "message was injected");
}
/**
* fleetd #571: the worker is idle, so the poller attempts delivery, but the {@code agent.prompt}
* call itself fails with a herdr error that is not a confirmed absence — {@link
* dev.ltms.fleet.inject.Injector} marks this {@code ATTEMPTED}, meaning the call was made and
* whether it reached the pane is unknown. {@code writeReply}'s default arm must map this to its
* own {@code "unconfirmed"} status, not silently fall through to {@code "done"} (which would
* claim the delegation completed) nor collapse into {@code "queued"} (which would claim the
* message will never arrive, when it may already be sitting in the pane).
*/
@Test
void messageTimesOutUnconfirmedWhenDeliveryAttemptFails() throws Exception {
FakeHerdr herdr = new FakeHerdr().agentStatus("idle").agentSendFailsWith("send_failed");
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
HttpResponse<String> res = postMessage(port, "{\"content\":\"hi\",\"timeoutMs\":250}");
assertEquals(202, res.statusCode());
JsonNode body = mapper.readTree(res.body());
assertEquals("unconfirmed", body.get("status").asText(),
"an ATTEMPTED delivery must report its own status, not \"queued\" or \"done\"");
assertTrue(herdr.called("agent.prompt"), "delivery must have been attempted");
}
@Test
void messageRejectsBlankContent() throws Exception {
int port = startHealthy();
+54 -3
View File
@@ -169,7 +169,55 @@ hash256() {
# on PATH). "absent" must never be the answer for a file that exists — that conflation, on Linux,
# was the whole defect this ticket fixes.
jar_id() { local f="${1:-$JAR}"; [ -f "$f" ] && hash256 "$f" || echo "absent"; }
running_pid() { pgrep -f "$PATTERN" || true; }
# fleetd #593 — `pgrep -f "$PATTERN"` matches ANY process whose full command line CONTAINS the
# pattern text, and that is not the same thing as "is the daemon". A shell that merely embeds the
# pattern as literal text — a human typing this exact investigation by hand, an ssh-shaped
# `sh -c '...; ...'`, a pipeline, or any other non-exec'ing shell that never replaced itself with
# the pattern-holding command — still shows up in that match, and it is the INSTRUMENT, not the
# daemon. Measured live on this Mac: `sh -c 'echo "target/fleetd.jar" >/dev/null; sleep 30' &`
# leaves a real `sh` process alive (it forks for the `sleep`, it does not exec into it) whose own
# `ps -o args` is `sh -c echo "target/fleetd.jar" >/dev/null; sleep 30` — `pgrep -f "$PATTERN"`
# matches that line right alongside the real `java -jar target/fleetd.jar` process. `pgrep -c`
# (an in-one-call count) does not exist on BSD/macOS at all, so this cannot be fixed by switching
# pgrep flags — it has to filter what pgrep already found, after the fact, in a way that still
# runs on BSD.
#
# fleetd #593 CORRECTION 1 — the first cut of this filter kept everything whose `comm` was NOT a
# shell name (a denylist: sh/bash/zsh/dash/ksh). Two holes in that, both the same false-positive
# shape the ticket exists to remove in the first place:
# 1. a pid `pgrep` just listed can exit before the `ps -o comm=` lookup runs; on a gone pid `ps`
# prints nothing, `comm` ends up empty, and an empty string matches none of the denied shell
# names — so a pid that no longer exists was still counted.
# 2. the denylist only knows the shells someone thought to name. `ssh`, `perl`, `python3`,
# `ruby`, `tail` — anything else that carries the pattern in its own argv — was still
# counted right along with the real daemon, and the ticket names `ssh` as a live route.
# Both close with the same change: allowlist `comm = java` instead of denying shells. Measured on
# the live daemon: `pid=30224 comm=java`. An empty comm (hole 1) is not `java` either, so it is
# excluded for free — no separate "is this pid still alive" check needed.
#
# The objection, because it is real: an allowlist can UNDER-count. If fleetd ever stops being
# launched as `java -jar ...` — a native image, a renamed launcher — `running_pid()` silently
# returns nothing and `assert_single_daemon` stops noticing a second daemon at all. For a guard,
# that false-negative direction is the worse one to be wrong in. This is not a new assumption,
# though: `PATTERN='target/fleetd.jar'` two lines up already assumes the daemon is a jar, which
# is only ever run by `java`. If that launch method changes, `PATTERN` stops matching anything
# before this allowlist would ever get the chance to be wrong — the allowlist rides on the same
# assumption that is already load-bearing, it does not add a new one. Whoever changes the launch
# method needs to update both `PATTERN` and this allowlist together.
running_pid() {
local pid comm out=''
for pid in $(pgrep -f "$PATTERN" 2>/dev/null || true); do
comm="$(ps -o comm= -p "$pid" 2>/dev/null || true)"
comm="${comm##*/}"
comm="${comm#-}"
# Allowlist, not a denylist of wrappers — see the CORRECTION 1 comment above. Anything that
# is not literally `java` is excluded, including an empty comm from a pid that already exited.
[ "$comm" = java ] || continue
out="$out$pid"$'\n'
done
printf '%s' "$out"
}
# fleetd #493 — three small, independently testable pieces of "never build into the path a
# running process holds":
@@ -507,8 +555,11 @@ assert_single_daemon() {
die "more than one fleetd process is running after this restart (pids: $(printf '%s' "$pids" | tr '\n' ' ')).
This is the exact failure a racing supervisor produces: the OLD jar was revived by its
supervisor while this script started a NEW copy. Two daemons on one herdr session kill
each other's members. Investigate with 'pgrep -f \"$PATTERN\"' and stop the wrong one by
hand — do not assume either pid is the one you want."
each other's members. Investigate with 'ps -eo pid,comm,args | grep -F \"$PATTERN\"' and
check the COMM column of each hit yourself before acting — a bare 'pgrep -f \"$PATTERN\"'
(fleetd #593) can match the very shell you type it into, not just the daemon, so it is not
safe remediation advice on its own. Stop the wrong one by hand — do not assume either pid
is the one you want."
fi
}
+130
View File
@@ -433,6 +433,130 @@ test_assert_single_daemon_rejects_two_pids() {
printf '%s' "$output" | grep -qF '4343' || fail "refusal message does not list the pids it found"
}
# fleetd #593 instance 2 — `running_pid()` used to be a bare `pgrep -f "$PATTERN"`, which matches
# ANY process whose full command line contains the pattern TEXT, including a shell that merely
# embeds it as literal text rather than being the daemon. Measured live on this Mac: `pgrep -c`
# (a one-call count) does not exist on BSD at all, and `bash -c "<single command>"` execs in place
# so no parent shell survives to hold the pattern — which is exactly why the defect did not
# reproduce from a plain script and needs a wrapper shaped like this instead. A `sh -c '...; ...'`
# with MORE THAN ONE statement does not get that exec-in-place treatment: the shell forks a child
# for the second statement and stays alive itself, holding the whole `-c` string — pattern text
# included — in its own `ps -o args`, for as long as it runs. That is the same shape an
# `ssh host "…; …"` wrapper or a hand-typed pipeline leaves behind. Before the fix this test would
# have found the wrapper's pid in running_pid()'s output; it must not.
test_running_pid_excludes_self_matching_wrapper_shell() {
local before after wrapper_pid
before="$(running_pid)"
sh -c 'echo "target/fleetd.jar" >/dev/null; sleep 20' &
wrapper_pid=$!
sleep 0.3
after="$(running_pid)"
kill "$wrapper_pid" 2>/dev/null || true
wait "$wrapper_pid" 2>/dev/null || true
[ "$after" = "$before" ] \
|| fail "running_pid() counted a self-matching wrapper shell (pid $wrapper_pid, holding the pattern as literal text in its own argv, not the daemon): before=[$before] after=[$after]"
}
# fleetd #593 CORRECTION 1 — the round-1 version of this test gave its standin an argv[0]
# containing the pattern text (via `exec -a`) and left `comm` as whatever that override produced,
# which was never `java`. That was fine for a denylist-of-shells filter, but the allowlist below
# now requires `comm = java` specifically, so the standin here must actually carry that comm, not
# just avoid being a shell. `exec -a java` overrides argv[0] to `java` while the process itself
# stays a genuine, harmless `sh`; combining it with the same non-exec'ing multi-statement shape
# the wrapper-shell test above uses keeps the pattern text in the process's own `ps -o args` for
# as long as it runs. Measured live on this Mac (BSD/macOS: `ps -o comm=` here reflects argv[0]):
# `comm=java`, `args` contains the pattern, `pgrep -f "$PATTERN"` finds it. Copying a real system
# binary into a scratch path and executing it from there was tried first, for a more literal
# stand-in daemon, and the OS killed it outright (SIGKILL, exit 137 — almost certainly a
# code-signing check on a relocated binary); `exec -a` needs no binary of its own and nothing
# under a scratch directory, and it is the technique CORRECTION 1 names as the right one.
#
# This is the one live-process test in this file whose result could differ on Linux: Linux sets
# `comm` from the actually-executed binary's own path, not from `exec -a`'s argv[0] override (BSD
# ties `comm` to argv[0], which is what makes this technique work here) — so on Linux this
# specific fixture might report `comm=sh`, not `comm=java`, even though the REAL daemon (a literal
# `java -jar target/fleetd.jar` process, never fabricated) is unaffected either way. I could not
# verify this fixture's behavior on Linux, so test_running_pid_counts_a_pid_whose_comm_is_java
# below backstops the same claim (the allowlist admits a pid whose comm is `java`) with a stubbed
# `ps`, which is identical bash on every platform and carries no such platform question.
test_running_pid_finds_a_real_java_named_second_process() {
local before after standin_pid
before="$(running_pid)"
( exec -a java sh -c 'echo "target/fleetd.jar" >/dev/null; sleep 20' ) &
standin_pid=$!
sleep 0.3
after="$(running_pid)"
kill "$standin_pid" 2>/dev/null || true
wait "$standin_pid" 2>/dev/null || true
printf '%s\n' "$after" | grep -qxF "$standin_pid" \
|| fail "running_pid() did not find a real second process (pid $standin_pid, comm forced to 'java' via exec -a) whose own argv holds the pattern: before=[$before] after=[$after]"
}
# fleetd #593 CORRECTION 1, hole 2 — the round-1 filter denied known shell names (sh/bash/zsh/
# dash/ksh) and counted everything else. `ssh`, `perl`, `python3`, `ruby`, `tail` — anything not on
# that list, carrying the pattern in its own argv — was still counted right alongside the real
# daemon, and the ticket names `ssh` as a live route. Stubbing `pgrep`/`ps` (rather than spawning a
# real perl/ssh process) pins the exact discriminator this correction is about — comm, not the
# caller's shape — deterministically on every platform, with no dependency on perl/python3/ruby
# being installed in whatever environment runs this suite, and no dependency on how a given OS
# derives `comm` for a fabricated process (see the comment above
# test_running_pid_finds_a_real_java_named_second_process for why that matters here).
test_running_pid_drops_a_pid_whose_comm_is_not_java() {
pgrep() { printf '4242\n'; }
ps() { printf 'perl\n'; }
local found
found="$(running_pid)"
unset -f pgrep ps
[ -z "$found" ] \
|| fail "running_pid() counted pid 4242 whose comm is 'perl', not 'java' — denying known shell names does not exclude a non-shell wrapper such as ssh or perl (fleetd #593 CORRECTION 1): found=[$found]"
}
# fleetd #593 CORRECTION 1, hole 1 — pgrep can list a pid that exits before the following
# `ps -o comm=` lookup runs; on a gone pid `ps` prints nothing, so `comm` comes back empty. Under
# the round-1 denylist an empty string matched none of the denied shell names, so the dead pid was
# still counted — the exact false-positive shape the ticket exists to remove, just rarer. The
# allowlist fixes this for free: an empty comm is not `java` either.
test_running_pid_drops_a_pid_that_exited_before_the_comm_lookup() {
pgrep() { printf '4242\n'; }
ps() { :; } # a pid that no longer exists: the real `ps -p <gone>` prints nothing and this mirrors that
local found
found="$(running_pid)"
unset -f pgrep ps
[ -z "$found" ] \
|| fail "running_pid() counted pid 4242 whose comm lookup came back empty (the pid had already exited before the lookup ran) — an empty comm must not pass the allowlist (fleetd #593 CORRECTION 1): found=[$found]"
}
# The positive backstop for both stubbed tests above, and for
# test_running_pid_finds_a_real_java_named_second_process on whatever platform that live fixture
# does not itself carry comm=java: the allowlist must still ADMIT the one comm value the real
# daemon actually has. Measured on the real, currently-running daemon on this Mac: `comm=java`.
test_running_pid_counts_a_pid_whose_comm_is_java() {
pgrep() { printf '4242\n'; }
ps() { printf 'java\n'; }
local found
found="$(running_pid)"
unset -f pgrep ps
printf '%s\n' "$found" | grep -qxF '4242' \
|| fail "running_pid() did not count pid 4242 whose comm is 'java' — the daemon's own name must pass the allowlist: found=[$found]"
}
# fleetd #593 instance 3 — assert_single_daemon's refusal message used to tell the operator to
# "Investigate with 'pgrep -f \"\$PATTERN\"'", which — typed by hand or over ssh — is precisely the
# self-matching invocation instance 2 above fixes. A source-text check, the same technique
# test_no_error_lines_message_gated_by_drain_state uses: this is prose inside a die() call, never
# reached by sourcing (the SOURCED guard stops before the main flow, and this text only prints
# from inside a call assert_single_daemon makes when it is already refusing).
test_die_message_does_not_recommend_bare_pgrep_as_remediation() {
local src="$ROOT/scripts/redeploy-fleetd.sh" block bad
block="$(grep -A6 -F 'racing supervisor produces' "$src" || true)"
[ -n "$block" ] || fail "could not find the assert_single_daemon refusal message in redeploy-fleetd.sh"
bad="$(printf '%s' "$block" | grep -F "Investigate with 'pgrep -f" || true)"
[ -z "$bad" ] \
|| fail "assert_single_daemon's die message still hands the operator a bare 'pgrep -f \"\$PATTERN\"' as remediation (fleetd #593) — that is exactly the self-matching invocation"
printf '%s' "$block" | grep -qF 'fleetd #593' \
|| fail "assert_single_daemon's die message does not say in words that a pattern can match the caller (fleetd #593)"
}
# fleetd #511 — jar_id()'s no-argument default was unpinned by any test: nothing proved it reports
# $JAR (the live path) rather than $JAR_STAGED. Both halves matter, so this pins both: the bare call
# must hash the live jar, and an explicit path argument must hash THAT file, not fall back to $JAR.
@@ -1894,6 +2018,12 @@ test_require_drivable_supervisor_accepts_known_kinds
test_count_daemon_pids
test_assert_single_daemon_accepts_one_pid
test_assert_single_daemon_rejects_two_pids
test_running_pid_excludes_self_matching_wrapper_shell
test_running_pid_finds_a_real_java_named_second_process
test_running_pid_drops_a_pid_whose_comm_is_not_java
test_running_pid_drops_a_pid_that_exited_before_the_comm_lookup
test_running_pid_counts_a_pid_whose_comm_is_java
test_die_message_does_not_recommend_bare_pgrep_as_remediation
test_jar_id_defaults_to_live_and_reports_explicit_path
test_hash256_computes_a_real_sha256
test_jar_id_reports_absent_for_missing_file