fleetd: report a lead's live context usage in fleet_list
CI / shell-tests (pull_request) Failing after 7s
CI / build (pull_request) Failing after 1m37s
CI / contract (pull_request) Successful in 2m7s

fleetd had no way to see how full a lead's context window is. On this host a
lead auto-compacted 30 times in one session, discarding roughly 250,000 tokens
and costing 46s to 3m16s each time, and nothing could see it coming.

LeadContextGauge reads the transcript Claude Code itself writes, never the
lead's pane. It finds <sessionId>.jsonl by NAME under <configDir>/projects/
rather than deriving the project slug, which is an undocumented internal.

Three states, not two: OK, HIGH, UNKNOWN. Every path that cannot positively
establish a token count reports UNKNOWN with no number, so a lead is never
told it is fine when the honest answer is "I could not look".

Bounded two ways: TAIL_BYTES caps bytes read off disk, and a 5s cache TTL caps
how often that read happens, because fleet_list is polled constantly.

Recovered by the lead: the authoring member ended on a backend error (DNS
ENOTFOUND) with this work uncommitted and unpushed in its worktree. Verified
before committing: mvn -o clean install exit 0, 1813 tests, 0 failures,
0 errors, 0 skipped, 144 reports; LeadContextGaugeTest 8/8.

KNOWN INCOMPLETE - see the PR. The fleet_list call site passes configDir=null,
which falls back to ~/.claude, but this host's lead profile sets configDir to
an override. Measured: 19 transcripts under the real configDir, 0 under the
fallback. The gauge is therefore INERT on this fleet until that is wired.
This commit is contained in:
Dai Ha
2026-09-20 16:05:15 +07:00
parent 6eb34a654f
commit 3ed7bfca67
3 changed files with 561 additions and 11 deletions
@@ -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).
*
* <p><strong>Why this exists.</strong> 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.
*
* <p><strong>The route.</strong> Claude Code appends one JSON object per line to
* {@code <configDir>/projects/<slug>/<sessionId>.jsonl}. {@code <slug>} is an undocumented,
* internal encoding of the working directory — this class never derives it. Instead it lists the
* one-level-deep subdirectories of {@code <configDir>/projects/} and looks for
* {@code <sessionId>.jsonl} by name, so the slug rule can change without breaking this reader.
*
* <ul>
* <li><strong>Live context</strong> is read off the last record in the read window that carries
* a {@code message.usage} object: {@code input_tokens + cache_read_input_tokens +
* cache_creation_input_tokens}. This is what actually fills the window — a plain
* {@code input_tokens} count alone understates it once the conversation has any cached
* prefix, which on a long-lived lead is always.</li>
* <li><strong>Compaction history</strong> is a count of {@code subtype: "compact_boundary"}
* records seen in the same read window — see {@link Reading#compactions()}. It is a count
* within the window this reader actually looked at, not a lifetime total: a session with
* more compactions than fit in {@link #TAIL_BYTES} of transcript will undercount. That
* trade-off is deliberate — see {@link #TAIL_BYTES}.</li>
* </ul>
*
* <p><strong>Three states, not two (OK / HIGH / UNKNOWN).</strong> Every path that cannot
* positively establish the live token count — a missing file, an unreadable one, a peer that
* is not a Claude backend, or a last line that fails to parse as JSON — returns
* {@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.
*
* <p><strong>Bounded cost.</strong> {@code fleet_list} is polled constantly, so every read is
* capped two ways: {@link #TAIL_BYTES} bounds how much of the transcript is ever read from disk
* (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<String, CacheEntry> 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 <user.home>/.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 <sessionId>.jsonl} under {@code <base>/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<String> 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;
}
}
}
@@ -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<String, Integer> liveCount, Function<String, Integer> 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<String, String> 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<String, String> 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<String, String> 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<String, String> 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<String, String> 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<String, String> 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<String, String> leads, String selfTerm,
LeadSeatSource leadSeats, LeadContextGauge contextGauge,
Map<String, String> leads, String selfTerm,
CoordinationSource coordination, boolean callerIsPrimary) {
try {
Map<String, Agent> live = workers.list().stream()
@@ -1698,7 +1719,7 @@ public final class FleetMcp {
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
List<Map<String, Object>> leadRows = leads.entrySet().stream()
.sorted(Map.Entry.comparingByValue())
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm))
.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.
*
* <p>{@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.
*
* <p>{@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<String, Object> leadView(String terminal, String name, Agent live,
String selfTerm) {
String selfTerm, LeadContextGauge contextGauge) {
Map<String, Object> 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<String, Object> 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<String, Object> 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)) {
@@ -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):
* <ol>
* <li>{@link #tokenCountTracksTheLastUsageRecordAndChangesWithIt()}</li>
* <li>{@link #compactionCountTracksCompactBoundaryRecordsAndChangesWithIt()}</li>
* <li>{@link #missingFileIsUnknown()}, {@link #unreadableFileIsUnknown()},
* {@link #truncatedOrInvalidLastLineIsUnknown()}</li>
* <li>{@link #readNeverExceedsTheTailBound()}</li>
* <li>{@link #secondReadInsideTtlDoesNotTouchDiskAgain()}</li>
* </ol>
*
* <p>Every fixture lives under {@code @TempDir} — never the operator's real config directory (see
* the ticket's hard constraint on this).
*/
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 <configDir>/projects/<anySlug>/<sessionId>.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());
}
}