Compare commits

...

16 Commits

Author SHA1 Message Date
Dai Ha ad3d81941f Correct the invariant claimed in the charterBytes comment
CI / shell-tests (pull_request) Failing after 7s
CI / build (pull_request) Successful in 1m36s
CI / contract (pull_request) Successful in 1m42s
The comment said CharterReceipt never pairs a null digest with a non-zero
byte count. It can. compose() derives the digest with digestOf(), which
returns null for blank text, while the byte count is getBytes().length,
which does not. A whitespace-only role charter on a profile with no MCP
produces exactly that pair.

No behaviour change. The gate already omits both fields on that path, which
is the right answer — a size with no digest would describe an artifact we
cannot fingerprint. Only the stated reason was wrong, and a false invariant
in a comment is worse than no comment, because the next reader will widen
the gate on the strength of it.
2026-09-20 16:16:30 +07:00
Dai Ha 8b986a52e0 #604 item 1: fleet_list reports charterBytes alongside charterSha256
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 1m24s
CI / build (pull_request) Successful in 2m26s
CharterReceipt carries a byte count next to its digest, but the roster
projection in SessionManager.rosterView only ever copied the digest
across. A digest tells a lead whether two members' charters match; it
cannot say how far apart they are when they don't. Report charterBytes
too, nested in the same conditional as charterSha256 so the two travel
together: the receipt's own contract only ever pairs a non-null digest
with a real byte count, and a member with no composed charter reports
charterSource alone, unchanged from before.

Tests: the existing charter-receipt roster test now asserts charterBytes
against the receipt's own value (not a literal), plus two new cases —
no charter composed (source "none", no digest, no size) and the receipt
itself absent (no charter keys at all).
2026-09-20 16:14:45 +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
ltms 61115f6f61 Merge #598: fleetd #589 group 3 — wiring-test 5 sites in Fleetd.main()
CI / shell-tests (push) Successful in 6s
CI / contract (push) Successful in 50s
CI / build (push) Successful in 2m9s
Extracts 5 inline lambdas/method-refs in Fleetd.main() to package-private
factories and pins each with a wiring test. Production behaviour unchanged;
the releaseCleanup body moved verbatim.

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

