Merge #610: nudge an idle lead to hand over when its own context reads HIGH (fleetd #609)
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 52s
CI / build (push) Failing after 1m51s

Closes fleetd #609. Completes the second half of the context work: #602/#606 could detect a full lead context, and nothing acted on it. LeadHeartbeatLoop now offers a handover when the lead's own gauge reads HIGH.

Never rolls a pane by itself. The nudge is text only; the lead still has to call fleet_handover, and that still needs operatorConfirmed.

Verified by me on a scratch worktree merging 89cb8ff onto fa62e99:
- 1864 tests, 0 failures, 0 errors, 149 surefire reports, mvn exit 0.
- Four mutations, all killed: the two the implementer ran (latch back in applyDecision: 4 failures; call site drops the latch argument: 2 failures), one of my own at the line the logic moved TO (latch set regardless of send outcome: 1 failure), and a control on an untouched line (quiet-cap boundary < to <=: 4 failures). The control is what makes the other kills evidence.
- The new tests assert on herdr.sentTexts() — what actually reached the fake pane — not on source text. That is the right observable for a defect whose essence was "the latch says told, the pane got nothing".

The review blocker from the first round is fixed: the latch used to be committed by applyDecision before injectNudge tried to send, and injectNudge swallows its own RuntimeException. On quietNudgeCap: 0, which is this host's configuration, that was the normal path and not an edge case. The latch is now set only when a notice was actually included and the send returned.

