fleetd: report a lead's live context usage in fleet_list #602
@@ -672,6 +672,10 @@ public final class Fleetd {
|
|||||||
leadMailbox,
|
leadMailbox,
|
||||||
outageSource,
|
outageSource,
|
||||||
new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)),
|
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
|
// 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()),
|
// 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
|
// 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
|
* 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
|
* 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.MemberPresence;
|
||||||
import dev.ltms.fleet.inject.CompletionResolver;
|
import dev.ltms.fleet.inject.CompletionResolver;
|
||||||
import dev.ltms.fleet.herdr.HerdrException;
|
import dev.ltms.fleet.herdr.HerdrException;
|
||||||
|
import dev.ltms.fleet.lead.LeadContextGauge;
|
||||||
import dev.ltms.fleet.lead.LeadRollover;
|
import dev.ltms.fleet.lead.LeadRollover;
|
||||||
import dev.ltms.fleet.msg.LeadChannel;
|
import dev.ltms.fleet.msg.LeadChannel;
|
||||||
import dev.ltms.fleet.msg.LeadMessage;
|
import dev.ltms.fleet.msg.LeadMessage;
|
||||||
@@ -115,6 +116,8 @@ public final class FleetMcp {
|
|||||||
private final OutageSource outage;
|
private final OutageSource outage;
|
||||||
/** fleetd #176: SEPARATE from both of the above — see {@link LeadSeatSource}'s doc. */
|
/** fleetd #176: SEPARATE from both of the above — see {@link LeadSeatSource}'s doc. */
|
||||||
private final LeadSeatSource leadSeats;
|
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. */
|
/** CB-637: this daemon's lead-to-lead channel; {@code null} when no coordinator is configured. */
|
||||||
private final LeadChannel leadChannel;
|
private final LeadChannel leadChannel;
|
||||||
/** fleetd #361: {@code coordinator.peers} — see {@link CoordinationSource}. Empty when unset. */
|
/** 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}.
|
* clean {@code NOT_CONFIGURED} refusal rather than throwing. See {@link #handover}.
|
||||||
*/
|
*/
|
||||||
private final LeadRollover leadRollover;
|
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. */
|
/** 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,
|
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); }
|
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
|
* 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
|
* 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) {
|
LeadSeatSource leadSeats, List<String> peers, LeadRollover leadRollover) {
|
||||||
this(messages, workers, sessions, identity, presence, primaryRegistry, callers, authorizationMode, metrics,
|
this(messages, workers, sessions, identity, presence, primaryRegistry, callers, authorizationMode, metrics,
|
||||||
capacity, healthCoverage, LoopHealthSource.none(), quarantine, leadChannel, outage, leadSeats,
|
capacity, healthCoverage, LoopHealthSource.none(), quarantine, leadChannel, outage, leadSeats,
|
||||||
peers, leadRollover);
|
LeadConfigDirSource.none(), peers, leadRollover);
|
||||||
}
|
}
|
||||||
|
|
||||||
public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
|
public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
|
||||||
@@ -354,7 +391,8 @@ public final class FleetMcp {
|
|||||||
CallerResolver callers, AuthorizationMode authorizationMode, Metrics metrics,
|
CallerResolver callers, AuthorizationMode authorizationMode, Metrics metrics,
|
||||||
CapacitySource capacity, HealthCoverageSource healthCoverage, LoopHealthSource loopHealth,
|
CapacitySource capacity, HealthCoverageSource healthCoverage, LoopHealthSource loopHealth,
|
||||||
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage,
|
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");
|
Objects.requireNonNull(callers, "callers");
|
||||||
this.authorizationEnforced = Objects.requireNonNull(authorizationMode, "authorizationMode")
|
this.authorizationEnforced = Objects.requireNonNull(authorizationMode, "authorizationMode")
|
||||||
== AuthorizationMode.ENFORCED;
|
== AuthorizationMode.ENFORCED;
|
||||||
@@ -364,6 +402,7 @@ public final class FleetMcp {
|
|||||||
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
|
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
|
||||||
this.outage = Objects.requireNonNull(outage, "outage");
|
this.outage = Objects.requireNonNull(outage, "outage");
|
||||||
this.leadSeats = Objects.requireNonNull(leadSeats, "leadSeats");
|
this.leadSeats = Objects.requireNonNull(leadSeats, "leadSeats");
|
||||||
|
this.leadConfigDirs = Objects.requireNonNull(leadConfigDirs, "leadConfigDirs");
|
||||||
this.healthCoverage = healthCoverage;
|
this.healthCoverage = healthCoverage;
|
||||||
this.loopHealth = Objects.requireNonNull(loopHealth, "loopHealth");
|
this.loopHealth = Objects.requireNonNull(loopHealth, "loopHealth");
|
||||||
this.leadRollover = leadRollover;
|
this.leadRollover = leadRollover;
|
||||||
@@ -498,7 +537,7 @@ public final class FleetMcp {
|
|||||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_list", Map.of()), null);
|
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_list", Map.of()), null);
|
||||||
if (denied != null) return denied;
|
if (denied != null) return denied;
|
||||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage,
|
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage,
|
||||||
leadSeats, callers.leads(),
|
leadSeats, leadContextGauge, leadConfigDirs, callers.leads(),
|
||||||
callerTerminal(exchange),
|
callerTerminal(exchange),
|
||||||
new CoordinationSource(leadChannel, peers),
|
new CoordinationSource(leadChannel, peers),
|
||||||
coordinatorVisibleTo(principal(exchange)));
|
coordinatorVisibleTo(principal(exchange)));
|
||||||
@@ -1601,7 +1640,7 @@ public final class FleetMcp {
|
|||||||
Map<String, String> leads, String selfTerm,
|
Map<String, String> leads, String selfTerm,
|
||||||
CoordinationSource coordination) {
|
CoordinationSource coordination) {
|
||||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine,
|
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}). */
|
/** 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,
|
QuarantineSource quarantine, OutageSource outage,
|
||||||
Map<String, String> leads, String selfTerm) {
|
Map<String, String> leads, String selfTerm) {
|
||||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
|
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,
|
QuarantineSource quarantine, Map<String, String> leads, String selfTerm,
|
||||||
CoordinationSource coordination) {
|
CoordinationSource coordination) {
|
||||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, OutageSource.none(),
|
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}). */
|
/** 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,
|
QuarantineSource quarantine, OutageSource outage,
|
||||||
Map<String, String> leads, String selfTerm, CoordinationSource coordination) {
|
Map<String, String> leads, String selfTerm, CoordinationSource coordination) {
|
||||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
|
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,
|
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
|
||||||
CoordinationSource coordination) {
|
CoordinationSource coordination) {
|
||||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
|
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,
|
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
|
||||||
CoordinationSource coordination, boolean callerIsPrimary) {
|
CoordinationSource coordination, boolean callerIsPrimary) {
|
||||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine,
|
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,
|
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
||||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||||
LoopHealthSource loopHealth,
|
LoopHealthSource loopHealth,
|
||||||
QuarantineSource quarantine, OutageSource outage,
|
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) {
|
CoordinationSource coordination, boolean callerIsPrimary) {
|
||||||
try {
|
try {
|
||||||
Map<String, Agent> live = workers.list().stream()
|
Map<String, Agent> live = workers.list().stream()
|
||||||
@@ -1698,7 +1747,8 @@ public final class FleetMcp {
|
|||||||
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
|
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
|
||||||
List<Map<String, Object>> leadRows = leads.entrySet().stream()
|
List<Map<String, Object>> leadRows = leads.entrySet().stream()
|
||||||
.sorted(Map.Entry.comparingByValue())
|
.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();
|
.toList();
|
||||||
// fleetd #209: this is the caller-driven fleet_list read that actually reports
|
// fleetd #209: this is the caller-driven fleet_list read that actually reports
|
||||||
// agentSessionId (via memberCapacityView -> SessionManager.rosterView), so it uses the
|
// 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
|
* <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
|
* 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,
|
* 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.
|
* 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,
|
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<>();
|
Map<String, Object> m = new LinkedHashMap<>();
|
||||||
m.put("sessionId", terminal);
|
m.put("sessionId", terminal);
|
||||||
m.put("name", name);
|
m.put("name", name);
|
||||||
@@ -2010,9 +2070,32 @@ public final class FleetMcp {
|
|||||||
if (terminal.equals(selfTerm)) {
|
if (terminal.equals(selfTerm)) {
|
||||||
m.put("self", true);
|
m.put("self", true);
|
||||||
}
|
}
|
||||||
|
String configDir = leadConfigDirs.configDirFor().apply(name);
|
||||||
|
m.put("context", contextView(contextGauge, live, configDir));
|
||||||
return m;
|
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. */
|
/** {@code fleet_stop}: tear a worker down by its pane id. */
|
||||||
static McpSchema.CallToolResult stop(SessionManager sessions, String paneId) {
|
static McpSchema.CallToolResult stop(SessionManager sessions, String paneId) {
|
||||||
if (isBlank(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);
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user