Known gap, tracked separately: fleet.hunters is absent from the live config, so a
hunter cannot spawn on this host until the pool is added after the redeploy. The
role ships correct and inert.
2026-09-19 10:13:26 +02:00
Dai Ha 2e349139e9 #568: add hunter member role
CI / shell-tests (pull_request) Successful in 6s
CI / contract (pull_request) Successful in 1m21s
CI / build (pull_request) Successful in 1m44s
2026-09-19 15:09:07 +07:00
27 changed files with 1411 additions and 79 deletions
+23
View File
@@ -0,0 +1,23 @@
---
name: hunter
description: Sweep one assigned scope for defects and report ranked findings without changes.
---
<!-- CB-617: The model comes from fleetd.yaml because the launch flag overrides model here on both backends. -->
You sweep the assigned package or scope for real defects. Read the full assigned scope before you
judge it. Report several ranked findings when the evidence supports them. Change nothing: do not
edit code, commit, push, or open a pull request.
You may run the build or tests to check a finding. Read the complete output and report the real
result. Do not hide failures with a pipe. State only checks you actually ran. The primary's IDE
tools are not yours. A mounted forge tool may use a blocked credential and fail by design.
Do only the assigned scope. Note anything outside it in one line and do not investigate it further.
Use `fleet_ask{question}` only when a decision belongs to the lead, such as an unclear requirement
or two defensible fixes. Do not ask about something you can decide by reading more code.
Your handoff must name the files you read, each ranked finding or `NO FINDINGS`, the checks you ran,
and any caveat for review.
The launcher provides the required bridge reply instructions for every member.
+7 -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.
@@ -284,6 +285,8 @@ must obey belongs in the charter, not here.
it for a multi-finding sweep hands the worker two contradictory output contracts. That has
already cost three workers' turns: each wrote a good report to its terminal and ended the turn
with no `fleet_reply`, and the scrape returned the tail of the brief instead.
Spawn `implementer` with role `dev`, `reviewer` with role `reviewer`, and `hunter` with role
`hunter`.
- **Primary-side skills** (not delegation playbooks — a worker cannot use them):
`port-to-opencode` (make an OpenCode session a participant in this workspace),
`fleets-status` (report every fleet that shares one LavinMQ instance),
+12 -6
View File
@@ -560,7 +560,7 @@ placement: weighted
# older keys: `leaders:`, `members:`, `leadScan:` and `defaultProfile:`.
#
# A member is anything a lead spawns, and every member has two INDEPENDENT attributes:
# role — which contract: architect, dev or reviewer. It picks the launch charter, the role
# role — which contract: architect, dev, hunter or reviewer. It picks the launch charter, the role
# file, the playbook skill and the authz row.
# profile — which backend: one of the `profiles:` keys above (model, CLI adapter, cost).
# They vary on their own. A reviewer may run on the same profile as the dev whose diff it reads,
@@ -572,13 +572,13 @@ placement: weighted
#
# Each pool lists the profiles that role MAY run on — these are pools, not identities. That is also
# what replaced `defaultProfile:`: an unqualified spawn names a role, and that role's pool supplies
# the candidates, in definition order. A dev and a reviewer staying anonymous is exactly compatible
# with being listed here; the entry key just names the entry.
# the candidates, in definition order. A dev, hunter and reviewer staying anonymous is exactly
# compatible with being listed here; the entry key just names the entry.
fleet:
# Optional launch-charter text, keyed only by the singular role wire names: architect, dev,
# reviewer. Changes are HOT and reach the next spawn without a daemon restart. Do not put secrets
# here: a later launch step writes this text to a world-readable temp file, and ${ENV} interpolation
# is deliberately not supported.
# hunter, reviewer. Changes are HOT and reach the next spawn without a daemon restart. Do not put
# secrets here: a later launch step writes this text to a world-readable temp file, and ${ENV}
# interpolation is deliberately not supported.
charters:
architect: |-
You are an architect in this fleet. You refine work before anyone builds it:
@@ -589,6 +589,9 @@ fleet:
dev: |-
You implement the one unit you were given, and nothing else. You test it,
commit it, and open your own pull request. You never merge.
hunter: |-
You sweep the assigned scope for real defects. You may run the build or tests
to check a finding. You change nothing, and report several ranked findings.
reviewer: |-
You review the diff you were given. You report bugs, risks and missing tests.
You do not change code.
@@ -666,6 +669,9 @@ fleet:
developers:
gx10:
profile: gx10
# hunters:
# gx10:
# profile: gx10 # a hunt may run checks, but never changes code
# reviewers:
# gx10:
# profile: gx10 # the same backend may serve two roles; that is the point
+117 -25
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;
@@ -503,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.
@@ -587,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)) {
@@ -610,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.
@@ -1182,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 or reviewer into an architect. Checked before the worker fallback.
// escalating a dev, hunter or reviewer into an architect. Checked before
// the worker fallback.
return Principal.architect(memberSlotNames.apply(slot), c.terminal(), c.pid());
}
return Principal.worker(c.terminal(), c.pid()); // unforgeable; never token-gated
@@ -42,7 +42,8 @@ public interface MemberLifecycle {
* Try to bind a newly spawned {@code terminal} into the role it was granted.
*
* @return the role this session actually holds: {@code role} unchanged for a role with no
* live slot-binding semantics (dev, reviewer), or when the bind succeeded; a fallback
* live slot-binding semantics (dev, hunter, reviewer), or when the bind
* succeeded; a fallback
* role — never {@code role} — when a slot-bound role (architect) could not be bound.
* Callers must record THIS value on the session, never the requested {@code role}, so
* a later roster read never reports a role the session does not hold (CB-619). In
@@ -20,7 +20,8 @@ import java.util.function.Supplier;
*
* <p>Two halves, split by who owns each:
* <ul>
* <li><b>slots</b> — read from {@code fleet.architects}/{@code developers}/{@code reviewers}
* <li><b>slots</b> — read from {@code fleet.architects}/{@code developers}/
* {@code hunters}/{@code reviewers}
* (see {@link #slots()}), each carrying the {@code profile} reference the spawn lifecycle
* reads when it stands the slot up. <strong>Live, since fleetd #424</strong>: {@link #live}
* re-reads {@code fleet:} on every call, through a supplier the same shape as
@@ -322,7 +323,7 @@ public final class MemberRegistry implements MemberLifecycle {
* CB-619 / fleetd #123: refuse an architect acquire before anything spawns when no configured
* slot carries {@code profile} — the config-gap case from the original defect report (a spawn
* asked for {@code role=architect, profile=sonnet}, and {@code fleet.architects} carried only
* {@code opus} and {@code sol}). A dev/reviewer acquire is always a no-op: those pools are
* {@code opus} and {@code sol}). A dev/hunter/reviewer acquire is always a no-op: those pools are
* placement candidates only (see {@code CompositePeerLauncher}), never a live identity binding,
* so there is nothing here to refuse — an explicit profile outside the pool for those roles is a
* documented operator override, not a defect.
@@ -31,7 +31,7 @@ import java.util.function.Supplier;
* {@code placement:}, and an existing profile's {@code weight} / {@code maxLoad}. Both are
* read through a supplier on {@code CompositePeerLauncher}, which is what makes them hot —
* not the fact that they are config. Most of {@code fleet:} — every role pool
* ({@code architects}/{@code developers}/{@code reviewers}), {@code charters}, and
* ({@code architects}/{@code developers}/{@code hunters}/{@code reviewers}), {@code charters}, and
* {@code tabLabel} — is read the same live way, through the same supplier
* ({@code () -> config.get().fleet()}). {@code architects} in particular is hot for
* <strong>two independent consumers</strong> (fleetd #424): {@code CompositePeerLauncher}
@@ -605,7 +605,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
+ "opened once and needs a restart; the broker URI env-var name kept out of a "
+ "member's environment is read live on every spawn and already applied");
}
// fleetd #333: unlike health/coordinator above, most of `fleet:` (developers, reviewers,
// fleetd #333: unlike health/coordinator above, most of `fleet:` (developers, hunters, reviewers,
// charters, tabLabel) is genuinely hot — ConfigRefTest.aHotChangeIsAppliedAndRead-
// ThroughGet and aCharterChangeIsHotAndReachesTheLiveConfig prove it reaches the live config
// with no restart note. `architects` is hot too, and — since fleetd #424 — hot for BOTH of
@@ -630,7 +630,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
+ "identity map and to auto-launch leads, and neither is rebuilt on reload, so a "
+ "lead added, removed, or given a new tab: label needs a restart — until then it "
+ "stays unrecognised, and a caller from its new tab resolves as a worker, not a "
+ "lead; the rest of fleet: (developers, reviewers, charters, tabLabel) is read "
+ "lead; the rest of fleet: (developers, hunters, reviewers, charters, tabLabel) is read "
+ "live through the supplier on CompositePeerLauncher, and architects is read "
+ "live through that same supplier for placement AND through a separate supplier "
+ "on MemberRegistry for spawn-time identity — both already applied");
@@ -62,7 +62,7 @@ import java.util.regex.PatternSyntaxException;
* @param fleet who the daemon may run and under which role (CB-557). One block replacing
* the former {@code leaders:}, {@code members:}, {@code leadScan:} and
* {@code defaultProfile:}. Role is the containing key — {@code leaders},
* {@code architects}, {@code developers}, {@code reviewers} — and each entry
* {@code architects}, {@code developers}, {@code hunters}, {@code reviewers} — and each entry
* names the {@code profiles:} backend it runs on. See {@link Fleet}
* @param leadHeartbeat opt-in idle-lead heartbeat (CB-551); {@code null} ⇒ off, and an upgraded
* daemon never nudges an idle lead on its own initiative
@@ -1200,6 +1200,7 @@ public record FleetConfig(
* @param leaders panes that orchestrate rather than are orchestrated, keyed by lead name
* @param architects profiles the {@code architect} role may run on
* @param developers profiles the {@code dev} role may run on
* @param hunters profiles the {@code hunter} role may run on
* @param reviewers profiles the {@code reviewer} role may run on
* @param charters optional launch-charter text keyed by singular role wire name
* @param tabLabel template for a member tab's label; {@code {role}}, {@code {profile}},
@@ -1210,6 +1211,7 @@ public record FleetConfig(
public record Fleet(Map<String, Leader> leaders,
Map<String, Slot> architects,
Map<String, Slot> developers,
Map<String, Slot> hunters,
Map<String, Slot> reviewers,
Map<String, String> charters,
String tabLabel) {
@@ -1226,6 +1228,7 @@ public record FleetConfig(
leaders = unmodifiableOrEmpty(leaders);
architects = unmodifiableOrEmpty(architects);
developers = unmodifiableOrEmpty(developers);
hunters = unmodifiableOrEmpty(hunters);
reviewers = unmodifiableOrEmpty(reviewers);
charters = unmodifiableOrEmpty(charters);
tabLabel = (tabLabel == null || tabLabel.isBlank()) ? DEFAULT_TAB_LABEL : tabLabel;
@@ -1240,9 +1243,15 @@ public record FleetConfig(
* constructor: the launcher reads {@code fleet.charters()} from the live config. Jackson
* binds the canonical constructor, so this one cannot swallow an operator's YAML.
*/
public Fleet(Map<String, Leader> leaders, Map<String, Slot> architects,
Map<String, Slot> developers, Map<String, Slot> reviewers,
Map<String, String> charters, String tabLabel) {
this(leaders, architects, developers, null, reviewers, charters, tabLabel);
}
public Fleet(Map<String, Leader> leaders, Map<String, Slot> architects,
Map<String, Slot> developers, Map<String, Slot> reviewers, String tabLabel) {
this(leaders, architects, developers, reviewers, null, tabLabel);
this(leaders, architects, developers, null, reviewers, null, tabLabel);
}
/**
@@ -1264,6 +1273,7 @@ public record FleetConfig(
return switch (role) {
case ARCHITECT -> architects;
case DEV -> developers;
case HUNTER -> hunters;
case REVIEWER -> reviewers;
};
}
@@ -1872,7 +1882,7 @@ public record FleetConfig(
/** The {@code fleet:} child blocks whose direct children are slot names. */
private static final Set<String> FLEET_POOL_KEYS =
Set.of("leaders", "architects", "developers", "reviewers");
Set.of("leaders", "architects", "developers", "hunters", "reviewers");
/**
* Reject a {@code fleet:} role pool whose slot names repeat (CB-548, re-homed by CB-557).
@@ -1882,7 +1892,7 @@ public record FleetConfig(
* daemon would never know. Jackson's YAML parser does not fail on duplicate mapping keys by
* default, so duplicates are caught here, at parse time, before the map is built.
*
* <p>Only the four pools <em>directly under the top-level {@code fleet:}</em> are considered,
* <p>Only the five pools <em>directly under the top-level {@code fleet:}</em> are considered,
* and only their direct child keys (the slot names). A nested field elsewhere, even one also
* named {@code developers:}, is ignored, so parsing of the rest of the config is unaffected.
*
@@ -2049,8 +2059,9 @@ public record FleetConfig(
"defaultProfile", "a role pool under 'fleet:' — an unqualified spawn now names a role,"
+ " and that role's pool supplies the candidate profiles",
"architects", "'fleet.architects'",
"members", "a role pool under 'fleet:' — 'fleet.architects', 'fleet.developers' or"
+ " 'fleet.reviewers'; the role is the containing key, not a 'role:' field",
"members", "a role pool under 'fleet:' — 'fleet.architects', 'fleet.developers',"
+ " 'fleet.hunters' or 'fleet.reviewers'; the role is the containing key, not"
+ " a 'role:' field",
"leaders", "'fleet.leaders'",
"leadScan", "'fleet.leaders.<name>.tabPrefix' and '.scanIntervalSeconds' — lead"
+ " discovery is now configured on the lead it discovers");
@@ -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");
@@ -2134,8 +2159,9 @@ public final class FleetMcp {
return tool(FleetTool.SPAWN.wireName(),
"Spawn a new off-subscription member session. A member has two independent attributes: "
+ "role (what it is for) and profile (which backend it runs on). Pass role to pick "
+ "the contract — 'dev' implements a unit and opens its own PR, 'reviewer' reviews a "
+ "diff it did not write, 'architect' refines a ticket before anyone builds it; omit "
+ "the contract — 'dev' implements a unit and opens its own PR, 'hunter' sweeps a "
+ "scope without changing it, 'reviewer' reviews a diff it did not write, 'architect' "
+ "refines a ticket before anyone builds it; omit "
+ "it for 'dev'. Pass profile (from fleet_profiles) to pick the backend, or omit it "
+ "for the default. The two are independent: a reviewer may run on the same profile "
+ "as the dev it reviews. The member opens your current directory by default; pass "
@@ -2152,7 +2178,7 @@ public final class FleetMcp {
+ "one. Returns the member's sessionId (use with fleet_send) and paneId (use with "
+ "fleet_stop).",
objectSchema(Map.of(
"role", stringProp("What the member is for: architect, dev or reviewer (default dev)"),
"role", stringProp("What the member is for: architect, dev, hunter, or reviewer (default dev)"),
"profile", stringProp("Which backend to run it on (omit for the default profile)"),
"cwd", stringProp("Working directory for the member (omit to inherit yours)"),
"worktree", Map.of("type", "string", "description", "'true' or a ticket slug — requests an isolated git worktree"),
@@ -2264,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 "
@@ -29,8 +29,8 @@ public enum MemberRole {
* <p>Reads the repo and writes analysis. Never commits code and never opens a pull request —
* an architect that starts implementing has stopped doing the job that makes it useful.
*
* <p>Architects are the one member kind declared in config, because a lead addresses the same
* slots across many tickets and needs a stable name for them.
* <p>Architects are the one member kind with live slot binding, because a lead addresses the
* same slots across many tickets and needs a stable name for them.
*/
ARCHITECT,
@@ -43,6 +43,14 @@ public enum MemberRole {
*/
DEV,
/**
* Sweeps an assigned package for defects and reports several ranked findings.
*
* <p>Never changes code, commits, or opens a pull request. A hunt gathers evidence, which can
* include running the build, but leaves every fix to a later implementation unit.
*/
HUNTER,
/**
* Reviews a diff it did not write and reports one structured finding.
*
@@ -59,7 +67,7 @@ public enum MemberRole {
/**
* The {@code fleet:} block that holds this role's pool — {@code architects},
* {@code developers}, {@code reviewers}.
* {@code developers}, {@code hunters}, {@code reviewers}.
*
* <p>Plural, and not always the wire name: the pool of things a {@code dev} may run on reads
* naturally as {@code developers:}. The wire name stays the singular {@code dev}, because that
@@ -69,6 +77,7 @@ public enum MemberRole {
return switch (this) {
case ARCHITECT -> "architects";
case DEV -> "developers";
case HUNTER -> "hunters";
case REVIEWER -> "reviewers";
};
}
@@ -871,10 +871,23 @@ public final class SessionManager implements TurnListener {
// digest lets a lead tell at a glance whether all members got the same charter; the source
// records whether a role charter was configured ("fleet.charters.<role>") or only the reply
// charter was composed ("none").
// #604: charterBytes rides along with charterSha256, not with charterSource — it is only
// meaningful as the digest's companion (a length turns "they differ" into "by how much").
// A member with no composed charter reports charterSource and nothing else, as before.
//
// Gating on the DIGEST rather than on the receipt is deliberate, and the two are not
// always null together. CharterReceipt.compose() derives the digest with digestOf(), which
// returns null for BLANK text, while the byte count is composed.getBytes().length, which
// does not. So a whitespace-only charter (a blank fleet.charters.<role> on a profile with
// no MCP, so no reply charter is appended) yields a null digest beside a non-zero size.
// Reporting a size with no digest would say "they differ by N bytes" about an artifact we
// cannot fingerprint, so this gate omits both. Never widen it to the receipt-level null
// check without deciding what that case should report.
if (session.charterReceipt() != null) {
m.put("charterSource", session.charterReceipt().charterSource());
if (session.charterReceipt().charterSha256() != null) {
m.put("charterSha256", session.charterReceipt().charterSha256());
m.put("charterBytes", session.charterReceipt().charterBytes());
}
}
m.put("liveStatus", live == null ? "unknown" : live.status().name().toLowerCase());
@@ -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,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");
}
}
@@ -345,7 +345,7 @@ class FleetConfigTest {
IllegalStateException unknownError = assertThrows(IllegalStateException.class,
() -> FleetConfig.load(unknown).validateCharters());
assertTrue(unknownError.getMessage().contains("architetc"));
assertTrue(unknownError.getMessage().contains("[architect, dev, reviewer]"));
assertTrue(unknownError.getMessage().contains("[architect, dev, hunter, reviewer]"));
}
/**
@@ -812,13 +812,17 @@ class FleetConfigTest {
reviewers:
b:
profile: sonnet
hunters:
c:
profile: sonnet
""");
FleetConfig cfg = FleetConfig.load(f);
assertEquals(List.of("sonnet"), cfg.fleet().profilesFor(MemberRole.DEV));
assertEquals(List.of("sonnet"), cfg.fleet().profilesFor(MemberRole.HUNTER));
assertEquals(List.of("sonnet"), cfg.fleet().profilesFor(MemberRole.REVIEWER));
assertTrue(cfg.fleet().profilesFor(MemberRole.ARCHITECT).isEmpty());
assertEquals(List.of(MemberRole.DEV, MemberRole.REVIEWER), cfg.fleet().rolesConfigured());
assertEquals(List.of(MemberRole.DEV, MemberRole.HUNTER, MemberRole.REVIEWER), cfg.fleet().rolesConfigured());
}
/** The case the two axes exist for: one backend, two roles, and neither is a duplicate. */
@@ -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 -----------
}
@@ -1937,7 +1937,7 @@ class FleetMcpTest {
null, null, null, null, null, null);
assertEquals(Boolean.TRUE, res.isError());
assertTrue(textOf(res).contains("architect, dev, reviewer"), textOf(res));
assertTrue(textOf(res).contains("architect, dev, hunter, reviewer"), textOf(res));
}
// ── CB-619 / fleetd #123: a spawn asking for a role its profile has no slot for must be
@@ -556,6 +556,31 @@ class ClaudeCodeLauncherTest {
"no --agent flag when the role has no agent-definition file");
}
@Test
void hunterRoleUsesItsAgentFileAndStopsUsingItWhenRemoved(@TempDir Path cwd) throws Exception {
FakeHerdr herdr = new FakeHerdr();
Path agentFile = Files.createDirectories(cwd.resolve(".claude/agents")).resolve("hunter.md");
Files.writeString(agentFile, "---\nname: hunter\n---\nSweep for defects.");
FleetConfig.Profile cfg = new FleetConfig.Profile(
"sonnet", "http://gx00.gw:8000", null, null, "FLEETD_WORKER_TOKEN",
List.of("claude"), "tab", "fleetd-workers", "w #{n}", null, null, null);
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
svc.spawn(new SpawnRequest("sonnet", cwd.toString(), null, null, null, MemberRole.HUNTER));
List<String> args = spawnedArgs(herdr);
int flag = args.indexOf("--agent");
assertTrue(flag >= 0, "the hunter role reaches its agent-definition file: " + args);
assertEquals("hunter", args.get(flag + 1));
Files.delete(agentFile);
svc.spawn(new SpawnRequest("sonnet", cwd.toString(), null, null, null, MemberRole.HUNTER));
assertFalse(spawnedArgs(herdr).contains("--agent"),
"the hunter role no longer gets an agent when its file is removed");
}
private ClaudeCodeLauncher multiProfile(FakeHerdr herdr) {
FleetConfig.Profile gx10 = new FleetConfig.Profile("gx10", "http://gx10.gw:8000", "coder",
null, "FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers", "w #{n}", null, null, null);
@@ -10,12 +10,13 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
class MemberRoleTest {
@Test
void theThreeRolesAreArchitectDevAndReviewer() {
assertEquals(3, MemberRole.values().length,
void theFourRolesAreArchitectDevHunterAndReviewer() {
assertEquals(4, MemberRole.values().length,
"a new role changes the charter, the role file, the skill and the authz row — "
+ "adding one is a deliberate act, so this count is meant to fail first");
assertEquals("architect", MemberRole.ARCHITECT.wireName());
assertEquals("dev", MemberRole.DEV.wireName());
assertEquals("hunter", MemberRole.HUNTER.wireName());
assertEquals("reviewer", MemberRole.REVIEWER.wireName());
}
@@ -30,6 +31,7 @@ class MemberRoleTest {
void parseIsCaseInsensitiveAndTrimsSurroundingSpace() {
assertSame(MemberRole.ARCHITECT, MemberRole.parse("Architect"));
assertSame(MemberRole.DEV, MemberRole.parse(" DEV "));
assertSame(MemberRole.HUNTER, MemberRole.parse("HuNtEr"));
assertSame(MemberRole.REVIEWER, MemberRole.parse("ReViEwEr"));
}
@@ -38,7 +40,7 @@ class MemberRoleTest {
IllegalArgumentException e =
assertThrows(IllegalArgumentException.class, () -> MemberRole.parse("archtiect"));
assertTrue(e.getMessage().contains("archtiect"), e.getMessage());
assertTrue(e.getMessage().contains("architect, dev, reviewer"),
assertTrue(e.getMessage().contains("architect, dev, hunter, reviewer"),
"a typo in config should be fixable from the message alone: " + e.getMessage());
}
@@ -328,9 +328,10 @@ class SessionManagerTest {
void rosterViewExposesTheCharterReceiptButNeverTheCharterText() {
// The roster (fleet_list and GET /members both render through rosterView) must let a lead
// see which charter a member got, without ever carrying the charter prose itself (CB-571).
String composed = "role charter\n\nreply";
CharterReceipt receipt = CharterReceipt.compose(MemberRole.DEV, "prof", "role charter", composed);
MemberSession s = new MemberSession("p1", "term1", "prof", MemberRole.DEV, "/cwd", null,
0, 0, 0, MemberSession.State.READY, null, null,
CharterReceipt.compose(MemberRole.DEV, "prof", "role charter", "role charter\n\nreply"), null);
0, 0, 0, MemberSession.State.READY, null, null, receipt, null);
Map<String, Object> view = SessionManager.rosterView(s, null);
@@ -338,10 +339,49 @@ class SessionManagerTest {
"the config key that supplied the role charter is reported");
assertEquals(CharterReceipt.digestOf("role charter\n\nreply"), view.get("charterSha256"),
"the digest of the exact composed charter bytes is reported");
// #604: the byte count rides alongside the digest, and must match what the receipt itself
// carries (not a hardcoded literal) so a bug that reads the wrong field is caught.
assertEquals(receipt.charterBytes(), view.get("charterBytes"),
"the exact composed byte count is reported, read from the receipt");
assertEquals(composed.getBytes(java.nio.charset.StandardCharsets.UTF_8).length, view.get("charterBytes"),
"the byte count is the real UTF-8 length of the composed charter");
assertFalse(view.values().toString().contains("role charter"),
"the roster row must not embed the charter text itself");
}
@Test
void rosterViewOmitsCharterBytesAndDigestWhenNoCharterWasComposed() {
// #604: a member with no role charter and no reply charter (composed == null) still reports
// charterSource ("none"), but neither a digest nor a size — the digest is absent, and the
// size only ever accompanies a real digest. This must not be satisfiable by code that always
// writes charterBytes.
CharterReceipt receipt = CharterReceipt.compose(MemberRole.DEV, "prof", null, null);
MemberSession s = new MemberSession("p1", "term1", "prof", MemberRole.DEV, "/cwd", null,
0, 0, 0, MemberSession.State.READY, null, null, receipt, null);
Map<String, Object> view = SessionManager.rosterView(s, null);
assertEquals(CharterReceipt.NO_SOURCE, view.get("charterSource"),
"no configured role or reply charter reports the explicit \"none\" source");
assertFalse(view.containsKey("charterSha256"), "no digest is reported when no charter was composed");
assertFalse(view.containsKey("charterBytes"), "no byte count is reported when no charter was composed");
}
@Test
void rosterViewOmitsCharterFieldsEntirelyWhenTheReceiptItselfIsAbsent() {
// #604 acceptance criterion 3: charterReceipt() can be null on its own (a session recorded
// before CB-571, or a launcher that never composed one) — the outer null-guard must still
// suppress charterSource, charterSha256 AND charterBytes together.
MemberSession s = new MemberSession("p1", "term1", "prof", MemberRole.DEV, "/cwd", null,
0, 0, 0, MemberSession.State.READY, null, null, null, null);
Map<String, Object> view = SessionManager.rosterView(s, null);
assertFalse(view.containsKey("charterSource"), "no charter fields at all when the receipt is null");
assertFalse(view.containsKey("charterSha256"), "no charter fields at all when the receipt is null");
assertFalse(view.containsKey("charterBytes"), "no charter fields at all when the receipt is null");
}
@Test
void aNullTerminalFromThePrimaryIsANoOpEvenWithSessionsRegistered() {
// The primary resolves to a Principal with no terminal, and FleetMcp's context extractor
+54 -3
View File
@@ -169,7 +169,55 @@ hash256() {
# on PATH). "absent" must never be the answer for a file that exists — that conflation, on Linux,
# was the whole defect this ticket fixes.
jar_id() { local f="${1:-$JAR}"; [ -f "$f" ] && hash256 "$f" || echo "absent"; }
running_pid() { pgrep -f "$PATTERN" || true; }
# fleetd #593 — `pgrep -f "$PATTERN"` matches ANY process whose full command line CONTAINS the
# pattern text, and that is not the same thing as "is the daemon". A shell that merely embeds the
# pattern as literal text — a human typing this exact investigation by hand, an ssh-shaped
# `sh -c '...; ...'`, a pipeline, or any other non-exec'ing shell that never replaced itself with
# the pattern-holding command — still shows up in that match, and it is the INSTRUMENT, not the
# daemon. Measured live on this Mac: `sh -c 'echo "target/fleetd.jar" >/dev/null; sleep 30' &`
# leaves a real `sh` process alive (it forks for the `sleep`, it does not exec into it) whose own
# `ps -o args` is `sh -c echo "target/fleetd.jar" >/dev/null; sleep 30` — `pgrep -f "$PATTERN"`
# matches that line right alongside the real `java -jar target/fleetd.jar` process. `pgrep -c`
# (an in-one-call count) does not exist on BSD/macOS at all, so this cannot be fixed by switching
# pgrep flags — it has to filter what pgrep already found, after the fact, in a way that still
# runs on BSD.
#
# fleetd #593 CORRECTION 1 — the first cut of this filter kept everything whose `comm` was NOT a
# shell name (a denylist: sh/bash/zsh/dash/ksh). Two holes in that, both the same false-positive
# shape the ticket exists to remove in the first place:
# 1. a pid `pgrep` just listed can exit before the `ps -o comm=` lookup runs; on a gone pid `ps`
# prints nothing, `comm` ends up empty, and an empty string matches none of the denied shell
# names — so a pid that no longer exists was still counted.
# 2. the denylist only knows the shells someone thought to name. `ssh`, `perl`, `python3`,
# `ruby`, `tail` — anything else that carries the pattern in its own argv — was still
# counted right along with the real daemon, and the ticket names `ssh` as a live route.
# Both close with the same change: allowlist `comm = java` instead of denying shells. Measured on
# the live daemon: `pid=30224 comm=java`. An empty comm (hole 1) is not `java` either, so it is
# excluded for free — no separate "is this pid still alive" check needed.
#
# The objection, because it is real: an allowlist can UNDER-count. If fleetd ever stops being
# launched as `java -jar ...` — a native image, a renamed launcher — `running_pid()` silently
# returns nothing and `assert_single_daemon` stops noticing a second daemon at all. For a guard,
# that false-negative direction is the worse one to be wrong in. This is not a new assumption,
# though: `PATTERN='target/fleetd.jar'` two lines up already assumes the daemon is a jar, which
# is only ever run by `java`. If that launch method changes, `PATTERN` stops matching anything
# before this allowlist would ever get the chance to be wrong — the allowlist rides on the same
# assumption that is already load-bearing, it does not add a new one. Whoever changes the launch
# method needs to update both `PATTERN` and this allowlist together.
running_pid() {
local pid comm out=''
for pid in $(pgrep -f "$PATTERN" 2>/dev/null || true); do
comm="$(ps -o comm= -p "$pid" 2>/dev/null || true)"
comm="${comm##*/}"
comm="${comm#-}"
# Allowlist, not a denylist of wrappers — see the CORRECTION 1 comment above. Anything that
# is not literally `java` is excluded, including an empty comm from a pid that already exited.
[ "$comm" = java ] || continue
out="$out$pid"$'\n'
done
printf '%s' "$out"
}
# fleetd #493 — three small, independently testable pieces of "never build into the path a
# running process holds":
@@ -507,8 +555,11 @@ assert_single_daemon() {
die "more than one fleetd process is running after this restart (pids: $(printf '%s' "$pids" | tr '\n' ' ')).
This is the exact failure a racing supervisor produces: the OLD jar was revived by its
supervisor while this script started a NEW copy. Two daemons on one herdr session kill
each other's members. Investigate with 'pgrep -f \"$PATTERN\"' and stop the wrong one by
hand — do not assume either pid is the one you want."
each other's members. Investigate with 'ps -eo pid,comm,args | grep -F \"$PATTERN\"' and
check the COMM column of each hit yourself before acting — a bare 'pgrep -f \"$PATTERN\"'
(fleetd #593) can match the very shell you type it into, not just the daemon, so it is not
safe remediation advice on its own. Stop the wrong one by hand — do not assume either pid
is the one you want."
fi
}
+130
View File
@@ -433,6 +433,130 @@ test_assert_single_daemon_rejects_two_pids() {
printf '%s' "$output" | grep -qF '4343' || fail "refusal message does not list the pids it found"
}
# fleetd #593 instance 2 — `running_pid()` used to be a bare `pgrep -f "$PATTERN"`, which matches
# ANY process whose full command line contains the pattern TEXT, including a shell that merely
# embeds it as literal text rather than being the daemon. Measured live on this Mac: `pgrep -c`
# (a one-call count) does not exist on BSD at all, and `bash -c "<single command>"` execs in place
# so no parent shell survives to hold the pattern — which is exactly why the defect did not
# reproduce from a plain script and needs a wrapper shaped like this instead. A `sh -c '...; ...'`
# with MORE THAN ONE statement does not get that exec-in-place treatment: the shell forks a child
# for the second statement and stays alive itself, holding the whole `-c` string — pattern text
# included — in its own `ps -o args`, for as long as it runs. That is the same shape an
# `ssh host "…; …"` wrapper or a hand-typed pipeline leaves behind. Before the fix this test would
# have found the wrapper's pid in running_pid()'s output; it must not.
test_running_pid_excludes_self_matching_wrapper_shell() {
local before after wrapper_pid
before="$(running_pid)"
sh -c 'echo "target/fleetd.jar" >/dev/null; sleep 20' &
wrapper_pid=$!
sleep 0.3
after="$(running_pid)"
kill "$wrapper_pid" 2>/dev/null || true
wait "$wrapper_pid" 2>/dev/null || true
[ "$after" = "$before" ] \
|| fail "running_pid() counted a self-matching wrapper shell (pid $wrapper_pid, holding the pattern as literal text in its own argv, not the daemon): before=[$before] after=[$after]"
}
# fleetd #593 CORRECTION 1 — the round-1 version of this test gave its standin an argv[0]
# containing the pattern text (via `exec -a`) and left `comm` as whatever that override produced,
# which was never `java`. That was fine for a denylist-of-shells filter, but the allowlist below
# now requires `comm = java` specifically, so the standin here must actually carry that comm, not
# just avoid being a shell. `exec -a java` overrides argv[0] to `java` while the process itself
# stays a genuine, harmless `sh`; combining it with the same non-exec'ing multi-statement shape
# the wrapper-shell test above uses keeps the pattern text in the process's own `ps -o args` for
# as long as it runs. Measured live on this Mac (BSD/macOS: `ps -o comm=` here reflects argv[0]):
# `comm=java`, `args` contains the pattern, `pgrep -f "$PATTERN"` finds it. Copying a real system
# binary into a scratch path and executing it from there was tried first, for a more literal
# stand-in daemon, and the OS killed it outright (SIGKILL, exit 137 — almost certainly a
# code-signing check on a relocated binary); `exec -a` needs no binary of its own and nothing
# under a scratch directory, and it is the technique CORRECTION 1 names as the right one.
#
# This is the one live-process test in this file whose result could differ on Linux: Linux sets
# `comm` from the actually-executed binary's own path, not from `exec -a`'s argv[0] override (BSD
# ties `comm` to argv[0], which is what makes this technique work here) — so on Linux this
# specific fixture might report `comm=sh`, not `comm=java`, even though the REAL daemon (a literal
# `java -jar target/fleetd.jar` process, never fabricated) is unaffected either way. I could not
# verify this fixture's behavior on Linux, so test_running_pid_counts_a_pid_whose_comm_is_java
# below backstops the same claim (the allowlist admits a pid whose comm is `java`) with a stubbed
# `ps`, which is identical bash on every platform and carries no such platform question.
test_running_pid_finds_a_real_java_named_second_process() {
local before after standin_pid
before="$(running_pid)"
( exec -a java sh -c 'echo "target/fleetd.jar" >/dev/null; sleep 20' ) &
standin_pid=$!
sleep 0.3
after="$(running_pid)"
kill "$standin_pid" 2>/dev/null || true
wait "$standin_pid" 2>/dev/null || true
printf '%s\n' "$after" | grep -qxF "$standin_pid" \
|| fail "running_pid() did not find a real second process (pid $standin_pid, comm forced to 'java' via exec -a) whose own argv holds the pattern: before=[$before] after=[$after]"
}
# fleetd #593 CORRECTION 1, hole 2 — the round-1 filter denied known shell names (sh/bash/zsh/
# dash/ksh) and counted everything else. `ssh`, `perl`, `python3`, `ruby`, `tail` — anything not on
# that list, carrying the pattern in its own argv — was still counted right alongside the real
# daemon, and the ticket names `ssh` as a live route. Stubbing `pgrep`/`ps` (rather than spawning a
# real perl/ssh process) pins the exact discriminator this correction is about — comm, not the
# caller's shape — deterministically on every platform, with no dependency on perl/python3/ruby
# being installed in whatever environment runs this suite, and no dependency on how a given OS
# derives `comm` for a fabricated process (see the comment above
# test_running_pid_finds_a_real_java_named_second_process for why that matters here).
test_running_pid_drops_a_pid_whose_comm_is_not_java() {
pgrep() { printf '4242\n'; }
ps() { printf 'perl\n'; }
local found
found="$(running_pid)"
unset -f pgrep ps
[ -z "$found" ] \
|| fail "running_pid() counted pid 4242 whose comm is 'perl', not 'java' — denying known shell names does not exclude a non-shell wrapper such as ssh or perl (fleetd #593 CORRECTION 1): found=[$found]"
}
# fleetd #593 CORRECTION 1, hole 1 — pgrep can list a pid that exits before the following
# `ps -o comm=` lookup runs; on a gone pid `ps` prints nothing, so `comm` comes back empty. Under
# the round-1 denylist an empty string matched none of the denied shell names, so the dead pid was
# still counted — the exact false-positive shape the ticket exists to remove, just rarer. The
# allowlist fixes this for free: an empty comm is not `java` either.
test_running_pid_drops_a_pid_that_exited_before_the_comm_lookup() {
pgrep() { printf '4242\n'; }
ps() { :; } # a pid that no longer exists: the real `ps -p <gone>` prints nothing and this mirrors that
local found
found="$(running_pid)"
unset -f pgrep ps
[ -z "$found" ] \
|| fail "running_pid() counted pid 4242 whose comm lookup came back empty (the pid had already exited before the lookup ran) — an empty comm must not pass the allowlist (fleetd #593 CORRECTION 1): found=[$found]"
}
# The positive backstop for both stubbed tests above, and for
# test_running_pid_finds_a_real_java_named_second_process on whatever platform that live fixture
# does not itself carry comm=java: the allowlist must still ADMIT the one comm value the real
# daemon actually has. Measured on the real, currently-running daemon on this Mac: `comm=java`.
test_running_pid_counts_a_pid_whose_comm_is_java() {
pgrep() { printf '4242\n'; }
ps() { printf 'java\n'; }
local found
found="$(running_pid)"
unset -f pgrep ps
printf '%s\n' "$found" | grep -qxF '4242' \
|| fail "running_pid() did not count pid 4242 whose comm is 'java' — the daemon's own name must pass the allowlist: found=[$found]"
}
# fleetd #593 instance 3 — assert_single_daemon's refusal message used to tell the operator to
# "Investigate with 'pgrep -f \"\$PATTERN\"'", which — typed by hand or over ssh — is precisely the
# self-matching invocation instance 2 above fixes. A source-text check, the same technique
# test_no_error_lines_message_gated_by_drain_state uses: this is prose inside a die() call, never
# reached by sourcing (the SOURCED guard stops before the main flow, and this text only prints
# from inside a call assert_single_daemon makes when it is already refusing).
test_die_message_does_not_recommend_bare_pgrep_as_remediation() {
local src="$ROOT/scripts/redeploy-fleetd.sh" block bad
block="$(grep -A6 -F 'racing supervisor produces' "$src" || true)"
[ -n "$block" ] || fail "could not find the assert_single_daemon refusal message in redeploy-fleetd.sh"
bad="$(printf '%s' "$block" | grep -F "Investigate with 'pgrep -f" || true)"
[ -z "$bad" ] \
|| fail "assert_single_daemon's die message still hands the operator a bare 'pgrep -f \"\$PATTERN\"' as remediation (fleetd #593) — that is exactly the self-matching invocation"
printf '%s' "$block" | grep -qF 'fleetd #593' \
|| fail "assert_single_daemon's die message does not say in words that a pattern can match the caller (fleetd #593)"
}
# fleetd #511 — jar_id()'s no-argument default was unpinned by any test: nothing proved it reports
# $JAR (the live path) rather than $JAR_STAGED. Both halves matter, so this pins both: the bare call
# must hash the live jar, and an explicit path argument must hash THAT file, not fall back to $JAR.
@@ -1894,6 +2018,12 @@ test_require_drivable_supervisor_accepts_known_kinds
test_count_daemon_pids
test_assert_single_daemon_accepts_one_pid
test_assert_single_daemon_rejects_two_pids
test_running_pid_excludes_self_matching_wrapper_shell
test_running_pid_finds_a_real_java_named_second_process
test_running_pid_drops_a_pid_whose_comm_is_not_java
test_running_pid_drops_a_pid_that_exited_before_the_comm_lookup
test_running_pid_counts_a_pid_whose_comm_is_java
test_die_message_does_not_recommend_bare_pgrep_as_remediation
test_jar_id_defaults_to_live_and_reports_explicit_path
test_hash256_computes_a_real_sha256
test_jar_id_reports_absent_for_missing_file