Compare commits

..

8 Commits

Author SHA1 Message Date
Dai Ha 388ef5a3c3 fleetd #615: write FAILED instead of leaving status(token) stuck at IN_PROGRESS
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Failing after 2m25s
LeadRollover.runRollover made two unwrapped agents.send calls. HerdrException
is unchecked, and the production continuationRunner is a bare virtual thread
with no uncaught-exception handler, so a throw from either send call killed
the continuation silently — confirm() had already written IN_PROGRESS into
outcomes before scheduling it, and nothing ever overwrote that entry with a
terminal state.

Wrap the whole continuation body in one try/catch(RuntimeException), matching
the local convention already used around agents.status in
waitUntilAtTurnBoundary. On a throw, write a new terminal RollState.FAILED
entry naming the exception, in the same diagnostic style as
TURN_NEVER_SETTLED and CLEAR_NEVER_SETTLED.

Two new tests make send() throw on the /clear call and on the bootstrap-text
call respectively, each asserting status(token) reports FAILED, not
IN_PROGRESS. Reverting only the production catch (keeping the tests) turns
both red; restoring it turns them green again.
2026-09-22 10:17:12 +07:00
ltms 076cc43f7b Merge #614: skip unreadableFileIsUnknown honestly when root ignores the read bit
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 49s
CI / build (push) Failing after 2m20s
CI has been red on main itself since #602/#606, on this one test, so the build has been giving no second opinion on any PR. Cause: the CI job runs in a container as root. setReadable(false) really does clear the read bit, so the test's own setup guard passes, but root opens the file anyway and the gauge correctly returns OK. The test was asserting on a condition the environment never created.

The fix adds an assumeFalse(Files.isReadable(file), ...) after the chmod and before the gauge is built, inside the existing try, so the finally still restores the bit on a skip.

Verified by me on a scratch worktree merging this onto 955b9ea:
- 1864 tests, 0 failures, 0 errors, 0 skipped, 149 surefire reports, mvn exit 0. The suite-wide skipped=0 is the point: the fix did not quietly turn the test into a permanent skip.
- LeadContextGaugeTest on this non-root Mac: 9 tests, 0 skipped, and unreadableFileIsUnknown present in the report. The assumption does not fire here, so developers keep the coverage.
- Mutation: made the IOException path return OK instead of UNKNOWN. unreadableFileIsUnknown failed with "expected: <UNKNOWN> but was: <OK>". The test still has teeth. Production file reverted, git diff clean before merge.

Known trade, recorded rather than hidden: under root this case is now covered by nothing at all. A skip is honest about that, which an assertion on an unreachable state was not. The durable fix is to run the CI build as a non-root user; that is a CI configuration change and out of scope here.
2026-09-20 12:43:43 +02:00
Dai Ha bad47a8444 fleetd CI: skip unreadableFileIsUnknown honestly when root ignores the read bit
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Successful in 2m27s
The test set the file's read bit off via setReadable(false), but on the
Gitea CI runner (root inside the container) the OS ignores that bit and
opens the file anyway, so the test asserted on a condition it never
actually created (LeadContextGaugeTest.java:142 UNKNOWN vs OK, CI run
1887 job 3104, commit fa62e99 on main).

Add Files.isReadable(file) after setReadable(false) and before the gauge
runs, and assumeFalse on it: a skip means "could not set up the case",
never "the behaviour is fine". Restores the read bit either way so
@TempDir cleanup still works.
2026-09-20 17:40:47 +07:00
ltms 955b9ea013 Merge #610: nudge an idle lead to hand over when its own context reads HIGH (fleetd #609)
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 52s
CI / build (push) Failing after 1m51s
Closes fleetd #609. Completes the second half of the context work: #602/#606 could detect a full lead context, and nothing acted on it. LeadHeartbeatLoop now offers a handover when the lead's own gauge reads HIGH.

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

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

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

Known and deliberately not blocked: Fleetd.main's own one-line call to leadContextSource is not pinned by a test. That is pre-existing and class-wide, tracked in #612.
2026-09-20 12:37:43 +02:00
Dai Ha 89cb8ff79b fleetd #609 review: the context latch must mean the notice reached the pane
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Failing after 2m6s
Fixes the PR #610 review blocker: LeadHeartbeatLoop committed contextNotified
before injectNudge attempted the send, so a transient herdr failure marked the
lead as told when nothing reached its pane, and contextNotice() carried no
latch at all, so a pending-driven INJECT re-appended the notice on every tick
while the context stayed HIGH.

