fleetd #609: nudge an idle lead to hand over when its own context reads HIGH
LeadHeartbeatLoop can now append a text-only notice to its nudge when the lead's own LeadContextGauge reading is HIGH and leadHeartbeat.contextHighNudge is on. Fires once per HIGH stretch (a latch, cleared only by a later OK reading; UNKNOWN neither sets nor clears it), never spends the quietNudgeCap budget, and never rolls a pane itself — only the operator can approve a handover. - LeadContextGauge.Reading.unknown() widened to public for LeadContextSource.none() - FleetConfig.LeadHeartbeat gains contextHighNudge (null/false = off, unchanged default) - LeadHeartbeatLoop.decide gains context/contextNotified; Fleetd wires a new leadContextLookup/leadContextSource factory pair (LeadHeartbeatLoop.LeadContextSource) - fleetd.example.yaml documents the new key
This commit is contained in:
@@ -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;
|
||||
@@ -1632,6 +1642,57 @@ public final class Fleetd {
|
||||
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;
|
||||
|
||||
@@ -120,7 +120,13 @@ public final class LeadContextGauge {
|
||||
* trust this number either way
|
||||
*/
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 ----------------------------------------------------------------------------------
|
||||
@@ -213,22 +281,29 @@ public final class LeadHeartbeatLoop {
|
||||
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 INJECT -> injectNudge(fleet, reading);
|
||||
case QUIET_DONE -> countNudge("exhausted");
|
||||
case WAIT_IDLE, LEAD_BUSY, STAND_DOWN -> { /* nothing to inject, nothing to count */ }
|
||||
}
|
||||
@@ -239,26 +314,64 @@ public final class LeadHeartbeatLoop {
|
||||
private void applyDecision(Decision d) {
|
||||
idleSinceNanos = d.idleSinceNanos() == null ? NOT_IDLE : d.idleSinceNanos();
|
||||
quietCount = d.quietCount();
|
||||
contextNotified = d.contextNotified();
|
||||
}
|
||||
|
||||
/** 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. */
|
||||
private void injectNudge(FleetState fleet, LeadContextGauge.Reading reading) {
|
||||
var lead = primaryRegistry.primaryTerminal();
|
||||
if (lead.isEmpty()) {
|
||||
return; // the lead disappeared between the decision and the injection
|
||||
}
|
||||
String leadTerminal = lead.get();
|
||||
String notice = contextNotice(contextHighNudge, reading);
|
||||
String text = fleet.nudgeText() + 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");
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("idle-heartbeat: failed to nudge lead {}: {}", leadTerminal, e.toString());
|
||||
countNudge("failed");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 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) {
|
||||
if (!enabled || 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 {
|
||||
// tokens() can be null even at HIGH is not true today, but do not assume it — leave the
|
||||
// token clause out rather than print "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,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,7 @@
|
||||
package dev.ltms.fleet.msg;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.lead.LeadContextGauge;
|
||||
import dev.ltms.fleet.mcp.PrimaryRegistry;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
@@ -17,11 +18,17 @@ 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 +63,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 +82,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 +94,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 +105,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 +115,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 +126,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 +135,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 +156,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 +167,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 +179,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 +237,178 @@ 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);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user