Compare commits

...

4 Commits

Author SHA1 Message Date
ltms 3762aca307 Merge #606: wire the lead context gauge to the real configDir, and stop flapping on a torn line
CI / shell-tests (pull_request) Failing after 6s
CI / build (pull_request) Failing after 1m21s
CI / contract (pull_request) Successful in 1m50s
Fixes the two defects recorded in #602's own body.

1. FleetMcp.contextView passed configDir=null, so the gauge read <user.home>/
   .claude while this host's lead profile sets an override. Measured before the
   fix: 19 transcripts under the real directory, 0 under the fallback. The gauge
   would have deployed green and reported UNKNOWN forever, for every lead.
   Now threaded via FleetMcp.LeadConfigDirSource, built by Fleetd
   .leadConfigDirSource, following fleet.leaders.<name>.profile to that
   profile's configDir and reading config.get() live inside the lambda.

2. A torn final line no longer means UNKNOWN. fleet_list reads a transcript
   Claude Code may be mid-write on, so the last line can be cut. The old code
   treated that as fatal, which would make the gauge flap at random. The stated
   reason ("a format change should show as UNKNOWN") does not hold: a real
   format change makes EVERY line unparseable, and that case is still caught.

Verified by the lead on the combined state with current main merged in, not on
the branch alone: mvn -o clean install exit 0, 147 reports, 1841 tests, 0
failures, 0 errors, 0 skipped.

Mutation-checked independently by the lead:
  Fleetd.leadConfigDirSource body -> none()   -> KILLED (1 failure)
  FleetMcp.contextView configDir -> null      -> KILLED (worker-measured)

KNOWN RESIDUAL, documented rather than overclaimed. main's own one-line call to
leadConfigDirSource could be swapped for none() and the suite stays green. Every
member of this wiring-test family (loopHealthSource, capacitySource,
healthCoverageSource) has the identical gap — no test runs Fleetd.main far enough
to observe which factory it called. Filed separately as a class-wide problem
rather than patched here. The empirical close is the dogfood check after redeploy.
2026-09-20 11:38:22 +02:00
Dai Ha 4e27bde2d7 fleetd #602 gauge-wiring follow-up: pin Fleetd.main's LeadConfigDirSource wiring
Extract the inline new FleetMcp.LeadConfigDirSource(leadConfigDirLookup(...))
construction in Fleetd.main into a package-private factory,
Fleetd.leadConfigDirSource, mirroring loopHealthSource/capacitySource/
healthCoverageSource. Add FleetdLeadConfigDirSourceWiringTest, which calls the
factory directly with real Profile/Leader fixtures and asserts the returned
source resolves a real configDir -- a property that is false if the factory's
body is mutated to return LeadConfigDirSource.none().

Neither FleetMcpLeadContextGaugeWiringTest nor FleetdLeadConfigDirLookupTest
could catch main losing this wiring: each builds its own instance instead of
calling what main calls. This closes that gap at the factory level, matching
the standard already accepted for loopHealthSource's own wiring test.
2026-09-20 16:34:49 +07:00
Dai Ha d345b14e43 fleetd #602 gauge-wiring: thread a lead's configured configDir into the context gauge
FleetMcp.contextView hardcoded LeadContextGauge.read(null, ...), so a lead whose
profile sets its own CLAUDE_CONFIG_DIR always read the wrong transcript directory
and reported UNKNOWN forever, with no error anywhere.

- Add FleetMcp.LeadConfigDirSource (same idiom as LeadSeatSource) and thread it
  through the constructor / listFleet overload chain / leadView / contextView.
- Add Fleetd.leadConfigDirLookup, wired at construction, following the same
  fleet.leaders.<name>.profile link leadSeatLookup already uses, one step
  further to that profile's own configDir.
- LeadContextGauge.parse: a single unparseable line (typically the final one,
  torn by a write this read raced) is now skipped rather than forcing UNKNOWN;
  only when every line in the read window fails to parse does it report
  UNKNOWN, which is the real format-change signal.
- Tests: FleetdLeadConfigDirLookupTest (lookup logic), FleetMcpLeadContextGaugeWiringTest
  (end-to-end: config naming directory A vs B decides which is read; a lead with
  no configured dir degrades without throwing), and two replacement properties in
  LeadContextGaugeTest for the torn-line fix plus its all-unparseable control.
2026-09-20 16:19:37 +07:00
Dai Ha 3ed7bfca67 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.
2026-09-20 16:05:15 +07: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);
}
}