- injectNudge now reports whether agents.send succeeded and persists
  contextNotified only when a notice was actually included in the text and
  the send did not throw. The latch is split out of applyDecision (kept for
  idleSinceNanos/quietCount, applied unconditionally as before) so it is
  written on the success path only, once per branch in tick().
- contextNotice gained an overloaded 3-arg form gated on the latch as it
  stood before the tick's decision; the existing 2-arg form delegates to it
  with alreadyNotified=false, so all pre-existing callers/tests are unchanged.
- tick() is now package-private (mirrors ReplyPushLoop#tick(String)) so tests
  can drive the real send path with a fake AgentControl instead of only the
  pure decide() function.
- Added tests I-L covering: a failed send does not consume the notice and
  retries; a successful send does; the text is gated when the latch is
  already set; and the notice appears exactly once across three differently
  driven INJECTs.

Both required mutations verified red and reverted:
1. Setting the latch from the Decision regardless of send outcome -> test I
   (iAFailedSendDoesNotConsumeTheNotice) fails.
2. Dropping the latch argument at the contextNotice call site -> tests K
   (kAPendingDrivenInjectWithTheLatchAlreadySetSendsNoNotice) and L
   (lTheNoticeAppearsExactlyOnceAcrossThreeDifferentlyDrivenInjects) fail.
2026-09-20 17:34:05 +07:00
ltms fa62e9906d Merge #611: make the collected-ticket nudge test deterministic (fleetd #608)
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 49s
CI / build (push) Failing after 1m48s
Replaces a wall-clock bet with a manually-driven scheduler, so the tick runs
only when the test runs it. Also closes a second, smaller race the brief did
not name: waiting on Phase.DONE is not enough, because complete() can publish
isDone() before every whenComplete dependent has run.

Verified by the lead: full suite 1841/1841 green in a clean worktree, and an
independent mutation (hasTicketWork forced true) diagnosed as an equivalent
mutant — ReplyPushLoop.injectNudge re-reads the pending collections and returns
early, so that line cannot reach agent.prompt. The worker's own mutation
(ticketCollected made a no-op) is the one that reaches the observable, and it
killed.
2026-09-20 12:26:16 +02:00
Dai Ha d7390ccd37 fleetd #609 review: repair a garbled comment carried over from the brief
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 1m32s
CI / build (pull_request) Failing after 1m41s
The brief's sentence about a null token count at HIGH was broken, and the
worker copied it into the source verbatim. The code was already right; only
the comment was unreadable.

Says what is actually true: a HIGH reading always carries a non-null token
count today, because LeadContextGauge only reaches HIGH by comparing a number
against HIGH_THRESHOLD_TOKENS. That invariant lives in another class and
nothing asserts it, so the branch stays.
2026-09-20 17:20:48 +07:00
Dai Ha 60496831c2 fleetd #609: nudge an idle lead to hand over when its own context reads HIGH
CI / shell-tests (pull_request) Failing after 6s
CI / contract (pull_request) Successful in 48s
CI / build (pull_request) Failing after 2m5s
LeadHeartbeatLoop can now append a text-only notice to its nudge when the lead's
own LeadContextGauge reading is HIGH and leadHeartbeat.contextHighNudge is on.
Fires once per HIGH stretch (a latch, cleared only by a later OK reading; UNKNOWN
neither sets nor clears it), never spends the quietNudgeCap budget, and never
rolls a pane itself — only the operator can approve a handover.

- LeadContextGauge.Reading.unknown() widened to public for LeadContextSource.none()
- FleetConfig.LeadHeartbeat gains contextHighNudge (null/false = off, unchanged default)
- LeadHeartbeatLoop.decide gains context/contextNotified; Fleetd wires a new
  leadContextLookup/leadContextSource factory pair (LeadHeartbeatLoop.LeadContextSource)
- fleetd.example.yaml documents the new key
2026-09-20 17:11:54 +07:00
11 changed files with 1116 additions and 50 deletions
+10 -1
View File
@@ -77,16 +77,25 @@ bind:
# a lead turn nobody asked for), so upgrading the daemon must never switch it on for you. Absent
# 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);
}
}
@@ -231,7 +231,9 @@ public final class LeadRollover {
* #status} could wrongly answer {@link #UNKNOWN} ("nothing was ever requested") for a roll
* that is, in fact, actively running. This is not sticky: the deferred continuation
* overwrites this same entry with a terminal state ({@link #ROLLED}, {@link
* #TURN_NEVER_SETTLED}, or {@link #CLEAR_NEVER_SETTLED}) once it finishes.
* #TURN_NEVER_SETTLED}, {@link #CLEAR_NEVER_SETTLED}, or {@link #FAILED}) once it finishes
* — including by throwing, which fleetd #615's catch in {@link #runRollover} now turns into
* {@link #FAILED} instead of leaving this entry stuck forever.
*/
IN_PROGRESS,
/**
@@ -253,6 +255,19 @@ public final class LeadRollover {
* within {@code clearSettleSeconds} — {@code bootstrapText} was never sent.
*/
CLEAR_NEVER_SETTLED,
/**
* fleetd #615: the deferred continuation threw a {@link RuntimeException} — most likely a
* {@link dev.ltms.fleet.herdr.HerdrException} out of one of the two unwrapped {@code
* agents.send} calls in {@link #runRollover} — and the continuation thread died with it.
* Before this state existed, that throw left {@link #outcomes} holding {@link #IN_PROGRESS}
* forever, because the production {@code continuationRunner} is a bare virtual thread with
* no uncaught-exception handler and nothing downstream of the throw ever ran to write a
* terminal outcome. {@code detail} names the exception, so a reader has something to act on
* — the same diagnostic style as {@link #TURN_NEVER_SETTLED} and {@link
* #CLEAR_NEVER_SETTLED}. The roll is dead at this point and does not retry itself; a stuck
* lead must {@link #open} a fresh request.
*/
FAILED,
/**
* {@code token} names nothing this instance currently knows about: never issued by {@link
* #open}, dropped by {@link #cancel}, or aged out of {@link #outcomes}'s bounded history.
@@ -493,8 +508,42 @@ public final class LeadRollover {
* entirely after {@link #confirm} has returned to its caller — see this class's javadoc for the
* four-step order. There is no result to return to by this point, so every outcome is logged
* only.
*
* <p><strong>fleetd #615 — the whole body is wrapped in one {@code try}.</strong> The two {@code
* agents.send} calls below are not wrapped individually: {@code send} → {@code agentCall} →
* {@code herdr.call} can throw an unchecked {@link dev.ltms.fleet.herdr.HerdrException} (see
* {@code AgentControl.java}), and the production {@code continuationRunner} is a bare virtual
* thread with no uncaught-exception handler (see this class's public constructor). Before this
* fix, either throw killed the continuation thread silently, leaving the {@link
* RollState#IN_PROGRESS} entry {@link #confirm} wrote at hand-off stuck forever — {@link
* #status} had no way to tell a dead roll from one still genuinely running. The {@code catch}
* below is scoped to the method body rather than to each {@code send} call individually, so it
* also covers anything else added to this continuation later, not just today's two call sites —
* the same reasoning that put the write-a-terminal-outcome step at each of this method's other
* exits (see the {@link RollState#TURN_NEVER_SETTLED} and {@link RollState#CLEAR_NEVER_SETTLED}
* branches below) rather than inside the helpers that detect them.</p>
*
* <p>Only {@link RuntimeException} is caught, matching the local convention {@link
* #waitUntilAtTurnBoundary} already set around its own {@code agents.status} call — not the
* broader {@link Exception} or {@link Throwable}, which would also swallow something like an
* {@link OutOfMemoryError} this continuation has no business handling.</p>
*/
private void runRollover(PendingRollover p, FleetConfig.LeadRollover cfg) {
try {
runRolloverUnguarded(p, cfg);
} catch (RuntimeException e) {
log.warn("lead-rollover: continuation for token={} lead={} threw {} — the roll is dead; "
+ "no further step in this continuation will run",
p.token(), p.leadTerminal(), e.toString(), e);
outcomes.put(p.token(), new RollStatus(RollState.FAILED,
"the roll's continuation threw " + e.toString() + " — the roll is dead and will "
+ "not retry itself; check the daemon log for the stack trace, then open() "
+ "a fresh rollover request"));
}
}
/** The actual body of {@link #runRollover}, unwrapped — see that method's javadoc for the catch. */
private void runRolloverUnguarded(PendingRollover p, FleetConfig.LeadRollover cfg) {
String lead = p.leadTerminal();
long rollStartMillis = nowMillis.getAsLong();
TurnSettleResult turnResult = waitUntilAtTurnBoundary(lead, cfg.turnSettleSeconds());
@@ -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());
}
}
@@ -14,6 +14,7 @@ import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.junit.jupiter.api.Assumptions.assumeFalse;
/**
* Ticket "lead context gauge" — fleetd could not see how full a lead's own Claude Code context
@@ -137,6 +138,15 @@ class LeadContextGaugeTest {
Path file = tmp.resolve("projects").resolve("some-project-slug").resolve(SESSION_ID + ".jsonl");
assertTrue(file.toFile().setReadable(false), "test setup: must be able to revoke read permission");
try {
// setReadable(false) really did clear the read bit (asserted above), but that alone
// does not prove the file is UNREADABLE: running as root (e.g. a CI container) ignores
// the read bit and opens the file anyway. Files.isReadable checks what actually happens
// on open, not the bit. When it still reports readable, this test cannot create the
// condition it needs on this machine, so it skips honestly instead of asserting on a
// state that was never reached. A skip here means "I could not set up the case", NOT
// "the UNKNOWN behaviour is fine" -- it is not evidence either way.
assumeFalse(Files.isReadable(file),
"runs as root (CI container): the read bit does not stop root, so this case cannot be set up here");
LeadContextGauge gauge = new LeadContextGauge();
LeadContextGauge.Reading reading = gauge.read(configDir, SESSION_ID, "claude");
assertEquals(LeadContextGauge.State.UNKNOWN, reading.state());
@@ -1359,4 +1359,92 @@ class LeadRolloverTest {
+ "it were the measured wait duration: " + message);
}
}
// ---- fleetd #615: a HerdrException out of either unwrapped agents.send call must leave a ----
// ---- TERMINAL FAILED outcome, never a stuck IN_PROGRESS ---------------------------------------
@Test
@DisplayName("[fleetd #615 — 1] send() throwing on the /clear call leaves status(token) "
+ "reporting FAILED, not stuck at IN_PROGRESS")
void sendThrowingOnClearLeavesStatusReportingFailed() throws IOException {
FakeHerdr fake = new FakeHerdr(); // default idle — the turn-settle wait passes immediately
HerdrClient throwsOnClear = new HerdrClient() {
@Override
public JsonNode call(String method, Object params) throws HerdrException {
if ("agent.prompt".equals(method) && String.valueOf(params).contains("/clear")) {
throw new HerdrException("simulated herdr transport failure sending /clear");
}
return fake.call(method, params);
}
@Override
public void close() {
fake.close();
}
};
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(throwsOnClear, cfg(handover.toString()), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "every synchronous gate passes; the throw happens only "
+ "inside the deferred continuation, which this test's synchronous runner has "
+ "already run to completion by the time confirm() returns");
LeadRollover.RollStatus status = rollover.status(pending.token());
assertEquals(LeadRollover.RollState.FAILED, status.state(),
"a HerdrException out of the /clear send must leave a TERMINAL FAILED outcome — "
+ "before fleetd #615's fix, the continuation thread died silently and "
+ "status() was stuck reporting the IN_PROGRESS confirm() wrote at hand-off, "
+ "forever: got " + status.state() + " / " + status.detail());
assertNotEquals(LeadRollover.RollState.IN_PROGRESS, status.state());
assertTrue(status.detail().contains("HerdrException"), "the detail must name the exception "
+ "so an operator reading status() has something to act on: " + status.detail());
}
@Test
@DisplayName("[fleetd #615 — 2] send() throwing on the bootstrap-text call (after /clear "
+ "succeeded and the pane settled) also leaves status(token) reporting FAILED — a "
+ "DIFFERENT exit from the /clear-throw case above")
void sendThrowingOnBootstrapTextLeavesStatusReportingFailed() throws IOException {
FakeHerdr fake = new FakeHerdr(); // default idle throughout — both settle waits pass promptly
HerdrClient throwsOnBootstrapText = new HerdrClient() {
@Override
public JsonNode call(String method, Object params) throws HerdrException {
if ("agent.prompt".equals(method) && String.valueOf(params).contains("read the handover file")) {
throw new HerdrException("simulated herdr transport failure sending bootstrapText");
}
return fake.call(method, params);
}
@Override
public void close() {
fake.close();
}
};
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(throwsOnBootstrapText, cfg(handover.toString()), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "every synchronous gate passes; the throw happens only "
+ "inside the deferred continuation, which this test's synchronous runner has "
+ "already run to completion by the time confirm() returns");
assertEquals(1, promptCallCount(fake), "sanity: /clear was sent and settled — only the "
+ "SECOND agent.prompt call (bootstrapText) threw");
LeadRollover.RollStatus status = rollover.status(pending.token());
assertEquals(LeadRollover.RollState.FAILED, status.state(),
"a HerdrException out of the bootstrapText send — a DIFFERENT exit from the /clear "
+ "throw, reached only after /clear already succeeded and the pane already "
+ "settled — must also leave a TERMINAL FAILED outcome, not a stuck "
+ "IN_PROGRESS: got " + status.state() + " / " + status.detail());
assertNotEquals(LeadRollover.RollState.IN_PROGRESS, status.state());
assertTrue(status.detail().contains("HerdrException"), "the detail must name the exception "
+ "so an operator reading status() has something to act on: " + status.detail());
}
}
@@ -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() {
}
}
}