diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index 76ed1de..d06c01f 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 leadConfigDirSource's own doc for why this, not a hardcoded + // null, is what fleet_list's context row now reads. + 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 @@ -1576,6 +1580,58 @@ 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 #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/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..fc16d01 --- /dev/null +++ b/fleetd/src/main/java/dev/ltms/fleet/lead/LeadContextGauge.java @@ -0,0 +1,307 @@ +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. + * + *

+ * + *

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 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 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 + * (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). 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); + 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(); + } + Long tokens = null; + int compactions = 0; + 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) + + 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 (!anyLineParsed) { + return Reading.unknown(); + } + 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..f0341db 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; @@ -115,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. */ @@ -126,6 +129,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, @@ -246,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 @@ -346,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, @@ -354,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; @@ -364,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; @@ -498,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, callers.leads(), + leadSeats, leadContextGauge, leadConfigDirs, callers.leads(), callerTerminal(exchange), new CoordinationSource(leadChannel, peers), coordinatorVisibleTo(principal(exchange))); @@ -1601,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(), 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}). */ @@ -1610,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(), leads, selfTerm, CoordinationSource.none(), false); + LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, CoordinationSource.none(), false); } /** @@ -1627,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(), 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}). */ @@ -1636,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(), leads, selfTerm, coordination, false); + LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false); } /** @@ -1658,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, leads, selfTerm, coordination, false); + leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false); } /** @@ -1682,14 +1721,24 @@ 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(), LeadConfigDirSource.none(), 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, + LeadConfigDirSource leadConfigDirs, + Map leads, String selfTerm, CoordinationSource coordination, boolean callerIsPrimary) { try { Map live = workers.list().stream() @@ -1698,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)) + .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 @@ -1993,15 +2043,25 @@ 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, + LeadConfigDirSource leadConfigDirs) { Map m = new LinkedHashMap<>(); m.put("sessionId", terminal); m.put("name", name); @@ -2010,9 +2070,32 @@ public final class FleetMcp { if (terminal.equals(selfTerm)) { m.put("self", true); } + 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 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, String configDir) { + String sessionId = live == null ? null : live.sessionId(); + String agentType = live == null ? null : live.agentType(); + LeadContextGauge.Reading reading = contextGauge.read(configDir, 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/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/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")); + } +} 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..3eb41be --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/lead/LeadContextGaugeTest.java @@ -0,0 +1,238 @@ +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()}
  6. + *
  7. {@link #tornFinalLineFallsBackToTheLastGoodReading()}, + * {@link #everyLineUnparseableIsUnknown()} (fleetd #602 gauge-wiring, finding 2)
  8. + *
  9. {@link #readNeverExceedsTheTailBound()}
  10. + *
  11. {@link #secondReadInsideTtlDoesNotTouchDiskAgain()}
  12. + *
+ * + *

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 { + + 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: two 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("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, + "{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(), + "every line unparseable is the real format-change signal and must still report UNKNOWN"); + 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()); + } +} 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); + } +}