Compare commits

...

16 Commits

Author SHA1 Message Date
Dai Ha 856dfc6318 fleetd #603 review: close the untested fall-through path
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Successful in 2m15s
PR review (comment 17358) found a real gap via mutation testing: replacing
the "pid never found" warn with die "no process appeared" left the whole
suite green, because neither existing test drove the case the fall-through
exists for -- running_pid() never finds anything (as its own doc comment
says it eventually will) while /healthz answers anyway.

Adds test_await_daemon_started_pid_never_found_but_healthy_warns_and_survives:
running_pid always empty, poll_health_body succeeds. Asserts DIED_CALLED=0,
the warn line is emitted, and NEW_PID stays empty (the honest "could not
establish this" answer, never a guessed pid).

Verified both halves myself: reverting the warn to die "no process appeared"
turns this one test red (FAIL: await_daemon_started must not die...);
restoring it returns the suite to green.
2026-09-20 16:28:53 +07:00
Dai Ha 81c1d8e91c fleetd #603: share HEALTH_WAIT between the pid poll and the health check
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 1m18s
CI / build (pull_request) Successful in 1m47s
The start step gave the new process its own short, fixed 10s budget before a
hard die, while the health check right after it waits a full HEALTH_WAIT
(60s) for the same daemon. Under launchd, launchctl load returns before the
java process exists, and on a slow host that took longer than 10s -- so the
script reported "no process appeared" on a deploy that had fully succeeded.

wait_for_new_pid/await_daemon_started fold the pid poll and the health check
into one decision: the pid poll now shares HEALTH_WAIT instead of its own
shorter budget, and a miss there falls through to the health check (direct
proof the daemon is up) instead of killing the run. A genuine failure still
dies, and still prints the log tail.

Adds behavioural tests for both acceptance criteria (slow start succeeds,
genuine failure still fails and prints the log tail) in one suite run, plus
unit tests for wait_for_new_pid and a call-site test for the new function.
2026-09-20 16:21:52 +07:00
Dai Ha f5c6a0e4fc CLAUDE.md: architects settled the two invented specifics at line 144
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 47s
CI / build (push) Successful in 3m0s
Both specifics in the "consult architects" paragraph were mine, not the
operator's. The operator declined twice to rule on them and directed the lead
to consult architects instead, so two architects on different models settled
them over two rounds.

"after two rounds" is gone. It was a ceiling nobody had evidence for, and it
implied a counter fleetd does not have - nothing in the daemon counts rounds.
The bound is now expressed as a shape: form independent positions, then
compare. That is a floor of two without naming a number.

The three-item operator list read as complete, so a lead hitting anything not
on it would conclude it must not ask. It is now explicitly examples, and
"granting access" replaces "credentials" - the case that motivated this was a
forge merge refusal on a protected branch, which "credentials" covers only
awkwardly.

Canonical block and the wiki template updated together; sync check passes.
2026-09-19 23:32:41 +07:00
Dai Ha a7aee5b982 Merge #600: fleetd lead-rollover outcomes readable after confirm()
Adds LeadRollover.status() and a 'status' action on fleet_handover, so a lead
can find out what happened to its own roll. Every failure past confirm() was a
log.warn the lead cannot read.

Five states. IN_PROGRESS is written at the confirm hand-off, BEFORE the token
leaves 'pending', and status() reads 'outcomes' first — so there is no window
in which an in-flight roll reports UNKNOWN.

Gated by the lead: mvn -o clean install exit 0, 1819 tests, 0 failures,
0 errors, 0 skipped, 143 reports; LeadRolloverTest 43, FleetMcpHandoverTest 12.
Diff read in full. 0 source-text assertions in the new tests.
2026-09-19 23:32:32 +07:00
Dai Ha bf895616a5 fleetd: distinguish an in-flight roll from an unknown token
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 1m32s
CI / build (pull_request) Successful in 2m8s
confirm() removed the token from pending before handing the roll to the
continuation, and outcomes was only written when runRollover reached an exit.
For the whole duration of the roll the token was in neither map, so status()
answered UNKNOWN - documented as "never issued, cancelled, or aged out". A lead
polling right after its own confirm was told the roll had never been requested.

Adds RollState.IN_PROGRESS, written at the confirm hand-off rather than at the
roll's end, so there is no gap. status() now reads outcomes before pending, so
the hand-off write cannot race the removal.

Also corrects the class javadoc, which still claimed nothing calls this class.

Recovered by the lead: the authoring member ended on a backend error (the host
slept mid-response) with this work uncommitted in its worktree. Verified before
committing: mvn -o clean install exit 0, 1819 tests, 0 failures, 0 errors,
0 skipped, 143 reports; LeadRolloverTest 43 (was 39).
2026-09-19 22:31:47 +07:00
Dai Ha b874afb0af fleetd: make lead-rollover outcomes readable after confirm()
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 1m31s
CI / build (pull_request) Successful in 2m14s
LeadRollover previously logged every post-confirm() failure only — a lead
has no way to read the daemon log, so a roll that timed out because its
own turn never settled (or /clear never re-settled) was invisible; the
lead would carry on believing a fresh session was coming.

Add a bounded (cap=200) token -> outcome record, written at each of the
three exits in runRollover (ROLLED, TURN_NEVER_SETTLED, CLEAR_NEVER_SETTLED),
and a read-only LeadRollover#status(token) accessor. The TURN_NEVER_SETTLED
detail names turnSettleSeconds explicitly so a reader knows what to raise.

Wire a "status" action onto the fleet_handover MCP tool (handler + schema);
it never schedules, cancels, or retries anything — confirm() remains the
only path that can ever cause a /clear.

Extends LeadRolloverTest (33 -> 39 tests) covering the six acceptance
properties, and FleetMcpHandoverTest (8 -> 12) for the new tool action.
2026-09-19 16:17:39 +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
23 changed files with 2093 additions and 94 deletions
+5 -4
View File
@@ -141,10 +141,11 @@ 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
authorized to settle it, not only to advise. Architects first form independent positions, then
compare them. If they still disagree after that comparison, they return both positions and their
checked evidence; the lead decides. Go to the operator only for an action the fleet has no
authority to take, such as spending money, granting access, or making a promise 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.
+253 -45
View File
@@ -24,6 +24,7 @@ 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;
@@ -80,7 +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;
@@ -218,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
@@ -408,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.
@@ -461,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
@@ -501,21 +506,26 @@ public final class Fleetd {
// 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.
@@ -585,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)) {
@@ -608,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.
@@ -1009,6 +1005,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
@@ -1066,6 +1176,104 @@ public final class Fleetd {
() -> 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).
@@ -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, hunter, 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, hunter, 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,8 +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 hunters}/
* {@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
@@ -2059,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.hunters' 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");
@@ -10,6 +10,8 @@ import java.io.IOException;
import java.io.UncheckedIOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Collections;
import java.util.LinkedHashMap;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
@@ -25,9 +27,11 @@ import java.util.function.Supplier;
* every gate ({@link #confirm}'s own checks) has passed — a deferred, single-shot continuation
* clears the lead's own pane and bootstraps a fresh session against that file.
*
* <p>This is the executor only. Nothing in this ticket wires an MCP tool onto {@link #open}/
* {@link #confirm}/{@link #cancel} — that is a separate, later unit; until it lands, nothing calls
* this class at all.
* <p>This is the executor behind the {@code fleet_handover} MCP tool ({@code
* dev.ltms.fleet.mcp.FleetMcp#handover}), which drives {@link #open}, {@link #confirm}, {@link
* #cancel}, and {@link #status} from a tool call — wired in fleetd #480 Unit C. <strong>An earlier
* version of this paragraph said nothing called this class at all; that stopped being true once
* that unit landed, and this correction exists so the javadoc does not go on claiming it.</strong>
*
* <p><strong>{@code confirm()} cannot roll inline — a fleetd #480 correction.</strong> The first
* version of this class called {@code agents.send(lead, "/clear")} directly from inside {@code
@@ -181,6 +185,91 @@ public final class LeadRollover {
}
}
/**
* How many tokens {@link #outcomes} remembers before it starts evicting the oldest — bounded
* so a long-running daemon never grows this map without limit. Chosen generously rather than
* tightly: production rolls are rare (this class's own ticket found exactly ONE completed roll
* ever logged on this host), and each entry is a handful of short strings, so even a full cap
* costs a few tens of kilobytes — nowhere near a reason to make it configurable. 200 entries
* comfortably outlasts any operator's own memory of "did that roll I asked for actually
* happen", which is the whole reason {@link #status} exists.
*
* <p><strong>This cap counts {@link RollState#IN_PROGRESS} entries exactly the same as
* finished ones.</strong> There is only the one bounded map: {@link #confirm} writes an {@link
* RollState#IN_PROGRESS} entry into {@link #outcomes} at hand-off, and the deferred
* continuation later overwrites that SAME key with a terminal state — it never inserts a
* second entry. An approved roll therefore occupies one slot in this map for its entire
* lifetime, from the moment {@link #confirm} hands off, not only once it finishes; a
* confirmed-but-not-yet-finished roll counts against the cap exactly like a finished one. The
* alternative (a separate, uncapped in-flight map) would let a burst of confirmed-but-stuck
* rolls grow without bound — the exact failure this cap exists to prevent — so it was rejected.
*/
static final int OUTCOME_HISTORY_CAP = 200;
/**
* What is known about one token, right now — the answer {@link #status} gives. Distinguishes
* three terminal outcomes an approved roll can finish with, one in-flight outcome for a roll
* that has been approved but has not finished yet, and two answers for a token that names no
* active work at all: still pending confirmation, or nothing known about this token at all.
*/
public enum RollState {
/**
* {@code token} is still open: either {@link #open} was called and {@link #confirm} has not
* been (or not successfully) yet, or a {@link #confirm} call failed one of its gate checks
* and left the token pending for a retry — see {@link #confirm}'s javadoc ("token stays
* pending"). Indistinguishable from a genuinely fresh request; a caller wanting to know
* WHICH gate most recently refused should read the {@link RollDecision} that {@link
* #confirm} itself returned, not this status. <strong>Never the state of an APPROVED
* roll</strong> — see {@link #IN_PROGRESS}, which {@link #confirm} records at the moment it
* hands off, before this token is even removed from the pending set.
*/
PENDING,
/**
* {@link #confirm} approved this roll and handed it to the deferred continuation, which has
* not finished yet. Recorded by {@link #confirm} itself, at hand-off — <strong>before</strong>
* {@code token} is removed from the pending set — so there is never a gap in which {@link
* #status} could wrongly answer {@link #UNKNOWN} ("nothing was ever requested") for a roll
* that is, in fact, actively running. This is not sticky: the deferred continuation
* overwrites this same entry with a terminal state ({@link #ROLLED}, {@link
* #TURN_NEVER_SETTLED}, or {@link #CLEAR_NEVER_SETTLED}) once it finishes.
*/
IN_PROGRESS,
/**
* {@link #confirm} was approved and the deferred continuation completed the entire roll:
* the calling lead's turn settled, {@code /clear} was sent and settled, and {@code
* bootstrapText} was sent.
*/
ROLLED,
/**
* {@link #confirm} was approved, but the calling lead's own turn never reached a boundary
* (IDLE or DONE) within {@code turnSettleSeconds} — no {@code /clear} was ever sent, at
* all. This is the branch the fleetd #480 correction exists to make safe, and the one this
* status exists to make VISIBLE: before this, a lead that hit this case had no way to find
* out, and would carry on believing it was about to be replaced. See this class's javadoc.
*/
TURN_NEVER_SETTLED,
/**
* {@link #confirm} was approved and {@code /clear} was sent, but the pane never re-settled
* within {@code clearSettleSeconds} — {@code bootstrapText} was never sent.
*/
CLEAR_NEVER_SETTLED,
/**
* {@code token} names nothing this instance currently knows about: never issued by {@link
* #open}, dropped by {@link #cancel}, or aged out of {@link #outcomes}'s bounded history.
* These three causes are not distinguished — all of them mean "there is nothing to tell
* you", which is the entire content of a clean answer here.
*/
UNKNOWN
}
/**
* The answer {@link #status} gives for one token: a {@link RollState} and a human-readable
* {@code detail}. For {@link RollState#TURN_NEVER_SETTLED}, {@code detail} names {@code
* turnSettleSeconds} and its configured value explicitly, so a reader who sees this knows what
* to raise.
*/
public record RollStatus(RollState state, String detail) {}
private final AgentControl agents;
private final Supplier<FleetConfig.LeadRollover> configSupplier;
/**
@@ -201,6 +290,23 @@ public final class LeadRollover {
*/
private final Consumer<Runnable> continuationRunner;
private final Map<String, PendingRollover> pending = new ConcurrentHashMap<>();
/**
* Finished tokens → what actually happened, for {@link #status}. Bounded by {@link
* #OUTCOME_HISTORY_CAP}, oldest evicted first ({@code removeEldestEntry} on an insertion-order
* {@link LinkedHashMap}). Wrapped in {@link Collections#synchronizedMap} because entries are
* written from whatever thread {@code continuationRunner} runs the roll on (a fresh virtual
* thread in production, the calling test thread under {@code Runnable::run}) and read from
* whatever thread calls {@link #status} (the MCP handler thread) — a plain {@code
* LinkedHashMap} is not safe for that, and {@code removeEldestEntry} additionally requires
* external synchronization even for a thread-safe map that merely wraps it.
*/
private final Map<String, RollStatus> outcomes = Collections.synchronizedMap(
new LinkedHashMap<>(16, 0.75f, false) {
@Override
protected boolean removeEldestEntry(Map.Entry<String, RollStatus> eldest) {
return size() > OUTCOME_HISTORY_CAP;
}
});
/** Production constructor — wall clock, real sleep between settle polls, a real virtual thread. */
public LeadRollover(AgentControl agents, Supplier<FleetConfig.LeadRollover> configSupplier,
@@ -366,6 +472,15 @@ public final class LeadRollover {
return docCheck;
}
// Record IN_PROGRESS BEFORE removing from `pending` — see RollState#IN_PROGRESS and
// OUTCOME_HISTORY_CAP's javadoc. This ordering means `token` is written into `outcomes`
// while it is STILL present in `pending`; status() checks `outcomes` first (see that
// method), so it reports IN_PROGRESS immediately, not the brief-but-real gap a
// remove-then-put ordering would leave in which the token is in neither map.
outcomes.put(token, new RollStatus(RollState.IN_PROGRESS,
"confirm() approved this roll and handed it to the deferred continuation; it has "
+ "not finished yet — still waiting for the calling turn to settle, for "
+ "/clear to be sent and settle, or for bootstrapText to be sent"));
pending.remove(token);
log.info("lead-rollover: confirmed token={} lead={} — roll scheduled once the calling turn ends",
token, callerTerminal);
@@ -393,6 +508,11 @@ public final class LeadRollover {
+ "turn is still live and clearing it now would destroy live context "
+ "(token={}, configured={}s elapsed={}ms)",
lead, p.token(), cfg.turnSettleSeconds(), turnResult.elapsedMillis());
outcomes.put(p.token(), new RollStatus(RollState.TURN_NEVER_SETTLED,
"the calling lead's own turn never reached a boundary (IDLE or DONE) within "
+ "turnSettleSeconds=" + cfg.turnSettleSeconds() + "s (measured elapsed="
+ turnResult.elapsedMillis() + "ms) — no /clear was ever sent. If this "
+ "keeps happening, raise turnSettleSeconds in fleetd.yaml"));
return;
}
@@ -411,11 +531,17 @@ public final class LeadRollover {
+ "elapsed={}ms nudges={})",
lead, p.token(), cfg.clearSettleSeconds(), clearResult.elapsedMillis(),
clearResult.nudges());
outcomes.put(p.token(), new RollStatus(RollState.CLEAR_NEVER_SETTLED,
"/clear was sent, but the pane never re-settled within clearSettleSeconds="
+ cfg.clearSettleSeconds() + "s (measured elapsed=" + clearResult.elapsedMillis()
+ "ms, nudges=" + clearResult.nudges() + ") — bootstrapText was never sent"));
return;
}
agents.send(lead, cfg.bootstrapTextFor(p.handoverPath()));
long rollElapsedMillis = nowMillis.getAsLong() - rollStartMillis;
log.info("lead-rollover: rolled token={} lead={} elapsedMs={}", p.token(), lead, rollElapsedMillis);
outcomes.put(p.token(), new RollStatus(RollState.ROLLED,
"rolled successfully in " + rollElapsedMillis + "ms"));
}
/** Drop a pending request without rolling. @return whether a pending request existed for {@code token} */
@@ -423,6 +549,47 @@ public final class LeadRollover {
return pending.remove(token) != null;
}
/**
* Read-only: what is currently known about {@code token}. <strong>Never sends anything, never
* schedules, cancels, or retries a roll</strong> — a caller may poll this as often as it likes
* with no side effect at all, which is exactly why it exists: every failure past {@link
* #confirm} used to be a {@code log.warn} a lead can never read (see this class's javadoc), and
* this is the only route back.
*
* @param token the token {@link #open} returned; {@code null} or blank is a clean {@link
* RollState#UNKNOWN}, never a {@link NullPointerException} — {@link #pending} is a
* {@link ConcurrentHashMap}, which throws on a {@code null} key lookup, so this
* short-circuits before ever reaching it
* @return {@link RollState#IN_PROGRESS} for an approved roll whose continuation has not
* finished yet, or a terminal state once it has (both read from {@link #outcomes} —
* checked FIRST, see below); {@link RollState#PENDING} while {@code token} is still
* open and has not yet been approved (including one left pending by a {@link #confirm}
* gate refusal — see that method's javadoc); or {@link RollState#UNKNOWN} for a token
* never issued, cancelled, or aged out of the bounded history
*/
public RollStatus status(String token) {
if (token == null || token.isBlank()) {
return new RollStatus(RollState.UNKNOWN, "no token given");
}
// `outcomes` is checked BEFORE `pending`, deliberately: `confirm` writes an IN_PROGRESS
// entry into `outcomes` before it removes `token` from `pending` (see `confirm`'s own
// comment at that call site), so for the brief window where a token is present in BOTH
// maps, this order reports the more accurate answer (IN_PROGRESS, already approved) rather
// than the stale one (PENDING, not yet approved) a pending-first check would give.
RollStatus recorded = outcomes.get(token);
if (recorded != null) {
return recorded;
}
if (pending.containsKey(token)) {
return new RollStatus(RollState.PENDING, "open() has been called for this token and "
+ "it has not yet been confirmed — or a confirm() gate check failed and left it "
+ "pending, so the same token may be retried once the problem is fixed");
}
return new RollStatus(RollState.UNKNOWN, "token names no pending or finished rollover "
+ "request known to this instance — never issued, cancelled, or aged out of the "
+ "bounded history (cap=" + OUTCOME_HISTORY_CAP + ")");
}
/**
* The three handover-file checks, in order: exists, not empty, fresh (modified after
* {@link #open}'s timestamp and not older than {@code maxDocAgeSeconds}). Stats {@code
@@ -1248,14 +1248,16 @@ public final class FleetMcp {
Map<String, Object> args) {
String action = str(args, "action");
if (isBlank(action)) {
return error("action is required: \"open\", \"confirm\" or \"cancel\"");
return error("action is required: \"open\", \"confirm\", \"cancel\" or \"status\"");
}
return switch (action) {
case "open" -> handoverOpen(leadRollover, callerTerminal, str(args, "reason"));
case "confirm" -> handoverConfirm(leadRollover, callerTerminal, str(args, "token"),
truthy(args, "operatorConfirmed"));
case "cancel" -> handoverCancel(leadRollover, str(args, "token"));
default -> error("unknown action \"" + action + "\" — must be \"open\", \"confirm\" or \"cancel\"");
case "status" -> handoverStatus(leadRollover, str(args, "token"));
default -> error("unknown action \"" + action
+ "\" — must be \"open\", \"confirm\", \"cancel\" or \"status\"");
};
}
@@ -1325,6 +1327,29 @@ public final class FleetMcp {
return text(json(m));
}
/**
* {@code action: "status"}. Read-only — see {@link LeadRollover#status}: never schedules,
* cancels, or retries anything, and it is the only way for a lead to find out what happened to
* a token past {@code confirm()}, since every outcome after that point is otherwise logged only
* (see {@link LeadRollover}'s class javadoc).
*/
private static McpSchema.CallToolResult handoverStatus(LeadRollover leadRollover, String token) {
if (leadRollover == null) {
Map<String, Object> m = new LinkedHashMap<>();
m.put("state", "NOT_CONFIGURED");
m.put("detail", "leadRollover: is not configured");
return text(json(m));
}
if (isBlank(token)) {
return error("token is required for action \"status\"");
}
LeadRollover.RollStatus s = leadRollover.status(token);
Map<String, Object> m = new LinkedHashMap<>();
m.put("state", s.state().name());
m.put("detail", s.detail());
return text(json(m));
}
/** The one shared {@code NOT_CONFIGURED} refusal shape for {@code open}/{@code confirm}. */
private static McpSchema.CallToolResult notConfigured() {
return refusalJson(false, "NOT_CONFIGURED", "leadRollover: is not configured");
@@ -2265,20 +2290,26 @@ public final class FleetMcp {
return tool(FleetTool.HANDOVER.wireName(),
"Replace your OWN lead session once its context is full: write a handover file, "
+ "then use this to have fleetd clear your pane and bootstrap a fresh lead "
+ "session against it. Three actions: 'open' (requests a token and the "
+ "session against it. Four actions: 'open' (requests a token and the "
+ "handoverPath you must write the handover file to before confirming), "
+ "'confirm' (validates every gate and — only if every one passes — schedules "
+ "the roll; it does NOT itself clear the pane, the roll runs once this call's "
+ "own turn ends), and 'cancel' (drops a pending request without rolling). "
+ "Primary-only. There is deliberately no terminal/session/leadTerminal "
+ "parameter: the pane to roll is always resolved from YOUR OWN connection, "
+ "never a value you pass, so you can only ever roll yourself — never another "
+ "lead. Requires leadRollover: to be configured; when it is not, every action "
+ "returns a clean refusal naming NOT_CONFIGURED instead of failing.",
+ "own turn ends), 'cancel' (drops a pending request without rolling), and "
+ "'status' (read-only: what happened to a token after 'confirm' — still "
+ "running (approved but not finished yet), the roll completed, the calling "
+ "turn never settled within turnSettleSeconds so no /clear was ever sent, or "
+ "/clear itself never settled so bootstrapText was never sent; never "
+ "schedules, cancels or retries anything). Primary-only. "
+ "There is deliberately no terminal/session/leadTerminal parameter: the pane "
+ "to roll is always resolved from YOUR OWN connection, never a value you "
+ "pass, so you can only ever roll yourself — never another lead. Requires "
+ "leadRollover: to be configured; when it is not, every action returns a "
+ "clean refusal naming NOT_CONFIGURED instead of failing.",
objectSchema(Map.of(
"action", stringProp("\"open\", \"confirm\" or \"cancel\""),
"action", stringProp("\"open\", \"confirm\", \"cancel\" or \"status\""),
"reason", stringProp("Free-text audit note for \"open\" (optional, logged only)"),
"token", stringProp("The token \"open\" returned — required for \"confirm\" and \"cancel\""),
"token", stringProp("The token \"open\" returned — required for \"confirm\", "
+ "\"cancel\" and \"status\""),
"operatorConfirmed", Map.of("type", "boolean",
"description", "For \"confirm\": your answer to \"has the human operator "
+ "confirmed this wipe\" (default false; only consulted when "
@@ -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,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,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,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");
}
}
@@ -16,6 +16,7 @@ import org.junit.jupiter.api.io.TempDir;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Function;
@@ -1004,6 +1005,324 @@ class LeadRolloverTest {
}
}
// ---- CB-... : LeadRollover#status makes the outcome of a confirmed roll readable -----------
@Test
@DisplayName("[STATUS 1] after the calling turn never settles, status() reports "
+ "TURN_NEVER_SETTLED for that token — this assertion could not even be written before "
+ "status() existed")
void statusReportsTurnNeverSettledAfterTheRollIsAbandoned() throws IOException {
FakeHerdr herdr = new FakeHerdr();
herdr.agentStatus("working"); // the calling lead's own pane — never goes idle in this test
Path handover = writeHandover("handover contents");
FleetConfig.LeadRollover config =
new FleetConfig.LeadRollover(handover.toString(), true, 3600, 1 /*turnSettleSeconds*/, 20, "text");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, config, () -> clock.addAndGet(500));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "every synchronous gate should pass; the refusal happens "
+ "only inside the deferred continuation, which this test's synchronous runner has "
+ "already run to completion by the time confirm() returns");
assertEquals(0, promptCallCount(herdr), "sanity: /clear was never sent");
LeadRollover.RollStatus status = rollover.status(pending.token());
assertEquals(LeadRollover.RollState.TURN_NEVER_SETTLED, status.state());
assertTrue(status.detail().contains("turnSettleSeconds"), "the detail must name the knob a "
+ "reader needs to raise: " + status.detail());
assertTrue(status.detail().toLowerCase().contains("no /clear"), "the detail must say plainly "
+ "that no /clear was ever sent: " + status.detail());
}
@Test
@DisplayName("[STATUS 2] a roll that completes reports ROLLED for its token")
void statusReportsRolledForACompletedRoll() throws IOException {
FakeHerdr herdr = new FakeHerdr(); // default idle — a full successful roll
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, cfg(handover.toString()), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason() + " / " + decision.detail());
assertEquals(2, promptCallCount(herdr), "sanity: /clear then bootstrapText were both sent");
LeadRollover.RollStatus status = rollover.status(pending.token());
assertEquals(LeadRollover.RollState.ROLLED, status.state());
assertNotNull(status.detail());
}
@Test
@DisplayName("[STATUS 3] a roll where /clear never settles reports CLEAR_NEVER_SETTLED — "
+ "distinct from both ROLLED and TURN_NEVER_SETTLED")
void statusReportsClearNeverSettledDistinctFromTheOtherTwoStates() throws IOException {
// Idle until /clear is sent, then permanently working — the SECOND wait never settles.
FakeHerdr fake = new FakeHerdr();
HerdrClient flipsAfterClear = new HerdrClient() {
@Override
public JsonNode call(String method, Object params) throws HerdrException {
JsonNode result = fake.call(method, params);
if ("agent.prompt".equals(method) && String.valueOf(params).contains("/clear")) {
fake.agentStatus("working");
}
return result;
}
@Override
public void close() {
fake.close();
}
};
Path handover = writeHandover("handover contents");
FleetConfig.LeadRollover config =
new FleetConfig.LeadRollover(handover.toString(), true, 3600, 20, 1 /*clearSettleSeconds*/, "boot text");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(flipsAfterClear, config, () -> clock.addAndGet(500));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "every synchronous gate passes; the refusal is logged only, "
+ "deep inside the deferred continuation");
assertEquals(1, promptCallCount(fake), "sanity: only /clear was sent, never bootstrapText");
LeadRollover.RollStatus status = rollover.status(pending.token());
assertEquals(LeadRollover.RollState.CLEAR_NEVER_SETTLED, status.state());
assertNotEquals(LeadRollover.RollState.ROLLED, status.state());
assertNotEquals(LeadRollover.RollState.TURN_NEVER_SETTLED, status.state());
assertTrue(status.detail().contains("clearSettleSeconds"), status.detail());
}
@Test
@DisplayName("[STATUS 4] a token that was never issued, or was cancelled, gives a clean "
+ "UNKNOWN answer rather than an exception or a false ROLLED")
void statusOnUnissuedOrCancelledTokenIsCleanNotAnExceptionOrFalseSuccess() throws IOException {
FakeHerdr herdr = new FakeHerdr();
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, cfg(handover.toString()), fixedClock(clock));
// never issued at all
LeadRollover.RollStatus neverIssued = assertDoesNotThrow(() -> rollover.status("no-such-token"));
assertEquals(LeadRollover.RollState.UNKNOWN, neverIssued.state());
assertNotEquals(LeadRollover.RollState.ROLLED, neverIssued.state());
// null / blank must not throw either — pending is a ConcurrentHashMap, which throws on a
// null-key lookup unless status() guards it first
assertEquals(LeadRollover.RollState.UNKNOWN, assertDoesNotThrow(() -> rollover.status(null)).state());
assertEquals(LeadRollover.RollState.UNKNOWN, assertDoesNotThrow(() -> rollover.status(" ")).state());
// opened, then cancelled
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
assertTrue(rollover.cancel(pending.token()));
LeadRollover.RollStatus cancelled = assertDoesNotThrow(() -> rollover.status(pending.token()));
assertEquals(LeadRollover.RollState.UNKNOWN, cancelled.state());
assertNotEquals(LeadRollover.RollState.ROLLED, cancelled.state());
}
@Test
@DisplayName("[STATUS 5] the bounded outcome history never grows past its cap")
void statusHistoryDoesNotGrowPastItsCap() throws IOException {
FakeHerdr herdr = new FakeHerdr(); // default idle — every roll completes
Path handover = writeHandover("handover contents"); // one file, reused by every roll below —
// checkHandover only compares its mtime against each open()'s OWN requestedAtMillis (the
// fake clock, in the low thousands), and the file's real (wall-clock) mtime is always far
// larger than that, so freshness passes on every iteration without rewriting the file.
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, cfg(handover.toString()), () -> clock.addAndGet(1));
int rolls = LeadRollover.OUTCOME_HISTORY_CAP + 50;
String[] tokens = new String[rolls];
for (int i = 0; i < rolls; i++) {
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full #" + i);
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "roll #" + i + " should have been approved: " + decision.detail());
tokens[i] = pending.token();
}
assertEquals(LeadRollover.RollState.ROLLED, rollover.status(tokens[rolls - 1]).state(),
"the most recently finished roll's outcome must still be in the bounded history");
// There is no direct size accessor for the bounded history, so boundedness is asserted
// indirectly and behaviourally: the OLDEST finished roll's outcome must have been evicted
// (reads back as a clean UNKNOWN, exactly like a token that was never issued) once more than
// OUTCOME_HISTORY_CAP rolls have gone through this instance. If the cap were not enforced,
// tokens[0] would still read back ROLLED here, and this assertion would fail.
assertEquals(LeadRollover.RollState.UNKNOWN, rollover.status(tokens[0]).state(),
"the oldest finished roll's outcome must have been evicted once the cap was "
+ "exceeded — otherwise the bounded history is not actually bounded");
}
@Test
@DisplayName("[STATUS 6] a status() call against a pending roll is read-only — no herdr call is "
+ "made and the pending request is left untouched")
void statusCallAgainstAPendingRollIsReadOnlyAndTouchesNothing() throws IOException {
FakeHerdr herdr = new FakeHerdr();
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, cfg(handover.toString()), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollStatus status = rollover.status(pending.token());
assertEquals(LeadRollover.RollState.PENDING, status.state());
assertEquals(0, herdr.calls.size(), "status() must never make any herdr call at all — not "
+ "just no agent.prompt — since it must never schedule, cancel, or retry anything");
// the pending request must be left exactly as it was: the SAME token can still be confirmed
// afterwards, as if status() had never been called.
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "status() must not have consumed or otherwise disturbed the "
+ "pending request: " + decision.reason() + " / " + decision.detail());
}
// ---- PR #600 review round 2: IN_PROGRESS — the gap between confirm() handing off and the ----
// ---- continuation finishing must never read back as UNKNOWN ("nothing was ever requested") --
/**
* A {@code continuationRunner} that CAPTURES the roll instead of running it, so a test can
* observe {@link LeadRollover#status} in the window between {@link LeadRollover#confirm}
* handing off and the roll actually finishing — the window the synchronous {@code
* Runnable::run} runner used everywhere else in this class collapses to nothing. Call {@link
* #runNext()} to finish exactly one held roll, once the test is done observing the in-flight
* state.
*/
private static final class HoldingRunner implements java.util.function.Consumer<Runnable> {
private final List<Runnable> held = new ArrayList<>();
@Override
public void accept(Runnable runnable) {
held.add(runnable);
}
int heldCount() {
return held.size();
}
/** Runs (and removes) the oldest held roll — FIFO, matching confirm() call order. */
void runNext() {
held.remove(0).run();
}
}
private static LeadRollover newRolloverWithHoldingRunner(HerdrClient herdr,
FleetConfig.LeadRollover config, LongSupplier nowMillis, HoldingRunner runner) {
AgentControl agents = new AgentControl(herdr);
return new LeadRollover(agents, () -> config, _ -> null, nowMillis, () -> { }, runner);
}
@Test
@DisplayName("[IN-PROGRESS 1] between an approved confirm() and the continuation finishing, "
+ "status() reports IN_PROGRESS — not UNKNOWN, and not PENDING")
void statusReportsInProgressBetweenConfirmAndTheContinuationFinishing() throws IOException {
FakeHerdr herdr = new FakeHerdr(); // default idle — the held roll WOULD complete once run
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
HoldingRunner runner = new HoldingRunner();
LeadRollover rollover = newRolloverWithHoldingRunner(herdr, cfg(handover.toString()),
fixedClock(clock), runner);
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason() + " / " + decision.detail());
assertEquals(1, runner.heldCount(), "sanity: the roll must have been handed to the "
+ "continuation runner and held there, not run yet");
assertEquals(0, promptCallCount(herdr), "sanity: the held continuation has not run, so "
+ "nothing has been sent to the pane yet");
LeadRollover.RollStatus status = rollover.status(pending.token());
assertEquals(LeadRollover.RollState.IN_PROGRESS, status.state(),
"a lead calling status() right after confirm() returned approved, while the roll "
+ "is still running, must be told IN_PROGRESS — not UNKNOWN (\"nothing was "
+ "ever requested\", which would wrongly invite it to call open() again "
+ "mid-roll) and not PENDING (\"not yet approved\", which is simply false "
+ "here): got " + status.state() + " / " + status.detail());
}
@Test
@DisplayName("[IN-PROGRESS 2] once the continuation finishes, the same token reports its "
+ "terminal state — IN_PROGRESS is not sticky")
void inProgressStateIsNotStickyOnceTheContinuationFinishes() throws IOException {
FakeHerdr herdr = new FakeHerdr(); // default idle — the held roll completes once run
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
HoldingRunner runner = new HoldingRunner();
LeadRollover rollover = newRolloverWithHoldingRunner(herdr, cfg(handover.toString()),
fixedClock(clock), runner);
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason() + " / " + decision.detail());
assertEquals(LeadRollover.RollState.IN_PROGRESS, rollover.status(pending.token()).state(),
"sanity: must be IN_PROGRESS before the held continuation is run");
runner.runNext(); // finish the held roll now
assertEquals(2, promptCallCount(herdr), "sanity: the roll actually ran to completion "
+ "once released — /clear then bootstrapText");
LeadRollover.RollStatus after = rollover.status(pending.token());
assertEquals(LeadRollover.RollState.ROLLED, after.state(),
"the SAME token must now report its terminal state — IN_PROGRESS must not still be "
+ "reported once the roll has actually finished");
}
@Test
@DisplayName("[IN-PROGRESS 4] IN_PROGRESS is distinct from every other RollState")
void inProgressStateIsDistinctFromAllOtherStates() throws IOException {
FakeHerdr herdr = new FakeHerdr();
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
HoldingRunner runner = new HoldingRunner();
LeadRollover rollover = newRolloverWithHoldingRunner(herdr, cfg(handover.toString()),
fixedClock(clock), runner);
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason() + " / " + decision.detail());
LeadRollover.RollState inProgress = rollover.status(pending.token()).state();
assertEquals(LeadRollover.RollState.IN_PROGRESS, inProgress);
for (LeadRollover.RollState other : LeadRollover.RollState.values()) {
if (other == LeadRollover.RollState.IN_PROGRESS) {
continue;
}
assertNotEquals(other, inProgress, "IN_PROGRESS must be distinct from " + other);
}
}
@Test
@DisplayName("[IN-PROGRESS 5] eviction counts IN_PROGRESS entries toward the cap exactly like "
+ "finished ones — a burst of confirmed-but-not-yet-finished rolls still ages out")
void evictionCountsInProgressEntriesTowardTheCap() throws IOException {
FakeHerdr herdr = new FakeHerdr();
Path handover = writeHandover("handover contents"); // reused by every roll — see
// statusHistoryDoesNotGrowPastItsCap for why one shared file is enough for freshness.
AtomicLong clock = new AtomicLong(1_000);
HoldingRunner runner = new HoldingRunner(); // nothing run below — every roll stays IN_PROGRESS
LeadRollover rollover = newRolloverWithHoldingRunner(herdr, cfg(handover.toString()),
() -> clock.addAndGet(1), runner);
int rolls = LeadRollover.OUTCOME_HISTORY_CAP + 50;
String[] tokens = new String[rolls];
for (int i = 0; i < rolls; i++) {
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full #" + i);
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "roll #" + i + " should have been approved: " + decision.detail());
tokens[i] = pending.token();
}
assertEquals(rolls, runner.heldCount(), "sanity: none of these rolls have been run — every "
+ "one of them is sitting in outcomes as IN_PROGRESS, not a separate uncapped map");
assertEquals(LeadRollover.RollState.IN_PROGRESS, rollover.status(tokens[rolls - 1]).state(),
"the most recently confirmed (still in-flight) roll must still be in the bounded "
+ "history");
assertEquals(LeadRollover.RollState.UNKNOWN, rollover.status(tokens[0]).state(),
"the oldest confirmed roll's IN_PROGRESS entry must have been evicted once the cap "
+ "was exceeded, exactly like a finished entry would be — proving IN_PROGRESS "
+ "entries share the SAME bounded map and count against the SAME cap, rather "
+ "than living in a second, uncapped in-flight map");
}
@Test
@DisplayName("[fleetd #494 follow-up] the turn-settle timeout warn line prints the MEASURED "
+ "elapsed time next to the configured budget, never the configured value alone")
@@ -258,5 +258,51 @@ class FleetMcpHandoverTest {
assertTrue(textOf(r).contains("\"cancelled\":true"), textOf(r));
}
// --- new: action "status" — makes the outcome of a confirm() readable through the tool ----
@Test
@DisplayName("status on a null LeadRollover is a clean NOT_CONFIGURED refusal, never a throw")
void statusWithNullLeadRolloverRefusesCleanly() {
McpSchema.CallToolResult r = assertDoesNotThrow(() -> FleetMcp.handover(null, LEAD,
Map.of("action", "status", "token", "whatever")));
assertFalse(r.isError());
assertTrue(textOf(r).contains("NOT_CONFIGURED"), textOf(r));
}
@Test
@DisplayName("status on a token that was never opened reports UNKNOWN")
void statusOnUnknownTokenReportsUnknown() {
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
McpSchema.CallToolResult r = FleetMcp.handover(rollover, LEAD,
Map.of("action", "status", "token", "does-not-exist"));
assertFalse(r.isError());
assertTrue(textOf(r).contains("\"state\":\"UNKNOWN\""), textOf(r));
}
@Test
@DisplayName("status on a token that is still pending (opened, not confirmed) reports PENDING")
void statusOnPendingTokenReportsPending() {
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
String token = extractToken(textOf(FleetMcp.handover(rollover, LEAD, Map.of("action", "open"))));
McpSchema.CallToolResult r = FleetMcp.handover(rollover, LEAD,
Map.of("action", "status", "token", token));
assertFalse(r.isError());
assertTrue(textOf(r).contains("\"state\":\"PENDING\""), textOf(r));
}
@Test
@DisplayName("status is registered on the tool's schema and the schema still names no caller-identity parameter")
void statusActionIsAdvertisedOnTheSchema() {
FleetMcp m = mcp(null);
McpSchema.Tool tool = m.registeredTools().stream()
.filter(t -> "fleet_handover".equals(t.name()))
.findFirst()
.orElseThrow(() -> new AssertionError("fleet_handover was not registered"));
assertTrue(tool.description().contains("'status'"),
"the tool's own description must advertise the 'status' action: " + tool.description());
}
// --- acceptance 7 (wiring) is covered by FleetdLeadRolloverWiringTest, unchanged -----------
}
+131 -23
View File
@@ -50,6 +50,14 @@
# absence is the real signal, because a drain that dies on its first session prints nothing
# else either. Warns loudly; never fails the redeploy, because by the time this is detectable
# the new daemon is already up and healthy.
# 10. fleetd #603 — the same shape as trap 3 above, through a different door: the step that waited
# for the NEW process to appear gave it its own short, fixed 10s budget, then hard-`die`d,
# while the health check right after it waits a full $HEALTH_WAIT (60s) for the same daemon to
# answer. Under launchd, `launchctl load` returns as soon as launchd accepts the job, before the
# java process exists, and on a slow host that took longer than 10s — so the script died with
# "no process appeared" on a deploy that had fully succeeded. The pid poll now shares
# $HEALTH_WAIT instead of a separate, shorter budget, and a miss there falls through to the
# health check (the truer signal: is it actually answering?) instead of killing the run.
#
# Usage:
# scripts/redeploy-fleetd.sh # build, confirm, restart, verify
@@ -79,7 +87,9 @@ OUT="$MODULE/fleetd.out"
PATTERN='target/fleetd.jar'
HEALTH='http://127.0.0.1:8765/healthz'
STOP_WAIT=30 # seconds to wait for a clean exit before reporting failure
HEALTH_WAIT=60 # seconds to wait for /healthz to answer after start
HEALTH_WAIT=60 # seconds to wait for /healthz to answer after start — fleetd #603: also the pid-
# poll budget below (wait_for_new_pid/await_daemon_started), so the two checks
# share one named budget instead of the pid poll holding its own shorter one
# CB-594: the launchd agent this script must not fight with (see trap 6 above).
LAUNCHD_LABEL='dev.ltms.fleetd'
@@ -169,7 +179,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 +565,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
}
@@ -1016,6 +1077,67 @@ $(tail -30 "$out_file" 2>/dev/null)"
fi
}
# fleetd #603 — the dual of wait_for_daemon_exit above: waits for a pid to APPEAR instead of
# disappear. Used to be a bare `for _ in $(seq 10)` sitting directly in the main flow, with its own
# short, fixed budget that had nothing to do with $HEALTH_WAIT (60s) — the budget the health check
# right after it gets for the very same daemon. Under launchd, `launchctl load` returns as soon as
# launchd accepts the job, before the java process exists, and on a slow host that took longer than
# 10s — so the script died with "no process appeared" on a deploy that had fully succeeded (the
# operator confirmed pid, healthz, and a fresh log line all present, by hand, right afterwards).
# Sharing $HEALTH_WAIT here removes the extra, shorter magic number without inventing a new one.
wait_for_new_pid() {
local timeout="$1" _i
for _i in $(seq "$timeout"); do
[ -n "$(running_pid)" ] && return 0
sleep 1
done
[ -n "$(running_pid)" ]
}
# fleetd #603 — the pid-appeared check and the healthz check, folded into one decision. Same shape
# as swap_if_built/refuse_drain_gate/report_shutdown_drain above (#521/#528/#512): the main flow
# calls this ONE function unconditionally, so there is no bare guard left for a future edit to
# invert independently of it. That matters more here than for most of those: a plain grep of this
# script's source cannot tell "a pid miss falls through to the health check" from "a pid miss still
# dies" apart, because both read as the same two lines of text with only the runtime branch
# changed — a source-text test could pass on either behavior. A real behavioural test on this
# function is the only thing that can actually tell them apart, which is why one exists below.
#
# wait_for_new_pid's miss is no longer fatal by itself: it falls through to the health check, which
# is direct proof the new daemon is up (/healthz answers 200) rather than a proxy for it (a process
# merely existing under a name running_pid() recognises). A genuine failure still dies here: it
# misses the pid poll AND the health poll, and report_health's own die() still prints the tail of
# $out_file, exactly as before this fix.
#
# Sets NEW_PID (global — the caller's "pid N, jar ..." result line reads it afterwards) and
# HEALTH_BODY/HEALTH_CODE (globals, the same reason report_health already needs them handed back).
# Only ever called as a bare statement in the main flow below, never from inside a `$( )`: a die()
# reached from inside a command substitution only kills that subshell, not the whole script, which
# would silently turn a genuine failure back into a false "succeeded" exit (see poll_health_body's
# own `|| true` idiom for the same hazard from the other direction).
await_daemon_started() {
local health_wait="$1" old_pid="$2" health_url="$3" out_file="$4"
NEW_PID=""
if wait_for_new_pid "$health_wait"; then
NEW_PID="$(running_pid)"
[ "$NEW_PID" != "${old_pid:-}" ] || die "pid unchanged ($NEW_PID) — the old daemon never died"
ok "started, pid $NEW_PID"
else
warn "no process matched $PATTERN within ${health_wait}s of starting — falling through to the health check, which is the more truthful signal"
fi
HEALTH_BODY="$(poll_health_body "$health_url" "$health_wait")" || true
HEALTH_CODE="000"
[ -n "$HEALTH_BODY" ] || HEALTH_CODE="$(curl -s -o /dev/null -w '%{http_code}' --max-time 2 "$health_url" 2>/dev/null || echo 000)"
report_health "$HEALTH_BODY" "$HEALTH_CODE" "$out_file" "$health_wait"
# By now /healthz has answered (report_health above would have died otherwise), so the daemon is
# confirmed up even if the pid poll never matched it — see running_pid()'s own comment on
# under-counting if the launch method ever stops being a plain `java -jar`. Fill NEW_PID in for the
# result line rather than leave it blank on an otherwise fully successful redeploy.
[ -n "$NEW_PID" ] || NEW_PID="$(running_pid)"
}
# Ticket item 8 — `HAD_OLD_PID=0; [ -n "$OLD_PID" ] && HAD_OLD_PID=1`. Measured safe under `set -e`
# at both bash 3.2.57 and 5.x (see the header comment trap 9 discussion in the ticket) — not a `set
# -e` hazard, but still an untested computation feeding report_shutdown_drain's own four-way
@@ -1234,28 +1356,14 @@ say "start"
# there now, so this line is the only thing left in the main flow to get wrong.
dispatch_start "$SUPERVISOR_KIND"
for _ in $(seq 10); do
NEW_PID="$(running_pid)"
[ -n "$NEW_PID" ] && break
sleep 1
done
[ -n "${NEW_PID:-}" ] || die "no process appeared. Last lines of $OUT:
$(tail -20 "$OUT" 2>/dev/null)"
[ "$NEW_PID" != "${OLD_PID:-}" ] || die "pid unchanged ($NEW_PID) — the old daemon never died"
ok "started, pid $NEW_PID"
# ------------------------------------------------------------------ verify
say "verify"
# fleetd #555: poll_health_body/health_is_up/report_health above. HEALTH_CODE is only ever
# consulted by report_health when the body came back empty; `|| true` on both assignments is the
# same "an absent/failing command substitution must not kill the script under set -e" idiom the
# swap/drain helpers already rely on (see poll_health_body's own comment).
HEALTH_BODY="$(poll_health_body "$HEALTH" "$HEALTH_WAIT")" || true
HEALTH_CODE="000"
[ -n "$HEALTH_BODY" ] || HEALTH_CODE="$(curl -s -o /dev/null -w '%{http_code}' --max-time 2 "$HEALTH" 2>/dev/null || echo 000)"
report_health "$HEALTH_BODY" "$HEALTH_CODE" "$OUT" "$HEALTH_WAIT"
# fleetd #603: await_daemon_started above folds the pid-appeared check and the healthz check into
# one decision — see its own comment for why a bare `if` here, split back into the old two pieces,
# would put back an untestable branch this ticket exists to close.
await_daemon_started "$HEALTH_WAIT" "$OLD_PID" "$HEALTH" "$OUT"
# A fresh listening line, strictly after the restart mark. An old daemon that never died would
# otherwise let an old line pass for a new one.
@@ -1310,7 +1418,7 @@ report_shutdown_drain "$FRESH_LOG" "$HAD_OLD_PID"
assert_single_daemon "$(running_pid)"
say "result"
ok "pid $NEW_PID, jar $(jar_id)"
ok "pid ${NEW_PID:-unknown}, jar $(jar_id)"
if [ "$REDEPLOY_AMQP_CHECK_SKIPPED" -eq 1 ]; then
# fleetd #552: the fourth reader of the fresh-log region. Without this branch REDEPLOY_ERROR_COUNT
# stays at its untouched 0 (classify_amqp_connection_errors never got a file to read) and
+271 -2
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.
@@ -646,6 +770,137 @@ test_wait_for_daemon_exit_times_out_if_pid_never_clears() {
source "$ROOT/scripts/redeploy-fleetd.sh" # restore the real running_pid/sleep for later tests
}
# fleetd #603 — wait_for_new_pid is the dual of wait_for_daemon_exit above: it must not report
# success while running_pid() still answers empty, and must report success the moment a pid
# appears. Same counter-file idiom as test_wait_for_daemon_exit_returns_true_once_pid_clears above,
# for the same reason (running_pid() runs inside a `$(...)` subshell on every call).
test_wait_for_new_pid_returns_true_once_pid_appears() {
local counter_file="$TMP/wait-new-pid-calls" final_calls
printf '0' > "$counter_file"
running_pid() {
local n
n="$(cat "$counter_file")"
n=$((n + 1))
printf '%s' "$n" > "$counter_file"
if [ "$n" -lt 3 ]; then printf ''; else printf '4242'; fi
}
sleep() { :; }
wait_for_new_pid 10 || fail "wait_for_new_pid did not report success once the pid appeared"
final_calls="$(cat "$counter_file")"
[ "$final_calls" -ge 3 ] || fail "wait_for_new_pid returned before actually re-checking running_pid"
source "$ROOT/scripts/redeploy-fleetd.sh" # restore the real running_pid/sleep for later tests
}
test_wait_for_new_pid_times_out_if_pid_never_appears() {
local rc=0
running_pid() { printf ''; }
sleep() { :; }
wait_for_new_pid 3 || rc=$?
[ "$rc" -ne 0 ] || fail "wait_for_new_pid reported success while the pid never appeared"
source "$ROOT/scripts/redeploy-fleetd.sh" # restore the real running_pid/sleep for later tests
}
# fleetd #603 — await_daemon_started folds the pid-appeared check and the healthz check into one
# decision (see its own comment in redeploy-fleetd.sh for why a source-text grep cannot tell the two
# possible behaviors apart here). These two tests are the ticket's own acceptance criteria, run
# together in this one suite invocation so neither can be satisfied by code that never fails at all:
#
# 1. a slow start must still succeed — running_pid mimics a process that does not appear until
# well after the OLD, buggy 10-second budget, but does appear, and healthz answers.
# 2. a genuine failure must still fail, and still print the log tail — running_pid and the health
# check both report nothing at all, ever.
test_await_daemon_started_slow_pid_then_healthy_succeeds() {
source "$ROOT/scripts/redeploy-fleetd.sh"
stub_die_recorder
local counter_file="$TMP/await-slow-pid-calls" output
printf '0' > "$counter_file"
running_pid() {
local n
n="$(cat "$counter_file")"
n=$((n + 1))
printf '%s' "$n" > "$counter_file"
# Stays empty well past the old 10-second budget, then appears — the exact shape #603 reports.
if [ "$n" -lt 12 ]; then printf ''; else printf '4242'; fi
}
sleep() { :; }
poll_health_body() { printf '{"status":"ok"}'; return 0; }
# NOT `output="$(await_daemon_started ...)"`: that would run the call in a subshell, and
# NEW_PID — a plain global assignment inside the function, by design (see its own comment) — would
# die with that subshell instead of reaching this test's own shell. Redirect to a file instead, the
# same hazard await_daemon_started's own comment warns die() itself is subject to.
await_daemon_started 60 "" "http://ignored/healthz" "$TMP/await-slow-pid.out" \
> "$TMP/await-slow-pid-output.log" 2>&1
output="$(cat "$TMP/await-slow-pid-output.log")"
[ "$DIED_CALLED" = 0 ] \
|| fail "await_daemon_started must not die on a slow-but-real start: $DIED_MESSAGE"
assert_equals "4242" "$NEW_PID" "await_daemon_started NEW_PID after a slow-but-real start"
printf '%s' "$output" | grep -qF 'started, pid 4242' \
|| fail "await_daemon_started did not report the pid once it finally appeared"
source "$ROOT/scripts/redeploy-fleetd.sh"
}
test_await_daemon_started_never_appears_dies_with_log_tail() {
source "$ROOT/scripts/redeploy-fleetd.sh"
stub_die_recorder
local out_file="$TMP/await-never-appears.out"
printf 'boot line one\nboot line two\n' > "$out_file"
running_pid() { printf ''; }
sleep() { :; }
poll_health_body() { return 1; }
await_daemon_started 2 "" "http://127.0.0.1:1/healthz" "$out_file" > /dev/null 2>&1
[ "$DIED_CALLED" = 1 ] \
|| fail "await_daemon_started must die when the daemon never appears and never becomes healthy"
printf '%s' "$DIED_MESSAGE" | grep -qF 'never answered' \
|| fail "await_daemon_started die message does not say healthz never answered"
printf '%s' "$DIED_MESSAGE" | grep -qF 'boot line two' \
|| fail "await_daemon_started die message does not include the log tail"
source "$ROOT/scripts/redeploy-fleetd.sh"
}
# fleetd #603 review — the path the fall-through actually exists for, and the one gap a lead
# mutation found in the first version of this test file: running_pid() NEVER finds anything (as its
# own doc comment says it eventually will, once the daemon stops being launched as a plain
# `java -jar` its allowlist recognises), while /healthz answers anyway. Neither of the two tests
# above drives this: the slow-pid test has the pid appear, so the `else` branch never runs, and the
# never-appears test fails BOTH checks, so it dies either way and cannot tell which branch fired.
# This must not die, must warn (so the operator is told the pid could not be identified), and must
# leave NEW_PID empty — the honest "could not establish this" answer, never a guessed pid, which is
# what the final result line's `${NEW_PID:-unknown}` fallback exists to print truthfully.
#
# Proof this actually pins the behavior, not just the source text (paste from a real run, not
# claimed): reverting the `warn` below back to `die "no process appeared"` (the old fleetd #603
# defect, reintroduced) turns this one test red —
# FAIL: await_daemon_started must not die when the pid is never found but healthz answers
# — and restoring `warn` turns the whole suite green again. Both halves observed, not asserted.
test_await_daemon_started_pid_never_found_but_healthy_warns_and_survives() {
source "$ROOT/scripts/redeploy-fleetd.sh"
stub_die_recorder
local out_file="$TMP/await-pid-never-found.out" output
printf 'boot line\n' > "$out_file"
running_pid() { printf ''; }
sleep() { :; }
poll_health_body() { printf '{"status":"ok"}'; return 0; }
# NOT `output="$(await_daemon_started ...)"` — see the slow-pid test above for why that would
# drop NEW_PID's assignment in a subshell instead of reaching this test's own shell.
await_daemon_started 3 "" "http://ignored/healthz" "$out_file" \
> "$TMP/await-pid-never-found-output.log" 2>&1
output="$(cat "$TMP/await-pid-never-found-output.log")"
[ "$DIED_CALLED" = 0 ] \
|| fail "await_daemon_started must not die when the pid is never found but healthz answers: $DIED_MESSAGE"
printf '%s' "$output" | grep -qF 'falling through to the health check' \
|| fail "await_daemon_started did not warn that the pid could not be identified"
assert_equals "" "$NEW_PID" \
"await_daemon_started NEW_PID when the pid is never found but healthz answers — must stay empty, never a guessed pid"
source "$ROOT/scripts/redeploy-fleetd.sh"
}
test_await_daemon_started_call_site_present() {
local src="$ROOT/scripts/redeploy-fleetd.sh" call_line
call_line="$(grep -Fn 'await_daemon_started "$HEALTH_WAIT" "$OLD_PID" "$HEALTH" "$OUT"' "$src" | head -1 | cut -d: -f1 || true)"
[ -n "$call_line" ] \
|| fail "could not find the main flow's await_daemon_started call site in redeploy-fleetd.sh"
}
# fleetd #521 — the swap step's guard, at two levels.
#
# The first two tests call the predicate should_swap() directly. They pin its logic, and that is all
@@ -1226,9 +1481,11 @@ test_poll_health_body_returns_nonzero_when_unreachable() {
test_report_health_call_site_present() {
local src="$ROOT/scripts/redeploy-fleetd.sh" call_line
call_line="$(grep -Fn 'report_health "$HEALTH_BODY" "$HEALTH_CODE" "$OUT" "$HEALTH_WAIT"' "$src" | head -1 | cut -d: -f1 || true)"
# fleetd #603 — moved from a literal main-flow call into await_daemon_started (see its own
# comment above for why); this now finds the call inside that function instead.
call_line="$(grep -Fn 'report_health "$HEALTH_BODY" "$HEALTH_CODE" "$out_file" "$health_wait"' "$src" | head -1 | cut -d: -f1 || true)"
[ -n "$call_line" ] \
|| fail "could not find the main flow's report_health call site in redeploy-fleetd.sh"
|| fail "could not find await_daemon_started's report_health call site in redeploy-fleetd.sh"
}
# fleetd #555 item 8 — `HAD_OLD_PID=0; [ -n "$OLD_PID" ] && HAD_OLD_PID=1`. Measured safe under
@@ -1894,6 +2151,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
@@ -1911,6 +2174,12 @@ test_require_no_build_jar_dies_when_absent
test_require_no_build_jar_accepts_present_jar
test_wait_for_daemon_exit_returns_true_once_pid_clears
test_wait_for_daemon_exit_times_out_if_pid_never_clears
test_wait_for_new_pid_returns_true_once_pid_appears
test_wait_for_new_pid_times_out_if_pid_never_appears
test_await_daemon_started_slow_pid_then_healthy_succeeds
test_await_daemon_started_never_appears_dies_with_log_tail
test_await_daemon_started_pid_never_found_but_healthy_warns_and_survives
test_await_daemon_started_call_site_present
test_swap_ordered_after_wait_and_before_start
test_drain_gate_abort_message_says_no_no_build
test_drain_gate_refusal_build_ran_staged_present