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).
This commit is contained in:
Dai Ha
2026-09-10 19:00:07 +07:00
parent 30d6872779
commit 7772b41993
2 changed files with 251 additions and 53 deletions
+76 -53
View File
@@ -427,64 +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 #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 (right 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
// 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 #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);
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. Built by
// usageLimitFixWarning/usageLimitFixWarningNoModel (extracted fleetd #446 follow-up)
// so a test can pin the exact text without capturing real log output — see those
// methods' javadoc and FleetdUsageLimitFixWarningTest.
String model = profile.model();
log.warn(model != null && !model.isBlank()
? usageLimitFixWarning(profile.profile(), model)
: usageLimitFixWarningNoModel(profile.profile(), cfg.quarantineCooldownSeconds()));
};
// 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);
@@ -832,6 +788,73 @@ public final class Fleetd {
return target -> presence.isPresent(target) || leads.get().containsKey(target);
}
/**
* 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
@@ -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);
}
}