Known and deliberately not blocked: Fleetd.main's own one-line call to leadContextSource is not pinned by a test. That is pre-existing and class-wide, tracked in #612.
This commit was merged in pull request #610.
This commit is contained in:
2026-09-20 12:37:43 +02:00
8 changed files with 968 additions and 49 deletions
+10 -1
View File
@@ -77,16 +77,25 @@ bind:
# a lead turn nobody asked for), so upgrading the daemon must never switch it on for you. Absent # a lead turn nobody asked for), so upgrading the daemon must never switch it on for you. Absent
# block = feature off, exactly as before. # 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 — # 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) # # 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) # 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 # 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) # # 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: # leadHeartbeat:
# idleAfterSeconds: 300 # idleAfterSeconds: 300
# backoffMs: 60000 # backoffMs: 60000
# quietNudgeCap: 3 # quietNudgeCap: 3
# contextHighNudge: false
# Lead rollover (fleetd #480): replace a lead session that has decided it is ready to be replaced, # 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 # 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.ConfigRef;
import dev.ltms.fleet.config.ConfigWatcher; import dev.ltms.fleet.config.ConfigWatcher;
import dev.ltms.fleet.guard.SubscriptionGuard; import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.AgentControl; import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.HerdrClient; import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.HerdrException; import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.herdr.HerdrRouter; import dev.ltms.fleet.herdr.HerdrRouter;
import dev.ltms.fleet.herdr.LeadTabScanner; import dev.ltms.fleet.herdr.LeadTabScanner;
import dev.ltms.fleet.lead.LeadContextGauge;
import dev.ltms.fleet.lead.LeadLauncher; import dev.ltms.fleet.lead.LeadLauncher;
import dev.ltms.fleet.lead.LeadRollover; import dev.ltms.fleet.lead.LeadRollover;
import dev.ltms.fleet.herdr.PaneLocator; import dev.ltms.fleet.herdr.PaneLocator;
@@ -563,10 +565,18 @@ public final class Fleetd {
Thread.ofVirtual().name("bridge-heartbeat-").unstarted(r)); Thread.ofVirtual().name("bridge-heartbeat-").unstarted(r));
if (cfg.leadHeartbeat() != null) { if (cfg.leadHeartbeat() != null) {
var hb = cfg.leadHeartbeat(); 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, heartbeat = new LeadHeartbeatLoop(primaryRegistry, router.leadAgents(), replyInbox, sessions::roster,
pushLoop, heartbeatScheduler, System::nanoTime, pushLoop, heartbeatScheduler, System::nanoTime,
TimeUnit.SECONDS.toNanos(hb.idleAfterSeconds()), hb.backoffMs(), hb.quietNudgeCap(), 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(); heartbeat.start();
} else { } else {
heartbeat = null; heartbeat = null;
@@ -1632,6 +1642,57 @@ public final class Fleetd {
return new FleetMcp.LeadConfigDirSource(leadConfigDirLookup(profiles, 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 * fleetd #248 / fleetd#201 Unit 5: package-private factory for the per-target backend-error
* pattern lookup {@link CompletionResolver} classifies a pane scrape against. Closes over the * pattern lookup {@link CompletionResolver} classifies a pane scrape against. Closes over the
@@ -1329,9 +1329,20 @@ public record FleetConfig(
* each such nudge costs the lead a turn just to read "nothing pending"; * 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 * 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. * 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) @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 { public LeadHeartbeat {
idleAfterSeconds = (idleAfterSeconds == null || idleAfterSeconds <= 0) ? 300 : idleAfterSeconds; idleAfterSeconds = (idleAfterSeconds == null || idleAfterSeconds <= 0) ? 300 : idleAfterSeconds;
backoffMs = (backoffMs == null || backoffMs <= 0) ? 60_000L : backoffMs; backoffMs = (backoffMs == null || backoffMs <= 0) ? 60_000L : backoffMs;
@@ -120,7 +120,13 @@ public final class LeadContextGauge {
* trust this number either way * trust this number either way
*/ */
public record Reading(State state, Long tokens, int compactions) { public record Reading(State state, Long tokens, int compactions) {
static Reading unknown() { /**
* 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); return new Reading(State.UNKNOWN, null, 0);
} }
} }
@@ -2,6 +2,7 @@ package dev.ltms.fleet.msg;
import dev.ltms.fleet.herdr.AgentControl; import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus; import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.lead.LeadContextGauge;
import dev.ltms.fleet.mcp.PrimaryRegistry; import dev.ltms.fleet.mcp.PrimaryRegistry;
import dev.ltms.fleet.metrics.FleetMetrics; import dev.ltms.fleet.metrics.FleetMetrics;
import dev.ltms.fleet.metrics.Metrics; import dev.ltms.fleet.metrics.Metrics;
@@ -13,6 +14,7 @@ import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.function.Function;
import java.util.function.LongSupplier; import java.util.function.LongSupplier;
import java.util.function.Supplier; 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 * loop stands down, so two competing injections never start two turns in the same pane
* (constraint 6).</li> * (constraint 6).</li>
* </ol> * </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 { public final class LeadHeartbeatLoop {
@@ -63,11 +74,15 @@ public final class LeadHeartbeatLoop {
private final long backoffMs; private final long backoffMs;
private final int quietNudgeCap; private final int quietNudgeCap;
private final Metrics metrics; // CB-512 pattern: nullable — no registry in unit tests 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. */ /** When the current idle stretch began (nanos), or {@link #NOT_IDLE}. Single scheduler thread only. */
private long idleSinceNanos = NOT_IDLE; private long idleSinceNanos = NOT_IDLE;
/** Consecutive nudges that found no pending fleet state. Single scheduler thread only. */ /** Consecutive nudges that found no pending fleet state. Single scheduler thread only. */
private int quietCount = 0; 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). */ /** Constructor with an injectable clock and no metric registry (unit tests, or wiring that opts out). */
public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox, public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
@@ -75,7 +90,7 @@ public final class LeadHeartbeatLoop {
ScheduledExecutorService scheduler, LongSupplier clock, ScheduledExecutorService scheduler, LongSupplier clock,
long idleAfterNanos, long backoffMs, int quietNudgeCap) { long idleAfterNanos, long backoffMs, int quietNudgeCap) {
this(primaryRegistry, agents, inbox, roster, pushLoop, scheduler, clock, 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. */ /** 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, Supplier<List<MemberSession>> roster, ReplyPushLoop pushLoop,
ScheduledExecutorService scheduler, LongSupplier clock, ScheduledExecutorService scheduler, LongSupplier clock,
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics) { 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.primaryRegistry = primaryRegistry;
this.agents = agents; this.agents = agents;
this.inbox = inbox; this.inbox = inbox;
@@ -94,6 +123,19 @@ public final class LeadHeartbeatLoop {
this.backoffMs = backoffMs; this.backoffMs = backoffMs;
this.quietNudgeCap = quietNudgeCap; this.quietNudgeCap = quietNudgeCap;
this.metrics = metrics; 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 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 * 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 pushLoopActive whether {@link ReplyPushLoop} is currently nudging some target (constraint 6)
* @param leadKnown whether a lead terminal is known to nudge at all * @param leadKnown whether a lead terminal is known to nudge at all
* @param fleet a snapshot of the pending fleet state (constraint 5) * @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 * @return the action to take and the state to persist
*/ */
Decision decide(long nowNanos, Long idleSinceNanos, int quietCount, AgentStatus status, 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, // 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 // 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 // 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 // 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) { 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 // 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 // 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 // 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. // lead was / may be active, so the next idle stretch must count its own quiet period fresh.
if (status == null || !status.injectable()) { 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) { if (idleSinceNanos == null) {
// The lead just became injectable — record the start of an idle stretch and wait out the // The lead just became injectable — record the start of an idle stretch and wait out the
// debounce quiet period before ever nudging (constraint 3). // 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) { if (nowNanos - idleSinceNanos < idleAfterNanos) {
// Still within the quiet period: the lead that just finished a turn sits momentarily idle // 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. // 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 // Past the quiet period with an injectable lead: it is a genuine candidate for a nudge. Two
// gating facts decide whether and how: // 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 // 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 // to nudge. Keep waiting; the window stays open so discovery re-arms it without a fresh
// quiet period. // quiet period.
return new Decision(Action.WAIT_IDLE, idleSinceNanos, quietCount); return new Decision(Action.WAIT_IDLE, idleSinceNanos, quietCount, latch);
} }
if (fleet.hasPending()) { if (fleet.hasPending()) {
// Real fleet state is waiting — a worker reply or a DONE session. This is new state, so // 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. // it resets the quiet counter (constraint 4) and the lead is nudged to go collect it. The
return new Decision(Action.INJECT, idleSinceNanos, 0); // 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) { if (quietCount < quietNudgeCap) {
// Nothing is pending, but the cap is not exhausted: nudge anyway, telling the lead // 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 // exactly that nothing is waiting so it can choose to stand down rather than hunt
// (constraint 5). Count it toward the consecutive-quiet cap. // (constraint 5). Count it toward the consecutive-quiet cap. This nudge also carries the
return new Decision(Action.INJECT, idleSinceNanos, quietCount + 1); // 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 // 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 // (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. // 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 ---------------------------------------------------------------------------------- // --- loop ----------------------------------------------------------------------------------
@@ -208,57 +276,152 @@ public final class LeadHeartbeatLoop {
scheduler.schedule(this::tick, backoffMs, TimeUnit.MILLISECONDS); scheduler.schedule(this::tick, backoffMs, TimeUnit.MILLISECONDS);
} }
/** One loop tick, every {@link #backoffMs} — the thin scheduler around {@link #decide}. */ /** One loop tick, every {@link #backoffMs} — the thin scheduler around {@link #decide}. Package-private
private void tick() { * (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(); boolean leadKnown = primaryRegistry.primaryTerminal().isPresent();
FleetState fleet = snapshot(inbox, roster); FleetState fleet = snapshot(inbox, roster);
AgentStatus status = AgentStatus.UNKNOWN; AgentStatus status = AgentStatus.UNKNOWN;
LeadContextGauge.Reading reading = LeadContextGauge.Reading.unknown();
if (leadKnown) { if (leadKnown) {
String leadTerminal = primaryRegistry.primaryTerminal().orElseThrow();
try { try {
status = agents.status(primaryRegistry.primaryTerminal().orElseThrow()); status = agents.status(leadTerminal);
} catch (RuntimeException e) { } catch (RuntimeException e) {
// A failed status read degrades to "unknown" — decide() treats that like a busy lead // 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. // 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()); 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(), Decision d = decide(clock.getAsLong(),
idleSinceNanos == NOT_IDLE ? null : idleSinceNanos, idleSinceNanos == NOT_IDLE ? null : idleSinceNanos,
quietCount, status, pushLoop.isActive(), leadKnown, fleet); quietCount, status, pushLoop.isActive(), leadKnown, fleet,
reading.state(), contextNotified);
applyDecision(d); applyDecision(d);
switch (d.action()) { switch (d.action()) {
case INJECT -> injectNudge(fleet); case INJECT -> injectNudge(d, fleet, reading);
case QUIET_DONE -> countNudge("exhausted"); case QUIET_DONE -> {
case WAIT_IDLE, LEAD_BUSY, STAND_DOWN -> { /* nothing to inject, nothing to count */ } countNudge("exhausted");
contextNotified = d.contextNotified();
}
case WAIT_IDLE, LEAD_BUSY, STAND_DOWN -> contextNotified = d.contextNotified();
} }
scheduleNext(); 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) { private void applyDecision(Decision d) {
idleSinceNanos = d.idleSinceNanos() == null ? NOT_IDLE : d.idleSinceNanos(); idleSinceNanos = d.idleSinceNanos() == null ? NOT_IDLE : d.idleSinceNanos();
quietCount = d.quietCount(); 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(); var lead = primaryRegistry.primaryTerminal();
if (lead.isEmpty()) { boolean sent = lead.isPresent() && trySend(lead.get(), fleet.nudgeText() + notice, notice);
return; // the lead disappeared between the decision and the injection // 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
String leadTerminal = lead.get(); // 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 { 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: {})", log.debug("idle-heartbeat: nudge sent to lead {} (quiet nudges so far in this stretch: {})",
leadTerminal, quietCount); 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) { } catch (RuntimeException e) {
log.warn("idle-heartbeat: failed to nudge lead {}: {}", leadTerminal, e.toString()); log.warn("idle-heartbeat: failed to nudge lead {}: {}", leadTerminal, e.toString());
countNudge("failed"); 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. */ /** Schedule the next tick on the scheduler thread pool. */
private void scheduleNext() { private void scheduleNext() {
scheduler.schedule(this::tick, backoffMs, TimeUnit.MILLISECONDS); scheduler.schedule(this::tick, backoffMs, TimeUnit.MILLISECONDS);
@@ -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());
}
}
@@ -1,6 +1,11 @@
package dev.ltms.fleet.msg; 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.AgentStatus;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.lead.LeadContextGauge;
import dev.ltms.fleet.mcp.PrimaryRegistry; import dev.ltms.fleet.mcp.PrimaryRegistry;
import dev.ltms.fleet.peer.MemberRole; import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.session.MemberSession; 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.BeforeEach;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.Map;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.*; import static org.junit.jupiter.api.Assertions.*;
/** /**
* Unit tests for the CB-551 idle-lead heartbeat: the pure {@link LeadHeartbeatLoop#decide} decision * 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 * <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 * clock — no sleeping, no scheduler races. This mirrors how {@code ReplyPushLoopTest} pins the pure
* decision before exercising the loop. * 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 { class LeadHeartbeatLoopTest {
@@ -56,12 +70,18 @@ class LeadHeartbeatLoopTest {
return new LeadHeartbeatLoop.FleetState(2, 1, 3, List.of("term_a", "term_b")); 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) { 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( return new LeadHeartbeatLoop(
new PrimaryRegistry("term_lead"), null /*agents — unused on the decide path*/, new PrimaryRegistry("term_lead"), null /*agents — unused on the decide path*/,
null /*inbox*/, List::of, null /*pushLoop*/, null /*scheduler*/, () -> 0L, 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 ─────────────────────────────────────────────────── // ── (a) a working lead is never injected ───────────────────────────────────────────────────
@@ -69,7 +89,8 @@ class LeadHeartbeatLoopTest {
@Test @Test
void aWorkingLeadIsNeverInjected() { void aWorkingLeadIsNeverInjected() {
LeadHeartbeatLoop.Decision d = loop(3).decide( 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(), assertEquals(LeadHeartbeatLoop.Action.LEAD_BUSY, d.action(),
"a WORKING lead is making progress and must not be touched"); "a WORKING lead is making progress and must not be touched");
assertNull(d.idleSinceNanos(), "a busy lead resets the idle window"); assertNull(d.idleSinceNanos(), "a busy lead resets the idle window");
@@ -80,7 +101,8 @@ class LeadHeartbeatLoopTest {
void anUnreadableStatusIsNeverInjectedEither() { void anUnreadableStatusIsNeverInjectedEither() {
// A failed status read (or a gone agent) must degrade to "do not inject", never hammer the pane. // A failed status read (or a gone agent) must degrade to "do not inject", never hammer the pane.
LeadHeartbeatLoop.Decision d = loop(3).decide( 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(), assertEquals(LeadHeartbeatLoop.Action.LEAD_BUSY, d.action(),
"never inject into a state the loop cannot read"); "never inject into a state the loop cannot read");
} }
@@ -90,7 +112,8 @@ class LeadHeartbeatLoopTest {
@Test @Test
void justBecameIdleStartsTheDebounceWindow() { void justBecameIdleStartsTheDebounceWindow() {
LeadHeartbeatLoop.Decision d = loop(3).decide( 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(), assertEquals(LeadHeartbeatLoop.Action.WAIT_IDLE, d.action(),
"the first injectable tick only records the start of the idle stretch"); "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"); assertEquals(NOW, d.idleSinceNanos(), "the idle window opens at the moment the lead became injectable");
@@ -99,7 +122,8 @@ class LeadHeartbeatLoopTest {
@Test @Test
void idleWithinQuietPeriodIsNotInjected() { void idleWithinQuietPeriodIsNotInjected() {
LeadHeartbeatLoop.Decision d = loop(3).decide( 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(), assertEquals(LeadHeartbeatLoop.Action.WAIT_IDLE, d.action(),
"a lead idle for 10s (< 300s) has just finished a turn — do not re-prompt it"); "a lead idle for 10s (< 300s) has just finished a turn — do not re-prompt it");
} }
@@ -109,7 +133,8 @@ class LeadHeartbeatLoopTest {
@Test @Test
void idlePastQuietPeriodIsInjected() { void idlePastQuietPeriodIsInjected() {
LeadHeartbeatLoop.Decision d = loop(3).decide( 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(), assertEquals(LeadHeartbeatLoop.Action.INJECT, d.action(),
"a lead continuously idle past the quiet period is the reason to nudge"); "a lead continuously idle past the quiet period is the reason to nudge");
} }
@@ -117,14 +142,17 @@ class LeadHeartbeatLoopTest {
@Test @Test
void blockedAndDoneAreInjectableViewsOfIdle() { void blockedAndDoneAreInjectableViewsOfIdle() {
assertEquals(LeadHeartbeatLoop.Action.INJECT, loop(3).decide( 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( 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 @Test
void nothingIsInjectedWhenNoLeadIsKnown() { 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(), assertEquals(LeadHeartbeatLoop.Action.WAIT_IDLE, d.action(),
"with no known lead there is nobody to nudge — keep waiting until one is discovered"); "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"); assertEquals(IDLE_PAST, d.idleSinceNanos(), "the idle window stays open so discovery re-arms it");
@@ -135,7 +163,8 @@ class LeadHeartbeatLoopTest {
@Test @Test
void quietNudgeCapStopsTheLoopWhenNothingIsPending() { void quietNudgeCapStopsTheLoopWhenNothingIsPending() {
LeadHeartbeatLoop.Decision d = loop(3).decide( 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(), assertEquals(LeadHeartbeatLoop.Action.QUIET_DONE, d.action(),
"3 consecutive nothing-pending nudges have already happened — stop nagging an empty fleet"); "3 consecutive nothing-pending nudges have already happened — stop nagging an empty fleet");
} }
@@ -145,7 +174,8 @@ class LeadHeartbeatLoopTest {
@Test @Test
void newPendingStateResetsTheQuietCap() { void newPendingStateResetsTheQuietCap() {
LeadHeartbeatLoop.Decision d = loop(3).decide( 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(), assertEquals(LeadHeartbeatLoop.Action.INJECT, d.action(),
"real state appearing re-arms the loop past an exhausted cap"); "real state appearing re-arms the loop past an exhausted cap");
assertEquals(0, d.quietCount(), "the pending state resets the consecutive-quiet counter"); assertEquals(0, d.quietCount(), "the pending state resets the consecutive-quiet counter");
@@ -156,7 +186,8 @@ class LeadHeartbeatLoopTest {
@Test @Test
void standsDownWhileReplyPushLoopIsActive() { void standsDownWhileReplyPushLoopIsActive() {
LeadHeartbeatLoop.Decision d = loop(3).decide( 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(), assertEquals(LeadHeartbeatLoop.Action.STAND_DOWN, d.action(),
"a second injection would start a competing turn — stand aside instead"); "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"); 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()); assertFalse(fs.hasPending());
assertEquals(0, fs.liveWorkers()); 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() {
}
}
} }