Compare commits

..

10 Commits

Author SHA1 Message Date
Dai Ha 7772b41993 fleetd #446 round 3: extract exhaustionSink and pin its caller (Cell A/B)
CI / contract (pull_request) Successful in 55s
CI / build (pull_request) Successful in 2m13s
Round 2 pinned usageLimitFixWarning/usageLimitFixWarningNoModel's TEXT via
FleetdUsageLimitFixWarningTest, but a mutation battery against the merged PR
proved two gaps in the caller that builds main()'s real quarantine
ExhaustionSink: nothing proved the sink's log.warn actually invokes either
method (Cell A), and nothing proved it picks the right one for a profile
with vs without a configured model: (Cell B).

Extract the inline lambda into a new static Fleetd.exhaustionSink(...)
factory (same refactor-for-testability class the lead approved in round 2
for the two static warning methods), behaviourally unchanged from the
lambda it replaces. FleetdExhaustionSinkWarningTest drives this factory's
return value directly and asserts on the real text a ListAppender attached
to Fleetd's own logger captures, covering both cells from one mechanism:

- Cell A (ternary result replaced by a literal string): confirmed red,
  1594 run / 1 failure, restored, confirmed green (1594/0).
- Cell B (ternary's two branches swapped): confirmed red, 1594 run /
  2 failures (both test methods independently caught it), restored,
  confirmed green (1594/0).
2026-09-10 19:00:07 +07:00
Dai Ha 30d6872779 fleetd #446 follow-up: pin the WARNING text and the fleet_profiles model/reason fields
CI / contract (pull_request) Successful in 1m14s
CI / build (pull_request) Successful in 2m4s
A mutation battery run against merged PR #457 (d02dd1b) proved criteria 2 and 3
shipped without a test that could catch them breaking: deleting either
row.put("model", model) or row.put("reason", reason) in FleetMcp.profilesView
left all 1586 tests green (rc=0), and renaming the "usage-limit fix:" log tag
to something meaningless did too. Criterion 1's own mutation (LiveExhaustedPatterns
snapshotting instead of reading live) was correctly killed by the existing
FleetdExhaustionDetectionArmedWiringTest/LiveExhaustedPatternsTest — only 2 and 3
were unguarded.

- Extracted the exhaustionSink WARNING text out of two inline SLF4J {}-placeholder
  log.warn calls into two static methods, Fleetd.usageLimitFixWarning(profile, model)
  and Fleetd.usageLimitFixWarningNoModel(profile, quarantineCooldownSeconds), the
  same extracted-static-method + dedicated-test idiom as
  exhaustedPatternCoverageLine/errorPatternCoverageLine. log.warn is now called with
  each method's return value as a single already-formatted argument, so the string a
  test asserts on is byte-identical to what fleetd.out receives. No behaviour change:
  same text, same two branches, same call site.
- New FleetdUsageLimitFixWarningTest pins the leading "usage-limit fix:" grep tag and
  three specific facts (profile name, model name, "enabled: false" under
  models.allow) rather than the whole sentence, plus the no-model fallback's profile
  name and cooldown-seconds substitution and its explicit absence of "enabled: false".
- New FleetProfilesQuarantineModelReasonFieldsTest exercises FleetMcp.QuarantineSource
  with a modelFor/reasonFor that actually return values (every existing test used
  QuarantineSource.none() or a 3-arg form defaulting both to null), asserting the
  quarantined row's model/reason are present when supplied and absent (not null, not
  blank) when modelFor/reasonFor return null or a blank string.
- Each of the three: implemented, broken by hand (row.put deleted / tag renamed),
  confirmed the new test goes red, restored, confirmed green again. See PR body and
  this ticket's fleet_reply for the verbatim failure output of all three.
