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

+ * + *

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()); + } +}