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 merging89cb8ffontofa62e99: - 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:
@@ -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 ----------------------------------------------------------------------------------
|
||||
@@ -208,57 +276,152 @@ public final class LeadHeartbeatLoop {
|
||||
scheduler.schedule(this::tick, backoffMs, TimeUnit.MILLISECONDS);
|
||||
}
|
||||
|
||||
/** One loop tick, every {@link #backoffMs} — the thin scheduler around {@link #decide}. */
|
||||
private void tick() {
|
||||
/** One loop tick, every {@link #backoffMs} — the thin scheduler around {@link #decide}. Package-private
|
||||
* (mirroring {@link ReplyPushLoop#tick(String)}) so tests can drive it directly with a fake clock and a
|
||||
* fake {@link AgentControl} instead of racing the scheduler thread. */
|
||||
void tick() {
|
||||
boolean leadKnown = primaryRegistry.primaryTerminal().isPresent();
|
||||
FleetState fleet = snapshot(inbox, roster);
|
||||
AgentStatus status = AgentStatus.UNKNOWN;
|
||||
LeadContextGauge.Reading reading = LeadContextGauge.Reading.unknown();
|
||||
if (leadKnown) {
|
||||
String leadTerminal = primaryRegistry.primaryTerminal().orElseThrow();
|
||||
try {
|
||||
status = agents.status(primaryRegistry.primaryTerminal().orElseThrow());
|
||||
status = agents.status(leadTerminal);
|
||||
} catch (RuntimeException e) {
|
||||
// A failed status read degrades to "unknown" — decide() treats that like a busy lead
|
||||
// and never injects into a state it cannot read. Retry on the next backoff.
|
||||
log.debug("idle-heartbeat: status check failed for lead, will retry: {}", e.toString());
|
||||
}
|
||||
// fleetd #609: read the lead's own context regardless of status — decide() still gates on
|
||||
// status first (constraint 2), so this is harmless work on a WORKING lead and lets the
|
||||
// latch state stay accurate for whenever the lead does go idle.
|
||||
reading = contextSource.readingFor().apply(leadTerminal);
|
||||
}
|
||||
|
||||
Decision d = decide(clock.getAsLong(),
|
||||
idleSinceNanos == NOT_IDLE ? null : idleSinceNanos,
|
||||
quietCount, status, pushLoop.isActive(), leadKnown, fleet);
|
||||
quietCount, status, pushLoop.isActive(), leadKnown, fleet,
|
||||
reading.state(), contextNotified);
|
||||
applyDecision(d);
|
||||
switch (d.action()) {
|
||||
case INJECT -> injectNudge(fleet);
|
||||
case QUIET_DONE -> countNudge("exhausted");
|
||||
case WAIT_IDLE, LEAD_BUSY, STAND_DOWN -> { /* nothing to inject, nothing to count */ }
|
||||
case INJECT -> injectNudge(d, fleet, reading);
|
||||
case QUIET_DONE -> {
|
||||
countNudge("exhausted");
|
||||
contextNotified = d.contextNotified();
|
||||
}
|
||||
case WAIT_IDLE, LEAD_BUSY, STAND_DOWN -> contextNotified = d.contextNotified();
|
||||
}
|
||||
scheduleNext();
|
||||
}
|
||||
|
||||
/** Persist the state a decision returned, so the next tick starts from it. */
|
||||
/**
|
||||
* Persist the idle/quiet state a decision returned, so the next tick starts from it.
|
||||
*
|
||||
* <p>fleetd #609 review: the context latch ({@link #contextNotified}) is deliberately <em>not</em>
|
||||
* set here any more. Setting it from the decision unconditionally — before {@link #injectNudge} even
|
||||
* tries to send — is exactly the review's blocker: a decision to notify is not the same fact as "the
|
||||
* notice reached the pane". Every branch of {@link #tick} now assigns {@link #contextNotified} itself,
|
||||
* once it knows whether a send happened and whether it carried the notice (see {@link #injectNudge}).
|
||||
*/
|
||||
private void applyDecision(Decision d) {
|
||||
idleSinceNanos = d.idleSinceNanos() == null ? NOT_IDLE : d.idleSinceNanos();
|
||||
quietCount = d.quietCount();
|
||||
}
|
||||
|
||||
/** Send the nudge to the known lead. */
|
||||
private void injectNudge(FleetState fleet) {
|
||||
/**
|
||||
* Send the nudge to the known lead, with the fleetd #609 context notice appended when it applies, and
|
||||
* persist the context latch based on what actually happened this tick — not merely what {@code d}
|
||||
* chose to attempt.
|
||||
*/
|
||||
private void injectNudge(Decision d, FleetState fleet, LeadContextGauge.Reading reading) {
|
||||
// fleetd #609 review: build the notice from the latch as it stood BEFORE this tick's decision —
|
||||
// d.contextNotified() is the value to persist once delivery is confirmed, not the value the text
|
||||
// itself should be built from. Otherwise a HIGH stretch that is still latched would never see the
|
||||
// notice at all, defeating the very check this fixes.
|
||||
String notice = contextNotice(contextHighNudge, reading, contextNotified);
|
||||
var lead = primaryRegistry.primaryTerminal();
|
||||
if (lead.isEmpty()) {
|
||||
return; // the lead disappeared between the decision and the injection
|
||||
}
|
||||
String leadTerminal = lead.get();
|
||||
boolean sent = lead.isPresent() && trySend(lead.get(), fleet.nudgeText() + notice, notice);
|
||||
// The latch becomes true only when all three hold: decide() chose to notify, a notice was
|
||||
// actually included in the text, and the send reached the pane without throwing. Whenever no
|
||||
// notice was attempted (disabled, not HIGH, or already latched), nothing was promised to the lead
|
||||
// this tick, so apply the decision's own carried-forward value unconditionally — that is how the
|
||||
// OK-only re-arm rule and STAND_DOWN's "don't burn the notice" rule keep working through this path
|
||||
// too. A lead that disappeared between the decision and the send (lead.isEmpty()) is treated the
|
||||
// same as a failed send: nothing reached the pane, so the latch must not be set.
|
||||
contextNotified = notice.isEmpty() ? d.contextNotified() : sent;
|
||||
}
|
||||
|
||||
/**
|
||||
* Attempt one herdr send and count its outcome. Returns whether {@code agents.send} returned without
|
||||
* throwing — the caller ({@link #injectNudge}) needs this to decide whether the fleetd #609 context
|
||||
* latch may be persisted as set.
|
||||
*/
|
||||
private boolean trySend(String leadTerminal, String text, String notice) {
|
||||
try {
|
||||
agents.send(leadTerminal, fleet.nudgeText());
|
||||
agents.send(leadTerminal, text);
|
||||
log.debug("idle-heartbeat: nudge sent to lead {} (quiet nudges so far in this stretch: {})",
|
||||
leadTerminal, quietCount);
|
||||
countNudge("sent");
|
||||
// fleetd #609: a nudge that carries the context notice is counted under its own outcome so
|
||||
// it is visible in /metrics — one count per nudge either way, never two.
|
||||
countNudge(notice.isEmpty() ? "sent" : "sent_context");
|
||||
return true;
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("idle-heartbeat: failed to nudge lead {}: {}", leadTerminal, e.toString());
|
||||
countNudge("failed");
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #609: the text appended to a nudge when the lead's own context is full — {@code ""}
|
||||
* whenever the notice does not apply, so callers can unconditionally append this without an extra
|
||||
* branch. Wording stays plain (CEFR B1) and honest that only the operator approves a roll — this
|
||||
* loop only ever prints text, it never calls {@code fleet_handover} itself.
|
||||
*
|
||||
* @param enabled the {@code leadHeartbeat.contextHighNudge} config flag
|
||||
* @param reading the lead's current {@link LeadContextGauge} reading
|
||||
* @return the notice text (starting with a leading space, to append directly after {@link
|
||||
* FleetState#nudgeText()}), or {@code ""} when disabled or the state is not {@code HIGH}
|
||||
*/
|
||||
static String contextNotice(boolean enabled, LeadContextGauge.Reading reading) {
|
||||
return contextNotice(enabled, reading, false);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #609 review: as {@link #contextNotice(boolean, LeadContextGauge.Reading)}, but also gated on
|
||||
* {@code alreadyNotified} — the context latch as it stood <em>before</em> the current tick's decision.
|
||||
* Without this gate, every pending-driven {@code INJECT} that lands while the context stays {@code
|
||||
* HIGH} would re-append the full notice on top of an already-latched stretch, making the notice's own
|
||||
* closing sentence ("You will not be told again until your context reads ok.") false. {@link
|
||||
* #injectNudge} is the only caller that passes a non-default {@code alreadyNotified}.
|
||||
*
|
||||
* @param alreadyNotified whether the lead has already been told about the current HIGH stretch
|
||||
*/
|
||||
static String contextNotice(boolean enabled, LeadContextGauge.Reading reading, boolean alreadyNotified) {
|
||||
if (!enabled || alreadyNotified || reading.state() != LeadContextGauge.State.HIGH) {
|
||||
return "";
|
||||
}
|
||||
StringBuilder sb = new StringBuilder(" Your own context is nearly full");
|
||||
String compactionWord = reading.compactions() == 1 ? "compaction" : "compactions";
|
||||
if (reading.tokens() != null) {
|
||||
sb.append(": ").append(reading.tokens()).append(" tokens used, ")
|
||||
.append(reading.compactions()).append(' ').append(compactionWord).append(" so far.");
|
||||
} else {
|
||||
// A HIGH reading always carries a non-null token count today: LeadContextGauge only
|
||||
// reaches HIGH by comparing a number against HIGH_THRESHOLD_TOKENS. That invariant
|
||||
// lives in another class and nothing asserts it, so this branch does not rely on it —
|
||||
// it drops the token clause rather than printing "null tokens".
|
||||
sb.append(" (").append(reading.compactions()).append(' ').append(compactionWord)
|
||||
.append(" so far).");
|
||||
}
|
||||
sb.append(" A fresh session would work better. To hand over: call fleet_handover(action=\"open\"), "
|
||||
+ "write the file it names, ask the operator, then call fleet_handover(action=\"confirm\", "
|
||||
+ "token, operatorConfirmed). Only the operator can approve the roll. You will not be told "
|
||||
+ "again until your context reads ok.");
|
||||
return sb.toString();
|
||||
}
|
||||
|
||||
/** Schedule the next tick on the scheduler thread pool. */
|
||||
private void scheduleNext() {
|
||||
scheduler.schedule(this::tick, backoffMs, TimeUnit.MILLISECONDS);
|
||||
|
||||
@@ -0,0 +1,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;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.lead.LeadContextGauge;
|
||||
import dev.ltms.fleet.mcp.PrimaryRegistry;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
@@ -8,20 +13,29 @@ import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
/**
|
||||
* Unit tests for the CB-551 idle-lead heartbeat: the pure {@link LeadHeartbeatLoop#decide} decision
|
||||
* function, the nudge text, and the {@link LeadHeartbeatLoop#snapshot} fleet snapshot.
|
||||
* function, the nudge text, and the {@link LeadHeartbeatLoop#snapshot} fleet snapshot. Also covers
|
||||
* fleetd #609's context-high notice: the {@code context}/{@code contextNotified} parameters {@code
|
||||
* decide} gained, and the standalone {@link LeadHeartbeatLoop#contextNotice} text builder.
|
||||
*
|
||||
* <p>All decision tests call {@code decide} directly with explicit nanoTime values from an injected
|
||||
* clock — no sleeping, no scheduler races. This mirrors how {@code ReplyPushLoopTest} pins the pure
|
||||
* decision before exercising the loop.
|
||||
*
|
||||
* <p>Every pre-#609 test below passes {@code LeadContextGauge.State.UNKNOWN, false} for the two new
|
||||
* {@code decide} parameters — the same as a lead whose context could not be read and has never been
|
||||
* notified — so each one still pins exactly the behaviour it pinned before this ticket.
|
||||
*/
|
||||
class LeadHeartbeatLoopTest {
|
||||
|
||||
@@ -56,12 +70,18 @@ class LeadHeartbeatLoopTest {
|
||||
return new LeadHeartbeatLoop.FleetState(2, 1, 3, List.of("term_a", "term_b"));
|
||||
}
|
||||
|
||||
/** The loop under test; the scheduler is never invoked on the pure decide path. */
|
||||
/** The loop under test, context notice off; the scheduler is never invoked on the pure decide path. */
|
||||
private static LeadHeartbeatLoop loop(int quietCap) {
|
||||
return loop(quietCap, false);
|
||||
}
|
||||
|
||||
/** As above, with the fleetd #609 {@code contextHighNudge} flag set explicitly. */
|
||||
private static LeadHeartbeatLoop loop(int quietCap, boolean contextHighNudge) {
|
||||
return new LeadHeartbeatLoop(
|
||||
new PrimaryRegistry("term_lead"), null /*agents — unused on the decide path*/,
|
||||
null /*inbox*/, List::of, null /*pushLoop*/, null /*scheduler*/, () -> 0L,
|
||||
IDLE_AFTER_NANOS, 1_000L, quietCap);
|
||||
IDLE_AFTER_NANOS, 1_000L, quietCap, null,
|
||||
LeadHeartbeatLoop.LeadContextSource.none(), contextHighNudge);
|
||||
}
|
||||
|
||||
// ── (a) a working lead is never injected ───────────────────────────────────────────────────
|
||||
@@ -69,7 +89,8 @@ class LeadHeartbeatLoopTest {
|
||||
@Test
|
||||
void aWorkingLeadIsNeverInjected() {
|
||||
LeadHeartbeatLoop.Decision d = loop(3).decide(
|
||||
NOW, IDLE_PAST, 0, AgentStatus.WORKING, false, true, quietFleet());
|
||||
NOW, IDLE_PAST, 0, AgentStatus.WORKING, false, true, quietFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.LEAD_BUSY, d.action(),
|
||||
"a WORKING lead is making progress and must not be touched");
|
||||
assertNull(d.idleSinceNanos(), "a busy lead resets the idle window");
|
||||
@@ -80,7 +101,8 @@ class LeadHeartbeatLoopTest {
|
||||
void anUnreadableStatusIsNeverInjectedEither() {
|
||||
// A failed status read (or a gone agent) must degrade to "do not inject", never hammer the pane.
|
||||
LeadHeartbeatLoop.Decision d = loop(3).decide(
|
||||
NOW, IDLE_PAST, 0, AgentStatus.UNKNOWN, false, true, pendingFleet());
|
||||
NOW, IDLE_PAST, 0, AgentStatus.UNKNOWN, false, true, pendingFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.LEAD_BUSY, d.action(),
|
||||
"never inject into a state the loop cannot read");
|
||||
}
|
||||
@@ -90,7 +112,8 @@ class LeadHeartbeatLoopTest {
|
||||
@Test
|
||||
void justBecameIdleStartsTheDebounceWindow() {
|
||||
LeadHeartbeatLoop.Decision d = loop(3).decide(
|
||||
NOW, null, 0, AgentStatus.IDLE, false, true, quietFleet());
|
||||
NOW, null, 0, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.WAIT_IDLE, d.action(),
|
||||
"the first injectable tick only records the start of the idle stretch");
|
||||
assertEquals(NOW, d.idleSinceNanos(), "the idle window opens at the moment the lead became injectable");
|
||||
@@ -99,7 +122,8 @@ class LeadHeartbeatLoopTest {
|
||||
@Test
|
||||
void idleWithinQuietPeriodIsNotInjected() {
|
||||
LeadHeartbeatLoop.Decision d = loop(3).decide(
|
||||
NOW, IDLE_WITHIN, 0, AgentStatus.IDLE, false, true, quietFleet());
|
||||
NOW, IDLE_WITHIN, 0, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.WAIT_IDLE, d.action(),
|
||||
"a lead idle for 10s (< 300s) has just finished a turn — do not re-prompt it");
|
||||
}
|
||||
@@ -109,7 +133,8 @@ class LeadHeartbeatLoopTest {
|
||||
@Test
|
||||
void idlePastQuietPeriodIsInjected() {
|
||||
LeadHeartbeatLoop.Decision d = loop(3).decide(
|
||||
NOW, IDLE_PAST, 0, AgentStatus.IDLE, false, true, quietFleet());
|
||||
NOW, IDLE_PAST, 0, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.INJECT, d.action(),
|
||||
"a lead continuously idle past the quiet period is the reason to nudge");
|
||||
}
|
||||
@@ -117,14 +142,17 @@ class LeadHeartbeatLoopTest {
|
||||
@Test
|
||||
void blockedAndDoneAreInjectableViewsOfIdle() {
|
||||
assertEquals(LeadHeartbeatLoop.Action.INJECT, loop(3).decide(
|
||||
NOW, IDLE_PAST, 0, AgentStatus.BLOCKED, false, true, quietFleet()).action());
|
||||
NOW, IDLE_PAST, 0, AgentStatus.BLOCKED, false, true, quietFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false).action());
|
||||
assertEquals(LeadHeartbeatLoop.Action.INJECT, loop(3).decide(
|
||||
NOW, IDLE_PAST, 0, AgentStatus.DONE, false, true, quietFleet()).action());
|
||||
NOW, IDLE_PAST, 0, AgentStatus.DONE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false).action());
|
||||
}
|
||||
|
||||
@Test
|
||||
void nothingIsInjectedWhenNoLeadIsKnown() {
|
||||
var d = loop(3).decide(NOW, IDLE_PAST, 0, AgentStatus.IDLE, false, false, pendingFleet());
|
||||
var d = loop(3).decide(NOW, IDLE_PAST, 0, AgentStatus.IDLE, false, false, pendingFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.WAIT_IDLE, d.action(),
|
||||
"with no known lead there is nobody to nudge — keep waiting until one is discovered");
|
||||
assertEquals(IDLE_PAST, d.idleSinceNanos(), "the idle window stays open so discovery re-arms it");
|
||||
@@ -135,7 +163,8 @@ class LeadHeartbeatLoopTest {
|
||||
@Test
|
||||
void quietNudgeCapStopsTheLoopWhenNothingIsPending() {
|
||||
LeadHeartbeatLoop.Decision d = loop(3).decide(
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet());
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.QUIET_DONE, d.action(),
|
||||
"3 consecutive nothing-pending nudges have already happened — stop nagging an empty fleet");
|
||||
}
|
||||
@@ -145,7 +174,8 @@ class LeadHeartbeatLoopTest {
|
||||
@Test
|
||||
void newPendingStateResetsTheQuietCap() {
|
||||
LeadHeartbeatLoop.Decision d = loop(3).decide(
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, pendingFleet());
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, pendingFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.INJECT, d.action(),
|
||||
"real state appearing re-arms the loop past an exhausted cap");
|
||||
assertEquals(0, d.quietCount(), "the pending state resets the consecutive-quiet counter");
|
||||
@@ -156,7 +186,8 @@ class LeadHeartbeatLoopTest {
|
||||
@Test
|
||||
void standsDownWhileReplyPushLoopIsActive() {
|
||||
LeadHeartbeatLoop.Decision d = loop(3).decide(
|
||||
NOW, IDLE_PAST, 2, AgentStatus.IDLE, true, true, quietFleet());
|
||||
NOW, IDLE_PAST, 2, AgentStatus.IDLE, true, true, quietFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.STAND_DOWN, d.action(),
|
||||
"a second injection would start a competing turn — stand aside instead");
|
||||
assertEquals(0, d.quietCount(), "the active push is real state, so it re-arms the cap");
|
||||
@@ -213,4 +244,378 @@ class LeadHeartbeatLoopTest {
|
||||
assertFalse(fs.hasPending());
|
||||
assertEquals(0, fs.liveWorkers());
|
||||
}
|
||||
|
||||
// ── fleetd #609: context-high notice ───────────────────────────────────────────────────────
|
||||
|
||||
/** One gate scenario, keyed by name, over which the off-state/context-independence property is checked. */
|
||||
private record Gate(String name, Long idleSince, int quietCount, AgentStatus status,
|
||||
boolean pushLoopActive, boolean leadKnown, LeadHeartbeatLoop.FleetState fleet) {}
|
||||
|
||||
private static List<Gate> allEightGates() {
|
||||
return List.of(
|
||||
new Gate("1-standDown", IDLE_PAST, 2, AgentStatus.IDLE, true, true, quietFleet()),
|
||||
new Gate("2-leadBusy", IDLE_PAST, 0, AgentStatus.WORKING, false, true, quietFleet()),
|
||||
new Gate("3-justBecameIdle", null, 0, AgentStatus.IDLE, false, true, quietFleet()),
|
||||
new Gate("4-withinQuietPeriod", IDLE_WITHIN, 0, AgentStatus.IDLE, false, true, quietFleet()),
|
||||
new Gate("5-noLeadKnown", IDLE_PAST, 0, AgentStatus.IDLE, false, false, pendingFleet()),
|
||||
new Gate("6-pending", IDLE_PAST, 3, AgentStatus.IDLE, false, true, pendingFleet()),
|
||||
new Gate("7-quietNotExhausted", IDLE_PAST, 0, AgentStatus.IDLE, false, true, quietFleet()),
|
||||
new Gate("8-quietExhausted", IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet()));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aOffFlagIgnoresContextAcrossAllEightGates() {
|
||||
// Property A: with contextHighNudge off, decide() must not depend on the context reading at
|
||||
// all — not just "usually agrees", but identical Action/idleSinceNanos/quietCount whatever
|
||||
// context state is passed, and contextNotified must stay false throughout. This is the proof
|
||||
// that fleetd #609 is opt-in: a daemon upgraded to carry this code, but never configuring
|
||||
// `contextHighNudge: true`, behaves exactly as it did before this ticket for every one of the
|
||||
// 8 gates the class javadoc numbers.
|
||||
LeadHeartbeatLoop offLoop = loop(3, false);
|
||||
for (Gate g : allEightGates()) {
|
||||
var withHigh = offLoop.decide(NOW, g.idleSince(), g.quietCount(), g.status(),
|
||||
g.pushLoopActive(), g.leadKnown(), g.fleet(), LeadContextGauge.State.HIGH, false);
|
||||
var withUnknown = offLoop.decide(NOW, g.idleSince(), g.quietCount(), g.status(),
|
||||
g.pushLoopActive(), g.leadKnown(), g.fleet(), LeadContextGauge.State.UNKNOWN, false);
|
||||
var withOk = offLoop.decide(NOW, g.idleSince(), g.quietCount(), g.status(),
|
||||
g.pushLoopActive(), g.leadKnown(), g.fleet(), LeadContextGauge.State.OK, false);
|
||||
|
||||
assertEquals(withUnknown.action(), withHigh.action(), g.name() + ": action must not depend on context");
|
||||
assertEquals(withUnknown.idleSinceNanos(), withHigh.idleSinceNanos(), g.name());
|
||||
assertEquals(withUnknown.quietCount(), withHigh.quietCount(), g.name());
|
||||
assertEquals(withUnknown.action(), withOk.action(), g.name() + ": nor on an OK reading");
|
||||
assertFalse(withHigh.contextNotified(), g.name() + ": contextNotified must stay false when the flag is off");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void bHighContextFiresEvenPastAnExhaustedQuietCapWithoutSpendingIt() {
|
||||
LeadHeartbeatLoop on = loop(3, true);
|
||||
LeadHeartbeatLoop.Decision d = on.decide(
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.HIGH, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.INJECT, d.action(),
|
||||
"an idle, quiet, HIGH-context lead is exactly the case the exhausted cap must not swallow");
|
||||
assertEquals(3, d.quietCount(), "the context notice is an event notice, not a quiet nudge — it must not "
|
||||
+ "spend or grow the consecutive-quiet budget");
|
||||
assertTrue(d.contextNotified(), "firing the notice sets the latch");
|
||||
}
|
||||
|
||||
@Test
|
||||
void cASecondTickWithTheLatchAlreadySetDoesNotTellTheLeadAgain() {
|
||||
LeadHeartbeatLoop on = loop(3, true);
|
||||
LeadHeartbeatLoop.Decision d = on.decide(
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.HIGH, true);
|
||||
assertEquals(LeadHeartbeatLoop.Action.QUIET_DONE, d.action(),
|
||||
"the lead was already told once about this HIGH stretch — telling it every tick would be nagging, "
|
||||
+ "not a notice");
|
||||
}
|
||||
|
||||
@Test
|
||||
void dFlappingBetweenHighAndUnknownNeverReInjectsWhileLatched() {
|
||||
LeadHeartbeatLoop on = loop(3, true);
|
||||
// The lead was already told once (latch = true from a prior HIGH tick).
|
||||
LeadHeartbeatLoop.Decision afterUnknown = on.decide(
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.UNKNOWN, true);
|
||||
assertTrue(afterUnknown.contextNotified(),
|
||||
"UNKNOWN means 'I could not look', not 'it got better' — it must not clear the latch");
|
||||
assertEquals(LeadHeartbeatLoop.Action.QUIET_DONE, afterUnknown.action(),
|
||||
"a still-latched, still-quiet tick must not inject just because the reading is UNKNOWN");
|
||||
|
||||
LeadHeartbeatLoop.Decision afterHighAgain = on.decide(
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.HIGH, afterUnknown.contextNotified());
|
||||
assertEquals(LeadHeartbeatLoop.Action.QUIET_DONE, afterHighAgain.action(),
|
||||
"HIGH -> UNKNOWN -> HIGH with the latch already set must never inject again — this is the "
|
||||
+ "flapping case that would otherwise cost a full lead its remaining turns");
|
||||
}
|
||||
|
||||
@Test
|
||||
void eAnOkReadingClearsTheLatchSoALaterHighInjectsAgain() {
|
||||
LeadHeartbeatLoop on = loop(3, true);
|
||||
LeadHeartbeatLoop.Decision afterOk = on.decide(
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.OK, true);
|
||||
assertFalse(afterOk.contextNotified(), "an OK reading clears the latch — the lead's context recovered");
|
||||
|
||||
LeadHeartbeatLoop.Decision afterHighAgain = on.decide(
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet(),
|
||||
LeadContextGauge.State.HIGH, afterOk.contextNotified());
|
||||
assertEquals(LeadHeartbeatLoop.Action.INJECT, afterHighAgain.action(),
|
||||
"the cleared latch lets a genuinely new HIGH stretch notify again");
|
||||
}
|
||||
|
||||
@Test
|
||||
void fStandingDownDoesNotBurnTheOneContextNotice() {
|
||||
LeadHeartbeatLoop on = loop(3, true);
|
||||
LeadHeartbeatLoop.Decision d = on.decide(
|
||||
NOW, IDLE_PAST, 2, AgentStatus.IDLE, true, true, quietFleet(),
|
||||
LeadContextGauge.State.HIGH, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.STAND_DOWN, d.action(), "ReplyPushLoop is active — stand aside");
|
||||
assertFalse(d.contextNotified(),
|
||||
"standing down must not spend the one notice this HIGH stretch gets — the latch stays clear so a "
|
||||
+ "later tick can still fire it");
|
||||
}
|
||||
|
||||
@Test
|
||||
void gAWorkingLeadWithHighContextIsStillLeadBusy() {
|
||||
LeadHeartbeatLoop on = loop(3, true);
|
||||
LeadHeartbeatLoop.Decision d = on.decide(
|
||||
NOW, IDLE_PAST, 0, AgentStatus.WORKING, false, true, quietFleet(),
|
||||
LeadContextGauge.State.HIGH, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.LEAD_BUSY, d.action(), "the status gate wins over the context notice");
|
||||
}
|
||||
|
||||
@Test
|
||||
void hAPendingDrivenInjectWhileHighSetsTheLatchTooSoTheLeadIsNotToldTwice() {
|
||||
LeadHeartbeatLoop on = loop(3, true);
|
||||
LeadHeartbeatLoop.Decision d = on.decide(
|
||||
NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, pendingFleet(),
|
||||
LeadContextGauge.State.HIGH, false);
|
||||
assertEquals(LeadHeartbeatLoop.Action.INJECT, d.action());
|
||||
assertTrue(d.contextNotified(),
|
||||
"the pending-driven nudge carries the context notice too (see contextNotice), so it must set the "
|
||||
+ "latch — otherwise the lead could be told twice by two different routes");
|
||||
}
|
||||
|
||||
// ── contextNotice text builder ─────────────────────────────────────────────────────────────
|
||||
|
||||
@Test
|
||||
void contextNoticeIsEmptyWhenDisabled() {
|
||||
var reading = new LeadContextGauge.Reading(LeadContextGauge.State.HIGH, 260_771L, 1);
|
||||
assertEquals("", LeadHeartbeatLoop.contextNotice(false, reading));
|
||||
}
|
||||
|
||||
@Test
|
||||
void contextNoticeIsEmptyWhenStateIsOk() {
|
||||
var reading = new LeadContextGauge.Reading(LeadContextGauge.State.OK, 1_000L, 0);
|
||||
assertEquals("", LeadHeartbeatLoop.contextNotice(true, reading));
|
||||
}
|
||||
|
||||
@Test
|
||||
void contextNoticeIsEmptyWhenStateIsUnknown() {
|
||||
assertEquals("", LeadHeartbeatLoop.contextNotice(true, LeadContextGauge.Reading.unknown()));
|
||||
}
|
||||
|
||||
@Test
|
||||
void contextNoticeNamesTheTokenCountWhenHighAndEnabled() {
|
||||
var reading = new LeadContextGauge.Reading(LeadContextGauge.State.HIGH, 260_771L, 1);
|
||||
String notice = LeadHeartbeatLoop.contextNotice(true, reading);
|
||||
assertFalse(notice.isEmpty());
|
||||
assertTrue(notice.contains("260771"), notice);
|
||||
assertTrue(notice.contains("1 compaction"), notice);
|
||||
assertTrue(notice.contains("fleet_handover"), notice);
|
||||
assertTrue(notice.contains("operator"), notice);
|
||||
}
|
||||
|
||||
@Test
|
||||
void contextNoticeOmitsTheTokenClauseRatherThanPrintingNull() {
|
||||
var reading = new LeadContextGauge.Reading(LeadContextGauge.State.HIGH, null, 2);
|
||||
String notice = LeadHeartbeatLoop.contextNotice(true, reading);
|
||||
assertFalse(notice.isEmpty());
|
||||
assertFalse(notice.toLowerCase().contains("null"), notice);
|
||||
assertTrue(notice.contains("2 compactions"), notice);
|
||||
}
|
||||
|
||||
// ── fleetd #609 review: the latch must mean "the notice reached the pane" ────────────────────
|
||||
//
|
||||
// These four drive LeadHeartbeatLoop.tick() directly (package-private, same reasoning as
|
||||
// ReplyPushLoop#tick(String) being directly testable) against a real AgentControl wrapping a
|
||||
// FailableHerdrClient, so the send path (agents.send -> herdr -> possible throw) is exercised
|
||||
// for real rather than assumed from decide()'s Decision alone.
|
||||
|
||||
private static final String LEAD = "term_lead";
|
||||
private static final String WORKER = "term_w1";
|
||||
|
||||
/** A HIGH reading with a fixed token/compaction count, for the four tests below. */
|
||||
private static LeadContextGauge.Reading highReading() {
|
||||
return new LeadContextGauge.Reading(LeadContextGauge.State.HIGH, 260_771L, 1);
|
||||
}
|
||||
|
||||
/**
|
||||
* Builds a real {@link LeadHeartbeatLoop} wired to {@code herdr} via a real {@link AgentControl},
|
||||
* a mutable fake clock, and a mutable roster so a test can change fleet state between ticks. The
|
||||
* lead is always reported IDLE by {@code herdr}, so every tick's outcome is governed only by the
|
||||
* idle-window/quiet-cap/context gates under test.
|
||||
*/
|
||||
private static LeadHeartbeatLoop tickableLoop(FailableHerdrClient herdr, AtomicLong now,
|
||||
List<MemberSession>[] rosterBox, InMemoryReplyInbox inbox,
|
||||
int quietNudgeCap, ScheduledExecutorService scheduler) {
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
PrimaryRegistry registry = new PrimaryRegistry(LEAD);
|
||||
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, agents, inbox, scheduler, 5, 100_000);
|
||||
return new LeadHeartbeatLoop(registry, agents, inbox, () -> rosterBox[0], pushLoop, scheduler,
|
||||
now::get, IDLE_AFTER_NANOS, 100_000L, quietNudgeCap, null,
|
||||
new LeadHeartbeatLoop.LeadContextSource(t -> highReading()), true);
|
||||
}
|
||||
|
||||
@Test
|
||||
void iAFailedSendDoesNotConsumeTheNotice() {
|
||||
var herdr = new FailableHerdrClient(LEAD);
|
||||
var now = new AtomicLong(NOW);
|
||||
@SuppressWarnings("unchecked")
|
||||
List<MemberSession>[] rosterBox = new List[]{List.of()}; // quiet: nothing pending
|
||||
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
LeadHeartbeatLoop loop = tickableLoop(herdr, now, rosterBox, inbox, 0, scheduler);
|
||||
|
||||
loop.tick(); // first injectable tick: only opens the idle window (WAIT_IDLE)
|
||||
now.addAndGet(TimeUnit.SECONDS.toNanos(400)); // now clearly past the quiet period
|
||||
|
||||
herdr.throwOnNextSend();
|
||||
loop.tick(); // quiet fleet, quiet cap exhausted (0), context HIGH, latch clear -> INJECT, send throws
|
||||
|
||||
assertEquals(0, herdr.sentTexts().size(), "the failed send must not have recorded any text");
|
||||
|
||||
loop.tick(); // same inputs — the latch must still be clear, so this must INJECT and send again
|
||||
assertEquals(1, herdr.sentTexts().size(),
|
||||
"a retried tick with the latch still clear must attempt the send again");
|
||||
assertTrue(herdr.sentTexts().get(0).contains("Your own context is nearly full"),
|
||||
"the retried, successful send must carry the notice: " + herdr.sentTexts().get(0));
|
||||
}
|
||||
|
||||
@Test
|
||||
void jASuccessfulSendDoesConsumeIt() {
|
||||
var herdr = new FailableHerdrClient(LEAD);
|
||||
var now = new AtomicLong(NOW);
|
||||
@SuppressWarnings("unchecked")
|
||||
List<MemberSession>[] rosterBox = new List[]{List.of()};
|
||||
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
LeadHeartbeatLoop loop = tickableLoop(herdr, now, rosterBox, inbox, 0, scheduler);
|
||||
|
||||
loop.tick(); // opens the idle window
|
||||
now.addAndGet(TimeUnit.SECONDS.toNanos(400));
|
||||
|
||||
loop.tick(); // INJECT, send succeeds -> latch set
|
||||
assertEquals(1, herdr.sentTexts().size());
|
||||
assertTrue(herdr.sentTexts().get(0).contains("Your own context is nearly full"), herdr.sentTexts().get(0));
|
||||
|
||||
loop.tick(); // same inputs — the lead was already told this stretch
|
||||
assertEquals(1, herdr.sentTexts().size(),
|
||||
"the next tick with the same inputs must not send a second notice");
|
||||
}
|
||||
|
||||
@Test
|
||||
void kAPendingDrivenInjectWithTheLatchAlreadySetSendsNoNotice() {
|
||||
var herdr = new FailableHerdrClient(LEAD);
|
||||
var now = new AtomicLong(NOW);
|
||||
@SuppressWarnings("unchecked")
|
||||
List<MemberSession>[] rosterBox = new List[]{List.of()};
|
||||
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
LeadHeartbeatLoop loop = tickableLoop(herdr, now, rosterBox, inbox, 0, scheduler);
|
||||
|
||||
loop.tick();
|
||||
now.addAndGet(TimeUnit.SECONDS.toNanos(400));
|
||||
loop.tick(); // latches the notice (quiet, HIGH, cap exhausted -> the forced context INJECT)
|
||||
assertEquals(1, herdr.sentTexts().size());
|
||||
assertTrue(herdr.sentTexts().get(0).contains("Your own context is nearly full"));
|
||||
|
||||
// Now make the fleet have real pending state, so the NEXT INJECT is pending-driven, not the
|
||||
// forced-context route — with the latch already set from the tick above.
|
||||
inbox.own(WORKER);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
rosterBox[0] = List.of(new MemberSession("p1", WORKER, "prof", MemberRole.DEV, "/cwd", null,
|
||||
0, 0, 0, MemberSession.State.READY, null, null));
|
||||
|
||||
loop.tick();
|
||||
|
||||
assertEquals(2, herdr.sentTexts().size(), "the pending-driven tick must still send a nudge");
|
||||
assertFalse(herdr.sentTexts().get(1).contains("Your own context is nearly full"),
|
||||
"a pending-driven INJECT while the latch is already set must carry no context notice: "
|
||||
+ herdr.sentTexts().get(1));
|
||||
}
|
||||
|
||||
@Test
|
||||
void lTheNoticeAppearsExactlyOnceAcrossThreeDifferentlyDrivenInjects() {
|
||||
var herdr = new FailableHerdrClient(LEAD);
|
||||
var now = new AtomicLong(NOW);
|
||||
@SuppressWarnings("unchecked")
|
||||
List<MemberSession>[] rosterBox = new List[]{List.of()};
|
||||
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
// quietNudgeCap=1 so a still-not-exhausted quiet nudge is available as the "forced" route below,
|
||||
// distinct from both the pending-driven route and the exhausted-cap forced-context route.
|
||||
LeadHeartbeatLoop loop = tickableLoop(herdr, now, rosterBox, inbox, 1, scheduler);
|
||||
|
||||
loop.tick(); // opens the idle window
|
||||
now.addAndGet(TimeUnit.SECONDS.toNanos(400));
|
||||
|
||||
// 1) pending-driven INJECT: real fleet state present. Sets the latch and carries the notice.
|
||||
inbox.own(WORKER);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
rosterBox[0] = List.of(new MemberSession("p1", WORKER, "prof", MemberRole.DEV, "/cwd", null,
|
||||
0, 0, 0, MemberSession.State.READY, null, null));
|
||||
loop.tick();
|
||||
assertEquals(1, herdr.sentTexts().size());
|
||||
assertTrue(herdr.sentTexts().get(0).contains("Your own context is nearly full"), herdr.sentTexts().get(0));
|
||||
|
||||
// 2) "forced" INJECT: nothing pending, but the quiet cap (1) is not yet exhausted, so decide()
|
||||
// nudges anyway. The latch is already set, so no notice.
|
||||
inbox.ack(WORKER, "m1");
|
||||
rosterBox[0] = List.of();
|
||||
loop.tick();
|
||||
assertEquals(2, herdr.sentTexts().size(), "the quiet-cap-not-yet-exhausted nudge must still fire");
|
||||
assertFalse(herdr.sentTexts().get(1).contains("Your own context is nearly full"), herdr.sentTexts().get(1));
|
||||
|
||||
// 3) pending-driven INJECT again. Still latched, still no notice.
|
||||
inbox.publish(WORKER, "m2", "hello again");
|
||||
rosterBox[0] = List.of(new MemberSession("p1", WORKER, "prof", MemberRole.DEV, "/cwd", null,
|
||||
0, 0, 0, MemberSession.State.READY, null, null));
|
||||
loop.tick();
|
||||
assertEquals(3, herdr.sentTexts().size());
|
||||
assertFalse(herdr.sentTexts().get(2).contains("Your own context is nearly full"), herdr.sentTexts().get(2));
|
||||
|
||||
long noticeCount = herdr.sentTexts().stream()
|
||||
.filter(t -> t.contains("Your own context is nearly full")).count();
|
||||
assertEquals(1, noticeCount,
|
||||
"the notice text must appear exactly once across all three sends: " + herdr.sentTexts());
|
||||
}
|
||||
|
||||
/**
|
||||
* Fake herdr client for the four tests above: always reports {@code lead} as IDLE, records the
|
||||
* {@code text} of every {@code agent.prompt} call, and can be told to throw on the very next
|
||||
* {@code agent.prompt} call — standing in for one transient herdr send failure.
|
||||
*/
|
||||
private static final class FailableHerdrClient implements HerdrClient {
|
||||
private static final ObjectMapper MAPPER = new ObjectMapper();
|
||||
private final String lead;
|
||||
private final List<String> sentTexts = new ArrayList<>();
|
||||
private boolean throwOnNextSend = false;
|
||||
|
||||
FailableHerdrClient(String lead) {
|
||||
this.lead = lead;
|
||||
}
|
||||
|
||||
void throwOnNextSend() {
|
||||
throwOnNextSend = true;
|
||||
}
|
||||
|
||||
List<String> sentTexts() {
|
||||
return List.copyOf(sentTexts);
|
||||
}
|
||||
|
||||
@Override
|
||||
@SuppressWarnings("unchecked")
|
||||
public JsonNode call(String method, Object params) {
|
||||
if ("agent.get".equals(method)) {
|
||||
return MAPPER.createObjectNode()
|
||||
.set("agent", MAPPER.createObjectNode()
|
||||
.put("terminal_id", lead)
|
||||
.put("agent_status", "idle"));
|
||||
}
|
||||
if ("agent.prompt".equals(method)) {
|
||||
if (throwOnNextSend) {
|
||||
throwOnNextSend = false;
|
||||
throw new RuntimeException("simulated transient herdr send failure");
|
||||
}
|
||||
Map<String, Object> p = params instanceof Map ? (Map<String, Object>) params : Map.of();
|
||||
sentTexts.add(String.valueOf(p.get("text")));
|
||||
}
|
||||
return MAPPER.createObjectNode();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user