Compare commits

..

18 Commits

Author SHA1 Message Date
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
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 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 ed2027b202 fleetd #435: make FixedPlacementPolicy honor maxLoad
CI / contract (pull_request) Successful in 45s
CI / build (pull_request) Successful in 1m54s
FixedPlacementPolicy (the default placement policy) never consulted
maxLoad, so an at-cap default was chosen anyway on every unqualified
spawn -- the cap was advisory, not enforced, for the one policy every
config uses by default. weighted/round-robin already gated on it via
PlacementPolicyUtil.available().

Extract the "at cap" predicate into PlacementPolicyUtil.atCap(ctx, c)
so all three policies share one definition, and consult it at both of
FixedPlacementPolicy's filter sites (the default fast path and the
candidate walk), mirroring the existing weightExcluded pattern. An
at-cap default now falls through to the next candidate instead of
refusing the spawn -- only when every candidate is unusable does the
policy still throw, naming the cap in the message. Update the class
javadoc (five exceptions -> six) and the reason-priority comments to
match CompositePeerLauncher's explicit-spawn order (quarantine,
cooling off, max load, model-off).
2026-09-10 13:46:06 +07:00
ltms 6b7caba248 Merge #434: make the model gate's own state observable
CI / contract (push) Successful in 47s
CI / build (push) Successful in 1m49s
Verified in the worker's tree at 7fd914d: 1544 tests green (base 1535 + 9),
0 compile errors, unpiped mvn clean install. merge-tree against f8b0d42 reports
no conflicts and the two sides share no files.

Read all four production diffs. The design is right: one ModelGateState record
carrying both "armed" and "off" from a single models0() read, so the startup
log line, fleet_profiles' modelGateArmed and the spawn gate cannot disagree —
the fleetd #404 lesson applied properly. disabledModels() now delegates to it
rather than being a second independent read.

Mutated the two subtlest lines, with a control in the same script and the
changed line echoed back with its number:

- modelGateState(): sentinel identity check replaced by the naive
  "m.offIds().isEmpty()" -> 2 failures. FleetProfilesModelGateStateTest
  .modelsBlockWithNothingOffReportsGateArmedAndZeroOff:86 and
  CompositePeerLauncherTest.modelGateStateIsHotReloadedThroughARealConfigRef
  :1519, both "expected: <true> but was: <false>". So the one state this
  ticket exists to expose — a block present with nothing off — is pinned.
- notConfigured() returning a non-empty off set, breaking the invariant the
  record's javadoc states but does not enforce -> 3 failures, including
  disabledModelsIsEmptyWithNoModelsConfigured:1581. So the invariant is
  observable even though the constructor does not check it.

Control run unmutated: 1544 green.

Follow-up on me, not a merge blocker: fleet_profiles gains an operator-visible
field, so this needs a wiki/11-Features.md entry. Workers cannot commit the
wiki submodule, so I am adding it.
2026-09-10 08:32:54 +02:00
Dai Ha f8b0d42a5c fleetd #431 follow-up: profileForSlot's javadoc named a caller that does not exist
CI / contract (push) Successful in 1m28s
CI / build (push) Failing after 1m53s
The javadoc said "what the spawn lifecycle reads". Nothing in src/main calls
profileForSlot at all, in either the ".profileForSlot(" or the
"::profileForSlot" form. My own #431 ticket text repeated that sentence as a
fact and ranked the three accessors by it, and the #432 worker copied it into
the test file's comment and one assertion message. Corrected in all three
places; the ticket correction is posted on #431.

What the spawn lifecycle actually reads for an architect's profile is the
SlotReservation that reserve() returns — SessionManager.java:225,
"reservation == null ? profile : reservation.profile()".

The corrected ranking, measured rather than read off the javadoc:
- nameForSlot is wired, at CallerResolver.java:137 (method reference, which is
  why a ".nameForSlot(" grep missed it)
- isSlot is reached through bind, called at MemberRegistry.java:235 and :376
- profileForSlot has no caller at all

Prose only. No behaviour change.
2026-09-10 13:25:19 +07:00
ltms d1e7d71eee Merge #432: pin profileForSlot, isSlot and nameForSlot against a live reload
CI / contract (push) Successful in 1m29s
CI / build (push) Successful in 1m34s
Verified in my own tree at a196d34: 1539 tests green, 0 compile errors, 0 files
changed under src/main (test-only, as reported).

Mutated two halves the worker's own proof did not cover, with a control in the
same script and the changed line echoed back each time:

- roleForSlot returning ARCHITECT for any configured slot (flatten is
  role-blind) -> 2 failures, incl. CallerResolverTest
  .aBoundNonArchitectSlotStillResolvesAsAWorker:454. So the role-blind flatten
  cannot grant ARCHITECT through a non-architect pool; that was already pinned.
- nameForSlot parsing the key suffix instead of reading config -> 1 failure,
  MemberRegistryLiveTest.nameForSlotReflectsANameChangedByReload:255. The new
  test has teeth beyond the freeze the worker ran.

Control run unmutated: green.

