Compare commits

..

21 Commits

Author SHA1 Message Date
Dai Ha d4f93a7b13 fleetd #449: fix stale herdr protocol 14 javadocs/assertion, diagnose and fix the timing-raced AgentControlContractTest, select contract tests by tag in CI
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Successful in 2m9s
- HerdrClient.java, HerdrCodec.java, HerdrContractTest.java: the herdr port to
  protocol 19 (CB-521) left the client javadoc and the contract test's own
  assertion still saying protocol 14 / herdr 0.7.0. Updated to 19 / 0.8.0 and
  renamed pingReturnsProtocol14 -> pingReturnsProtocol19. Verified the
  assertion is real by temporarily changing the expected value to 20 (fails),
  then restoring 19 (passes).

- AgentControlContractTest.java: tabCreateInjectsEnvIntoTheSeedShell was
  failing, not skipping, on a host with a live herdr socket. Diagnosed with a
  temporary instrumented run (not committed) that polled the pane every
  200ms before and after sending input: the seed shell reliably takes ~2.5s
  to reach its prompt (measured 3x), while the test's fixed 1000ms sleep
  raced that startup. Input typed too early was swallowed by the shell's own
  startup, leaving the typed line followed by the "Restored session" banner
  and no command output — indistinguishable at a glance from the env map
  never reaching the shell. Once the shell was actually ready, the injected
  env value showed up in ~200ms, ruling out an env-seam defect. Replaced both
  fixed sleeps with bounded polling on the actual conditions (pane text
  settling, then the expected output appearing). Ran the fixed test 3x
  standalone, all green.

- .gitea/workflows/ci.yml: the "Contract tests" step ran exactly one class by
  name (-Dtest=AmqpReplyInboxContractTest), silently excluding every other
  @Tag("contract") test from CI including the herdr ones above -- which is
  how the stale protocol 14 assertion went unnoticed. Changed to
  -Dgroups=contract, which selects the whole tagged group and picks up
  future contract tests automatically.
2026-09-10 17:14:09 +07:00
ltms 822327eed5 Merge #447: pin the place()-to-spawn() window PlacementDecision closes (fleetd #444)
CI / contract (push) Successful in 1m33s
CI / build (push) Successful in 1m36s
Verified by the lead on the exact tree that lands (head e4c703a, base 82fae94 is
an ancestor, so this is the tree I measured):

  FULL BUILD  Tests run: 1575, Failures: 0, Errors: 0, Skipped: 0  BUILD SUCCESS
              compile errors: 0
  CONTROL     CompositePeerLauncherTest  Tests run: 76, Failures: 0  -> GREEN

  M1  the 2-arg spawn re-enters the 1-arg spawn (the inherited default this
      ticket forbids for a multi-profile launcher)
      -> KILLED  Errors: 1
      CompositePeerLauncherTest
        .spawnHonorsAPlacementDecisionEvenAfterItsProfileIsQuarantinedInTheWindowAfterPlace

  M2  drop the stamping: route to the decided profile but do not carry it
      (SpawnRequest routedReq = req)
      -> KILLED  Failures: 1
      same test method

Both mutations proved applied by printing the mutated method, and the tree was
restored clean after each (git status --porcelain empty).

Round 1 of this PR had a fixture weakness I found by mutation: StubLauncher's own
fallback default was "sol", the same profile place() decides, so an unstamped
request landed on spawnCount("sol") by coincidence and M2 survived. e4c703a gives
the adapter "b" as its fallback instead. One fixture now kills both mutations.

src/main is javadoc-only in this PR: 12 added lines, 0 added code lines, measured
by filtering the main-side diff.
2026-09-10 12:00:59 +02:00
Dai Ha e4c703a51a fleetd #444: separate the adapter's fallback default from the decided profile
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Successful in 1m56s
Review found the fixture's StubLauncher fell back to 'sol' too — the
same profile place() decides — so an UNSTAMPED request could land on
spawnCount('sol') by coincidence, and the assertion's claim that the
request 'actually carried sol' was unproven. Dropping the stamping
(SpawnRequest routedReq = req) while keeping the routing survived the
test unchanged.

Fix: give the adapter 'b' as its own fallback default instead, so an
unstamped request counts against 'b', not 'sol'. Verified both
mutations against the single test in isolation:
  - drop-stamping (routedReq = req): RED, expected <sol> but was <b>
  - re-entering (return spawn(req.withProfile(decision.profile())))
    i.e. M1 from the first round: still RED, PlacementException
    naming the now-quarantined 'sol'
Restored both; full suite green at 1575 tests.

No changes to src/main — PeerLauncher's javadoc from the first round
is unchanged.
2026-09-10 16:50:39 +07:00
Dai Ha 3f036b2a62 fleetd #444: pin the place()-to-spawn() window PlacementDecision closes
CI / contract (pull_request) Successful in 1m4s
CI / build (pull_request) Successful in 1m53s
Add a test that resolves place(role) while nothing is quarantined,
then quarantines the resolved profile's credential BEFORE spawning
against the held PlacementDecision. CompositePeerLauncher.spawn(req,
decision) must still honor the decision and land on the quarantined
profile, since it never re-runs the explicit-profile enforce* checks.

Verified the test kills the regression: with the override's body
replaced by the re-entering spawn(req.withProfile(...)) form, this
exact test goes RED with a PlacementException naming the now-
quarantined profile; restored, the full suite is green (1575 tests).

Also documents on PeerLauncher's default spawn(req, decision) that a
launcher routing across more than one profile MUST override it,
naming the four enforce* checks the default's re-entry re-applies.
2026-09-10 16:41:52 +07:00
ltms 82fae94c55 Merge #445: pin every startup report call in Fleetd.main (fleetd #442)
CI / contract (push) Successful in 48s
CI / build (push) Successful in 1m59s
Test written by a worker whose backend died before it could report; evidence
re-run by the lead against the merged tree.

Verified: merge of current main clean (0 conflicts); full build 1574 tests,
0 failures, BUILD SUCCESS, 0 compile errors; control green; deleting each of
reportGitHostShape, reportMemberTrustModel, reportMemberCredentialsGap and
reportExhaustedPatternGap from main() is KILLED by
mainReportsEveryStartupGapBeforeValidationAborts.
2026-09-10 11:30:31 +02:00
Dai Ha b1f34c2e6b fleetd #442: drop the unused java.util.List import
CI / contract (pull_request) Successful in 47s
CI / build (pull_request) Successful in 2m8s
The new test never names List — only ListAppender, which has its own
import. An unused import is an IDE warning, and this repo treats
warnings as gates. No behaviour change: FleetdStartupReportTest still
runs 1 test, 0 failures, BUILD SUCCESS, 0 compile errors.
2026-09-10 16:29:56 +07:00
ltms 3f807d9f1b Merge #443: derive coordinator.heldDurable from queue durability + ack mode (fleetd #440)
CI / contract (push) Successful in 51s
CI / build (push) Successful in 1m35s
Found by the fleet01 lead reviewing #438 after I had merged it. Verified
independently before merging.

The implementation choice is the load-bearing part: LeadMailbox.own() now
assigns queueDeclare's durable flag and basicConsume's autoAck flag to named
locals, passes those SAME locals into the two real AMQP calls (:203/:204), and
derives heldDurable from them (:207). So the reported fact cannot drift from a
duplicate constant - a mutation to either call's argument moves the behaviour
and the report together. LeadChannel.heldDurable() is abstract, so a future
implementer gets a compile error rather than a silent default.

My own battery, merged tree, control green, tree restored clean:
- M1 revert to the literal true -> KILLED by
  FleetMcpTest.listReportsHeldDurableFalseWhenTheChannelSaysMailIsNotDurable
- M3 heldCount forced to 0 (a half the worker did not touch) -> KILLED by
  FleetMcpTest.listReportsAnHonestHeldCountAndDurabilityNotJustPendingZero
- M2 break the derivation itself -> SURVIVED under plain clean install, exactly
  as the worker reported. LeadMailbox needs a real broker, so the only test that
  reaches the real queueDeclare/basicConsume is @Tag("contract"), excluded from
  the default build. The worker ran that arm with -Pcontract and got 3 reds
  including its own new test. Pre-existing structural limit of this class, not
  introduced here, and the worker flagged it rather than hiding it.

Full build, my own run: Tests run: 1573, Failures: 0, Errors: 0, Skipped: 0 -
BUILD SUCCESS, 0 compile errors.

Not blocking, noted for a possible follow-up: one boolean over two independent
facts cannot say WHICH fact was lost. The merged field is still strictly better
than the literal it replaces, because it can now go false at all.
2026-09-10 11:21:41 +02:00
Dai Ha e70263062c fleetd #442: pin startup report calls
CI / contract (pull_request) Successful in 1m15s
CI / build (pull_request) Successful in 1m36s
2026-09-10 14:41:10 +07:00
ltms 1a1e586b62 Merge #433: carry the PlacementDecision instead of re-resolving the profile (fleetd #425)
CI / contract (push) Successful in 45s
CI / build (push) Successful in 1m47s
Round 4 pins the fix. Verified independently before merging.

My own battery in the worker's tree (control green, tree restored clean):
- M1 revert SessionManager:677 to launcher.spawn(spawnReq) -> KILLED by
  SessionManagerTest.acquireWithWorktreeSpawnsOnTheSameProfileItProvisionedThe
  WorktreeForUnderARotatingPolicy. This is the ticket's own deliverable and it
  survived 186 tests in round 3.
- M3 drop the withProfile stamping in the 2-arg spawn -> KILLED by 3 tests.
- M2 make the 2-arg spawn re-enter the refusing branch -> SURVIVED, but it is a
  near-equivalent mutant, not a gap in this work. Post-#435 all four enforce*
  conditions are ones the routing branch already filtered on, so the two paths
  differ only if placement state moves between place() and spawn(). Filed
  separately.

Full build, my own run this turn, whole log redirected and grepped:
Tests run: 1572, Failures: 0, Errors: 0, Skipped: 0 - BUILD SUCCESS, 0 compile
errors. Branch already contains current main.

Read the src/main diff. The new test uses roundRobin() (stateful) and asserts
agreement between the overlay profile and the spawned profile, rather than a
hardcoded name, which is the right shape - select() is stateful, so two calls
disagree by design.
2026-09-10 09:35:02 +02:00
Dai Ha c16d118f09 fleetd #440: derive coordinator.heldDurable from queue durability + ack mode
CI / contract (pull_request) Successful in 47s
CI / build (pull_request) Successful in 1m36s
FleetMcp.coordinatorView wrote heldDurable as a literal true, so a change
that broke either the durable queue declare or the manual-ack consume in
LeadMailbox would leave the field, and the full suite, green.

- LeadChannel gets a new heldDurable() method: the conclusion of a durable
  queue declare AND a manual-ack consumer, derived by the implementation
  from what it actually did, never asserted.
- LeadMailbox.own() captures the exact booleans it passes to
  queueDeclare/basicConsume and stores their conjunction; heldDurable()
  returns it.
- FleetMcp.coordinatorView now reads channel.heldDurable() instead of a
  literal; updated the javadoc to say where the fact comes from.
- FakeLeadChannel gets a heldDurable field (default true) + withHeldDurable
  setter so FleetMcpTest can prove the field goes false.
- FleetMcpTest: new test asserts heldDurable:false when the channel says so.
- LeadMailboxTest (contract, real broker): new test asserts heldDurable()
  true against a real LeadMailbox. Verified by hand that flipping own()'s
  autoAck local to true turns this test (and two pre-existing redelivery
  tests) red, and restoring it turns them green again.
2026-09-10 14:31:44 +07:00
Dai Ha 4b10d02207 fleetd #425 rework round 4: mutation-pinning test for the dropped PlacementDecision
CI / contract (pull_request) Successful in 1m10s
CI / build (pull_request) Successful in 2m3s
SessionManager.acquireWithWorktree's unqualified branch must carry the
PlacementDecision it already resolved via launcher.place() into
launcher.spawn(spawnReq, decision) rather than re-deriving it through a
blank-profile launcher.spawn(spawnReq). Every existing test in this file uses
PlacementPolicies.fixed(), which answers select() the same way on every call,
so dropping the decision (handle = launcher.spawn(spawnReq);) was invisible:
186 tests stayed green under that mutation.

acquireWithWorktreeSpawnsOnTheSameProfileItProvisionedTheWorktreeForUnderARotatingPolicy
uses PlacementPolicies.roundRobin() instead — deterministic AND stateful, so
two select() calls on the same policy instance disagree (index 0 then index 1
across a two-profile pool). It asserts AGREEMENT between the profile the
worktree's parity overlay was provisioned for and the profile the member
actually spawned on, never a hardcoded expected profile name.

Verified as a real mutation, not a no-op: applying the exact mutation
(handle = launcher.spawn(spawnReq);) turns it red — expected [b.mcp.json]
but was [a.mcp.json] — and reverting turns it green again. Full build:
1572 tests, 0 failures, 0 errors, BUILD SUCCESS.
2026-09-10 14:24:14 +07:00
Dai Ha c5fbfdbf4a Merge main into #425 rework branch (brings #438 held-peer-mail read) 2026-09-10 14:11:37 +07:00
ltms 12cff28abb Merge #438: let a lead read its own held peer mail, primary-only (fleetd #421)
CI / contract (push) Successful in 1m15s
CI / build (push) Successful in 2m6s
fleet_poll{coordId} peeks this daemon's own held lead-to-lead mail and
returns full bodies without acking. New Authz.Action COORD_READ, primary
only — not the architect, which holds READ today. pollAction is now
argument-derived over both target and coordId.

heldView and HELD_PREVIEW_MAX_CHARS are untouched: fleet_list stays a
cheap always-safe scan, and the full read is a separately authorized call.

Verified by me, not taken from the report: merged tree builds 1562 green
(main 1555 + 7 new methods), 0 compile errors, no merge conflicts.
Mutation battery on lines the worker did NOT mutate, control 111 green:
  - COORD_READ widened to the architect  KILLED (AuthzTest + FleetMcpAuthzTest)
  - self-coord-id guard removed          KILLED (FleetMcpTest)
  - full body swapped for the preview    KILLED (FleetMcpTest)

The third mutation targets the ticket's own deliverable, and it is pinned.
Authz.permits has no default, so a new action is a compile error rather
than a silently unhandled case.

Follow-up filed as #439: fleet_list's coordinator row is still READ-gated,
so a worker sees peer coord-ids and 80-char previews of lead-to-lead
bodies. Pre-existing; the implementer flagged it and left it alone.
2026-09-10 09:08:17 +02:00
Dai Ha 77a6a7142e Merge main into #421 branch (brings #434 model-gate observability and #436 fixed-placement cap)
CI / contract (pull_request) Successful in 58s
CI / build (pull_request) Successful in 1m48s
2026-09-10 14:06:34 +07:00
Dai Ha 9f3671b801 fleetd #425 rework round 3: rewrite prose after #435 made fixed honour maxLoad
CI / contract (pull_request) Successful in 48s
CI / build (pull_request) Successful in 1m56s
fleetd #435 (merged to main) made FixedPlacementPolicy evaluate maxLoad during automatic
selection, the same way weighted/round-robin already did. Six comments across
CompositePeerLauncher.java, PeerLauncher.java, PlacementDecision.java, and SessionManager.java
justified round 2's place()/spawn(req, decision) mechanism by saying fixed "deliberately never
evaluates maxLoad" — that claim is now false, and needed restating, not just deleting.

The honest case after #435: the two-path shape (routing branch falls through an excluded
candidate; explicit-profile branch refuses on it) is still real and still deliberate — an
operator who names a profile should get a refusal, not a silent substitution. What round 1 got
wrong, and what round 2 still needs to prevent, is turning a fall-through into a refusal by
accident: resolving a name via place() and then feeding it back to spawn(SpawnRequest) as an
explicit profile. Before #435 that accident was reachable through maxLoad specifically, because
fixed never evaluated it; #435 closed that specific gap, so a PlacementDecision can no longer be
at-cap in the first place. What survives as the justification for spawn(req, decision): it never
re-evaluates a condition place() already decided, and it closes the window between that decision
and the spawn in which the underlying state could otherwise move — not a failure #435 already
prevents.