2026-09-10 18:22:51 +07:00
Dai Ha 45aca9eb3e fleetd #446: make exhaustedPattern hot, name the fix in the warning, report it in fleet_profiles
CI / contract (pull_request) Successful in 1m30s
CI / build (pull_request) Successful in 1m34s
The model gate can be turned off at runtime (models.allow[].enabled: false, hot
since fleetd #422), but the usage-limit detector it's meant to react to was
compiled once at Fleetd.main startup into a frozen Map<String,Pattern> — arming
or disarming exhaustedPattern needed a daemon restart. Backwards for a feature
meant to react live.

- New LiveExhaustedPatterns: reads exhaustedPattern off the live config supplier
  per lookup (matching CompositePeerLauncher#models0's live-supplier pattern),
  caching compiled Pattern objects by PROFILE NAME (not pattern text — pattern
  text would grow unboundedly as an operator tunes a regex across reloads;
  profile names are bounded by the small, restart-gated set of configured
  profiles). patternFor() backs CompletionResolver's classification; armed()
  backs fleet_profiles' exhaustionDetectionArmed — both read the same object,
  the fleetd #404 single-accessor rule CompositePeerLauncher.modelGateState()
  established for the model gate.
- Fleetd.java: on BACKEND_EXHAUSTED, log a WARNING naming the profile's model
  and the exact fix (enabled: false under models.allow, hot, no restart; remove
  it again once the window resets) — or, when the profile has no model:
  configured, say quarantine is the only thing keeping spawns off it.
- FleetMcp.QuarantineSource gains modelFor/reasonFor; fleet_profiles'
  quarantined rows gain model/reason fields so a lead can see why without
  reading the daemon log. capacityView (fleet_list) intentionally untouched —
  scoped to fleet_profiles only.
- ConfigRef/FleetConfig docs + fleetd.example.yaml updated: exhaustedPattern
  moves from Deferred to Hot. errorPattern stays deferred on purpose (out of
  scope for this ticket).
- Tests: LiveExhaustedPatternsTest (new, unit-level hotness/caching proof),
  FleetdExhaustionDetectionArmedWiringTest (rewritten — same QuarantineSource
  object, read before and after a reload, asserts the answer flips with no
  restart), ConfigRefTest/ConfigRefProfileCoverageTest updated for the new
  Hot/Deferred classification.
2026-09-10 18:02:00 +07:00
Dai Ha 20c1094cbf fleetd #449: say what the timing fix actually proved, not what it assumed
CI / contract (push) Successful in 1m22s
CI / build (push) Successful in 1m33s
The polling fix that landed in #452 is right, but its javadoc named a
mechanism nobody measured: that input typed before the shell's prompt was
swallowed by the shell's own startup.

I mutated the settle poll away — SHELL_READY_TIMEOUT_MS = 0, so input is
typed at once with no wait — and the test passed 3 of 3. So waitForText is
the load-bearing half, and the proven cause is the old 800ms READ deadline,
not the 1000ms write delay.

The direction is the point: typing at 0ms works where typing at 1000ms
failed. If early input were swallowed, 0ms would be worse than 1000ms. It is
better, so the swallow explanation is unsupported.

waitUntilSettled stays as cheap insurance, now labelled as insurance rather
than as the fix. Comment-only; AgentControlContractTest still green.
2026-09-10 17:36:22 +07:00
ltms 9011c59b9f Merge #452: run the contract tag in CI, fix the stale protocol 14 assertion (fleetd #449)
CI / build (push) Successful in 1m28s
CI / contract (push) Successful in 1m33s
Verified on a locally built merge onto c11ad71 (the PR was branched from 822327e, before #448):
1578 unit tests green, then `-Dgroups=contract` gives 30 tests / 0 failures / 0 skipped here.

CI's own run of the tag on d4f93a7: 29 tests, 0 failures, 6 skipped — all 6 herdr tests skip on
the runner, and LeadMailboxTest's 15 tests run there for the first time.

Two mutations killed: restoring `assertEquals(14` fails the test, and pointing the expected env
value at a wrong URL fails it with the real pane text in the message.
2026-09-10 12:35:10 +02:00
ltms c11ad71ed0 Merge #448: fleet_ack errors instead of claiming success on a miss (fleetd #437)
CI / contract (push) Successful in 52s
CI / build (push) Successful in 1m33s
Verified by the lead on head 5289eb5, and then again on the MERGE.

Battery on the branch tip:

  FULL BUILD  Tests run: 1577, Failures: 0, Errors: 0, Skipped: 0  BUILD SUCCESS
              compile errors: 0
  CONTROL     contract test, real broker: Tests run: 9, Failures: 0, Skipped: 0

  M3a  AmqpReplyInbox.ack's `h == null` branch reports true
       -> KILLED  AmqpReplyInboxContractTest.ackReportsHitVsMissAgainstARealBroker
  M3b  AmqpReplyInbox.ack's `perTarget == null` branch reports true
       -> KILLED  same method

M3a is the gap I found in round 1: reporting true for something never held
survived both the default suite and -Pcontract against a real broker. It is now
closed in the adapter this daemon actually runs.

Both branches die, and they die through DIFFERENT assertions, which the worker
worked out and I confirmed by reading the code. held.get(target) is populated by
the deliver callback, not by own(), so after own() with nothing ever delivered
the map entry is still null — the "never held" assertion therefore exercises
`perTarget == null`, and only the double-ack assertion reaches `h == null`. Both
assertions are load-bearing; neither is redundant.

This PR was pushed before #451 landed, so the battery above tested the branch,
not the result. Merging main (bdcf285) in gave 0 conflicts and #448 adds no
PeerLauncher implementer, but a clean auto-merge is not a compiling merge, so I
built the merged tree 0829542:

  FULL BUILD               Tests run: 1578, Failures: 0  BUILD SUCCESS, 0 compile errors
  contract, real broker    Tests run: 9,    Failures: 0
  CompositePeerLauncherTest (#447's guarantee)  Tests run: 76, Failures: 0

The worker also declined to point the contract test at the shared local LavinMQ
broker, because this adapter never deletes queues and a run would leave orphaned
durable queues on the instance backing the live fleet. It started a disposable
rabbitmq:3.13-management container on a throwaway port instead, then removed it.
That was its own judgment and it was correct.

Held peer mail is deliberately NOT ackable: fleet_ack against a coord-id errors
and names fleet_poll{coordId}. See #437 for why refusing is the right answer
today, and why that is a policy choice rather than a structural one.
2026-09-10 12:20:49 +02:00
ltms bdcf285265 Merge #451: make PeerLauncher.spawn(req, decision) abstract (fleetd #450)
CI / contract (push) Successful in 1m17s
CI / build (push) Successful in 1m33s
Verified by the lead on head cfebc57 (base 822327e is current main).

  FULL BUILD  Tests run: 1575, Failures: 0, Errors: 0, Skipped: 0  BUILD SUCCESS
              compile errors: 0

  M1  delete HerdrPeerLauncher's new override
      -> BUILD FAILURE, 1 COMPILATION ERROR block, naming exactly
         ClaudeCodeLauncher.java:[49,14] and OpenCodeLauncher.java:[59,14]

  M2  CONTROL for M1: put the interface method back to a `default` AND delete
      the override
      -> BUILD SUCCESS, 0 compile errors
      This is the row that makes M1 mean something. Without it, M1 only shows
      that the build broke; with it, the break is attributable to the method
      being abstract rather than to anything else the edit disturbed.

  M3  CompositePeerLauncher's override reverts to the re-entering form
      -> KILLED, Errors: 1
      CompositePeerLauncherTest
        .spawnHonorsAPlacementDecisionEvenAfterItsProfileIsQuarantinedInTheWindowAfterPlace
      So #447's guarantee survives this refactor of the interface it rests on.

Tree restored clean after each mutation (git status --porcelain empty).

The worker corrected my ticket, and it was right. My #450 body listed five
src/main implementers of PeerLauncher and quoted `grep -rln 'implements
PeerLauncher'` as the source; that command returns two files. I had run a wider
pattern that also matched a comment in ConfigRef and the `extends
HerdrPeerLauncher` line in two subclasses, then quoted the narrow command beside
the wide command's output. Ground truth: ConfigRef implements
Supplier<FleetConfig>; ClaudeCodeLauncher and OpenCodeLauncher extend
HerdrPeerLauncher. So one override in that parent serves both, which is what the
worker built. The ticket body is corrected.

Out-of-scope note carried forward from the worker: defaultProfileFor(MemberRole)
and place(MemberRole) are two more default methods with the same shape. Noted,
not fixed here.
2026-09-10 12:16:47 +02:00
Dai Ha cfebc575ea fleetd #450: make PeerLauncher.spawn(SpawnRequest, PlacementDecision) abstract
CI / contract (pull_request) Successful in 47s
CI / build (pull_request) Successful in 2m5s
The default re-entered the single-argument spawn(SpawnRequest), which re-runs
checks that can refuse the profile place() just chose (#444's window). Only
CompositePeerLauncher overrode it; a future placement-doing launcher could
have inherited the wrong body silently.

Give every current implementer an explicit override, chosen by what it does:
- HerdrPeerLauncher (base of ClaudeCodeLauncher/OpenCodeLauncher, neither of
  which overrides spawn(req) or place()) does no placement filtering of its
  own, so it gets the re-entering form.
- CompositePeerLauncher's routing-form override is untouched.
- 5 test-fake PeerLauncher implementers (SessionManagerTest, FleetdBackendErrorSinkTest)
  get overrides matching their existing spawn(SpawnRequest) shape: delegating
  wrappers delegate, unreachable stubs throw, the single-profile fake re-enters.

ConfigRef does not implement PeerLauncher at all (confirmed in this tree at
822327e) despite the ticket listing it as a src/main implementer.
2026-09-10 17:10:44 +07:00
Dai Ha 5289eb509f fleetd #437: pin the ack hit/miss contract in the AMQP contract test
CI / contract (pull_request) Successful in 1m19s
CI / build (pull_request) Successful in 1m34s
AmqpReplyInboxContractTest is the one contract-group class CI actually
runs, and it never asserted on ack()'s return value at all — so the
exact defect this ticket fixes (reporting success for an ack that
removed nothing) was unpinned in the adapter fleetd runs live.

Add ackReportsHitVsMissAgainstARealBroker: a msgId never held for an
owned target returns false without throwing, a real held reply returns
true and is removed, and acking the same msgId again returns false.
Ran against both broker modes the class supports: Testcontainers
(AMQP_URI unset) and an external broker via AMQP_URI (the CI shape,
using a disposable container — not the shared local LavinMQ instance).
2026-09-10 16:58:55 +07:00
Dai Ha 703a05db41 fleetd #437: fleet_ack errors instead of claiming success on a miss
CI / contract (pull_request) Successful in 1m25s
CI / build (pull_request) Successful in 1m37s
ReplyInbox.ack now returns boolean (true = removed, false = nothing to
remove) instead of void, so FleetMcp.ack can finally tell a hit from a
miss. FleetMcp.ack returns an error when the boolean is false, naming
fleet_poll{coordId} for held peer mail, which has no route through this
call. MessageService.ackReply propagates the boolean; drainReplies keeps
ignoring it (its own javadoc already documents that loss window as
deliberate). Updated the tool schema's target description to match.

Rewrote FleetMcpTest's ack tests to publish a real message before
asserting success, and added tests for a never-queued id and a coord-id
target, both now erroring. Added boolean assertions to
InMemoryReplyInboxTest's existing ack cases.
2026-09-10 16:44:08 +07:00
25 changed files with 1150 additions and 179 deletions
+17 -11
View File
@@ -207,8 +207,12 @@ herdrSocket: ~/.config/herdr/herdr.sock
# omit and this profile's completion fallback behaves exactly as before.
# Every backend words its refusal differently, so this is config, never a
# vendor string baked into fleetd itself.
# DEFERRED: compiled once into a startup pattern map — editing it needs a
# daemon restart, same as this profile's model/baseUrl/argv.
# HOT (fleetd #446): read live, cached by profile name, at every
# completion-fallback check AND by fleet_profiles' exhaustionDetectionArmed —
# editing it and reloading arms or disarms usage-limit detection for this
# profile with no daemon restart. (Before fleetd #446 this was DEFERRED,
# compiled once into a startup pattern map like model/baseUrl/argv still are —
# see errorPattern below, which is still deferred that way on purpose.)
# credentialId → CB-578 stage B: the credential this profile quarantines WITH when a
# BACKEND_EXHAUSTED classification fires. Two profiles that set the SAME
# credentialId share one quarantine — the case this exists for is two models
@@ -227,8 +231,9 @@ herdrSocket: ~/.config/herdr/herdr.sock
# happens, just without a profile-specific match; every backend words its
# failure differently, so a hardcoded sentence would only ever match one
# of them.
# DEFERRED: compiled once into a startup pattern map, same as exhaustedPattern
# — editing it needs a daemon restart.
# DEFERRED: compiled once into a startup pattern map — editing it needs a
# daemon restart. Unlike exhaustedPattern above (made hot by fleetd #446),
# errorPattern was scoped out of that ticket on purpose and stays deferred.
# # errorPattern: "503 Service Unavailable" # opt-in: classify a backend outage
#
# What happens once a match fires (BackendOutagePolicy, credentialId-keyed,
@@ -468,9 +473,10 @@ placement: weighted
# when the reload happens — not about how important the key is:
# HOT → takes effect on the next spawn: the whole `fleet:` block (every role pool,
# `charters`, and `tabLabel`), `placement:`, and an existing profile's weight / maxLoad
# / credentialId. Those are hot because the placement policy (and, for credentialId,
# the CB-578 stage B quarantine check) reads them through a supplier — being config is
# not by itself enough to make a key hot.
# / credentialId / exhaustedPattern. Those are hot because the placement policy (and,
# for credentialId, the CB-578 stage B quarantine check; for exhaustedPattern, fleetd
# #446's LiveExhaustedPatterns) reads them through a supplier — being config is not by
# itself enough to make a key hot.
# EXCEPT `fleet.leaders`: Fleetd.main reads it once at startup to build the lead tab
# scanner and launcher, and neither is rebuilt on reload. A changed/added/removed
# `fleet.leaders` entry is silently accepted — the reload reports "config reloaded"
@@ -482,10 +488,10 @@ placement: weighted
# stage B — baked once into the quarantine tracker built at startup), ADDING or
# REMOVING a profile (a new backend needs its own launcher, and launchers are built
# once), AND an existing profile's launch settings — model, baseUrl, argv, env,
# configDir, mcpUrl, tabLabel, exhaustedPattern, errorPattern (fleetd #201 / #227 —
# compiled once into a startup pattern map the same way exhaustedPattern is). The
# launcher takes a copy of `profiles:` at startup and resolves every spawn out of
# that copy, so those never
# configDir, mcpUrl, tabLabel, errorPattern (fleetd #201 / #227 — compiled once into
# a startup pattern map; exhaustedPattern used to be compiled the same way until
# fleetd #446 made it hot — see above). The launcher takes a copy of `profiles:` at
# startup and resolves every spawn out of that copy, so those never
# reach a launch until you restart. The reload logs them by name rather than
# pretending they applied.
# COLD → cannot change at all: `bind:`, `herdrSocket:`, `broker:` and `auth:`. The socket is
+162 -61
View File
@@ -18,6 +18,7 @@ import dev.ltms.fleet.inject.BackendErrorSink;
import dev.ltms.fleet.inject.CompletionResolver;
import dev.ltms.fleet.inject.ExhaustedPatternLookup;
import dev.ltms.fleet.inject.ExhaustionSink;
import dev.ltms.fleet.inject.LiveExhaustedPatterns;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.inject.StatusPoller;
import dev.ltms.fleet.inject.TurnListener;
@@ -65,11 +66,13 @@ import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.TreeSet;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
@@ -78,6 +81,7 @@ import java.util.function.Function;
import java.util.function.Predicate;
import java.util.function.Supplier;
import java.util.regex.Pattern;
import java.util.stream.Collectors;
/**
* {@code fleetd} entry point. Wires the real herdr socket client to the REST app and
@@ -378,28 +382,34 @@ public final class Fleetd {
Rendezvous rendezvous = new Rendezvous();
// CB-578 stage A: classify a completion-fallback scrape that matches a profile's configured
// usage-limit refusal as BACKEND_EXHAUSTED rather than handing it back as a real answer.
// Compiled once at startup, keyed by profile name; a profile with no exhaustedPattern is
// simply absent here, so its workers keep today's completion-fallback behaviour unchanged.
Map<String, Pattern> exhaustedPatternsByProfile = new LinkedHashMap<>();
cfg.profiles().forEach((name, profile) -> {
if (profile.hasExhaustedPattern()) {
exhaustedPatternsByProfile.put(name, Pattern.compile(profile.exhaustedPattern()));
}
});
// fleetd #446: read live off `config` per lookup, cached by profile name — see
// LiveExhaustedPatterns's class doc for why this replaces the old compiled-once-at-startup
// map. A profile with no exhaustedPattern simply returns null here, so its workers keep
// today's completion-fallback behaviour unchanged.
LiveExhaustedPatterns liveExhaustedPatterns = new LiveExhaustedPatterns(() -> config.get().profiles());
ExhaustedPatternLookup exhaustedPatterns = target -> sessions.roster().stream()
.filter(session -> target.equals(session.terminalId()))
.findFirst()
.map(session -> exhaustedPatternsByProfile.get(session.profile()))
.map(session -> liveExhaustedPatterns.patternFor(session.profile()))
.orElse(null);
// The startup coverage line still reports the boot-time snapshot only — it is printed once,
// here, and a reload no longer needs to change what it said; exhaustionDetectionArmed (via
// liveExhaustedPatterns.armed, wired into quarantineSource below) is what stays live.
Set<String> exhaustedConfiguredAtStartup = cfg.profiles().entrySet().stream()
.filter(e -> e.getValue().hasExhaustedPattern())
.map(Map.Entry::getKey)
.collect(Collectors.toCollection(LinkedHashSet::new));
log.info("backend-exhausted classification (CB-578 stage A): {}",
exhaustedPatternCoverageLine(cfg.profiles().keySet(), exhaustedPatternsByProfile.keySet()));
exhaustedPatternCoverageLine(cfg.profiles().keySet(), exhaustedConfiguredAtStartup));
// 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
// name, mirroring exhaustedPatternsByProfile above — a profile with no configured
// errorPattern is simply absent here, so CompletionResolver falls back to its built-in
// narrow {@code (?i)\bAPI Error\s*:} compatibility pattern for that profile's targets
// (BackendErrorPatternLookup#legacy's contract — see backendErrorPatterns below).
// rather than handing it back as a real answer. Still compiled once at startup, keyed by
// profile name — unlike exhaustedPattern above (fleetd #446), errorPattern was left DEFERRED
// on purpose: the ticket that made exhaustedPattern hot scoped errorPattern/cooling-off out
// explicitly. A profile with no configured errorPattern is simply absent here, so
// CompletionResolver falls back to its built-in narrow {@code (?i)\bAPI Error\s*:}
// compatibility pattern for that profile's targets (BackendErrorPatternLookup#legacy's
// contract — see backendErrorPatterns below).
Map<String, Pattern> errorPatternsByProfile = new LinkedHashMap<>();
cfg.profiles().forEach((name, profile) -> {
if (profile.hasErrorPattern()) {
@@ -417,46 +427,20 @@ public final class Fleetd {
// 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
// the profile config live off `config`, so a credentialId edit is hot: no restart needed.
//
// fleetd #175: this sink is now also the quarantine target for OpenCodeLauncher's
// model-mismatch check — it needs nothing profile-specific from the caller beyond `target`
// (a herdr terminal id) and `reason`, so reusing it here is exactly "the existing
// ExhaustionSink path", not a new mechanism.
//
// fleetd #234: that check fires from SessionAwareHandle.agentSessionId(), which runs during
// SessionManager.acquire() BEFORE this session is registered in sessions.roster() — so the
// roster-only lookup below used to find nothing, .ifPresent silently no-op'd, and the
// ERROR the check had just logged ("quarantining this profile's credential") was a lie:
// nothing was quarantined, and nothing said so. Two changes: (1) OpenCodeLauncher now
// passes its OWN profile name via ExhaustionSink's 3-arg overload — it already has the
// FleetConfig.Profile in hand and does not need the roster at all — used here as a
// fallback whenever the roster lookup misses; (2) if a profile still cannot be resolved
// (neither the roster nor the hint names a configured one), this logs loudly at ERROR
// instead of silently doing nothing — a control that cannot act must say so.
// fleetd #234, round 4: the 3-arg overload is now ExhaustionSink's single abstract method,
// so this is safely a lambda — there is no separate 2-arg overload left for it to bind to
// instead and silently drop profileHint (that was rounds 1-3's whole hazard).
ExhaustionSink exhaustionSink = (target, reason, profileHint) -> {
String profileName = sessions.roster().stream()
.filter(session -> target.equals(session.terminalId()))
.findFirst()
.map(MemberSession::profile)
.orElse(profileHint);
FleetConfig.Profile profile = profileName == null ? null : config.get().profiles().get(profileName);
if (profile == null) {
log.error("quarantine requested for target '{}' ({}) but no profile could be "
+ "resolved — the target is not (yet) in the roster, and {} — "
+ "credential NOT quarantined (fleetd #234)",
target, reason,
profileHint == null ? "no profile hint was given"
: "the hinted profile '" + profileHint + "' is not configured");
return;
}
String credentialId = profile.effectiveCredentialId();
quarantine.quarantine(credentialId);
log.warn("credential '{}' quarantined for {}s (profile '{}'): {}", credentialId,
cfg.quarantineCooldownSeconds(), profile.profile(), reason);
};
// fleetd #446 criterion 3: the backend text that triggered the most recent quarantine of
// each credential, so fleet_profiles can report WHY a limit was hit, not only that it was.
// Keyed by credential id — the same key BackendQuarantine's own remainingSeconds uses —
// and written at the one call site that actually quarantines (inside exhaustionSink below),
// so a reason can never be reported for a quarantine that never happened. Bounded the same
// way BackendQuarantine's own internal map is documented to be: by the number of distinct
// credentials ever exhausted, not by how often the config is edited.
Map<String, String> quarantineReasonByCredential = new ConcurrentHashMap<>();
// fleetd #446 follow-up (round 3): extracted into exhaustionSink(...) below — see that
// method's javadoc for the full fleetd #175/#234/#446 history this used to carry inline —
// so a dedicated test can drive the exact ExhaustionSink main() builds, not a hand-rebuilt
// copy of its shape.
ExhaustionSink exhaustionSink = exhaustionSink(sessions, config, quarantine,
quarantineReasonByCredential, cfg);
// fleetd #175: point the forwarding sink handed to OpenCodeLauncher above at the real one,
// now that `sessions` exists to resolve target -> session -> profile.
exhaustionSinkRef.set(exhaustionSink);
@@ -673,7 +657,7 @@ public final class Fleetd {
// / BackendOutagePolicy would still be able to drift (e.g. a future edit to the credentialIdFor
// closure in only one of the two places), exactly the shape #284 was.
FleetMcp.QuarantineSource quarantineSource = quarantineSource(config, quarantine,
exhaustedPatternsByProfile);
liveExhaustedPatterns, quarantineReasonByCredential);
FleetMcp.OutageSource outageSource = new FleetMcp.OutageSource(profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.effectiveCredentialId();
@@ -805,16 +789,99 @@ public final class Fleetd {
}
/**
* fleetd #404: production source for quarantine reporting. Credential IDs are hot, but
* exhausted patterns are compiled once at startup for {@link CompletionResolver}, so the armed
* field must use that same compiled map until restart.
* fleetd #446 follow-up (round 3): the quarantine {@link ExhaustionSink} main() actually wires
* — extracted out of {@code main} for the same reason {@link #quarantineSource} below and
* {@link #usageLimitFixWarning}/{@link #usageLimitFixWarningNoModel} were. Round 2 pinned the
* WARNING text by having {@code FleetdUsageLimitFixWarningTest} call those two methods
* directly, but a mutation battery proved that left two gaps: nothing proved this sink's
* {@code log.warn} call actually uses either method (replacing the whole ternary with a
* literal string), and nothing proved it picks the right one for a profile WITH a configured
* {@code model:} versus one WITHOUT (swapping the ternary's two branches) — both mutations
* left round 2's full 1592-test suite green. {@code FleetdExhaustionSinkWarningTest} drives
* THIS factory's return value directly and asserts on the real text a {@code ListAppender}
* attached to this class's own logger captures — the only way to prove the caller, not just
* the two callees in isolation.
*
* <p>Behaviourally unchanged from the inline lambda this replaces:
* <ul>
* <li>fleetd #175: also the quarantine target for {@code OpenCodeLauncher}'s model-mismatch
* check — it needs nothing profile-specific beyond {@code target}/{@code reason}, so
* reusing this sink is the existing {@code ExhaustionSink} path, not a new mechanism;</li>
* <li>fleetd #234 (round 4): resolves {@code target} to a profile via the live session
* roster first, falling back to {@code profileHint} — {@code
* SessionAwareHandle.agentSessionId()} fires this sink during {@code
* SessionManager.acquire()}, <em>before</em> that session is registered, so the
* roster-only lookup alone used to silently miss it. A profile resolved neither way
* logs loudly at ERROR instead of doing nothing;</li>
* <li>fleetd #446 criterion 3: writes {@code quarantineReasonByCredential} at the one call
* site that actually quarantines, keyed by credential id — see the field's declaration
* in {@code main} for the bound on its size.</li>
* </ul>
*
* <p>{@code cfg} is the startup {@link FleetConfig} snapshot, read only for {@code
* quarantineCooldownSeconds()} in the log text — a Cold key (see {@code ConfigRef}'s class
* doc), so reading it off the snapshot rather than {@code config.get()} makes no observable
* difference and matches what the inline version already did.
*/
static ExhaustionSink exhaustionSink(SessionManager sessions, ConfigRef config, BackendQuarantine quarantine,
Map<String, String> quarantineReasonByCredential, FleetConfig cfg) {
return (target, reason, profileHint) -> {
String profileName = sessions.roster().stream()
.filter(session -> target.equals(session.terminalId()))
.findFirst()
.map(MemberSession::profile)
.orElse(profileHint);
FleetConfig.Profile profile = profileName == null ? null : config.get().profiles().get(profileName);
if (profile == null) {
log.error("quarantine requested for target '{}' ({}) but no profile could be "
+ "resolved — the target is not (yet) in the roster, and {} — "
+ "credential NOT quarantined (fleetd #234)",
target, reason,
profileHint == null ? "no profile hint was given"
: "the hinted profile '" + profileHint + "' is not configured");
return;
}
String credentialId = profile.effectiveCredentialId();
quarantine.quarantine(credentialId);
quarantineReasonByCredential.put(credentialId, reason);
log.warn("credential '{}' quarantined for {}s (profile '{}'): {}", credentialId,
cfg.quarantineCooldownSeconds(), profile.profile(), reason);
// fleetd #446 criterion 2: name the fix, not just the fact — an operator reading this
// should not have to work out which of several configured models to touch.
String model = profile.model();
log.warn(model != null && !model.isBlank()
? usageLimitFixWarning(profile.profile(), model)
: usageLimitFixWarningNoModel(profile.profile(), cfg.quarantineCooldownSeconds()));
};
}
/**
* fleetd #404, superseded by fleetd #446: production source for quarantine reporting.
* Credential IDs were already hot; {@code exhaustedPattern} used to be compiled once at
* startup for {@link CompletionResolver}, so the armed field had to read that same frozen map
* until restart — the exact asymmetry #446 exists to close. Both {@code exhaustedPatternArmed}
* here and {@link CompletionResolver}'s own classification now read the ONE live {@link
* LiveExhaustedPatterns} instance ({@code exhaustedPatterns.armed}/{@code .patternFor}), so a
* reload that arms or disarms a profile's detection is visible to both at once — never two
* independently-updated copies that could disagree, the fleetd #404 rule this method used to
* violate on purpose and now upholds.
*
* <p>{@code modelFor} and {@code reasonFor} back fleetd #446 criterion 3 — {@code
* fleet_profiles}'s quarantined row naming which model a quarantined profile runs
* ({@code modelFor}, read live off {@code config} the same way {@code credentialIdFor} already
* is) and the backend text that triggered the most recent quarantine of that credential
* ({@code reasonFor}, backed by {@code reasonByCredential} — see its call site in {@code main}
* for where that map is written, at the one place a credential is actually quarantined).
*/
static FleetMcp.QuarantineSource quarantineSource(ConfigRef config, BackendQuarantine quarantine,
Map<String, Pattern> startupExhaustedPatterns) {
LiveExhaustedPatterns exhaustedPatterns, Map<String, String> reasonByCredential) {
return new FleetMcp.QuarantineSource(profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.effectiveCredentialId();
}, quarantine, profile -> startupExhaustedPatterns.containsKey(profile));
}, quarantine, exhaustedPatterns::armed, profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.model();
}, reasonByCredential::get);
}
/**
@@ -833,6 +900,40 @@ public final class Fleetd {
* with 0 errors and left all 1506 existing tests green before {@code
* FleetdPatternCoverageLineTest} was added to catch exactly that swap.
*/
/**
* fleetd #446 follow-up: the criterion-2 WARNING text — "name the fix, not just the fact" —
* for a profile whose {@code model:} is configured. Extracted out of the {@code
* exhaustionSink} lambda the same way {@link #exhaustedPatternCoverageLine} was extracted out
* of {@code main}: the inline version compiled, ran, and read correctly, but nothing pinned
* its text, so a later edit could silently stop naming the fix and every test would stay
* green (measured: renaming the leading {@code "usage-limit fix:"} tag left all 1586
* pre-follow-up tests passing). {@link FleetdUsageLimitFixWarningTest} asserts on the return
* value of this method directly, which is exactly what the caller logs — {@code log.warn} is
* called with this method's result as a single, already-formatted argument, so the string a
* test sees here is byte-identical to what {@code fleetd.out} receives.
*
* <p>Deliberately asserts on three substrings, not the whole sentence (the profile name, the
* model name, and the literal {@code enabled: false} / {@code models.allow} action) — see that
* test's class doc for why a whole-sentence pin is the wrong granularity here.
*/
static String usageLimitFixWarning(String profile, String model) {
return "usage-limit fix: profile '" + profile + "' runs model '" + model + "' — set "
+ "`enabled: false` on that model's entry under models.allow in fleetd.yaml to "
+ "stop new spawns landing on it (models: is hot, no restart needed); remove the "
+ "line again once the subscription window resets";
}
/**
* fleetd #446 follow-up: the {@link #usageLimitFixWarning} counterpart for a profile with no
* {@code model:} configured — {@code models.allow} gates by model name, so there is nothing
* for an operator to flip, and this says so instead of naming a fix that does not exist.
*/
static String usageLimitFixWarningNoModel(String profile, Integer quarantineCooldownSeconds) {
return "usage-limit fix: profile '" + profile + "' has no model: configured, so "
+ "models.allow cannot gate it by name — the " + quarantineCooldownSeconds
+ "s quarantine above is the only thing keeping new spawns off it for now";
}
static String exhaustedPatternCoverageLine(Set<String> allProfiles, Set<String> configuredProfiles) {
return CompletionResolver.coverage("exhaustedPattern", CompletionResolver.UnsetMeaning.OFF,
allProfiles, configuredProfiles);
@@ -51,7 +51,17 @@ import java.util.function.Supplier;
* 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>
* to report; it moved here from deferred rather than joining split. An existing profile's
* {@code exhaustedPattern} (fleetd #446) joined this class the same way: it used to be
* compiled once into {@code Fleetd.main}'s startup pattern map (see the Deferred bullet's old
* wording, and {@code LiveExhaustedPatterns}'s class doc for the history), and is now read
* live, cached by profile name, by both {@code CompletionResolver}'s classification (via
* {@code LiveExhaustedPatterns.patternFor}) and {@code fleet_profiles}'s {@code
* exhaustionDetectionArmed} (via {@code LiveExhaustedPatterns.armed}) — the one live object
* both read, so a reload that arms or disarms a profile's usage-limit detection takes effect
* on the next check with no restart. {@code errorPattern}, {@code exhaustedPattern}'s sibling
* key for backend-error (not usage-limit) classification, was deliberately left OUT of this
* fleetd #446 change and stays deferred below — the ticket scoped it out explicitly.</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
@@ -81,15 +91,17 @@ import java.util.function.Supplier;
* (if any) simply keeps its old settings), adding or removing a profile (a new backend needs its own launcher,
* which is constructed once), <em>and an existing profile's launch settings</em> —
* {@code model}, {@code baseUrl}, {@code argv}, {@code env}, {@code mcpUrl},
* {@code exhaustedPattern} (CB-578 stage A — compiled once into {@code Fleetd.main}'s
* pattern map at startup), {@code errorPattern} (fleetd #201 Unit 5 — compiled once into
* {@code Fleetd.main}'s backend-error pattern map at startup, the same way),
* {@code errorPattern} (fleetd #201 Unit 5 — compiled once into
* {@code Fleetd.main}'s backend-error pattern map at startup; deliberately NOT made hot
* alongside {@code exhaustedPattern} by fleetd #446 — that ticket scoped {@code errorPattern}
* and cooling-off out explicitly),
* {@code ideProjectDir} / {@code ideOpenCommand} / {@code autoCompactWindow} (fleetd #323
* instance 1 — all three are read at spawn off the same frozen profile map and were missing
* from {@link #sameLaunchSettings}), and the rest of {@link #sameLaunchSettings}.
* {@code credentialId} (CB-578 stage B) is NOT on
* this list — it is read live off the config supplier at every quarantine check and
* exhaustion event, exactly like {@code weight} / {@code maxLoad}, so it is hot instead.
* {@code credentialId} (CB-578 stage B) and {@code exhaustedPattern} (fleetd #446) are NOT on
* this list — both are read live off the config supplier at every quarantine check and
* exhaustion event, exactly like {@code weight} / {@code maxLoad}, so both are hot instead
* (see the Hot bullet above for {@code exhaustedPattern}'s history).
* {@code HerdrPeerLauncher} takes {@code Map.copyOf(profiles)} at construction and resolves
* each spawn out of that copy, so those never reach a launch until the daemon restarts. A
* reload logs these rather than pretending they applied.</li>
@@ -583,12 +595,16 @@ public final class ConfigRef implements Supplier<FleetConfig> {
* {@link #sameLaunchSettings} because they are read <em>live</em>, not baked in at spawn — see
* the class doc's <em>Hot</em> bullet. {@code weight} and {@code maxLoad} are read live by the
* placement policy on every spawn; {@code credentialId} is read live by
* {@code CompositePeerLauncher} and the CB-578 stage B exhaustion sink. Nothing else is
* {@code CompositePeerLauncher} and the CB-578 stage B exhaustion sink; {@code exhaustedPattern}
* (fleetd #446) is read live, cached by profile name, by {@code LiveExhaustedPatterns} — see
* that class's doc and the class doc's <em>Hot</em> bullet for the history (it used to be
* compared here, deferred, like its sibling {@code errorPattern} still is). Nothing else is
* excluded — see {@code sameLaunchSettingsComparesEveryProfileComponentOrExcludesIt} in
* {@code ConfigRefProfileCoverageTest}, which enumerates every {@code Profile} record component
* by reflection and fails the build if one is neither compared below nor named here.
*/
static final Set<String> LAUNCH_SETTINGS_EXCLUDED = Set.of("weight", "maxLoad", "credentialId");
static final Set<String> LAUNCH_SETTINGS_EXCLUDED =
Set.of("weight", "maxLoad", "credentialId", "exhaustedPattern");
/**
* Whether two versions of a profile would launch a peer identically.
@@ -625,13 +641,15 @@ public final class ConfigRef implements Supplier<FleetConfig> {
&& Objects.equals(a.kind(), b.kind())
&& Objects.equals(a.env(), b.env())
&& Objects.equals(a.subscription(), b.subscription())
// CB-578 stage B: exhaustedPattern is compiled once into Fleetd.main's pattern map
// at startup (see ExhaustedPatternLookup wiring) — a reload never re-reads it, so a
// changed pattern must be reported as deferred, exactly like model/baseUrl/argv.
&& Objects.equals(a.exhaustedPattern(), b.exhaustedPattern())
// fleetd #446: exhaustedPattern moved to LAUNCH_SETTINGS_EXCLUDED — it is now read
// live, cached by profile name, through LiveExhaustedPatterns (see that class's doc
// and ConfigRef's class doc Hot bullet), so it must NOT be compared here any more: a
// reload that changes only exhaustedPattern must report "config reloaded", not
// "these changes need a restart".
// fleetd #201 Unit 5: errorPattern is compiled once into Fleetd.main's backend-error
// pattern map at startup (see BackendErrorPatternLookup wiring), the same way
// exhaustedPattern is — a reload never re-reads it either.
// pattern map at startup (see BackendErrorPatternLookup wiring) — unlike its sibling
// exhaustedPattern above (fleetd #446), a reload still never re-reads it; fleetd #446
// scoped errorPattern out on purpose (see the class doc's Hot bullet).
&& Objects.equals(a.errorPattern(), b.errorPattern())
// fleetd #323 instance 1: ideProjectDir and ideOpenCommand are read at spawn off the
// same frozen profile map as ideMcpUrl above (ClaudeCodeLauncher.java:267/269,
@@ -424,6 +424,12 @@ public record FleetConfig(
* {@code null}/blank ⇒ the classification never fires for this profile and
* today's completion-fallback behaviour is unchanged. Every backend words
* its refusal differently, so this is config, never a vendor string in code.
* Read live off the current config through {@code LiveExhaustedPatterns},
* cached by profile name (fleetd #446), so it is HOT: an operator can arm
* or disarm this profile's usage-limit detection by editing this key and
* reloading, with no restart. Before fleetd #446 this was compiled once at
* daemon startup, the same deferred shape its sibling {@code errorPattern}
* (below) still has.
* @param credentialId the shared account this profile authenticates as (CB-578 stage B). Two
* or more profiles setting the <em>same</em> non-blank value are quarantined
* together by one {@code exhaustedPattern} classification on any one of
@@ -0,0 +1,91 @@
package dev.ltms.fleet.inject;
import dev.ltms.fleet.config.FleetConfig;
import java.util.Map;
import java.util.Objects;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.Supplier;
import java.util.regex.Pattern;
/**
* fleetd #446: the single live source for "is {@code exhaustedPattern} configured for this
* profile, and what does it compile to". Read fresh off the config supplier on every call — the
* same reason {@code CompositePeerLauncher#models0} is a live supplier read rather than a value
* captured at construction (see that class's doc) — so an operator can arm or disarm usage-limit
* detection for a profile by editing {@code exhaustedPattern} and reloading, with no restart.
*
* <p>Before this class, {@code exhaustedPattern} was compiled once into a {@code Map<String,
* Pattern>} built inside {@code Fleetd.main} at startup ({@code ConfigRef}'s class doc used to
* list it under <em>Deferred</em>, CB-578 stage A) — the asymmetry fleetd #446 exists to close:
* an operator could turn a model off at runtime ({@code models.allow}'s {@code enabled: false},
* hot since fleetd #422) but could not arm the detector that would tell them to, without a
* restart. That was backwards for a feature whose whole point is to react while the fleet runs.
*
* <p>{@link #patternFor} is what {@link CompletionResolver} enforces on (via the {@link
* ExhaustedPatternLookup} production wiring in {@code Fleetd.main}); {@link #armed} is what
* {@code fleet_profiles}/{@code GET /profiles} report as {@code exhaustionDetectionArmed}. Both
* read this ONE object, so the report can never disagree with the behaviour — the same rule
* {@code CompositePeerLauncher#modelGateState()}'s javadoc states for the model gate's own
* armed/off pair (fleetd #404): "armed" and "which models are off" must come from one read of the
* same accessor the gate enforces on.
*
* <h2>Cache eviction</h2>
* Cached by profile NAME, not by pattern text. A {@code Map<String, Pattern>} keyed by the
* pattern STRING would grow by one entry per distinct regex ever typed for any profile across the
* daemon's uptime — unbounded in practice, because tuning a regex to match a backend's exact
* wording is exactly the kind of edit an operator makes several times while getting it right, and
* every edit-and-reload cycle would leave the previous attempt's compiled {@link Pattern} behind
* forever. Keyed by profile name instead, this cache holds at most one entry per profile name that
* has ever been looked up — and that key space is bounded by the (small, human-authored) set of
* configured profiles, which changes only on a restart: adding or removing a profile is itself a
* deferred key (a new backend needs its own launcher, built once — see {@code ConfigRef}'s class
* doc), so profile names do not churn the way pattern text does. Re-editing an EXISTING profile's
* {@code exhaustedPattern} — the case this class exists to make hot — simply overwrites that
* profile's one cache entry; it never adds a new one.
*/
public final class LiveExhaustedPatterns {
/** A compiled pattern paired with the source string it was compiled from, for change detection. */
private record Cached(String source, Pattern compiled) {
}
private final Supplier<Map<String, FleetConfig.Profile>> profiles;
private final ConcurrentHashMap<String, Cached> cache = new ConcurrentHashMap<>();
public LiveExhaustedPatterns(Supplier<Map<String, FleetConfig.Profile>> profiles) {
this.profiles = Objects.requireNonNull(profiles, "profiles");
}
/**
* The compiled {@code exhaustedPattern} currently configured for {@code profileName}, or
* {@code null} when that profile is unknown or has none configured. Recompiles only when the
* live pattern text differs from what is cached for this profile name; {@code
* FleetConfig#rejectMalformedProfilePatterns} already refuses a config (at load and at reload)
* whose {@code exhaustedPattern} does not compile, so this is not expected to throw in
* production — it is not defended against here for that reason, the same trust
* {@code CompositePeerLauncher#models0} places in config validation having already run.
*/
public Pattern patternFor(String profileName) {
if (profileName == null) {
return null;
}
FleetConfig.Profile profile = profiles.get().get(profileName);
if (profile == null || !profile.hasExhaustedPattern()) {
return null;
}
String source = profile.exhaustedPattern();
Cached cached = cache.get(profileName);
if (cached != null && cached.source().equals(source)) {
return cached.compiled();
}
Cached fresh = new Cached(source, Pattern.compile(source));
cache.put(profileName, fresh);
return fresh.compiled();
}
/** Whether {@code profileName} currently has a usage-limit pattern configured (live). */
public boolean armed(String profileName) {
return patternFor(profileName) != null;
}
}
@@ -121,17 +121,46 @@ public final class FleetMcp {
* CB-578 stage B quarantine facts used by {@code fleet_profiles}: a profile → credential id
* lookup, plus the shared {@link BackendQuarantine} to read remaining cooldowns off.
*
* @param exhaustedPatternArmed fleetd #395: profile → whether that profile's {@code
* exhaustedPattern} is configured (see {@code
* @param exhaustedPatternArmed fleetd #395, made live by fleetd #446: profile → whether that
* profile's {@code exhaustedPattern} is configured right now (see {@code
* FleetConfig.Profile#hasExhaustedPattern}), i.e. whether a backend refusal
* on it can EVER be classified {@code BACKEND_EXHAUSTED} and quarantine its
* credential. Bundled here, not a separate Source, because it answers the
* exact question {@code fleet_profiles}'s quarantine facts already answer
* for a QUARANTINED profile — "can this profile's usage limit ever be
* caught?" — just for every profile, not only one currently caught.
* caught?" — just for every profile, not only one currently caught. In
* production this reads {@link dev.ltms.fleet.inject.LiveExhaustedPatterns#armed} —
* the SAME live accessor {@link dev.ltms.fleet.inject.CompletionResolver}
* enforces on — so a reload that arms or disarms detection is reported
* correctly on the very next call, no restart (see that class's doc).
* @param modelFor fleetd #446 criterion 3: profile → the {@code model:} it runs, or
* {@code null} when the profile names none. Read live off the current
* config, the same way {@code credentialIdFor} already is. Used to name,
* in a QUARANTINED profile's row, which of the operator's configured models
* the fix (\"set {@code enabled: false} on it under {@code models.allow}\")
* actually applies to.
* @param reasonFor fleetd #446 criterion 3: credential id → the backend text that triggered
* its most recent quarantine, or {@code null} when none is known (e.g. an
* inert source, or a quarantine recorded before this field existed). Never
* consulted on its own to decide whether a credential is quarantined —
* callers gate on {@link BackendQuarantine#remainingSeconds} first, exactly
* like {@code credentialId} itself, so a stale reason left behind after a
* quarantine expires is never surfaced.
*/
public record QuarantineSource(Function<String, String> credentialIdFor, BackendQuarantine quarantine,
Function<String, Boolean> exhaustedPatternArmed) {
Function<String, Boolean> exhaustedPatternArmed,
Function<String, String> modelFor,
Function<String, String> reasonFor) {
/**
* Backward-compatible 3-arg form, before fleetd #446 added {@code modelFor}/{@code
* reasonFor} — neither is ever reported. Keeps every pre-existing call site (production and
* test) compiling and behaving identically for the quarantine facts they actually asked for.
*/
public QuarantineSource(Function<String, String> credentialIdFor, BackendQuarantine quarantine,
Function<String, Boolean> exhaustedPatternArmed) {
this(credentialIdFor, quarantine, exhaustedPatternArmed, _ -> null, _ -> null);
}
/**
* Backward-compatible 2-arg form, before fleetd #395 added {@code exhaustedPatternArmed} —
* reports every profile unarmed. Keeps every pre-existing call site (production and test)
@@ -994,7 +1023,11 @@ public final class FleetMcp {
if (isBlank(target) || isBlank(msgId)) {
return error("target and msgId are required");
}
messages.ackReply(target, msgId);
if (!messages.ackReply(target, msgId)) {
return error(msgId + " is not in " + target + "'s reply inbox (wrong id, wrong target, "
+ "or already acked). Held lead-to-lead (peer) mail cannot be acked this way — "
+ "read it with fleet_poll{coordId}.");
}
return text("acknowledged " + msgId);
}
@@ -1206,13 +1239,21 @@ public final class FleetMcp {
* #284 was, where one rule computed in two places was widened in only one and a single response
* contradicted itself. Shared inputs do not make duplicated computation safe.
*
* <p>fleetd #395: also reports {@code exhaustionDetectionArmed}, one boolean per configured
* profile — {@code true} when that profile's {@code exhaustedPattern} is set, {@code false}
* when it is not, so an operator can tell "this profile is healthy" from "nothing can ever
* quarantine this profile" without reading {@code fleetd.yaml}. Unlike {@code quarantined}/
* {@code coolingOff}, this map always names every profile: an unarmed profile never enters a
* transient state to be absent from, so silence here would read as "healthy" rather than "not
* being watched at all".
* <p>fleetd #395, made live by fleetd #446: also reports {@code exhaustionDetectionArmed}, one
* boolean per configured profile — {@code true} when that profile's {@code exhaustedPattern} is
* set RIGHT NOW, {@code false} when it is not, so an operator can tell "this profile is
* healthy" from "nothing can ever quarantine this profile" without reading {@code fleetd.yaml}.
* Unlike {@code quarantined}/{@code coolingOff}, this map always names every profile: an
* unarmed profile never enters a transient state to be absent from, so silence here would read
* as "healthy" rather than "not being watched at all". Hot since fleetd #446: a reload that
* arms or disarms a profile's {@code exhaustedPattern} changes this map's answer on the very
* next call, no restart — see {@link dev.ltms.fleet.inject.LiveExhaustedPatterns}'s class doc.
*
* <p>fleetd #446 criterion 3: each {@code quarantined} row also names {@code model} (the
* profile's configured {@code model:}, omitted when the profile names none) and {@code reason}
* (the backend text that triggered the most recent quarantine of that credential, omitted when
* none is known) — so a lead can see WHICH model to turn off and WHY, without reading the
* daemon log.
*/
public static Map<String, Object> profilesView(PeerLauncher workers, QuarantineSource quarantine, OutageSource outage) {
Map<String, Object> result = new LinkedHashMap<>();
@@ -1229,6 +1270,14 @@ public final class FleetMcp {
Map<String, Object> row = new LinkedHashMap<>();
row.put("credentialId", credentialId);
row.put("quarantinedForSeconds", remaining);
String model = quarantine.modelFor().apply(profile);
if (model != null && !model.isBlank()) {
row.put("model", model);
}
String reason = quarantine.reasonFor().apply(credentialId);
if (reason != null && !reason.isBlank()) {
row.put("reason", reason);
}
quarantined.put(profile, row);
});
}
@@ -1766,7 +1815,10 @@ public final class FleetMcp {
+ "has processed a reply and wants to confirm it, leaving other pending replies "
+ "in the inbox for later drain.",
objectSchema(Map.of(
"target", stringProp("Worker session id whose inbox to ack from"),
"target", stringProp("Worker session id whose inbox to ack from. Must name a "
+ "reply actually queued for it — an id in no inbox, or a coord-id "
+ "(peer held mail, read with fleet_poll{coordId} instead), errors "
+ "rather than reporting a false success"),
"msgId", stringProp("The message id to acknowledge")),
List.of("target", "msgId")));
}
@@ -1808,11 +1860,16 @@ public final class FleetMcp {
"List the configured worker profiles (backends) and which one fleet_spawn uses by "
+ "default. A 'quarantined' map is present when a backend-exhausted refusal put "
+ "a profile's credential on cooldown — fleet_spawn onto it is refused until "
+ "quarantinedForSeconds elapses; a profile sharing that credential is listed too. "
+ "quarantinedForSeconds elapses; a profile sharing that credential is listed too, "
+ "each row naming 'model' (the model that profile runs, when configured) and "
+ "'reason' (the backend text that triggered the quarantine, when known) — the fix "
+ "is usually `enabled: false` on that model under models.allow. "
+ "'exhaustionDetectionArmed' reports, per profile, whether a usage-limit refusal "
+ "on it can EVER be classified and quarantined (its exhaustedPattern is "
+ "configured) — false means that profile's credential can never be quarantined "
+ "by this mechanism, however many usage-limit refusals it sees.",
+ "by this mechanism, however many usage-limit refusals it sees. Both "
+ "exhaustedPattern and models.allow's on/off state are hot: editing fleetd.yaml "
+ "and reloading arms/disarms detection or flips a model off with no restart.",
objectSchema(Map.of(), List.of()));
}
@@ -17,6 +17,7 @@ import dev.ltms.fleet.peer.PeerHandle;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.peer.PeerUnreachableException;
import dev.ltms.fleet.peer.SpawnRequest;
import dev.ltms.fleet.placement.PlacementDecision;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -594,6 +595,23 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
req.sessionName(), spawned.agentSessionId(), spawned.receipt());
}
/**
* {@inheritDoc}
*
* <p>fleetd #450: re-enters {@link #spawn(SpawnRequest)} with {@code decision}'s profile named
* explicitly. This is the re-entering form the interface javadoc describes for a launcher with
* no placement concept of its own — an instance of this class spawns a single adapter's own
* profile set by explicit name only ({@link #place}/{@link #defaultProfileFor} are unoverridden
* here and just wrap {@link #defaultProfile()}); it does no quarantine/cool-off/maxLoad/model-off
* filtering of its own to re-apply. That filtering lives one layer up, in {@code
* CompositePeerLauncher}, which is the launcher that routes across more than one profile and
* therefore overrides this method with the routing form instead.
*/
@Override
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
return spawn(req.withProfile(decision.profile()));
}
/** The herdr daemon that owns this launcher's pane coordinates. */
public HerdrClient herdr() {
return agents.herdr();
@@ -394,17 +394,17 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
}
@Override
public void ack(String target, String msgId) {
public boolean ack(String target, String msgId) {
var perTarget = held.get(target);
if (perTarget == null || perTarget == RELEASED) {
return;
return false;
}
Held h;
synchronized (perTarget) {
h = perTarget.remove(msgId);
}
if (h == null) {
return; // never held (or already acked) — no-op
return false; // never held (or already acked) — no-op
}
try {
synchronized (channelLock) {
@@ -418,6 +418,7 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
}
throw new IllegalStateException("cannot ack reply " + msgId + " on " + queueName(target), e);
}
return true;
}
private DeliverCallback deliverCallback(String target) {
@@ -59,16 +59,17 @@ public final class InMemoryReplyInbox implements ReplyInbox {
}
@Override
public void ack(String target, String msgId) {
public boolean ack(String target, String msgId) {
if (!owned.contains(target)) {
return;
return false;
}
var perTarget = store.get(target);
if (perTarget != null) {
//noinspection SynchronizationOnLocalVariableOrMethodParameter
synchronized (perTarget) {
perTarget.remove(msgId);
}
if (perTarget == null) {
return false;
}
//noinspection SynchronizationOnLocalVariableOrMethodParameter
synchronized (perTarget) {
return perTarget.remove(msgId) != null;
}
}
}
@@ -826,9 +826,13 @@ public final class MessageService {
/**
* Acknowledge a specific reply by {@code msgId} for {@code target}. Removes it from the inbox
* so that a subsequent drain or peek no longer returns it.
*
* @return {@code true} if an entry was actually removed, {@code false} if {@code msgId} was not
* in {@code target}'s inbox (wrong id, wrong target, or already acked). The caller —
* {@link dev.ltms.fleet.mcp.FleetMcp#ack} — must not report success on {@code false}.
*/
public void ackReply(String target, String msgId) {
inbox.ack(target, msgId);
public boolean ackReply(String target, String msgId) {
return inbox.ack(target, msgId);
}
/**
@@ -46,6 +46,13 @@ public interface ReplyInbox {
/** Non-destructive snapshot of pending replies for {@code target} (FIFO), empty list if none. */
List<InboxMessage> peek(String target);
/** Remove the reply {@code msgId} for {@code target} once the primary has taken it. No-op if absent. */
void ack(String target, String msgId);
/**
* Remove the reply {@code msgId} for {@code target} once the primary has taken it.
*
* @return {@code true} if an entry was actually removed, {@code false} if there was nothing to
* remove (unknown {@code target}, unowned {@code target}, or a {@code msgId} not held for
* it). A {@code false} is not an error — acking a {@code target} this daemon does not own is
* part of the normal contract, not a failure.
*/
boolean ack(String target, String msgId);
}
@@ -253,26 +253,27 @@ public interface PeerLauncher {
* matching profile; passing a request that names a <em>different</em>, explicit profile than
* the decision it is paired with is a caller bug this method does not attempt to detect.
*
* <p>Default implementation for a launcher with no placement concept of its own: delegates to
* {@link #spawn(SpawnRequest)} with the decision's profile named explicitly — its only spawn
* contract, since there is no separate routing path to honor. This default is correct ONLY for
* a launcher that spawns a single profile of its own (e.g. {@code HerdrPeerLauncher}), where
* the explicit-profile branch it re-enters and the routing branch {@link #place} would have
* used are the same thing. <strong>A launcher that routes across more than one profile — the
* way {@code CompositePeerLauncher} routes across every configured adapter — MUST override
* this method instead of inheriting this default.</strong> Re-entering {@link
* #spawn(SpawnRequest)} re-applies that single-argument method's explicit-profile checks
* ({@code enforceNotQuarantined}, {@code enforceNotCoolingOff}, {@code enforceMaxLoad}, {@code
* enforceModelEnabled} in {@code CompositePeerLauncher}), which can refuse the very profile
* {@link #place} just chose, if the underlying placement state moved in the window between the
* {@link #place} call and this one — the exact window this method and {@link PlacementDecision}
* exist to close (fleetd #444).
* <p>No default implementation (fleetd #450): the two correct bodies disagree on purpose, so an
* implementer must choose one rather than silently inherit whichever this interface happened to
* provide. An implementer with no placement concept of its own — spawns a single profile, e.g.
* {@code HerdrPeerLauncher} — should delegate to {@link #spawn(SpawnRequest)} with the decision's
* profile named explicitly, since there is no separate routing path to honor there: the
* explicit-profile branch it re-enters and the routing branch {@link #place} would have used are
* the same thing. <strong>A launcher that routes across more than one profile — the way {@code
* CompositePeerLauncher} routes across every configured adapter — MUST NOT re-enter {@link
* #spawn(SpawnRequest)}.</strong> Doing so re-applies that single-argument method's
* explicit-profile checks ({@code enforceNotQuarantined}, {@code enforceNotCoolingOff}, {@code
* enforceMaxLoad}, {@code enforceModelEnabled} in {@code CompositePeerLauncher}), which can
* refuse the very profile {@link #place} just chose, if the underlying placement state moved in
* the window between the {@link #place} call and this one — the exact window this method and
* {@link PlacementDecision} exist to close (fleetd #444). Before #450 this was a {@code default}
* method that only {@code CompositePeerLauncher} overrode; a future placement-doing launcher
* could have inherited the re-entering body silently and never known. Making it abstract turns
* that silent inheritance into a compile error.
*
* @throws IllegalArgumentException if the decision names an unknown profile
*/
default PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
return spawn(req.withProfile(decision.profile()));
}
PeerHandle spawn(SpawnRequest req, PlacementDecision decision);
/**
* Resolve the effective working directory for a spawn {@code req} without actually spawning.
@@ -20,6 +20,7 @@ import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.peer.SpawnRequest;
import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementDecision;
import dev.ltms.fleet.placement.PlacementPolicies;
import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.SessionManager;
@@ -123,6 +124,11 @@ class FleetdBackendErrorSinkTest {
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
}
@Override
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
}
@Override
public Set<String> profiles() {
return Set.of();
@@ -2,6 +2,7 @@ package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.inject.LiveExhaustedPatterns;
import dev.ltms.fleet.mcp.FleetMcp;
import dev.ltms.fleet.placement.BackendQuarantine;
import org.junit.jupiter.api.DisplayName;
@@ -11,16 +12,31 @@ import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Map;
import java.util.regex.Pattern;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #404: {@code exhaustionDetectionArmed} must describe the startup pattern map, not the
* reloaded config snapshot.
* fleetd #446: {@code exhaustionDetectionArmed} must describe the LIVE config, not a startup
* snapshot — the opposite of what fleetd #404 asked for. #404 pinned {@code exhaustedPattern} as
* compiled once at startup (matching {@code CompletionResolver}'s then-frozen behaviour); #446's
* whole point is that both sides — the classification {@code CompletionResolver} enforces and the
* {@code exhaustionDetectionArmed} field this reports — now read the SAME live {@link
* LiveExhaustedPatterns} instance, so a reload arms or disarms detection with no restart, and the
* report can never disagree with the behaviour (the fleetd #404 rule, now upheld for real instead
* of by freezing both sides).
*
* <p>This test needs a reload. At startup the two snapshots agree, so a test of only a newly
* started daemon would not detect a live {@code config.get()} lookup in the report field.
*
* <p><b>The hotness proof criterion 1 demands:</b> run this against the fixed code (green) and
* against the pre-#446 compile site — {@code Fleetd.main}'s {@code exhaustedPatternsByProfile}
* compiled once into a {@code Map<String, Pattern>} at startup, handed to {@code quarantineSource}
* — and show it fails there. See the class doc history above: that shape is exactly what
* {@code reloadedPatternArmsDetectionWithNoRestart} below is written to catch, and it is the test
* that used to assert the opposite (see git history for this file's pre-#446 version, which
* asserted {@code assertFalse(...)} on the identical scenario this now asserts {@code assertTrue}
* on).
*/
class FleetdExhaustionDetectionArmedWiringTest {
@@ -54,41 +70,68 @@ class FleetdExhaustionDetectionArmedWiringTest {
""";
@Test
@DisplayName("reloading an exhaustedPattern does not arm the startup detection source")
void reloadedPatternDoesNotChangeTheArmedFieldUntilRestart(@TempDir Path dir) throws Exception {
@DisplayName("reloading a profile's exhaustedPattern IN arms exhaustionDetectionArmed, no restart")
void reloadedPatternArmsDetectionWithNoRestart(@TempDir Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, NO_PATTERN);
ConfigRef config = new ConfigRef(file, FleetConfig.load(file));
FleetMcp.QuarantineSource before = Fleetd.quarantineSource(config, BackendQuarantine.none(),
new LiveExhaustedPatterns(() -> config.get().profiles()), Map.of());
assertFalse(before.exhaustedPatternArmed().apply("terra"),
"no exhaustedPattern configured yet — must report unarmed");
Files.writeString(file, WITH_PATTERN);
assertTrue(config.reload().applied());
assertTrue(config.get().profiles().get("terra").hasExhaustedPattern());
FleetMcp.QuarantineSource source = Fleetd.quarantineSource(config, BackendQuarantine.none(),
Map.of());
assertFalse(source.exhaustedPatternArmed().apply("terra"),
"exhaustionDetectionArmed must use the startup pattern map, not config.get()");
// Re-read exhaustedPatternArmed off the SAME QuarantineSource built BEFORE the reload — no
// new object, no new wiring — to prove the field itself is live, not merely that a freshly
// built source would be. This is the exact axis the pre-#446 code fails on: the same
// Function<String,Boolean> lambda, called again after a reload, sees the new answer only if
// it re-reads config.get() on every call rather than a value captured earlier.
assertTrue(before.exhaustedPatternArmed().apply("terra"),
"exhaustionDetectionArmed must read the LIVE config on every call — fleetd #446 made "
+ "this hot; it must reflect a reload with no restart and no new QuarantineSource");
}
@Test
@DisplayName("a profile in the startup pattern map is reported as armed")
void aProfileInTheStartupMapIsArmed(@TempDir Path dir) throws Exception {
// fleetd #404, second direction. The test above only ever passes an EMPTY startup map, so
// it cannot tell a correct lookup from one that is permanently off. Measured: replacing the
// armed lambda with `profile -> false` left the whole suite green at 1475 tests. That
// mutation would make #395's visibility feature dead — an operator fixing a detection gap
// would be told the gap is still open after fixing it, forever. Both directions are needed:
// this test is the only thing that fails when the field stops reporting armed at all.
@DisplayName("reloading exhaustedPattern OUT disarms exhaustionDetectionArmed, no restart")
void reloadedPatternRemovalDisarmsDetectionWithNoRestart(@TempDir Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, WITH_PATTERN);
ConfigRef config = new ConfigRef(file, FleetConfig.load(file));
FleetMcp.QuarantineSource source = Fleetd.quarantineSource(config, BackendQuarantine.none(),
Map.of("terra", Pattern.compile("usage limit")));
new LiveExhaustedPatterns(() -> config.get().profiles()), Map.of());
assertTrue(source.exhaustedPatternArmed().apply("terra"),
"exhaustedPattern is configured from the start — must report armed");
Files.writeString(file, NO_PATTERN);
assertTrue(config.reload().applied());
assertFalse(source.exhaustedPatternArmed().apply("terra"),
"removing exhaustedPattern and reloading must disarm detection with no restart — an "
+ "operator turning detection off (e.g. while debugging a false positive) "
+ "must see that reflected immediately, exactly like arming it is");
}
@Test
@DisplayName("a profile with no exhaustedPattern configured is reported unarmed")
void aProfileWithNoPatternIsUnarmed(@TempDir Path dir) throws Exception {
// fleetd #404's second direction, carried forward: this test alone fails if
// exhaustedPatternArmed degenerates to a constant `profile -> true` — it needs at least one
// profile that is genuinely unarmed to catch that.
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, WITH_PATTERN);
ConfigRef config = new ConfigRef(file, FleetConfig.load(file));
FleetMcp.QuarantineSource source = Fleetd.quarantineSource(config, BackendQuarantine.none(),
new LiveExhaustedPatterns(() -> config.get().profiles()), Map.of());
assertTrue(source.exhaustedPatternArmed().apply("terra"),
"a profile whose pattern was compiled at startup must report armed");
"a profile whose exhaustedPattern is configured must report armed");
assertFalse(source.exhaustedPatternArmed().apply("sonnet"),
"a profile absent from the startup map must not report armed");
"a profile absent from config entirely must not report armed");
}
}
@@ -0,0 +1,175 @@
package dev.ltms.fleet;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.ExhaustionSink;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.session.SessionManager;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.slf4j.LoggerFactory;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #446 follow-up (round 3): round 2's {@link FleetdUsageLimitFixWarningTest} calls {@link
* Fleetd#usageLimitFixWarning}/{@link Fleetd#usageLimitFixWarningNoModel} directly, which pins
* the two methods' TEXT but cannot see whether {@link Fleetd#exhaustionSink}'s {@code log.warn}
* call actually invokes either one, or invokes the right one. A mutation battery run against
* round 2's merge (1592 green tests) proved both gaps real:
*
* <ul>
* <li><b>Cell A</b> — replacing the whole ternary result with a literal string
* ({@code "usage-limit fix: MUTANT"}) left every test green;</li>
* <li><b>Cell B</b> — swapping the ternary's two branches, so a profile WITH a {@code model:}
* gets the no-model fallback message and vice versa, was measured unmeasured by the lead
* but the existing test file shows no call that would notice it either.</li>
* </ul>
*
* <p>This class drives {@link Fleetd#exhaustionSink} — the actual factory {@code main} wires,
* since round 3's extraction pulled it out of the inline lambda for exactly this reason — and
* asserts on the REAL text a {@link ListAppender} attached to {@code Fleetd}'s own logger
* captures, mirroring {@code ExhaustedPatternGapReportTest}'s idiom. Chosen over a source-text
* assertion (the {@code FleetMcpAuthzTest} {@code everyRegisteredToolHasItsHandlerActionPinned}
* idiom) because that route can prove Cell A (the ternary is still there, still built from the
* two named methods) but not Cell B (which branch a given profile actually reaches) — this route
* answers both from one mechanism, since it reads what the sink actually logged for each shape of
* profile.
*
* <p>Each test filters for the WARN line starting {@code "usage-limit fix:"} specifically — {@link
* Fleetd#exhaustionSink} also logs a separate {@code "credential '...' quarantined for..."} WARN
* on every call, and asserting against the wrong one would pass or fail for the wrong reason. That
* prefix survives even under the M5 mutation above (the mutant keeps the tag, only the body
* becomes {@code "MUTANT"}), so the filter itself is not what either mutation defeats.
*/
class FleetdExhaustionSinkWarningTest {
private static final String YAML = """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
terra:
baseUrl: http://gx00.gw:8000
model: claude-opus-5
gx:
baseUrl: http://gx00.gw:8000
guard:
offSubscriptionHosts:
- gx00.gw
""";
private static SessionManager emptyRosterSessions() {
FakeHerdr h = new FakeHerdr();
FleetConfig.Profile dummy = new FleetConfig.Profile(
"dummy", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
ClaudeCodeLauncher launcher = new ClaudeCodeLauncher(new AgentControl(h), new WorkspaceControl(h),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(dummy.profile(), dummy), dummy.profile(), _ -> "tok");
// Never acquires a session — exhaustionSink's roster lookup is expected to miss and fall
// back to the profileHint argument, exactly like OpenCodeLauncher's real call site does
// (fleetd #234) — so the sink under test never needs a populated roster.
return new SessionManager(launcher);
}
private static ListAppender<ILoggingEvent> attach() {
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
return appender;
}
private static void detach(ListAppender<ILoggingEvent> appender) {
((Logger) LoggerFactory.getLogger(Fleetd.class)).detachAppender(appender);
}
/** The one WARN line {@code exhaustionSink} builds from the ternary under test, or {@code null}. */
private static String fixWarning(ListAppender<ILoggingEvent> appender) {
List<String> warns = appender.list.stream()
.filter(e -> e.getLevel() == Level.WARN)
.map(ILoggingEvent::getFormattedMessage)
.filter(m -> m.startsWith("usage-limit fix:"))
.toList();
assertTrue(warns.size() == 1,
"expected exactly one 'usage-limit fix:' WARN per onExhausted call, got " + warns.size()
+ ": " + warns);
return warns.getFirst();
}
@Test
@DisplayName("a profile WITH a configured model gets the actionable models.allow fix, not the fallback")
void profileWithModelGetsTheActionableFix(@TempDir Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, YAML);
FleetConfig cfg = FleetConfig.load(file);
ConfigRef config = ConfigRef.fixed(cfg);
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
Map<String, String> reasonByCredential = new HashMap<>();
ExhaustionSink sink = Fleetd.exhaustionSink(emptyRosterSessions(), config, quarantine,
reasonByCredential, cfg);
ListAppender<ILoggingEvent> appender = attach();
try {
sink.onExhausted("term_x", "The usage limit has been reached", "terra");
} finally {
detach(appender);
}
String warning = fixWarning(appender);
assertTrue(warning.contains("terra"), "must name the profile: " + warning);
assertTrue(warning.contains("claude-opus-5"), "must name the model: " + warning);
assertTrue(warning.contains("enabled: false"), "must name the action: " + warning);
assertTrue(warning.contains("models.allow"), "must name where the action goes: " + warning);
assertFalse(warning.contains("has no model: configured"),
"a profile WITH a model must not get the no-model fallback text: " + warning);
}
@Test
@DisplayName("a profile with NO configured model gets the fallback, never a fix that does not exist")
void profileWithNoModelGetsTheFallbackNotAFix(@TempDir Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, YAML);
FleetConfig cfg = FleetConfig.load(file);
ConfigRef config = ConfigRef.fixed(cfg);
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
Map<String, String> reasonByCredential = new HashMap<>();
ExhaustionSink sink = Fleetd.exhaustionSink(emptyRosterSessions(), config, quarantine,
reasonByCredential, cfg);
ListAppender<ILoggingEvent> appender = attach();
try {
sink.onExhausted("term_y", "The usage limit has been reached", "gx");
} finally {
detach(appender);
}
String warning = fixWarning(appender);
assertTrue(warning.contains("gx"), "must name the profile: " + warning);
assertTrue(warning.contains("has no model: configured"),
"must say there is no model to gate by name: " + warning);
assertFalse(warning.contains("enabled: false"),
"a profile with no model: configured has no models.allow entry to flip — must "
+ "not claim one exists: " + warning);
}
}
@@ -0,0 +1,75 @@
package dev.ltms.fleet;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #446 follow-up (criterion 2): {@code exhaustionSink} used to build the "name the fix"
* WARNING inline with two SLF4J {@code {}}-placeholder {@code log.warn} calls. Nothing pinned
* that text — a mutation battery run against the merged PR #457 proved it: renaming the leading
* {@code "usage-limit fix:"} tag to {@code "MUTANT no fix named:"} at both call sites left all
* 1586 existing tests green. This class is what makes that mutation red.
*
* <p>{@link Fleetd#usageLimitFixWarning} and {@link Fleetd#usageLimitFixWarningNoModel} are the
* extracted call sites {@code exhaustionSink} actually invokes — {@code log.warn} is called with
* each method's return value as a single, already-formatted argument, so the string these tests
* assert on is byte-identical to what {@code fleetd.out} receives. Same extracted-static-method +
* dedicated-test idiom as {@link Fleetd#exhaustedPatternCoverageLine} /
* {@link FleetdPatternCoverageLineTest}, chosen over the {@link ch.qos.logback.core.read.ListAppender}
* idiom {@code ExhaustedPatternGapReportTest} uses — this warning is built once, at one call site,
* from a single method with no branching inside the message itself, so there is no real log
* plumbing left to prove; capturing an appender would only add setup/teardown around the same
* assertion.
*
* <p>Each assertion below checks for a specific substring — the leading {@code "usage-limit fix:"}
* tag an operator would grep the log for, the profile name, the model name, and the literal
* action ({@code enabled: false} under {@code models.allow}) — rather than the whole sentence.
* Criterion 2 promised those facts, not a particular wording of the prose that carries them; a
* whole-sentence {@code assertEquals} breaks on every copy-edit of that surrounding prose and
* teaches the next person to delete the test rather than read why it failed. The leading tag is
* pinned on purpose, unlike the rest of the sentence: it is the identifying label the fleetd #446
* follow-up's own mutation battery renamed ({@code "usage-limit fix:"} → {@code "MUTANT no fix
* named:"}) to prove this class did not yet catch a tag rename — only the three-fact substrings
* above would have missed it, since none of them mention the tag itself.
*/
class FleetdUsageLimitFixWarningTest {
@Test
@DisplayName("usageLimitFixWarning names the profile, the model, and the models.allow fix")
void usageLimitFixWarningNamesProfileModelAndFix() {
String warning = Fleetd.usageLimitFixWarning("terra", "claude-opus-5");
assertTrue(warning.startsWith("usage-limit fix:"),
"must carry the grep-able identifying tag: " + warning);
assertTrue(warning.contains("terra"), "must name the profile: " + warning);
assertTrue(warning.contains("claude-opus-5"), "must name the model: " + warning);
assertTrue(warning.contains("enabled: false"), "must name the action: " + warning);
assertTrue(warning.contains("models.allow"), "must name where the action goes: " + warning);
}
@Test
@DisplayName("usageLimitFixWarning does not cross-name a different profile or model")
void usageLimitFixWarningDoesNotCrossName() {
String warning = Fleetd.usageLimitFixWarning("terra", "claude-opus-5");
// Guards against the mutation of swapping the two profile/model arguments at the call
// site — a test that only checks "some name appears" cannot catch that.
assertTrue(!warning.contains("sonnet"), "must not name an unrelated model: " + warning);
}
@Test
@DisplayName("usageLimitFixWarningNoModel names the profile but no fix, since none exists")
void usageLimitFixWarningNoModelNamesProfileNotAFix() {
String warning = Fleetd.usageLimitFixWarningNoModel("gx", 1800);
assertTrue(warning.startsWith("usage-limit fix:"),
"must carry the grep-able identifying tag: " + warning);
assertTrue(warning.contains("gx"), "must name the profile: " + warning);
assertTrue(warning.contains("1800"), "must name the quarantine duration it falls back to: " + warning);
assertTrue(!warning.contains("enabled: false"),
"a profile with no model: configured has no models.allow entry to flip — the "
+ "fallback message must not claim one exists: " + warning);
}
}
@@ -161,7 +161,15 @@ class ConfigRefProfileCoverageTest {
// the #323 bug. Only ideProjectDir and worktreeGroup were saved by a behavioural test in
// ConfigRefTest; the other two had none. So: growing this set now requires editing this
// line as well, which is a visible, deliberate diff rather than a quiet one.
assertEquals(Set.of("weight", "maxLoad", "credentialId"), excluded,
//
// fleetd #446: "exhaustedPattern" joined the set legitimately — it moved from compiled
// once at Fleetd.main startup to read live, cached by profile name, through
// LiveExhaustedPatterns (see that class's doc and ConfigRef's class doc). Unlike the #323
// hazard this comment warns about, it is pinned by its OWN behavioural tests:
// FleetdExhaustionDetectionArmedWiringTest proves the reload takes effect with no restart,
// and the same test rerun against the pre-#446 compile site (see that ticket's report) is
// what proves the assertion would have failed before the fix.
assertEquals(Set.of("weight", "maxLoad", "credentialId", "exhaustedPattern"), excluded,
"ConfigRef.LAUNCH_SETTINGS_EXCLUDED changed. A component belongs in it ONLY if it "
+ "is read live off the config supplier, not baked into a launcher at "
+ "startup. If you are adding one to silence this test, that is fleetd #323 "
@@ -339,12 +339,18 @@ class ConfigRefTest {
}
/**
* CB-578 stage B: exhaustedPattern is compiled once into Fleetd.main's pattern map at startup
* (see ExhaustedPatternLookup), so a reload never re-reads it — a changed pattern must be
* reported deferred exactly like model/baseUrl, not silently claimed as applied.
* fleetd #446: exhaustedPattern moved from deferred to hot — {@code LiveExhaustedPatterns}
* reads {@code config.get()} fresh on every lookup, cached by profile name (not baked into a
* frozen map inside {@code Fleetd.main} any more, the way it was before this ticket, and the
* way its sibling {@code errorPattern} still is — see
* {@code changingAProfilesErrorPatternIsReportedAsDeferred} right below for the contrast).
* So a reload that changes ONLY exhaustedPattern must report a plain "config reloaded", with
* nothing deferred — the opposite of what this test asserted before fleetd #446, when
* ConfigRef.sameLaunchSettings still compared exhaustedPattern and this reload reported it as
* one of the changes needing a restart.
*/
@Test
void changingAProfilesExhaustedPatternIsReportedAsDeferred(@TempDir Path dir) throws Exception {
void changingAProfilesExhaustedPatternTakesEffectWithNoRestart(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
@@ -379,17 +385,21 @@ class ConfigRefTest {
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertEquals(1, out.deferred().size(), out.deferred().toString());
assertTrue(out.deferred().getFirst().contains("sonnet"), out.deferred().toString());
assertTrue(out.deferred().getFirst().contains("launch settings"), out.deferred().toString());
// The snapshot still carries the new value — a restart is what makes it take effect.
assertEquals(0, out.deferred().size(),
"exhaustedPattern is hot since fleetd #446 — changing only this key must not be "
+ "reported as needing a restart: " + out.deferred());
assertEquals(0, out.split().size(), out.split().toString());
assertEquals("config reloaded", out.summary());
// The snapshot carries the new value immediately — this IS the live value now, not
// merely what a restart would eventually pick up.
assertEquals("rate limit exceeded", ref.get().profiles().get("sonnet").exhaustedPattern());
}
/**
* fleetd #201 Unit 5: errorPattern is compiled once into Fleetd.main's backend-error pattern
* map at startup (see BackendErrorPatternLookup), the same way exhaustedPattern is above — a
* reload never re-reads it either, so a changed value must be reported deferred.
* map at startup (see BackendErrorPatternLookup) — unlike its sibling exhaustedPattern above,
* which fleetd #446 made hot, errorPattern was scoped OUT of that ticket on purpose and stays
* deferred: a reload still never re-reads it, so a changed value must be reported deferred.
*/
@Test
void changingAProfilesErrorPatternIsReportedAsDeferred(@TempDir Path dir) throws Exception {
@@ -18,12 +18,24 @@ import static org.junit.jupiter.api.Assumptions.assumeTrue;
* the throwaway space down.
*
* <p>The seed shell's own startup (restoring its session, printing its banner) is asynchronous
* and its length is not a fleetd contract — measured here at ~2.5s on one host (fleetd #449). A
* fixed sleep before typing raced that startup: input typed before the shell reached its prompt
* was swallowed by the shell's own startup, and the pane showed the typed line followed by the
* startup banner with no command output at all — indistinguishable, at a glance, from the env
* map never reaching the shell. So this polls for a real signal (the pane's visible text
* settling, then the expected output appearing) instead of guessing a sleep length.
* and its length is not a fleetd contract — measured here at ~2.5s on one host (fleetd #449).
* The old version used two fixed sleeps: 1000ms before typing, then 800ms before reading. It
* failed, and the pane showed the typed line followed by the startup banner with no command
* output at all — which looks, at a glance, exactly like the env map never reaching the shell.
* So this polls for a real signal instead of guessing a sleep length.
*
* <p><strong>What was measured, and what was not.</strong> Polling fixes it: 5 standalone runs
* green. The load-bearing half is {@link #waitForText}. With {@link #SHELL_READY_TIMEOUT_MS}
* set to 0 — so input is typed at once, with no settle wait at all — the test still passed 3 of
* 3. So the proven cause is the 800ms READ deadline being too short, not the 1000ms write delay.
* Note the direction, because it matters: typing at 0ms works where typing at 1000ms failed. The
* earlier explanation for this test — that input typed before the prompt is swallowed by the
* shell's startup — is therefore NOT supported by any measurement here. Please do not repeat it
* as the reason; if it were true, 0ms would be worse than 1000ms, and it is better.
*
* <p>{@link #waitUntilSettled} is kept as cheap insurance against that swallow case, not because
* anyone showed it was needed. If you want to delete it, the honest test is whether you can make
* this test fail by typing early. Nobody has managed that yet.
*
* <p>Tagged {@code contract}; run with {@code mvn test -Pcontract}.
*/
@@ -48,8 +60,9 @@ class AgentControlContractTest {
/**
* Poll {@code pane.read} until two consecutive reads come back identical — the shell's own
* startup output (restore banner, prompt) has stopped changing — or {@code timeoutMs} elapses.
* Never asserts by itself; the caller's own assertion is what actually verifies the outcome,
* this only avoids sending input into a shell still mid-startup.
* Never asserts by itself; the caller's own assertion is what actually verifies the outcome.
* Setting {@code timeoutMs} to 0 skips the wait entirely and the test still passes here, so
* treat this as insurance rather than as the fix — see the class javadoc.
*/
private static String waitUntilSettled(UnixSocketHerdrClient herdr, String paneId, long timeoutMs)
throws InterruptedException {
@@ -0,0 +1,138 @@
package dev.ltms.fleet.inject;
import dev.ltms.fleet.config.FleetConfig;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import java.util.Map;
import java.util.regex.Pattern;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotSame;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #446: {@code LiveExhaustedPatterns} is the mechanism that makes {@code exhaustedPattern}
* hot — read fresh off a config supplier per lookup, cached by profile name. This is the direct,
* minimal unit test of that mechanism: read the pattern once, mutate the backing "config", read
* again, and assert the second read saw the new value — the exact hotness proof shape.
*
* <p>This class alone does not prove {@code Fleetd.main} actually wires the live source into
* production — {@link dev.ltms.fleet.FleetdExhaustionDetectionArmedWiringTest} and {@code
* CompletionResolverTest}'s {@code ExhaustedPatternLookup} tests do that at the wiring level. What
* this class proves is that the mechanism itself is genuinely live and genuinely cached, not that
* it is used.
*/
class LiveExhaustedPatternsTest {
private static FleetConfig.Profile profileWithPattern(String pattern) {
return new FleetConfig.Profile("terra", "http://gx00.gw:8000", "terra",
null, null, null, "tab", "fleetd-workers", "w #{n}", null, null, null,
null, null, null, null, null, null, null, pattern);
}
/** Simple mutable holder standing in for {@code ConfigRef} — swapped, never mutated in place. */
private static final class MutableProfiles {
private volatile Map<String, FleetConfig.Profile> profiles;
MutableProfiles(Map<String, FleetConfig.Profile> initial) {
this.profiles = initial;
}
Map<String, FleetConfig.Profile> get() {
return profiles;
}
void set(Map<String, FleetConfig.Profile> fresh) {
this.profiles = fresh;
}
}
@Test
@DisplayName("patternFor reads the live config: read, change, read again, see the change — no restart")
void patternForIsHotAcrossAConfigChange() {
MutableProfiles live = new MutableProfiles(Map.of("terra", profileWithPattern("usage limit")));
LiveExhaustedPatterns patterns = new LiveExhaustedPatterns(live::get);
Pattern before = patterns.patternFor("terra");
assertTrue(before.matcher("your usage limit has been reached").find());
// Change the "config" — the exact thing a reload does to ConfigRef in production.
live.set(Map.of("terra", profileWithPattern("rate limit exceeded")));
Pattern after = patterns.patternFor("terra");
assertTrue(after.matcher("429: rate limit exceeded, try later").find(),
"the second read must see the NEW pattern text with no restart");
assertFalse(after.matcher("your usage limit has been reached").find(),
"the second read must have actually recompiled — not just reused a stale match");
}
@Test
@DisplayName("armed follows the same live read: on, then off, after a config change — no restart")
void armedIsHotAcrossAConfigChange() {
MutableProfiles live = new MutableProfiles(Map.of("terra", profileWithPattern("usage limit")));
LiveExhaustedPatterns patterns = new LiveExhaustedPatterns(live::get);
assertTrue(patterns.armed("terra"));
live.set(Map.of("terra", profileWithPattern(null)));
assertFalse(patterns.armed("terra"),
"removing exhaustedPattern and 'reloading' must disarm detection with no restart");
}
@Test
@DisplayName("an unchanged pattern string is not recompiled — the cache is reused")
void unchangedPatternTextReusesTheCompiledInstance() {
MutableProfiles live = new MutableProfiles(Map.of("terra", profileWithPattern("usage limit")));
LiveExhaustedPatterns patterns = new LiveExhaustedPatterns(live::get);
Pattern first = patterns.patternFor("terra");
// A NEW Profile object, same pattern text — a reload of an unrelated key produces a fresh
// FleetConfig even though this profile's own text did not change.
live.set(Map.of("terra", profileWithPattern("usage limit")));
Pattern second = patterns.patternFor("terra");
assertSame(first, second,
"identical pattern text must reuse the cached compiled Pattern, not recompile it — "
+ "recompiling a regex on every check is exactly the cost the cache exists to avoid");
}
@Test
@DisplayName("a changed pattern string is recompiled, not silently reused from the cache")
void changedPatternTextIsRecompiled() {
MutableProfiles live = new MutableProfiles(Map.of("terra", profileWithPattern("usage limit")));
LiveExhaustedPatterns patterns = new LiveExhaustedPatterns(live::get);
Pattern first = patterns.patternFor("terra");
live.set(Map.of("terra", profileWithPattern("rate limit")));
Pattern second = patterns.patternFor("terra");
assertNotSame(first, second, "changed pattern text must produce a freshly compiled Pattern");
}
@Test
void unknownProfileReturnsNullAndUnarmed() {
LiveExhaustedPatterns patterns = new LiveExhaustedPatterns(Map::of);
assertNull(patterns.patternFor("nope"));
assertFalse(patterns.armed("nope"));
}
@Test
void nullProfileNameReturnsNullAndUnarmed() {
LiveExhaustedPatterns patterns = new LiveExhaustedPatterns(
() -> Map.of("terra", profileWithPattern("usage limit")));
assertNull(patterns.patternFor(null));
assertFalse(patterns.armed(null));
}
@Test
void profileWithNoPatternConfiguredReturnsNull() {
LiveExhaustedPatterns patterns = new LiveExhaustedPatterns(
() -> Map.of("terra", profileWithPattern(null)));
assertNull(patterns.patternFor("terra"));
assertFalse(patterns.armed("terra"));
}
}
@@ -352,7 +352,7 @@ class FleetMcpTest {
fail("a lead fleet_reply must not publish to the worker inbox");
}
@Override public List<InboxMessage> peek(String target) { return List.of(); }
@Override public void ack(String target, String msgId) { }
@Override public boolean ack(String target, String msgId) { return false; }
};
MessageService leadMessages = new MessageService(agents, new Injector(agents), new Rendezvous(),
inboxThatRejectsPublishes);
@@ -1386,6 +1386,10 @@ class FleetMcpTest {
@Test
void bridgeAckReturnsConfirmationForValidArgs() {
// fleet_ack only reports success for a msgId actually queued in the target's inbox
// (fleetd #437) — publish one via the inbox directly rather than asserting on a
// fabricated id nothing ever queued.
inbox.publish("term_a", "msg-1", "queued reply");
McpSchema.CallToolResult res = FleetMcp.ack(messages, "term_a", "msg-1");
assertNotEquals(Boolean.TRUE, res.isError());
assertTrue(textOf(res).contains("msg-1"), "response should mention the msgId");
@@ -1398,6 +1402,26 @@ class FleetMcpTest {
assertTrue(FleetMcp.ack(messages, " ", "msg-1").isError());
}
@Test
void bridgeAckOfAnIdInNoInboxIsAnError() {
// fleetd #437: fleet_ack used to say "acknowledged <msgId>" for a message it never
// touched, because nothing in the chain reported hit vs. miss. "never-queued" is in no
// inbox at all, so this must error rather than claim success.
McpSchema.CallToolResult res = FleetMcp.ack(messages, "term_a", "never-queued");
assertTrue(res.isError());
assertTrue(textOf(res).contains("never-queued"), textOf(res));
}
@Test
void bridgeAckOfACoordIdTargetIsAnErrorNamingFleetPoll() {
// A coord-id names a peer lead's held mailbox (LeadChannel/LeadMailbox), never a
// worker's ReplyInbox — fleet_ack has no route to it and must say so, pointing at
// fleet_poll{coordId} instead of reporting a false "acknowledged".
McpSchema.CallToolResult res = FleetMcp.ack(messages, "coord-some-peer", "msg-1");
assertTrue(res.isError());
assertTrue(textOf(res).contains("fleet_poll{coordId}"), textOf(res));
}
@Test
void bridgeAckRemovesSpecificReply() {
// Queue a reply and capture its msgId.
@@ -1411,8 +1435,18 @@ class FleetMcpTest {
var peeked = messages.drainReplies("term_a");
assertEquals(1, peeked.size(), "one fresh reply in the inbox");
// ackReply works (no-op since published with a different UUID, but callable).
assertDoesNotThrow(() -> messages.ackReply("term_a", msgId));
// fleetd #437: msgId was already drained above (a fresh UUID each publish), so it is no
// longer in the inbox — ackReply must now report that miss instead of pretending to ack.
assertFalse(messages.ackReply("term_a", msgId));
}
@Test
void bridgeAckRemovingARealQueuedReplyReportsSuccessAndRemovesIt() {
// The worker path must not change behaviour: acking a reply that IS still in the inbox
// still succeeds and still removes it (fleetd #437).
inbox.publish("term_a", "real-1", "still queued");
assertTrue(messages.ackReply("term_a", "real-1"), "ack of a real queued reply must report true");
assertTrue(inbox.peek("term_a").isEmpty(), "the acked reply must be gone from the inbox");
}
@Test
@@ -0,0 +1,111 @@
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.peer.PeerLauncher;
import dev.ltms.fleet.placement.BackendQuarantine;
import io.modelcontextprotocol.spec.McpSchema;
import org.junit.jupiter.api.Test;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #446 follow-up (criterion 3): a mutation battery run against the merged PR #457 proved
* that deleting either {@code row.put("model", model);} or {@code row.put("reason", reason);}
* from {@link FleetMcp#profilesView} left all 1586 existing tests green — nothing exercised a
* {@link FleetMcp.QuarantineSource} whose {@code modelFor}/{@code reasonFor} actually return a
* value. This class is what makes both deletions red, and also pins the two "must be omitted, not
* emitted as null/blank" branches those two {@code if} guards exist for — {@link
* FleetMcpTest#profilesReportsAQuarantinedCredential} covers the credential/quarantine-duration
* shape but supplies neither function, so it cannot distinguish "the field is missing" from "the
* field was never asked for".
*/
class FleetProfilesQuarantineModelReasonFieldsTest {
private static final String PROFILE = "terra";
private static final String CREDENTIAL = "cred-terra";
private static PeerLauncher launcher(FakeHerdr h) {
FleetConfig.Profile profile = new FleetConfig.Profile(
PROFILE, "http://gx00.gw:8000", "claude-opus-5", null, "FLEETD_WORKER_TOKEN", null,
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
return new ClaudeCodeLauncher(new AgentControl(h), new WorkspaceControl(h),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(PROFILE, profile), PROFILE, _ -> "tok");
}
private static BackendQuarantine quarantined(String credentialId) {
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
quarantine.quarantine(credentialId);
return quarantine;
}
private static String textOf(McpSchema.CallToolResult r) {
return ((McpSchema.TextContent) r.content().getFirst()).text();
}
@Test
void quarantinedRowNamesTheModelAndTheReason() {
FakeHerdr h = new FakeHerdr();
FleetMcp.QuarantineSource source = new FleetMcp.QuarantineSource(
_ -> CREDENTIAL, quarantined(CREDENTIAL), _ -> true,
_ -> "claude-opus-5", _ -> "The usage limit has been reached");
McpSchema.CallToolResult res = FleetMcp.profiles(launcher(h), source);
assertNotEquals(Boolean.TRUE, res.isError());
String out = textOf(res);
assertTrue(out.contains("\"quarantined\""), out);
assertTrue(out.contains("\"model\":\"claude-opus-5\""),
"the quarantined row must name the model the fix applies to: " + out);
assertTrue(out.contains("\"reason\":\"The usage limit has been reached\""),
"the quarantined row must name why it was quarantined: " + out);
}
@Test
void quarantinedRowOmitsModelWhenNoneIsConfigured() {
FakeHerdr h = new FakeHerdr();
// null and "" (blank) both exercise the same production guard (`!model.isBlank()`) —
// covered together since both must produce the identical outcome: the key absent.
for (String noModel : new String[] {null, "", " "}) {
FleetMcp.QuarantineSource source = new FleetMcp.QuarantineSource(
_ -> CREDENTIAL, quarantined(CREDENTIAL), _ -> true,
_ -> noModel, _ -> "some reason");
McpSchema.CallToolResult res = FleetMcp.profiles(launcher(h), source);
String out = textOf(res);
assertTrue(out.contains("\"quarantined\""), out);
assertFalse(out.contains("\"model\""),
"modelFor returned " + (noModel == null ? "null" : "'" + noModel + "'")
+ " — the key must be ABSENT, not emitted as null or an empty string: " + out);
}
}
@Test
void quarantinedRowOmitsReasonWhenNoneIsKnown() {
FakeHerdr h = new FakeHerdr();
for (String noReason : new String[] {null, "", " "}) {
FleetMcp.QuarantineSource source = new FleetMcp.QuarantineSource(
_ -> CREDENTIAL, quarantined(CREDENTIAL), _ -> true,
_ -> "claude-opus-5", _ -> noReason);
McpSchema.CallToolResult res = FleetMcp.profiles(launcher(h), source);
String out = textOf(res);
assertTrue(out.contains("\"quarantined\""), out);
assertFalse(out.contains("\"reason\""),
"reasonFor returned " + (noReason == null ? "null" : "'" + noReason + "'")
+ " — the key must be ABSENT, not emitted as null or an empty string: " + out);
}
}
}
@@ -16,6 +16,7 @@ import java.util.List;
import java.util.concurrent.TimeUnit;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -221,6 +222,31 @@ class AmqpReplyInboxContractTest {
}
}
@Test
void ackReportsHitVsMissAgainstARealBroker() throws Exception {
// fleetd #437: fleet_ack said "acknowledged <msgId>" for a message it never touched,
// because ReplyInbox.ack() (void) could not tell a hit from a miss. Pin the fixed
// boolean contract against a real broker — the adapter fleetd actually runs live.
String target = "worker-ack-contract-" + System.nanoTime();
try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
inbox.own(target);
// Never held for this target at all: must report false, not throw.
assertFalse(inbox.ack(target, "never-held"),
"acking a msgId never held for an owned target must report false");
// A real message: first ack removes it and reports true...
inbox.publish(target, "m1", "ack me");
assertEquals(1, awaitPeek(inbox, target).size(), "the published reply should be held");
assertTrue(inbox.ack(target, "m1"), "acking a held reply must report true");
assertTrue(inbox.peek(target).isEmpty(), "an acked reply is dropped");
// ...and the second ack of the SAME msgId has nothing left to remove: false.
assertFalse(inbox.ack(target, "m1"),
"acking the same msgId twice must report false the second time");
}
}
@Test
void confirmedPublishDeliversNormally() throws Exception {
String target = "worker-confirm-" + System.nanoTime();
@@ -41,20 +41,20 @@ class InMemoryReplyInboxTest {
@Test
void ackRemovesTheMessage() {
inbox.publish("term_a", "m1", "hello");
inbox.ack("term_a", "m1");
assertTrue(inbox.ack("term_a", "m1"), "fleetd #437: ack of a real entry must report true");
assertTrue(inbox.peek("term_a").isEmpty(), "after ack, the message is gone");
}
@Test
void ackForUnknownMsgIdIsNoOp() {
inbox.publish("term_a", "m1", "hello");
inbox.ack("term_a", "no-such-id"); // no-op
assertFalse(inbox.ack("term_a", "no-such-id"), "fleetd #437: a miss must report false"); // no-op
assertEquals(1, inbox.peek("term_a").size(), "the published message is still there");
}
@Test
void ackForUnknownTargetIsNoOp() {
inbox.ack("no-such-target", "m1"); // no-op, should not throw
assertFalse(inbox.ack("no-such-target", "m1"), "fleetd #437: a miss must report false"); // no-op, should not throw
}
@Test
@@ -167,7 +167,7 @@ class InMemoryReplyInboxTest {
@Test
void peekAndAckAreNoOpsForUnownedTarget() {
assertTrue(inbox.peek("term_not_owned").isEmpty());
inbox.ack("term_not_owned", "m1"); // no-op, should not throw
assertFalse(inbox.ack("term_not_owned", "m1"), "fleetd #437: a miss must report false"); // no-op, should not throw
}
@Test
@@ -23,6 +23,7 @@ import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.peer.PeerUnreachableException;
import dev.ltms.fleet.peer.SpawnRequest;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementDecision;
import dev.ltms.fleet.placement.PlacementPolicies;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
@@ -1036,6 +1037,11 @@ class SessionManagerTest {
return delegate.spawn(req);
}
@Override
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
return delegate.spawn(req, decision);
}
@Override
public Set<String> profiles() {
return delegate.profiles();
@@ -1597,6 +1603,11 @@ class SessionManagerTest {
throw new UnsupportedOperationException("not reachable — the capability check refuses first");
}
@Override
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
throw new UnsupportedOperationException("not reachable — the capability check refuses first");
}
@Override
public Set<String> profiles() {
return Set.of("stub-profile");
@@ -1661,6 +1672,11 @@ class SessionManagerTest {
return delegate.spawn(req);
}
@Override
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
return delegate.spawn(req, decision);
}
@Override
public Set<String> profiles() {
return delegate.profiles();
@@ -1867,6 +1883,11 @@ class SessionManagerTest {
return handle;
}
@Override
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
return spawn(req.withProfile(decision.profile()));
}
@Override
public Set<String> profiles() {
return Set.of("lazy");