The ticket's severity ranking was wrong and I corrected it on #431. profileForSlot
has no caller in src/main in either the "." or "::" form, so its javadoc ("what
the spawn lifecycle reads") names a caller that does not exist; nameForSlot is
wired at CallerResolver.java:137; isSlot is reached through bind, at :235 and
:376. A follow-up commit fixes the test prose that repeated my claim.
2026-09-10 08:24:26 +02:00
Dai Ha 7fd914df1a fleetd #422 follow-up: make the model gate's own state observable
CI / contract (pull_request) Successful in 45s
CI / build (pull_request) Successful in 1m53s
PeerLauncher.disabledModels() reported an empty set both when no
models: block exists and when a block exists with nothing off, so
fleet_profiles/GET /profiles and the startup log could not tell an
inert gate from an armed one reporting zero. Add
PeerLauncher.ModelGateState (configured + off), a modelGateState()
default method disabledModels() now delegates to, and a
CompositePeerLauncher override that reads models0() once and
distinguishes the null-supplier case (no models: block) from a real,
config-supplied block via identity against the NO_MODELS_CONFIGURED
sentinel — reusing the exact accessor the spawn gate itself reads, per
the fleetd #404 lesson.

Wires the state into a new startup log line (Fleetd.modelGateCoverageLine)
and a new modelGateArmed field in FleetMcp.profilesView, reported
unconditionally alongside the existing modelsOff set.
2026-09-10 13:18:52 +07:00
Dai Ha a196d34455 fleetd #431: pin profileForSlot, isSlot and nameForSlot against a live reload
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Successful in 2m7s
#424 made MemberRegistry.slots() re-read fleet: on every call, but only
roleForSlot was tested against a real reload. profileForSlot, isSlot and
nameForSlot all have the same live-read line and none was pinned — proved by
freezing each to a construction-time snapshot and watching the full suite
stay green.

Adds 4 tests to MemberRegistryLiveTest, each driving a real ConfigRef.reload()
against a @TempDir config file (never two frozen registries compared in
memory, which would test the constructor instead of the reload):
- profileForSlotReflectsAProfileChangedByReload
- isSlotStopsReportingASlotRemovedByReload / isSlotStartsReportingASlotAddedByReload
- nameForSlotReflectsANameChangedByReload

No production change. profileForSlot has no call site anywhere in src/main
yet, so there is no spawn-lifecycle seam to drive the test through beyond the
accessor itself.
2026-09-10 13:06:56 +07:00
Dai Ha 766772763f Merge #428: revoke the ARCHITECT privilege on reload, not just future spawns
CI / contract (push) Successful in 1m27s
CI / build (push) Successful in 1m32s
fleetd #424. MemberRegistry.slots() now re-reads fleet: through a supplier,
so removing an architect slot demotes the bound pane on its very next request.
The boundEntries cache and entryFor() fallback from the first round are gone:
the slot OCCUPANCY (terminalToSlot) survives a reload, the ARCHITECT role does
not. That split was my own ticket wording's fault -- I asked for a test that a
bound architect "survives the rebuild", which conflated the binding with the
privilege.

Conflict resolved by hand in ConfigRef.java: #422 (models:) and #424
(architects) both rewrote the same Hot bullet. Kept both.

Also corrected two claims #424's own second commit left stale -- 7f672f0
reversed the behaviour but never touched ConfigRef, whose whole job is to tell
the operator what a reload does:
  - the Hot bullet said MemberRegistry's rule "keeps a live session's identity
    even after its slot is removed from config"
  - the reload-report comment said "only a NEW bind is refused"
Both now say what the code does: removal revokes ARCHITECT on the next
request, and only the slot occupancy survives.

Verified by the lead: 1535 tests, 0 failures, 0 compile errors, BUILD SUCCESS
on the merged tree.

Mutation of three halves the worker's own proof did not cover -- profileForSlot,
nameForSlot and isSlot each pointed at a frozen snapshot taken at construction
(live readers 5 -> 4, each mutation naming its method and line). All three
PASSED at 1535. The ticket's own fix is well pinned; these three sibling live
reads are not. Follow-up filed.
2026-09-10 12:54:35 +07:00
ltms eab8185d7b Merge #429: enforce the models.allow on/off gate at spawn, hot — including under fixed placement
CI / contract (push) Successful in 1m29s
CI / build (push) Successful in 1m33s
fleetd #422 + follow-up. Verified by the lead: 1524 tests green, 0 compile errors.

Mutation proof of three halves the worker's own proof did not cover:
- modelOffProfiles() -> Set.of() (the feed into ctx.modelOff): kills 3, incl. fixedPlacementSkipsAnOffModelProfileToo
- enforceModelEnabled() -> no-op (the explicit-spawn gate): kills 3, incl. the real-ConfigRef hot-reload test
- disabledModels() -> Set.of() (the reporting accessor): kills 1, alone
2026-09-10 07:44:35 +02:00
Dai Ha 7f672f0fb8 fleetd #424: revoke the ARCHITECT privilege on reload, not just future spawns
CI / contract (pull_request) Successful in 43s
CI / build (pull_request) Successful in 2m3s
Correction to the #424 fix in PR #428: the ticket asked to revoke a removed
architect slot, but the previous change (boundEntries) kept BOTH the binding
and the ARCHITECT privilege alive for an already-bound session after its slot
left config. That left the ticket's actual headline defect half-open.

The corrected rule: config governs both what may be bound next AND what a
bound slot still grants. Removing a slot now demotes its bound session to
worker on the very next request (roleForSlot/nameForSlot read slots() with no
cache, so CallerResolver.resolve falls through to Principal.worker(...)). The
terminalToSlot binding itself is untouched by a reload, on purpose: dropping
it would double-book the slot key and break unbind's compare-safe contract.

- Delete boundEntries and entryFor; profileForSlot/roleForSlot/nameForSlot/
  isSlot all read slots() directly, live, with no cache.
- Rewrite the class doc's binding rule for the corrected semantic.
- Replace the old "survives removal" test with anArchitectAlreadyBoundToASlotIsDemotedByReload,
  asserted through a real CallerResolver.resolve (not the roleForSlot seam),
  plus two tests for what must NOT change: the binding still occupies the
  slot after removal (a second terminal cannot claim it, even once the slot
  returns to config), and unbind still succeeds for the original terminal.

Verified snapshot()/CallerResolver.members() need no change: snapshot() only
ever reported raw terminalToSlot occupancy, and CallerResolver.members() has
no production caller.
2026-09-10 12:29:35 +07:00
Dai Ha ce74e164c6 fleetd #424: make architect-slot identity checks read fleet.architects live
CI / contract (pull_request) Successful in 1m15s
CI / build (pull_request) Successful in 1m32s
MemberRegistry used to flatten fleet.architects into an unmodifiable map at
construction, so removing (revoking) an architect slot from config never
took effect: reserve()/requireSlotFor() kept granting spawns against the
frozen snapshot forever, while ConfigRef told the operator "already
applied" for the wrong consumer.

- MemberRegistry gains a live constructor (MemberRegistry.live(Supplier))
  that re-flattens fleet.architects/developers/reviewers on every
  slots()/slotsFor() call, so reserve() and requireSlotFor() (which both
  read through slotsFor) govern the NEXT spawn with no restart. The frozen
  single-arg constructor is kept for tests and fixed/code-built configs.
- Binding rule: config governs what may be bound next; it never
  retroactively unbinds a live session. A slot removed from config while a
  terminal is bound to it keeps that binding. To keep the bound terminal's
  IDENTITY too (CallerResolver.resolve reads roleForSlot/nameForSlot on
  every request), every successful bind now caches the slot's Entry into a
  new boundEntries map; roleForSlot/nameForSlot/profileForSlot/isSlot fall
  back to it when the slot is no longer live, and unbind clears it in the
  same critical section it clears the binding.
- Fleetd.java now wires MemberRegistry.live(() -> config.get().fleet())
  instead of the frozen constructor.
- ConfigRef: corrected the fleet.leaders split-key message and the class
  doc's Hot bullet — architects is now hot for two independent consumers
  (CompositePeerLauncher for placement, MemberRegistry for identity), not
  only the one the message used to name. An architects-only edit still
  reports nothing beyond "config reloaded", which is now honest since the
  key really is fully hot for both consumers.

Added MemberRegistryLiveTest: real ConfigRef.reload() against a @TempDir
file, both directions (slot removed / slot added) for requireSlotFor and
reserve tested separately, plus a bound-architect-survives-removal test
that checks the binding AND the identity (roleForSlot/nameForSlot).

Mutation-tested: reverting requireSlotFor to a frozen snapshot fails
requireSlotForRefusesAProfileWhoseSlotWasRemovedByReload and its mirror;
reverting reserve the same way fails the two reserve tests; dropping the
boundEntries fallback fails the survives-removal test's roleForSlot
assertion. All three restored before commit.

mvn clean install: Tests run: 1510, Failures: 0, Errors: 0, Skipped: 0 —
BUILD SUCCESS.
2026-09-10 12:05:18 +07:00
ltms e60f892efd Merge pull request 'fleetd #415: split coverage() feature-state wording by pattern fallback semantics' (#423) from worker/415-coverage-wording-2cbf9c-5 into main
CI / contract (push) Successful in 1m13s
CI / build (push) Successful in 1m37s
2026-09-10 06:57:33 +02:00
Dai Ha ce05886831 fleetd #415: pin which UnsetMeaning Fleetd pairs with which pattern key
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Successful in 1m54s
Review found a gap: the earlier tests all called CompletionResolver.coverage()
directly, supplying the UnsetMeaning themselves — proving the enum's wording,
never that Fleetd's two call sites pair the right meaning with the right key.
Swapping the two UnsetMeaning arguments at those call sites (recreating #415's
defect with exhaustedPattern and errorPattern exchanged) compiled with 0 errors
and left all 1506 tests green.

Extract the two coverage-line call sites out of main() into package-private
static factories (Fleetd.exhaustedPatternCoverageLine /
errorPatternCoverageLine), the same pattern already used for capacitySource
and worktreeBranchLookup. Add FleetdPatternCoverageLineTest, which calls both
factories directly and asserts the actual wording each produces for the same
empty-coverage input, including that the two differ.

Also recorded the swap-mutation measurement (0 errors, 1506 green) in
UnsetMeaning's javadoc so a future reader does not delete the new test as
redundant with CompletionResolverTest.
2026-09-10 11:54:02 +07:00
Dai Ha be123d0ac7 fleetd #415: split coverage() feature-state wording by pattern-key fallback semantics
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Successful in 2m4s
CompletionResolver.coverage() measured pattern coverage (how many profiles set
a key) but its 'off' wording read as feature state. That is false for
errorPattern: an unset errorPattern still runs the classification against the
built-in BACKEND_ERROR pattern (CompletionResolver.java:84), so the empty case
is not off.

Add CompletionResolver.UnsetMeaning (OFF / BUILT_IN_DEFAULT), a required
parameter every coverage() call must supply — no defaulted overload, so a
future third pattern key cannot compile without stating what unset means for
it. Fleetd.java now passes UnsetMeaning.OFF for exhaustedPattern (no fallback
exists) and UnsetMeaning.BUILT_IN_DEFAULT for errorPattern.

Tests: updated the three existing empty/full/partial cases to pass the new
parameter, corrected the one test that pinned the old (wrong) errorPattern
wording, and added a test that asserts the same empty-coverage input produces
different wording for the two keys.
2026-09-10 11:40:32 +07:00
25 changed files with 1636 additions and 136 deletions
+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}` |
@@ -69,6 +69,7 @@ import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.TreeSet;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
@@ -231,6 +232,13 @@ public final class Fleetd {
profileName -> liveCountRef.get().apply(profileName),
quarantine,
outagePolicy);
// fleetd #422 follow-up: say which of the three model-gate states the daemon booted into —
// no models: block at all, a block armed with nothing off, or a block with N off — the same
// way exhaustedPatternCoverageLine/errorPatternCoverageLine report CB-578 stage A/fleetd
// #201 Unit 5 coverage just below. Read from workers.modelGateState() (never a separate
// config.get().models() here) so this line and fleet_profiles' modelGateArmed can never
// disagree about what CompositePeerLauncher's spawn gate actually enforces.
log.info("model gate (fleetd #422): {}", modelGateCoverageLine(workers.modelGateState()));
// CB-504: under supervision (launchd/systemd) fleetd can start before herdr's socket
// exists. The client itself is lazy — it connects per call — but the orphan reap below is
// the first thing that actually talks to herdr, so without this wait a boot-order race
@@ -352,7 +360,11 @@ public final class Fleetd {
// pane resolves to an architect until the later spawn lifecycle binds one. The registry is
// what CallerResolver resolves against and what that lifecycle will read profiles from;
// nothing here spawns a slot.
MemberRegistry members = new MemberRegistry(cfg.fleet());
// fleetd #424: MemberRegistry.live re-reads fleet.architects through `config` on every
// reserve/requireSlotFor call, so a reload that removes or adds an architect slot governs
// the next spawn with no restart — the frozen `new MemberRegistry(cfg.fleet())` this used
// to be let a "revoked" slot keep granting new architect spawns forever.
MemberRegistry members = MemberRegistry.live(() -> config.get().fleet());
sessions.setMemberLifecycle(members);
if (!members.slots().isEmpty()) {
log.info("member slots: {} configured {} — none bound yet (a slot is idle until the "
@@ -380,8 +392,7 @@ public final class Fleetd {
.map(session -> exhaustedPatternsByProfile.get(session.profile()))
.orElse(null);
log.info("backend-exhausted classification (CB-578 stage A): {}",
CompletionResolver.coverage("exhaustedPattern", cfg.profiles().keySet(),
exhaustedPatternsByProfile.keySet()));
exhaustedPatternCoverageLine(cfg.profiles().keySet(), exhaustedPatternsByProfile.keySet()));
// fleetd #201 Unit 5: classify a completion-fallback scrape that matches a profile's
// configured backend-error refusal (a credential outage, a provider 5xx) as a backend error
// rather than handing it back as a real answer. Compiled once at startup, keyed by profile
@@ -401,8 +412,7 @@ public final class Fleetd {
BackendErrorPatternLookup backendErrorPatterns = backendErrorPatternLookup(sessions::roster,
errorPatternsByProfile);
log.info("backend-error classification (fleetd #201 Unit 5): {}",
CompletionResolver.coverage("errorPattern", cfg.profiles().keySet(),
errorPatternsByProfile.keySet()));
errorPatternCoverageLine(cfg.profiles().keySet(), errorPatternsByProfile.keySet()));
// CB-578 stage B: on a classification that actually wins, quarantine the exhausted profile's
// CREDENTIAL — not the profile name — so a profile sharing that credential (e.g. two models
// on one OpenAI account) is refused too, not just the one that happened to report it. Reads
@@ -807,6 +817,66 @@ public final class Fleetd {
}, quarantine, profile -> startupExhaustedPatterns.containsKey(profile));
}
/**
* fleetd #415 (review follow-up): package-private factory for the CB-578 stage A {@code
* exhaustedPattern} startup coverage line, paired explicitly with {@link
* CompletionResolver.UnsetMeaning#OFF} — {@code exhaustedPattern} has no fallback, so a
* profile with none configured really does have the classification off.
*
* <p>Extracted out of {@code main} for the same reason {@link #capacitySource} and {@link
* #worktreeBranchLookup} were: {@code coverage()}'s own tests ({@code CompletionResolverTest})
* prove it words {@code OFF} and {@link CompletionResolver.UnsetMeaning#BUILT_IN_DEFAULT}
* correctly when a test supplies the meaning itself — they cannot prove {@code main} pairs the
* right meaning with the right key, which is the actual fleetd #415 defect. <b>Measured:</b>
* swapping the {@code UnsetMeaning} arguments between this method and {@link
* #errorPatternCoverageLine} — recreating #415's defect with the two keys exchanged — compiled
* with 0 errors and left all 1506 existing tests green before {@code
* FleetdPatternCoverageLineTest} was added to catch exactly that swap.
*/
static String exhaustedPatternCoverageLine(Set<String> allProfiles, Set<String> configuredProfiles) {
return CompletionResolver.coverage("exhaustedPattern", CompletionResolver.UnsetMeaning.OFF,
allProfiles, configuredProfiles);
}
/**
* fleetd #415 (review follow-up): the {@code errorPattern} counterpart of {@link
* #exhaustedPatternCoverageLine}, paired explicitly with {@link
* CompletionResolver.UnsetMeaning#BUILT_IN_DEFAULT} — an unset {@code errorPattern} still runs
* backend-error classification against {@code CompletionResolver}'s built-in {@code
* BACKEND_ERROR} pattern, so the empty case is not "off". See {@link
* #exhaustedPatternCoverageLine}'s javadoc for the measured swap mutation this pairing guards
* against.
*/
static String errorPatternCoverageLine(Set<String> allProfiles, Set<String> configuredProfiles) {
return CompletionResolver.coverage("errorPattern", CompletionResolver.UnsetMeaning.BUILT_IN_DEFAULT,
allProfiles, configuredProfiles);
}
/**
* fleetd #422 follow-up: package-private factory for the startup line reporting which of the
* three central {@code models.allow:} gate states the daemon booted into. Extracted the same
* way {@link #exhaustedPatternCoverageLine}/{@link #errorPatternCoverageLine} are, so a
* dedicated test can call it directly rather than parsing log output, and so {@code main}'s
* only source for this line is {@link PeerLauncher#modelGateState()} — never a second,
* independently-derived read of {@code cfg.models()} that could disagree with what {@code
* CompositePeerLauncher}'s spawn gate actually enforces (the fleetd #404 lesson).
*
* <p>Unlike the two pattern-key lines above, there is no {@code UnsetMeaning} choice to make
* here: {@link PeerLauncher.ModelGateState#configured()} already states, unambiguously, whether
* an empty {@link PeerLauncher.ModelGateState#off()} means "no {@code models:} block to gate
* with" or "a block armed and currently reporting zero off" — the exact two states a bare
* {@code disabledModels()} read could not tell apart before this ticket.
*/
static String modelGateCoverageLine(PeerLauncher.ModelGateState state) {
if (!state.configured()) {
return "not configured (no models: block — nothing is gated, and nothing can be)";
}
Set<String> off = state.off();
return off.isEmpty()
? "armed (models: block present; 0 models currently turned off)"
: "armed (" + off.size() + " model(s) turned off: " + new TreeSet<>(off) + ")";
}
/**
* fleetd #416: production source for {@code fleet_list}'s per-profile capacity facts.
*
@@ -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();
};
}
@@ -11,6 +11,7 @@ import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.function.Supplier;
/**
* The architect-slot registry (CB-548): every gateway-local architect name and the strong-model
@@ -19,16 +20,43 @@ import java.util.Objects;
*
* <p>Two halves, split by who owns each:
* <ul>
* <li><b>slots</b> — configured once, keyed by the gateway-local unique name; each carries the
* {@code profile} reference the spawn lifecycle reads when it stands the slot up. A read-only
* snapshot taken at construction.</li>
* <li><b>terminal bindings</b> — owned by this registry and initially <em>empty</em>. Config
* declares no architect terminal, so at startup every slot is idle and nothing resolves to an
* architect; a session only becomes one when the spawn lifecycle {@linkplain #bind(String,
* String) binds} its terminal to a slot. {@link CallerResolver} reads this through
* {@link #snapshot()} to turn a pane into an {@link Role#ARCHITECT}.</li>
* <li><b>slots</b> — read from {@code fleet.architects}/{@code developers}/{@code reviewers}
* (see {@link #slots()}), each carrying the {@code profile} reference the spawn lifecycle
* reads when it stands the slot up. <strong>Live, since fleetd #424</strong>: {@link #live}
* re-reads {@code fleet:} on every call, through a supplier the same shape as
* {@code CompositePeerLauncher}'s (see {@code ConfigRef}'s class doc) — so a config reload
* that removes or adds an architect slot governs the <em>next</em> spawn with no restart.
* Only {@link #MemberRegistry(FleetConfig.Fleet)} freezes the pool at construction, and that
* constructor exists for tests and for the (rare) case of wiring a fixed, code-built config.</li>
* <li><b>terminal bindings</b> — owned by this registry, initially <em>empty</em>, and
* <strong>never</strong> touched by a reload. Config declares no architect terminal, so at
* startup every slot is idle and nothing resolves to an architect; a session only becomes one
* when the spawn lifecycle {@linkplain #bind(String, String) binds} its terminal to a slot.
* {@link CallerResolver} reads this through {@link #snapshot()} to turn a pane into an
* {@link Role#ARCHITECT}.</li>
* </ul>
*
* <p><strong>The binding rule (fleetd #424): config governs what a bound slot still grants, as
* well as what may be bound next.</strong> Removing a slot from config revokes it — that is the
* ticket's entire point ("Revoking an architect slot does not revoke it"). Revoking it means an
* architect already bound to that slot loses the ARCHITECT privilege on its very next request:
* {@link #roleForSlot} and {@link #nameForSlot} read {@link #slots()} directly, with no cache, so
* the moment a slot drops out of config, {@link CallerResolver#resolve} (which calls both on every
* request from a bound pane, {@code CallerResolver.java:220}) can no longer confirm the pane's slot
* is an architect slot, and the pane falls through to {@code Principal.worker(...)}. What does
* <em>not</em> change is the {@code terminalToSlot} <em>occupancy</em> — the binding created by
* {@link #bind} is untouched by a reload, on purpose: unbinding it here would double-book the slot
* key (a second terminal could then bind to the "freed" key while the first is still the terminal
* the operator actually meant to demote) and would silently break {@link #unbind}'s compare-safe
* contract, which needs the original {@code terminal → slot} pair intact to remove it cleanly. So
* the demoted session keeps occupying its slot — {@link #slotForTerminal} and {@link #snapshot()}
* still name it — it just no longer resolves as an architect through that occupancy, and a fresh
* spawn still cannot bind to the same key while it is occupied ({@link #reserve}/
* {@link #requireSlotFor} refuse it anyway, since it is gone from {@link #slots()}). The demoted
* session's own turn is unaffected: {@code fleet_reply}'s authorization
* ({@code Authz.Action.REPLY}) is {@code caller.ownsSession(targetSession)} — identity by terminal,
* not by role — so a demoted architect can still end its own turn normally.
*
* <p>Spawning/lifecycle is deliberately a separate unit: this class only owns the bindings and
* exposes the map the resolver resolves against plus the profile lookup lifecycle will call.
* Nothing here creates or manages an architect session.
@@ -55,14 +83,42 @@ public final class MemberRegistry implements MemberLifecycle {
}
}
private final Map<String, Entry> slots;
private final Supplier<FleetConfig.Fleet> fleet;
/** Live {@code terminal_id → qualified slot key}; guarded by {@code terminalToSlot}. */
private final Map<String, String> terminalToSlot = new HashMap<>();
/** Slot keys held between reservation and the terminal binding. Guarded by terminalToSlot. */
private final java.util.Set<String> reservedSlots = new java.util.HashSet<>();
/** Flatten every role pool in {@code fleet} into one registry. Leaders are not members. */
/**
* Freeze the pool at construction — for tests, and for the rare case of wiring a fixed,
* code-built config. Production wiring should prefer {@link #live}, which re-reads
* {@code fleet:} on every call.
*/
public MemberRegistry(FleetConfig.Fleet fleet) {
this(() -> fleet);
}
private MemberRegistry(Supplier<FleetConfig.Fleet> fleet) {
this.fleet = fleet;
}
/**
* Live variant (fleetd #424): {@code fleet} is read fresh on every {@link #slots()} call — pass
* {@code () -> config.get().fleet()}, the same supplier shape {@code CompositePeerLauncher}
* already uses for placement — so a reload that adds or removes an architect slot governs the
* next spawn's {@link #reserve}/{@link #requireSlotFor} check with no restart. A separate,
* private constructor rather than a same-arity public overload of
* {@link #MemberRegistry(FleetConfig.Fleet)}: a {@code FleetConfig.Fleet} and a
* {@code Supplier<FleetConfig.Fleet>} overload are ambiguous for a literal {@code null} — the
* same reason {@code CallerResolver.withLeads} is a static factory rather than a fourth
* constructor overload.
*/
public static MemberRegistry live(Supplier<FleetConfig.Fleet> fleet) {
return new MemberRegistry(Objects.requireNonNull(fleet, "fleet"));
}
/** Flatten every role pool in {@code fleet} into one map. Leaders are not members. */
private static Map<String, Entry> flatten(FleetConfig.Fleet fleet) {
Map<String, Entry> flat = new LinkedHashMap<>();
if (fleet != null) {
for (MemberRole role : MemberRole.values()) {
@@ -74,18 +130,22 @@ public final class MemberRegistry implements MemberLifecycle {
});
}
}
this.slots = Collections.unmodifiableMap(flat);
return Collections.unmodifiableMap(flat);
}
/** The configured slots, keyed by qualified {@link Entry#key()}. Unmodifiable snapshot. */
/**
* The configured slots, keyed by qualified {@link Entry#key()}. Unmodifiable snapshot of
* {@code fleet:} <em>as of this call</em> — see the class doc for which constructor makes that
* live versus frozen.
*/
public Map<String, Entry> slots() {
return slots;
return flatten(fleet.get());
}
/** The slots belonging to {@code role}, in definition order. */
/** The slots belonging to {@code role}, in definition order, as of this call. */
public Map<String, Entry> slotsFor(MemberRole role) {
Map<String, Entry> out = new LinkedHashMap<>();
slots.forEach((key, e) -> {
slots().forEach((key, e) -> {
if (e.role() == role) {
out.put(key, e);
}
@@ -117,31 +177,49 @@ public final class MemberRegistry implements MemberLifecycle {
}
/**
* The strong-model profile a slot runs under — what the spawn lifecycle reads.
* The strong-model profile a slot runs under, as of this call.
*
* <p>Nothing in {@code src/main} calls this (fleetd #431 — grepped both the {@code
* .profileForSlot(} and the {@code ::profileForSlot} form). This javadoc used to say "what the
* spawn lifecycle reads", and that seam does not exist: the spawn lifecycle takes its profile
* from the {@link MemberLifecycle.SlotReservation} that {@code reserve} returns, never from
* here. Kept and pinned rather than deleted because it is the natural accessor for that seam
* if one is added; live for the same reason as {@link #roleForSlot}, so a reload cannot leave
* it answering for the old config.
*
* @return the slot's configured {@code profile}, or {@code null} if the slot is unknown or
* declares none
*/
public String profileForSlot(String slotName) {
Entry e = slots.get(slotName);
Entry e = slots().get(slotName);
return (e == null || e.profile() == null) ? null : e.profile();
}
/** The role a qualified slot key belongs to, or {@code null} when the key is unknown. */
/**
* The role a qualified slot key belongs to, or {@code null} when the key is not currently
* configured. Deliberately live, with no cache (fleetd #424, see the class doc's binding rule):
* removing a slot from config must make {@link CallerResolver#resolve} stop granting the
* ARCHITECT role for it on the very next request from a terminal that was bound to it, which is
* the ticket's whole point — revoking a slot must actually revoke it, not just refuse the next
* spawn.
*/
public MemberRole roleForSlot(String slotName) {
Entry e = slots.get(slotName);
Entry e = slots().get(slotName);
return e == null ? null : e.role();
}
/** The unqualified configured name for a slot, or {@code null} if it is unknown. */
/**
* The unqualified configured name for a slot, or {@code null} if it is not currently configured.
* Live for the same reason as {@link #roleForSlot} — see the class doc's binding rule.
*/
public String nameForSlot(String slotName) {
Entry e = slots.get(slotName);
Entry e = slots().get(slotName);
return e == null ? null : e.name();
}
/** True when {@code slotName} is a configured architect slot. */
public boolean isSlot(String slotName) {
return slots.containsKey(slotName);
return slots().containsKey(slotName);
}
/**
@@ -32,18 +32,26 @@ import java.util.function.Supplier;
* not the fact that they are config. Most of {@code fleet:} — every role pool
* ({@code architects}/{@code developers}/{@code reviewers}), {@code charters}, and
* {@code tabLabel} — is read the same live way, through the same supplier
* ({@code () -> config.get().fleet()}). <strong>But {@code fleet:} as a whole is NOT in this
* class</strong>: {@code fleet.leaders} inside the same key is frozen, which is exactly what
* makes {@code fleet:} split rather than hot — see below. {@code models:} (fleetd #422) joined
* this class whole: {@link FleetConfig#validateModels()} re-runs fully against the fresh
* config on every {@link #reload()} (via {@link FleetConfig#validateAll()}), refusing a bad
* edit outright rather than caching a stale copy anywhere, and the on/off half added by
* fleetd #422 is read live both by {@code CompositePeerLauncher}'s spawn gate
* ({@code enforceModelEnabled} and its candidate filter) and by {@code fleet_profiles}/
* {@code GET /profiles} (via {@code PeerLauncher.disabledModels()}). Nothing about
* {@code models:} is baked into an object built at startup, so — unlike the deferred keys
* below — there is no frozen half left to report; it moved here from deferred rather than
* joining split.</li>
* ({@code () -> config.get().fleet()}). {@code architects} in particular is hot for
* <strong>two independent consumers</strong> (fleetd #424): {@code CompositePeerLauncher}
* reads it live for placement (which profile an unqualified architect spawn may land on), and
* {@code MemberRegistry} separately reads it live, through its own instance of the same
* supplier shape, for identity — both which slot a spawn may bind to <em>and</em> what a slot
* already bound still grants. Removing an architect slot from config therefore revokes the
* {@link dev.ltms.fleet.auth.Role#ARCHITECT} role on the bound pane's very next request; only
* the slot <em>occupancy</em> survives, so the demoted session still holds its slot key until
* it unbinds. See {@code MemberRegistry}'s class doc for that binding rule.
* <strong>But {@code fleet:} as a whole is NOT in this class</strong>: {@code fleet.leaders}
* inside the same key is frozen, which is exactly what makes {@code fleet:} split rather than
* hot — see below. {@code models:} (fleetd #422) joined this class whole: {@link
* FleetConfig#validateModels()} re-runs fully against the fresh config on every {@link
* #reload()} (via {@link FleetConfig#validateAll()}), refusing a bad edit outright rather than
* caching a stale copy anywhere, and the on/off half added by fleetd #422 is read live both by
* {@code CompositePeerLauncher}'s spawn gate ({@code enforceModelEnabled} and its candidate
* filter) and by {@code fleet_profiles}/{@code GET /profiles} (via
* {@code PeerLauncher.disabledModels()}). Nothing about {@code models:} is baked into an
* object built at startup, so — unlike the deferred keys below — there is no frozen half left
* to report; it moved here from deferred rather than joining split.</li>
* <li><strong>Deferred</strong> — accepted into the new snapshot, but the wiring built at startup
* keeps the old value until a restart: {@code lifecycle:}, {@code leadHeartbeat:},
* {@code idleSleepGuard:} ({@code Fleetd.java} reads it once, at startup, to decide whether
@@ -525,26 +533,35 @@ public final class ConfigRef implements Supplier<FleetConfig> {
+ "opened once and needs a restart; the broker URI env-var name kept out of a "
+ "member's environment is read live on every spawn and already applied");
}
// fleetd #333: unlike health/coordinator above, most of `fleet:` (architects, developers,
// reviewers, charters, tabLabel) is genuinely hot — ConfigRefTest.aHotChangeIsAppliedAndRead-
// fleetd #333: unlike health/coordinator above, most of `fleet:` (developers, reviewers,
// charters, tabLabel) is genuinely hot — ConfigRefTest.aHotChangeIsAppliedAndRead-
// ThroughGet and aCharterChangeIsHotAndReachesTheLiveConfig prove it reaches the live config
// with no restart note. Only fleet.leaders is frozen (Fleetd.java:281 reads
// cfg.fleet().leaders() off the startup snapshot to build both the LeadTabScanner's
// tab-label-to-name map, wired into CallerResolver.withLeadsAndMembers at Fleetd.java:620/624,
// and — when herdr answered — LeadLauncher(...).ensureLeads() at Fleetd.java:315, which
// auto-launches each lead up to its `instances` count; neither is rebuilt on reload). So this
// compares fleet.leaders alone, not the whole Fleet record: comparing the whole record would
// report "split" for a tabLabel-only or charters-only change that is actually fully hot,
// which is the over-claim mirror of the under-claim bug this class exists to prevent.
// with no restart note. `architects` is hot too, and — since fleetd #424 — hot for BOTH of
// its consumers, not just the one this comment used to name: CompositePeerLauncher reads it
// live for PLACEMENT through the () -> config.get().fleet() supplier named in the class doc's
// Hot bullet, and MemberRegistry separately reads it live for IDENTITY (which slot a spawn
// may bind to, AND what a slot already bound still grants) through its own instance of that
// same supplier shape — see MemberRegistry.live and its class doc for the binding rule:
// removing a slot revokes ARCHITECT on the bound pane's very next request, and only the slot
// OCCUPANCY survives, so the demoted session keeps its slot key until it unbinds. Only
// fleet.leaders is frozen (Fleetd.java:281 reads cfg.fleet().leaders() off the startup
// snapshot to build both the LeadTabScanner's tab-label-to-name map, wired into
// CallerResolver.withLeadsAndMembers at Fleetd.java:620/624, and — when herdr answered —
// LeadLauncher(...).ensureLeads() at Fleetd.java:315, which auto-launches each lead up to its
// `instances` count; neither is rebuilt on reload). So this compares fleet.leaders alone, not
// the whole Fleet record: comparing the whole record would report "split" for a tabLabel-only
// or architects-only change that is actually fully hot, which is the over-claim mirror of the
// under-claim bug this class exists to prevent.
if (!Objects.equals(leadersOf(old), leadersOf(fresh))) {
changed.add("fleet: fleet.leaders (each lead's tab, workspace, cwd, profile and "
+ "instances count) is read once at startup to build the LeadTabScanner's "
+ "identity map and to auto-launch leads, and neither is rebuilt on reload, so a "
+ "lead added, removed, or given a new tab: label needs a restart — until then it "
+ "stays unrecognised, and a caller from its new tab resolves as a worker, not a "
+ "lead; the rest of fleet: (architects, developers, reviewers, charters, "
+ "tabLabel) is read live through the supplier on CompositePeerLauncher and "
+ "already applied");
+ "lead; the rest of fleet: (developers, reviewers, charters, tabLabel) is read "
+ "live through the supplier on CompositePeerLauncher, and architects is read "
+ "live through that same supplier for placement AND through a separate supplier "
+ "on MemberRegistry for spawn-time identity — both already applied");
}
// Kept in step with SPLIT_KEYS the same way changedColdKeys is kept in step with COLD_KEYS —
// every message here must be traceable to one of the split keys the class doc documents.
@@ -698,17 +698,57 @@ public final class CompletionResolver implements TurnListener {
}
/**
* Coverage summary for the CB-578 stage A exhausted-pattern classification, logged at startup
* What an unset pattern key means for the classification it configures (fleetd#415).
* {@code coverage()} cannot infer this from the key's name — the two keys it currently
* describes disagree on it, and a string comparison on the name would just move the same bug
* to a new spot — so every caller must state it explicitly.
*
* <p><strong>This alone does not prove a caller passes the right one for its key.</strong> A
* test that calls {@code coverage()} directly and supplies the meaning itself only proves this
* enum is worded correctly, never that {@code Fleetd}'s two call sites pair each key with its
* true meaning — that pairing is #415's actual defect. Measured on review: swapping the two
* {@code UnsetMeaning} arguments at those call sites (giving {@code exhaustedPattern} the
* built-in-default wording and {@code errorPattern} the off wording — #415's exact defect with
* the keys exchanged) compiled with 0 errors and left all 1506 existing tests green. See
* {@code dev.ltms.fleet.Fleetd#exhaustedPatternCoverageLine}/{@code #errorPatternCoverageLine}
* and {@code FleetdPatternCoverageLineTest}, which exists specifically to catch that swap.
*/
public enum UnsetMeaning {
/** No fallback exists: a profile with no configured pattern truly has this classification off. */
OFF,
/** A built-in pattern applies when unset: the classification still runs for that profile. */
BUILT_IN_DEFAULT
}
/**
* Coverage summary for a fleetd#201/CB-578-style pattern-key classification, logged at startup
* the way {@link dev.ltms.fleet.health.FleetHealthMonitor#coverage} is — so an operator can
* see whether the classification is on, and for which profiles, without reading every
* profile's config by hand.
*
* <p>fleetd#415: this method measures <em>pattern coverage</em> — how many profiles set the
* key — which is not the same thing as <em>feature state</em> for a key with a fallback. For
* {@code errorPattern}, an empty {@code configuredProfiles} still runs the classification
* against {@code CompletionResolver}'s built-in compatibility pattern ({@link #BACKEND_ERROR}
* at line ~84); for {@code exhaustedPattern} there is no fallback, so empty really does mean
* off. {@code unsetMeaning} is the single, required source of that fact — see
* {@link dev.ltms.fleet.config.FleetConfig#rejectMalformedProfilePatterns} lines ~2029-2032 for
* where it is documented for config authors. It is a required parameter, not a defaulted
* overload: a third pattern key added later must supply one to compile at all, rather than
* silently inheriting whichever wording this method happened to default to.
*
* @param allProfiles every configured profile name
* @param configuredProfiles the subset of {@code allProfiles} that carry an exhausted pattern
* @param configuredProfiles the subset of {@code allProfiles} that carry the pattern
*/
public static String coverage(String patternKey, Set<String> allProfiles, Set<String> configuredProfiles) {
public static String coverage(String patternKey, UnsetMeaning unsetMeaning, Set<String> allProfiles,
Set<String> configuredProfiles) {
if (configuredProfiles.isEmpty()) {
return "off (no profile has an " + patternKey + " configured; profiles: " + sorted(allProfiles) + ")";
return switch (unsetMeaning) {
case OFF -> "off (no profile has an " + patternKey + " configured; profiles: "
+ sorted(allProfiles) + ")";
case BUILT_IN_DEFAULT -> "built-in default for all profiles (no profile customises "
+ patternKey + "; profiles: " + sorted(allProfiles) + ")";
};
}
Set<String> unconfigured = new TreeSet<>(allProfiles);
unconfigured.removeAll(configuredProfiles);
@@ -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).
@@ -1181,12 +1250,21 @@ public final class FleetMcp {
result.put("coolingOff", coolingOff);
}
// fleetd #422: read the exact same accessor CompositePeerLauncher's spawn gate reads
// (PeerLauncher.disabledModels(), which for the composite is models0().offIds()) — never a
// (PeerLauncher.modelGateState(), which for the composite is models0() read live) — never a
// separately-derived answer, so this status can never overstate or understate what the gate
// actually enforces (the fleetd #404 lesson).
Set<String> modelsOff = workers.disabledModels();
if (!modelsOff.isEmpty()) {
result.put("modelsOff", new ArrayList<>(modelsOff));
//
// fleetd #422 follow-up: "armed" and "off" come from the ONE modelGateState() call below,
// never two independent reads of the gate — a reload landing between two separate reads
// could otherwise make them disagree. modelGateArmed is reported unconditionally (never
// omitted like quarantined/coolingOff above) precisely so a lead can tell "no models: block
// at all" (false) apart from "a models: block with nothing currently off" (true, with
// modelsOff simply absent below) — the two states PeerLauncher.disabledModels() alone
// cannot distinguish, both reporting an empty set.
PeerLauncher.ModelGateState modelGate = workers.modelGateState();
result.put("modelGateArmed", modelGate.configured());
if (!modelGate.off().isEmpty()) {
result.put("modelsOff", new ArrayList<>(modelGate.off()));
}
return result;
}
@@ -1314,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();
@@ -1321,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;
}
@@ -1654,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()));
}
@@ -644,14 +644,29 @@ public final class CompositePeerLauncher implements PeerLauncher {
}
/**
* fleetd #422: read live off {@link #models0()} — the exact same accessor {@link
* fleetd #422 follow-up: the single live read that answers both "is the models.allow: gate
* armed" and "which models are off", off the exact same accessor ({@link #models0()}) {@link
* #enforceModelEnabled} and {@link #modelOffProfiles} read — so {@code fleet_profiles}/{@code
* GET /profiles} can never report a different answer than the gate enforces (the fleetd #404
* lesson).
* GET /profiles} (via {@link PeerLauncher#disabledModels()}, which now delegates here) can
* never report a different answer than the gate enforces (the fleetd #404 lesson), and the
* startup log line built from this can never disagree with either.
*
* <p>{@link #models0()} itself normalizes a {@code null} {@link #models} read to the shared
* {@link #NO_MODELS_CONFIGURED} sentinel — deliberately the one object no config-supplied
* {@code Models} instance can ever be identical to, since it is private to this class — so
* comparing by reference here recovers exactly the fact {@code models0()}'s normalization
* would otherwise erase: whether the live source was {@code null} (no {@code models:} block,
* armed = false) or a real, config-supplied block (armed = true, even one whose {@code allow:}
* is itself empty or absent — {@link FleetConfig.Models}'s "absent or empty allow: is off"
* wording governs config-load validation, a distinct question from whether this gate is armed
* for reporting).
*/
@Override
public Set<String> disabledModels() {
return models0().offIds();
public PeerLauncher.ModelGateState modelGateState() {
FleetConfig.Models m = models0();
return m == NO_MODELS_CONFIGURED
? PeerLauncher.ModelGateState.notConfigured()
: PeerLauncher.ModelGateState.armed(m.offIds());
}
@Override
@@ -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.
@@ -202,8 +202,56 @@ public interface PeerLauncher {
* can drift from what the gate ({@code CompositePeerLauncher.enforceModelEnabled} and its
* candidate filter) actually enforces. A default of {@code Set.of()} keeps every other {@link
* PeerLauncher} implementer (the herdr adapters, and the two test-fake implementers) unchanged.
*
* <p>fleetd #422 follow-up: this alone cannot tell "no {@code models:} block at all" from "a
* {@code models:} block where nothing is currently off" — both report an empty set here. Delegates
* to {@link #modelGateState()} so the two facts always come from the one read {@link
* #modelGateState()}'s implementer makes; do not override this method separately from that one.
*/
default Set<String> disabledModels() {
return Set.of();
return modelGateState().off();
}
/**
* Whether the central {@code models.allow:} gate (fleetd #422) is armed at all, together with
* which model ids are currently off — fleetd #422 follow-up. {@link #disabledModels()} alone
* cannot distinguish two states that both report an empty set: a host with no {@code models:}
* block (nothing is gated, and nothing can be) and a host WITH a {@code models:} block where
* nothing is currently turned off (the gate is armed and reporting zero). This method exists so
* a caller — the startup log, {@code fleet_profiles}/{@code GET /profiles} — can tell the two
* apart, the same reason {@code CompletionResolver.UnsetMeaning} exists: an accessor that can
* legitimately report "empty" must never let a caller guess why.
*
* <p>Default {@link ModelGateState#notConfigured()} — every launcher without a {@code models:}
* block to read from (the herdr adapters, and the two test-fake implementers), matching {@link
* #disabledModels()}'s own default of an empty set.
*/
default ModelGateState modelGateState() {
return ModelGateState.notConfigured();
}
/**
* fleetd #422 follow-up: the result of {@link #modelGateState()} — see that method's javadoc
* for why "armed" and "off" must be reported together from one read rather than as two
* separately-derived facts that a reload landing between them could make disagree.
*
* @param configured {@code true} when a {@code models:} block exists at all (armed), regardless
* of whether anything in it is currently turned off; {@code false} when there
* is no block to gate against
* @param off the model ids currently turned off; always empty when {@code configured} is
* {@code false}
*/
record ModelGateState(boolean configured, Set<String> off) {
public ModelGateState {
off = Set.copyOf(off);
}
public static ModelGateState notConfigured() {
return new ModelGateState(false, Set.of());
}
public static ModelGateState armed(Set<String> off) {
return new ModelGateState(true, off);
}
}
}
@@ -5,14 +5,12 @@ import java.util.List;
/**
* Backward-compatible placement: an unqualified spawn always resolves to the configured default
* profile, exactly as {@code CompositePeerLauncher} did before CB-518. This ignores caps
* ({@code maxLoad}) so that a pre-existing config behaves identically after upgrade — capacity
* gating for automatic placement is deliberately out of scope for {@code fixed}, exactly as it
* always has been. Reachability is a narrower exception (fleetd #315, below): a profile is never
* checked for reachability up front, only skipped once it has already failed in <em>this same</em>
* spawn call's retry loop — see the unreachable case below.
* profile, exactly as {@code CompositePeerLauncher} did before CB-518. Reachability is a narrower
* exception (fleetd #315, below): a profile is never checked for reachability up front, only
* skipped once it has already failed in <em>this same</em> spawn call's retry loop — see the
* unreachable case below.
*
* <p>Five exceptions walk past the default instead of returning it unconditionally:
* <p>Six exceptions walk past the default instead of returning it unconditionally:
* <ul>
* <li>Quarantine (CB-578 stage B): a quarantined default is a credential that just refused on
* a usage limit, not a transient capacity or reachability concern.
@@ -20,14 +18,24 @@ import java.util.List;
* ({@code BackendOutagePolicy}) — a separate, shorter-lived source from quarantine. When a
* profile is both quarantined and cooling off, only the quarantine reason is reported
* (exhaustion takes priority), matching {@code CompositePeerLauncher}'s explicit-spawn order.
* <li>At cap (fleetd #435): a profile whose live count has reached its {@code maxLoad}
* ({@link PlacementPolicyUtil#atCap}) — a documented, unconditional capacity limit (see
* {@code FleetConfig.Profile#maxLoad}), so {@code fixed} must gate on it exactly as {@code
* weighted}/{@code round-robin} already do via {@link PlacementPolicyUtil#available}. Before
* this fix {@code fixed} built its own {@link PlacementCandidate} for the default with {@code
* maxLoad} forced to {@code null}, so a capped default was chosen anyway on every unqualified
* spawn — the cap was advisory, not enforced, for the one placement policy every config uses
* by default. Reported only when quarantine and cooling off are both absent, matching {@code
* CompositePeerLauncher}'s explicit-spawn check order (quarantine, then cooling off, then max
* load, then model-off).
* <li>Model off (fleetd #422): a profile whose {@code model:} the operator has turned off in
* {@code models.allow:} — an operator decision, never a backend-reported outage, so it is a
* FOURTH, independent source from both quarantine and cooling off (never merged with either),
* exactly as {@code CompositePeerLauncher.enforceModelEnabled} and {@link
* PlacementPolicyUtil#available} treat it. When a profile is model-off <em>and</em> quarantined
* or cooling off, only the quarantine/cooling-off reason is reported — those still take
* priority, matching {@code CompositePeerLauncher}'s explicit-spawn check order (quarantine,
* then cooling off, then max load, then model-off).
* fifth, independent source from quarantine, cooling off, and at-cap (never merged with any
* of them), exactly as {@code CompositePeerLauncher.enforceModelEnabled} and {@link
* PlacementPolicyUtil#available} treat it. When a profile is model-off <em>and</em> quarantined,
* cooling off, or at cap, only the higher-priority reason is reported, matching {@code
* CompositePeerLauncher}'s explicit-spawn check order (quarantine, then cooling off, then max
* load, then model-off).
* <li>Unreachable (fleetd #315): {@code CompositePeerLauncher.spawn} retries a failed candidate
* on the next one and rebuilds the {@link PlacementContext} so {@code ctx.unreachable()}
* names every profile that already failed with {@code PeerUnreachableException} in this same
@@ -40,10 +48,10 @@ import java.util.List;
* {@code weighted}/{@code round-robin} skip it — an explicit {@code fleet_spawn} naming
* the profile is unaffected, only this automatic fallback walk.
* </ul>
* A fleet where nothing is ever quarantined, cooling off, model-off, unreachable, or weight-0 never
* exercises any of these paths, so today's behaviour is unchanged — in particular, the very first
* selection of a spawn call always sees an empty {@code unreachable} set, so the first choice is
* untouched.
* A fleet where nothing is ever quarantined, cooling off, at cap, model-off, unreachable, or
* weight-0 never exercises any of these paths, so today's behaviour is unchanged — in particular,
* the very first selection of a spawn call always sees an empty {@code unreachable} set, so the
* first choice is untouched.
*/
final class FixedPlacementPolicy implements PlacementPolicy {
@@ -51,13 +59,14 @@ final class FixedPlacementPolicy implements PlacementPolicy {
public PlacementCandidate select(PlacementContext ctx) {
String d = ctx.defaultProfile();
if (d != null && !d.isBlank() && !ctx.quarantined().contains(d) && !ctx.coolingOff().contains(d)
&& !ctx.modelOff().contains(d) && !ctx.unreachable().contains(d) && !weightExcluded(ctx, d)) {
&& !ctx.modelOff().contains(d) && !ctx.unreachable().contains(d) && !weightExcluded(ctx, d)
&& !capExcluded(ctx, d)) {
return new PlacementCandidate(d, null, 1.0f, null);
}
for (PlacementCandidate c : ctx.candidates()) {
if (!ctx.quarantined().contains(c.profile()) && !ctx.coolingOff().contains(c.profile())
&& !ctx.modelOff().contains(c.profile()) && !ctx.unreachable().contains(c.profile())
&& !c.excluded()) {
&& !c.excluded() && !PlacementPolicyUtil.atCap(ctx, c)) {
return new PlacementCandidate(c.profile(), null, c.weight(), c.maxLoad());
}
}
@@ -66,14 +75,18 @@ final class FixedPlacementPolicy implements PlacementPolicy {
// Exhaustion quarantine takes priority: reported only when quarantine is absent, so the
// message never claims "cooling off" for a profile that is really backend-exhausted.
boolean dCoolingOff = !dQuarantined && ctx.coolingOff().contains(d);
// fleetd #422: model-off is a fourth, independent source (an operator decision) — but
// quarantine/cooling-off still take priority when more than one applies, matching
// fleetd #435: at-cap sits between cooling off and model-off, matching
// CompositePeerLauncher's explicit-spawn check order (quarantine, cooling off, max load,
// then model-off) — reported only when quarantine/cooling-off are both absent.
boolean dAtCap = !dQuarantined && !dCoolingOff && capExcluded(ctx, d);
// fleetd #422: model-off is a fifth, independent source (an operator decision) — but
// quarantine/cooling-off/at-cap still take priority when more than one applies, matching
// CompositePeerLauncher's explicit-spawn check order (quarantine, cooling off, max load,
// then model-off).
boolean dModelOff = !dQuarantined && !dCoolingOff && ctx.modelOff().contains(d);
boolean dModelOff = !dQuarantined && !dCoolingOff && !dAtCap && ctx.modelOff().contains(d);
boolean dUnreachable = ctx.unreachable().contains(d);
boolean dWeightExcluded = weightExcluded(ctx, d);
if (dQuarantined || dCoolingOff || dModelOff || dUnreachable || dWeightExcluded) {
if (dQuarantined || dCoolingOff || dAtCap || dModelOff || dUnreachable || dWeightExcluded) {
List<String> reasons = new ArrayList<>();
if (dQuarantined) {
reasons.add("is quarantined (backend exhausted)");
@@ -81,6 +94,11 @@ final class FixedPlacementPolicy implements PlacementPolicy {
if (dCoolingOff) {
reasons.add("is cooling off after repeated backend errors");
}
if (dAtCap) {
PlacementCandidate c = candidateFor(ctx, d);
int live = ctx.liveCount().apply(d);
reasons.add("is at maxLoad (" + live + " live >= " + c.maxLoad() + " cap)");
}
if (dModelOff) {
reasons.add("names a model the operator has turned off in models.allow");
}
@@ -96,18 +114,38 @@ final class FixedPlacementPolicy implements PlacementPolicy {
}
if (!ctx.candidates().isEmpty()) {
throw new PlacementException("all worker profiles are excluded from automatic "
+ "selection (quarantined, cooling off, model-off, unreachable, or weight-0)");
+ "selection (quarantined, cooling off, at cap, model-off, unreachable, or weight-0)");
}
throw new PlacementException("no worker profiles configured");
}
/** Whether {@code profile} carries {@code weight <= 0} (CB-554) among {@code ctx}'s candidates. */
private static boolean weightExcluded(PlacementContext ctx, String profile) {
PlacementCandidate c = candidateFor(ctx, profile);
return c != null && c.excluded();
}
/**
* Whether {@code profile} has reached its {@code maxLoad} cap (fleetd #435), using the shared
* {@link PlacementPolicyUtil#atCap} definition — the same one {@code weighted}/{@code
* round-robin} already consult via {@link PlacementPolicyUtil#available}. Looked up by name,
* the same way {@link #weightExcluded} is: the default fast path above builds its own {@link
* PlacementCandidate} with {@code maxLoad} forced to {@code null} (it carries no cap of its
* own), so the candidate actually configured for {@code profile} has to be found in {@code
* ctx.candidates()} first.
*/
private static boolean capExcluded(PlacementContext ctx, String profile) {
PlacementCandidate c = candidateFor(ctx, profile);
return c != null && PlacementPolicyUtil.atCap(ctx, c);
}
/** The configured candidate named {@code profile} in {@code ctx}, or {@code null} if none. */
private static PlacementCandidate candidateFor(PlacementContext ctx, String profile) {
for (PlacementCandidate c : ctx.candidates()) {
if (c.profile().equals(profile)) {
return c.excluded();
return c;
}
}
return false;
return null;
}
}
@@ -11,14 +11,29 @@ final class PlacementPolicyUtil {
private PlacementPolicyUtil() {
}
/**
* True when {@code c} has reached its {@code maxLoad} cap: {@code liveCount(c.profile()) >=
* c.maxLoad()}. A {@code null} maxLoad means unlimited, so it is never at cap.
*
* <p>Extracted as the single shared definition of "at cap" (fleetd #435): before this fix it
* was computed inline in both {@link #available} and {@link #emptyException}, and {@code
* FixedPlacementPolicy} — not a caller of either — quietly kept its own {@code select} free of
* any cap check at all, so a capped default profile was chosen anyway under the default
* placement policy. Every automatic policy must call this, not re-derive it.
*/
static boolean atCap(PlacementContext ctx, PlacementCandidate c) {
Integer cap = c.maxLoad();
return cap != null && ctx.liveCount().apply(c.profile()) >= cap;
}
/**
* Candidates that are not weight-excluded (CB-554: explicit {@code weight <= 0}, checked
* first because it is a static config choice rather than transient state), not
* known-unreachable, not quarantined (CB-578 stage B), not cooling off after repeated backend
* errors (fleetd #201 Unit 5 — a separate, shorter-lived source from quarantine), not naming a
* model the operator has turned off (fleetd #422 — a third, independent source: an operator
* decision, never a backend-reported outage), and have not reached their maxLoad. A {@code
* null} maxLoad means unlimited.
* decision, never a backend-reported outage), and have not reached their maxLoad (see {@link
* #atCap}). A {@code null} maxLoad means unlimited.
*/
static List<PlacementCandidate> available(PlacementContext ctx) {
List<PlacementCandidate> out = new ArrayList<>();
@@ -26,16 +41,10 @@ final class PlacementPolicyUtil {
if (c.excluded() || ctx.unreachable().contains(c.profile())
|| ctx.quarantined().contains(c.profile())
|| ctx.coolingOff().contains(c.profile())
|| ctx.modelOff().contains(c.profile())) {
|| ctx.modelOff().contains(c.profile())
|| atCap(ctx, c)) {
continue;
}
Integer cap = c.maxLoad();
if (cap != null) {
int live = ctx.liveCount().apply(c.profile());
if (live >= cap) {
continue;
}
}
out.add(c);
}
return out;
@@ -59,7 +68,6 @@ final class PlacementPolicyUtil {
int coolingOff = 0;
int modelOff = 0;
for (PlacementCandidate c : ctx.candidates()) {
Integer cap = c.maxLoad();
if (c.excluded()) {
weightExcluded++;
} else if (ctx.quarantined().contains(c.profile())) {
@@ -70,7 +78,7 @@ final class PlacementPolicyUtil {
modelOff++;
} else if (ctx.unreachable().contains(c.profile())) {
unreachable++;
} else if (cap != null && ctx.liveCount().apply(c.profile()) >= cap) {
} else if (atCap(ctx, c)) {
atCap++;
}
}
@@ -0,0 +1,62 @@
package dev.ltms.fleet;
import dev.ltms.fleet.peer.PeerLauncher;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import java.util.Set;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
/**
* fleetd #422 follow-up: {@code Fleetd.modelGateCoverageLine} is the startup-log counterpart of
* {@code exhaustedPatternCoverageLine}/{@code errorPatternCoverageLine} — see {@code
* FleetdPatternCoverageLineTest} for the identical shape this follows — except here there is no
* {@code UnsetMeaning} choice for a caller to get backwards: {@link
* PeerLauncher.ModelGateState#configured()} already states, unambiguously, whether an empty
* {@link PeerLauncher.ModelGateState#off()} means "no {@code models:} block to gate with at all"
* or "a block armed and currently reporting zero off". This class proves {@code
* modelGateCoverageLine} words those two states — plus the third, N off — distinctly, so a
* mutation that made it ignore {@code configured()} either way is caught here.
*/
class FleetdModelGateCoverageLineTest {
@Test
@DisplayName("no models: block reports not configured, distinct from armed-with-zero")
void noModelsBlockReportsNotConfigured() {
String line = Fleetd.modelGateCoverageLine(PeerLauncher.ModelGateState.notConfigured());
assertEquals("not configured (no models: block — nothing is gated, and nothing can be)", line);
}
@Test
@DisplayName("a models: block armed with nothing off reports armed, distinct from not configured")
void armedWithNothingOffReportsArmed() {
String line = Fleetd.modelGateCoverageLine(PeerLauncher.ModelGateState.armed(Set.of()));
assertEquals("armed (models: block present; 0 models currently turned off)", line);
}
@Test
@DisplayName("a models: block with N off names the off models")
void armedWithModelsOffNamesThem() {
String line = Fleetd.modelGateCoverageLine(
PeerLauncher.ModelGateState.armed(Set.of("deepseek-v4-flash", "claude-opus-9000")));
assertEquals("armed (2 model(s) turned off: [claude-opus-9000, deepseek-v4-flash])", line);
}
@Test
@DisplayName("the three states produce pairwise-distinct wording for the same empty-looking input")
void theThreeStatesProduceDistinctWording() {
String notConfigured = Fleetd.modelGateCoverageLine(PeerLauncher.ModelGateState.notConfigured());
String armedZero = Fleetd.modelGateCoverageLine(PeerLauncher.ModelGateState.armed(Set.of()));
String armedOne = Fleetd.modelGateCoverageLine(PeerLauncher.ModelGateState.armed(Set.of("x")));
// Pinned individually above; restated here so this test alone still catches a regression
// even if one of the three tests above were ever deleted — the exact FleetdPatternCoverageLineTest
// pattern, adapted from "two keys" to "three states of one gate".
assertNotEquals(notConfigured, armedZero,
"collapsing 'no models: block' into 'armed, zero off' is fleetd #422 follow-up's exact defect");
assertNotEquals(armedZero, armedOne);
assertNotEquals(notConfigured, armedOne);
}
}
@@ -0,0 +1,67 @@
package dev.ltms.fleet;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import java.util.Set;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
/**
* fleetd #415 (review follow-up): {@code CompletionResolverTest} proves {@code coverage()} words
* {@code UnsetMeaning.OFF} and {@code UnsetMeaning.BUILT_IN_DEFAULT} correctly — but every one of
* those tests supplies the meaning itself. That proves the enum's wording, never that {@code
* Fleetd} pairs the right meaning with the right pattern key. That pairing is #415's actual
* defect: {@code coverage()} had no way to know what unset meant for its key, so the fix moved
* the fact to the caller — and nothing yet proved the caller states it correctly.
*
* <p><b>Measured or it didn't happen:</b> swapping the two {@code UnsetMeaning} arguments at
* {@code Fleetd}'s two coverage call sites — giving {@code exhaustedPattern} the built-in-default
* wording and {@code errorPattern} the off wording, #415's exact defect with the keys exchanged —
* compiled with 0 errors and left all 1506 existing tests green. This class exists to turn that
* swap red.
*
* <p>It calls {@link Fleetd#exhaustedPatternCoverageLine} and {@link Fleetd#errorPatternCoverageLine}
* directly rather than reading {@code Fleetd.java} as source text (the shape {@code
* FleetdCompletionResolverWiringTest} uses for a different wiring gap): those two methods are the
* extracted call sites {@code main} actually invokes, following the same {@code static} factory +
* dedicated-test pattern as {@link Fleetd#capacitySource} and {@link Fleetd#worktreeBranchLookup}.
*/
class FleetdPatternCoverageLineTest {
private static final Set<String> PROFILES = Set.of("terra", "gx10");
@Test
@DisplayName("exhaustedPatternCoverageLine says off when no profile configures exhaustedPattern")
void exhaustedPatternCoverageLineSaysOffWhenNoProfileConfiguresIt() {
assertEquals("off (no profile has an exhaustedPattern configured; profiles: [gx10, terra])",
Fleetd.exhaustedPatternCoverageLine(PROFILES, Set.of()));
}
@Test
@DisplayName("errorPatternCoverageLine says built-in default when no profile configures errorPattern")
void errorPatternCoverageLineSaysBuiltInDefaultWhenNoProfileConfiguresIt() {
assertEquals("built-in default for all profiles (no profile customises errorPattern; "
+ "profiles: [gx10, terra])",
Fleetd.errorPatternCoverageLine(PROFILES, Set.of()));
}
@Test
@DisplayName("the two keys produce different wording for the identical empty-coverage input")
void theTwoKeysProduceDifferentWordingForTheSameEmptyInput() {
String exhaustedLine = Fleetd.exhaustedPatternCoverageLine(PROFILES, Set.of());
String errorLine = Fleetd.errorPatternCoverageLine(PROFILES, Set.of());
// Pinned individually above; restated here so this test alone still catches a swap even
// if one of the two tests above were ever deleted.
assertEquals("off (no profile has an exhaustedPattern configured; profiles: [gx10, terra])",
exhaustedLine);
assertEquals("built-in default for all profiles (no profile customises errorPattern; "
+ "profiles: [gx10, terra])", errorLine);
assertNotEquals(exhaustedLine, errorLine,
"swapping which UnsetMeaning pairs with which pattern key at Fleetd's call sites "
+ "must be caught here — that pairing, not coverage()'s own wording in isolation, "
+ "is fleetd #415's actual defect");
}
}
@@ -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
@@ -0,0 +1,360 @@
package dev.ltms.fleet.auth;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.PaneLocator;
import dev.ltms.fleet.mcp.ConnectionIdentity;
import dev.ltms.fleet.peer.MemberRole;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Map;
import static org.junit.jupiter.api.Assertions.*;
/**
* fleetd #424 — revoking (or granting) an architect slot must take effect on the next spawn with
* no restart. A session already bound to a slot keeps its <em>binding</em> (the {@code
* terminalToSlot} occupancy) even after that slot drops out of config, but NOT the ARCHITECT
* <em>privilege</em> the slot used to grant — that is revoked on the bound session's very next
* request. See {@link MemberRegistry}'s class doc for the exact rule: "config governs what a
* bound slot still grants, as well as what may be bound next."
*
* <p>Every test here drives a REAL {@link ConfigRef#reload()} against a {@code @TempDir} file and
* asserts {@link ConfigRef.Outcome#applied()}, rather than comparing two frozen
* {@code MemberRegistry} instances in memory — the defect this ticket fixes is specifically that
* {@link MemberRegistry} used to ignore a live reload, so a test that never reloads cannot tell the
* fixed registry from the broken one. {@code requireSlotFor} and {@code reserve} are pinned in
* separate tests, in both directions (removed and added), so a registry that simply refuses (or
* simply allows) everything cannot pass by accident — see {@link MemberRegistryTest} for the
* registry's other invariants (bind/unbind cardinality, thread-safety), which are unaffected by
* this ticket and still exercised against the frozen constructor.
*/
class MemberRegistryLiveTest {
private static String yaml(String fleetBlock) {
return """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
sonnet:
baseUrl: http://gx00.gw:8000
model: sonnet
opus:
baseUrl: http://gx00.gw:8001
model: opus
guard:
offSubscriptionHosts:
- gx00.gw
""" + fleetBlock;
}
private static final String WITH_SONNET_SLOT = """
fleet:
architects:
designer:
profile: sonnet
""";
/** No architect pool at all — developers is unrelated dead data for this registry (#424 out of scope). */
private static final String WITHOUT_ARCHITECT_SLOTS = """
fleet:
developers:
dev1:
profile: sonnet
""";
/** Same slot name ({@code designer}) as {@link #WITH_SONNET_SLOT}, repointed to a different profile. */
private static final String WITH_OPUS_SLOT = """
fleet:
architects:
designer:
profile: opus
""";
/** Same profile ({@code sonnet}) as {@link #WITH_SONNET_SLOT}, but the pool key is renamed. */
private static final String WITH_RENAMED_SLOT = """
fleet:
architects:
architect-lead:
profile: sonnet
""";
private static ConfigRef refFor(Path f) {
return new ConfigRef(f, FleetConfig.load(f));
}
// ── requireSlotFor is live (criteria 1, 2, 4) ──────────────────────────────────────────────
@Test
void requireSlotForRefusesAProfileWhoseSlotWasRemovedByReload(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml(WITH_SONNET_SLOT));
ConfigRef ref = refFor(f);
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
assertDoesNotThrow(() -> registry.requireSlotFor(MemberRole.ARCHITECT, "sonnet"),
"the slot is configured before the reload");
Files.writeString(f, yaml(WITHOUT_ARCHITECT_SLOTS));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
assertThrows(IllegalArgumentException.class,
() -> registry.requireSlotFor(MemberRole.ARCHITECT, "sonnet"),
"revoking the slot must refuse the NEXT spawn that names it");
}
@Test
void requireSlotForAllowsAProfileWhoseSlotWasAddedByReload(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml(WITHOUT_ARCHITECT_SLOTS));
ConfigRef ref = refFor(f);
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
assertThrows(IllegalArgumentException.class,
() -> registry.requireSlotFor(MemberRole.ARCHITECT, "sonnet"),
"no architect slot is configured yet");
Files.writeString(f, yaml(WITH_SONNET_SLOT));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
assertDoesNotThrow(() -> registry.requireSlotFor(MemberRole.ARCHITECT, "sonnet"),
"a slot added by reload must be usable with no restart");
}
// ── reserve is live too — tested separately from requireSlotFor (criteria 1, 2, 4) ────────
@Test
void reserveRefusesAProfileWhoseSlotWasRemovedByReload(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml(WITH_SONNET_SLOT));
ConfigRef ref = refFor(f);
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
MemberLifecycle.SlotReservation before = registry.reserve(MemberRole.ARCHITECT, "sonnet");
assertEquals("architect:designer", before.slot());
registry.release(before); // free it back up so the reload-side reserve below starts clean
Files.writeString(f, yaml(WITHOUT_ARCHITECT_SLOTS));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
assertThrows(IllegalArgumentException.class,
() -> registry.reserve(MemberRole.ARCHITECT, "sonnet"),
"revoking the slot must refuse the NEXT reservation for it");
}
@Test
void reserveAllowsAProfileWhoseSlotWasAddedByReload(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml(WITHOUT_ARCHITECT_SLOTS));
ConfigRef ref = refFor(f);
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
assertThrows(IllegalArgumentException.class,
() -> registry.reserve(MemberRole.ARCHITECT, "sonnet"),
"no architect slot is configured yet");
Files.writeString(f, yaml(WITH_SONNET_SLOT));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
MemberLifecycle.SlotReservation after = registry.reserve(MemberRole.ARCHITECT, "sonnet");
assertEquals("architect:designer", after.slot(),
"a slot added by reload must be reservable with no restart");
}
// ── profileForSlot, isSlot and nameForSlot are live too (fleetd #431) ─────────────────────────
// #424 pinned roleForSlot against a live reload but left these three untested — proved by
// mutating each to read a snapshot flattened once at construction: the full suite stayed green
// for all three.
//
// The three differ in how much production behaviour depends on them, and the ticket first got
// this ranking wrong. nameForSlot is wired: CallerResolver passes members::nameForSlot, next to
// members::roleForSlot. isSlot is reached through bind, which calls it to refuse an unknown
// slot. profileForSlot has NO caller in src/main at all — grepped both the ".profileForSlot("
// and the "::profileForSlot" form — so there is no seam to drive its test through beyond the
// accessor itself, and its own javadoc ("what the spawn lifecycle reads") describes a caller
// that does not exist. These tests pin the accessors as they are; whether profileForSlot should
// be wired or deleted is a separate question.
@Test
void profileForSlotReflectsAProfileChangedByReload(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml(WITH_SONNET_SLOT));
ConfigRef ref = refFor(f);
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
assertEquals("sonnet", registry.profileForSlot("architect:designer"),
"the slot's profile before the reload");
Files.writeString(f, yaml(WITH_OPUS_SLOT));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
assertEquals("opus", registry.profileForSlot("architect:designer"),
"repointing the slot to a different profile must take effect with no restart");
}
@Test
void isSlotStopsReportingASlotRemovedByReload(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml(WITH_SONNET_SLOT));
ConfigRef ref = refFor(f);
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
assertTrue(registry.isSlot("architect:designer"), "the slot is configured before the reload");
Files.writeString(f, yaml(WITHOUT_ARCHITECT_SLOTS));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
assertFalse(registry.isSlot("architect:designer"),
"removing the slot from config must make isSlot say so on the very next call, "
+ "with no restart");
}
@Test
void isSlotStartsReportingASlotAddedByReload(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml(WITHOUT_ARCHITECT_SLOTS));
ConfigRef ref = refFor(f);
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
assertFalse(registry.isSlot("architect:designer"), "no architect slot is configured yet");
Files.writeString(f, yaml(WITH_SONNET_SLOT));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
assertTrue(registry.isSlot("architect:designer"),
"a slot added by reload must be visible to isSlot with no restart");
}
@Test
void nameForSlotReflectsANameChangedByReload(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml(WITH_SONNET_SLOT));
ConfigRef ref = refFor(f);
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
assertEquals("designer", registry.nameForSlot("architect:designer"),
"the configured name before the reload");
Files.writeString(f, yaml(WITH_RENAMED_SLOT));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
assertNull(registry.nameForSlot("architect:designer"),
"the old key no longer names a configured slot — it was renamed away by the reload");
assertEquals("architect-lead", registry.nameForSlot("architect:architect-lead"),
"the new name must be visible under its new qualified key with no restart");
}
// ── a bound architect is demoted, but the binding itself is not touched (fleetd #424) ───────
// The lead's corrected ruling: the PRIVILEGE a slot grants is revoked on the bound session's
// very next request, but the terminalToSlot BINDING itself is untouched by a reload — dropping
// it would double-book the slot key and break unbind's compare-safe contract. See the class
// doc's binding rule.
/** A caller identity resolving the one canned pane (terminal {@code term_a}) in {@link FakeHerdr}. */
private static ConnectionIdentity boundPaneIdentity() {
return new ConnectionIdentity(new PaneLocator(new FakeHerdr()), _ -> FakeHerdr.WORKER_PID);
}
@Test
void anArchitectAlreadyBoundToASlotIsDemotedByReload(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml(WITH_SONNET_SLOT));
ConfigRef ref = refFor(f);
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
MemberLifecycle.SlotReservation reservation = registry.reserve(MemberRole.ARCHITECT, "sonnet");
assertTrue(registry.bind(reservation, "term_a"));
assertEquals("architect:designer", registry.slotForTerminal("term_a"));
// Drive the real caller path, not the roleForSlot seam directly: CallerResolver.resolve is
// what a live request actually goes through (CallerResolver.java:220), and a resolver that
// ignored roleForSlot entirely would still pass a test that only checked the seam.
CallerResolver resolver = CallerResolver.withLeadsAndMembers(
boundPaneIdentity(), false, null, Map::of, registry);
Principal before = resolver.resolve("127.0.0.1", 42, null);
assertEquals(Role.ARCHITECT, before.role(), "sanity check: the harness binds term_a as an architect");
assertEquals("designer", before.name());
Files.writeString(f, yaml(WITHOUT_ARCHITECT_SLOTS));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
Principal after = resolver.resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, after.role(),
"removing the slot from config must demote the bound session to worker on its "
+ "NEXT request — this is the ticket's whole point");
assertEquals("term_a", after.terminal(), "same pane, same terminal — only the role changed");
}
@Test
void theOriginalBindingStillOccupiesTheRemovedSlotSoASecondTerminalCannotClaimIt(@TempDir Path dir)
throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml(WITH_SONNET_SLOT));
ConfigRef ref = refFor(f);
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
MemberLifecycle.SlotReservation reservation = registry.reserve(MemberRole.ARCHITECT, "sonnet");
assertTrue(registry.bind(reservation, "term_a"));
Files.writeString(f, yaml(WITHOUT_ARCHITECT_SLOTS));
ConfigRef.Outcome removed = ref.reload();
assertTrue(removed.applied(), "the reload must actually take effect: " + removed.summary());
// The binding survives the removal untouched.
assertEquals("architect:designer", registry.slotForTerminal("term_a"),
"a live binding must never be retroactively unbound by a config edit");
assertEquals(Map.of("term_a", "architect:designer"), registry.snapshot());
// Bring the slot back into config. If the binding had been silently dropped by the removal
// (rather than merely losing the privilege it grants), a second terminal could now claim
// the "freed" key — the exact double-booking the class doc's binding rule rules out.
Files.writeString(f, yaml(WITH_SONNET_SLOT));
ConfigRef.Outcome restored = ref.reload();
assertTrue(restored.applied(), "the reload must actually take effect: " + restored.summary());
assertFalse(registry.bind("architect:designer", "term_b"),
"the slot is still occupied by term_a — a second terminal must not bind to it");
assertThrows(IllegalArgumentException.class,
() -> registry.reserve(MemberRole.ARCHITECT, "sonnet"),
"the slot is still occupied by term_a — a fresh reservation must not find it free");
assertEquals("architect:designer", registry.slotForTerminal("term_a"),
"the original binding is unchanged throughout");
}
@Test
void unbindStillSucceedsForTheOriginalTerminalAfterItsSlotIsRemoved(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml(WITH_SONNET_SLOT));
ConfigRef ref = refFor(f);
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
MemberLifecycle.SlotReservation reservation = registry.reserve(MemberRole.ARCHITECT, "sonnet");
assertTrue(registry.bind(reservation, "term_a"));
Files.writeString(f, yaml(WITHOUT_ARCHITECT_SLOTS));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
assertTrue(registry.unbind("architect:designer", "term_a"),
"unbind must still work for a slot that config has since removed, or a session "
+ "that outlives its slot's removal could never release it");
assertNull(registry.slotForTerminal("term_a"));
assertEquals(Map.of(), registry.snapshot());
}
}
@@ -20,6 +20,7 @@ import java.util.regex.Pattern;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/** Unit behaviour of the CB-106 completion resolver in isolation from the injector. */
@@ -884,19 +885,22 @@ class CompletionResolverTest {
@Test
void coverageIsOffWhenNoProfileHasAPatternConfigured() {
assertEquals("off (no profile has an exhaustedPattern configured; profiles: [terra])",
CompletionResolver.coverage("exhaustedPattern", Set.of("terra"), Set.of()));
CompletionResolver.coverage("exhaustedPattern", CompletionResolver.UnsetMeaning.OFF,
Set.of("terra"), Set.of()));
}
@Test
void coverageIsFullWhenEveryProfileHasAPatternConfigured() {
assertEquals("full (all profiles configured: [gx10, terra])",
CompletionResolver.coverage("exhaustedPattern", Set.of("terra", "gx10"), Set.of("terra", "gx10")));
CompletionResolver.coverage("exhaustedPattern", CompletionResolver.UnsetMeaning.OFF,
Set.of("terra", "gx10"), Set.of("terra", "gx10")));
}
@Test
void coverageIsPartialAndNamesWhichProfilesAreConfigured() {
assertEquals("partial (configured: [terra]; not configured: [gx10])",
CompletionResolver.coverage("exhaustedPattern", Set.of("terra", "gx10"), Set.of("terra")));
CompletionResolver.coverage("exhaustedPattern", CompletionResolver.UnsetMeaning.OFF,
Set.of("terra", "gx10"), Set.of("terra")));
}
/**
@@ -910,11 +914,42 @@ class CompletionResolverTest {
*
* <p>Every earlier test here passed the exhaustion case only, so none of them could see it. This
* one pins that the message names the key the caller actually meant.
*
* <p>fleetd#415: the expected wording changed here too. {@code errorPattern} has a built-in
* fallback ({@link CompletionResolver#BACKEND_ERROR}), so an empty {@code configuredProfiles}
* for it is not "off" — see {@link #coverageDistinguishesOffFromBuiltInDefaultForTheSameEmptyInput}
* for the test built specifically to pin that distinction.
*/
@Test
void coverageNamesTheConfigKeyItsCallerMeansRatherThanAlwaysSayingExhaustedPattern() {
assertEquals("off (no profile has an errorPattern configured; profiles: [gx10, terra])",
CompletionResolver.coverage("errorPattern", Set.of("terra", "gx10"), Set.of()));
assertEquals("built-in default for all profiles (no profile customises errorPattern; "
+ "profiles: [gx10, terra])",
CompletionResolver.coverage("errorPattern", CompletionResolver.UnsetMeaning.BUILT_IN_DEFAULT,
Set.of("terra", "gx10"), Set.of()));
}
/**
* fleetd#415: {@code coverage()} measures pattern coverage (how many profiles set the key), but
* for {@code errorPattern} the empty case is not the feature-off state — a profile with no
* configured {@code errorPattern} still runs the classification against
* {@link CompletionResolver#BACKEND_ERROR}. For {@code exhaustedPattern} there is no fallback,
* so empty really is off. Same shape of input (empty {@code configuredProfiles}, one profile),
* different {@link CompletionResolver.UnsetMeaning} — the wording must differ, or this method is
* back to conflating pattern coverage with feature state for the one key where they disagree.
*/
@Test
void coverageDistinguishesOffFromBuiltInDefaultForTheSameEmptyInput() {
String exhaustedLine = CompletionResolver.coverage("exhaustedPattern",
CompletionResolver.UnsetMeaning.OFF, Set.of("gx10", "terra"), Set.of());
String errorLine = CompletionResolver.coverage("errorPattern",
CompletionResolver.UnsetMeaning.BUILT_IN_DEFAULT, Set.of("gx10", "terra"), Set.of());
assertEquals("off (no profile has an exhaustedPattern configured; profiles: [gx10, terra])",
exhaustedLine);
assertEquals("built-in default for all profiles (no profile customises errorPattern; "
+ "profiles: [gx10, terra])", errorLine);
assertNotEquals(exhaustedLine, errorLine,
"the same empty-coverage input must not read as the same feature state for both keys");
}
// --- fleetd#201 Unit 1: target-keyed backend-error pattern + typed sink ----------------------
@@ -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,109 @@
package dev.ltms.fleet.mcp;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.member.CompositePeerLauncher;
import dev.ltms.fleet.member.HerdrPeerLauncher;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementPolicies;
import org.junit.jupiter.api.Test;
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.assertFalse;
/**
* fleetd #422 follow-up: {@code fleet_profiles}/{@code GET /profiles} — the reporting surface a
* lead actually reads — must let it tell apart the three states {@link
* dev.ltms.fleet.peer.PeerLauncher#disabledModels()} alone collapses into one empty set: no
* {@code models:} block at all, a block armed with nothing currently off, and a block with N
* models off. See {@link dev.ltms.fleet.peer.PeerLauncher.ModelGateState}'s javadoc for why a
* bare {@code disabledModels()} read cannot make this distinction, and {@link
* FleetMcp#profilesView} for where {@code modelGateArmed} is added alongside the existing {@code
* modelsOff} key.
*
* <p>Every assertion here goes through {@link FleetMcp#profilesView}, never {@code
* PeerLauncher.modelGateState()} directly — {@code CompositePeerLauncherTest} already proves the
* accessor itself; this class proves the surface a lead reads (fleet_profiles / GET /profiles)
* renders what that accessor reports.
*/
class FleetProfilesModelGateStateTest {
private static FleetConfig.Profile profile(String name, String model) {
return new FleetConfig.Profile(name, "http://gx00.gw:8000", model, null, "FLEETD_WORKER_TOKEN",
null, "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
}
private static FleetMcp.QuarantineSource noQuarantine() {
return new FleetMcp.QuarantineSource(_ -> null, BackendQuarantine.none(), _ -> false);
}
private static HerdrPeerLauncher claudeAdapter(FakeHerdr h, Map<String, FleetConfig.Profile> profiles) {
return new ClaudeCodeLauncher(new AgentControl(h), new WorkspaceControl(h),
new SubscriptionGuard(Set.of("gx00.gw")), profiles, "local", _ -> "tok");
}
/**
* State 1: no {@code models:} block at all — a plain {@code ClaudeCodeLauncher} (no {@code
* models:} supplier exists for it to read) has nothing to gate against, matching the fleet01
* host measured for this ticket: {@code grep -c '^models:' fleetd.yaml} returns 0 there.
*/
@Test
void noModelsBlockReportsGateNotArmed() {
FakeHerdr h = new FakeHerdr();
Map<String, FleetConfig.Profile> profiles = Map.of("local", profile("local", "deepseek-v4-flash"));
PeerLauncher workers = claudeAdapter(h, profiles);
Map<String, Object> view = FleetMcp.profilesView(workers, noQuarantine(), FleetMcp.OutageSource.none());
assertEquals(Boolean.FALSE, view.get("modelGateArmed"),
"no models: block to read from — nothing is gated, and nothing can be");
assertFalse(view.containsKey("modelsOff"), "nothing configured, so no off set to report either");
}
/** State 2: a {@code models:} block is present, but nothing in it is currently turned off. */
@Test
void modelsBlockWithNothingOffReportsGateArmedAndZeroOff() {
FakeHerdr h = new FakeHerdr();
Map<String, FleetConfig.Profile> profiles = Map.of("local", profile("local", "deepseek-v4-flash"));
FleetConfig.Models models = new FleetConfig.Models(
List.of(new FleetConfig.Models.ModelEntry("deepseek-v4-flash", true)));
PeerLauncher workers = new CompositePeerLauncher(List.of(claudeAdapter(h, profiles)), "local", profiles,
PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(),
new BackendOutagePolicy(System::nanoTime), () -> models);
Map<String, Object> view = FleetMcp.profilesView(workers, noQuarantine(), FleetMcp.OutageSource.none());
assertEquals(Boolean.TRUE, view.get("modelGateArmed"),
"a models: block is present, so the gate is armed even though nothing is off yet");
assertFalse(view.containsKey("modelsOff"),
"nothing is off, so the key stays absent — an empty list here would be indistinguishable "
+ "from today's modelsOff omission, exactly the ambiguity modelGateArmed exists to remove");
}
/** State 3: a {@code models:} block is present with one model currently turned off. */
@Test
void modelsBlockWithModelsOffReportsGateArmedAndTheOffSet() {
FakeHerdr h = new FakeHerdr();
Map<String, FleetConfig.Profile> profiles = Map.of("local", profile("local", "deepseek-v4-flash"));
FleetConfig.Models models = new FleetConfig.Models(
List.of(new FleetConfig.Models.ModelEntry("deepseek-v4-flash", false)));
PeerLauncher workers = new CompositePeerLauncher(List.of(claudeAdapter(h, profiles)), "local", profiles,
PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(),
new BackendOutagePolicy(System::nanoTime), () -> models);
Map<String, Object> view = FleetMcp.profilesView(workers, noQuarantine(), FleetMcp.OutageSource.none());
assertEquals(Boolean.TRUE, view.get("modelGateArmed"));
assertEquals(List.of("deepseek-v4-flash"), view.get("modelsOff"));
}
}
@@ -536,6 +536,31 @@ class CompositePeerLauncherTest {
assertEquals("claude", h.profile(), "the returned handle carries the resolved default profile");
}
/**
* fleetd #435: {@code FixedPlacementPolicy} — the default placement policy every config uses
* unless {@code placement:} is set — never consulted {@code maxLoad}, so an unqualified spawn
* (a blank profile, the normal delegation path) landed on a capped default anyway. Measured on
* 7667727: a single dev profile at {@code maxLoad: 1} with 1 live, under {@code fixed()},
* returned "SPAWNED on profile=a". This test goes through {@code CompositePeerLauncher.spawn}
* with a blank profile, not the policy in isolation, so it proves the caller actually reaches
* the fixed default's new cap check rather than only the {@code select} method.
*/
@Test
void fixedPolicyGatesDefaultProfileAtMaxLoadOnUnqualifiedSpawn() {
FakeHerdr herdr = new FakeHerdr();
Map<String, FleetConfig.Profile> profiles = ordered(
"a", stubWorker("a", 1.0f, 1),
"b", stubWorker("b", 1.0f, null));
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "a", Set.of());
CompositePeerLauncher composite = new CompositePeerLauncher(
List.of(adapter), "a", profiles, PlacementPolicies.fixed(), name -> "a".equals(name) ? 1 : 0);
PeerHandle h = composite.spawn(new SpawnRequest(null, null, null));
assertEquals("b", h.profile(),
"the default profile a is at maxLoad, so fixed placement must fall through to b");
assertEquals(0, adapter.spawnCount("a"), "a is never spawned — it is already at cap");
}
@Test
void weightedPolicyGatesProfileAtMaxLoad() {
FakeHerdr herdr = new FakeHerdr();
@@ -1453,6 +1478,108 @@ class CompositePeerLauncherTest {
assertEquals(2, adapter.spawnCount("local"));
}
/**
* fleetd #422 follow-up, acceptance criterion 2: {@link CompositePeerLauncher#modelGateState()}
* is LIVE — no restart — proven through a REAL {@link ConfigRef#reload()}, exactly like {@link
* #modelOnOffIsHotReloadedThroughARealConfigRef} above proves for the on/off gate itself. This
* single reload sequence walks through all three states the ticket asks for: no {@code models:}
* block, a block armed with nothing off, and a block with one model off — so a reload that flips
* between any of the three is proven live, not just the on/off edit within an already-armed block.
*/
@Test
void modelGateStateIsHotReloadedThroughARealConfigRef(@TempDir Path dir) throws Exception {
Path yaml = dir.resolve("fleetd.yaml");
Files.writeString(yaml, """
bind:
port: 8080
profiles:
local:
baseUrl: http://local.gw:8000
model: deepseek-v4-flash
""");
FleetConfig initial = FleetConfig.load(yaml);
ConfigRef configRef = new ConfigRef(yaml, initial);
FakeHerdr herdr = new FakeHerdr();
StubLauncher adapter = new StubLauncher("claude", herdr,
Map.of("local", stubWorker("local")), "local", Set.of());
CompositePeerLauncher composite = new CompositePeerLauncher(
List.of(adapter), "local", configRef, _ -> 0, BackendQuarantine.none(), NO_OUTAGE);
PeerLauncher.ModelGateState notConfigured = composite.modelGateState();
assertFalse(notConfigured.configured(), "no models: block in the config at all");
assertEquals(Set.of(), notConfigured.off());
Files.writeString(yaml, """
bind:
port: 8080
profiles:
local:
baseUrl: http://local.gw:8000
model: deepseek-v4-flash
models:
allow:
- model: deepseek-v4-flash
enabled: false
""");
assertTrue(configRef.reload().applied(), "adding a models: block must apply live, no restart");
PeerLauncher.ModelGateState armedWithOneOff = composite.modelGateState();
assertTrue(armedWithOneOff.configured(), "a models: block now exists — the gate is armed");
assertEquals(Set.of("deepseek-v4-flash"), armedWithOneOff.off());
Files.writeString(yaml, """
bind:
port: 8080
profiles:
local:
baseUrl: http://local.gw:8000
model: deepseek-v4-flash
models:
allow:
- model: deepseek-v4-flash
enabled: true
""");
assertTrue(configRef.reload().applied(), "flipping the entry back on must apply live too");
PeerLauncher.ModelGateState armedWithZeroOff = composite.modelGateState();
assertTrue(armedWithZeroOff.configured(),
"the block is still present — armed and reporting zero, not the same as no block at all");
assertEquals(Set.of(), armedWithZeroOff.off());
}
/**
* fleetd #422 follow-up, acceptance criterion 3: the invariant is that an absent {@code
* models:} block stays permitted and must never be fatal. Proved, not assumed — a config
* without one loads, validates, reports the gate as not configured, AND still spawns normally
* (no {@link PlacementException} from a gate that has nothing to check against), using the same
* production-shaped {@code Supplier<FleetConfig>} wiring {@code Fleetd.main} actually uses.
*/
@Test
void noModelsBlockConfigStillLoadsAndSpawnsNormally(@TempDir Path dir) throws Exception {
Path yaml = dir.resolve("fleetd.yaml");
Files.writeString(yaml, """
bind:
port: 8080
profiles:
local:
baseUrl: http://local.gw:8000
model: deepseek-v4-flash
""");
FleetConfig cfg = FleetConfig.load(yaml);
assertDoesNotThrow(cfg::validateAll, "a config with no models: block must load and validate cleanly");
ConfigRef configRef = new ConfigRef(yaml, cfg);
FakeHerdr herdr = new FakeHerdr();
StubLauncher adapter = new StubLauncher("claude", herdr,
Map.of("local", stubWorker("local")), "local", Set.of());
CompositePeerLauncher composite = new CompositePeerLauncher(
List.of(adapter), "local", configRef, _ -> 0, BackendQuarantine.none(), NO_OUTAGE);
assertFalse(composite.modelGateState().configured());
assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("local", null, null)),
"no models: block means nothing to gate against — the spawn must go through");
assertEquals(1, adapter.spawnCount("local"));
}
/** {@code fleet_profiles}/{@code GET /profiles} must read the exact same live source the gate reads. */
@Test
void disabledModelsReportsWhatTheGateActuallyEnforces() {
@@ -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");
@@ -461,6 +461,93 @@ class PlacementPolicyTest {
assertTrue(e.getMessage().contains("weight 0"), e.getMessage());
}
// --- fleetd #435: FixedPlacementPolicy must consult maxLoad too, at BOTH filter sites -------
/**
* The default-profile fast path must skip a capped default. Measured on 7667727 before this
* fix: a single dev profile at {@code maxLoad: 1} with 1 live, under {@code fixed()}, still
* returned "SPAWNED on profile=a" — the cap was advisory for every unqualified spawn.
*/
@Test
void fixedSkipsCappedDefault() {
PlacementPolicy policy = PlacementPolicies.fixed();
PlacementContext ctx = new PlacementContext("b",
List.of(PlacementCandidate.profile("a", 1.0f, null),
PlacementCandidate.profile("b", 1.0f, 1)),
name -> "b".equals(name) ? 1 : 0, Set.of(), Set.of(), Set.of());
assertEquals("a", policy.select(ctx).profile(),
"the default 'b' is at its maxLoad cap, so fixed falls through to the free candidate 'a'");
}
/**
* The fallback walk must skip a capped candidate too — exercised independently of the
* default-profile fast path by using no default at all, so this is the only filter that runs.
*/
@Test
void fixedFallbackWalkSkipsCappedCandidate() {
PlacementPolicy policy = PlacementPolicies.fixed();
PlacementContext ctx = new PlacementContext(null,
List.of(PlacementCandidate.profile("a", 1.0f, 1),
PlacementCandidate.profile("b", 1.0f, null)),
name -> "a".equals(name) ? 1 : 0, Set.of(), Set.of(), Set.of());
assertEquals("b", policy.select(ctx).profile(),
"candidate 'a' is at its maxLoad cap, so the fallback walk skips it and picks 'b'");
}
@Test
void fixedThrowsWhenDefaultAndEveryCandidateAtCap() {
PlacementPolicy policy = PlacementPolicies.fixed();
PlacementContext ctx = new PlacementContext("b",
List.of(PlacementCandidate.profile("a", 1.0f, 1),
PlacementCandidate.profile("b", 1.0f, 1)),
_ -> 1, Set.of(), Set.of(), Set.of());
PlacementException e = assertThrows(PlacementException.class, () -> policy.select(ctx));
assertTrue(e.getMessage().contains("maxLoad"), "message names the cap: " + e.getMessage());
assertTrue(e.getMessage().contains("1 cap"), "message names the cap value: " + e.getMessage());
assertFalse(e.getMessage().contains("quarantined"), "must not read like quarantine: " + e.getMessage());
assertFalse(e.getMessage().contains("turned off"), "must not read like model-off: " + e.getMessage());
}
/** The mirror: an uncapped default is still chosen, so the new term cannot exclude everything. */
@Test
void fixedStillReturnsUncappedDefault() {
PlacementPolicy policy = PlacementPolicies.fixed();
PlacementContext ctx = new PlacementContext("b",
List.of(PlacementCandidate.profile("a", 1.0f, null),
PlacementCandidate.profile("b", 1.0f, null)),
noSessions(), Set.of(), Set.of(), Set.of());
assertEquals("b", policy.select(ctx).profile(), "an uncapped default is returned unconditionally");
}
/** CB-585: an explicit {@code maxLoad: 0} on the default caps it at zero live members. */
@Test
void fixedSkipsMaxLoadZeroDefaultEvenWithZeroLiveWorkers() {
PlacementPolicy policy = PlacementPolicies.fixed();
PlacementContext ctx = new PlacementContext("b",
List.of(PlacementCandidate.profile("a", 1.0f, null),
PlacementCandidate.profile("b", 1.0f, 0)),
_ -> 0, Set.of(), Set.of(), Set.of());
assertEquals("a", policy.select(ctx).profile(),
"the default 'b' has maxLoad 0, so it is already at its cap with nobody live");
}
/**
* Quarantine still wins when a profile is both quarantined and at cap (mirrors {@code
* fixedReportsQuarantineNotModelOffWhenBothApply}'s priority over model-off).
*/
@Test
void fixedReportsQuarantineNotAtCapWhenBothApply() {
PlacementPolicy policy = PlacementPolicies.fixed();
PlacementContext ctx = new PlacementContext("b",
List.of(PlacementCandidate.profile("a", 1.0f, 1),
PlacementCandidate.profile("b", 1.0f, 1)),
_ -> 1, Set.of(), Set.of("a", "b"), Set.of());
PlacementException e = assertThrows(PlacementException.class, () -> policy.select(ctx));
assertTrue(e.getMessage().contains("quarantined"), e.getMessage());
assertFalse(e.getMessage().contains("maxLoad"),
"quarantine takes priority over at-cap in the message: " + e.getMessage());
}
@Test
void unknownPolicyNameThrows() {
assertThrows(IllegalArgumentException.class, () -> PlacementPolicies.fromName("random"));