Re-measured the sibling paths this round exists to keep in agreement (one profile at maxLoad: 1,
liveCount pinned at 1, PlacementPolicies.fixed(), unqualified spawn): both the with-worktree and
without-worktree paths now throw the identical PlacementException — "worker profile 'a' is at
maxLoad (1 live >= 1 cap), and no available candidate remains" — closed upstream by #435, at
place()/select(), before either path ever reaches a spawn call. The observable asymmetry this PR
was filed to fix is gone; what remains is the structural argument above.

No behavior change: place()/PlacementDecision/spawn(req, decision) are untouched, and
FixedPlacementPolicy/PlacementPolicyUtil are taken wholesale from main's merge.
2026-09-10 14:06:10 +07:00
Dai Ha 84034b34d1 Merge remote-tracking branch 'origin/main' into worker/425-rework-placement-resolve-c58ba1-9 2026-09-10 13:56:54 +07:00
Dai Ha 1e9b2c9b7e fleetd #421: let a lead peek its own held peer mail, primary-only
CI / contract (pull_request) Successful in 1m27s
CI / build (pull_request) Successful in 1m36s
fleet_list truncated held lead-to-lead messages to an 80-char preview with
no way to read the full body, and fleet_poll{target} drained the wrong
inbox (a worker's reply queue, not the coordinator mailbox) -- it silently
returned []. fleet_ack would have destroyed the message unread.

Add a non-destructive read: fleet_poll{coordId} peeks (never acks) this
daemon's own held mail via LeadChannel.peek(). The coordId must equal the
caller's own selfCoordId -- passing a peer's id is refused with a reason,
instead of repeating the original silent-[] confusion.

This is authorization-sensitive: mapping it to the existing READ action
would let any worker read every peer lead's mail in full. READ's openness
rests on "the roster carries no secrets" (Authz.java), which does not hold
for lead-to-lead coordination bodies. Added Authz.Action.COORD_READ,
primary-only (not even the architect, which holds READ today), and made
pollAction's signature depend on both target and coordId so every call
site states explicitly what it passes.

Also fixes fleet_list's "pending: 0" trap: mailbox.pending only counts
broker-ready messages, so a healthy held mailbox reads as empty. Added
heldCount/heldDurable beside held[] so the durability fact isn't implied
only by reading the code.

Mutation-tested: pollAction's COORD_READ->READ mapping, the peek->ack
substitution, the 80-char preview cap widened to 81, and the Authz case
widened to include caller.isWorker() -- each breaks exactly its matching
test and nothing else. The first attempt at the preview-cap test used a
homogeneous "x"*200 body, which a widened cap slipped through unnoticed
(contains() found a shifted match); replaced with a sentinel character at
index 80 to actually pin the boundary.

Updates CLAUDE.md's intent->tool table for fleet_poll's new coordId
semantics, per this repo's own "prompt is part of the product" rule.
wiki/ is a submodule and not committable from a worker's worktree --
wiki-bound content is in the PR body instead.
2026-09-10 13:54:50 +07:00
ltms 5d422f85fa Merge #436: make fixed placement honour maxLoad (fleetd #435)
CI / contract (push) Successful in 48s
CI / build (push) Successful in 1m37s
fixed was the one automatic policy that ignored maxLoad, and it is the
default for an absent placement: key. It now gates on the same shared
atCap predicate weighted and round-robin use, and falls through to the
next candidate rather than refusing.

Verified here: 1555 tests green (main was 1548, +7 new methods), 0
compile errors. Mutation battery, control 113 green:
  - atCap boundary >= -> >        KILLED (15 tests, across all 3 policies)
  - drop cap check in fixed walk  KILLED (3 tests)
  - weightExcluded -> false       KILLED (2 pre-existing tests)

Neither live host is affected today: both run placement: weighted (Mac
fleetd.yaml:210, fleet01 fleetd.yaml:172), and weighted already skipped
at-cap candidates via PlacementPolicyUtil.available. This aligns fixed
with the other policies and with maxLoad's documented contract.
2026-09-10 08:54:18 +02:00
Dai Ha 6b0a99b2b7 fleetd #425 rework round 2: stop routedProfileFor's caller re-entering the throwing branch
CI / contract (pull_request) Successful in 1m24s
CI / build (pull_request) Successful in 1m34s
Round 1 closed quarantine/cool-off/model-off routing for acquireWithWorktree by resolving the
profile through routedProfileFor(role) and handing that name back to launcher.spawn(SpawnRequest)
as an EXPLICIT profile. That re-resolution has a cost the lead measured directly: naming a
profile explicitly makes CompositePeerLauncher.spawn take its THROWING branch (enforceMaxLoad
included), while the routing branch a blank spawn takes never calls enforceMaxLoad at all, and
FixedPlacementPolicy (the default) deliberately never evaluates maxLoad during automatic
selection. So an at-cap pool-first profile that placement itself would have picked for a plain
unqualified spawn could die at enforceMaxLoad one call later, purely because the worktree path's
route to the spawn passed through an explicit profile name — a new failure a worktree-less
unqualified spawn never hits.

This closes the two-path shape instead of moving it: PeerLauncher gains place(role), returning an
opaque PlacementDecision, and spawn(req, decision), which honors that decision through the SAME
routing branch a blank spawn uses — no enforce* check is newly applied. SessionManager.
acquireWithWorktree now keeps the PlacementDecision from place() and hands it to
spawn(req, decision) for an unqualified request, instead of re-resolving through an explicit
profile name. An explicitly-named profile is unaffected: it still goes through spawn(req) and its
throwing branch, exactly as before.

Also corrects the acquireWithWorktree comment's false claim that round 1 "loses nothing else" —
maxLoad was lost too, as a new hard failure, not a retry. The comment now names it explicitly.

Kept the four round-1 tests (still pass — routedProfileFor now just delegates to place()). Added
one class asserting the invariant itself: an unqualified spawn on a maxLoad-capped profile must
land the same outcome with and without a worktree, asserting on the pair rather than a hardcoded
direction, so it stays correct however fleetd #435 (not this ticket) resolves whether maxLoad
should gate an unqualified spawn at all.
2026-09-10 13:42:06 +07:00
Dai Ha b066eb1903 fleetd #425 rework: resolve acquireWithWorktree through real placement, not a blind pool-first read
CI / contract (pull_request) Successful in 1m14s
CI / build (pull_request) Successful in 2m13s
e1d7dde (PR #430) kept two good fixes and one regressed one. Kept: (1)
CompositePeerLauncher.defaultProfile() delegating to
defaultProfileFor(MemberRole.DEV) so fleet_profiles' "default" tracks a live
reload, and (2) PeerLauncher.defaultProfileFor(MemberRole). Redone:
acquireWithWorktree's profile pre-resolution.

The regression: acquireWithWorktree pre-resolved via
launcher.defaultProfileFor(memberRole), which just returns the role pool's
FIRST entry, blind to quarantine/cool-off/model-off. That name was then
passed to launcher.spawn as an EXPLICIT profile, which takes
CompositePeerLauncher.spawn's THROWING branch (enforceNotQuarantined /
enforceMaxLoad / enforceModelEnabled) instead of the ROUTING branch a blank
profile gets. So a quarantined or model-off pool-first profile turned a
routine unqualified spawn into a hard PlacementException -- undermining
fleetd #429's "the fleet keeps working when a model is turned off"
guarantee for every worktree spawn.

Fix: add PeerLauncher.routedProfileFor(MemberRole), the profile an
unqualified spawn of that role would actually be routed to right now --
same candidate list, same quarantined/coolingOff/modelOff filtering, same
PlacementPolicy spawn() itself consults. CompositePeerLauncher implements it
by extracting spawn()'s context-building into a shared private
placementContextFor(role, unreachable), so spawn() and routedProfileFor()
can never disagree about which conditions apply to which candidate.
acquireWithWorktree now calls routedProfileFor once and reuses that name for
repoRoot, parityOverlay, and the spawn -- the fleetd #425 defect (the three
disagreeing) stays fixed, now on the routed path instead of the blind one.

An explicit profile named by the caller is untouched -- it still hits the
throwing branch, which is correct for an operator override.

Trade-off carried over from e1d7dde, now precisely scoped: an unqualified
worktree spawn still loses CompositePeerLauncher's cross-candidate retry on
a live PeerUnreachableException (a transport failure at spawn time, which
placement cannot see in advance) -- but NOT the quarantine/cool-off/
model-off routing, which routedProfileFor already resolved before spawn
ever runs. Accepted: a worktree provisioned for the wrong backend is worse
than a spawn that fails cleanly and can be retried by the caller.

Tests: CompositePeerLauncherTest gains
routedProfileForSkipsAQuarantinedPoolFirstProfileUnderFixedPolicy and
...ModelOff..., both under PlacementPolicies.fixed() (the default policy,
not weighted() -- the previous round's tests all used weighted() and never
exercised FixedPlacementPolicy's own inline filter, which is exactly what
regressed). SessionManagerTest gains
acquireWithWorktreeRoutesAroundAQuarantinedPoolFirstProfile, proving
repoRoot/parityOverlay/spawn agree on the ROUTED profile, not just the
pool-reordered-by-reload one the existing #425 tests already covered.
2026-09-10 13:10:30 +07:00
Dai Ha 051d320ea0 fleetd #425: fleet_profiles' default and worktree provisioning must read live placement
fleet_profiles' "default" was CompositePeerLauncher.defaultProfile, a value
frozen at construction from cfg.effectiveDefaultProfile(). An unqualified
fleet_spawn instead resolves the dev pool live via defaultProfileFor(DEV) on
every call, so reordering fleet.developers and reloading changed where a
spawn landed without ever changing what fleet_profiles reported.

- CompositePeerLauncher.defaultProfile() now delegates to
  defaultProfileFor(MemberRole.DEV) -- the same live, reload-aware pool read
  placement already uses -- falling back to the frozen field only when no
  profiles are configured at all.
- PeerLauncher gains a default defaultProfileFor(MemberRole) method so a
  generic PeerLauncher reference can ask for a role's live default; the
  default implementation delegates to defaultProfile() for launchers with no
  pool concept of their own.
