From 3ed7bfca678ed1cd42cd45338bb1c71700dd4024 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 20 Sep 2026 16:05:15 +0700 Subject: [PATCH 1/3] fleetd: report a lead's live context usage in fleet_list fleetd had no way to see how full a lead's context window is. On this host a lead auto-compacted 30 times in one session, discarding roughly 250,000 tokens and costing 46s to 3m16s each time, and nothing could see it coming. LeadContextGauge reads the transcript Claude Code itself writes, never the lead's pane. It finds .jsonl by NAME under /projects/ rather than deriving the project slug, which is an undocumented internal. Three states, not two: OK, HIGH, UNKNOWN. Every path that cannot positively establish a token count reports UNKNOWN with no number, so a lead is never told it is fine when the honest answer is "I could not look". Bounded two ways: TAIL_BYTES caps bytes read off disk, and a 5s cache TTL caps how often that read happens, because fleet_list is polled constantly. Recovered by the lead: the authoring member ended on a backend error (DNS ENOTFOUND) with this work uncommitted and unpushed in its worktree. Verified before committing: mvn -o clean install exit 0, 1813 tests, 0 failures, 0 errors, 0 skipped, 144 reports; LeadContextGaugeTest 8/8. KNOWN INCOMPLETE - see the PR. The fleet_list call site passes configDir=null, which falls back to ~/.claude, but this host's lead profile sets configDir to an override. Measured: 19 transcripts under the real configDir, 0 under the fallback. The gauge is therefore INERT on this fleet until that is wired. --- .../dev/ltms/fleet/lead/LeadContextGauge.java | 295 ++++++++++++++++++ .../java/dev/ltms/fleet/mcp/FleetMcp.java | 73 ++++- .../ltms/fleet/lead/LeadContextGaugeTest.java | 204 ++++++++++++ 3 files changed, 561 insertions(+), 11 deletions(-) create mode 100644 fleetd/src/main/java/dev/ltms/fleet/lead/LeadContextGauge.java create mode 100644 fleetd/src/test/java/dev/ltms/fleet/lead/LeadContextGaugeTest.java diff --git a/fleetd/src/main/java/dev/ltms/fleet/lead/LeadContextGauge.java b/fleetd/src/main/java/dev/ltms/fleet/lead/LeadContextGauge.java new file mode 100644 index 0000000..8b11260 --- /dev/null +++ b/fleetd/src/main/java/dev/ltms/fleet/lead/LeadContextGauge.java @@ -0,0 +1,295 @@ +package dev.ltms.fleet.lead; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; + +import java.io.IOException; +import java.io.RandomAccessFile; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.LongSupplier; + +/** + * Reads how full a lead's own Claude Code context window is, from the transcript Claude Code + * itself writes — never from the lead's pane (fleetd has {@code AgentControl.read} for that, and + * this must not use it: a pane holds terminal text, not the structured usage numbers a transcript + * carries, and scraping it would also race the lead's own rendering). + * + *

Why this exists. A lead auto-compacts when its context fills — on the host + * this was built for, that happened 30 times in one session, discarding roughly 250,000 tokens + * and costing 46 seconds to 3 minutes each time, and fleetd had no way to see it coming. This + * class is the first thing that looks. + * + *

The route. Claude Code appends one JSON object per line to + * {@code /projects//.jsonl}. {@code } is an undocumented, + * internal encoding of the working directory — this class never derives it. Instead it lists the + * one-level-deep subdirectories of {@code /projects/} and looks for + * {@code .jsonl} by name, so the slug rule can change without breaking this reader. + * + *

    + *
  • Live context is read off the last record in the read window that carries + * a {@code message.usage} object: {@code input_tokens + cache_read_input_tokens + + * cache_creation_input_tokens}. This is what actually fills the window — a plain + * {@code input_tokens} count alone understates it once the conversation has any cached + * prefix, which on a long-lived lead is always.
  • + *
  • Compaction history is a count of {@code subtype: "compact_boundary"} + * records seen in the same read window — see {@link Reading#compactions()}. It is a count + * within the window this reader actually looked at, not a lifetime total: a session with + * more compactions than fit in {@link #TAIL_BYTES} of transcript will undercount. That + * trade-off is deliberate — see {@link #TAIL_BYTES}.
  • + *
+ * + *

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 + * {@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. + * + *

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 + * (never the whole 52 MB a long-lived transcript reaches on the host this was measured on), and + * {@link #DEFAULT_CACHE_TTL_MILLIS} bounds how often that bounded read actually happens — a burst + * of {@code fleet_list} calls inside one TTL window reads the file once. One instance's cache is + * keyed by {@code (configDir, sessionId)}, so it is safe to share across every lead a single + * {@code fleet_list} call reports on. + */ +public final class LeadContextGauge { + + private static final String PROJECTS_DIR = "projects"; + private static final ObjectMapper MAPPER = new ObjectMapper(); + + /** + * How many trailing bytes of a transcript a single read ever pulls off disk. Chosen so one + * read comfortably spans many recent turns — each usage or compact_boundary record is at most + * a few KB — while staying nowhere near the 52 MB a long session's real transcript reaches on + * the host this was built for; reading that whole file on every {@code fleet_list} call is + * exactly the cost this bound exists to avoid. 2 MiB holds on the order of hundreds of recent + * lines even when a turn's tool output is unusually large, which is far more than needed to + * find the most recent usage record and any recent compaction. + */ + static final int TAIL_BYTES = 2 * 1024 * 1024; + + /** + * How long a {@link Reading} is served from cache before the file is read again. + * {@code fleet_list} is called constantly (by design — it is the fleet's own status probe), so + * without a TTL a burst of calls would re-read the transcript tail once per call. 5 seconds is + * short enough that a caller watching for a state change never waits long, and long enough that + * a poll loop calling every second or two only touches disk once per window. + */ + static final long DEFAULT_CACHE_TTL_MILLIS = 5_000; + + /** + * Live tokens at or above this count report {@link State#HIGH}. On the host this was measured + * on, auto-compaction actually fires around 267,000–270,000 tokens, but the point of a HIGH + * state is to warn before that happens, not at it — 200,000 is the standard Claude context + * window size and a sensible built-in default: no config key is required to pick it, and a + * lead crossing it is already deep enough into its window that a compaction is foreseeable. + */ + static final long HIGH_THRESHOLD_TOKENS = 200_000; + + /** The only peer kind this reader understands ({@code Agent.agentType()}'s wire value). */ + private static final String CLAUDE_AGENT_TYPE = "claude"; + + public enum State { OK, HIGH, UNKNOWN } + + /** + * @param state {@link State#UNKNOWN} whenever {@code tokens} could not be established + * @param tokens live context tokens, or {@code null} exactly when {@code state} is + * {@link State#UNKNOWN} + * @param compactions {@code compact_boundary} records seen in the read window (see class + * javadoc) — {@code 0} both for "genuinely none seen" and for "unknown", + * since a caller that already sees {@code state: UNKNOWN} has no reason to + * trust this number either way + */ + public record Reading(State state, Long tokens, int compactions) { + static Reading unknown() { + return new Reading(State.UNKNOWN, null, 0); + } + } + + private record CacheEntry(Reading reading, long readAtMillis) { + } + + private final LongSupplier clock; + private final long ttlMillis; + private final Map cache = new ConcurrentHashMap<>(); + /** Test seam only (package-private) — counts real disk reads, i.e. cache misses. */ + private final AtomicInteger diskReads = new AtomicInteger(); + + public LeadContextGauge() { + this(System::currentTimeMillis, DEFAULT_CACHE_TTL_MILLIS); + } + + /** Test seam: an injectable clock and TTL so cache expiry is provable without sleeping. */ + LeadContextGauge(LongSupplier clock, long ttlMillis) { + this.clock = clock; + this.ttlMillis = ttlMillis; + } + + /** How many times this instance has actually read a transcript off disk — test seam only. */ + int diskReadCount() { + return diskReads.get(); + } + + /** + * @param configDir the lead's {@code CLAUDE_CONFIG_DIR}, or {@code null}/blank to use the + * default {@code /.claude} — the right answer for the common case + * where the lead's profile sets no {@code configDir} override + * @param sessionId the lead's own Claude session id ({@code Agent.sessionId()}), or + * {@code null} when herdr has not resolved one yet + * @param agentType the detected peer kind ({@code Agent.agentType()}); anything other than + * {@code "claude"} (including {@code null}, meaning undetected) reports + * {@link State#UNKNOWN} — this reader only understands Claude Code's own + * transcript format + */ + public Reading read(String configDir, String sessionId, String agentType) { + if (sessionId == null || sessionId.isBlank()) { + return Reading.unknown(); + } + if (!CLAUDE_AGENT_TYPE.equalsIgnoreCase(agentType)) { + return Reading.unknown(); + } + String base = (configDir == null || configDir.isBlank()) + ? System.getProperty("user.home") + "/.claude" + : configDir; + String cacheKey = base + '\u0000' + sessionId; + long now = clock.getAsLong(); + CacheEntry cached = cache.get(cacheKey); + if (cached != null && now - cached.readAtMillis() < ttlMillis) { + return cached.reading(); + } + Reading fresh = readUncached(base, sessionId); + cache.put(cacheKey, new CacheEntry(fresh, now)); + return fresh; + } + + private Reading readUncached(String base, String sessionId) { + diskReads.incrementAndGet(); + Path file = findTranscript(base, sessionId); + if (file == null) { + return Reading.unknown(); + } + TailRead tail; + try { + tail = tailBytes(file, TAIL_BYTES); + } catch (IOException e) { + return Reading.unknown(); + } + return parse(tail); + } + + /** + * Finds {@code .jsonl} under {@code /projects/}, one level deep — never by + * deriving the slug directory from a working directory (see class javadoc). Bounded to a + * single {@code list()} of {@code projects/} itself: it never recurses further, so the cost is + * the number of project directories, not the size of any transcript inside them. + */ + private Path findTranscript(String base, String sessionId) { + Path projectsDir = Path.of(base, PROJECTS_DIR); + if (!Files.isDirectory(projectsDir)) { + return null; + } + String filename = sessionId + ".jsonl"; + Path direct = projectsDir.resolve(filename); + if (Files.isRegularFile(direct)) { + return direct; + } + try (var children = Files.list(projectsDir)) { + return children.filter(Files::isDirectory) + .map(dir -> dir.resolve(filename)) + .filter(Files::isRegularFile) + .findFirst() + .orElse(null); + } catch (IOException e) { + return null; + } + } + + /** Package-private (not {@code private}): {@link #tailBytes} is a test seam, see its javadoc. */ + record TailRead(byte[] bytes, boolean fromStart) { + } + + /** + * Reads at most {@code maxBytes} trailing bytes of {@code file}. Package-private (not + * {@code private}) so a test can assert directly on the returned array's length — "bytes + * actually read", not on any parsed answer — without needing a file anywhere near + * {@link #TAIL_BYTES} in size to prove the cap holds. + */ + static TailRead tailBytes(Path file, int maxBytes) throws IOException { + try (RandomAccessFile raf = new RandomAccessFile(file.toFile(), "r")) { + long length = raf.length(); + long start = Math.max(0, length - maxBytes); + raf.seek(start); + byte[] buf = new byte[(int) (length - start)]; + raf.readFully(buf); + return new TailRead(buf, start == 0); + } + } + + /** + * 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. + */ + private Reading parse(TailRead tail) { + String text = new String(tail.bytes(), StandardCharsets.UTF_8); + List lines = new ArrayList<>(List.of(text.split("\n", -1))); + if (!lines.isEmpty() && lines.get(lines.size() - 1).isEmpty()) { + lines.remove(lines.size() - 1); // trailing newline leaves a phantom empty last element + } + if (!tail.fromStart() && !lines.isEmpty()) { + lines.remove(0); // first line is a fragment cut by our own tail bound, not real data + } + 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)); + if (node == null) { + continue; + } + JsonNode usage = node.path("message").path("usage"); + if (usage.isObject()) { + tokens = usage.path("input_tokens").asLong(0) + + usage.path("cache_read_input_tokens").asLong(0) + + usage.path("cache_creation_input_tokens").asLong(0); + } + if ("compact_boundary".equals(node.path("subtype").asText(null))) { + compactions++; + } + } + if (tokens == null) { + return new Reading(State.UNKNOWN, null, compactions); + } + State state = tokens >= HIGH_THRESHOLD_TOKENS ? State.HIGH : State.OK; + return new Reading(state, tokens, compactions); + } + + private JsonNode tryParse(String line) { + try { + return MAPPER.readTree(line); + } catch (IOException e) { + return null; + } + } +} 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 6fa0559..b475773 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -12,6 +12,7 @@ import dev.ltms.fleet.metrics.Metrics; import dev.ltms.fleet.inject.MemberPresence; import dev.ltms.fleet.inject.CompletionResolver; import dev.ltms.fleet.herdr.HerdrException; +import dev.ltms.fleet.lead.LeadContextGauge; import dev.ltms.fleet.lead.LeadRollover; import dev.ltms.fleet.msg.LeadChannel; import dev.ltms.fleet.msg.LeadMessage; @@ -126,6 +127,17 @@ public final class FleetMcp { * clean {@code NOT_CONFIGURED} refusal rather than throwing. See {@link #handover}. */ private final LeadRollover leadRollover; + /** + * "Lead context gauge": how full each lead's own Claude Code context window is, reported on + * {@code fleet_list}'s {@code leads} rows (see {@link #contextView}). Built unconditionally, in + * the field initializer rather than a constructor parameter — this is not a togglable feature + * with an on/off config knob the way {@link OutageSource}/{@link LeadSeatSource} are: it needs + * no config at all (see {@link LeadContextGauge}'s own javadoc for the built-in defaults), so + * there is no "off" value to thread through every existing constructor call site. One instance + * per daemon so its read cache (keyed by session, TTL'd) is actually shared across + * {@code fleet_list} calls rather than rebuilt — and therefore useless — on every call. + */ + private final LeadContextGauge leadContextGauge = new LeadContextGauge(); /** Capacity facts used by {@code fleet_list}; production must supply the placement live count. */ public record CapacitySource(Function liveCount, Function maxLoad, @@ -498,7 +510,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, callers.leads(), + leadSeats, leadContextGauge, callers.leads(), callerTerminal(exchange), new CoordinationSource(leadChannel, peers), coordinatorVisibleTo(principal(exchange))); @@ -1601,7 +1613,7 @@ public final class FleetMcp { Map leads, String selfTerm, CoordinationSource coordination) { return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, - OutageSource.none(), LeadSeatSource.none(), leads, selfTerm, coordination, false); + OutageSource.none(), LeadSeatSource.none(), new LeadContextGauge(), leads, selfTerm, coordination, false); } /** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */ @@ -1610,7 +1622,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(), leads, selfTerm, CoordinationSource.none(), false); + LeadSeatSource.none(), new LeadContextGauge(), leads, selfTerm, CoordinationSource.none(), false); } /** @@ -1627,7 +1639,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(), leads, selfTerm, coordination, false); + LeadSeatSource.none(), new LeadContextGauge(), leads, selfTerm, coordination, false); } /** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */ @@ -1636,7 +1648,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(), leads, selfTerm, coordination, false); + LeadSeatSource.none(), new LeadContextGauge(), leads, selfTerm, coordination, false); } /** @@ -1658,7 +1670,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, leads, selfTerm, coordination, false); + leadSeats, new LeadContextGauge(), leads, selfTerm, coordination, false); } /** @@ -1682,14 +1694,23 @@ 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, leads, selfTerm, coordination, callerIsPrimary); + outage, leadSeats, new LeadContextGauge(), leads, selfTerm, coordination, callerIsPrimary); } + /** + * The canonical implementation. {@code contextGauge} is the "lead context gauge" (see + * {@link LeadContextGauge}) — every wrapper overload above passes a freshly constructed one, + * which is correct for them (none of them exercise repeated calls where a shared cache would + * matter); the one caller that matters for caching, {@code fleet_list}'s MCP handler, passes + * its own single long-lived instance instead (see {@code FleetMcp}'s {@code leadContextGauge} + * field). + */ static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages, CapacitySource capacity, HealthCoverageSource healthCoverage, LoopHealthSource loopHealth, QuarantineSource quarantine, OutageSource outage, - LeadSeatSource leadSeats, Map leads, String selfTerm, + LeadSeatSource leadSeats, LeadContextGauge contextGauge, + Map leads, String selfTerm, CoordinationSource coordination, boolean callerIsPrimary) { try { Map live = workers.list().stream() @@ -1698,7 +1719,7 @@ 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)) + .map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm, contextGauge)) .toList(); // fleetd #209: this is the caller-driven fleet_list read that actually reports // agentSessionId (via memberCapacityView -> SessionManager.rosterView), so it uses the @@ -1993,15 +2014,24 @@ public final class FleetMcp { } /** - * One lead's row: its address, its name, and whether it can be reached right now. + * One lead's row: its address, its name, whether it can be reached right now, and how full its + * own Claude Code context window is. * *

{@code status} is herdr's live view, and {@code unknown} when herdr is not tracking that * pane as an agent — the honest answer, and the one that matters: a lead whose pane herdr cannot * see is a lead a {@code fleet_send} cannot be typed into. It is reported rather than hidden, * because a peer that has gone unreachable is exactly what the sender needs to know. + * + *

{@code context} is the "lead context gauge" (fleetd's context-usage visibility ticket): + * {@code {state: "ok"|"high"|"unknown", tokens?: number, compactions: number}}, read from the + * lead's transcript — see {@link LeadContextGauge}. Kept small on purpose (a token count, a + * state, and a compaction count) rather than echoing the whole reading history: this is a + * roster row a caller glances at, not a diagnostics dump. {@code tokens} is present only when + * {@code state} is not {@code "unknown"} — never a stale or default number standing in for "I + * could not tell". */ private static Map leadView(String terminal, String name, Agent live, - String selfTerm) { + String selfTerm, LeadContextGauge contextGauge) { Map m = new LinkedHashMap<>(); m.put("sessionId", terminal); m.put("name", name); @@ -2010,9 +2040,30 @@ public final class FleetMcp { if (terminal.equals(selfTerm)) { m.put("self", true); } + m.put("context", contextView(contextGauge, live)); 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. + */ + private static Map contextView(LeadContextGauge contextGauge, Agent live) { + String sessionId = live == null ? null : live.sessionId(); + String agentType = live == null ? null : live.agentType(); + LeadContextGauge.Reading reading = contextGauge.read(null, sessionId, agentType); + Map c = new LinkedHashMap<>(); + c.put("state", reading.state().name().toLowerCase()); + if (reading.tokens() != null) { + c.put("tokens", reading.tokens()); + } + c.put("compactions", reading.compactions()); + return c; + } + /** {@code fleet_stop}: tear a worker down by its pane id. */ static McpSchema.CallToolResult stop(SessionManager sessions, String paneId) { if (isBlank(paneId)) { diff --git a/fleetd/src/test/java/dev/ltms/fleet/lead/LeadContextGaugeTest.java b/fleetd/src/test/java/dev/ltms/fleet/lead/LeadContextGaugeTest.java new file mode 100644 index 0000000..505ca35 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/lead/LeadContextGaugeTest.java @@ -0,0 +1,204 @@ +package dev.ltms.fleet.lead; + +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.io.RandomAccessFile; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.concurrent.atomic.AtomicLong; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Ticket "lead context gauge" — fleetd could not see how full a lead's own Claude Code context + * window was, so a lead auto-compacting was always a surprise. {@link LeadContextGauge} reads that + * off the lead's own transcript. Five acceptance properties from the ticket, one test class each + * (plus the three separate UNKNOWN cases the ticket calls out by name): + *

    + *
  1. {@link #tokenCountTracksTheLastUsageRecordAndChangesWithIt()}
  2. + *
  3. {@link #compactionCountTracksCompactBoundaryRecordsAndChangesWithIt()}
  4. + *
  5. {@link #missingFileIsUnknown()}, {@link #unreadableFileIsUnknown()}, + * {@link #truncatedOrInvalidLastLineIsUnknown()}
  6. + *
  7. {@link #readNeverExceedsTheTailBound()}
  8. + *
  9. {@link #secondReadInsideTtlDoesNotTouchDiskAgain()}
  10. + *
+ * + *

Every fixture lives under {@code @TempDir} — never the operator's real config directory (see + * the ticket's hard constraint on this). + */ +class LeadContextGaugeTest { + + private static final String SESSION_ID = "11111111-1111-1111-1111-111111111111"; + + /** One line Claude Code would write for a turn with the given live-context total. */ + private static String usageLine(long inputTokens, long cacheRead, long cacheCreation) { + return "{\"type\":\"assistant\",\"message\":{\"role\":\"assistant\",\"usage\":{" + + "\"input_tokens\":" + inputTokens + "," + + "\"cache_read_input_tokens\":" + cacheRead + "," + + "\"cache_creation_input_tokens\":" + cacheCreation + "}}}"; + } + + /** One line Claude Code writes when an auto-compaction happens. */ + private static String compactionLine() { + return "{\"type\":\"system\",\"subtype\":\"compact_boundary\"," + + "\"compactMetadata\":{\"preTokens\":250000,\"postTokens\":5000," + + "\"trigger\":\"auto\",\"durationMs\":54000}}"; + } + + /** Lays out {@code /projects//.jsonl} and writes {@code lines}. */ + private static String writeTranscript(Path configDir, String sessionId, String... lines) throws IOException { + Path projectDir = configDir.resolve("projects").resolve("some-project-slug"); + Files.createDirectories(projectDir); + Path file = projectDir.resolve(sessionId + ".jsonl"); + StringBuilder sb = new StringBuilder(); + for (String line : lines) { + sb.append(line).append('\n'); + } + Files.writeString(file, sb.toString(), StandardCharsets.UTF_8); + return configDir.toString(); + } + + // --- property 1: live token count tracks the last usage record ----------------------------- + + @Test + @DisplayName("reported tokens equal the LAST usage record's total, and change when it does") + void tokenCountTracksTheLastUsageRecordAndChangesWithIt(@TempDir Path tmp) throws IOException { + String configDir = writeTranscript(tmp, SESSION_ID, + usageLine(1_000, 0, 0), + usageLine(40_000, 5_000, 3_000)); // last record: 48,000 + LeadContextGauge gauge = new LeadContextGauge(); + + LeadContextGauge.Reading first = gauge.read(configDir, SESSION_ID, "claude"); + assertEquals(48_000L, first.tokens(), "must total input+cache_read+cache_creation of the LAST usage record"); + assertEquals(LeadContextGauge.State.OK, first.state()); + + // Change N in the fixture (a fresh session id avoids the cache) -- the reported number + // must change with it, not stay pinned to the first fixture's total. + String otherSession = "22222222-2222-2222-2222-222222222222"; + writeTranscript(tmp, otherSession, usageLine(100_000, 50_000, 50_000)); // last record: 200,000 + LeadContextGauge.Reading second = gauge.read(configDir, otherSession, "claude"); + assertEquals(200_000L, second.tokens()); + assertTrue(second.tokens() != first.tokens(), "changing N in the fixture must change the reported number"); + } + + // --- property 2: compaction count tracks compact_boundary records -------------------------- + + @Test + @DisplayName("reported compaction count equals K compact_boundary records, and changes with K") + void compactionCountTracksCompactBoundaryRecordsAndChangesWithIt(@TempDir Path tmp) throws IOException { + String sessionTwoCompactions = "33333333-3333-3333-3333-333333333333"; + writeTranscript(tmp, sessionTwoCompactions, + usageLine(1_000, 0, 0), + compactionLine(), + usageLine(2_000, 0, 0), + compactionLine(), + usageLine(3_000, 0, 0)); + LeadContextGauge gauge = new LeadContextGauge(); + LeadContextGauge.Reading twoCompactions = gauge.read(tmp.toString(), sessionTwoCompactions, "claude"); + assertEquals(2, twoCompactions.compactions()); + + String sessionZeroCompactions = "44444444-4444-4444-4444-444444444444"; + writeTranscript(tmp, sessionZeroCompactions, usageLine(3_000, 0, 0)); + LeadContextGauge.Reading zeroCompactions = gauge.read(tmp.toString(), sessionZeroCompactions, "claude"); + assertEquals(0, zeroCompactions.compactions(), "changing K in the fixture must change the reported count"); + } + + // --- property 3: three separate UNKNOWN cases ----------------------------------------------- + + @Test + @DisplayName("a missing transcript file reports UNKNOWN with no token number") + void missingFileIsUnknown(@TempDir Path tmp) { + LeadContextGauge gauge = new LeadContextGauge(); + LeadContextGauge.Reading reading = gauge.read(tmp.toString(), SESSION_ID, "claude"); + assertEquals(LeadContextGauge.State.UNKNOWN, reading.state()); + assertNull(reading.tokens()); + } + + @Test + @DisplayName("an unreadable transcript file reports UNKNOWN with no token number") + void unreadableFileIsUnknown(@TempDir Path tmp) throws IOException { + String configDir = writeTranscript(tmp, SESSION_ID, usageLine(1_000, 0, 0)); + Path file = tmp.resolve("projects").resolve("some-project-slug").resolve(SESSION_ID + ".jsonl"); + assertTrue(file.toFile().setReadable(false), "test setup: must be able to revoke read permission"); + try { + LeadContextGauge gauge = new LeadContextGauge(); + LeadContextGauge.Reading reading = gauge.read(configDir, SESSION_ID, "claude"); + assertEquals(LeadContextGauge.State.UNKNOWN, reading.state()); + assertNull(reading.tokens()); + } finally { + file.toFile().setReadable(true); // so @TempDir cleanup can delete it + } + } + + @Test + @DisplayName("a truncated or invalid last line reports UNKNOWN with no token number") + void truncatedOrInvalidLastLineIsUnknown(@TempDir Path tmp) throws IOException { + String configDir = writeTranscript(tmp, SESSION_ID, + usageLine(1_000, 0, 0), + "{\"type\":\"assistant\",\"message\":{\"role\":\"assistant\",\"usage\":{\"input_tok"); // cut mid-write + 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"); + assertNull(reading.tokens()); + } + + // --- property 4: the tail bound is actually enforced ---------------------------------------- + + @Test + @DisplayName("reading a transcript far larger than the tail bound reads no more than the bound") + void readNeverExceedsTheTailBound(@TempDir Path tmp) throws IOException { + int smallBound = 4_096; // exercise the mechanism without writing a multi-MB fixture + Path file = tmp.resolve("big.jsonl"); + // One line far bigger than smallBound, repeated, so the whole file is many times the bound. + String line = usageLine(42, 0, 0) + " ".repeat(500); + StringBuilder sb = new StringBuilder(); + for (int i = 0; i < 50; i++) { + sb.append(line).append('\n'); + } + Files.writeString(file, sb.toString(), StandardCharsets.UTF_8); + long fileLength = Files.size(file); + assertTrue(fileLength > (long) smallBound * 5, "test setup: fixture must genuinely dwarf the bound"); + + var tail = LeadContextGauge.tailBytes(file, smallBound); + assertEquals(smallBound, tail.bytes().length, + "must read exactly the bound, not the whole " + fileLength + "-byte file — assert on bytes actually read"); + } + + // --- property 5: the cache TTL means a burst of calls reads the file once ------------------- + + @Test + @DisplayName("the same (configDir, sessionId) read twice inside the TTL touches disk once") + void secondReadInsideTtlDoesNotTouchDiskAgain(@TempDir Path tmp) throws IOException { + String configDir = writeTranscript(tmp, SESSION_ID, usageLine(1_000, 0, 0)); + AtomicLong now = new AtomicLong(0); + LeadContextGauge gauge = new LeadContextGauge(now::get, 5_000); + + gauge.read(configDir, SESSION_ID, "claude"); + gauge.read(configDir, SESSION_ID, "claude"); // still inside the TTL window + assertEquals(1, gauge.diskReadCount(), "two reads inside the TTL must touch disk once"); + + now.set(6_000); // past the TTL + gauge.read(configDir, SESSION_ID, "claude"); + assertEquals(2, gauge.diskReadCount(), "a read past the TTL must touch disk again"); + } + + // --- non-Claude peer / unresolved session ----------------------------------------------------- + + @Test + @DisplayName("a non-claude agent type or a null session id is UNKNOWN, never OK") + void nonClaudePeerOrUnresolvedSessionIsUnknown(@TempDir Path tmp) throws IOException { + String configDir = writeTranscript(tmp, SESSION_ID, usageLine(1_000, 0, 0)); + LeadContextGauge gauge = new LeadContextGauge(); + + assertEquals(LeadContextGauge.State.UNKNOWN, gauge.read(configDir, SESSION_ID, "opencode").state()); + assertEquals(LeadContextGauge.State.UNKNOWN, gauge.read(configDir, SESSION_ID, null).state()); + assertEquals(LeadContextGauge.State.UNKNOWN, gauge.read(configDir, null, "claude").state()); + } +} -- 2.52.0 From d345b14e43bd190907a9632e2c984b8cf7291ce5 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 20 Sep 2026 16:19:37 +0700 Subject: [PATCH 2/3] 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..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. --- .../src/main/java/dev/ltms/fleet/Fleetd.java | 35 +++++ .../dev/ltms/fleet/lead/LeadContextGauge.java | 46 ++++--- .../java/dev/ltms/fleet/mcp/FleetMcp.java | 70 +++++++--- .../fleet/FleetdLeadConfigDirLookupTest.java | 100 ++++++++++++++ .../ltms/fleet/lead/LeadContextGaugeTest.java | 50 +++++-- .../FleetMcpLeadContextGaugeWiringTest.java | 126 ++++++++++++++++++ 6 files changed, 383 insertions(+), 44 deletions(-) create mode 100644 fleetd/src/test/java/dev/ltms/fleet/FleetdLeadConfigDirLookupTest.java create mode 100644 fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpLeadContextGaugeWiringTest.java 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); + } +} -- 2.52.0 From 4e27bde2d7f8dacc2cdc3f7d884ae1c3f662404e Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 20 Sep 2026 16:34:49 +0700 Subject: [PATCH 3/3] fleetd #602 gauge-wiring follow-up: pin Fleetd.main's LeadConfigDirSource wiring Extract the inline new FleetMcp.LeadConfigDirSource(leadConfigDirLookup(...)) construction in Fleetd.main into a package-private factory, Fleetd.leadConfigDirSource, mirroring loopHealthSource/capacitySource/ healthCoverageSource. Add FleetdLeadConfigDirSourceWiringTest, which calls the factory directly with real Profile/Leader fixtures and asserts the returned source resolves a real configDir -- a property that is false if the factory's body is mutated to return LeadConfigDirSource.none(). Neither FleetMcpLeadContextGaugeWiringTest nor FleetdLeadConfigDirLookupTest could catch main losing this wiring: each builds its own instance instead of calling what main calls. This closes that gap at the factory level, matching the standard already accepted for loopHealthSource's own wiring test. --- .../src/main/java/dev/ltms/fleet/Fleetd.java | 25 ++++- .../FleetdLeadConfigDirSourceWiringTest.java | 98 +++++++++++++++++++ 2 files changed, 121 insertions(+), 2 deletions(-) create mode 100644 fleetd/src/test/java/dev/ltms/fleet/FleetdLeadConfigDirSourceWiringTest.java diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index ecf73e0..d06c01f 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -673,9 +673,9 @@ public final class Fleetd { 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 + // context gauge — see leadConfigDirSource'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)), + leadConfigDirSource(() -> 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 @@ -1611,6 +1611,27 @@ public final class Fleetd { }; } + /** + * fleetd #602 gauge-wiring follow-up (PR #606 review comment 17353): {@code main} used to build + * {@code new FleetMcp.LeadConfigDirSource(leadConfigDirLookup(...))} inline, with nothing a test + * could call directly. Measured on that shape: replacing the whole expression with {@code + * FleetMcp.LeadConfigDirSource.none()} at the call site compiled with 0 errors and left the full + * 1822-test suite green — the daemon could be changed to always report every lead's context as + * {@code UNKNOWN}, forever, while every test stayed green. That is the same hand-built-vs-wired + * shape as fleetd #561/#248/#426/#562 ({@link #loopHealthSource}). + * + *

The fix extracts the inline {@code new} into this factory, in the same style as {@link + * #loopHealthSource}/{@link #capacitySource}/{@link #healthCoverageSource} — which is exactly + * what makes it directly callable from {@code FleetdLeadConfigDirSourceWiringTest}. That test + * calls this factory with real {@link FleetConfig.Profile}/{@link FleetConfig.Leader} fixtures and + * asserts the returned source resolves a real {@code configDir} — a property that would be false + * if this method's body were mutated to {@code return FleetMcp.LeadConfigDirSource.none();}. + */ + static FleetMcp.LeadConfigDirSource leadConfigDirSource(Supplier> profiles, + Map leaders) { + return new FleetMcp.LeadConfigDirSource(leadConfigDirLookup(profiles, leaders)); + } + /** * 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/test/java/dev/ltms/fleet/FleetdLeadConfigDirSourceWiringTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdLeadConfigDirSourceWiringTest.java new file mode 100644 index 0000000..9f2872a --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdLeadConfigDirSourceWiringTest.java @@ -0,0 +1,98 @@ +package dev.ltms.fleet; + +import dev.ltms.fleet.config.FleetConfig; +import dev.ltms.fleet.mcp.FleetMcp; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; + +/** + * fleetd #602 gauge-wiring follow-up (PR #606 review comment 17353): {@code Fleetd.main}'s {@code + * LeadConfigDirSource} local used to be a bare {@code new FleetMcp.LeadConfigDirSource( + * leadConfigDirLookup(...))} built inline, with nothing a test could call directly. Measured on + * that shape: replacing the whole expression with {@code FleetMcp.LeadConfigDirSource.none()} at + * the call site compiled with 0 errors and left the full 1822-test suite green — the daemon could + * be changed to always report every lead's context as {@code UNKNOWN}, forever, and no test would + * notice. That is the same hand-built-vs-config-wired shape as fleetd #561/#248/#426/#562 + * ({@code FleetdLoopHealthSourceWiringTest}). + * + *

{@code FleetMcpLeadContextGaugeWiringTest} and {@code FleetdLeadConfigDirLookupTest} both + * predate this class and are both still correct — but neither can catch the mutation above. One + * builds its own {@code FleetMcp} and hands it its own {@code LeadConfigDirSource}; the other + * builds its own lookup and calls {@link Fleetd#leadConfigDirLookup} directly. Neither one ever + * calls the thing {@code Fleetd.main} actually calls. + * + *

The fix extracts the inline {@code new} into {@link Fleetd#leadConfigDirSource}, a + * package-private factory in the same style as {@link Fleetd#loopHealthSource}/{@link + * Fleetd#capacitySource}/{@link Fleetd#healthCoverageSource} — which is exactly what makes it + * directly callable here. This test calls that factory with real {@link FleetConfig.Profile}/ + * {@link FleetConfig.Leader} fixtures (the same shapes {@code FleetdLeadConfigDirLookupTest} + * already uses) and asserts the returned source resolves a real {@code configDir} — a property + * that would be false if {@link Fleetd#leadConfigDirSource} were mutated to {@code return + * FleetMcp.LeadConfigDirSource.none();}. Measured: mutating exactly that line makes + * {@link #resolvesTheRealConfiguredConfigDir()} fail ({@code expected: but was: + * }); restoring it makes the whole suite green again. + * + *

What this class does not and cannot cover. {@code main}'s own line — + * {@code leadConfigDirSource(() -> config.get().profiles(), leaders)} — could itself be swapped + * for a bare {@code FleetMcp.LeadConfigDirSource.none()}, bypassing this factory entirely. Measured: + * that exact mutation compiles with 0 errors and leaves every test in this file, and the full + * 1825-test suite, green. {@link Fleetd#loopHealthSource}'s own wiring test has the identical gap + * for its own one-line call in {@code main} — no test in this codebase calls {@code Fleetd.main} + * far enough to observe which factory call it made. This class narrows the gap from "nothing tests + * the wiring" (the pre-extraction state this ticket found) to "the factory's own logic is pinned, + * and main's call to it is a one-line, visually-verifiable delegation" — the same standard already + * accepted for {@code loopHealthSource}/{@code capacitySource}/{@code healthCoverageSource}. + */ +class FleetdLeadConfigDirSourceWiringTest { + + 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("the returned source resolves the lead's REAL configured configDir, not a hardcoded null") + void resolvesTheRealConfiguredConfigDir() { + Map profiles = Map.of("opus", profileWithConfigDir("opus", "/mnt/opus-claude")); + Map leaders = Map.of("primary", leadOnProfile("opus")); + + FleetMcp.LeadConfigDirSource source = Fleetd.leadConfigDirSource(() -> profiles, leaders); + + assertEquals("/mnt/opus-claude", source.configDirFor().apply("primary"), + "the configDirFor function must delegate to the real leadConfigDirLookup — mutating " + + "Fleetd.leadConfigDirSource's own body to `return FleetMcp.LeadConfigDirSource.none();` " + + "must fail this assertion (measured: it does — see this test's class javadoc for the " + + "companion measurement on main's one-line call to this factory, which this assertion " + + "does not and structurally cannot cover)"); + } + + @Test + @DisplayName("a lead on a profile with no configDir override still resolves to null, not a crash") + void leadWithNoConfigDirOverrideResolvesToNull() { + Map profiles = Map.of("opus", profileWithConfigDir("opus", null)); + Map leaders = Map.of("primary", leadOnProfile("opus")); + + FleetMcp.LeadConfigDirSource source = Fleetd.leadConfigDirSource(() -> profiles, leaders); + + assertNull(source.configDirFor().apply("primary"), + "no configDir: override configured ⇒ null, so LeadContextGauge falls back to its own default"); + } + + @Test + @DisplayName("an unrecognised lead name resolves to null, not a thrown exception") + void unrecognisedLeadNameResolvesToNull() { + FleetMcp.LeadConfigDirSource source = Fleetd.leadConfigDirSource(Map::of, Map.of()); + + assertNull(source.configDirFor().apply("ghost-lead")); + } +} -- 2.52.0