diff --git a/fleetd/fleetd.example.yaml b/fleetd/fleetd.example.yaml index ae16542..41f54dd 100644 --- a/fleetd/fleetd.example.yaml +++ b/fleetd/fleetd.example.yaml @@ -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 diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index d06c01f..00cf817 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -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. + * + *

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 leadContextLookup(LeadContextGauge gauge, AgentControl agents, + Supplier> liveLeadTerminals, Function 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> liveLeadTerminals, Function 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 diff --git a/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java b/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java index 6c4d213..2d25241 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java +++ b/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java @@ -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; diff --git a/fleetd/src/main/java/dev/ltms/fleet/lead/LeadContextGauge.java b/fleetd/src/main/java/dev/ltms/fleet/lead/LeadContextGauge.java index fc16d01..b24d48e 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/lead/LeadContextGauge.java +++ b/fleetd/src/main/java/dev/ltms/fleet/lead/LeadContextGauge.java @@ -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); } } diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadHeartbeatLoop.java b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadHeartbeatLoop.java index 55e4940..2636d69 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadHeartbeatLoop.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadHeartbeatLoop.java @@ -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). * + * + *

fleetd #609 — context-high notice. 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> 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> 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 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. + * + *

fleetd #609 review: the context latch ({@link #contextNotified}) is deliberately not + * 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 before 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); diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdLeadContextLookupTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdLeadContextLookupTest.java new file mode 100644 index 0000000..29fea46 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdLeadContextLookupTest.java @@ -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. + * + *

{@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 /projects//.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 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 liveLeadTerminals = Map.of(LEAD_TERMINAL, LEAD_NAME); + Function 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 liveLeadTerminals = Map.of(LEAD_TERMINAL, LEAD_NAME); + AgentControl agents = agentControlStub(SESSION_ID, "claude", "idle"); + Function 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 liveLeadTerminals = Map.of(LEAD_TERMINAL, LEAD_NAME); + AgentControl agents = agentControlStub(SESSION_ID, "opencode", "idle"); + Function 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"); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdLeadContextSourceWiringTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdLeadContextSourceWiringTest.java new file mode 100644 index 0000000..6d79735 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdLeadContextSourceWiringTest.java @@ -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). + * + *

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();}. + * + *

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 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()); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/LeadHeartbeatLoopTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/LeadHeartbeatLoopTest.java index 1e7dc29..d327861 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/LeadHeartbeatLoopTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/LeadHeartbeatLoopTest.java @@ -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. * *

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. + * + *

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 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[] 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[] 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[] 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[] 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[] 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 sentTexts = new ArrayList<>(); + private boolean throwOnNextSend = false; + + FailableHerdrClient(String lead) { + this.lead = lead; + } + + void throwOnNextSend() { + throwOnNextSend = true; + } + + List 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 p = params instanceof Map ? (Map) params : Map.of(); + sentTexts.add(String.valueOf(p.get("text"))); + } + return MAPPER.createObjectNode(); + } + + @Override + public void close() { + } + } }