- SessionManager.acquireWithWorktree resolved a profile via
  launcher.defaultProfile() (DEV-only) to provision repoRoot/parityOverlay,
  then spawned with the original (possibly blank) profile, which re-resolves
  independently through placement -- for any non-DEV role, or across a config
  reload between the two reads, the two resolutions could disagree and
  provision a worktree for a profile the member never runs on. Fixed by
  resolving once, through defaultProfileFor(the caller's actual role), and
  reusing that same resolved name for repoRoot, parityOverlay, and the spawn
  itself. Trade-off: this path now spawns with an explicit profile rather
  than a blank one, so it loses CompositePeerLauncher's cross-candidate retry
  on PeerUnreachableException -- accepted because a worktree provisioned for
  the wrong backend is worse than a spawn that fails cleanly and can be
  retried.

Tests: CompositePeerLauncherTest (live dev-pool reorder + empty-pool
fallback), FleetProfilesLiveDefaultTest (drives FleetMcp.profilesView
directly), SessionManagerTest (worktree overlay follows a reorder, and a
non-DEV role's worktree spawn uses that role's pool, not DEV's).
2026-09-10 12:58:20 +07:00
23 changed files with 1549 additions and 64 deletions
+10 -4
View File
@@ -87,12 +87,18 @@ jobs:
apt-get update && apt-get install -y --no-install-recommends maven
mvn -version
# The `contract` profile clears the default-excludes group, so the @Tag("contract") AMQP test
# runs against the RabbitMQ service container (AMQP_URI). Pinned to the one contract test to
# avoid re-running the unit suite already covered by the `build` job.
# The `contract` profile clears the default-excludes group, so `-Dgroups=contract` runs every
# @Tag("contract") test and nothing from the unit suite the `build` job already covered — a
# tag selects the whole group, so a test added to it later runs here automatically. A prior
# version of this step pinned `-Dtest=AmqpReplyInboxContractTest` by class name instead: that
# silently excluded every other contract test (including the herdr ones) from CI, and nobody
# noticed until the herdr protocol drifted out from under a test that never ran here
# (fleetd #449). If this runner has no herdr socket, the herdr-backed tests in the group
# skip on their own `assumeTrue` and only the broker-backed ones actually run — check the
# step output rather than assuming which.
- name: Contract tests
working-directory: fleetd
run: mvn -B -Pcontract test -Dtest=AmqpReplyInboxContractTest
run: mvn -B -Pcontract test -Dgroups=contract
- name: Failing test output
if: failure()
+1
View File
@@ -138,6 +138,7 @@ the merge — and merging on a reviewer's word is delegating it by proxy.
| Message a **peer lead** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task |
| Message a **peer lead** on another daemon or host | `fleet_send{coordId: <their coord-id>, content}` — needs a `coordinator:` block; your own coord-id is in `fleet_list`. Coordination only, **never** a task |
| Answer a peer lead that messaged you | `fleet_send{coordId}` — or `{sessionId}` if they are on this host. **Not** `fleet_reply`: it has no peer route and the publish is refused |
| Read your own held lead-to-lead mail (no ack) | `fleet_poll{coordId: <your own coord-id, from fleet_list's coordinator.selfId>}` — primary-only; never acks, so `fleet_list`'s `held[]` still shows it after. `fleet_list`'s `held[]` gives only a truncated preview — this is the only way to read the full body |
| Collect a held reply | `fleet_poll{target}` · then `fleet_ack{target, msgId}` |
| Tear down a member | `fleet_stop{paneId}` |
@@ -30,6 +30,16 @@ public final class Authz {
DRAIN,
/** Read-only observation: status, roster, profiles, task polling. */
READ,
/**
* Read (never ack) this daemon's own held lead-to-lead coordination mail (fleetd #421).
*
* <p>Deliberately <strong>not</strong> folded into {@link #READ}. {@code READ}'s grant
* rests on "the roster carries no secrets" (see its case below) — a lead-to-lead body is
* not the roster; it is where leads discuss host shapes, credentials and unmerged work.
* Mapping this to {@code READ} would let any worker read every peer lead's mail in full
* and would silently falsify that comment for every other {@code READ} caller.
*/
COORD_READ,
/** Scrape the metrics endpoint. */
METRICS
}
@@ -68,6 +78,11 @@ public final class Authz {
// Observation is open to every authenticated role: a worker legitimately polls its own
// status, and the roster carries no secrets.
case READ, METRICS -> caller.isPrimary() || caller.isWorker() || caller.isArchitect();
// fleetd #421: reading held lead-to-lead mail is the primary's alone. An architect
// holds READ today (CB-548), so "not primary" must mean not-architect here too — this
// is coordination between leads, not observation of the roster.
case COORD_READ -> caller.isPrimary();
};
}
@@ -3,7 +3,7 @@ package dev.ltms.fleet.herdr;
import com.fasterxml.jackson.databind.JsonNode;
/**
* Client face onto the herdr daemon (protocol 14, herdr 0.7.0).
* Client face onto the herdr daemon (protocol 19, herdr 0.8.0).
*
* <p>This is the ONLY thing in {@code fleetd} that speaks to herdr. Every method
* maps to a herdr JSON-RPC call over its Unix domain socket. Requests are
@@ -8,7 +8,7 @@ import com.fasterxml.jackson.databind.node.ObjectNode;
import java.nio.charset.StandardCharsets;
/**
* Wire codec for herdr's newline-delimited JSON-RPC (protocol 14).
* Wire codec for herdr's newline-delimited JSON-RPC (protocol 19).
*
* <p>Split out from the socket so the framing rules — the ones that actually bit us
* during the spike (id MUST be a string; response carries {@code result} or
@@ -386,10 +386,11 @@ public final class FleetMcp {
(exchange, req) -> {
Map<String, Object> a = req.arguments();
String target = str(a, "target");
String coordId = str(a, "coordId");
// The action depends on the ARGUMENTS, not on the tool name -- see pollAction.
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_poll", a), target);
if (denied != null) return denied;
return poll(messages, str(a, "ticket"), target);
return poll(messages, leadChannel, str(a, "ticket"), target, coordId);
};
// CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox).
// Acking removes a reply from the inbox, so it is a drain, not a read.
@@ -811,29 +812,42 @@ public final class FleetMcp {
/**
* Which authorization action a {@code fleet_poll} call needs, decided by its arguments
* (fleetd #272).
* (fleetd #272, widened by fleetd #421).
*
* <p>{@code fleet_poll} is <strong>two operations behind one tool name</strong>. With {@code
* ticket} it observes an async delegation and changes nothing, which is a {@link
* <p>{@code fleet_poll} is now <strong>three operations behind one tool name</strong>. With
* {@code ticket} it observes an async delegation and changes nothing, which is a {@link
* Authz.Action#READ}. With {@code target} it calls {@link MessageService#drainReplies} on that
* session -- the replies are removed from the inbox and a second call returns nothing -- so it
* is a {@link Authz.Action#DRAIN}, the same gate {@code fleet_ack} already uses for removing a
* single message, and the same one the REST path uses at {@code FleetApp.drainReplies}.
* single message, and the same one the REST path uses at {@code FleetApp.drainReplies}. With
* {@code coordId} it reads (never acks) this daemon's own held lead-to-lead mail, which is a
* {@link Authz.Action#COORD_READ} -- <strong>not</strong> {@code READ}, even though nothing is
* consumed: {@code READ}'s grant is open to every authenticated role on the premise that the
* roster carries no secrets, and a lead-to-lead body is not the roster. Mapping a non-destructive
* peer-mail read to {@code READ} would let any worker read every peer lead's mail in full.
*
* <p>Until this method existed the handler passed a constant {@code READ} for both branches.
* {@code READ} is open to every authenticated role, so any worker could read a peer's id out of
* {@code fleet_list} and destroy the replies that peer had queued for the primary. The gate
* failed open, and it did so because the required action is a function of the arguments while
* the handler chose it before looking at them.
* <p>Before this method existed (fleetd #272) the handler passed a constant {@code READ} for
* both of the original branches. {@code READ} is open to every authenticated role, so any
* worker could read a peer's id out of {@code fleet_list} and destroy the replies that peer had
* queued for the primary. The gate failed open, and it did so because the required action is a
* function of the arguments while the handler chose it before looking at them.
*
* <p>The choice lives in this method, and not inline in the handler, so that a test can assert
* the mapping the handler actually uses. {@code FleetMcpAuthzTest} already checked every
* {@link Authz.Action} against every {@link Role} and passed throughout -- it tested the policy
* table, which was correct, while the defect was in which action the caller handed it.
*
* @param target the {@code target} argument of the call, or {@code null}/blank when absent
* <p>Checked first, and exclusively of {@code target}: a call naming {@code coordId} is reading
* a different inbox entirely (this daemon's own lead channel, never a worker's), so it takes
* priority over whatever {@code target} might also say.
*
* @param target the {@code target} argument of the call, or {@code null}/blank when absent
* @param coordId the {@code coordId} argument of the call, or {@code null}/blank when absent
*/
static Authz.Action pollAction(String target) {
static Authz.Action pollAction(String target, String coordId) {
if (!isBlank(coordId)) {
return Authz.Action.COORD_READ;
}
return isBlank(target) ? Authz.Action.READ : Authz.Action.DRAIN;
}
@@ -848,7 +862,7 @@ public final class FleetMcp {
case "fleet_reply" -> Authz.Action.REPLY;
case "fleet_ask" -> Authz.Action.ASK;
case "fleet_status", "fleet_list", "fleet_profiles", "fleet_whoami" -> Authz.Action.READ;
case "fleet_poll" -> pollAction(str(arguments, "target"));
case "fleet_poll" -> pollAction(str(arguments, "target"), str(arguments, "coordId"));
case "fleet_ack" -> Authz.Action.DRAIN;
case "fleet_spawn" -> Authz.Action.SPAWN;
case "fleet_stop" -> Authz.Action.STOP;
@@ -858,6 +872,21 @@ public final class FleetMcp {
/** {@code fleet_poll}: check an async delegation by ticket, or drain a worker's inbox by target. */
static McpSchema.CallToolResult poll(MessageService messages, String ticket, String target) {
return poll(messages, null, ticket, target, null);
}
/**
* As above, plus (fleetd #421) a held-peer-mail read when {@code coordId} is present: returns
* this daemon's own held lead-to-lead messages, in full, without acking them. Checked first and
* exclusively of {@code ticket}/{@code target} — see {@link #pollAction}'s javadoc for why this
* is a different inbox (this daemon's own {@link LeadChannel}) that authorizes differently
* ({@link Authz.Action#COORD_READ}, primary-only) from either of the original two branches.
*/
static McpSchema.CallToolResult poll(MessageService messages, LeadChannel leadChannel, String ticket,
String target, String coordId) {
if (!isBlank(coordId)) {
return pollHeldPeerMail(leadChannel, coordId);
}
if (!isBlank(target)) {
var replies = messages.drainReplies(target);
if (replies.isEmpty()) {
@@ -885,6 +914,46 @@ public final class FleetMcp {
};
}
/**
* fleetd #421: a lead's own held lead-to-lead mail, read without consuming it.
*
* <p>{@link LeadChannel#peek} is non-destructive, so calling this twice returns the same
* bodies, and {@code fleet_list}'s {@code coordinator.held[]} is unaffected — this adds a
* read, it never acks. Never let this short-circuit the delivery contract: {@code
* LeadCoordLoop}'s javadoc explains why a message must stay unacked until actually delivered,
* and that is unchanged here.
*
* <p>{@code coordId} must be THIS daemon's own coord-id ({@code fleet_list}'s
* {@code coordinator.selfId}) — there is no route here to read a PEER's outbound mail, only
* your own inbound mail. Requiring the caller to echo its own id catches the exact confusion
* that opened fleetd #421: the original failed attempt passed a PEER's id ("fleet01") as
* {@code fleet_poll}'s {@code target}, expecting to read that peer's messages, and got the same
* silent {@code []} as a genuinely empty worker inbox. This refuses the same mistake here with a
* reason, instead of a second silent wrong answer.
*/
static McpSchema.CallToolResult pollHeldPeerMail(LeadChannel leadChannel, String coordId) {
if (leadChannel == null) {
return error("lead coordination is not configured (no coordinator: block) — there is "
+ "no held peer mail to read.");
}
String selfId = leadChannel.selfCoordId();
if (!coordId.equals(selfId)) {
return error("coordId \"" + coordId + "\" is not this daemon's own coord-id (\"" + selfId
+ "\"). fleet_poll reads only YOUR OWN held mail — pass your own coordId "
+ "(fleet_list's coordinator.selfId), not a peer's.");
}
return text(json(leadChannel.peek().stream().map(FleetMcp::heldMailView).toList()));
}
/** The full body of one held lead-to-lead message — never truncated, unlike {@link #heldView}. */
private static Map<String, Object> heldMailView(LeadMessage m) {
Map<String, Object> row = new LinkedHashMap<>();
row.put("msgId", m.msgId());
row.put("from", m.from());
row.put("content", m.content());
return row;
}
/**
* {@code fleet_reply}: the worker returns its structured answer, resolving the awaiting send
* or — when no send is open — queueing the reply in the inbox for later drain (CB-307).
@@ -1323,6 +1392,16 @@ public final class FleetMcp {
* counts, it can never make {@code fleet_list} itself slow or fail. {@code held} comes from
* {@link LeadChannel#peek}, a pure in-memory read with no broker round trip, so it is never
* subject to that bound.
*
* <p><strong>fleetd #421: {@code heldCount}/{@code heldDurable} fix the "pending: 0" trap.</strong>
* {@code mailbox.pending} counts only broker-<em>ready</em> messages; a held message is already
* an unacked delivery sitting with this consumer, so the normal, healthy state of a blocked lead
* is {@code "pending": 0} next to a non-empty {@code held[]} — which invites the false reading
* "these are only in memory, a restart will lose them". {@code heldCount} is the honest second
* number beside {@code pending} ({@code held.size()}, not left for the reader to count the
* array). {@code heldDurable} comes straight from {@link LeadChannel#heldDurable}, which the
* channel implementation derives from what it actually did when it declared and consumed its own
* queue (fleetd #440) — this method never asserts the fact itself.
*/
private static Map<String, Object> coordinatorView(CoordinationSource coordination) {
LeadChannel channel = coordination.leadChannel();
@@ -1330,11 +1409,14 @@ public final class FleetMcp {
return null;
}
String selfId = channel.selfCoordId();
List<LeadMessage> held = channel.peek();
Map<String, Object> row = new LinkedHashMap<>();
row.put("selfId", selfId);
row.put("configured", true);
row.put("mailbox", mailboxView(probe(channel, selfId)));
row.put("held", channel.peek().stream().map(FleetMcp::heldView).toList());
row.put("heldCount", held.size());
row.put("heldDurable", channel.heldDurable());
row.put("held", held.stream().map(FleetMcp::heldView).toList());
row.put("peers", coordination.peers().stream().map(p -> peerView(channel, p)).toList());
return row;
}
@@ -1663,10 +1745,18 @@ public final class FleetMcp {
"Check an async delegation (a fleet_send with wait:false) by its ticket: "
+ "pending, done (with the worker's reply), or failed. When target (a worker "
+ "session id) is present instead of ticket, drain that worker's inbox of "
+ "replies delivered when no send was open.",
+ "replies delivered when no send was open. When coordId is present instead, "
+ "read (never consume) your own held lead-to-lead mail — primary-only.",
objectSchema(Map.of(
"ticket", stringProp("The ticket returned by fleet_send wait:false"),
"target", stringProp("Worker session id to drain pending replies from (optional)")),
"target", stringProp("Worker session id to drain pending replies from (optional)"),
"coordId", stringProp("Your own coord-id (fleet_list's coordinator.selfId) — "
+ "reads every message currently held[] for you in full, without "
+ "acking. Read twice, get the same bodies both times; fleet_list's "
+ "held[] still reports them afterward. Primary-only, and this can "
+ "read only YOUR OWN mailbox — there is no route to a peer's outbound "
+ "mail, so passing a peer's coordId here is refused rather than "
+ "silently returning the wrong thing (or nothing).")),
List.of()));
}
@@ -14,6 +14,7 @@ import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementCandidate;
import dev.ltms.fleet.placement.PlacementContext;
import dev.ltms.fleet.placement.PlacementDecision;
import dev.ltms.fleet.placement.PlacementException;
import dev.ltms.fleet.placement.PlacementPolicies;
import dev.ltms.fleet.placement.PlacementPolicy;
@@ -402,21 +403,10 @@ public final class CompositePeerLauncher implements PeerLauncher {
// the whole profile list. An EXPLICIT profile (above) is left alone on purpose — it is the
// operator overriding, and refusing it would break `fleet_spawn{profile:"opus"}`, which
// carries no role and so would be judged against the dev pool it was never meant for.
List<PlacementCandidate> candidates = candidates(req.role());
String roleDefault = defaultProfileFor(req.role());
Set<String> unreachable = new HashSet<>();
// CB-578 stage B: computed once up front — a quarantine's expiry cannot pass within one spawn
// call, so re-deriving it per retry would only cost work, never change the answer.
Set<String> quarantined = quarantinedProfiles(candidates);
// fleetd #201 Unit 5: a distinct set from quarantined — see PlacementContext.coolingOff.
Set<String> coolingOff = coolingOffProfiles(candidates);
// fleetd #422: read live per spawn, same as quarantined/coolingOff above — a config reload
// that flips a model's enabled state is visible to the very next unqualified spawn.
Set<String> modelOff = modelOffProfiles(candidates);
PlacementContext ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable,
quarantined, coolingOff, modelOff);
PlacementContext ctx = placementContextFor(req.role(), unreachable);
int maxAttempts = candidates.isEmpty() ? 1 : candidates.size();
int maxAttempts = ctx.candidates().isEmpty() ? 1 : ctx.candidates().size();
for (int attempt = 0; attempt < maxAttempts; attempt++) {
// Deliberately uncaught: when no candidate is left (all at cap, or all unreachable) the
// policy already throws a clear message. Catching it to rethrow a generic
@@ -443,8 +433,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
chosen.profile(), e.getMessage());
unreachable.add(chosen.profile());
// Update the context for the next selection so the policy excludes this profile.
ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable,
quarantined, coolingOff, modelOff);
ctx = placementContextFor(req.role(), unreachable);
}
}
@@ -615,8 +604,23 @@ public final class CompositePeerLauncher implements PeerLauncher {
return known.isEmpty() ? List.copyOf(configured.keySet()) : known;
}
/** The profile an unqualified spawn for {@code role} falls back to under {@code fixed} placement. */
private String defaultProfileFor(MemberRole role) {
/**
* {@inheritDoc}
*
* <p>Live: reads {@link #poolFor}, which reads {@link #profileConfigs} and {@link #fleet} fresh
* on every call, so a config reload is visible without a restart (fleetd #425) — unlike {@link
* #defaultProfile}, the field captured once at construction, which this falls back to only when
* {@link #poolFor} has nothing to offer at all (no profiles configured for this composite).
*
* <p>Exact only under the {@code fixed} placement policy — the one that reads this value
* ({@code FixedPlacementPolicy}, package-private, hence not linked) as its first, preferred
* candidate. {@code weighted}/{@code round-robin} placement can choose a different candidate
* from {@code role}'s pool even on the very first spawn; this method does not simulate that
* choice, matching what the {@code defaultProfile:}-derived reporting this replaces has always
* done.
*/
@Override
public String defaultProfileFor(MemberRole role) {
List<String> pool = poolFor(role);
return pool.isEmpty() ? defaultProfile : pool.getFirst();
}
@@ -633,6 +637,142 @@ public final class CompositePeerLauncher implements PeerLauncher {
return out;
}
/**
* Build the {@link PlacementContext} an unqualified spawn of {@code role} would be judged
* against right now — the single source both {@link #spawn} and {@link #place} read, so the two
* can never disagree about which conditions (quarantine, cool-off, model-off) apply to which
* candidate (fleetd #425 rework: round 1 duplicated this into a second, blind resolver —
* {@link #defaultProfileFor} — which is why it regressed; round 2 found that even a single
* shared resolver is not enough on its own if the CALLER re-resolves through an explicit
* profile afterwards — see {@link PlacementDecision}).
*
* @param unreachable the caller's mutable unreachable set; {@link #spawn} grows this across
* retries and rebuilds the context from it, {@link #place} passes a fresh
* empty one since it never retries
*/
private PlacementContext placementContextFor(MemberRole role, Set<String> unreachable) {
List<PlacementCandidate> candidates = candidates(role);
String roleDefault = defaultProfileFor(role);
// CB-578 stage B: computed once up front — a quarantine's expiry cannot pass within one spawn
// call, so re-deriving it per retry would only cost work, never change the answer.
Set<String> quarantined = quarantinedProfiles(candidates);
// fleetd #201 Unit 5: a distinct set from quarantined — see PlacementContext.coolingOff.
Set<String> coolingOff = coolingOffProfiles(candidates);
// fleetd #422: read live per spawn, same as quarantined/coolingOff above — a config reload
// that flips a model's enabled state is visible to the very next unqualified spawn.
Set<String> modelOff = modelOffProfiles(candidates);
return new PlacementContext(roleDefault, candidates, liveCount, unreachable,
quarantined, coolingOff, modelOff);
}
/**
* {@inheritDoc}
*
* <p>fleetd #425 rework, round 2: runs the exact same selection {@link #spawn} uses for a
* blank-profile request — {@link #placementContextFor} plus one {@link PlacementPolicy#select}
* — rather than {@link #defaultProfileFor}'s blind "pool's first entry", so a quarantined,
* cooling-off, or model-off pool-first candidate is routed around here exactly as it would be
* by a real spawn. Unlike {@link #spawn}, this never retries on {@link
* PeerUnreachableException}: there is no spawn attempt to fail, so "unreachable" never grows
* past the empty set it starts with, and a single {@link PlacementPolicy#select} call already
* reflects the live quarantine/cool-off/model-off state.
*
* <p>Deliberately does <em>not</em> apply {@link #enforceMaxLoad} (or any of the other three
* {@code enforce*} checks): those belong to {@link #spawn}'s EXPLICIT-profile branch, the
* operator-override path, and this method answers a different question — "where would an
* UNQUALIFIED spawn land". That is not the same as {@code select} ignoring these conditions —
* every condition {@code select} filters on (quarantine, cooling off, {@code maxLoad} under
* every placement policy including the default {@code fixed}, since fleetd #435, model-off,
* unreachable, weight-0) is already reflected in the {@link PlacementDecision} this method
* returns, because {@code select} walked past every excluded candidate to find it. What this
* method's caller must not do is take that resolved name and hand it back to {@link
* #spawn(SpawnRequest)} as an explicit profile: the explicit-profile branch treats the same
* exclusion conditions as a reason to REFUSE, where {@code select} had already treated them as
* a reason to fall through — round 1 of this fix did exactly that, turning a fall-through this
* method had already resolved around into a refusal one call later. Round 2 fixes that at the
* caller: {@link #spawn(SpawnRequest, PlacementDecision)} carries this exact decision to the
* spawn without re-resolving or re-checking it, through the same routing path {@code select}
* itself was consulted from.
*
* @throws PlacementException if no candidate in {@code role}'s pool is currently placeable
* (mirrors what an actual unqualified spawn would throw)
*/
@Override
public PlacementDecision place(MemberRole role) {
PlacementContext ctx = placementContextFor(role, new HashSet<>());
return new PlacementDecision(placementPolicy.get().select(ctx).profile());
}
/**
* {@inheritDoc}
*
* <p>Delegates to {@link #place}, so the two can never disagree about the answer for the same
* {@code role} at the same instant — kept as a convenience for a caller that only wants the
* resolved name (a status report, a log line), never for a caller that will act on it by
* spawning: that caller must hold the {@link PlacementDecision} itself and pass it to {@link
* #spawn(SpawnRequest, PlacementDecision)} — see {@link PlacementDecision}'s javadoc for why
* resolving here and spawning separately, with the name fed back in as an explicit profile,
* regressed fleetd #425 twice.
*/
@Override
public String routedProfileFor(MemberRole role) {
return place(role).profile();
}
/**
* {@inheritDoc}
*
* <p>Routes {@code decision.profile()} directly to its owning delegate — the identical
* {@code d.spawn(routedReq)} call {@link #spawn(SpawnRequest)}'s blank-profile branch makes for
* its first pick — WITHOUT re-running {@link #enforceNotQuarantined}, {@link
* #enforceNotCoolingOff}, {@link #enforceMaxLoad}, or {@link #enforceModelEnabled}: those are
* the EXPLICIT-profile branch's checks, and {@code decision} did not come from an operator
* naming a profile — it came from {@link #place}, which already applied whichever of these
* conditions {@link PlacementPolicy#select} actually filters on (fleetd #425 rework, round 2).
*
* <p>The two branches disagree on purpose about what an excluded profile means, and that
* disagreement is not what this method removes. The blank-profile routing branch (and
* {@link #place}) treats a quarantined/cooling-off/at-cap/model-off/unreachable/weight-0 profile
* as a reason to fall through to the next candidate; the EXPLICIT-profile branch treats naming
* that same profile as a reason to refuse outright — someone who names a profile should get a
* refusal, not a silent substitution onto a different backend. That is still correct after
* fleetd #435. What round 1 got wrong, and what this method exists to stop happening again, is
* turning a fall-through into a refusal by accident: resolving a name via {@link #place} and
* then handing that same name back to {@link #spawn(SpawnRequest)} as an explicit profile takes
* the refusing branch on a decision the routing branch had already approved by falling through
* past everything else.
*
* <p>Before fleetd #435, this exact accident was reachable through {@code maxLoad} specifically:
* {@code FixedPlacementPolicy} — the default policy — did not evaluate {@code maxLoad} at all
* for automatic selection, so {@link #place} could approve an at-cap profile that {@link
* #enforceMaxLoad} would then refuse one call later. fleetd #435 closed that: {@code
* FixedPlacementPolicy} now walks past an at-cap candidate exactly like {@code weighted}/
* {@code round-robin} already did, so {@link #place} can no longer return one, and this specific
* failure — an approved placement dying at {@code enforceMaxLoad} — cannot happen any more.
* What this method still buys, now that {@code maxLoad} can no longer cause it: it never
* re-evaluates a condition {@link #place} already decided, and it closes the window between
* that decision and the spawn in which the underlying state (another spawn landing on the same
* profile, a config reload) could otherwise move and make a stale explicit re-check wrong.
*
* <p>Deliberately does not retry on {@link PeerUnreachableException} across candidates the way
* {@link #spawn(SpawnRequest)}'s blank-profile branch does: retrying here would silently
* re-place the caller onto a different profile than the one {@code decision} named, behind the
* back of a caller that may already have provisioned something (a worktree's {@code repoRoot},
* parity overlay) specifically for that name. A caller that wants the composite's own failover
* should call {@link #spawn(SpawnRequest)} with a blank profile directly, not resolve through
* {@link #place} first. Losing that retry on a resolve-then-spawn path is an accepted, unrelated
* cost — see {@code SessionManager.acquireWithWorktree}'s own comment on it — never widened by
* this round to include {@code maxLoad}, which is what round 1 actually lost.
*/
@Override
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
HerdrPeerLauncher d = route(decision.profile());
SpawnRequest routedReq = req.withProfile(decision.profile());
PeerHandle handle = d.spawn(routedReq);
spawnedBy.put(handle.id(), d);
return handle;
}
@Override
public String effectiveCwd(SpawnRequest req) {
return route(req.profileName()).effectiveCwd(req);
@@ -780,9 +920,21 @@ public final class CompositePeerLauncher implements PeerLauncher {
return byProfile.keySet();
}
/**
* {@inheritDoc}
*
* <p>fleetd #425: reports the <em>live</em> {@code dev} pool's first entry — the same value
* {@link #defaultProfileFor} computes for {@link MemberRole#DEV} — not the {@link
* #defaultProfile} field captured at construction. An unqualified {@code fleet_spawn} defaults
* to {@code MemberRole#DEV} (see {@link dev.ltms.fleet.peer.SpawnRequest}), so "the dev pool's
* live first entry" is exactly the profile such a spawn actually lands on right now — the
* question {@code fleet_profiles}' {@code "default"} field exists to answer. The frozen field is
* a role-agnostic fallback used only when {@link #poolFor} has nothing to report at all (no
* profiles configured), which {@link #defaultProfileFor} already handles.
*/
@Override
public String defaultProfile() {
return defaultProfile;
return defaultProfileFor(MemberRole.DEV);
}
/**
@@ -42,6 +42,18 @@ public interface LeadChannel {
/** This daemon's own lead coordination id — the mailbox it owns, and the {@code from} it sends as. */
String selfCoordId();
/**
* Whether a message sitting in {@link #peek}'s held set (fetched but not yet {@link #ack}ed) is
* still safe if this daemon crashes or restarts right now — the conclusion of two independent
* facts about how this channel owns its own queue: the queue was declared <em>durable</em>, and
* the consumer that filled {@code held} uses <em>manual ack</em>, so an unacked delivery is still
* owned by the broker rather than only in this process's memory. Both must hold for {@code true};
* an implementation must derive this from what it actually did when it declared and consumed its
* queue, never return a literal — fleetd #440 found {@code FleetMcp}'s {@code heldDurable} field
* doing exactly that, unable to ever report {@code false} even after the fact stopped being true.
*/
boolean heldDurable();
/**
* A non-destructive look at {@code coordId}'s mailbox — does it exist, how many messages are
* waiting on it, and how many consumers are attached — without owning, consuming, or otherwise
@@ -86,6 +86,12 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
private final Object channelLock = new Object();
/** msgId → held delivery, for this mailbox's own queue only (there is exactly one). */
private final LinkedHashMap<String, Held> held = new LinkedHashMap<>();
/**
* fleetd #440: the answer to {@link #heldDurable()}, set once by {@link #own()} from the exact
* booleans it passed to {@code queueDeclare}/{@code basicConsume} — never a separate literal that
* could drift from what those calls actually did.
*/
private boolean heldDurable;
/** Successful broker acks on this connection, retained only to make a repeated caller ack quiet. */
private final LinkedHashMap<String, Boolean> recentlyAcked = new LinkedHashMap<>();
/** Bounds {@link #recentlyAcked}: it is only an idempotency aid, never delivery state. */
@@ -191,13 +197,22 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
/** Declare + consume this daemon's own {@code lead.<selfCoordId>.inbox}. Called once, at construction. */
private void own() throws IOException {
String queue = queueName(selfCoordId);
boolean durableQueue = true; // durable, non-exclusive, keep on idle
boolean autoAck = false; // manual ack
synchronized (channelLock) {
channel.queueDeclare(queue, true, false, false, null); // durable, non-exclusive, keep on idle
channel.basicConsume(queue, false, deliverCallback(), _ -> { }); // autoAck=false: manual ack
channel.queueDeclare(queue, durableQueue, false, false, null);
channel.basicConsume(queue, autoAck, deliverCallback(), _ -> { });
}
// fleetd #440: held mail is durable only while both hold — a durable queue AND manual ack.
this.heldDurable = durableQueue && !autoAck;
log.debug("lead mailbox owns queue {} for coord-id {}", queue, selfCoordId);
}
@Override
public boolean heldDurable() {
return heldDurable;
}
/**
* Publish {@code msg} to {@code toCoordId}'s mailbox and block until the broker's publisher
* confirm for it lands. Does <em>not</em> imply owning or consuming {@code toCoordId}'s queue.
@@ -1,5 +1,7 @@
package dev.ltms.fleet.peer;
import dev.ltms.fleet.placement.PlacementDecision;
import java.nio.file.Path;
import java.util.List;
import java.util.Set;
@@ -143,9 +145,135 @@ public interface PeerLauncher {
/**
* The profile a no-argument {@link #spawn(SpawnRequest)} uses, or {@code null} if none is configured.
*
* <p>fleetd #425: for an implementation with role pools (a no-argument spawn is read as {@link
* MemberRole#DEV}, see {@link SpawnRequest}), this must be the profile a live spawn of that role
* would actually be placed on right now, not a value captured once at startup — a caller such as
* {@code fleet_profiles} relies on this to report a live, not frozen, fact.
*/
String defaultProfile();
/**
* The profile an unqualified spawn of {@code role} would resolve to right now — the role-aware,
* live counterpart of {@link #defaultProfile()} (fleetd #425).
*
* <p>A caller that must provision something profile-specific (working directory, parity overlay
* files) <em>before</em> the actual spawn — {@code SessionManager.acquireWithWorktree} is the one
* that exists today — needs the exact profile that spawn will use, for the caller's real role,
* not a role-agnostic guess. Calling {@link #defaultProfile()} for that purpose reads {@code
* MemberRole#DEV}'s answer regardless of the caller's actual role, which is wrong for any other
* role and can provision for a profile the spawn never lands on.
*
* <p>Default implementation returns {@link #defaultProfile()}, ignoring {@code role} — the right
* answer for a launcher with no role-pool concept of its own (e.g. a single {@code
* HerdrPeerLauncher} adapter, which is never reached this way in production: {@code
* CompositePeerLauncher} always fronts it and resolves roles itself).
*/
default String defaultProfileFor(MemberRole role) {
return defaultProfile();
}
/**
* The profile an <em>unqualified</em> spawn of {@code role} would actually be routed to right
* now — the same candidate list, the same {@code quarantined}/{@code coolingOff}/{@code
* modelOff} filtering, and the same {@code PlacementPolicy} that {@link #spawn} itself
* consults for a blank-profile request (fleetd #425 rework).
*
* <p>This is <em>not</em> {@link #defaultProfileFor}: that method answers "what is first in
* {@code role}'s pool", blind to quarantine, cool-off, and the model on/off gate — the right
* answer for a role-agnostic, best-effort report ({@code fleet_profiles}' {@code "default"}
* field), but the wrong one for a caller that needs the profile a spawn will actually land on.
* A quarantined or model-off pool-first profile makes {@link #defaultProfileFor} return a name
* an unqualified spawn will never be routed to.
*
* <p>Just the resolved name, not the full {@link PlacementDecision} — a caller that only wants
* to know the answer (a status report, a log line) can call this; a caller that will later
* <em>act</em> on the answer by spawning — provisioning a worktree for a specific profile
* before the peer exists is the one that matters — must call {@link #place} and carry the
* {@link PlacementDecision} itself through to {@link #spawn(SpawnRequest, PlacementDecision)}
* instead of calling this method and feeding the string back in as an explicit profile. Doing
* that re-enters {@link #spawn(SpawnRequest)}'s explicit-profile branch, which disagrees with
* the routing branch on purpose about what an excluded profile means: the routing branch (and
* {@link #place}) falls through a quarantined/cooling-off/at-cap/model-off/unreachable/weight-0
* profile to the next candidate, while the explicit branch refuses outright — correct for an
* operator who named that profile on purpose, wrong for a name that only ever came from placement
* itself. That accidental refusal is exactly the regression fleetd #425 rework round 2 fixes:
* the default implementation below delegates to {@link #place}, so the two can never drift apart,
* but a caller that resolves through this method alone and spawns separately can still recreate
* the round-1 defect for itself. (Before fleetd #435, this accident was also reachable through
* {@code maxLoad} specifically, because {@code FixedPlacementPolicy} — the default policy — did
* not evaluate it at all for automatic selection; #435 closed that gap, so a placement decision
* can no longer be at cap in the first place. The refusal-vs-fall-through disagreement above is
* the part that was never about {@code maxLoad} and is still real.)
*
* @throws RuntimeException (implementation-specific, typically a placement exception) if no
* candidate in {@code role}'s pool is currently placeable
*/
default String routedProfileFor(MemberRole role) {
return place(role).profile();
}
/**
* Resolve, <em>without spawning</em>, the {@link PlacementDecision} an unqualified spawn of
* {@code role} would make right now — the same candidate list, the same {@code
* quarantined}/{@code coolingOff}/{@code modelOff} filtering, and the same {@code
* PlacementPolicy} {@link #spawn(SpawnRequest)}'s blank-profile branch itself consults (fleetd
* #425 rework).
*
* <p>Pair this with {@link #spawn(SpawnRequest, PlacementDecision)}, never with {@link
* #spawn(SpawnRequest)} fed the decision's profile as an explicit name — see {@link
* PlacementDecision}'s own javadoc for why that second form regressed.
*
* <p>Default implementation wraps {@link #defaultProfile()}, ignoring {@code role} and every
* placement condition — the right answer for a launcher with no pool or placement-policy
* concept of its own, matching {@link #defaultProfileFor}'s own default.
*
* @throws RuntimeException (implementation-specific, typically a placement exception) if no
* candidate in {@code role}'s pool is currently placeable
*/
default PlacementDecision place(MemberRole role) {
return new PlacementDecision(defaultProfile());
}
/**
* Spawn against an already-resolved {@link PlacementDecision} from {@link #place}, honoring it
* completely: none of the conditions {@link #place} already applied — quarantine, cooling off,
* {@code maxLoad} (evaluated by every placement policy including the default {@code fixed},
* since fleetd #435), model-off — are re-evaluated here; {@code decision} already reflects them.
* This is not skipping a check {@code place} left undone; it is not repeating one {@code place}
* already did, and not re-opening the window between that decision and this spawn in which the
* underlying state could otherwise move. This is what lets a resolve-then-spawn caller
* ({@code SessionManager.acquireWithWorktree}, which must know the profile before it can
* provision a worktree for it) and a plain blank-profile {@link #spawn(SpawnRequest)} caller
* land on the exact same outcome for the exact same placement state (fleetd #425 rework,
* round 2).
*
* <p>{@code req}'s own {@link SpawnRequest#profileName()} is ignored in favor of {@code
* decision.profile()} — the caller is expected to have built {@code req} with a blank or
* matching profile; passing a request that names a <em>different</em>, explicit profile than
* the decision it is paired with is a caller bug this method does not attempt to detect.
*
* <p>Default implementation for a launcher with no placement concept of its own: delegates to
* {@link #spawn(SpawnRequest)} with the decision's profile named explicitly — its only spawn
* contract, since there is no separate routing path to honor. This default is correct ONLY for
* a launcher that spawns a single profile of its own (e.g. {@code HerdrPeerLauncher}), where
* the explicit-profile branch it re-enters and the routing branch {@link #place} would have
* used are the same thing. <strong>A launcher that routes across more than one profile — the
* way {@code CompositePeerLauncher} routes across every configured adapter — MUST override
* this method instead of inheriting this default.</strong> Re-entering {@link
* #spawn(SpawnRequest)} re-applies that single-argument method's explicit-profile checks
* ({@code enforceNotQuarantined}, {@code enforceNotCoolingOff}, {@code enforceMaxLoad}, {@code
* enforceModelEnabled} in {@code CompositePeerLauncher}), which can refuse the very profile
* {@link #place} just chose, if the underlying placement state moved in the window between the
* {@link #place} call and this one — the exact window this method and {@link PlacementDecision}
* exist to close (fleetd #444).
*
* @throws IllegalArgumentException if the decision names an unknown profile
*/
default PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
return spawn(req.withProfile(decision.profile()));
}
/**
* Resolve the effective working directory for a spawn {@code req} without actually spawning.
* Resolution order: requestedCwd → profile cwd → callerCwd → daemon cwd.
@@ -0,0 +1,49 @@
package dev.ltms.fleet.placement;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.peer.SpawnRequest;
/**
* An already-completed placement choice — the outcome of one {@link PeerLauncher#place} call,
* carried forward so a later {@link PeerLauncher#spawn(SpawnRequest, PlacementDecision)} can honor
* it directly instead of re-resolving the profile a second time (fleetd #425 rework, round 2).
*
* <p>The problem this exists to close: a caller that must know the profile <em>before</em> it can
* spawn — {@code SessionManager.acquireWithWorktree} provisions a worktree's {@code repoRoot} and
* parity overlay for a specific profile before the peer process exists — used to resolve that name
* with {@code PeerLauncher.routedProfileFor(role)} and then hand the SAME string back to {@link
* PeerLauncher#spawn(SpawnRequest)} as an EXPLICIT profile. That re-resolution is not free: naming
* a profile explicitly makes {@code CompositePeerLauncher.spawn} take its THROWING branch
* ({@code enforceNotQuarantined}/{@code enforceNotCoolingOff}/{@code enforceMaxLoad}/{@code
* enforceModelEnabled}), while an unqualified spawn's ROUTING branch never runs those checks at
* all — it instead FALLS THROUGH to the next candidate on exactly the same conditions the throwing
* branch refuses on. That disagreement is deliberate: an operator who names a profile should get a
* refusal, not a silent substitution. The bug is turning the fall-through into a refusal by
* accident — resolving a name through the routing side and then re-entering the refusing side with
* it, for a decision the routing side had already approved by walking past everything else.
* Before fleetd #435, this accident was also reachable through {@code maxLoad} specifically: the
* default {@code fixed} placement policy did not evaluate {@code maxLoad} at all for automatic
* selection, so a profile placement itself just approved could still die at {@code enforceMaxLoad}
* one call later, purely because the caller's route to the spawn passed through an explicit
* profile name instead of the routing branch — a failure a worktree-less unqualified spawn would
* never hit. fleetd #435 closed that specific gap ({@code fixed} now evaluates {@code maxLoad}
* exactly like every other placement policy), so a {@link PlacementDecision} can no longer be
* at-cap in the first place — but the refusal-vs-fall-through disagreement above was never about
* {@code maxLoad}, and resolving a name and re-entering the refusing branch with it is still wrong
* for every OTHER condition placement filters on.
*
* <p>{@link PeerLauncher#spawn(SpawnRequest, PlacementDecision)} closes that by spawning through
* the identical code path the routing branch itself uses, keyed off the SAME decision {@link
* PeerLauncher#place} returned — no re-checking of any condition placement already evaluated. A
* resolve-then-spawn caller and a blank-profile {@link PeerLauncher#spawn(SpawnRequest)} caller can
* then never disagree about which conditions apply to the same placement state, and neither one
* re-opens the window between the placement decision and the spawn in which the underlying state
* could otherwise move.
*
* @param profile the profile this decision resolved to (may be {@code null} only when no profile is
* configured at all — the same corner case {@link PeerLauncher#defaultProfile()}
* already tolerates)
*/
public record PlacementDecision(String profile) {
}
@@ -10,6 +10,7 @@ import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.peer.PeerHandle;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.peer.SpawnRequest;
import dev.ltms.fleet.placement.PlacementDecision;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -584,8 +585,65 @@ public final class SessionManager implements TurnListener {
String ownerTerminal, WorktreeRequest wt,
String sessionName, String resumeSessionId,
MemberLifecycle.SlotReservation reservation) {
String preResolvedProfile = (profile == null || profile.isBlank())
? launcher.defaultProfile() : profile;
// fleetd #425 rework (round 2): resolved through launcher.place(memberRole) — the same
// candidate list, quarantine/cool-off/model-off filtering, and PlacementPolicy an unqualified
// spawn of this role is actually judged against right now — never launcher.defaultProfile()
// (MemberRole.DEV only, wrong for any other role) and never launcher.defaultProfileFor()
// (the role's pool FIRST entry, blind to quarantine/cool-off/model-off: a first-round fix
// used exactly this and regressed fleetd #429's "the fleet keeps working when a model is
// turned off" guarantee — a quarantined or model-off pool-first profile made this throw
// instead of routing around it, which an unqualified spawn is supposed to do). This same
// resolved decision is reused below for repoRoot, parityOverlay, AND the spawn itself so the
// worktree is always provisioned for the profile the member actually runs on — the two could
// disagree before fleetd #425: this name picked repoRoot/overlay, but the spawn below passed
// the ORIGINAL (blank) profile through to placement, which re-resolves live and can pick a
// different profile if the pool changed between the two reads, or a genuinely different one
// under weighted/round-robin placement.
//
// Round 1 of this rework fed the resolved name back into launcher.spawn(SpawnRequest) as an
// EXPLICIT profile. That was a mistake this round corrects, and the mistake is not that the
// two branches apply different checks — they are SUPPOSED to disagree: the routing branch a
// blank spawn takes treats a quarantined/cooling-off/at-cap/model-off/unreachable/weight-0
// profile as a reason to fall through to the next candidate, while CompositePeerLauncher's
// THROWING branch (enforceNotQuarantined/enforceNotCoolingOff/enforceMaxLoad/
// enforceModelEnabled) treats naming that same profile explicitly as a reason to refuse
// outright. That is correct: an operator who names a profile should get a refusal, not a
// silent substitution onto a different backend. The mistake was turning a fall-through into
// a refusal by accident — resolving a name via the routing side and then re-entering the
// refusing side with it, for a placement the routing side had already approved by walking
// past everything else.
//
// Before fleetd #435, this accident was reachable through maxLoad specifically: the default
// `fixed` placement policy did not evaluate maxLoad at all for automatic selection, so an
// at-cap pool-first profile that placement itself would have picked for a plain unqualified
// spawn could die at enforceMaxLoad one call later, purely because this method's route to
// the spawn passed through an explicit profile name — a failure a worktree-less unqualified
// spawn never hit. fleetd #435 closed that gap (`fixed` now evaluates maxLoad exactly like
// every other placement policy), so that specific failure can no longer happen — a
// PlacementDecision this method resolves can no longer be at-cap in the first place. What
// this round's fix still buys, now that maxLoad can no longer cause the accident: it keeps
// the PlacementDecision from place() and hands it to launcher.spawn(SpawnRequest,
// PlacementDecision) for an unqualified request, which spawns through the SAME routing
// branch a blank spawn uses — no enforce* check is newly applied, and the window between the
// placement decision and the spawn (in which the pool, a config reload, or another spawn
// landing on the same profile could otherwise move the state) never reopens. An
// explicitly-named profile still goes through launcher.spawn(SpawnRequest) and its throwing
// branch, unchanged — that caller asked for one profile by name and still gets everything
// enforceNotQuarantined/enforceNotCoolingOff/enforceMaxLoad/enforceModelEnabled decide about
// it, refusal included.
//
// The one cost that remains, unchanged from round 1: an unqualified worktree-provisioned
// spawn does not get CompositePeerLauncher's cross-candidate retry on a live
// PeerUnreachableException raised by the backend itself at spawn time (a transport-level
// failure placement cannot see in advance) — spawn(req, decision) commits to the one profile
// place() already chose, the same way an explicit-profile spawn commits to its one name. That
// trade is deliberate: a worktree provisioned for the wrong backend (the #425 hazard) is worse
// than a spawn that fails cleanly and can be retried by the caller. Nothing else is lost:
// maxLoad, quarantine, cool-off and model-off all behave identically whether or not a
// worktree was requested — that agreement is the invariant this rework exists to hold.
boolean unqualifiedProfile = profile == null || profile.isBlank();
PlacementDecision decision = unqualifiedProfile ? launcher.place(memberRole) : new PlacementDecision(profile);
String preResolvedProfile = decision.profile();
// CB-507: resolve through the launcher's CB-112 chain (requested → profile cwd → caller →
// daemon cwd → "."), never the raw args. A plain REST spawn supplies neither a requested
// nor a caller cwd, so taking the first non-blank of those two yielded null and put
@@ -608,7 +666,15 @@ public final class SessionManager implements TurnListener {
// copies more files into the worktree after add() returns, so sharing the group any earlier
// leaves those overlay files operator-owned and read-only for a different-uid member.
worktrees.shareWithGroup(repoRoot, path);
handle = launcher.spawn(new SpawnRequest(profile, path, callerCwd, sessionName, resumeSessionId, memberRole));
// fleetd #425: preResolvedProfile, not the original (possibly blank) profile — see the
// comment above where it is resolved. The overlay/repoRoot above and the spawn here must
// name the same profile. An unqualified request stays unqualified here and is honored via
// the PlacementDecision already captured above (spawn(req, decision) — the routing branch,
// no enforce* re-check); an explicitly-named profile still goes through the single-arg
// spawn(req) and its throwing branch, exactly as before this rework.
SpawnRequest spawnReq = new SpawnRequest(unqualifiedProfile ? null : preResolvedProfile,
path, callerCwd, sessionName, resumeSessionId, memberRole);
handle = unqualifiedProfile ? launcher.spawn(spawnReq, decision) : launcher.spawn(spawnReq);
} catch (RuntimeException e) {
log.warn("spawn failed for profile={} role={} branch={} path={}: {}",
preResolvedProfile, memberRole, branch, path, e.getMessage());
@@ -636,7 +702,7 @@ public final class SessionManager implements TurnListener {
}
throw e;
}
String resolvedProfile = resolveProfile(handle, profile);
String resolvedProfile = resolveProfile(handle, preResolvedProfile);
String cwd = launcher.effectiveCwd(new SpawnRequest(resolvedProfile, path, callerCwd));
long now = nowNanos.getAsLong();
// CB-619: see the no-worktree path above — bind before recording, and store the returned
@@ -0,0 +1,79 @@
package dev.ltms.fleet;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import dev.ltms.fleet.config.FleetConfig;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.slf4j.LoggerFactory;
import java.nio.file.Files;
import java.nio.file.Path;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* Proves {@link Fleetd#main(String[])} calls every startup report before validation aborts startup.
* The invalid non-loopback bind makes {@link FleetConfig#validateAll()} throw before {@code main}
* can open the herdr socket or bind a port. The fixture also triggers every report, so removing any
* one call from {@code main} leaves its expected log line absent.
*/
class FleetdStartupReportTest {
private static Level originalLevel;
private static ListAppender<ILoggingEvent> attach() {
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
originalLevel = logger.getLevel();
logger.setLevel(Level.INFO);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
return appender;
}
private static void detach(ListAppender<ILoggingEvent> appender) {
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
logger.detachAppender(appender);
logger.setLevel(originalLevel);
}
private static boolean contains(ListAppender<ILoggingEvent> appender, String fragment) {
return appender.list.stream()
.map(ILoggingEvent::getFormattedMessage)
.anyMatch(message -> message.contains(fragment));
}
@Test
void mainReportsEveryStartupGapBeforeValidationAborts(@TempDir Path dir) throws Exception {
Path config = dir.resolve("fleetd.yaml");
Files.writeString(config, """
bind:
host: 0.0.0.0
port: 8765
profiles:
worker:
baseUrl: https://llm.ltms.dev/v1
gitTokenEnv: GITEA_TOKEN
""");
ListAppender<ILoggingEvent> appender = attach();
try {
assertThrows(IllegalStateException.class, () -> Fleetd.main(new String[]{config.toString()}));
} finally {
detach(appender);
}
assertTrue(contains(appender, "startup git host GITEA_HOST:"),
"Fleetd.main must report the git host shape");
assertTrue(contains(appender, "member trust model: members run as the same OS user"),
"Fleetd.main must report the member trust model");
assertTrue(contains(appender, "memberCredentials: absent or empty"),
"Fleetd.main must report an absent memberCredentials policy");
assertTrue(contains(appender, "exhaustedPattern: profile(s) [worker] have no exhaustedPattern configured"),
"Fleetd.main must report profiles without exhaustedPattern");
}
}
@@ -113,6 +113,20 @@ class AuthzTest {
assertTrue(Authz.permits(WORKER_A, METRICS, null));
}
/**
* fleetd #421 — unlike READ (the case above), COORD_READ (a lead's own held peer mail) is the
* primary's alone. A worker or an architect reading it would disclose lead-to-lead
* coordination bodies, not the secret-free roster READ is open about.
*/
@Test
void coordReadIsThePrimarysAloneNotAWidenedRead() {
assertTrue(Authz.permits(PRIMARY, COORD_READ, null));
assertFalse(Authz.permits(WORKER_A, COORD_READ, null),
"a worker must not read held lead-to-lead mail");
assertFalse(Authz.permits(ARCH_DESIGN, COORD_READ, null),
"an architect holds READ today, but not-primary must mean not-architect here too");
}
@Test
void unauthenticatedIsDistinguishedFromMerelyForbidden() {
// Drives the 401-vs-403 split: a missing credential is fixable by the caller, a wrong role
@@ -17,15 +17,70 @@ import static org.junit.jupiter.api.Assumptions.assumeTrue;
* SHELL directly (never {@code claude}, so no subscription/token involvement) and always tears
* the throwaway space down.
*
* <p>The seed shell's own startup (restoring its session, printing its banner) is asynchronous
* and its length is not a fleetd contract — measured here at ~2.5s on one host (fleetd #449). A
* fixed sleep before typing raced that startup: input typed before the shell reached its prompt
* was swallowed by the shell's own startup, and the pane showed the typed line followed by the
* startup banner with no command output at all — indistinguishable, at a glance, from the env
* map never reaching the shell. So this polls for a real signal (the pane's visible text
* settling, then the expected output appearing) instead of guessing a sleep length.
*
* <p>Tagged {@code contract}; run with {@code mvn test -Pcontract}.
*/
@Tag("contract")
class AgentControlContractTest {
private static final long POLL_INTERVAL_MS = 150;
/** Bound for the seed shell to settle: observed ~2.5s three times running; this leaves headroom. */
private static final long SHELL_READY_TIMEOUT_MS = 8_000;
/** Bound for the typed command's output to appear once the shell is ready: observed ~0.2s. */
private static final long OUTPUT_TIMEOUT_MS = 5_000;
private boolean noSocket() {
return !Files.exists(UnixSocketHerdrClient.defaultSocketPath());
}
private static String readPane(UnixSocketHerdrClient herdr, String paneId) {
return herdr.call("pane.read", Map.of("pane_id", paneId, "source", "visible"))
.path("read").path("text").asText("");
}
/**
* Poll {@code pane.read} until two consecutive reads come back identical — the shell's own
* startup output (restore banner, prompt) has stopped changing — or {@code timeoutMs} elapses.
* Never asserts by itself; the caller's own assertion is what actually verifies the outcome,
* this only avoids sending input into a shell still mid-startup.
*/
private static String waitUntilSettled(UnixSocketHerdrClient herdr, String paneId, long timeoutMs)
throws InterruptedException {
long deadline = System.currentTimeMillis() + timeoutMs;
String previous = null;
while (System.currentTimeMillis() < deadline) {
Thread.sleep(POLL_INTERVAL_MS);
String current = readPane(herdr, paneId);
if (current.equals(previous) && !current.isBlank()) {
return current;
}
previous = current;
}
return previous == null ? "" : previous;
}
/** Poll {@code pane.read} until {@code needle} appears or {@code timeoutMs} elapses. */
private static String waitForText(UnixSocketHerdrClient herdr, String paneId, String needle, long timeoutMs)
throws InterruptedException {
long deadline = System.currentTimeMillis() + timeoutMs;
String last = "";
while (System.currentTimeMillis() < deadline) {
last = readPane(herdr, paneId);
if (last.contains(needle)) {
return last;
}
Thread.sleep(POLL_INTERVAL_MS);
}
return last;
}
@Test
void tabCreateInjectsEnvIntoTheSeedShell() throws Exception {
assumeTrue(!noSocket(), "no herdr socket — skipping");
@@ -36,15 +91,13 @@ class AgentControlContractTest {
Map.of("ANTHROPIC_BASE_URL", "http://gx00.gw:8000"));
try {
assertNotNull(tab.rootPaneId(), "tab.create must return the seed pane");
Thread.sleep(1000); // let the seed shell reach its prompt
waitUntilSettled(herdr, tab.rootPaneId(), SHELL_READY_TIMEOUT_MS);
herdr.call("pane.send_input", Map.of(
"pane_id", tab.rootPaneId(),
"text", "printf 'PROBE_BASE=[%s]\\n' \"$ANTHROPIC_BASE_URL\"",
"keys", List.of("enter")));
Thread.sleep(800);
String visible = herdr.call("pane.read",
Map.of("pane_id", tab.rootPaneId(), "source", "visible"))
.path("read").path("text").asText("");
String visible = waitForText(herdr, tab.rootPaneId(),
"PROBE_BASE=[http://gx00.gw:8000]", OUTPUT_TIMEOUT_MS);
assertTrue(visible.contains("PROBE_BASE=[http://gx00.gw:8000]"),
"env map must reach the seed shell; saw: " + visible);
} finally {
@@ -14,7 +14,7 @@ import static org.junit.jupiter.api.Assumptions.assumeTrue;
* Contract test against a REAL running herdr. Tagged {@code contract} so it is
* excluded from {@code mvn test}; run it with {@code mvn test -Pcontract}. It fails
* loudly if herdr drifts from the protocol {@code fleetd} was built against
* (0.7.0, protocol 14) — catching breakage that unit tests with canned frames cannot.
* (0.8.0, protocol 19) — catching breakage that unit tests with canned frames cannot.
*/
@Tag("contract")
class HerdrContractTest {
@@ -24,13 +24,13 @@ class HerdrContractTest {
}
@Test
void pingReturnsProtocol14() {
void pingReturnsProtocol19() {
assumeTrue(Files.exists(socket()), "no herdr socket at " + socket() + " — skipping");
try (UnixSocketHerdrClient herdr = UnixSocketHerdrClient.connect()) {
JsonNode pong = herdr.call("ping");
assertEquals("pong", pong.get("type").asText());
assertEquals(14, pong.get("protocol").asInt(),
"fleetd is built against herdr protocol 14");
assertEquals(19, pong.get("protocol").asInt(),
"fleetd is built against herdr protocol 19");
assertFalse(pong.get("version").asText().isBlank());
}
}
@@ -201,14 +201,34 @@ class FleetMcpAuthzTest {
*/
@Test
void pollingByTargetIsADrainAndPollingByTicketIsARead() {
assertEquals(Authz.Action.DRAIN, FleetMcp.pollAction("term_b"),
assertEquals(Authz.Action.DRAIN, FleetMcp.pollAction("term_b", null),
"poll by target removes the replies — that is a drain, not an observation");
assertEquals(Authz.Action.READ, FleetMcp.pollAction(null),
assertEquals(Authz.Action.READ, FleetMcp.pollAction(null, null),
"poll by ticket changes nothing");
assertEquals(Authz.Action.READ, FleetMcp.pollAction(" "),
assertEquals(Authz.Action.READ, FleetMcp.pollAction(" ", null),
"a blank target is an absent target");
}
/**
* fleetd #421: a coordId branch is a THIRD operation behind fleet_poll's one name, and it must
* map to {@link Authz.Action#COORD_READ} — never {@link Authz.Action#READ}, even though this
* branch also consumes nothing. READ's grant is open to every authenticated role on the premise
* that the roster carries no secrets; a lead-to-lead body is not the roster, so folding this
* branch into READ would let any worker read every peer lead's mail in full. coordId also takes
* priority over target when both happen to be present — it addresses a different inbox entirely.
*/
@Test
void pollingByCoordIdIsACoordReadNeverAPlainRead() {
assertEquals(Authz.Action.COORD_READ, FleetMcp.pollAction(null, "mac-opus"),
"reading held peer mail must not be mapped to the everyone-readable READ action");
assertEquals(Authz.Action.COORD_READ, FleetMcp.pollAction(" ", "mac-opus"),
"a blank target must not fall through to READ/DRAIN when coordId is present");
assertEquals(Authz.Action.READ, FleetMcp.pollAction(null, " "),
"a blank coordId is an absent coordId, same as target/ticket");
assertEquals(Authz.Action.COORD_READ, FleetMcp.pollAction("term_b", "mac-opus"),
"coordId takes priority over target — this is a different inbox, not a drain");
}
@Test
void everyRegisteredToolHasItsHandlerActionPinned() {
Set<String> registered = toolsTheServerRegisters();
@@ -230,6 +250,8 @@ class FleetMcpAuthzTest {
assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_whoami", Map.of()));
assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_poll", Map.of("ticket", "task")));
assertEquals(Authz.Action.DRAIN, FleetMcp.toolAction("fleet_poll", Map.of("target", "term_b")));
assertEquals(Authz.Action.COORD_READ,
FleetMcp.toolAction("fleet_poll", Map.of("coordId", "mac-opus")));
}
private static Set<String> toolsTheServerRegisters() {
@@ -249,11 +271,11 @@ class FleetMcpAuthzTest {
void aWorkerMayNotDrainAnotherSessionsInboxByPolling() {
FleetMcp m = mcp(true);
assertNotNull(m.denyFor(WORKER_A, FleetMcp.pollAction("term_b"), "term_b"),
assertNotNull(m.denyFor(WORKER_A, FleetMcp.pollAction("term_b", null), "term_b"),
"a worker draining a peer's inbox would destroy replies queued for the primary");
assertNotNull(m.denyFor(ARCH_DESIGN, FleetMcp.pollAction("term_b"), "term_b"),
assertNotNull(m.denyFor(ARCH_DESIGN, FleetMcp.pollAction("term_b", null), "term_b"),
"an architect has no lifecycle rights either — same gate as fleet_ack");
assertNull(m.denyFor(PRIMARY, FleetMcp.pollAction("term_b"), "term_b"),
assertNull(m.denyFor(PRIMARY, FleetMcp.pollAction("term_b", null), "term_b"),
"collecting a held reply is the primary's job");
}
@@ -265,11 +287,33 @@ class FleetMcpAuthzTest {
void pollingAnOwnTicketStaysOpenToWorkersAndArchitects() {
FleetMcp m = mcp(true);
assertNull(m.denyFor(WORKER_A, FleetMcp.pollAction(null), null));
assertNull(m.denyFor(ARCH_DESIGN, FleetMcp.pollAction(null), null),
assertNull(m.denyFor(WORKER_A, FleetMcp.pollAction(null, null), null));
assertNull(m.denyFor(ARCH_DESIGN, FleetMcp.pollAction(null, null), null),
"an architect delegates with wait:false, so it must be able to poll its ticket");
}
/**
* fleetd #421 acceptance: the whole ticket, pinned at the policy-table layer. A worker must be
* refused the coordId branch, the primary must be allowed, and (CB-548) an architect — which
* holds READ today — must be refused too, because "not primary" means not-architect here.
*/
@Test
void aWorkerAndAnArchitectMayNotReadHeldPeerMailOnlyThePrimaryMay() {
FleetMcp m = mcp(true);
Authz.Action coordRead = FleetMcp.pollAction(null, "mac-opus");
McpSchema.CallToolResult workerDenied = m.denyFor(WORKER_A, coordRead, null);
assertNotNull(workerDenied, "a worker must not read held lead-to-lead mail");
assertTrue(workerDenied.isError());
McpSchema.CallToolResult archDenied = m.denyFor(ARCH_DESIGN, coordRead, null);
assertNotNull(archDenied, "an architect holds READ today, but not-primary must mean "
+ "not-architect here too");
assertTrue(archDenied.isError());
assertNull(m.denyFor(PRIMARY, coordRead, null), "reading its own held mail is the primary's job");
}
// --- identity reconstruction from the transport context ------------------------------------
@Test
@@ -680,7 +680,11 @@ class FleetMcpTest {
void listReportsHeldMessagesWithATruncatedPreviewNeverTheFullBody() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
String longContent = "x".repeat(200);
// A homogeneous "x".repeat(200) body would not pin the exact 80-char boundary: a preview
// widened by one (81 chars) still CONTAINS "x".repeat(80) + "…" as a substring one position
// later, because every character is 'x'. Put a sentinel ("Y") exactly at index 80 — the
// first character a widened cap would leak — so any preview past 80 chars is caught.
String longContent = "x".repeat(80) + "Y" + "z".repeat(119);
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
.hold(new LeadMessage("m1", "fleet01-lead", "mac-opus", longContent));
@@ -694,9 +698,113 @@ class FleetMcpTest {
assertTrue(out.contains("\"msgId\":\"m1\""), out);
assertTrue(out.contains("\"from\":\"fleet01-lead\""), out);
assertFalse(out.contains(longContent), "fleet_list must never dump a held message's full body: " + out);
assertFalse(out.contains("Y"),
"the sentinel at index 80 must never appear — a preview past 80 chars leaked it: " + out);
assertTrue(out.contains("x".repeat(80) + "…"), "expected an 80-char preview with an ellipsis: " + out);
}
/**
* fleetd #421: {@code mailbox.pending} counts only broker-ready messages, so a blocked lead's
* normal, healthy state is {@code "pending": 0} next to a non-empty {@code held[]} — which
* invites the false reading that held mail is in-memory-only. {@code heldCount}/{@code
* heldDurable} put an honest second number and fact beside {@code pending} instead of leaving it
* as the only one.
*/
@Test
void listReportsAnHonestHeldCountAndDurabilityNotJustPendingZero() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
.withMailbox("mac-opus", LeadChannel.MailboxState.exists("mac-opus", 0, 1))
.hold(new LeadMessage("m1", "fleet01-lead", "mac-opus", "one"))
.hold(new LeadMessage("m2", "fleet01-lead", "mac-opus", "two"))
.hold(new LeadMessage("m3", "fleet01-lead", "mac-opus", "three"));
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), Map.of(), "",
new FleetMcp.CoordinationSource(channel, List.of()));
String out = textOf(res);
assertTrue(out.contains("\"pending\":0"), out);
assertTrue(out.contains("\"heldCount\":3"),
"the honest count beside pending: 0 — three messages really are held: " + out);
assertTrue(out.contains("\"heldDurable\":true"),
"must state the durability fact, not leave pending as the only number next to held[]: " + out);
}
/**
* fleetd #440: {@code heldDurable} must be a derived fact, not a literal — so it can report
* {@code false} when the channel behind it says held mail is not durable (a non-durable queue,
* or a consumer running with {@code autoAck=true}). A test that only ever asserts {@code true}
* repeats the defect this ticket fixes.
*/
@Test
void listReportsHeldDurableFalseWhenTheChannelSaysMailIsNotDurable() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
.withMailbox("mac-opus", LeadChannel.MailboxState.exists("mac-opus", 0, 1))
.withHeldDurable(false)
.hold(new LeadMessage("m1", "fleet01-lead", "mac-opus", "one"));
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), Map.of(), "",
new FleetMcp.CoordinationSource(channel, List.of()));
String out = textOf(res);
assertTrue(out.contains("\"heldDurable\":false"),
"heldDurable must follow the channel, not a hardcoded true: " + out);
}
// ── fleetd #421: a lead reads (never consumes) its own held peer mail ──────────────────────
@Test
void pollWithCoordIdReturnsTheFullBodyWithoutAckingAndLeavesItHeld() {
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
.hold(new LeadMessage("m1", "fleet01-lead", "mac-opus", "x".repeat(200)));
McpSchema.CallToolResult first = FleetMcp.poll(messages, channel, null, null, "mac-opus");
assertNotEquals(Boolean.TRUE, first.isError(), textOf(first));
String out1 = textOf(first);
assertTrue(out1.contains("\"msgId\":\"m1\""), out1);
assertTrue(out1.contains("\"from\":\"fleet01-lead\""), out1);
assertTrue(out1.contains("\"content\":\"" + "x".repeat(200) + "\""),
"the coordId route must return the FULL body, unlike fleet_list's preview: " + out1);
// Read again: identical bodies, and nothing was acked — peek() still holds it.
McpSchema.CallToolResult second = FleetMcp.poll(messages, channel, null, null, "mac-opus");
assertEquals(out1, textOf(second), "peek is non-destructive — reading twice must return the same bodies");
assertTrue(channel.acked().isEmpty(), "a read must never ack — that is the whole point of the ticket");
assertEquals(1, channel.peek().size(), "the message must still be held after being read");
}
@Test
void pollWithCoordIdRefusesAPeersCoordIdInsteadOfReturningTheWrongMailOrNothing() {
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
.hold(new LeadMessage("m1", "fleet01-lead", "mac-opus", "secret coordination body"));
// The ORIGINAL fleetd #421 confusion: passing a PEER's id where a self-address was meant.
McpSchema.CallToolResult res = FleetMcp.poll(messages, channel, null, null, "fleet01-lead");
assertEquals(Boolean.TRUE, res.isError());
String out = textOf(res);
assertFalse(out.contains("secret coordination body"),
"a refused read must never leak the body it refused to return: " + out);
assertTrue(out.contains("mac-opus"), "the error must name this daemon's own coordId: " + out);
}
@Test
void pollWithCoordIdErrorsHonestlyWhenLeadCoordinationIsNotConfigured() {
McpSchema.CallToolResult res = FleetMcp.poll(messages, null, null, null, "mac-opus");
assertEquals(Boolean.TRUE, res.isError());
assertTrue(textOf(res).toLowerCase().contains("coordinat"), textOf(res));
}
@Test
void listReportsEachDeclaredPeersLiveReachability() {
FakeHerdr h = new FakeHerdr();
@@ -752,6 +860,9 @@ class FleetMcpTest {
@Override
public String selfCoordId() { return "mac-opus"; }
@Override
public boolean heldDurable() { return true; }
@Override
public MailboxState inspect(String coordId) {
started.countDown();
@@ -0,0 +1,114 @@
package dev.ltms.fleet.mcp;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.member.CompositePeerLauncher;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.peer.SpawnRequest;
import dev.ltms.fleet.placement.BackendQuarantine;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import java.util.Map;
import java.util.Set;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #425: {@code fleet_profiles}' {@code "default"} field was captured once at boot
* ({@code cfg.effectiveDefaultProfile()}, frozen into {@code CompositePeerLauncher.defaultProfile}
* at construction) while an unqualified spawn resolves the same underlying key
* ({@code fleet.developers}' first entry) live, on every call. Reordering {@code fleet.developers}
* and reloading changed where a spawn landed without ever changing what {@code fleet_profiles}
* reported — a lead following {@code CLAUDE.md}'s "check {@code fleet_profiles} once per session"
* instruction was told a stale answer.
*
* <p>This test drives the exact caller {@code fleet_profiles} uses —
* {@link FleetMcp#profilesView(PeerLauncher, FleetMcp.QuarantineSource, FleetMcp.OutageSource)} —
* against a real, reloadable {@link ConfigRef}, so it fails if the reporting path is ever recoupled
* to a frozen value instead of {@link CompositePeerLauncher#defaultProfile()}'s live answer.
*/
class FleetProfilesLiveDefaultTest {
/** A minimal fleetd.yaml whose dev pool is {@code profilesInOrder}, in that definition order. */
private static String yamlWithDevPool(String... profilesInOrder) {
StringBuilder devPool = new StringBuilder();
for (int i = 0; i < profilesInOrder.length; i++) {
devPool.append(" slot").append(i).append(":\n profile: ")
.append(profilesInOrder[i]).append('\n');
}
return """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
opus:
baseUrl: http://gx00.gw:8000
model: opus-coder
sonnet:
baseUrl: http://gx00.gw:8000
model: sonnet-coder
guard:
offSubscriptionHosts:
- gx00.gw
fleet:
developers:
""" + devPool;
}
@Test
void fleetProfilesDefaultTracksALiveDevPoolReorderAfterReload(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yamlWithDevPool("opus", "sonnet"));
ConfigRef ref = new ConfigRef(f, FleetConfig.load(f));
Map<String, FleetConfig.Profile> profiles = Map.of(
"opus", new FleetConfig.Profile("opus", "http://gx00.gw:8000", "opus-coder", null,
"FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, null, null),
"sonnet", new FleetConfig.Profile("sonnet", "http://gx00.gw:8000", "sonnet-coder", null,
"FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, null, null));
FakeHerdr herdr = new FakeHerdr();
ClaudeCodeLauncher adapter = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), profiles, "opus", _ -> null);
PeerLauncher workers = new CompositePeerLauncher(
List.of(adapter), "opus", ref, _ -> 0, BackendQuarantine.none());
assertReportedDefaultMatchesAnUnqualifiedSpawn(workers, "opus");
Files.writeString(f, yamlWithDevPool("sonnet", "opus"));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied(), () -> "reload should apply cleanly: " + out.error());
assertReportedDefaultMatchesAnUnqualifiedSpawn(workers, "sonnet");
}
/**
* Asserts BOTH that {@code fleet_profiles}' {@code "default"} equals {@code expected}, AND that
* it equals what a real unqualified {@code MemberRole#DEV} spawn actually gets placed on right
* now — the two facts fleetd #425 found disagreeing.
*/
private static void assertReportedDefaultMatchesAnUnqualifiedSpawn(PeerLauncher workers, String expected) {
Map<String, Object> view = FleetMcp.profilesView(
workers, FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none());
assertEquals(expected, view.get("default"),
"fleet_profiles' \"default\" must be the live dev-pool answer, not a boot-time snapshot");
String placed = workers.spawn(
new SpawnRequest(null, null, null, null, null, MemberRole.DEV)).profile();
assertEquals(expected, placed,
"sanity: the profile an unqualified dev spawn actually lands on");
}
}
@@ -20,6 +20,7 @@ import dev.ltms.fleet.peer.PeerUnreachableException;
import dev.ltms.fleet.peer.SpawnRequest;
import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementDecision;
import dev.ltms.fleet.placement.PlacementException;
import dev.ltms.fleet.placement.PlacementPolicies;
import org.junit.jupiter.api.Test;
@@ -944,6 +945,89 @@ class CompositePeerLauncherTest {
assertTrue(e.getMessage().contains("maxLoad"), e.getMessage());
}
// ── fleetd #425: defaultProfile()/defaultProfileFor() must track a live reload ─────────────
/** A minimal fleetd.yaml whose dev pool is {@code profilesInOrder}, in that definition order. */
private static String yamlWithDevPool(String... profilesInOrder) {
StringBuilder devPool = new StringBuilder();
for (int i = 0; i < profilesInOrder.length; i++) {
devPool.append(" slot").append(i).append(":\n profile: ")
.append(profilesInOrder[i]).append('\n');
}
return """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
opus:
baseUrl: http://gx00.gw:8000
model: opus-coder
sonnet:
baseUrl: http://gx00.gw:8000
model: sonnet-coder
guard:
offSubscriptionHosts:
- gx00.gw
fleet:
developers:
""" + devPool;
}
/**
* Criterion 1 (fleetd #425): reorder {@code fleet.developers}, reload, and assert the reported
* default ({@link CompositePeerLauncher#defaultProfile()} — what {@code fleet_profiles}' {@code
* "default"} is built from, see {@code FleetMcp.profilesView}) matches what an unqualified
* {@code MemberRole#DEV} spawn is actually placed on, both before and after the reorder. Asserts
* {@code applied()} so the test proves the reload actually took, not that nothing changed.
*/
@Test
void defaultProfileTracksALiveDevPoolReorderAfterReload(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yamlWithDevPool("opus", "sonnet"));
ConfigRef ref = new ConfigRef(f, FleetConfig.load(f));
FakeHerdr herdr = new FakeHerdr();
StubLauncher adapter = new StubLauncher("claude", herdr, threeProfiles(), "opus", Set.of());
CompositePeerLauncher composite = new CompositePeerLauncher(
List.of(adapter), "opus", ref, _ -> 0, BackendQuarantine.none());
assertEquals("opus", composite.defaultProfile(),
"reported default starts at the dev pool's first entry");
assertEquals("opus", composite.spawn(
new SpawnRequest(null, null, null, null, null, MemberRole.DEV)).profile(),
"an unqualified dev spawn must land on the same profile that was just reported");
Files.writeString(f, yamlWithDevPool("sonnet", "opus"));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied(), () -> "reload should apply cleanly: " + out.error());
assertEquals("sonnet", composite.defaultProfile(),
"the reported default must follow the reorder with no daemon restart");
assertEquals("sonnet", composite.spawn(
new SpawnRequest(null, null, null, null, null, MemberRole.DEV)).profile(),
"and it must still be exactly what an unqualified spawn actually gets");
}
/**
* Criterion 2 — the mirror, and the load-bearing half (fleetd #425): with NOTHING configured (no
* profiles at all, hence an empty pool for every role), the frozen {@code defaultProfile} field
* is still what gets reported. A fix that always returns {@code poolFor(role).getFirst()} with no
* empty-pool fallback throws or returns the wrong thing here even though criterion 1 above still
* passes — this is the test that catches it.
*/
@Test
void defaultProfileFallsBackToTheFrozenFieldWhenNothingIsConfiguredAtAll() {
FakeHerdr herdr = new FakeHerdr();
StubLauncher adapter = new StubLauncher("claude", herdr, Map.of(), "opus", Set.of());
CompositePeerLauncher composite = new CompositePeerLauncher(
List.of(adapter), "opus", Map.of(), PlacementPolicies.fixed(), _ -> 0);
assertEquals("opus", composite.defaultProfile(),
"with no profiles configured at all, the frozen field is the only answer available");
assertEquals("opus", composite.defaultProfileFor(MemberRole.DEV));
}
// ── CB-578 stage B: a BACKEND_EXHAUSTED classification quarantines the credential ──────────
@Test
@@ -1010,6 +1094,98 @@ class CompositePeerLauncherTest {
assertEquals(0, adapter.spawnCount("sol"));
}
/**
* fleetd #425 rework, acceptance 1: {@link CompositePeerLauncher#routedProfileFor} must apply
* the SAME quarantine filtering {@link CompositePeerLauncher#spawn} does, under {@code fixed()}
* — the DEFAULT placement policy, deliberately not {@code weighted()} (which the regressed
* round's own tests all used, and which never exercises {@code FixedPlacementPolicy}'s own
* inline filter). This is the exact defect: the previous round's {@code defaultProfileFor}
* blindly returns the pool's first entry ("sol", quarantined here) with no awareness of
* quarantine at all, which is what turned a routine unqualified spawn into a hard throw once
* {@code acquireWithWorktree} pre-resolved through it.
*/
@Test
void routedProfileForSkipsAQuarantinedPoolFirstProfileUnderFixedPolicy() {
FakeHerdr herdr = new FakeHerdr();
Map<String, FleetConfig.Profile> profiles = ordered(
"sol", stubWorker("sol", "shared-openai"),
"b", stubWorker("b"));
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "sol", Set.of());
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
quarantine.quarantine("shared-openai");
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
PlacementPolicies.fixed(), _ -> 0, null, quarantine);
assertEquals("b", composite.routedProfileFor(MemberRole.DEV),
"sol (the pool's first entry) is quarantined, so the routed answer must be b");
assertEquals("sol", composite.defaultProfileFor(MemberRole.DEV),
"sanity: defaultProfileFor stays blind to quarantine — that's the gap routedProfileFor closes");
assertEquals(0, adapter.spawnCount("sol"), "routedProfileFor never spawns anything");
assertEquals(0, adapter.spawnCount("b"), "routedProfileFor never spawns anything");
}
/**
* fleetd #444: {@link PlacementDecision} exists to close the window between {@link
* CompositePeerLauncher#place} and {@link CompositePeerLauncher#spawn(SpawnRequest,
* PlacementDecision)} — the placement state must be free to move in that window without the
* held decision being re-checked against the new state. Every quarantine test above resolves
* and spawns in one call, so none of them ever open that window; this test is the one that
* does: "sol" is placed FIRST, while nothing is quarantined yet, and only THEN is its
* credential quarantined, before the held decision is spawned.
*
* <p>This is the test that tells the real override apart from the alternative body the ticket
* measured: routing {@code decision.profile()} straight to its adapter (the real override)
* never re-runs {@code enforceNotQuarantined}, so the spawn against the held decision still
* succeeds on sol. Re-entering {@code spawn(req.withProfile(decision.profile()))} instead
* lands in the explicit-profile branch, which refuses a now-quarantined sol outright — before
* this test existed, replacing the real override's body with that re-entering call left the
* whole suite green.
*/
@Test
void spawnHonorsAPlacementDecisionEvenAfterItsProfileIsQuarantinedInTheWindowAfterPlace() {
FakeHerdr herdr = new FakeHerdr();
Map<String, FleetConfig.Profile> profiles = ordered(
"sol", stubWorker("sol", "shared-openai"),
"b", stubWorker("b"));
// The adapter's OWN fallback default is "b", deliberately different from the profile place()
// decides ("sol") — see the note below on why this must not be "sol" too.
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "b", Set.of());
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
PlacementPolicies.fixed(), _ -> 0, null, quarantine);
// 1. Resolve BEFORE anything is quarantined — sol (definition order first, fixed policy) wins.
// composite's own defaultProfile ("sol", the constructor arg above) never enters this: the
// pool poolFor(DEV) resolves to is never empty here, so place() only ever reads that field as
// a fallback for an empty pool, which this test does not exercise.
PlacementDecision decision = composite.place(MemberRole.DEV);
assertEquals("sol", decision.profile(), "sanity: nothing is quarantined yet, so sol is placed");
// 2. Move the placement state IN THE WINDOW between place() and spawn() — sol's credential
// is now quarantined. A fresh place()/spawn(req) pair would fall through to b instead; the
// held decision must not be re-evaluated against this new state at all.
quarantine.quarantine("shared-openai");
// 3. Spawn against the HELD decision, not a fresh resolve.
SpawnRequest req = new SpawnRequest(null, null, null, null, null, MemberRole.DEV);
PeerHandle handle = composite.spawn(req, decision);
assertEquals("sol", handle.profile(),
"the decision from place() is honored even though sol is now quarantined");
// A fixture whose adapter falls back to "sol" too would let an UNSTAMPED request (one
// routed but never given req.withProfile("sol")) land on spawnCount("sol") == 1 by
// COINCIDENCE, since StubLauncher.spawn falls back to its own defaultProfile whenever
// req.profileName() is blank. Giving the adapter "b" as its fallback instead means only an
// actually-stamped request can produce this count — an unstamped one would count against
// "b" and this assertion would fail.
assertEquals(1, adapter.spawnCount("sol"),
"the request that reached the delegate actually carried sol as its profile "
+ "(the adapter's own fallback default is 'b', so this can't happen by accident)");
assertEquals(0, adapter.spawnCount("b"),
"b must never be touched — neither as the decision's profile nor as an unstamped "
+ "request's accidental fallback");
}
@Test
void aQuarantineLiftsOnTheInjectedClockAndTheProfileBecomesSpawnableAgain() {
FakeHerdr herdr = new FakeHerdr();
@@ -1403,6 +1579,37 @@ class CompositePeerLauncherTest {
assertEquals(0, adapter.spawnCount("local"));
}
/**
* fleetd #425 rework, acceptance 2: same shape as {@link #fixedPlacementSkipsAnOffModelProfileToo}
* above, but through {@link CompositePeerLauncher#routedProfileFor} rather than an actual
* {@link CompositePeerLauncher#spawn} — the exact call {@code SessionManager.acquireWithWorktree}
* makes to pre-resolve a profile for provisioning. This is the fleetd #429 case named in the
* ticket: an operator turns a model off, and an unqualified worktree spawn must still route
* around it instead of throwing "names model, which the operator has turned off" — the throw
* {@link CompositePeerLauncher#enforceModelEnabled} raises only on the EXPLICIT-profile branch,
* which is exactly the branch the regressed round accidentally routed every worktree spawn onto.
*/
@Test
void routedProfileForSkipsAModelOffPoolFirstProfileUnderFixedPolicy() {
FakeHerdr herdr = new FakeHerdr();
Map<String, FleetConfig.Profile> profiles = ordered(
"local", stubWorkerModel("local", "deepseek-v4-flash"),
"sonnet", stubWorkerModel("sonnet", "claude-sonnet-5"));
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "local", Set.of());
FleetConfig.Models models = new FleetConfig.Models(List.of(
new FleetConfig.Models.ModelEntry("deepseek-v4-flash", false)));
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "local", profiles,
PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(), NO_OUTAGE,
() -> models);
assertEquals("sonnet", composite.routedProfileFor(MemberRole.DEV),
"local (the pool's first entry) names an off model, so the routed answer must be sonnet");
assertEquals("local", composite.defaultProfileFor(MemberRole.DEV),
"sanity: defaultProfileFor stays blind to model-off — that's the gap routedProfileFor closes");
assertEquals(0, adapter.spawnCount("local"), "routedProfileFor never spawns anything");
assertEquals(0, adapter.spawnCount("sonnet"), "routedProfileFor never spawns anything");
}
/**
* Criterion 4: turning a model off/on is HOT — no restart — proven through a REAL
* {@code ConfigRef.reload()}, not a hand-rolled supplier swap. Also proves {@code models} is
@@ -28,11 +28,19 @@ public final class FakeLeadChannel implements LeadChannel {
private volatile IllegalStateException publishFailure;
/** Canned {@link #inspect} results by coord-id — absent for any coord-id not configured here. */
private final Map<String, MailboxState> mailboxes = new ConcurrentHashMap<>();
/** fleetd #440: matches {@link LeadMailbox}'s real default (durable queue + manual ack) unless overridden. */
private volatile boolean heldDurable = true;
public FakeLeadChannel(String selfCoordId) {
this.selfCoordId = selfCoordId;
}
/** Make {@link #heldDurable()} report {@code durable} — the fleetd #440 seam for the false case. */
public FakeLeadChannel withHeldDurable(boolean durable) {
this.heldDurable = durable;
return this;
}
/** Make {@link #inspect(String)} return {@code state} for {@code coordId} instead of "absent". */
public FakeLeadChannel withMailbox(String coordId, MailboxState state) {
mailboxes.put(coordId, state);
@@ -80,6 +88,11 @@ public final class FakeLeadChannel implements LeadChannel {
return mailboxes.getOrDefault(coordId, MailboxState.absent(coordId));
}
@Override
public boolean heldDurable() {
return heldDurable;
}
public List<LeadMessage> published() {
return List.copyOf(published);
}
@@ -207,6 +207,21 @@ class LeadMailboxTest {
}
}
/**
* fleetd #440: {@code heldDurable()} must be derived from what {@link LeadMailbox#own} actually
* did against the real broker — a durable queue declare plus a manual-ack consumer — not a
* hardcoded literal. This is the mutation-sensitive test: flip {@code own()}'s {@code autoAck}
* local to {@code true} (or its {@code durableQueue} local to {@code false}) and this must fail.
*/
@Test
void heldDurableReportsTrueBecauseTheQueueIsDurableAndTheConsumeIsManualAck() throws Exception {
String self = coordId("lead-held-durable");
try (LeadMailbox mailbox = LeadMailbox.open(uri(), self)) {
assertTrue(mailbox.heldDurable(),
"own() declares a durable queue and consumes with autoAck=false, so held mail is durable");
}
}
@Test
void inspectReportsAMissingMailboxAsAbsentRatherThanThrowing() throws Exception {
String nobody = coordId("lead-inspect-nobody");
@@ -6,12 +6,14 @@ import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import dev.ltms.fleet.auth.MemberRegistry;
import dev.ltms.fleet.auth.MemberLifecycle;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.member.CompositePeerLauncher;
import dev.ltms.fleet.msg.TestTurnTokens;
import dev.ltms.fleet.peer.Capability;
import dev.ltms.fleet.peer.CharterReceipt;
@@ -20,9 +22,15 @@ import dev.ltms.fleet.peer.PeerHandle;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.peer.PeerUnreachableException;
import dev.ltms.fleet.peer.SpawnRequest;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementPolicies;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.slf4j.LoggerFactory;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
@@ -33,6 +41,7 @@ import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.LongSupplier;
import java.util.function.Supplier;
import static org.junit.jupiter.api.Assertions.*;
@@ -1991,4 +2000,296 @@ class SessionManagerTest {
sessions.rosterResolved();
assertEquals(2, handle.callCount(), "once resolved, the id must not be looked up again");
}
// ── fleetd #425 criterion 3: acquireWithWorktree must provision for the profile it actually
// spawns, never a name resolved before a live pool change is accounted for ────────────────────
/** Two profiles with distinct {@code cwd}/{@code parityOverlay}, and a dev pool of {@code first,second}. */
private static String worktreeReorderYaml(String first, String second) {
return """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
a:
baseUrl: http://gx00.gw:8000
model: coder-a
b:
baseUrl: http://gx00.gw:8000
model: coder-b
guard:
offSubscriptionHosts:
- gx00.gw
fleet:
developers:
slot0:
profile: %s
slot1:
profile: %s
""".formatted(first, second);
}
@Test
void acquireWithWorktreeProvisionsTheOverlayForTheProfileActuallySpawned(
@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, worktreeReorderYaml("a", "b"));
ConfigRef ref = new ConfigRef(f, FleetConfig.load(f));
Map<String, FleetConfig.Profile> profiles = Map.of(
"a", new FleetConfig.Profile("a", "http://gx00.gw:8000", "coder-a", null,
"FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, "/repo/a", List.of("a.mcp.json")),
"b", new FleetConfig.Profile("b", "http://gx00.gw:8000", "coder-b", null,
"FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, "/repo/b", List.of("b.mcp.json")));
FakeHerdr herdr = new FakeHerdr();
ClaudeCodeLauncher adapter = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), profiles, "a", _ -> null);
PeerLauncher launcher = new CompositePeerLauncher(
List.of(adapter), "a", ref, _ -> 0, BackendQuarantine.none());
// The pool changes AFTER the composite/launcher is built, and BEFORE the unqualified
// worktree spawn — exactly the fleetd #425 scenario: the live pool's first entry is "b" by
// the time acquireWithWorktree runs, even though nothing here was rebuilt.
Files.writeString(f, worktreeReorderYaml("b", "a"));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied(), () -> "reload should apply cleanly: " + out.error());
FakeWorktrees worktrees = new FakeWorktrees();
SessionManager sessions = new SessionManager(launcher, worktrees, () -> 0L);
MemberSession s = sessions.acquire(null, null, "/caller",
null, new WorktreeRequest("fleetd-425", null));
assertEquals("b", s.profile(),
"the live dev pool now starts at b, so the unqualified spawn must land there");
FakeWorktrees.OverlayCall overlay = worktrees.lastOverlay();
assertNotNull(overlay, "overlayParity must have been called");
assertEquals(List.of("b.mcp.json"), overlay.requested(),
"the worktree must be provisioned with profile b's overlay — the one actually "
+ "spawned — never a's, the pool's stale first entry");
}
/**
* The deterministic, mutation-pinning half of criterion 3: {@code launcher.defaultProfile()}
* only ever answers for {@link MemberRole#DEV} (see {@link CompositePeerLauncher#defaultProfile()}),
* so resolving a worktree spawn's profile through it — instead of through {@link
* PeerLauncher#defaultProfileFor(MemberRole)}, resolved against the CALLER's actual role — picks
* the wrong pool's answer for any role other than DEV. No reload or race is needed to see it: an
* ARCHITECT pool and a DEV pool that simply disagree, held constant, are enough.
*/
@Test
void acquireWithWorktreeForANonDevRoleUsesThatRolesPoolNotTheDevPool() {
Map<String, FleetConfig.Profile> profiles = Map.of(
"a", new FleetConfig.Profile("a", "http://gx00.gw:8000", "coder-a", null,
"FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, "/repo/a", List.of("a.mcp.json")),
"b", new FleetConfig.Profile("b", "http://gx00.gw:8000", "coder-b", null,
"FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, "/repo/b", List.of("b.mcp.json")));
FakeHerdr herdr = new FakeHerdr();
ClaudeCodeLauncher adapter = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), profiles, "a", _ -> null);
// developers -> a (first/only entry); architects -> b (first/only entry). The two pools
// disagree on purpose, so a role-blind resolution (DEV's answer, "a") is visibly wrong for
// an ARCHITECT spawn, which must land on "b".
FleetConfig.Fleet fleet = new FleetConfig.Fleet(Map.of(),
Map.of("s0", new FleetConfig.Slot("b")),
Map.of("s0", new FleetConfig.Slot("a")),
Map.of(), null);
PeerLauncher launcher = new CompositePeerLauncher(List.of(adapter), "a", profiles,
PlacementPolicies.fixed(), _ -> 0, fleet);
FakeWorktrees worktrees = new FakeWorktrees();
SessionManager sessions = new SessionManager(launcher, worktrees, () -> 0L);
MemberSession s = sessions.acquire(null, MemberRole.ARCHITECT, null, "/caller",
null, new WorktreeRequest("fleetd-425b", null));
assertEquals("b", s.profile(),
"an unqualified ARCHITECT worktree spawn must land on the architect pool's profile");
FakeWorktrees.OverlayCall overlay = worktrees.lastOverlay();
assertNotNull(overlay, "overlayParity must have been called");
assertEquals(List.of("b.mcp.json"), overlay.requested(),
"the worktree must be provisioned with profile b's overlay — the ARCHITECT pool's "
+ "answer, the one actually spawned — never a's, the DEV pool's answer that "
+ "launcher.defaultProfile() alone would have given");
}
/**
* fleetd #425 rework, acceptance 3: repoRoot, parityOverlay, AND the actual spawn must all name
* the SAME routed profile, proven on the ROUTED path — a quarantine skips the pool's first entry
* — not just the "pool reordered by a live reload" path the two tests above already cover.
*
* <p>This is the exact regression the rework fixes: the first round resolved
* {@code acquireWithWorktree}'s profile through {@code launcher.defaultProfileFor(memberRole)},
* which is blind to quarantine and just returns the pool's first entry ("a" here, quarantined).
* That name went on to provision repoRoot/overlay for "a", and then the spawn itself — now an
* EXPLICIT-profile spawn naming "a" — hit {@code CompositePeerLauncher.enforceNotQuarantined}
* and threw, where the pre-fix code (a blank-profile spawn) would have routed around "a" onto
* "b" without any trouble. {@code launcher.routedProfileFor(memberRole)} closes that gap by
* running the SAME quarantine-aware selection {@code spawn} itself uses, so all three — repoRoot,
* overlay, and the spawn — land on "b" together.
*/
@Test
void acquireWithWorktreeRoutesAroundAQuarantinedPoolFirstProfile() {
Map<String, FleetConfig.Profile> profiles = new LinkedHashMap<>();
profiles.put("a", new FleetConfig.Profile("a", "http://gx00.gw:8000", "coder-a", null,
"FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, "/repo/a", List.of("a.mcp.json"),
null, null, null, null, null, null, null, null, "shared-cred", null));
profiles.put("b", new FleetConfig.Profile("b", "http://gx00.gw:8000", "coder-b", null,
"FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, "/repo/b", List.of("b.mcp.json")));
FakeHerdr herdr = new FakeHerdr();
ClaudeCodeLauncher adapter = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), profiles, "a", _ -> null);
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
quarantine.quarantine("shared-cred");
PeerLauncher launcher = new CompositePeerLauncher(List.of(adapter), "a", profiles,
PlacementPolicies.fixed(), _ -> 0, null, quarantine);
FakeWorktrees worktrees = new FakeWorktrees();
SessionManager sessions = new SessionManager(launcher, worktrees, () -> 0L);
MemberSession s = sessions.acquire(null, null, "/caller",
null, new WorktreeRequest("fleetd-425c", null));
assertEquals("b", s.profile(),
"a is quarantined, so the unqualified worktree spawn must route to b");
FakeWorktrees.OverlayCall overlay = worktrees.lastOverlay();
assertNotNull(overlay, "overlayParity must have been called");
assertEquals(List.of("b.mcp.json"), overlay.requested(),
"parityOverlay must be provisioned for b — the profile actually spawned, never a's, "
+ "the quarantined pool-first entry");
FakeWorktrees.RepoRootCall repoRootCall = worktrees.repoRootCalls().getLast();
assertTrue(repoRootCall.cwd().contains("/repo/b"),
"repoRoot must be resolved through b's effectiveCwd, not a's: " + repoRootCall.cwd());
}
/**
* fleetd #425 rework, round 2: this is the exact probe that found round 1's maxLoad
* regression. One dev profile ("a") is configured with {@code maxLoad: 1} and a liveCount
* pinned at 1 — permanently at cap — under {@code PlacementPolicies.fixed()}, the default
* policy, which deliberately never evaluates {@code maxLoad} during automatic selection (see
* {@code CompositePeerLauncher}'s own javadoc on {@code place}/{@code FixedPlacementPolicy}).
*
* <p>Round 1 resolved {@code acquireWithWorktree}'s profile through
* {@code launcher.routedProfileFor(memberRole)} and then fed that name back into
* {@code launcher.spawn(SpawnRequest)} as an EXPLICIT profile. Naming a profile explicitly
* takes {@code CompositePeerLauncher.spawn}'s THROWING branch, which calls
* {@code enforceMaxLoad} — so the worktree path died with a {@code PlacementException} while
* the exact same unqualified request, with no worktree, still spawned cleanly through the
* routing branch that never checks {@code maxLoad} at all. One intent, two different answers,
* depending only on whether a worktree was asked for — the #425 shape, moved to a different
* filter instead of closed.
*
* <p>This test does not hardcode which of the two outcomes is correct — whether an unqualified
* spawn SHOULD respect {@code maxLoad} is fleetd #435, a separate ticket. It only asserts that
* the WITH-worktree and WITHOUT-worktree paths agree: both spawn on the same profile, or both
* fail with the same exception type and message. That way this test stays correct however
* #435 is eventually resolved, and only breaks if the two paths disagree again.
*/
@Test
void unqualifiedAcquireAgreesWithAndWithoutAWorktreeWhenTheOnlyProfileIsAtMaxLoad() {
Map<String, FleetConfig.Profile> profiles = new LinkedHashMap<>();
profiles.put("a", new FleetConfig.Profile("a", "http://gx00.gw:8000", "coder-a", null,
"FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, "/repo/a", List.of("a.mcp.json"),
null, null, null, null, 1.0f, 1));
FakeHerdr herdr = new FakeHerdr();
ClaudeCodeLauncher adapter = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), profiles, "a", _ -> null);
// liveCount pinned at 1 for "a", exactly matching maxLoad — "a" is permanently at cap,
// regardless of how many times either branch below actually spawns.
PeerLauncher launcher = new CompositePeerLauncher(List.of(adapter), "a", profiles,
PlacementPolicies.fixed(), name -> "a".equals(name) ? 1 : 0);
Object without = attemptAcquire(() ->
new SessionManager(launcher, new FakeWorktrees(), () -> 0L)
.acquire(null, null, "/caller", null));
Object with = attemptAcquire(() ->
new SessionManager(launcher, new FakeWorktrees(), () -> 0L)
.acquire(null, null, "/caller", null, new WorktreeRequest("fleetd-425-maxload", null)));
assertEquals(without, with, "an unqualified spawn on a profile at maxLoad must agree "
+ "whether or not a worktree was requested — no condition may become newly fatal "
+ "on the worktree path alone (fleetd #425 rework, round 2)");
}
/**
* fleetd #425 rework, round 4: the exact regression a mutation test found that 186 green tests
* missed — {@code acquireWithWorktree} dropping the {@link
* dev.ltms.fleet.placement.PlacementDecision} it already resolved via {@code launcher.place},
* and letting the unqualified spawn re-run placement a second time (a blank-profile {@code
* launcher.spawn(spawnReq)}) instead of carrying that decision forward via {@code
* launcher.spawn(spawnReq, decision)}. Every earlier test in this file uses {@code
* PlacementPolicies.fixed()}, which returns the same answer on every {@code select()} call, so
* dropping the decision is invisible under it — two {@code select()} calls simply agree by
* accident. {@code PlacementPolicies.roundRobin()} is deterministic AND stateful: its {@code
* select()} advances an internal index on every call, so two consecutive calls for the SAME
* spawn (one from {@code place()} to provision the worktree, a second from a dropped-decision
* blank-profile {@code spawn(spawnReq)}) land on DIFFERENT profiles from a two-profile pool —
* index 0 ("a"), then index 1 ("b").
*
* <p>This test does not hardcode which profile wins — asserting one specific name would pass
* for the wrong reason the moment the rotation order changes (round-4 brief invariant 3). It
* asserts AGREEMENT instead: whichever profile the worktree's parity overlay was provisioned
* for must be the SAME profile the member actually spawned on. Each profile's overlay list is
* named after the profile itself ({@code "a.mcp.json"}/{@code "b.mcp.json"}), so comparing the
* recorded overlay against {@code s.profile() + ".mcp.json"} checks agreement without ever
* naming an expected winner.
*/
@Test
void acquireWithWorktreeSpawnsOnTheSameProfileItProvisionedTheWorktreeForUnderARotatingPolicy() {
Map<String, FleetConfig.Profile> profiles = new LinkedHashMap<>();
profiles.put("a", new FleetConfig.Profile("a", "http://gx00.gw:8000", "coder-a", null,
"FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, "/repo/a", List.of("a.mcp.json")));
profiles.put("b", new FleetConfig.Profile("b", "http://gx00.gw:8000", "coder-b", null,
"FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, "/repo/b", List.of("b.mcp.json")));
FakeHerdr herdr = new FakeHerdr();
ClaudeCodeLauncher adapter = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), profiles, "a", _ -> null);
// roundRobin is deterministic AND stateful: the first select() call picks index 0 ("a"),
// and the SAME policy instance's second select() call (reached only if the
// PlacementDecision is dropped) picks index 1 ("b") — the two-call disagreement this test
// needs to make a dropped decision observable, rather than merely probable.
PeerLauncher launcher = new CompositePeerLauncher(List.of(adapter), "a", profiles,
PlacementPolicies.roundRobin(), _ -> 0);
FakeWorktrees worktrees = new FakeWorktrees();
SessionManager sessions = new SessionManager(launcher, worktrees, () -> 0L);
MemberSession s = sessions.acquire(null, null, "/caller",
null, new WorktreeRequest("fleetd-425-round4", null));
FakeWorktrees.OverlayCall overlay = worktrees.lastOverlay();
assertNotNull(overlay, "overlayParity must have been called");
assertEquals(List.of(s.profile() + ".mcp.json"), overlay.requested(),
"the worktree must be provisioned for the SAME profile the member actually spawned "
+ "on — under a rotating policy, dropping the PlacementDecision makes the "
+ "second, spawn-time select() call disagree with the first, place()-time "
+ "call, so the member ends up on a profile whose worktree (repoRoot/parity "
+ "overlay) was built for a DIFFERENT profile (fleetd #425 rework, round 4)");
}
/**
* Reduce one {@code acquire(...)} attempt to a value comparable across the with-worktree and
* without-worktree paths: the spawned profile name on success, or the thrown exception's class
* and message on failure. Comparing THIS — instead of asserting "both spawn" or "both throw" as
* a hardcoded direction — is what keeps {@link
* #unqualifiedAcquireAgreesWithAndWithoutAWorktreeWhenTheOnlyProfileIsAtMaxLoad} valid whichever
* way fleetd #435 eventually resolves whether an unqualified spawn should respect maxLoad.
*/
private static Object attemptAcquire(Supplier<MemberSession> call) {
try {
return "spawned:" + call.get().profile();
} catch (RuntimeException e) {
return "threw:" + e.getClass().getName() + ":" + e.getMessage();
}
}
}