diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index 76ed1de..ecf73e0 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -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 /.claude} default. + * + *

Follows the same {@code fleet.leaders..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 leadConfigDirLookup(Supplier> profiles, + Map 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 diff --git a/fleetd/src/main/java/dev/ltms/fleet/lead/LeadContextGauge.java b/fleetd/src/main/java/dev/ltms/fleet/lead/LeadContextGauge.java index 8b11260..fc16d01 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/lead/LeadContextGauge.java +++ b/fleetd/src/main/java/dev/ltms/fleet/lead/LeadContextGauge.java @@ -47,13 +47,21 @@ import java.util.function.LongSupplier; * *

Three states, not two (OK / HIGH / UNKNOWN). 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". + * + *

A torn final line does not mean UNKNOWN. {@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 + * every 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 none 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". * *

Bounded cost. {@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); } diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index b475773..f0341db 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -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 /.claude} default. + * + *

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..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 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 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 peers, LeadRollover leadRollover) { + LeadSeatSource leadSeats, LeadConfigDirSource leadConfigDirs, List 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 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 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 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 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 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 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 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> 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 leadView(String terminal, String name, Agent live, - String selfTerm, LeadContextGauge contextGauge) { + String selfTerm, LeadContextGauge contextGauge, + LeadConfigDirSource leadConfigDirs) { Map 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..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 /.claude}). */ - private static Map contextView(LeadContextGauge contextGauge, Agent live) { + private static Map 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 c = new LinkedHashMap<>(); c.put("state", reading.state().name().toLowerCase()); if (reading.tokens() != null) { diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdLeadConfigDirLookupTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdLeadConfigDirLookupTest.java new file mode 100644 index 0000000..f591823 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdLeadConfigDirLookupTest.java @@ -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 /.claude} default (see {@code LeadContextGauge}). + * + *

{@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..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 profiles = Map.of("opus", profileWithConfigDir("opus", "/mnt/opus-claude")); + Map leaders = Map.of("primary", leadOnProfile("opus")); + Function 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> profilesRef = + new java.util.concurrent.atomic.AtomicReference<>( + Map.of("opus", profileWithConfigDir("opus", "/mnt/dir-a"))); + Map leaders = Map.of("primary", leadOnProfile("opus")); + Function 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 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 leaders = Map.of("primary", recogniseOnly); + Function 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 leaders = Map.of("primary", leadOnProfile("ghost-profile")); + Function 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 profiles = Map.of("opus", profileWithConfigDir("opus", null)); + Map leaders = Map.of("primary", leadOnProfile("opus")); + Function 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 lookup = Fleetd.leadConfigDirLookup(Map::of, Map.of()); + + assertNull(lookup.apply("ghost-lead")); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/lead/LeadContextGaugeTest.java b/fleetd/src/test/java/dev/ltms/fleet/lead/LeadContextGaugeTest.java index 505ca35..3eb41be 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/lead/LeadContextGaugeTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/lead/LeadContextGaugeTest.java @@ -23,14 +23,24 @@ import static org.junit.jupiter.api.Assertions.assertTrue; *

    *
  1. {@link #tokenCountTracksTheLastUsageRecordAndChangesWithIt()}
  2. *
  3. {@link #compactionCountTracksCompactBoundaryRecordsAndChangesWithIt()}
  4. - *
  5. {@link #missingFileIsUnknown()}, {@link #unreadableFileIsUnknown()}, - * {@link #truncatedOrInvalidLastLineIsUnknown()}
  6. + *
  7. {@link #missingFileIsUnknown()}, {@link #unreadableFileIsUnknown()}
  8. + *
  9. {@link #tornFinalLineFallsBackToTheLastGoodReading()}, + * {@link #everyLineUnparseableIsUnknown()} (fleetd #602 gauge-wiring, finding 2)
  10. *
  11. {@link #readNeverExceedsTheTailBound()}
  12. *
  13. {@link #secondReadInsideTtlDoesNotTouchDiskAgain()}
  14. *
* *

Every fixture lives under {@code @TempDir} — never the operator's real config directory (see * the ticket's hard constraint on this). + * + *

fleetd #602 gauge-wiring, finding 2. 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()); } diff --git a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpLeadContextGaugeWiringTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpLeadContextGaugeWiringTest.java new file mode 100644 index 0000000..956524c --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpLeadContextGaugeWiringTest.java @@ -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. + * + *

This class proves the directory {@code fleet.leaders..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-}. */ + 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 /projects//.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 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); + } +}