Compare commits
19 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 388ef5a3c3 | |||
| 076cc43f7b | |||
| bad47a8444 | |||
| 955b9ea013 | |||
| 89cb8ff79b | |||
| fa62e9906d | |||
| d7390ccd37 | |||
| aa517ae0ec | |||
| 60496831c2 | |||
| 9a992d0f70 | |||
| 3762aca307 | |||
| 4e27bde2d7 | |||
| b85d9b0e46 | |||
| 856dfc6318 | |||
| 81c1d8e91c | |||
| d345b14e43 | |||
| 9640deeffc | |||
| f336bcef39 | |||
| 3ed7bfca67 |
@@ -145,9 +145,10 @@ authorized to settle it, not only to advise. Architects first form independent p
|
||||
compare them. If they still disagree after that comparison, they return both positions and their
|
||||
checked evidence; the lead decides. Go to the operator only for an action the fleet has no
|
||||
authority to take, such as spending money, granting access, or making a promise to someone else.
|
||||
**Then write the decision on the ticket.** Taking the operator out of the loop also removes the signal they used to get, because
|
||||
that signal was the block itself — work stopped, so they found out. A ticket comment replaces it,
|
||||
and it reaches them whether or not they are at a terminal when you decide.
|
||||
**Then write the decision on the ticket.** Taking the operator out of the loop also removes the
|
||||
signal they used to get, because that signal was the block itself — work stopped, so they found
|
||||
out. A ticket comment replaces it, and it reaches them whether or not they are at a terminal when
|
||||
you decide.
|
||||
|
||||
| Intent | Tool |
|
||||
|---|---|
|
||||
|
||||
@@ -77,16 +77,25 @@ bind:
|
||||
# a lead turn nobody asked for), so upgrading the daemon must never switch it on for you. Absent
|
||||
# block = feature off, exactly as before.
|
||||
#
|
||||
# Three knobs, each with a default that errs on the side of not burning context:
|
||||
# Four knobs, each with a default that errs on the side of not burning context:
|
||||
# idleAfterSeconds: 300 # how long the lead must stay idle before the FIRST nudge (default 300 —
|
||||
# # absorbs normal post-turn pauses; re-prompting every pause burns context)
|
||||
# backoffMs: 60000 # re-check cadence / spacing between nudges past the quiet period (default 60000)
|
||||
# quietNudgeCap: 3 # cap on consecutive nudges that find NOTHING pending, then it stops
|
||||
# # until real state appears (default 3 — never nag an empty fleet forever)
|
||||
# contextHighNudge: false # fleetd #609 — when true, an idle lead whose OWN Claude Code context
|
||||
# # reads HIGH (see LeadContextGauge; fleet_list's context row) gets a text
|
||||
# # notice appended to its nudge telling it to consider fleet_handover. Text
|
||||
# # only — it never rolls a pane by itself, and only the operator can approve
|
||||
# # a roll. Fires once per HIGH stretch (a later OK reading re-arms it), and
|
||||
# # never spends the quietNudgeCap budget. Default false/absent = off, same
|
||||
# # as every other knob here — an upgraded daemon must not start telling
|
||||
# # leads to hand over on its own.
|
||||
# leadHeartbeat:
|
||||
# idleAfterSeconds: 300
|
||||
# backoffMs: 60000
|
||||
# quietNudgeCap: 3
|
||||
# contextHighNudge: false
|
||||
|
||||
# Lead rollover (fleetd #480): replace a lead session that has decided it is ready to be replaced,
|
||||
# without an operator doing it by hand. A lead writes a handover file, then asks fleetd to clear its
|
||||
|
||||
@@ -4,11 +4,13 @@ import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.ConfigWatcher;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.Agent;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.herdr.HerdrRouter;
|
||||
import dev.ltms.fleet.herdr.LeadTabScanner;
|
||||
import dev.ltms.fleet.lead.LeadContextGauge;
|
||||
import dev.ltms.fleet.lead.LeadLauncher;
|
||||
import dev.ltms.fleet.lead.LeadRollover;
|
||||
import dev.ltms.fleet.herdr.PaneLocator;
|
||||
@@ -563,10 +565,18 @@ public final class Fleetd {
|
||||
Thread.ofVirtual().name("bridge-heartbeat-").unstarted(r));
|
||||
if (cfg.leadHeartbeat() != null) {
|
||||
var hb = cfg.leadHeartbeat();
|
||||
// fleetd #609: own LeadContextGauge instance for the heartbeat loop — separate from the
|
||||
// one FleetMcp builds internally for fleet_list's context row. Each caches independently
|
||||
// (keyed by configDir+sessionId), so this costs at most one extra bounded tail read per
|
||||
// TTL window, never a shared-mutable-state hazard between the two callers.
|
||||
var leadContextGauge = new LeadContextGauge();
|
||||
heartbeat = new LeadHeartbeatLoop(primaryRegistry, router.leadAgents(), replyInbox, sessions::roster,
|
||||
pushLoop, heartbeatScheduler, System::nanoTime,
|
||||
TimeUnit.SECONDS.toNanos(hb.idleAfterSeconds()), hb.backoffMs(), hb.quietNudgeCap(),
|
||||
metrics);
|
||||
metrics,
|
||||
leadContextSource(leadContextGauge, router.leadAgents(), leads,
|
||||
leadConfigDirLookup(() -> config.get().profiles(), leaders)),
|
||||
Boolean.TRUE.equals(hb.contextHighNudge()));
|
||||
heartbeat.start();
|
||||
} else {
|
||||
heartbeat = null;
|
||||
@@ -672,6 +682,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 +1590,109 @@ 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 #609: per-terminal factory for {@link LeadHeartbeatLoop.LeadContextSource} — the lead's
|
||||
* own {@link LeadContextGauge} reading, so the heartbeat loop can tell an idle, HIGH-context lead
|
||||
* to consider a handover.
|
||||
*
|
||||
* <p>Three hops, each degrading to {@link LeadContextGauge.Reading#unknown()} rather than
|
||||
* throwing, since a herdr hiccup or an unrecognised terminal must never kill the heartbeat's own
|
||||
* tick: {@code liveLeadTerminals} (terminal id → lead name, the same live supplier {@link
|
||||
* #leadSeatLookup} and {@code LeadCoordLoop} already read) → {@code configDirForLeadName} (that
|
||||
* lead's {@code configDir}, normally {@link #leadConfigDirLookup}'s return) → {@code agents.get}
|
||||
* for the live {@link Agent#sessionId()}/{@link Agent#agentType()} the gauge itself needs.
|
||||
*
|
||||
* @param gauge the {@link LeadContextGauge} instance to read through — shares its
|
||||
* cache across every call this factory's function makes
|
||||
* @param agents the {@link AgentControl} instance that reaches the LEAD's pane
|
||||
* (not {@code memberAgents}), normally {@code router.leadAgents()}
|
||||
* @param liveLeadTerminals terminal id → lead name for every CURRENTLY recognised lead
|
||||
* @param configDirForLeadName lead name → {@code configDir}, normally {@link
|
||||
* #leadConfigDirLookup}'s return
|
||||
*/
|
||||
static Function<String, LeadContextGauge.Reading> leadContextLookup(LeadContextGauge gauge, AgentControl agents,
|
||||
Supplier<Map<String, String>> liveLeadTerminals, Function<String, String> configDirForLeadName) {
|
||||
return terminal -> {
|
||||
String leadName = liveLeadTerminals.get().get(terminal);
|
||||
if (leadName == null) {
|
||||
return LeadContextGauge.Reading.unknown();
|
||||
}
|
||||
String configDir = configDirForLeadName.apply(leadName);
|
||||
Agent live;
|
||||
try {
|
||||
live = agents.get(terminal);
|
||||
} catch (RuntimeException e) {
|
||||
return LeadContextGauge.Reading.unknown();
|
||||
}
|
||||
return gauge.read(configDir, live.sessionId(), live.agentType());
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #609: wraps {@link #leadContextLookup} into a {@link LeadHeartbeatLoop.LeadContextSource}
|
||||
* — the same hand-built-vs-wired shape as {@link #leadConfigDirSource}/{@link #loopHealthSource}/
|
||||
* {@link #capacitySource}/{@link #healthCoverageSource}. Extracted so a test can call the exact
|
||||
* factory {@code main} calls, rather than only a lookup nothing in {@code main} is proven to use
|
||||
* (see {@code FleetdLeadConfigDirSourceWiringTest}'s javadoc for the measured gap this shape closes).
|
||||
*/
|
||||
static LeadHeartbeatLoop.LeadContextSource leadContextSource(LeadContextGauge gauge, AgentControl agents,
|
||||
Supplier<Map<String, String>> liveLeadTerminals, Function<String, String> configDirForLeadName) {
|
||||
return new LeadHeartbeatLoop.LeadContextSource(
|
||||
leadContextLookup(gauge, agents, liveLeadTerminals, configDirForLeadName));
|
||||
}
|
||||
|
||||
/**
|
||||
* 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
|
||||
|
||||
@@ -1329,9 +1329,20 @@ public record FleetConfig(
|
||||
* each such nudge costs the lead a turn just to read "nothing pending";
|
||||
* three is enough to tell it it may stand down without nagging forever, and
|
||||
* it is the bound that stops an idle fleet from being a subscription burner.
|
||||
* @param contextHighNudge fleetd #609: when {@code true}, an idle lead whose own {@code
|
||||
* LeadContextGauge} reading is {@code HIGH} gets a text notice telling it
|
||||
* to consider a handover, appended to whatever heartbeat nudge the loop
|
||||
* already sends. Default {@code false} ({@code null} also means off) — an
|
||||
* upgraded daemon must not silently start telling leads to hand over.
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record LeadHeartbeat(Integer idleAfterSeconds, Long backoffMs, Integer quietNudgeCap) {
|
||||
public record LeadHeartbeat(Integer idleAfterSeconds, Long backoffMs, Integer quietNudgeCap,
|
||||
Boolean contextHighNudge) {
|
||||
/** Convenience constructor for every call site that predates fleetd #609: no context notice. */
|
||||
public LeadHeartbeat(Integer idleAfterSeconds, Long backoffMs, Integer quietNudgeCap) {
|
||||
this(idleAfterSeconds, backoffMs, quietNudgeCap, null);
|
||||
}
|
||||
|
||||
public LeadHeartbeat {
|
||||
idleAfterSeconds = (idleAfterSeconds == null || idleAfterSeconds <= 0) ? 300 : idleAfterSeconds;
|
||||
backoffMs = (backoffMs == null || backoffMs <= 0) ? 60_000L : backoffMs;
|
||||
|
||||
@@ -0,0 +1,313 @@
|
||||
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) {
|
||||
/**
|
||||
* fleetd #609: widened from package-private to public so {@code
|
||||
* dev.ltms.fleet.msg.LeadHeartbeatLoop.LeadContextSource.none()} (a different package) can
|
||||
* return the same inert "I could not look" reading the gauge itself uses, without inventing
|
||||
* a parallel unknown-reading constant. Behaviour of this class is otherwise unchanged.
|
||||
*/
|
||||
public 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;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -231,7 +231,9 @@ public final class LeadRollover {
|
||||
* #status} could wrongly answer {@link #UNKNOWN} ("nothing was ever requested") for a roll
|
||||
* that is, in fact, actively running. This is not sticky: the deferred continuation
|
||||
* overwrites this same entry with a terminal state ({@link #ROLLED}, {@link
|
||||
* #TURN_NEVER_SETTLED}, or {@link #CLEAR_NEVER_SETTLED}) once it finishes.
|
||||
* #TURN_NEVER_SETTLED}, {@link #CLEAR_NEVER_SETTLED}, or {@link #FAILED}) once it finishes
|
||||
* — including by throwing, which fleetd #615's catch in {@link #runRollover} now turns into
|
||||
* {@link #FAILED} instead of leaving this entry stuck forever.
|
||||
*/
|
||||
IN_PROGRESS,
|
||||
/**
|
||||
@@ -253,6 +255,19 @@ public final class LeadRollover {
|
||||
* within {@code clearSettleSeconds} — {@code bootstrapText} was never sent.
|
||||
*/
|
||||
CLEAR_NEVER_SETTLED,
|
||||
/**
|
||||
* fleetd #615: the deferred continuation threw a {@link RuntimeException} — most likely a
|
||||
* {@link dev.ltms.fleet.herdr.HerdrException} out of one of the two unwrapped {@code
|
||||
* agents.send} calls in {@link #runRollover} — and the continuation thread died with it.
|
||||
* Before this state existed, that throw left {@link #outcomes} holding {@link #IN_PROGRESS}
|
||||
* forever, because the production {@code continuationRunner} is a bare virtual thread with
|
||||
* no uncaught-exception handler and nothing downstream of the throw ever ran to write a
|
||||
* terminal outcome. {@code detail} names the exception, so a reader has something to act on
|
||||
* — the same diagnostic style as {@link #TURN_NEVER_SETTLED} and {@link
|
||||
* #CLEAR_NEVER_SETTLED}. The roll is dead at this point and does not retry itself; a stuck
|
||||
* lead must {@link #open} a fresh request.
|
||||
*/
|
||||
FAILED,
|
||||
/**
|
||||
* {@code token} names nothing this instance currently knows about: never issued by {@link
|
||||
* #open}, dropped by {@link #cancel}, or aged out of {@link #outcomes}'s bounded history.
|
||||
@@ -493,8 +508,42 @@ public final class LeadRollover {
|
||||
* entirely after {@link #confirm} has returned to its caller — see this class's javadoc for the
|
||||
* four-step order. There is no result to return to by this point, so every outcome is logged
|
||||
* only.
|
||||
*
|
||||
* <p><strong>fleetd #615 — the whole body is wrapped in one {@code try}.</strong> The two {@code
|
||||
* agents.send} calls below are not wrapped individually: {@code send} → {@code agentCall} →
|
||||
* {@code herdr.call} can throw an unchecked {@link dev.ltms.fleet.herdr.HerdrException} (see
|
||||
* {@code AgentControl.java}), and the production {@code continuationRunner} is a bare virtual
|
||||
* thread with no uncaught-exception handler (see this class's public constructor). Before this
|
||||
* fix, either throw killed the continuation thread silently, leaving the {@link
|
||||
* RollState#IN_PROGRESS} entry {@link #confirm} wrote at hand-off stuck forever — {@link
|
||||
* #status} had no way to tell a dead roll from one still genuinely running. The {@code catch}
|
||||
* below is scoped to the method body rather than to each {@code send} call individually, so it
|
||||
* also covers anything else added to this continuation later, not just today's two call sites —
|
||||
* the same reasoning that put the write-a-terminal-outcome step at each of this method's other
|
||||
* exits (see the {@link RollState#TURN_NEVER_SETTLED} and {@link RollState#CLEAR_NEVER_SETTLED}
|
||||
* branches below) rather than inside the helpers that detect them.</p>
|
||||
*
|
||||
* <p>Only {@link RuntimeException} is caught, matching the local convention {@link
|
||||
* #waitUntilAtTurnBoundary} already set around its own {@code agents.status} call — not the
|
||||
* broader {@link Exception} or {@link Throwable}, which would also swallow something like an
|
||||
* {@link OutOfMemoryError} this continuation has no business handling.</p>
|
||||
*/
|
||||
private void runRollover(PendingRollover p, FleetConfig.LeadRollover cfg) {
|
||||
try {
|
||||
runRolloverUnguarded(p, cfg);
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("lead-rollover: continuation for token={} lead={} threw {} — the roll is dead; "
|
||||
+ "no further step in this continuation will run",
|
||||
p.token(), p.leadTerminal(), e.toString(), e);
|
||||
outcomes.put(p.token(), new RollStatus(RollState.FAILED,
|
||||
"the roll's continuation threw " + e.toString() + " — the roll is dead and will "
|
||||
+ "not retry itself; check the daemon log for the stack trace, then open() "
|
||||
+ "a fresh rollover request"));
|
||||
}
|
||||
}
|
||||
|
||||
/** The actual body of {@link #runRollover}, unwrapped — see that method's javadoc for the catch. */
|
||||
private void runRolloverUnguarded(PendingRollover p, FleetConfig.LeadRollover cfg) {
|
||||
String lead = p.leadTerminal();
|
||||
long rollStartMillis = nowMillis.getAsLong();
|
||||
TurnSettleResult turnResult = waitUntilAtTurnBoundary(lead, cfg.turnSettleSeconds());
|
||||
|
||||
@@ -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)));
|
||||
@@ -1626,7 +1665,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}). */
|
||||
@@ -1635,7 +1674,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);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1652,7 +1691,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}). */
|
||||
@@ -1661,7 +1700,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);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1683,7 +1722,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);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1707,14 +1746,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()
|
||||
@@ -1723,7 +1772,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
|
||||
@@ -2018,15 +2068,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);
|
||||
@@ -2035,9 +2095,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)) {
|
||||
|
||||
@@ -2,6 +2,7 @@ package dev.ltms.fleet.msg;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.lead.LeadContextGauge;
|
||||
import dev.ltms.fleet.mcp.PrimaryRegistry;
|
||||
import dev.ltms.fleet.metrics.FleetMetrics;
|
||||
import dev.ltms.fleet.metrics.Metrics;
|
||||
@@ -13,6 +14,7 @@ import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.LongSupplier;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
@@ -44,6 +46,15 @@ import java.util.function.Supplier;
|
||||
* loop stands down, so two competing injections never start two turns in the same pane
|
||||
* (constraint 6).</li>
|
||||
* </ol>
|
||||
*
|
||||
* <p><b>fleetd #609 — context-high notice.</b> Optionally ({@code contextHighNudge}, opt-in like the
|
||||
* loop itself), a tick that finds the lead's own {@link LeadContextGauge} reading at {@link
|
||||
* LeadContextGauge.State#HIGH} appends a text notice to whatever nudge it sends, telling the lead to
|
||||
* consider {@code fleet_handover}. This is text only — it never rolls a pane itself. It fires once per
|
||||
* HIGH stretch (a latch, cleared only by a later {@code OK} reading — {@code UNKNOWN} neither sets nor
|
||||
* clears it, since "I could not look" must not be read as "it got better"), and it never spends the
|
||||
* quiet-nudge budget: an idle, quiet, HIGH-context lead is exactly the case {@link Action#QUIET_DONE}
|
||||
* would otherwise swallow, and it is the one case most worth interrupting the quiet cap for.
|
||||
*/
|
||||
public final class LeadHeartbeatLoop {
|
||||
|
||||
@@ -63,11 +74,15 @@ public final class LeadHeartbeatLoop {
|
||||
private final long backoffMs;
|
||||
private final int quietNudgeCap;
|
||||
private final Metrics metrics; // CB-512 pattern: nullable — no registry in unit tests
|
||||
private final LeadContextSource contextSource; // fleetd #609
|
||||
private final boolean contextHighNudge; // fleetd #609: opt-in, like the loop itself
|
||||
|
||||
/** When the current idle stretch began (nanos), or {@link #NOT_IDLE}. Single scheduler thread only. */
|
||||
private long idleSinceNanos = NOT_IDLE;
|
||||
/** Consecutive nudges that found no pending fleet state. Single scheduler thread only. */
|
||||
private int quietCount = 0;
|
||||
/** fleetd #609: latched "already told this lead about this HIGH stretch". Single scheduler thread only. */
|
||||
private boolean contextNotified = false;
|
||||
|
||||
/** Constructor with an injectable clock and no metric registry (unit tests, or wiring that opts out). */
|
||||
public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
||||
@@ -75,7 +90,7 @@ public final class LeadHeartbeatLoop {
|
||||
ScheduledExecutorService scheduler, LongSupplier clock,
|
||||
long idleAfterNanos, long backoffMs, int quietNudgeCap) {
|
||||
this(primaryRegistry, agents, inbox, roster, pushLoop, scheduler, clock,
|
||||
idleAfterNanos, backoffMs, quietNudgeCap, null);
|
||||
idleAfterNanos, backoffMs, quietNudgeCap, null, LeadContextSource.none(), false);
|
||||
}
|
||||
|
||||
/** As above, with a metric registry (the CB-512 pattern) so nudge outcomes are counted. */
|
||||
@@ -83,6 +98,20 @@ public final class LeadHeartbeatLoop {
|
||||
Supplier<List<MemberSession>> roster, ReplyPushLoop pushLoop,
|
||||
ScheduledExecutorService scheduler, LongSupplier clock,
|
||||
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics) {
|
||||
this(primaryRegistry, agents, inbox, roster, pushLoop, scheduler, clock,
|
||||
idleAfterNanos, backoffMs, quietNudgeCap, metrics, LeadContextSource.none(), false);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #609: as above, plus the lead's own context source and whether a HIGH reading should
|
||||
* append a hand-over notice to the loop's nudge. Pass {@link LeadContextSource#none()} and
|
||||
* {@code false} to keep the pre-#609 behaviour exactly (both existing public constructors do).
|
||||
*/
|
||||
public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
||||
Supplier<List<MemberSession>> roster, ReplyPushLoop pushLoop,
|
||||
ScheduledExecutorService scheduler, LongSupplier clock,
|
||||
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics,
|
||||
LeadContextSource contextSource, boolean contextHighNudge) {
|
||||
this.primaryRegistry = primaryRegistry;
|
||||
this.agents = agents;
|
||||
this.inbox = inbox;
|
||||
@@ -94,6 +123,19 @@ public final class LeadHeartbeatLoop {
|
||||
this.backoffMs = backoffMs;
|
||||
this.quietNudgeCap = quietNudgeCap;
|
||||
this.metrics = metrics;
|
||||
this.contextSource = contextSource;
|
||||
this.contextHighNudge = contextHighNudge;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #609: one lead's own context reading, keyed by its terminal id — the same injected-source
|
||||
* idiom {@code FleetMcp.LeadSeatSource}/{@code FleetMcp.LeadConfigDirSource} already use.
|
||||
*/
|
||||
public record LeadContextSource(Function<String, LeadContextGauge.Reading> readingFor) {
|
||||
/** Inert source — every lead reads UNKNOWN, so the context notice can never fire. */
|
||||
public static LeadContextSource none() {
|
||||
return new LeadContextSource(_ -> LeadContextGauge.Reading.unknown());
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -126,8 +168,11 @@ public final class LeadHeartbeatLoop {
|
||||
STAND_DOWN
|
||||
}
|
||||
|
||||
/** The outcome of one decision: the action plus the state to persist for the next tick. */
|
||||
record Decision(Action action, Long idleSinceNanos, int quietCount) {}
|
||||
/**
|
||||
* The outcome of one decision: the action, the state to persist for the next tick, and (fleetd
|
||||
* #609) whether the lead has now been told about the current HIGH context stretch.
|
||||
*/
|
||||
record Decision(Action action, Long idleSinceNanos, int quietCount, boolean contextNotified) {}
|
||||
|
||||
/**
|
||||
* Pure decision function: given the current loop state and fleet/lead facts, return what to do
|
||||
@@ -142,34 +187,47 @@ public final class LeadHeartbeatLoop {
|
||||
* @param pushLoopActive whether {@link ReplyPushLoop} is currently nudging some target (constraint 6)
|
||||
* @param leadKnown whether a lead terminal is known to nudge at all
|
||||
* @param fleet a snapshot of the pending fleet state (constraint 5)
|
||||
* @param context fleetd #609: the lead's own {@link LeadContextGauge} reading for this tick
|
||||
* @param contextNotified fleetd #609: whether the lead has already been told about the current HIGH
|
||||
* stretch — a latch, carried forward by {@link #applyDecision}
|
||||
* @return the action to take and the state to persist
|
||||
*/
|
||||
Decision decide(long nowNanos, Long idleSinceNanos, int quietCount, AgentStatus status,
|
||||
boolean pushLoopActive, boolean leadKnown, FleetState fleet) {
|
||||
boolean pushLoopActive, boolean leadKnown, FleetState fleet,
|
||||
LeadContextGauge.State context, boolean contextNotified) {
|
||||
// fleetd #609: re-arm the latch only on a positive OK reading. UNKNOWN means "I could not
|
||||
// look", not "it got better" — re-arming on UNKNOWN would let a flapping gauge (a transcript
|
||||
// read that misses one tick) nudge a full lead again on every recovery, defeating the "once
|
||||
// per HIGH stretch" promise. Computed once, up front, so every gate below carries it forward
|
||||
// unchanged unless it is the gate that actually discharges it.
|
||||
boolean latch = context == LeadContextGauge.State.OK ? false : contextNotified;
|
||||
boolean contextHigh = contextHighNudge && context == LeadContextGauge.State.HIGH;
|
||||
|
||||
// Constraint 6: while ReplyPushLoop is actively nudging the lead, injecting a second,
|
||||
// competing prompt into the same pane would start a second turn — racing loops multiply
|
||||
// turns and context burn. Stand aside, and treat the active push as real state (re-arm the
|
||||
// quiet counter), because the reply that drove it is exactly the kind of new state that
|
||||
// should reset the cap.
|
||||
// should reset the cap. The context latch is untouched: standing down must not spend the
|
||||
// one notice this stretch gets.
|
||||
if (pushLoopActive) {
|
||||
return new Decision(Action.STAND_DOWN, idleSinceNanos, 0);
|
||||
return new Decision(Action.STAND_DOWN, idleSinceNanos, 0, latch);
|
||||
}
|
||||
// Constraint 2: a WORKING lead is making progress and must NOT be touched; an unreadable
|
||||
// status (read failure, or the agent is gone) is safest treated the same way — never inject
|
||||
// into a state we cannot read. Either way, reset the idle window and the quiet counter: the
|
||||
// lead was / may be active, so the next idle stretch must count its own quiet period fresh.
|
||||
if (status == null || !status.injectable()) {
|
||||
return new Decision(Action.LEAD_BUSY, null, 0);
|
||||
return new Decision(Action.LEAD_BUSY, null, 0, latch);
|
||||
}
|
||||
if (idleSinceNanos == null) {
|
||||
// The lead just became injectable — record the start of an idle stretch and wait out the
|
||||
// debounce quiet period before ever nudging (constraint 3).
|
||||
return new Decision(Action.WAIT_IDLE, nowNanos, quietCount);
|
||||
return new Decision(Action.WAIT_IDLE, nowNanos, quietCount, latch);
|
||||
}
|
||||
if (nowNanos - idleSinceNanos < idleAfterNanos) {
|
||||
// Still within the quiet period: the lead that just finished a turn sits momentarily idle
|
||||
// and must not be re-prompted into every natural pause.
|
||||
return new Decision(Action.WAIT_IDLE, idleSinceNanos, quietCount);
|
||||
return new Decision(Action.WAIT_IDLE, idleSinceNanos, quietCount, latch);
|
||||
}
|
||||
// Past the quiet period with an injectable lead: it is a genuine candidate for a nudge. Two
|
||||
// gating facts decide whether and how:
|
||||
@@ -177,23 +235,33 @@ public final class LeadHeartbeatLoop {
|
||||
// No lead terminal is known yet (e.g. the registry has not learned one) — there is nobody
|
||||
// to nudge. Keep waiting; the window stays open so discovery re-arms it without a fresh
|
||||
// quiet period.
|
||||
return new Decision(Action.WAIT_IDLE, idleSinceNanos, quietCount);
|
||||
return new Decision(Action.WAIT_IDLE, idleSinceNanos, quietCount, latch);
|
||||
}
|
||||
if (fleet.hasPending()) {
|
||||
// Real fleet state is waiting — a worker reply or a DONE session. This is new state, so
|
||||
// it resets the quiet counter (constraint 4) and the lead is nudged to go collect it.
|
||||
return new Decision(Action.INJECT, idleSinceNanos, 0);
|
||||
// it resets the quiet counter (constraint 4) and the lead is nudged to go collect it. The
|
||||
// nudge text carries the context notice too when contextHigh — see injectNudge/contextNotice
|
||||
// — so this route discharges the same duty and must set the latch.
|
||||
return new Decision(Action.INJECT, idleSinceNanos, 0, latch || contextHigh);
|
||||
}
|
||||
if (contextHigh && !latch) {
|
||||
// fleetd #609: the lead is idle, its context is full, and nothing is pending. This is the
|
||||
// one case the quiet cap would otherwise swallow, and it is exactly when the lead most
|
||||
// needs to hear it. Fire once per HIGH stretch, and do NOT spend the quiet budget on it:
|
||||
// this is an event notice, not a "are you still there" nudge.
|
||||
return new Decision(Action.INJECT, idleSinceNanos, quietCount, true);
|
||||
}
|
||||
if (quietCount < quietNudgeCap) {
|
||||
// Nothing is pending, but the cap is not exhausted: nudge anyway, telling the lead
|
||||
// exactly that nothing is waiting so it can choose to stand down rather than hunt
|
||||
// (constraint 5). Count it toward the consecutive-quiet cap.
|
||||
return new Decision(Action.INJECT, idleSinceNanos, quietCount + 1);
|
||||
// (constraint 5). Count it toward the consecutive-quiet cap. This nudge also carries the
|
||||
// context notice when contextHigh (already latched above, or being latched now).
|
||||
return new Decision(Action.INJECT, idleSinceNanos, quietCount + 1, latch || contextHigh);
|
||||
}
|
||||
// Nothing pending and the cap is exhausted: stop nudging until real state appears again
|
||||
// (constraint 4). The loop still ticks on backoff so a genuinely new reply or session change
|
||||
// re-arms it — QUIET_DONE stops injection, not observation.
|
||||
return new Decision(Action.QUIET_DONE, idleSinceNanos, quietCount);
|
||||
return new Decision(Action.QUIET_DONE, idleSinceNanos, quietCount, latch);
|
||||
}
|
||||
|
||||
// --- loop ----------------------------------------------------------------------------------
|
||||
@@ -208,57 +276,152 @@ public final class LeadHeartbeatLoop {
|
||||
scheduler.schedule(this::tick, backoffMs, TimeUnit.MILLISECONDS);
|
||||
}
|
||||
|
||||
/** One loop tick, every {@link #backoffMs} — the thin scheduler around {@link #decide}. */
|
||||
private void tick() {
|
||||
/** One loop tick, every {@link #backoffMs} — the thin scheduler around {@link #decide}. Package-private
|
||||
* (mirroring {@link ReplyPushLoop#tick(String)}) so tests can drive it directly with a fake clock and a
|
||||
* fake {@link AgentControl} instead of racing the scheduler thread. */
|
||||
void tick() {
|
||||
boolean leadKnown = primaryRegistry.primaryTerminal().isPresent();
|
||||
FleetState fleet = snapshot(inbox, roster);
|
||||
AgentStatus status = AgentStatus.UNKNOWN;
|
||||
LeadContextGauge.Reading reading = LeadContextGauge.Reading.unknown();
|
||||
if (leadKnown) {
|
||||
String leadTerminal = primaryRegistry.primaryTerminal().orElseThrow();
|
||||
try {
|
||||
status = agents.status(primaryRegistry.primaryTerminal().orElseThrow());
|
||||
status = agents.status(leadTerminal);
|
||||
} catch (RuntimeException e) {
|
||||
// A failed status read degrades to "unknown" — decide() treats that like a busy lead
|
||||
// and never injects into a state it cannot read. Retry on the next backoff.
|
||||
log.debug("idle-heartbeat: status check failed for lead, will retry: {}", e.toString());
|
||||
}
|
||||
// fleetd #609: read the lead's own context regardless of status — decide() still gates on
|
||||
// status first (constraint 2), so this is harmless work on a WORKING lead and lets the
|
||||
// latch state stay accurate for whenever the lead does go idle.
|
||||
reading = contextSource.readingFor().apply(leadTerminal);
|
||||
}
|
||||
|
||||
Decision d = decide(clock.getAsLong(),
|
||||
idleSinceNanos == NOT_IDLE ? null : idleSinceNanos,
|
||||
quietCount, status, pushLoop.isActive(), leadKnown, fleet);
|
||||
quietCount, status, pushLoop.isActive(), leadKnown, fleet,
|
||||
reading.state(), contextNotified);
|
||||
applyDecision(d);
|
||||
switch (d.action()) {
|
||||
case INJECT -> injectNudge(fleet);
|
||||
case QUIET_DONE -> countNudge("exhausted");
|
||||
case WAIT_IDLE, LEAD_BUSY, STAND_DOWN -> { /* nothing to inject, nothing to count */ }
|
||||
case INJECT -> injectNudge(d, fleet, reading);
|
||||
case QUIET_DONE -> {
|
||||
countNudge("exhausted");
|
||||
contextNotified = d.contextNotified();
|
||||
}
|
||||
case WAIT_IDLE, LEAD_BUSY, STAND_DOWN -> contextNotified = d.contextNotified();
|
||||
}
|
||||
scheduleNext();
|
||||
}
|
||||
|
||||
/** Persist the state a decision returned, so the next tick starts from it. */
|
||||
/**
|
||||
* Persist the idle/quiet state a decision returned, so the next tick starts from it.
|
||||
*
|
||||
* <p>fleetd #609 review: the context latch ({@link #contextNotified}) is deliberately <em>not</em>
|
||||
* set here any more. Setting it from the decision unconditionally — before {@link #injectNudge} even
|
||||
* tries to send — is exactly the review's blocker: a decision to notify is not the same fact as "the
|
||||
* notice reached the pane". Every branch of {@link #tick} now assigns {@link #contextNotified} itself,
|
||||
* once it knows whether a send happened and whether it carried the notice (see {@link #injectNudge}).
|
||||
*/
|
||||
private void applyDecision(Decision d) {
|
||||
idleSinceNanos = d.idleSinceNanos() == null ? NOT_IDLE : d.idleSinceNanos();
|
||||
quietCount = d.quietCount();
|
||||
}
|
||||
|
||||
/** Send the nudge to the known lead. */
|
||||
private void injectNudge(FleetState fleet) {
|
||||
/**
|
||||
* Send the nudge to the known lead, with the fleetd #609 context notice appended when it applies, and
|
||||
* persist the context latch based on what actually happened this tick — not merely what {@code d}
|
||||
* chose to attempt.
|
||||
*/
|
||||
private void injectNudge(Decision d, FleetState fleet, LeadContextGauge.Reading reading) {
|
||||
// fleetd #609 review: build the notice from the latch as it stood BEFORE this tick's decision —
|
||||
// d.contextNotified() is the value to persist once delivery is confirmed, not the value the text
|
||||
// itself should be built from. Otherwise a HIGH stretch that is still latched would never see the
|
||||
// notice at all, defeating the very check this fixes.
|
||||
String notice = contextNotice(contextHighNudge, reading, contextNotified);
|
||||
var lead = primaryRegistry.primaryTerminal();
|
||||
if (lead.isEmpty()) {
|
||||
return; // the lead disappeared between the decision and the injection
|
||||
}
|
||||
String leadTerminal = lead.get();
|
||||
boolean sent = lead.isPresent() && trySend(lead.get(), fleet.nudgeText() + notice, notice);
|
||||
// The latch becomes true only when all three hold: decide() chose to notify, a notice was
|
||||
// actually included in the text, and the send reached the pane without throwing. Whenever no
|
||||
// notice was attempted (disabled, not HIGH, or already latched), nothing was promised to the lead
|
||||
// this tick, so apply the decision's own carried-forward value unconditionally — that is how the
|
||||
// OK-only re-arm rule and STAND_DOWN's "don't burn the notice" rule keep working through this path
|
||||
// too. A lead that disappeared between the decision and the send (lead.isEmpty()) is treated the
|
||||
// same as a failed send: nothing reached the pane, so the latch must not be set.
|
||||
contextNotified = notice.isEmpty() ? d.contextNotified() : sent;
|
||||
}
|
||||
|
||||
/**
|
||||
* Attempt one herdr send and count its outcome. Returns whether {@code agents.send} returned without
|
||||
* throwing — the caller ({@link #injectNudge}) needs this to decide whether the fleetd #609 context
|
||||
* latch may be persisted as set.
|
||||
*/
|
||||
private boolean trySend(String leadTerminal, String text, String notice) {
|
||||
try {
|
||||
agents.send(leadTerminal, fleet.nudgeText());
|
||||
agents.send(leadTerminal, text);
|
||||
log.debug("idle-heartbeat: nudge sent to lead {} (quiet nudges so far in this stretch: {})",
|
||||
leadTerminal, quietCount);
|
||||
countNudge("sent");
|
||||
// fleetd #609: a nudge that carries the context notice is counted under its own outcome so
|
||||
// it is visible in /metrics — one count per nudge either way, never two.
|
||||
countNudge(notice.isEmpty() ? "sent" : "sent_context");
|
||||
return true;
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("idle-heartbeat: failed to nudge lead {}: {}", leadTerminal, e.toString());
|
||||
countNudge("failed");
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #609: the text appended to a nudge when the lead's own context is full — {@code ""}
|
||||
* whenever the notice does not apply, so callers can unconditionally append this without an extra
|
||||
* branch. Wording stays plain (CEFR B1) and honest that only the operator approves a roll — this
|
||||
* loop only ever prints text, it never calls {@code fleet_handover} itself.
|
||||
*
|
||||
* @param enabled the {@code leadHeartbeat.contextHighNudge} config flag
|
||||
* @param reading the lead's current {@link LeadContextGauge} reading
|
||||
* @return the notice text (starting with a leading space, to append directly after {@link
|
||||
* FleetState#nudgeText()}), or {@code ""} when disabled or the state is not {@code HIGH}
|
||||
*/
|
||||
static String contextNotice(boolean enabled, LeadContextGauge.Reading reading) {
|
||||
return contextNotice(enabled, reading, false);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #609 review: as {@link #contextNotice(boolean, LeadContextGauge.Reading)}, but also gated on
|
||||
* {@code alreadyNotified} — the context latch as it stood <em>before</em> the current tick's decision.
|
||||
* Without this gate, every pending-driven {@code INJECT} that lands while the context stays {@code
|
||||
* HIGH} would re-append the full notice on top of an already-latched stretch, making the notice's own
|
||||
* closing sentence ("You will not be told again until your context reads ok.") false. {@link
|
||||
* #injectNudge} is the only caller that passes a non-default {@code alreadyNotified}.
|
||||
*
|
||||
* @param alreadyNotified whether the lead has already been told about the current HIGH stretch
|
||||
*/
|
||||
static String contextNotice(boolean enabled, LeadContextGauge.Reading reading, boolean alreadyNotified) {
|
||||
if (!enabled || alreadyNotified || reading.state() != LeadContextGauge.State.HIGH) {
|
||||
return "";
|
||||
}
|
||||
StringBuilder sb = new StringBuilder(" Your own context is nearly full");
|
||||
String compactionWord = reading.compactions() == 1 ? "compaction" : "compactions";
|
||||
if (reading.tokens() != null) {
|
||||
sb.append(": ").append(reading.tokens()).append(" tokens used, ")
|
||||
.append(reading.compactions()).append(' ').append(compactionWord).append(" so far.");
|
||||
} else {
|
||||
// A HIGH reading always carries a non-null token count today: LeadContextGauge only
|
||||
// reaches HIGH by comparing a number against HIGH_THRESHOLD_TOKENS. That invariant
|
||||
// lives in another class and nothing asserts it, so this branch does not rely on it —
|
||||
// it drops the token clause rather than printing "null tokens".
|
||||
sb.append(" (").append(reading.compactions()).append(' ').append(compactionWord)
|
||||
.append(" so far).");
|
||||
}
|
||||
sb.append(" A fresh session would work better. To hand over: call fleet_handover(action=\"open\"), "
|
||||
+ "write the file it names, ask the operator, then call fleet_handover(action=\"confirm\", "
|
||||
+ "token, operatorConfirmed). Only the operator can approve the roll. You will not be told "
|
||||
+ "again until your context reads ok.");
|
||||
return sb.toString();
|
||||
}
|
||||
|
||||
/** Schedule the next tick on the scheduler thread pool. */
|
||||
private void scheduleNext() {
|
||||
scheduler.schedule(this::tick, backoffMs, TimeUnit.MILLISECONDS);
|
||||
|
||||
@@ -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,156 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.lead.LeadContextGauge;
|
||||
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.function.Function;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
|
||||
/**
|
||||
* fleetd #609: {@link Fleetd#leadContextLookup} is the factory {@code Fleetd.main} wires into
|
||||
* {@code LeadHeartbeatLoop.LeadContextSource} so the idle-lead heartbeat can read a lead's own
|
||||
* {@link LeadContextGauge} reading by its terminal id. Three hops: terminal → lead name (unknown ⇒
|
||||
* UNKNOWN), lead name → configDir, and {@code agents.get(terminal)} for the live session id and
|
||||
* agent type the gauge itself needs — a herdr failure on that last hop must degrade to UNKNOWN, not
|
||||
* throw and kill the heartbeat's own tick.
|
||||
*
|
||||
* <p>{@code AgentControl.agentCall} resolves a {@code term_}-prefixed target's pane id via a first
|
||||
* {@code agent.list} round trip (see its own javadoc); this class's terminal id deliberately does
|
||||
* NOT start with {@code term_} so the stub {@link HerdrClient} below only needs to answer
|
||||
* {@code agent.get} — the one call this factory actually depends on.
|
||||
*/
|
||||
class FleetdLeadContextLookupTest {
|
||||
|
||||
private static final String LEAD_TERMINAL = "leadpane1";
|
||||
private static final String LEAD_NAME = "opus";
|
||||
private static final String SESSION_ID = "sess-609-happy-path";
|
||||
|
||||
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 void 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);
|
||||
}
|
||||
|
||||
/** An {@link AgentControl} whose every {@code agent.get} answers with the given session/type/status. */
|
||||
private static AgentControl agentControlStub(String sessionId, String agentType, String status) {
|
||||
HerdrClient client = new HerdrClient() {
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) {
|
||||
if (!"agent.get".equals(method)) {
|
||||
throw new HerdrException("stub has no canned response for " + method);
|
||||
}
|
||||
try {
|
||||
String sessionField = sessionId == null ? ""
|
||||
: ",\"agent_session\":{\"kind\":\"id\",\"value\":\"" + sessionId + "\"}";
|
||||
String agentField = agentType == null ? "null" : "\"" + agentType + "\"";
|
||||
return new ObjectMapper().readTree(("""
|
||||
{"type":"agent_info","agent":{"terminal_id":"%s","agent":%s,
|
||||
"agent_status":"%s"%s}}""")
|
||||
.formatted(LEAD_TERMINAL, agentField, status, sessionField));
|
||||
} catch (Exception e) {
|
||||
throw new HerdrException("stub decode failed", e);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
};
|
||||
return new AgentControl(client);
|
||||
}
|
||||
|
||||
/** An {@link AgentControl} whose every herdr call fails — models a herdr hiccup mid-tick. */
|
||||
private static AgentControl throwingAgentControl() {
|
||||
HerdrClient client = new HerdrClient() {
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) {
|
||||
throw new HerdrException("herdr unreachable (stub)");
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
};
|
||||
return new AgentControl(client);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("an unrecognised terminal resolves to UNKNOWN, not a thrown exception")
|
||||
void unrecognisedTerminalResolvesToUnknown() {
|
||||
Function<String, LeadContextGauge.Reading> lookup = Fleetd.leadContextLookup(
|
||||
new LeadContextGauge(), throwingAgentControl(), Map::of, name -> null);
|
||||
|
||||
LeadContextGauge.Reading reading = assertDoesNotThrow(() -> lookup.apply("ghost-terminal"));
|
||||
|
||||
assertEquals(LeadContextGauge.State.UNKNOWN, reading.state());
|
||||
assertNull(reading.tokens());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("agents.get throwing degrades to UNKNOWN, no exception escapes")
|
||||
void agentsGetThrowingDegradesToUnknown() {
|
||||
Map<String, String> liveLeadTerminals = Map.of(LEAD_TERMINAL, LEAD_NAME);
|
||||
Function<String, LeadContextGauge.Reading> lookup = Fleetd.leadContextLookup(
|
||||
new LeadContextGauge(), throwingAgentControl(), () -> liveLeadTerminals, name -> null);
|
||||
|
||||
LeadContextGauge.Reading reading = assertDoesNotThrow(() -> lookup.apply(LEAD_TERMINAL));
|
||||
|
||||
assertEquals(LeadContextGauge.State.UNKNOWN, reading.state(),
|
||||
"a herdr failure resolving the live agent must degrade to UNKNOWN, never kill the heartbeat's tick");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("the happy path resolves the configDir, sessionId and agentType through to the gauge")
|
||||
void happyPathPassesResolvedFactsThroughToTheGauge(@TempDir Path tmp) throws IOException {
|
||||
writeTranscript(tmp, SESSION_ID, 12_345);
|
||||
Map<String, String> liveLeadTerminals = Map.of(LEAD_TERMINAL, LEAD_NAME);
|
||||
AgentControl agents = agentControlStub(SESSION_ID, "claude", "idle");
|
||||
Function<String, LeadContextGauge.Reading> lookup = Fleetd.leadContextLookup(
|
||||
new LeadContextGauge(), agents, () -> liveLeadTerminals,
|
||||
name -> LEAD_NAME.equals(name) ? tmp.toString() : null);
|
||||
|
||||
LeadContextGauge.Reading reading = lookup.apply(LEAD_TERMINAL);
|
||||
|
||||
assertEquals(LeadContextGauge.State.OK, reading.state(),
|
||||
"the resolved configDir + the live agent's own sessionId/agentType must reach the gauge — "
|
||||
+ "a wrong hop anywhere in the chain would read no transcript and report UNKNOWN instead");
|
||||
assertEquals(12_345L, reading.tokens());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a lead whose agent type is not claude still resolves to UNKNOWN, never a crash")
|
||||
void nonClaudeAgentTypeResolvesToUnknown(@TempDir Path tmp) throws IOException {
|
||||
writeTranscript(tmp, SESSION_ID, 12_345);
|
||||
Map<String, String> liveLeadTerminals = Map.of(LEAD_TERMINAL, LEAD_NAME);
|
||||
AgentControl agents = agentControlStub(SESSION_ID, "opencode", "idle");
|
||||
Function<String, LeadContextGauge.Reading> lookup = Fleetd.leadContextLookup(
|
||||
new LeadContextGauge(), agents, () -> liveLeadTerminals,
|
||||
name -> LEAD_NAME.equals(name) ? tmp.toString() : null);
|
||||
|
||||
LeadContextGauge.Reading reading = lookup.apply(LEAD_TERMINAL);
|
||||
|
||||
assertEquals(LeadContextGauge.State.UNKNOWN, reading.state(),
|
||||
"the agentType hop must reach the gauge too — a non-claude peer must not be misread as claude");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,108 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.lead.LeadContextGauge;
|
||||
import dev.ltms.fleet.msg.LeadHeartbeatLoop;
|
||||
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 static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* fleetd #609: {@code Fleetd.main}'s {@code LeadHeartbeatLoop.LeadContextSource} local could be
|
||||
* swapped for a bare {@code LeadHeartbeatLoop.LeadContextSource.none()} at the call site — compiling
|
||||
* with 0 errors and leaving every pre-existing test green — exactly the shape #602/#606 already found
|
||||
* for {@code LeadConfigDirSource} (see {@code FleetdLeadConfigDirSourceWiringTest}'s own javadoc for
|
||||
* the measured version of that gap).
|
||||
*
|
||||
* <p>The fix follows the same pattern: {@link Fleetd#leadContextSource} is the extracted,
|
||||
* directly-callable factory {@code main} calls to build the source it hands {@code
|
||||
* LeadHeartbeatLoop}'s constructor. This test calls that exact factory and asserts it resolves a
|
||||
* REAL reading off a real transcript file — a property that would be false if {@link
|
||||
* Fleetd#leadContextSource} were mutated to {@code return LeadHeartbeatLoop.LeadContextSource.none();}.
|
||||
*
|
||||
* <p>What this class does not and cannot cover: {@code main}'s own one-line call to this factory
|
||||
* could itself be swapped for {@code LeadHeartbeatLoop.LeadContextSource.none()}, bypassing this
|
||||
* factory entirely — the same structural gap {@code FleetdLeadConfigDirSourceWiringTest} names for
|
||||
* its own factory, and for the same reason (no test in this codebase calls {@code Fleetd.main} far
|
||||
* enough to observe which factory call it made).
|
||||
*/
|
||||
class FleetdLeadContextSourceWiringTest {
|
||||
|
||||
/** Deliberately not {@code term_}-prefixed — see {@code FleetdLeadContextLookupTest}'s class doc. */
|
||||
private static final String LEAD_TERMINAL = "leadpane1";
|
||||
private static final String LEAD_NAME = "opus";
|
||||
private static final String SESSION_ID = "sess-609-wiring";
|
||||
|
||||
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}}}";
|
||||
}
|
||||
|
||||
private static void 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);
|
||||
}
|
||||
|
||||
private static AgentControl agentControlStub() {
|
||||
HerdrClient client = new HerdrClient() {
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) {
|
||||
if (!"agent.get".equals(method)) {
|
||||
throw new HerdrException("stub has no canned response for " + method);
|
||||
}
|
||||
try {
|
||||
return new ObjectMapper().readTree(("""
|
||||
{"type":"agent_info","agent":{"terminal_id":"%s","agent":"claude",
|
||||
"agent_status":"idle","agent_session":{"kind":"id","value":"%s"}}}""")
|
||||
.formatted(LEAD_TERMINAL, SESSION_ID));
|
||||
} catch (Exception e) {
|
||||
throw new HerdrException("stub decode failed", e);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
};
|
||||
return new AgentControl(client);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("main's factory resolves a REAL reading, not the inert none() answer")
|
||||
void resolvesARealReadingNotTheInertNoneAnswer(@TempDir Path tmp) throws IOException {
|
||||
writeTranscript(tmp, SESSION_ID, 54_321);
|
||||
Map<String, String> liveLeadTerminals = Map.of(LEAD_TERMINAL, LEAD_NAME);
|
||||
|
||||
LeadHeartbeatLoop.LeadContextSource source = Fleetd.leadContextSource(new LeadContextGauge(),
|
||||
agentControlStub(), () -> liveLeadTerminals, name -> LEAD_NAME.equals(name) ? tmp.toString() : null);
|
||||
|
||||
LeadContextGauge.Reading reading = source.readingFor().apply(LEAD_TERMINAL);
|
||||
|
||||
assertEquals(LeadContextGauge.State.OK, reading.state(),
|
||||
"mutating Fleetd.leadContextSource's own body to `return LeadHeartbeatLoop.LeadContextSource.none();` "
|
||||
+ "must fail this assertion");
|
||||
assertEquals(54_321L, reading.tokens());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("an unrecognised lead terminal resolves to UNKNOWN, not a thrown exception")
|
||||
void unrecognisedTerminalResolvesToUnknown() {
|
||||
LeadHeartbeatLoop.LeadContextSource source = Fleetd.leadContextSource(new LeadContextGauge(),
|
||||
agentControlStub(), Map::of, name -> null);
|
||||
|
||||
assertEquals(LeadContextGauge.State.UNKNOWN, source.readingFor().apply("ghost-terminal").state());
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,248 @@
|
||||
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;
|
||||
import static org.junit.jupiter.api.Assumptions.assumeFalse;
|
||||
|
||||
/**
|
||||
* 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 {
|
||||
// setReadable(false) really did clear the read bit (asserted above), but that alone
|
||||
// does not prove the file is UNREADABLE: running as root (e.g. a CI container) ignores
|
||||
// the read bit and opens the file anyway. Files.isReadable checks what actually happens
|
||||
// on open, not the bit. When it still reports readable, this test cannot create the
|
||||
// condition it needs on this machine, so it skips honestly instead of asserting on a
|
||||
// state that was never reached. A skip here means "I could not set up the case", NOT
|
||||
// "the UNKNOWN behaviour is fine" -- it is not evidence either way.
|
||||
assumeFalse(Files.isReadable(file),
|
||||
"runs as root (CI container): the read bit does not stop root, so this case cannot be set up here");
|
||||
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());
|
||||
}
|
||||
}
|
||||
@@ -1359,4 +1359,92 @@ class LeadRolloverTest {
|
||||
+ "it were the measured wait duration: " + message);
|
||||
}
|
||||
}
|
||||
|
||||
// ---- fleetd #615: a HerdrException out of either unwrapped agents.send call must leave a ----
|
||||
// ---- TERMINAL FAILED outcome, never a stuck IN_PROGRESS ---------------------------------------
|
||||
|
||||
@Test
|
||||
@DisplayName("[fleetd #615 — 1] send() throwing on the /clear call leaves status(token) "
|
||||
+ "reporting FAILED, not stuck at IN_PROGRESS")
|
||||
void sendThrowingOnClearLeavesStatusReportingFailed() throws IOException {
|
||||
FakeHerdr fake = new FakeHerdr(); // default idle — the turn-settle wait passes immediately
|
||||
HerdrClient throwsOnClear = new HerdrClient() {
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) throws HerdrException {
|
||||
if ("agent.prompt".equals(method) && String.valueOf(params).contains("/clear")) {
|
||||
throw new HerdrException("simulated herdr transport failure sending /clear");
|
||||
}
|
||||
return fake.call(method, params);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
fake.close();
|
||||
}
|
||||
};
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(throwsOnClear, cfg(handover.toString()), fixedClock(clock));
|
||||
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
|
||||
assertTrue(decision.accepted(), "every synchronous gate passes; the throw happens only "
|
||||
+ "inside the deferred continuation, which this test's synchronous runner has "
|
||||
+ "already run to completion by the time confirm() returns");
|
||||
|
||||
LeadRollover.RollStatus status = rollover.status(pending.token());
|
||||
assertEquals(LeadRollover.RollState.FAILED, status.state(),
|
||||
"a HerdrException out of the /clear send must leave a TERMINAL FAILED outcome — "
|
||||
+ "before fleetd #615's fix, the continuation thread died silently and "
|
||||
+ "status() was stuck reporting the IN_PROGRESS confirm() wrote at hand-off, "
|
||||
+ "forever: got " + status.state() + " / " + status.detail());
|
||||
assertNotEquals(LeadRollover.RollState.IN_PROGRESS, status.state());
|
||||
assertTrue(status.detail().contains("HerdrException"), "the detail must name the exception "
|
||||
+ "so an operator reading status() has something to act on: " + status.detail());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[fleetd #615 — 2] send() throwing on the bootstrap-text call (after /clear "
|
||||
+ "succeeded and the pane settled) also leaves status(token) reporting FAILED — a "
|
||||
+ "DIFFERENT exit from the /clear-throw case above")
|
||||
void sendThrowingOnBootstrapTextLeavesStatusReportingFailed() throws IOException {
|
||||
FakeHerdr fake = new FakeHerdr(); // default idle throughout — both settle waits pass promptly
|
||||
HerdrClient throwsOnBootstrapText = new HerdrClient() {
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) throws HerdrException {
|
||||
if ("agent.prompt".equals(method) && String.valueOf(params).contains("read the handover file")) {
|
||||
throw new HerdrException("simulated herdr transport failure sending bootstrapText");
|
||||
}
|
||||
return fake.call(method, params);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
fake.close();
|
||||
}
|
||||
};
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(throwsOnBootstrapText, cfg(handover.toString()), fixedClock(clock));
|
||||
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
|
||||
assertTrue(decision.accepted(), "every synchronous gate passes; the throw happens only "
|
||||
+ "inside the deferred continuation, which this test's synchronous runner has "
|
||||
+ "already run to completion by the time confirm() returns");
|
||||
assertEquals(1, promptCallCount(fake), "sanity: /clear was sent and settled — only the "
|
||||
+ "SECOND agent.prompt call (bootstrapText) threw");
|
||||
|
||||
LeadRollover.RollStatus status = rollover.status(pending.token());
|
||||
assertEquals(LeadRollover.RollState.FAILED, status.state(),
|
||||
"a HerdrException out of the bootstrapText send — a DIFFERENT exit from the /clear "
|
||||
+ "throw, reached only after /clear already succeeded and the pane already "
|
||||
+ "settled — must also leave a TERMINAL FAILED outcome, not a stuck "
|
||||
+ "IN_PROGRESS: got " + status.state() + " / " + status.detail());
|
||||
assertNotEquals(LeadRollover.RollState.IN_PROGRESS, status.state());
|
||||
assertTrue(status.detail().contains("HerdrException"), "the detail must name the exception "
|
||||
+ "so an operator reading status() has something to act on: " + status.detail());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
@@ -1,6 +1,11 @@
|
||||
package dev.ltms.fleet.msg;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.lead.LeadContextGauge;
|
||||
import dev.ltms.fleet.mcp.PrimaryRegistry;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
@@ -8,20 +13,29 @@ import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
/**
|
||||
* Unit tests for the CB-551 idle-lead heartbeat: the pure {@link LeadHeartbeatLoop#decide} decision
|
||||
* function, the nudge text, and the {@link LeadHeartbeatLoop#snapshot} fleet snapshot.
|
||||
* function, the nudge text, and the {@link LeadHeartbeatLoop#snapshot} fleet snapshot. Also covers
|
||||
* fleetd #609's context-high notice: the {@code context}/{@code contextNotified} parameters {@code
|
||||
* decide} gained, and the standalone {@link LeadHeartbeatLoop#contextNotice} text builder.
|
||||
*
|
||||
* <p>All decision tests call {@code decide} directly with explicit nanoTime values from an injected
|
||||
* clock — no sleeping, no scheduler races. This mirrors how {@code ReplyPushLoopTest} pins the pure
|
||||
* decision before exercising the loop.
|
||||
*
|
||||
* <p>Every pre-#609 test below passes {@code LeadContextGauge.State.UNKNOWN, false} for the two new
|
||||
* {@code decide} parameters — the same as a lead whose context could not be read and has never been
|
||||
* notified — so each one still pins exactly the behaviour it pinned before this ticket.
|
||||
*/
|
||||
class LeadHeartbeatLoopTest {
|
||||
|
||||
@@ -56,12 +70,18 @@ class LeadHeartbeatLoopTest {
|
||||
return new LeadHeartbeatLoop.FleetState(2, 1, 3, List.of("term_a", "term_b"));
|
||||
}
|
||||
|
||||
/** The loop under test; the scheduler is never invoked on the pure decide path. */
|
||||
/** The loop under test, context notice off; the scheduler is never invoked on the pure decide path. */
|
||||
private static LeadHeartbeatLoop loop(int quietCap) {
|
||||
return loop(quietCap, false);
|
||||
}
|
||||
|
||||
/** As above, with the fleetd #609 {@code contextHighNudge} flag set explicitly. */
|
||||
private static LeadHeartbeatLoop loop(int quietCap, boolean contextHighNudge) {
|
||||
return new LeadHeartbeatLoop(
|
||||
new PrimaryRegistry("term_lead"), null /*agents — unused on the decide path*/,
|
||||
null /*inbox*/, List::of, null /*pushLoop*/, null /*scheduler*/, () -> 0L,
|
||||
IDLE_AFTER_NANOS, 1_000L, quietCap);
|
||||
IDLE_AFTER_NANOS, 1_000L, quietCap, null,
|
||||
LeadHeartbeatLoop.LeadContextSource.none(), contextHighNudge);
|
||||
}
|
||||
|
||||
// ── (a) a working lead is never injected ───────────────────────────────────────────────────
|
||||
@@ -69,7 +89,8 @@ class LeadHeartbeatLoopTest {
|
||||
@Test
|
||||
void aWorkingLeadIsNeverInjected() {
|
||||
LeadHeartbeatLoop.Decision d = loop(3).decide(
|
||||
NOW, IDLE_PAST, 0, AgentStatus.WORKING, false, true, quietFleet());
|
||||
NOW, IDLE_PAST, 0, AgentStatus.WORKING, false, true, quietFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.LEAD_BUSY, d.action(),
|
||||
"a WORKING lead is making progress and must not be touched");
|
||||
assertNull(d.idleSinceNanos(), "a busy lead resets the idle window");
|
||||
@@ -80,7 +101,8 @@ class LeadHeartbeatLoopTest {
|
||||
void anUnreadableStatusIsNeverInjectedEither() {
|
||||
// A failed status read (or a gone agent) must degrade to "do not inject", never hammer the pane.
|
||||
LeadHeartbeatLoop.Decision d = loop(3).decide(
|
||||
NOW, IDLE_PAST, 0, AgentStatus.UNKNOWN, false, true, pendingFleet());
|
||||
NOW, IDLE_PAST, 0, AgentStatus.UNKNOWN, false, true, pendingFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.LEAD_BUSY, d.action(),
|
||||
"never inject into a state the loop cannot read");
|
||||
}
|
||||
@@ -90,7 +112,8 @@ class LeadHeartbeatLoopTest {
|
||||
@Test
|
||||
void justBecameIdleStartsTheDebounceWindow() {
|
||||
LeadHeartbeatLoop.Decision d = loop(3).decide(
|
||||
NOW, null, 0, AgentStatus.IDLE, false, true, quietFleet());
|
||||
NOW, null, 0, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.WAIT_IDLE, d.action(),
|
||||
"the first injectable tick only records the start of the idle stretch");
|
||||
assertEquals(NOW, d.idleSinceNanos(), "the idle window opens at the moment the lead became injectable");
|
||||
@@ -99,7 +122,8 @@ class LeadHeartbeatLoopTest {
|
||||
@Test
|
||||
void idleWithinQuietPeriodIsNotInjected() {
|
||||
LeadHeartbeatLoop.Decision d = loop(3).decide(
|
||||
NOW, IDLE_WITHIN, 0, AgentStatus.IDLE, false, true, quietFleet());
|
||||
NOW, IDLE_WITHIN, 0, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.WAIT_IDLE, d.action(),
|
||||
"a lead idle for 10s (< 300s) has just finished a turn — do not re-prompt it");
|
||||
}
|
||||
@@ -109,7 +133,8 @@ class LeadHeartbeatLoopTest {
|
||||
@Test
|
||||
void idlePastQuietPeriodIsInjected() {
|
||||
LeadHeartbeatLoop.Decision d = loop(3).decide(
|
||||
NOW, IDLE_PAST, 0, AgentStatus.IDLE, false, true, quietFleet());
|
||||
NOW, IDLE_PAST, 0, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.INJECT, d.action(),
|
||||
"a lead continuously idle past the quiet period is the reason to nudge");
|
||||
}
|
||||
@@ -117,14 +142,17 @@ class LeadHeartbeatLoopTest {
|
||||
@Test
|
||||
void blockedAndDoneAreInjectableViewsOfIdle() {
|
||||
assertEquals(LeadHeartbeatLoop.Action.INJECT, loop(3).decide(
|
||||
NOW, IDLE_PAST, 0, AgentStatus.BLOCKED, false, true, quietFleet()).action());
|
||||
NOW, IDLE_PAST, 0, AgentStatus.BLOCKED, false, true, quietFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false).action());
|
||||
assertEquals(LeadHeartbeatLoop.Action.INJECT, loop(3).decide(
|
||||
NOW, IDLE_PAST, 0, AgentStatus.DONE, false, true, quietFleet()).action());
|
||||
NOW, IDLE_PAST, 0, AgentStatus.DONE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false).action());
|
||||
}
|
||||
|
||||
@Test
|
||||
void nothingIsInjectedWhenNoLeadIsKnown() {
|
||||
var d = loop(3).decide(NOW, IDLE_PAST, 0, AgentStatus.IDLE, false, false, pendingFleet());
|
||||
var d = loop(3).decide(NOW, IDLE_PAST, 0, AgentStatus.IDLE, false, false, pendingFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.WAIT_IDLE, d.action(),
|
||||
"with no known lead there is nobody to nudge — keep waiting until one is discovered");
|
||||
assertEquals(IDLE_PAST, d.idleSinceNanos(), "the idle window stays open so discovery re-arms it");
|
||||
@@ -135,7 +163,8 @@ class LeadHeartbeatLoopTest {
|
||||
@Test
|
||||
void quietNudgeCapStopsTheLoopWhenNothingIsPending() {
|
||||
LeadHeartbeatLoop.Decision d = loop(3).decide(
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet());
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.QUIET_DONE, d.action(),
|
||||
"3 consecutive nothing-pending nudges have already happened — stop nagging an empty fleet");
|
||||
}
|
||||
@@ -145,7 +174,8 @@ class LeadHeartbeatLoopTest {
|
||||
@Test
|
||||
void newPendingStateResetsTheQuietCap() {
|
||||
LeadHeartbeatLoop.Decision d = loop(3).decide(
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, pendingFleet());
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, pendingFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.INJECT, d.action(),
|
||||
"real state appearing re-arms the loop past an exhausted cap");
|
||||
assertEquals(0, d.quietCount(), "the pending state resets the consecutive-quiet counter");
|
||||
@@ -156,7 +186,8 @@ class LeadHeartbeatLoopTest {
|
||||
@Test
|
||||
void standsDownWhileReplyPushLoopIsActive() {
|
||||
LeadHeartbeatLoop.Decision d = loop(3).decide(
|
||||
NOW, IDLE_PAST, 2, AgentStatus.IDLE, true, true, quietFleet());
|
||||
NOW, IDLE_PAST, 2, AgentStatus.IDLE, true, true, quietFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.STAND_DOWN, d.action(),
|
||||
"a second injection would start a competing turn — stand aside instead");
|
||||
assertEquals(0, d.quietCount(), "the active push is real state, so it re-arms the cap");
|
||||
@@ -213,4 +244,378 @@ class LeadHeartbeatLoopTest {
|
||||
assertFalse(fs.hasPending());
|
||||
assertEquals(0, fs.liveWorkers());
|
||||
}
|
||||
|
||||
// ── fleetd #609: context-high notice ───────────────────────────────────────────────────────
|
||||
|
||||
/** One gate scenario, keyed by name, over which the off-state/context-independence property is checked. */
|
||||
private record Gate(String name, Long idleSince, int quietCount, AgentStatus status,
|
||||
boolean pushLoopActive, boolean leadKnown, LeadHeartbeatLoop.FleetState fleet) {}
|
||||
|
||||
private static List<Gate> allEightGates() {
|
||||
return List.of(
|
||||
new Gate("1-standDown", IDLE_PAST, 2, AgentStatus.IDLE, true, true, quietFleet()),
|
||||
new Gate("2-leadBusy", IDLE_PAST, 0, AgentStatus.WORKING, false, true, quietFleet()),
|
||||
new Gate("3-justBecameIdle", null, 0, AgentStatus.IDLE, false, true, quietFleet()),
|
||||
new Gate("4-withinQuietPeriod", IDLE_WITHIN, 0, AgentStatus.IDLE, false, true, quietFleet()),
|
||||
new Gate("5-noLeadKnown", IDLE_PAST, 0, AgentStatus.IDLE, false, false, pendingFleet()),
|
||||
new Gate("6-pending", IDLE_PAST, 3, AgentStatus.IDLE, false, true, pendingFleet()),
|
||||
new Gate("7-quietNotExhausted", IDLE_PAST, 0, AgentStatus.IDLE, false, true, quietFleet()),
|
||||
new Gate("8-quietExhausted", IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet()));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aOffFlagIgnoresContextAcrossAllEightGates() {
|
||||
// Property A: with contextHighNudge off, decide() must not depend on the context reading at
|
||||
// all — not just "usually agrees", but identical Action/idleSinceNanos/quietCount whatever
|
||||
// context state is passed, and contextNotified must stay false throughout. This is the proof
|
||||
// that fleetd #609 is opt-in: a daemon upgraded to carry this code, but never configuring
|
||||
// `contextHighNudge: true`, behaves exactly as it did before this ticket for every one of the
|
||||
// 8 gates the class javadoc numbers.
|
||||
LeadHeartbeatLoop offLoop = loop(3, false);
|
||||
for (Gate g : allEightGates()) {
|
||||
var withHigh = offLoop.decide(NOW, g.idleSince(), g.quietCount(), g.status(),
|
||||
g.pushLoopActive(), g.leadKnown(), g.fleet(), LeadContextGauge.State.HIGH, false);
|
||||
var withUnknown = offLoop.decide(NOW, g.idleSince(), g.quietCount(), g.status(),
|
||||
g.pushLoopActive(), g.leadKnown(), g.fleet(), LeadContextGauge.State.UNKNOWN, false);
|
||||
var withOk = offLoop.decide(NOW, g.idleSince(), g.quietCount(), g.status(),
|
||||
g.pushLoopActive(), g.leadKnown(), g.fleet(), LeadContextGauge.State.OK, false);
|
||||
|
||||
assertEquals(withUnknown.action(), withHigh.action(), g.name() + ": action must not depend on context");
|
||||
assertEquals(withUnknown.idleSinceNanos(), withHigh.idleSinceNanos(), g.name());
|
||||
assertEquals(withUnknown.quietCount(), withHigh.quietCount(), g.name());
|
||||
assertEquals(withUnknown.action(), withOk.action(), g.name() + ": nor on an OK reading");
|
||||
assertFalse(withHigh.contextNotified(), g.name() + ": contextNotified must stay false when the flag is off");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void bHighContextFiresEvenPastAnExhaustedQuietCapWithoutSpendingIt() {
|
||||
LeadHeartbeatLoop on = loop(3, true);
|
||||
LeadHeartbeatLoop.Decision d = on.decide(
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.HIGH, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.INJECT, d.action(),
|
||||
"an idle, quiet, HIGH-context lead is exactly the case the exhausted cap must not swallow");
|
||||
assertEquals(3, d.quietCount(), "the context notice is an event notice, not a quiet nudge — it must not "
|
||||
+ "spend or grow the consecutive-quiet budget");
|
||||
assertTrue(d.contextNotified(), "firing the notice sets the latch");
|
||||
}
|
||||
|
||||
@Test
|
||||
void cASecondTickWithTheLatchAlreadySetDoesNotTellTheLeadAgain() {
|
||||
LeadHeartbeatLoop on = loop(3, true);
|
||||
LeadHeartbeatLoop.Decision d = on.decide(
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.HIGH, true);
|
||||
assertEquals(LeadHeartbeatLoop.Action.QUIET_DONE, d.action(),
|
||||
"the lead was already told once about this HIGH stretch — telling it every tick would be nagging, "
|
||||
+ "not a notice");
|
||||
}
|
||||
|
||||
@Test
|
||||
void dFlappingBetweenHighAndUnknownNeverReInjectsWhileLatched() {
|
||||
LeadHeartbeatLoop on = loop(3, true);
|
||||
// The lead was already told once (latch = true from a prior HIGH tick).
|
||||
LeadHeartbeatLoop.Decision afterUnknown = on.decide(
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, true);
|
||||
assertTrue(afterUnknown.contextNotified(),
|
||||
"UNKNOWN means 'I could not look', not 'it got better' — it must not clear the latch");
|
||||
assertEquals(LeadHeartbeatLoop.Action.QUIET_DONE, afterUnknown.action(),
|
||||
"a still-latched, still-quiet tick must not inject just because the reading is UNKNOWN");
|
||||
|
||||
LeadHeartbeatLoop.Decision afterHighAgain = on.decide(
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.HIGH, afterUnknown.contextNotified());
|
||||
assertEquals(LeadHeartbeatLoop.Action.QUIET_DONE, afterHighAgain.action(),
|
||||
"HIGH -> UNKNOWN -> HIGH with the latch already set must never inject again — this is the "
|
||||
+ "flapping case that would otherwise cost a full lead its remaining turns");
|
||||
}
|
||||
|
||||
@Test
|
||||
void eAnOkReadingClearsTheLatchSoALaterHighInjectsAgain() {
|
||||
LeadHeartbeatLoop on = loop(3, true);
|
||||
LeadHeartbeatLoop.Decision afterOk = on.decide(
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.OK, true);
|
||||
assertFalse(afterOk.contextNotified(), "an OK reading clears the latch — the lead's context recovered");
|
||||
|
||||
LeadHeartbeatLoop.Decision afterHighAgain = on.decide(
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.HIGH, afterOk.contextNotified());
|
||||
assertEquals(LeadHeartbeatLoop.Action.INJECT, afterHighAgain.action(),
|
||||
"the cleared latch lets a genuinely new HIGH stretch notify again");
|
||||
}
|
||||
|
||||
@Test
|
||||
void fStandingDownDoesNotBurnTheOneContextNotice() {
|
||||
LeadHeartbeatLoop on = loop(3, true);
|
||||
LeadHeartbeatLoop.Decision d = on.decide(
|
||||
NOW, IDLE_PAST, 2, AgentStatus.IDLE, true, true, quietFleet(),
|
||||
LeadContextGauge.State.HIGH, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.STAND_DOWN, d.action(), "ReplyPushLoop is active — stand aside");
|
||||
assertFalse(d.contextNotified(),
|
||||
"standing down must not spend the one notice this HIGH stretch gets — the latch stays clear so a "
|
||||
+ "later tick can still fire it");
|
||||
}
|
||||
|
||||
@Test
|
||||
void gAWorkingLeadWithHighContextIsStillLeadBusy() {
|
||||
LeadHeartbeatLoop on = loop(3, true);
|
||||
LeadHeartbeatLoop.Decision d = on.decide(
|
||||
NOW, IDLE_PAST, 0, AgentStatus.WORKING, false, true, quietFleet(),
|
||||
LeadContextGauge.State.HIGH, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.LEAD_BUSY, d.action(), "the status gate wins over the context notice");
|
||||
}
|
||||
|
||||
@Test
|
||||
void hAPendingDrivenInjectWhileHighSetsTheLatchTooSoTheLeadIsNotToldTwice() {
|
||||
LeadHeartbeatLoop on = loop(3, true);
|
||||
LeadHeartbeatLoop.Decision d = on.decide(
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, pendingFleet(),
|
||||
LeadContextGauge.State.HIGH, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.INJECT, d.action());
|
||||
assertTrue(d.contextNotified(),
|
||||
"the pending-driven nudge carries the context notice too (see contextNotice), so it must set the "
|
||||
+ "latch — otherwise the lead could be told twice by two different routes");
|
||||
}
|
||||
|
||||
// ── contextNotice text builder ─────────────────────────────────────────────────────────────
|
||||
|
||||
@Test
|
||||
void contextNoticeIsEmptyWhenDisabled() {
|
||||
var reading = new LeadContextGauge.Reading(LeadContextGauge.State.HIGH, 260_771L, 1);
|
||||
assertEquals("", LeadHeartbeatLoop.contextNotice(false, reading));
|
||||
}
|
||||
|
||||
@Test
|
||||
void contextNoticeIsEmptyWhenStateIsOk() {
|
||||
var reading = new LeadContextGauge.Reading(LeadContextGauge.State.OK, 1_000L, 0);
|
||||
assertEquals("", LeadHeartbeatLoop.contextNotice(true, reading));
|
||||
}
|
||||
|
||||
@Test
|
||||
void contextNoticeIsEmptyWhenStateIsUnknown() {
|
||||
assertEquals("", LeadHeartbeatLoop.contextNotice(true, LeadContextGauge.Reading.unknown()));
|
||||
}
|
||||
|
||||
@Test
|
||||
void contextNoticeNamesTheTokenCountWhenHighAndEnabled() {
|
||||
var reading = new LeadContextGauge.Reading(LeadContextGauge.State.HIGH, 260_771L, 1);
|
||||
String notice = LeadHeartbeatLoop.contextNotice(true, reading);
|
||||
assertFalse(notice.isEmpty());
|
||||
assertTrue(notice.contains("260771"), notice);
|
||||
assertTrue(notice.contains("1 compaction"), notice);
|
||||
assertTrue(notice.contains("fleet_handover"), notice);
|
||||
assertTrue(notice.contains("operator"), notice);
|
||||
}
|
||||
|
||||
@Test
|
||||
void contextNoticeOmitsTheTokenClauseRatherThanPrintingNull() {
|
||||
var reading = new LeadContextGauge.Reading(LeadContextGauge.State.HIGH, null, 2);
|
||||
String notice = LeadHeartbeatLoop.contextNotice(true, reading);
|
||||
assertFalse(notice.isEmpty());
|
||||
assertFalse(notice.toLowerCase().contains("null"), notice);
|
||||
assertTrue(notice.contains("2 compactions"), notice);
|
||||
}
|
||||
|
||||
// ── fleetd #609 review: the latch must mean "the notice reached the pane" ────────────────────
|
||||
//
|
||||
// These four drive LeadHeartbeatLoop.tick() directly (package-private, same reasoning as
|
||||
// ReplyPushLoop#tick(String) being directly testable) against a real AgentControl wrapping a
|
||||
// FailableHerdrClient, so the send path (agents.send -> herdr -> possible throw) is exercised
|
||||
// for real rather than assumed from decide()'s Decision alone.
|
||||
|
||||
private static final String LEAD = "term_lead";
|
||||
private static final String WORKER = "term_w1";
|
||||
|
||||
/** A HIGH reading with a fixed token/compaction count, for the four tests below. */
|
||||
private static LeadContextGauge.Reading highReading() {
|
||||
return new LeadContextGauge.Reading(LeadContextGauge.State.HIGH, 260_771L, 1);
|
||||
}
|
||||
|
||||
/**
|
||||
* Builds a real {@link LeadHeartbeatLoop} wired to {@code herdr} via a real {@link AgentControl},
|
||||
* a mutable fake clock, and a mutable roster so a test can change fleet state between ticks. The
|
||||
* lead is always reported IDLE by {@code herdr}, so every tick's outcome is governed only by the
|
||||
* idle-window/quiet-cap/context gates under test.
|
||||
*/
|
||||
private static LeadHeartbeatLoop tickableLoop(FailableHerdrClient herdr, AtomicLong now,
|
||||
List<MemberSession>[] rosterBox, InMemoryReplyInbox inbox,
|
||||
int quietNudgeCap, ScheduledExecutorService scheduler) {
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
PrimaryRegistry registry = new PrimaryRegistry(LEAD);
|
||||
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, agents, inbox, scheduler, 5, 100_000);
|
||||
return new LeadHeartbeatLoop(registry, agents, inbox, () -> rosterBox[0], pushLoop, scheduler,
|
||||
now::get, IDLE_AFTER_NANOS, 100_000L, quietNudgeCap, null,
|
||||
new LeadHeartbeatLoop.LeadContextSource(t -> highReading()), true);
|
||||
}
|
||||
|
||||
@Test
|
||||
void iAFailedSendDoesNotConsumeTheNotice() {
|
||||
var herdr = new FailableHerdrClient(LEAD);
|
||||
var now = new AtomicLong(NOW);
|
||||
@SuppressWarnings("unchecked")
|
||||
List<MemberSession>[] rosterBox = new List[]{List.of()}; // quiet: nothing pending
|
||||
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
LeadHeartbeatLoop loop = tickableLoop(herdr, now, rosterBox, inbox, 0, scheduler);
|
||||
|
||||
loop.tick(); // first injectable tick: only opens the idle window (WAIT_IDLE)
|
||||
now.addAndGet(TimeUnit.SECONDS.toNanos(400)); // now clearly past the quiet period
|
||||
|
||||
herdr.throwOnNextSend();
|
||||
loop.tick(); // quiet fleet, quiet cap exhausted (0), context HIGH, latch clear -> INJECT, send throws
|
||||
|
||||
assertEquals(0, herdr.sentTexts().size(), "the failed send must not have recorded any text");
|
||||
|
||||
loop.tick(); // same inputs — the latch must still be clear, so this must INJECT and send again
|
||||
assertEquals(1, herdr.sentTexts().size(),
|
||||
"a retried tick with the latch still clear must attempt the send again");
|
||||
assertTrue(herdr.sentTexts().get(0).contains("Your own context is nearly full"),
|
||||
"the retried, successful send must carry the notice: " + herdr.sentTexts().get(0));
|
||||
}
|
||||
|
||||
@Test
|
||||
void jASuccessfulSendDoesConsumeIt() {
|
||||
var herdr = new FailableHerdrClient(LEAD);
|
||||
var now = new AtomicLong(NOW);
|
||||
@SuppressWarnings("unchecked")
|
||||
List<MemberSession>[] rosterBox = new List[]{List.of()};
|
||||
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
LeadHeartbeatLoop loop = tickableLoop(herdr, now, rosterBox, inbox, 0, scheduler);
|
||||
|
||||
loop.tick(); // opens the idle window
|
||||
now.addAndGet(TimeUnit.SECONDS.toNanos(400));
|
||||
|
||||
loop.tick(); // INJECT, send succeeds -> latch set
|
||||
assertEquals(1, herdr.sentTexts().size());
|
||||
assertTrue(herdr.sentTexts().get(0).contains("Your own context is nearly full"), herdr.sentTexts().get(0));
|
||||
|
||||
loop.tick(); // same inputs — the lead was already told this stretch
|
||||
assertEquals(1, herdr.sentTexts().size(),
|
||||
"the next tick with the same inputs must not send a second notice");
|
||||
}
|
||||
|
||||
@Test
|
||||
void kAPendingDrivenInjectWithTheLatchAlreadySetSendsNoNotice() {
|
||||
var herdr = new FailableHerdrClient(LEAD);
|
||||
var now = new AtomicLong(NOW);
|
||||
@SuppressWarnings("unchecked")
|
||||
List<MemberSession>[] rosterBox = new List[]{List.of()};
|
||||
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
LeadHeartbeatLoop loop = tickableLoop(herdr, now, rosterBox, inbox, 0, scheduler);
|
||||
|
||||
loop.tick();
|
||||
now.addAndGet(TimeUnit.SECONDS.toNanos(400));
|
||||
loop.tick(); // latches the notice (quiet, HIGH, cap exhausted -> the forced context INJECT)
|
||||
assertEquals(1, herdr.sentTexts().size());
|
||||
assertTrue(herdr.sentTexts().get(0).contains("Your own context is nearly full"));
|
||||
|
||||
// Now make the fleet have real pending state, so the NEXT INJECT is pending-driven, not the
|
||||
// forced-context route — with the latch already set from the tick above.
|
||||
inbox.own(WORKER);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
rosterBox[0] = List.of(new MemberSession("p1", WORKER, "prof", MemberRole.DEV, "/cwd", null,
|
||||
0, 0, 0, MemberSession.State.READY, null, null));
|
||||
|
||||
loop.tick();
|
||||
|
||||
assertEquals(2, herdr.sentTexts().size(), "the pending-driven tick must still send a nudge");
|
||||
assertFalse(herdr.sentTexts().get(1).contains("Your own context is nearly full"),
|
||||
"a pending-driven INJECT while the latch is already set must carry no context notice: "
|
||||
+ herdr.sentTexts().get(1));
|
||||
}
|
||||
|
||||
@Test
|
||||
void lTheNoticeAppearsExactlyOnceAcrossThreeDifferentlyDrivenInjects() {
|
||||
var herdr = new FailableHerdrClient(LEAD);
|
||||
var now = new AtomicLong(NOW);
|
||||
@SuppressWarnings("unchecked")
|
||||
List<MemberSession>[] rosterBox = new List[]{List.of()};
|
||||
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
// quietNudgeCap=1 so a still-not-exhausted quiet nudge is available as the "forced" route below,
|
||||
// distinct from both the pending-driven route and the exhausted-cap forced-context route.
|
||||
LeadHeartbeatLoop loop = tickableLoop(herdr, now, rosterBox, inbox, 1, scheduler);
|
||||
|
||||
loop.tick(); // opens the idle window
|
||||
now.addAndGet(TimeUnit.SECONDS.toNanos(400));
|
||||
|
||||
// 1) pending-driven INJECT: real fleet state present. Sets the latch and carries the notice.
|
||||
inbox.own(WORKER);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
rosterBox[0] = List.of(new MemberSession("p1", WORKER, "prof", MemberRole.DEV, "/cwd", null,
|
||||
0, 0, 0, MemberSession.State.READY, null, null));
|
||||
loop.tick();
|
||||
assertEquals(1, herdr.sentTexts().size());
|
||||
assertTrue(herdr.sentTexts().get(0).contains("Your own context is nearly full"), herdr.sentTexts().get(0));
|
||||
|
||||
// 2) "forced" INJECT: nothing pending, but the quiet cap (1) is not yet exhausted, so decide()
|
||||
// nudges anyway. The latch is already set, so no notice.
|
||||
inbox.ack(WORKER, "m1");
|
||||
rosterBox[0] = List.of();
|
||||
loop.tick();
|
||||
assertEquals(2, herdr.sentTexts().size(), "the quiet-cap-not-yet-exhausted nudge must still fire");
|
||||
assertFalse(herdr.sentTexts().get(1).contains("Your own context is nearly full"), herdr.sentTexts().get(1));
|
||||
|
||||
// 3) pending-driven INJECT again. Still latched, still no notice.
|
||||
inbox.publish(WORKER, "m2", "hello again");
|
||||
rosterBox[0] = List.of(new MemberSession("p1", WORKER, "prof", MemberRole.DEV, "/cwd", null,
|
||||
0, 0, 0, MemberSession.State.READY, null, null));
|
||||
loop.tick();
|
||||
assertEquals(3, herdr.sentTexts().size());
|
||||
assertFalse(herdr.sentTexts().get(2).contains("Your own context is nearly full"), herdr.sentTexts().get(2));
|
||||
|
||||
long noticeCount = herdr.sentTexts().stream()
|
||||
.filter(t -> t.contains("Your own context is nearly full")).count();
|
||||
assertEquals(1, noticeCount,
|
||||
"the notice text must appear exactly once across all three sends: " + herdr.sentTexts());
|
||||
}
|
||||
|
||||
/**
|
||||
* Fake herdr client for the four tests above: always reports {@code lead} as IDLE, records the
|
||||
* {@code text} of every {@code agent.prompt} call, and can be told to throw on the very next
|
||||
* {@code agent.prompt} call — standing in for one transient herdr send failure.
|
||||
*/
|
||||
private static final class FailableHerdrClient implements HerdrClient {
|
||||
private static final ObjectMapper MAPPER = new ObjectMapper();
|
||||
private final String lead;
|
||||
private final List<String> sentTexts = new ArrayList<>();
|
||||
private boolean throwOnNextSend = false;
|
||||
|
||||
FailableHerdrClient(String lead) {
|
||||
this.lead = lead;
|
||||
}
|
||||
|
||||
void throwOnNextSend() {
|
||||
throwOnNextSend = true;
|
||||
}
|
||||
|
||||
List<String> sentTexts() {
|
||||
return List.copyOf(sentTexts);
|
||||
}
|
||||
|
||||
@Override
|
||||
@SuppressWarnings("unchecked")
|
||||
public JsonNode call(String method, Object params) {
|
||||
if ("agent.get".equals(method)) {
|
||||
return MAPPER.createObjectNode()
|
||||
.set("agent", MAPPER.createObjectNode()
|
||||
.put("terminal_id", lead)
|
||||
.put("agent_status", "idle"));
|
||||
}
|
||||
if ("agent.prompt".equals(method)) {
|
||||
if (throwOnNextSend) {
|
||||
throwOnNextSend = false;
|
||||
throw new RuntimeException("simulated transient herdr send failure");
|
||||
}
|
||||
Map<String, Object> p = params instanceof Map ? (Map<String, Object>) params : Map.of();
|
||||
sentTexts.add(String.valueOf(p.get("text")));
|
||||
}
|
||||
return MAPPER.createObjectNode();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,168 @@
|
||||
package dev.ltms.fleet.msg;
|
||||
|
||||
import java.util.ArrayDeque;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collection;
|
||||
import java.util.Deque;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.Callable;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.Future;
|
||||
import java.util.concurrent.RejectedExecutionException;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.ScheduledFuture;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
/**
|
||||
* A {@link ScheduledExecutorService} for tests that never runs a task on a timer — it records what
|
||||
* {@link ReplyPushLoop} schedules and only ever runs it when the test itself calls
|
||||
* {@link #runDueTasks()} (fleetd #608).
|
||||
*
|
||||
* <p>The test this exists for ({@code anAlreadyCollectedTicketProducesNoNudge}) used to wire
|
||||
* {@link ReplyPushLoop} to a real {@code Executors.newSingleThreadScheduledExecutor()} and bet a
|
||||
* 300ms backoff was "wide" enough that the test's own work (collecting the ticket) always won the
|
||||
* race against the scheduler's own timer firing the next tick. On an idle machine that held; under
|
||||
* a loaded full-suite run it did not, and the test went red on perfectly correct code. Replacing
|
||||
* the timer with this fake removes the race rather than widening it: nothing here ever runs on a
|
||||
* schedule of its own, so no backoff value — however small — can make the tick fire before the test
|
||||
* is ready for it.
|
||||
*
|
||||
* <p><strong>Deliberately narrow.</strong> {@link ReplyPushLoop} calls exactly two methods on its
|
||||
* {@link ScheduledExecutorService} field — {@link #schedule(Runnable, long, TimeUnit)} (every tick,
|
||||
* including the very first one {@code startOrCoalesce} kicks off) and {@link #shutdownNow()} (on
|
||||
* {@code ReplyPushLoop.stop()}/{@code close()}) — verified by reading every {@code scheduler.} call
|
||||
* site in that class. Every other {@link ScheduledExecutorService} method throws
|
||||
* {@link UnsupportedOperationException} rather than silently doing the wrong thing, so a future
|
||||
* change to {@code ReplyPushLoop} that starts calling one of them fails this fake loudly, on the
|
||||
* very first test that exercises it, instead of being quietly mishandled.
|
||||
*/
|
||||
final class ManualScheduler implements ScheduledExecutorService {
|
||||
|
||||
private final Deque<Runnable> pending = new ArrayDeque<>();
|
||||
private boolean shutdown = false;
|
||||
|
||||
/**
|
||||
* Run every task pending as of the START of this call — a snapshot taken before anything runs.
|
||||
* {@link ReplyPushLoop#tick} always reschedules its own next tick before returning (see
|
||||
* {@code scheduleNext} in its {@code INJECT}/{@code WAIT_BUSY} branches), so without the
|
||||
* snapshot a single call here would recurse forever. Taking it up front means one call is
|
||||
* exactly one tick, deterministically, no matter what that tick itself goes on to schedule.
|
||||
*
|
||||
* @return how many tasks actually ran
|
||||
*/
|
||||
synchronized int runDueTasks() {
|
||||
List<Runnable> due = new ArrayList<>(pending);
|
||||
pending.clear();
|
||||
for (Runnable task : due) {
|
||||
task.run();
|
||||
}
|
||||
return due.size();
|
||||
}
|
||||
|
||||
/** How many tasks are currently queued, without running any of them. */
|
||||
synchronized int pendingCount() {
|
||||
return pending.size();
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized ScheduledFuture<?> schedule(Runnable command, long delay, TimeUnit unit) {
|
||||
if (shutdown) {
|
||||
throw new RejectedExecutionException("ManualScheduler is shut down");
|
||||
}
|
||||
pending.add(command);
|
||||
return null; // ReplyPushLoop discards the return value of every schedule() call it makes.
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized List<Runnable> shutdownNow() {
|
||||
shutdown = true;
|
||||
List<Runnable> left = new ArrayList<>(pending);
|
||||
pending.clear();
|
||||
return left;
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized boolean isShutdown() {
|
||||
return shutdown;
|
||||
}
|
||||
|
||||
// --- everything below: not called by ReplyPushLoop, and not supported by this fake (see the
|
||||
// class javadoc's "deliberately narrow" note) -------------------------------------------------
|
||||
|
||||
@Override
|
||||
public ScheduledFuture<?> scheduleAtFixedRate(Runnable command, long initialDelay, long period, TimeUnit unit) {
|
||||
throw unsupported();
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledFuture<?> scheduleWithFixedDelay(Runnable command, long initialDelay, long delay, TimeUnit unit) {
|
||||
throw unsupported();
|
||||
}
|
||||
|
||||
@Override
|
||||
public <V> ScheduledFuture<V> schedule(Callable<V> callable, long delay, TimeUnit unit) {
|
||||
throw unsupported();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
throw unsupported();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isTerminated() {
|
||||
throw unsupported();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean awaitTermination(long timeout, TimeUnit unit) {
|
||||
throw unsupported();
|
||||
}
|
||||
|
||||
@Override
|
||||
public <T> Future<T> submit(Callable<T> task) {
|
||||
throw unsupported();
|
||||
}
|
||||
|
||||
@Override
|
||||
public <T> Future<T> submit(Runnable task, T result) {
|
||||
throw unsupported();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Future<?> submit(Runnable task) {
|
||||
throw unsupported();
|
||||
}
|
||||
|
||||
@Override
|
||||
public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks) {
|
||||
throw unsupported();
|
||||
}
|
||||
|
||||
@Override
|
||||
public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) {
|
||||
throw unsupported();
|
||||
}
|
||||
|
||||
@Override
|
||||
public <T> T invokeAny(Collection<? extends Callable<T>> tasks) throws ExecutionException {
|
||||
throw unsupported();
|
||||
}
|
||||
|
||||
@Override
|
||||
public <T> T invokeAny(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit)
|
||||
throws ExecutionException {
|
||||
throw unsupported();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void execute(Runnable command) {
|
||||
throw unsupported();
|
||||
}
|
||||
|
||||
private static UnsupportedOperationException unsupported() {
|
||||
return new UnsupportedOperationException(
|
||||
"ManualScheduler only supports schedule(Runnable, long, TimeUnit), shutdownNow() and isShutdown() "
|
||||
+ "— see the class javadoc");
|
||||
}
|
||||
}
|
||||
@@ -1841,6 +1841,36 @@ class MessageServiceTest {
|
||||
return new PushWiring(service, leadHerdr, scheduler);
|
||||
}
|
||||
|
||||
/**
|
||||
* A {@link MessageService} wired to a real {@link ReplyPushLoop}, like {@link #wireWithPushLoop},
|
||||
* but backed by {@link ManualScheduler} instead of a real timer (fleetd #608): its tick never
|
||||
* fires on its own — a test drives it explicitly via {@link ManualScheduler#runDueTasks()}.
|
||||
*/
|
||||
private record ManualPushWiring(MessageService service, FakeHerdr leadHerdr, ManualScheduler scheduler)
|
||||
implements AutoCloseable {
|
||||
@Override
|
||||
public void close() {
|
||||
scheduler.shutdownNow();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #wireWithPushLoop}, but the schedule never runs on its own. {@code backoffMs} is
|
||||
* still threaded through to {@link ReplyPushLoop}'s constructor (it takes one), but nothing here
|
||||
* ever waits it out, so its value cannot affect anything a test built on this observes — see
|
||||
* {@code anAlreadyCollectedTicketProducesNoNudge}, which sets it to 1 to prove exactly that.
|
||||
*/
|
||||
private ManualPushWiring wireWithManualScheduler(int maxReminders, long backoffMs) {
|
||||
PrimaryRegistry registry = new PrimaryRegistry(null);
|
||||
registry.recordDelegation(T, LEAD);
|
||||
FakeHerdr leadHerdr = new FakeHerdr();
|
||||
AgentControl leadAgents = new AgentControl(leadHerdr);
|
||||
ManualScheduler scheduler = new ManualScheduler();
|
||||
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, leadAgents, inbox, scheduler, maxReminders, backoffMs);
|
||||
MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop, null, System::nanoTime);
|
||||
return new ManualPushWiring(service, leadHerdr, scheduler);
|
||||
}
|
||||
|
||||
private void awaitNudge(FakeHerdr leadHerdr) throws InterruptedException {
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!leadHerdr.called("agent.prompt") && System.currentTimeMillis() < deadline) {
|
||||
@@ -1882,16 +1912,46 @@ class MessageServiceTest {
|
||||
|
||||
@Test
|
||||
void anAlreadyCollectedTicketProducesNoNudge() throws Exception {
|
||||
try (var wiring = wireWithPushLoop(1, 300)) { // wide backoff: poll before the first tick fires
|
||||
String ticket = wiring.service().sendAsync(T, "long task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "async result"));
|
||||
// fleetd #608: backoff is 1ms — the most hostile value there is, the tick due immediately —
|
||||
// and the test still must pass, because with ManualScheduler the tick never runs on a timer
|
||||
// at all; it only runs when this test calls runDueTasks() below. The old version bet a 300ms
|
||||
// backoff was "wide" enough that collecting the ticket always won the race against a real
|
||||
// scheduler's own timer; that held on an idle machine and failed under a loaded full-suite
|
||||
// run — a false red on correct code, since nothing here forced that ordering, it only made
|
||||
// it likely.
|
||||
try (var wiring = wireWithManualScheduler(1, 1)) {
|
||||
java.util.concurrent.CountDownLatch terminalReached = new java.util.concurrent.CountDownLatch(1);
|
||||
// finishAsyncTask's task.future.complete(...) runs every whenComplete registered on that
|
||||
// same future — including sendAsync's own hook that calls ReplyPushLoop.onTicketTerminal
|
||||
// — synchronously, before complete() returns (see finishAsyncTask's javadoc and fleetd
|
||||
// #399). This test-only hook fires right after that complete() call, on the very same
|
||||
// (async executor) thread, so waiting for it guarantees onTicketTerminal has already run
|
||||
// and the ticket is really sitting in the push loop's pending set — unlike waiting for
|
||||
// Phase.DONE via poll(), which CompletableFuture.complete() can make visible to another
|
||||
// thread before every whenComplete dependent has actually finished running (the same
|
||||
// publish-then-run-dependents gap awaitCompletionStamped exists to close elsewhere).
|
||||
String ticket;
|
||||
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(terminalReached::countDown);
|
||||
try {
|
||||
ticket = wiring.service().sendAsync(T, "long task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "async result"));
|
||||
assertTrue(terminalReached.await(5, TimeUnit.SECONDS),
|
||||
"the ticket never reached its terminal phase");
|
||||
} finally {
|
||||
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(null);
|
||||
}
|
||||
|
||||
// Collect it — the exact action the nudge must never follow.
|
||||
MessageService.TaskView view = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.DONE);
|
||||
assertEquals(MessageService.Phase.DONE, view.phase());
|
||||
|
||||
Thread.sleep(400); // let the scheduled tick run — it must find nothing pending
|
||||
// Now run the pending tick explicitly. onTicketTerminal scheduled exactly one (the ticket
|
||||
// is the only thing this lead has ever had pending); it must find nothing left pending —
|
||||
// the collection above already removed it — and send no nudge.
|
||||
assertEquals(1, wiring.scheduler().runDueTasks(),
|
||||
"expected exactly the one tick onTicketTerminal scheduled");
|
||||
assertFalse(wiring.leadHerdr().called("agent.prompt"),
|
||||
"a ticket the lead already polled must never be nudged");
|
||||
}
|
||||
|
||||
+77
-20
@@ -50,6 +50,14 @@
|
||||
# absence is the real signal, because a drain that dies on its first session prints nothing
|
||||
# else either. Warns loudly; never fails the redeploy, because by the time this is detectable
|
||||
# the new daemon is already up and healthy.
|
||||
# 10. fleetd #603 — the same shape as trap 3 above, through a different door: the step that waited
|
||||
# for the NEW process to appear gave it its own short, fixed 10s budget, then hard-`die`d,
|
||||
# while the health check right after it waits a full $HEALTH_WAIT (60s) for the same daemon to
|
||||
# answer. Under launchd, `launchctl load` returns as soon as launchd accepts the job, before the
|
||||
# java process exists, and on a slow host that took longer than 10s — so the script died with
|
||||
# "no process appeared" on a deploy that had fully succeeded. The pid poll now shares
|
||||
# $HEALTH_WAIT instead of a separate, shorter budget, and a miss there falls through to the
|
||||
# health check (the truer signal: is it actually answering?) instead of killing the run.
|
||||
#
|
||||
# Usage:
|
||||
# scripts/redeploy-fleetd.sh # build, confirm, restart, verify
|
||||
@@ -79,7 +87,9 @@ OUT="$MODULE/fleetd.out"
|
||||
PATTERN='target/fleetd.jar'
|
||||
HEALTH='http://127.0.0.1:8765/healthz'
|
||||
STOP_WAIT=30 # seconds to wait for a clean exit before reporting failure
|
||||
HEALTH_WAIT=60 # seconds to wait for /healthz to answer after start
|
||||
HEALTH_WAIT=60 # seconds to wait for /healthz to answer after start — fleetd #603: also the pid-
|
||||
# poll budget below (wait_for_new_pid/await_daemon_started), so the two checks
|
||||
# share one named budget instead of the pid poll holding its own shorter one
|
||||
|
||||
# CB-594: the launchd agent this script must not fight with (see trap 6 above).
|
||||
LAUNCHD_LABEL='dev.ltms.fleetd'
|
||||
@@ -1067,6 +1077,67 @@ $(tail -30 "$out_file" 2>/dev/null)"
|
||||
fi
|
||||
}
|
||||
|
||||
# fleetd #603 — the dual of wait_for_daemon_exit above: waits for a pid to APPEAR instead of
|
||||
# disappear. Used to be a bare `for _ in $(seq 10)` sitting directly in the main flow, with its own
|
||||
# short, fixed budget that had nothing to do with $HEALTH_WAIT (60s) — the budget the health check
|
||||
# right after it gets for the very same daemon. Under launchd, `launchctl load` returns as soon as
|
||||
# launchd accepts the job, before the java process exists, and on a slow host that took longer than
|
||||
# 10s — so the script died with "no process appeared" on a deploy that had fully succeeded (the
|
||||
# operator confirmed pid, healthz, and a fresh log line all present, by hand, right afterwards).
|
||||
# Sharing $HEALTH_WAIT here removes the extra, shorter magic number without inventing a new one.
|
||||
wait_for_new_pid() {
|
||||
local timeout="$1" _i
|
||||
for _i in $(seq "$timeout"); do
|
||||
[ -n "$(running_pid)" ] && return 0
|
||||
sleep 1
|
||||
done
|
||||
[ -n "$(running_pid)" ]
|
||||
}
|
||||
|
||||
# fleetd #603 — the pid-appeared check and the healthz check, folded into one decision. Same shape
|
||||
# as swap_if_built/refuse_drain_gate/report_shutdown_drain above (#521/#528/#512): the main flow
|
||||
# calls this ONE function unconditionally, so there is no bare guard left for a future edit to
|
||||
# invert independently of it. That matters more here than for most of those: a plain grep of this
|
||||
# script's source cannot tell "a pid miss falls through to the health check" from "a pid miss still
|
||||
# dies" apart, because both read as the same two lines of text with only the runtime branch
|
||||
# changed — a source-text test could pass on either behavior. A real behavioural test on this
|
||||
# function is the only thing that can actually tell them apart, which is why one exists below.
|
||||
#
|
||||
# wait_for_new_pid's miss is no longer fatal by itself: it falls through to the health check, which
|
||||
# is direct proof the new daemon is up (/healthz answers 200) rather than a proxy for it (a process
|
||||
# merely existing under a name running_pid() recognises). A genuine failure still dies here: it
|
||||
# misses the pid poll AND the health poll, and report_health's own die() still prints the tail of
|
||||
# $out_file, exactly as before this fix.
|
||||
#
|
||||
# Sets NEW_PID (global — the caller's "pid N, jar ..." result line reads it afterwards) and
|
||||
# HEALTH_BODY/HEALTH_CODE (globals, the same reason report_health already needs them handed back).
|
||||
# Only ever called as a bare statement in the main flow below, never from inside a `$( )`: a die()
|
||||
# reached from inside a command substitution only kills that subshell, not the whole script, which
|
||||
# would silently turn a genuine failure back into a false "succeeded" exit (see poll_health_body's
|
||||
# own `|| true` idiom for the same hazard from the other direction).
|
||||
await_daemon_started() {
|
||||
local health_wait="$1" old_pid="$2" health_url="$3" out_file="$4"
|
||||
NEW_PID=""
|
||||
if wait_for_new_pid "$health_wait"; then
|
||||
NEW_PID="$(running_pid)"
|
||||
[ "$NEW_PID" != "${old_pid:-}" ] || die "pid unchanged ($NEW_PID) — the old daemon never died"
|
||||
ok "started, pid $NEW_PID"
|
||||
else
|
||||
warn "no process matched $PATTERN within ${health_wait}s of starting — falling through to the health check, which is the more truthful signal"
|
||||
fi
|
||||
|
||||
HEALTH_BODY="$(poll_health_body "$health_url" "$health_wait")" || true
|
||||
HEALTH_CODE="000"
|
||||
[ -n "$HEALTH_BODY" ] || HEALTH_CODE="$(curl -s -o /dev/null -w '%{http_code}' --max-time 2 "$health_url" 2>/dev/null || echo 000)"
|
||||
report_health "$HEALTH_BODY" "$HEALTH_CODE" "$out_file" "$health_wait"
|
||||
|
||||
# By now /healthz has answered (report_health above would have died otherwise), so the daemon is
|
||||
# confirmed up even if the pid poll never matched it — see running_pid()'s own comment on
|
||||
# under-counting if the launch method ever stops being a plain `java -jar`. Fill NEW_PID in for the
|
||||
# result line rather than leave it blank on an otherwise fully successful redeploy.
|
||||
[ -n "$NEW_PID" ] || NEW_PID="$(running_pid)"
|
||||
}
|
||||
|
||||
# Ticket item 8 — `HAD_OLD_PID=0; [ -n "$OLD_PID" ] && HAD_OLD_PID=1`. Measured safe under `set -e`
|
||||
# at both bash 3.2.57 and 5.x (see the header comment trap 9 discussion in the ticket) — not a `set
|
||||
# -e` hazard, but still an untested computation feeding report_shutdown_drain's own four-way
|
||||
@@ -1285,28 +1356,14 @@ say "start"
|
||||
# there now, so this line is the only thing left in the main flow to get wrong.
|
||||
dispatch_start "$SUPERVISOR_KIND"
|
||||
|
||||
for _ in $(seq 10); do
|
||||
NEW_PID="$(running_pid)"
|
||||
[ -n "$NEW_PID" ] && break
|
||||
sleep 1
|
||||
done
|
||||
[ -n "${NEW_PID:-}" ] || die "no process appeared. Last lines of $OUT:
|
||||
$(tail -20 "$OUT" 2>/dev/null)"
|
||||
[ "$NEW_PID" != "${OLD_PID:-}" ] || die "pid unchanged ($NEW_PID) — the old daemon never died"
|
||||
ok "started, pid $NEW_PID"
|
||||
|
||||
# ------------------------------------------------------------------ verify
|
||||
|
||||
say "verify"
|
||||
|
||||
# fleetd #555: poll_health_body/health_is_up/report_health above. HEALTH_CODE is only ever
|
||||
# consulted by report_health when the body came back empty; `|| true` on both assignments is the
|
||||
# same "an absent/failing command substitution must not kill the script under set -e" idiom the
|
||||
# swap/drain helpers already rely on (see poll_health_body's own comment).
|
||||
HEALTH_BODY="$(poll_health_body "$HEALTH" "$HEALTH_WAIT")" || true
|
||||
HEALTH_CODE="000"
|
||||
[ -n "$HEALTH_BODY" ] || HEALTH_CODE="$(curl -s -o /dev/null -w '%{http_code}' --max-time 2 "$HEALTH" 2>/dev/null || echo 000)"
|
||||
report_health "$HEALTH_BODY" "$HEALTH_CODE" "$OUT" "$HEALTH_WAIT"
|
||||
# fleetd #603: await_daemon_started above folds the pid-appeared check and the healthz check into
|
||||
# one decision — see its own comment for why a bare `if` here, split back into the old two pieces,
|
||||
# would put back an untestable branch this ticket exists to close.
|
||||
await_daemon_started "$HEALTH_WAIT" "$OLD_PID" "$HEALTH" "$OUT"
|
||||
|
||||
# A fresh listening line, strictly after the restart mark. An old daemon that never died would
|
||||
# otherwise let an old line pass for a new one.
|
||||
@@ -1361,7 +1418,7 @@ report_shutdown_drain "$FRESH_LOG" "$HAD_OLD_PID"
|
||||
assert_single_daemon "$(running_pid)"
|
||||
|
||||
say "result"
|
||||
ok "pid $NEW_PID, jar $(jar_id)"
|
||||
ok "pid ${NEW_PID:-unknown}, jar $(jar_id)"
|
||||
if [ "$REDEPLOY_AMQP_CHECK_SKIPPED" -eq 1 ]; then
|
||||
# fleetd #552: the fourth reader of the fresh-log region. Without this branch REDEPLOY_ERROR_COUNT
|
||||
# stays at its untouched 0 (classify_amqp_connection_errors never got a file to read) and
|
||||
|
||||
@@ -770,6 +770,137 @@ test_wait_for_daemon_exit_times_out_if_pid_never_clears() {
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh" # restore the real running_pid/sleep for later tests
|
||||
}
|
||||
|
||||
# fleetd #603 — wait_for_new_pid is the dual of wait_for_daemon_exit above: it must not report
|
||||
# success while running_pid() still answers empty, and must report success the moment a pid
|
||||
# appears. Same counter-file idiom as test_wait_for_daemon_exit_returns_true_once_pid_clears above,
|
||||
# for the same reason (running_pid() runs inside a `$(...)` subshell on every call).
|
||||
test_wait_for_new_pid_returns_true_once_pid_appears() {
|
||||
local counter_file="$TMP/wait-new-pid-calls" final_calls
|
||||
printf '0' > "$counter_file"
|
||||
running_pid() {
|
||||
local n
|
||||
n="$(cat "$counter_file")"
|
||||
n=$((n + 1))
|
||||
printf '%s' "$n" > "$counter_file"
|
||||
if [ "$n" -lt 3 ]; then printf ''; else printf '4242'; fi
|
||||
}
|
||||
sleep() { :; }
|
||||
wait_for_new_pid 10 || fail "wait_for_new_pid did not report success once the pid appeared"
|
||||
final_calls="$(cat "$counter_file")"
|
||||
[ "$final_calls" -ge 3 ] || fail "wait_for_new_pid returned before actually re-checking running_pid"
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh" # restore the real running_pid/sleep for later tests
|
||||
}
|
||||
|
||||
test_wait_for_new_pid_times_out_if_pid_never_appears() {
|
||||
local rc=0
|
||||
running_pid() { printf ''; }
|
||||
sleep() { :; }
|
||||
wait_for_new_pid 3 || rc=$?
|
||||
[ "$rc" -ne 0 ] || fail "wait_for_new_pid reported success while the pid never appeared"
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh" # restore the real running_pid/sleep for later tests
|
||||
}
|
||||
|
||||
# fleetd #603 — await_daemon_started folds the pid-appeared check and the healthz check into one
|
||||
# decision (see its own comment in redeploy-fleetd.sh for why a source-text grep cannot tell the two
|
||||
# possible behaviors apart here). These two tests are the ticket's own acceptance criteria, run
|
||||
# together in this one suite invocation so neither can be satisfied by code that never fails at all:
|
||||
#
|
||||
# 1. a slow start must still succeed — running_pid mimics a process that does not appear until
|
||||
# well after the OLD, buggy 10-second budget, but does appear, and healthz answers.
|
||||
# 2. a genuine failure must still fail, and still print the log tail — running_pid and the health
|
||||
# check both report nothing at all, ever.
|
||||
test_await_daemon_started_slow_pid_then_healthy_succeeds() {
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
stub_die_recorder
|
||||
local counter_file="$TMP/await-slow-pid-calls" output
|
||||
printf '0' > "$counter_file"
|
||||
running_pid() {
|
||||
local n
|
||||
n="$(cat "$counter_file")"
|
||||
n=$((n + 1))
|
||||
printf '%s' "$n" > "$counter_file"
|
||||
# Stays empty well past the old 10-second budget, then appears — the exact shape #603 reports.
|
||||
if [ "$n" -lt 12 ]; then printf ''; else printf '4242'; fi
|
||||
}
|
||||
sleep() { :; }
|
||||
poll_health_body() { printf '{"status":"ok"}'; return 0; }
|
||||
# NOT `output="$(await_daemon_started ...)"`: that would run the call in a subshell, and
|
||||
# NEW_PID — a plain global assignment inside the function, by design (see its own comment) — would
|
||||
# die with that subshell instead of reaching this test's own shell. Redirect to a file instead, the
|
||||
# same hazard await_daemon_started's own comment warns die() itself is subject to.
|
||||
await_daemon_started 60 "" "http://ignored/healthz" "$TMP/await-slow-pid.out" \
|
||||
> "$TMP/await-slow-pid-output.log" 2>&1
|
||||
output="$(cat "$TMP/await-slow-pid-output.log")"
|
||||
[ "$DIED_CALLED" = 0 ] \
|
||||
|| fail "await_daemon_started must not die on a slow-but-real start: $DIED_MESSAGE"
|
||||
assert_equals "4242" "$NEW_PID" "await_daemon_started NEW_PID after a slow-but-real start"
|
||||
printf '%s' "$output" | grep -qF 'started, pid 4242' \
|
||||
|| fail "await_daemon_started did not report the pid once it finally appeared"
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
}
|
||||
|
||||
test_await_daemon_started_never_appears_dies_with_log_tail() {
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
stub_die_recorder
|
||||
local out_file="$TMP/await-never-appears.out"
|
||||
printf 'boot line one\nboot line two\n' > "$out_file"
|
||||
running_pid() { printf ''; }
|
||||
sleep() { :; }
|
||||
poll_health_body() { return 1; }
|
||||
await_daemon_started 2 "" "http://127.0.0.1:1/healthz" "$out_file" > /dev/null 2>&1
|
||||
[ "$DIED_CALLED" = 1 ] \
|
||||
|| fail "await_daemon_started must die when the daemon never appears and never becomes healthy"
|
||||
printf '%s' "$DIED_MESSAGE" | grep -qF 'never answered' \
|
||||
|| fail "await_daemon_started die message does not say healthz never answered"
|
||||
printf '%s' "$DIED_MESSAGE" | grep -qF 'boot line two' \
|
||||
|| fail "await_daemon_started die message does not include the log tail"
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
}
|
||||
|
||||
# fleetd #603 review — the path the fall-through actually exists for, and the one gap a lead
|
||||
# mutation found in the first version of this test file: running_pid() NEVER finds anything (as its
|
||||
# own doc comment says it eventually will, once the daemon stops being launched as a plain
|
||||
# `java -jar` its allowlist recognises), while /healthz answers anyway. Neither of the two tests
|
||||
# above drives this: the slow-pid test has the pid appear, so the `else` branch never runs, and the
|
||||
# never-appears test fails BOTH checks, so it dies either way and cannot tell which branch fired.
|
||||
# This must not die, must warn (so the operator is told the pid could not be identified), and must
|
||||
# leave NEW_PID empty — the honest "could not establish this" answer, never a guessed pid, which is
|
||||
# what the final result line's `${NEW_PID:-unknown}` fallback exists to print truthfully.
|
||||
#
|
||||
# Proof this actually pins the behavior, not just the source text (paste from a real run, not
|
||||
# claimed): reverting the `warn` below back to `die "no process appeared"` (the old fleetd #603
|
||||
# defect, reintroduced) turns this one test red —
|
||||
# FAIL: await_daemon_started must not die when the pid is never found but healthz answers
|
||||
# — and restoring `warn` turns the whole suite green again. Both halves observed, not asserted.
|
||||
test_await_daemon_started_pid_never_found_but_healthy_warns_and_survives() {
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
stub_die_recorder
|
||||
local out_file="$TMP/await-pid-never-found.out" output
|
||||
printf 'boot line\n' > "$out_file"
|
||||
running_pid() { printf ''; }
|
||||
sleep() { :; }
|
||||
poll_health_body() { printf '{"status":"ok"}'; return 0; }
|
||||
# NOT `output="$(await_daemon_started ...)"` — see the slow-pid test above for why that would
|
||||
# drop NEW_PID's assignment in a subshell instead of reaching this test's own shell.
|
||||
await_daemon_started 3 "" "http://ignored/healthz" "$out_file" \
|
||||
> "$TMP/await-pid-never-found-output.log" 2>&1
|
||||
output="$(cat "$TMP/await-pid-never-found-output.log")"
|
||||
[ "$DIED_CALLED" = 0 ] \
|
||||
|| fail "await_daemon_started must not die when the pid is never found but healthz answers: $DIED_MESSAGE"
|
||||
printf '%s' "$output" | grep -qF 'falling through to the health check' \
|
||||
|| fail "await_daemon_started did not warn that the pid could not be identified"
|
||||
assert_equals "" "$NEW_PID" \
|
||||
"await_daemon_started NEW_PID when the pid is never found but healthz answers — must stay empty, never a guessed pid"
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
}
|
||||
|
||||
test_await_daemon_started_call_site_present() {
|
||||
local src="$ROOT/scripts/redeploy-fleetd.sh" call_line
|
||||
call_line="$(grep -Fn 'await_daemon_started "$HEALTH_WAIT" "$OLD_PID" "$HEALTH" "$OUT"' "$src" | head -1 | cut -d: -f1 || true)"
|
||||
[ -n "$call_line" ] \
|
||||
|| fail "could not find the main flow's await_daemon_started call site in redeploy-fleetd.sh"
|
||||
}
|
||||
|
||||
# fleetd #521 — the swap step's guard, at two levels.
|
||||
#
|
||||
# The first two tests call the predicate should_swap() directly. They pin its logic, and that is all
|
||||
@@ -1350,9 +1481,11 @@ test_poll_health_body_returns_nonzero_when_unreachable() {
|
||||
|
||||
test_report_health_call_site_present() {
|
||||
local src="$ROOT/scripts/redeploy-fleetd.sh" call_line
|
||||
call_line="$(grep -Fn 'report_health "$HEALTH_BODY" "$HEALTH_CODE" "$OUT" "$HEALTH_WAIT"' "$src" | head -1 | cut -d: -f1 || true)"
|
||||
# fleetd #603 — moved from a literal main-flow call into await_daemon_started (see its own
|
||||
# comment above for why); this now finds the call inside that function instead.
|
||||
call_line="$(grep -Fn 'report_health "$HEALTH_BODY" "$HEALTH_CODE" "$out_file" "$health_wait"' "$src" | head -1 | cut -d: -f1 || true)"
|
||||
[ -n "$call_line" ] \
|
||||
|| fail "could not find the main flow's report_health call site in redeploy-fleetd.sh"
|
||||
|| fail "could not find await_daemon_started's report_health call site in redeploy-fleetd.sh"
|
||||
}
|
||||
|
||||
# fleetd #555 item 8 — `HAD_OLD_PID=0; [ -n "$OLD_PID" ] && HAD_OLD_PID=1`. Measured safe under
|
||||
@@ -2041,6 +2174,12 @@ test_require_no_build_jar_dies_when_absent
|
||||
test_require_no_build_jar_accepts_present_jar
|
||||
test_wait_for_daemon_exit_returns_true_once_pid_clears
|
||||
test_wait_for_daemon_exit_times_out_if_pid_never_clears
|
||||
test_wait_for_new_pid_returns_true_once_pid_appears
|
||||
test_wait_for_new_pid_times_out_if_pid_never_appears
|
||||
test_await_daemon_started_slow_pid_then_healthy_succeeds
|
||||
test_await_daemon_started_never_appears_dies_with_log_tail
|
||||
test_await_daemon_started_pid_never_found_but_healthy_warns_and_survives
|
||||
test_await_daemon_started_call_site_present
|
||||
test_swap_ordered_after_wait_and_before_start
|
||||
test_drain_gate_abort_message_says_no_no_build
|
||||
test_drain_gate_refusal_build_ran_staged_present
|
||||
|
||||
Reference in New Issue
Block a user