fleetd #602 gauge-wiring: thread a lead's configured configDir into the context gauge

FleetMcp.contextView hardcoded LeadContextGauge.read(null, ...), so a lead whose
profile sets its own CLAUDE_CONFIG_DIR always read the wrong transcript directory
and reported UNKNOWN forever, with no error anywhere.

- Add FleetMcp.LeadConfigDirSource (same idiom as LeadSeatSource) and thread it
  through the constructor / listFleet overload chain / leadView / contextView.
- Add Fleetd.leadConfigDirLookup, wired at construction, following the same
  fleet.leaders.<name>.profile link leadSeatLookup already uses, one step
  further to that profile's own configDir.
- LeadContextGauge.parse: a single unparseable line (typically the final one,
  torn by a write this read raced) is now skipped rather than forcing UNKNOWN;
  only when every line in the read window fails to parse does it report
  UNKNOWN, which is the real format-change signal.
- Tests: FleetdLeadConfigDirLookupTest (lookup logic), FleetMcpLeadContextGaugeWiringTest
  (end-to-end: config naming directory A vs B decides which is read; a lead with
  no configured dir degrades without throwing), and two replacement properties in
  LeadContextGaugeTest for the torn-line fix plus its all-unparseable control.
This commit is contained in:
Dai Ha
2026-09-20 16:19:37 +07:00
parent 3ed7bfca67
commit d345b14e43
6 changed files with 383 additions and 44 deletions
@@ -672,6 +672,10 @@ public final class Fleetd {
leadMailbox,
outageSource,
new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)),
// fleetd #602 gauge-wiring: threads each lead's configured configDir into the
// context gauge — see leadConfigDirLookup's own doc for why this, not a hardcoded
// null, is what fleet_list's context row now reads.
new FleetMcp.LeadConfigDirSource(leadConfigDirLookup(() -> config.get().profiles(), leaders)),
// fleetd #361: the operator-declared peers this daemon's fleet_list should try to
// reach. Read from the SAME snapshot leadMailbox itself opened from (cfg.coordinator()),
// not the live config.get() — coordinator wiring is already boot-time-fixed (see
@@ -1576,6 +1580,37 @@ public final class Fleetd {
};
}
/**
* fleetd #602 gauge-wiring: per-lead-name factory for {@link FleetMcp.LeadConfigDirSource} — the
* {@code CLAUDE_CONFIG_DIR} the named lead's own profile runs on, so {@code LeadContextGauge}
* reads the transcript directory that lead's Claude Code process actually writes to, not always
* the built-in {@code <user.home>/.claude} default.
*
* <p>Follows the same {@code fleet.leaders.<name>.profile} link {@link #leadSeatLookup} already
* uses to find a lead's profile, one step further to that profile's own {@code configDir:}. A
* lead entry that names no {@code profile:}, or whose named profile is not configured, or whose
* profile sets no {@code configDir:} override, returns {@code null} — {@code LeadContextGauge}
* then falls back to its own default, exactly as before this ticket.
*
* @param profiles the live profile map, normally {@code () -> config.get().profiles()} in
* {@code main} — read live, like every other {@code configDir} lookup, so a
* config reload takes effect on the very next {@code fleet_list} call
* @param leaders {@code fleet.leaders}, read once at startup like {@link #leadSeatLookup}'s own
* {@code leaders} parameter — passed as a plain map, never re-read from
* {@code config.get()}
*/
static Function<String, String> leadConfigDirLookup(Supplier<Map<String, FleetConfig.Profile>> profiles,
Map<String, FleetConfig.Leader> leaders) {
return leadName -> {
FleetConfig.Leader lead = leaders.get(leadName);
if (lead == null || lead.profile() == null || lead.profile().isBlank()) {
return null;
}
FleetConfig.Profile leadProfile = profiles.get().get(lead.profile());
return leadProfile == null ? null : leadProfile.configDir();
};
}
/**
* fleetd #248 / fleetd#201 Unit 5: package-private factory for the per-target backend-error
* pattern lookup {@link CompletionResolver} classifies a pane scrape against. Closes over the
@@ -47,13 +47,21 @@ import java.util.function.LongSupplier;
*
* <p><strong>Three states, not two (OK / HIGH / UNKNOWN).</strong> Every path that cannot
* positively establish the live token count — a missing file, an unreadable one, a peer that
* is not a Claude backend, or a last line that fails to parse as JSON — returns
* is not a Claude backend, or every line in the read window failing to parse as JSON — returns
* {@link State#UNKNOWN} with no token number, never a default "0" or "ok" that would read as
* "this lead is fine" when the honest answer is "I could not look". A last-line parse failure in
* particular is treated as UNKNOWN rather than falling back to the previous good line: Claude Code
* appends and flushes one line at a time, so a malformed last line means either a write caught
* mid-flush or a format this reader no longer understands — reporting the stale previous number as
* if it were current would be exactly the false confidence this whole feature exists to avoid.
* "this lead is fine" when the honest answer is "I could not look".
*
* <p><strong>A torn final line does not mean UNKNOWN.</strong> {@code fleet_list} reads this
* transcript while Claude Code may be mid-write on it, so the last line in the window can be cut
* off mid-flush — that is an ordinary, expected race, not a sign the format has changed. Earlier
* this class treated ANY unparseable last line as UNKNOWN, on the theory that "if the format
* changes, we should see UNKNOWN". That reasoning does not hold: a real format change makes
* <em>every</em> line in the window unparseable, not only the last one written. So a single
* malformed line (most often the final, torn one, but the check is not position-specific) is
* simply skipped rather than treated as fatal, and the reading is built from whatever lines in the
* window did parse. Only when <em>none</em> of them parse — the real format-change signal — does
* this return {@link State#UNKNOWN}, still with no stale number standing in for "I could not
* tell".
*
* <p><strong>Bounded cost.</strong> {@code fleet_list} is polled constantly, so every read is
* capped two ways: {@link #TAIL_BYTES} bounds how much of the transcript is ever read from disk
@@ -239,9 +247,10 @@ public final class LeadContextGauge {
/**
* Parses the tail into a {@link Reading}. The first line is dropped unconditionally whenever
* the tail is not the whole file (it starts mid-line, cut by {@link #TAIL_BYTES} — an expected
* artefact of the bound, not a data problem). The window's actual last line, by contrast, must
* parse cleanly: see the "three states" section of the class javadoc for why a bad last line
* means {@link State#UNKNOWN} rather than a fall-back to the previous good line.
* artefact of the bound, not a data problem). Every remaining line is then parsed on a
* best-effort basis: a line that fails to parse (most often the last one, torn by a write this
* read raced — see the "torn final line" section of the class javadoc) is skipped, not fatal.
* Only when none of the remaining lines parse does this report {@link State#UNKNOWN}.
*/
private Reading parse(TailRead tail) {
String text = new String(tail.bytes(), StandardCharsets.UTF_8);
@@ -255,19 +264,19 @@ public final class LeadContextGauge {
if (lines.isEmpty()) {
return Reading.unknown();
}
JsonNode lastNode;
try {
lastNode = MAPPER.readTree(lines.get(lines.size() - 1));
} catch (IOException e) {
return Reading.unknown();
}
Long tokens = null;
int compactions = 0;
for (int i = 0; i < lines.size(); i++) {
JsonNode node = i == lines.size() - 1 ? lastNode : tryParse(lines.get(i));
boolean anyLineParsed = false;
for (String line : lines) {
JsonNode node = tryParse(line);
if (node == null) {
// A malformed line — typically the last one, cut mid-flush by a write this read
// raced — is skipped rather than treated as fatal. See the class javadoc's "torn
// final line" section for why: a real format change makes EVERY line unparseable,
// not only this one, and that case is still caught below by anyLineParsed.
continue;
}
anyLineParsed = true;
JsonNode usage = node.path("message").path("usage");
if (usage.isObject()) {
tokens = usage.path("input_tokens").asLong(0)
@@ -278,6 +287,9 @@ public final class LeadContextGauge {
compactions++;
}
}
if (!anyLineParsed) {
return Reading.unknown();
}
if (tokens == null) {
return new Reading(State.UNKNOWN, null, compactions);
}
@@ -116,6 +116,8 @@ public final class FleetMcp {
private final OutageSource outage;
/** fleetd #176: SEPARATE from both of the above — see {@link LeadSeatSource}'s doc. */
private final LeadSeatSource leadSeats;
/** fleetd #602 gauge-wiring: see {@link LeadConfigDirSource}. */
private final LeadConfigDirSource leadConfigDirs;
/** CB-637: this daemon's lead-to-lead channel; {@code null} when no coordinator is configured. */
private final LeadChannel leadChannel;
/** fleetd #361: {@code coordinator.peers} — see {@link CoordinationSource}. Empty when unset. */
@@ -258,6 +260,29 @@ public final class FleetMcp {
public static LeadSeatSource none() { return new LeadSeatSource(_ -> 0); }
}
/**
* fleetd #602 gauge-wiring: a lead name's configured {@code CLAUDE_CONFIG_DIR} override, fed to
* {@link LeadContextGauge#read} so {@code fleet_list}'s {@code context} row reads the transcript
* directory the lead's own profile actually writes to, not always the built-in
* {@code <user.home>/.claude} default.
*
* <p>The same idiom as {@link LeadSeatSource} — a {@code FleetMcp} constructor field, not a
* lookup {@code contextView} performs itself, because {@code FleetMcp} holds no
* {@link dev.ltms.fleet.config.FleetConfig} and {@code contextView} is {@code static}. See
* {@code Fleetd.leadConfigDirLookup} for the derivation: the same
* {@code fleet.leaders.<name>.profile} link {@link LeadSeatSource} already follows, resolved to
* that profile's own {@code configDir:}.
*
* @param configDirFor lead name → {@code configDir}, or {@code null} when the lead's entry names
* no profile, or that profile sets no {@code configDir} override — either
* way {@link LeadContextGauge} then falls back to its own built-in default,
* exactly as before this ticket
*/
public record LeadConfigDirSource(Function<String, String> configDirFor) {
/** Inert source — every lead reads {@link LeadContextGauge}'s built-in default {@code configDir}. */
public static LeadConfigDirSource none() { return new LeadConfigDirSource(_ -> null); }
}
/**
* fleetd #361: peer-visibility facts for {@code fleet_list}'s {@code coordinator} row — this
* daemon's own {@link LeadChannel} (for its self mailbox state and held messages) plus the
@@ -358,7 +383,7 @@ public final class FleetMcp {
LeadSeatSource leadSeats, List<String> peers, LeadRollover leadRollover) {
this(messages, workers, sessions, identity, presence, primaryRegistry, callers, authorizationMode, metrics,
capacity, healthCoverage, LoopHealthSource.none(), quarantine, leadChannel, outage, leadSeats,
peers, leadRollover);
LeadConfigDirSource.none(), peers, leadRollover);
}
public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
@@ -366,7 +391,8 @@ public final class FleetMcp {
CallerResolver callers, AuthorizationMode authorizationMode, Metrics metrics,
CapacitySource capacity, HealthCoverageSource healthCoverage, LoopHealthSource loopHealth,
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage,
LeadSeatSource leadSeats, List<String> peers, LeadRollover leadRollover) {
LeadSeatSource leadSeats, LeadConfigDirSource leadConfigDirs, List<String> peers,
LeadRollover leadRollover) {
Objects.requireNonNull(callers, "callers");
this.authorizationEnforced = Objects.requireNonNull(authorizationMode, "authorizationMode")
== AuthorizationMode.ENFORCED;
@@ -376,6 +402,7 @@ public final class FleetMcp {
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
this.outage = Objects.requireNonNull(outage, "outage");
this.leadSeats = Objects.requireNonNull(leadSeats, "leadSeats");
this.leadConfigDirs = Objects.requireNonNull(leadConfigDirs, "leadConfigDirs");
this.healthCoverage = healthCoverage;
this.loopHealth = Objects.requireNonNull(loopHealth, "loopHealth");
this.leadRollover = leadRollover;
@@ -510,7 +537,7 @@ public final class FleetMcp {
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_list", Map.of()), null);
if (denied != null) return denied;
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage,
leadSeats, leadContextGauge, callers.leads(),
leadSeats, leadContextGauge, leadConfigDirs, callers.leads(),
callerTerminal(exchange),
new CoordinationSource(leadChannel, peers),
coordinatorVisibleTo(principal(exchange)));
@@ -1613,7 +1640,7 @@ public final class FleetMcp {
Map<String, String> leads, String selfTerm,
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine,
OutageSource.none(), LeadSeatSource.none(), new LeadContextGauge(), leads, selfTerm, coordination, false);
OutageSource.none(), LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
@@ -1622,7 +1649,7 @@ public final class FleetMcp {
QuarantineSource quarantine, OutageSource outage,
Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
LeadSeatSource.none(), new LeadContextGauge(), leads, selfTerm, CoordinationSource.none(), false);
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, CoordinationSource.none(), false);
}
/**
@@ -1639,7 +1666,7 @@ public final class FleetMcp {
QuarantineSource quarantine, Map<String, String> leads, String selfTerm,
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, OutageSource.none(),
LeadSeatSource.none(), new LeadContextGauge(), leads, selfTerm, coordination, false);
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
@@ -1648,7 +1675,7 @@ public final class FleetMcp {
QuarantineSource quarantine, OutageSource outage,
Map<String, String> leads, String selfTerm, CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
LeadSeatSource.none(), new LeadContextGauge(), leads, selfTerm, coordination, false);
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
}
/**
@@ -1670,7 +1697,7 @@ public final class FleetMcp {
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
leadSeats, new LeadContextGauge(), leads, selfTerm, coordination, false);
leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
}
/**
@@ -1694,7 +1721,7 @@ public final class FleetMcp {
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
CoordinationSource coordination, boolean callerIsPrimary) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine,
outage, leadSeats, new LeadContextGauge(), leads, selfTerm, coordination, callerIsPrimary);
outage, leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, callerIsPrimary);
}
/**
@@ -1710,6 +1737,7 @@ public final class FleetMcp {
LoopHealthSource loopHealth,
QuarantineSource quarantine, OutageSource outage,
LeadSeatSource leadSeats, LeadContextGauge contextGauge,
LeadConfigDirSource leadConfigDirs,
Map<String, String> leads, String selfTerm,
CoordinationSource coordination, boolean callerIsPrimary) {
try {
@@ -1719,7 +1747,8 @@ public final class FleetMcp {
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
List<Map<String, Object>> leadRows = leads.entrySet().stream()
.sorted(Map.Entry.comparingByValue())
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm, contextGauge))
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm, contextGauge,
leadConfigDirs))
.toList();
// fleetd #209: this is the caller-driven fleet_list read that actually reports
// agentSessionId (via memberCapacityView -> SessionManager.rosterView), so it uses the
@@ -2031,7 +2060,8 @@ public final class FleetMcp {
* could not tell".
*/
private static Map<String, Object> leadView(String terminal, String name, Agent live,
String selfTerm, LeadContextGauge contextGauge) {
String selfTerm, LeadContextGauge contextGauge,
LeadConfigDirSource leadConfigDirs) {
Map<String, Object> m = new LinkedHashMap<>();
m.put("sessionId", terminal);
m.put("name", name);
@@ -2040,21 +2070,23 @@ public final class FleetMcp {
if (terminal.equals(selfTerm)) {
m.put("self", true);
}
m.put("context", contextView(contextGauge, live));
String configDir = leadConfigDirs.configDirFor().apply(name);
m.put("context", contextView(contextGauge, live, configDir));
return m;
}
/**
* Reads the lead context gauge for one lead. {@code configDir} is always the built-in default
* (see {@link LeadContextGauge#read}) — this daemon does not currently thread a lead's
* configured {@code configDir} override through to {@code fleet_list}, so a lead whose profile
* sets one will read {@code unknown} rather than the wrong file. That is the honest degrade the
* three-state design exists for, not a defect: see this class's addendum note on scope.
* Reads the lead context gauge for one lead. {@code configDir} is this lead's configured
* {@code CLAUDE_CONFIG_DIR} override (see {@link LeadConfigDirSource}), derived from
* {@code fleet.leaders.<name>.profile} → that profile's own {@code configDir:} — or {@code null}
* when the lead's entry names no profile, or that profile sets no override, in which case
* {@link LeadContextGauge#read} falls back to its own built-in default
* ({@code <user.home>/.claude}).
*/
private static Map<String, Object> contextView(LeadContextGauge contextGauge, Agent live) {
private static Map<String, Object> contextView(LeadContextGauge contextGauge, Agent live, String configDir) {
String sessionId = live == null ? null : live.sessionId();
String agentType = live == null ? null : live.agentType();
LeadContextGauge.Reading reading = contextGauge.read(null, sessionId, agentType);
LeadContextGauge.Reading reading = contextGauge.read(configDir, sessionId, agentType);
Map<String, Object> c = new LinkedHashMap<>();
c.put("state", reading.state().name().toLowerCase());
if (reading.tokens() != null) {
@@ -0,0 +1,100 @@
package dev.ltms.fleet;
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.function.Function;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* fleetd #602 gauge-wiring: {@link Fleetd#leadConfigDirLookup} is the factory {@code Fleetd.main}
* wires into {@code FleetMcp.LeadConfigDirSource} so {@code fleet_list}'s {@code context} row reads
* the transcript directory a lead's OWN profile actually writes to, instead of always falling back
* to the built-in {@code <user.home>/.claude} default (see {@code LeadContextGauge}).
*
* <p>{@code FleetMcpLeadContextGaugeWiringTest} proves the directory this factory returns is what
* actually gets read; this class proves the factory's own matching logic — the same
* {@code fleet.leaders.<name>.profile} link {@code Fleetd.leadSeatLookup} already follows (see
* {@code FleetdLeadSeatLookupTest}), one step further to that profile's own {@code configDir:}.
*/
class FleetdLeadConfigDirLookupTest {
private static FleetConfig.Profile profileWithConfigDir(String name, String configDir) {
return new FleetConfig.Profile(name, null, "claude-sonnet-5", configDir, null, null,
"tab", "fleet", "w #{n}", null, null, null, null, null, null, null,
null, 3, true, null, null, null);
}
private static FleetConfig.Leader leadOnProfile(String profile) {
return new FleetConfig.Leader(profile, "lead: primary", 1, "lead:", 10, "claude", "claude-sonnet-5");
}
@Test
@DisplayName("a lead on a profile that sets configDir resolves to that directory")
void leadOnAProfileWithConfigDirResolvesToIt() {
Map<String, FleetConfig.Profile> profiles = Map.of("opus", profileWithConfigDir("opus", "/mnt/opus-claude"));
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
Function<String, String> lookup = Fleetd.leadConfigDirLookup(() -> profiles, leaders);
assertEquals("/mnt/opus-claude", lookup.apply("primary"));
}
@Test
@DisplayName("changing the config from directory A to directory B changes what the lookup reports")
void configChangeFromDirectoryAToDirectoryBChangesTheAnswer() {
java.util.concurrent.atomic.AtomicReference<Map<String, FleetConfig.Profile>> profilesRef =
new java.util.concurrent.atomic.AtomicReference<>(
Map.of("opus", profileWithConfigDir("opus", "/mnt/dir-a")));
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
Function<String, String> lookup = Fleetd.leadConfigDirLookup(profilesRef::get, leaders);
assertEquals("/mnt/dir-a", lookup.apply("primary"), "must read directory A before the config changes");
profilesRef.set(Map.of("opus", profileWithConfigDir("opus", "/mnt/dir-b")));
assertEquals("/mnt/dir-b", lookup.apply("primary"), "must read directory B once the LIVE config changes — "
+ "a lookup that snapshotted the profile map at construction would still answer directory A here");
}
@Test
@DisplayName("a lead entry with no `profile:` (recognise-only) resolves to null, not a thrown exception")
void recogniseOnlyLeadWithNoProfileResolvesToNull() {
Map<String, FleetConfig.Profile> profiles = Map.of("opus", profileWithConfigDir("opus", "/mnt/opus-claude"));
FleetConfig.Leader recogniseOnly = new FleetConfig.Leader(null, "lead: primary", 1, "lead:", 10,
"claude", "claude-sonnet-5");
Map<String, FleetConfig.Leader> leaders = Map.of("primary", recogniseOnly);
Function<String, String> lookup = Fleetd.leadConfigDirLookup(() -> profiles, leaders);
assertNull(lookup.apply("primary"));
}
@Test
@DisplayName("a lead naming a profile that is not configured resolves to null, not a thrown exception")
void leadOnAnUnconfiguredProfileResolvesToNull() {
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("ghost-profile"));
Function<String, String> lookup = Fleetd.leadConfigDirLookup(Map::of, leaders);
assertNull(lookup.apply("primary"));
}
@Test
@DisplayName("a lead on a profile that sets no configDir override resolves to null")
void leadOnAProfileWithNoConfigDirResolvesToNull() {
Map<String, FleetConfig.Profile> profiles = Map.of("opus", profileWithConfigDir("opus", null));
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
Function<String, String> lookup = Fleetd.leadConfigDirLookup(() -> profiles, leaders);
assertNull(lookup.apply("primary"));
}
@Test
@DisplayName("an unrecognised lead name resolves to null, not a thrown exception")
void unrecognisedLeadNameResolvesToNull() {
Function<String, String> lookup = Fleetd.leadConfigDirLookup(Map::of, Map.of());
assertNull(lookup.apply("ghost-lead"));
}
}
@@ -23,14 +23,24 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
* <ol>
* <li>{@link #tokenCountTracksTheLastUsageRecordAndChangesWithIt()}</li>
* <li>{@link #compactionCountTracksCompactBoundaryRecordsAndChangesWithIt()}</li>
* <li>{@link #missingFileIsUnknown()}, {@link #unreadableFileIsUnknown()},
* {@link #truncatedOrInvalidLastLineIsUnknown()}</li>
* <li>{@link #missingFileIsUnknown()}, {@link #unreadableFileIsUnknown()}</li>
* <li>{@link #tornFinalLineFallsBackToTheLastGoodReading()},
* {@link #everyLineUnparseableIsUnknown()} (fleetd #602 gauge-wiring, finding 2)</li>
* <li>{@link #readNeverExceedsTheTailBound()}</li>
* <li>{@link #secondReadInsideTtlDoesNotTouchDiskAgain()}</li>
* </ol>
*
* <p>Every fixture lives under {@code @TempDir} — never the operator's real config directory (see
* the ticket's hard constraint on this).
*
* <p><strong>fleetd #602 gauge-wiring, finding 2.</strong> The property above used to be "a
* truncated or invalid LAST line reports UNKNOWN" — on the theory that a bad last line signals a
* format change. That reasoning did not hold: {@code fleet_list} reads this transcript while Claude
* Code may be mid-write on it, so a torn LAST line is an ordinary race, not a format change, and a
* real format change makes EVERY line unparseable, not only the last one written. So the property
* is now split in two: {@link #tornFinalLineFallsBackToTheLastGoodReading} (the torn-line case must
* NOT destroy a good earlier reading) and its control, {@link #everyLineUnparseableIsUnknown} (only
* when NOTHING in the window parses does this report UNKNOWN).
*/
class LeadContextGaugeTest {
@@ -109,7 +119,7 @@ class LeadContextGaugeTest {
assertEquals(0, zeroCompactions.compactions(), "changing K in the fixture must change the reported count");
}
// --- property 3: three separate UNKNOWN cases -----------------------------------------------
// --- property 3: two separate UNKNOWN cases -------------------------------------------------
@Test
@DisplayName("a missing transcript file reports UNKNOWN with no token number")
@@ -137,15 +147,39 @@ class LeadContextGaugeTest {
}
@Test
@DisplayName("a truncated or invalid last line reports UNKNOWN with no token number")
void truncatedOrInvalidLastLineIsUnknown(@TempDir Path tmp) throws IOException {
@DisplayName("fleetd #602 finding 2: a torn final line does not destroy a good earlier reading")
void tornFinalLineFallsBackToTheLastGoodReading(@TempDir Path tmp) throws IOException {
// Deliberately NOT via writeTranscript: that helper appends a trailing newline after every
// line, including the last one, which would not model what a write caught mid-flush looks
// like on disk. Claude Code appends and flushes one line at a time, so the torn line here has
// no trailing newline at all -- exactly the shape a read racing an in-progress write sees.
Path projectDir = tmp.resolve("projects").resolve("some-project-slug");
Files.createDirectories(projectDir);
Path file = projectDir.resolve(SESSION_ID + ".jsonl");
String lastCompleteLine = usageLine(1_000, 2_000, 3_000); // last COMPLETE record: 6,000
String tornLine = "{\"type\":\"assistant\",\"message\":{\"role\":\"assistant\",\"usage\":{\"input_tok";
Files.writeString(file, lastCompleteLine + "\n" + tornLine, StandardCharsets.UTF_8);
LeadContextGauge gauge = new LeadContextGauge();
LeadContextGauge.Reading reading = gauge.read(tmp.toString(), SESSION_ID, "claude");
assertEquals(LeadContextGauge.State.OK, reading.state(),
"a torn final line must not turn a good earlier reading into UNKNOWN");
assertEquals(6_000L, reading.tokens(),
"must report the last COMPLETE line's total, not fail the whole read");
}
@Test
@DisplayName("fleetd #602 finding 2 control: when EVERY line is unparseable, the reading IS UNKNOWN")
void everyLineUnparseableIsUnknown(@TempDir Path tmp) throws IOException {
// The real signal a format change gives: not one bad line, but ALL of them. Without this
// control, code that never reports UNKNOWN at all would still pass the torn-line property.
String configDir = writeTranscript(tmp, SESSION_ID,
usageLine(1_000, 0, 0),
"{\"type\":\"assistant\",\"message\":{\"role\":\"assistant\",\"usage\":{\"input_tok"); // cut mid-write
"{this is not json at all",
"neither is this{{{");
LeadContextGauge gauge = new LeadContextGauge();
LeadContextGauge.Reading reading = gauge.read(configDir, SESSION_ID, "claude");
assertEquals(LeadContextGauge.State.UNKNOWN, reading.state(),
"a bad LAST line must not silently fall back to the previous good line's number");
"every line unparseable is the real format-change signal and must still report UNKNOWN");
assertNull(reading.tokens());
}
@@ -0,0 +1,126 @@
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.lead.LeadContextGauge;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.session.SessionManager;
import io.modelcontextprotocol.spec.McpSchema;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.atomic.AtomicReference;
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 #602 gauge-wiring: {@code FleetMcp}'s {@code fleet_list} handler shipped calling
* {@code LeadContextGauge.read(null, ...)} unconditionally, so every lead whose profile sets a
* {@code configDir:} override read the wrong transcript directory forever, with no error anywhere —
* {@link LeadContextGauge}'s own 8 tests all passed because none of them exercised THIS wiring; each
* one hands the gauge its own directory directly.
*
* <p>This class proves the directory {@code fleet.leaders.<name>.profile}'s own {@code configDir:}
* names is the one {@code fleet_list}'s {@code context} row actually reads — not merely that SOME
* directory got passed. If the wiring in {@code FleetMcp.leadView}/{@code contextView} is ever
* reverted to a hardcoded {@code null}, {@link #configuredDirectoryDecidesWhichTranscriptIsRead()}
* must go red: both directories below hold a REAL, DIFFERENT token count for the SAME session id,
* so a hardcoded {@code null} (always reading the built-in default, where neither transcript lives)
* would report {@code UNKNOWN} both times instead of the two distinct numbers this test asserts.
*/
class FleetMcpLeadContextGaugeWiringTest {
private static final String LEAD_TERMINAL = "term_lead";
private static final String LEAD_NAME = "opus";
/** {@code FakeHerdr.withAgent} always projects {@code agent_session.value} as {@code sess-<terminalId>}. */
private static final String LEAD_SESSION_ID = "sess-" + LEAD_TERMINAL;
private static ClaudeCodeLauncher workerService(FakeHerdr h) {
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", 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(cfg.profile(), cfg), cfg.profile(), _ -> "tok");
}
private static String textOf(McpSchema.CallToolResult r) {
return ((McpSchema.TextContent) r.content().getFirst()).text();
}
/** One line Claude Code would write for a turn with the given live-context total. */
private static String usageLine(long tokens) {
return "{\"type\":\"assistant\",\"message\":{\"role\":\"assistant\",\"usage\":{"
+ "\"input_tokens\":" + tokens + ",\"cache_read_input_tokens\":0,\"cache_creation_input_tokens\":0}}}";
}
/** Lays out {@code <configDir>/projects/<anySlug>/<sessionId>.jsonl} carrying one usage record. */
private static String writeTranscript(Path configDir, String sessionId, long tokens) throws IOException {
Path projectDir = configDir.resolve("projects").resolve("some-project-slug");
Files.createDirectories(projectDir);
Files.writeString(projectDir.resolve(sessionId + ".jsonl"), usageLine(tokens) + "\n", StandardCharsets.UTF_8);
return configDir.toString();
}
@Test
@DisplayName("the CONFIGURED directory decides which transcript is read, and follows a config change")
void configuredDirectoryDecidesWhichTranscriptIsRead(@TempDir Path tmp) throws IOException {
Path dirA = Files.createDirectory(tmp.resolve("dir-a"));
Path dirB = Files.createDirectory(tmp.resolve("dir-b"));
writeTranscript(dirA, LEAD_SESSION_ID, 11_000);
writeTranscript(dirB, LEAD_SESSION_ID, 22_000);
FakeHerdr herdr = new FakeHerdr().withAgent("lead-opus", LEAD_TERMINAL, "wL:p1", "wL:t1");
SessionManager sessions = new SessionManager(workerService(herdr));
LeadContextGauge contextGauge = new LeadContextGauge();
AtomicReference<String> configuredDir = new AtomicReference<>(dirA.toString());
FleetMcp.LeadConfigDirSource source = new FleetMcp.LeadConfigDirSource(name ->
LEAD_NAME.equals(name) ? configuredDir.get() : null);
String firstRead = textOf(listFleet(herdr, sessions, contextGauge, source));
assertTrue(firstRead.contains("\"tokens\":11000"),
"config names directory A ⇒ fleet_list must report A's token count: " + firstRead);
// The config changes to name directory B instead -- the very next read must follow it.
configuredDir.set(dirB.toString());
String secondRead = textOf(listFleet(herdr, sessions, contextGauge, source));
assertTrue(secondRead.contains("\"tokens\":22000"),
"config now names directory B ⇒ fleet_list must report B's token count, not A's stale one: "
+ secondRead);
}
@Test
@DisplayName("a lead whose config names no directory still degrades to the built-in default, never throws")
void aLeadWhoseConfigNamesNoDirectoryStillFallsBackWithoutThrowing() {
FakeHerdr herdr = new FakeHerdr().withAgent("lead-opus", LEAD_TERMINAL, "wL:p1", "wL:t1");
SessionManager sessions = new SessionManager(workerService(herdr));
LeadContextGauge contextGauge = new LeadContextGauge();
McpSchema.CallToolResult res = listFleet(herdr, sessions, contextGauge, FleetMcp.LeadConfigDirSource.none());
assertNotEquals(Boolean.TRUE, res.isError(), "a lead with no configured configDir must degrade, never throw");
String out = textOf(res);
assertTrue(out.contains("\"state\":\"unknown\""), "no transcript under the built-in default ⇒ UNKNOWN: " + out);
assertFalse(out.contains("\"tokens\""), "UNKNOWN must never carry a stale/default token number: " + out);
}
private static McpSchema.CallToolResult listFleet(FakeHerdr herdr, SessionManager sessions,
LeadContextGauge contextGauge, FleetMcp.LeadConfigDirSource leadConfigDirs) {
return FleetMcp.listFleet(workerService(herdr), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.LoopHealthSource.none(), FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
FleetMcp.LeadSeatSource.none(), contextGauge, leadConfigDirs,
Map.of(LEAD_TERMINAL, LEAD_NAME), LEAD_TERMINAL, FleetMcp.CoordinationSource.none(), false);
}
}