fleetd: report a lead's live context usage in fleet_list #602

Merged
ltms merged 4 commits from worker/lead-context-gauge-ad404f-1 into main 2026-09-20 11:38:31 +02:00
7 changed files with 1021 additions and 13 deletions
@@ -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 <user.home>/.claude} default.
*
* <p>Follows the same {@code fleet.leaders.<name>.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<String, String> leadConfigDirLookup(Supplier<Map<String, FleetConfig.Profile>> profiles,
Map<String, FleetConfig.Leader> 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}).
*
* <p>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<Map<String, FleetConfig.Profile>> profiles,
Map<String, FleetConfig.Leader> 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
@@ -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).
*
* <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 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".
*
* <p><strong>A torn final line does not mean UNKNOWN.</strong> {@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
* <em>every</em> 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 <em>none</em> 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".
*
* <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). 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<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();
}
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;
}
}
}
@@ -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<String, Integer> liveCount, Function<String, Integer> 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 <user.home>/.claude} default.
*
* <p>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.<name>.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<String, String> 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<String> 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<String> peers, LeadRollover leadRollover) {
LeadSeatSource leadSeats, LeadConfigDirSource leadConfigDirs, List<String> 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<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(), 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<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(), LeadConfigDirSource.none(), leads, selfTerm, CoordinationSource.none(), false);
}
/**
@@ -1627,7 +1666,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(), 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<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(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
}
/**
@@ -1658,7 +1697,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(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
}
/**
@@ -1682,14 +1721,24 @@ 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(), 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<String, String> leads, String selfTerm,
LeadSeatSource leadSeats, LeadContextGauge contextGauge,
LeadConfigDirSource leadConfigDirs,
Map<String, String> leads, String selfTerm,
CoordinationSource coordination, boolean callerIsPrimary) {
try {
Map<String, Agent> live = workers.list().stream()
@@ -1698,7 +1747,8 @@ 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,
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.
*
* <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,
LeadConfigDirSource leadConfigDirs) {
Map<String, Object> 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.<name>.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 <user.home>/.claude}).
*/
private static Map<String, Object> 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<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,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 <user.home>/.claude} default (see {@code LeadContextGauge}).
*
* <p>{@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.<name>.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<String, FleetConfig.Profile> profiles = Map.of("opus", profileWithConfigDir("opus", "/mnt/opus-claude"));
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
Function<String, String> 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<Map<String, FleetConfig.Profile>> profilesRef =
new java.util.concurrent.atomic.AtomicReference<>(
Map.of("opus", profileWithConfigDir("opus", "/mnt/dir-a")));
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
Function<String, String> 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<String, FleetConfig.Profile> 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<String, FleetConfig.Leader> leaders = Map.of("primary", recogniseOnly);
Function<String, String> 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<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("ghost-profile"));
Function<String, String> 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<String, FleetConfig.Profile> profiles = Map.of("opus", profileWithConfigDir("opus", null));
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
Function<String, String> 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<String, String> lookup = Fleetd.leadConfigDirLookup(Map::of, Map.of());
assertNull(lookup.apply("ghost-lead"));
}
}
@@ -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}).
*
* <p>{@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.
*
* <p>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: </mnt/opus-claude> but was:
* <null>}); restoring it makes the whole suite green again.
*
* <p><b>What this class does not and cannot cover.</b> {@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<String, FleetConfig.Profile> profiles = Map.of("opus", profileWithConfigDir("opus", "/mnt/opus-claude"));
Map<String, FleetConfig.Leader> 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<String, FleetConfig.Profile> profiles = Map.of("opus", profileWithConfigDir("opus", null));
Map<String, FleetConfig.Leader> 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"));
}
}
@@ -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):
* <ol>
* <li>{@link #tokenCountTracksTheLastUsageRecordAndChangesWithIt()}</li>
* <li>{@link #compactionCountTracksCompactBoundaryRecordsAndChangesWithIt()}</li>
* <li>{@link #missingFileIsUnknown()}, {@link #unreadableFileIsUnknown()}</li>
* <li>{@link #tornFinalLineFallsBackToTheLastGoodReading()},
* {@link #everyLineUnparseableIsUnknown()} (fleetd #602 gauge-wiring, finding 2)</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).
*
* <p><strong>fleetd #602 gauge-wiring, finding 2.</strong> 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 <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: 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());
}
}
@@ -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.
*
* <p>This class proves the directory {@code fleet.leaders.<name>.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-<terminalId>}. */
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 <configDir>/projects/<anySlug>/<sessionId>.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<String> 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);
}
}