Compare commits

..

16 Commits

Author SHA1 Message Date
Dai Ha 6cb31a10e4 fleetd #618: fix the third stale spot the brief missed
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 1m30s
CI / build (pull_request) Failing after 2m6s
The method-level javadoc on FleetConfig.warnConflictingAutoCompactWindows
(above the log.warn call) still claimed the autoCompactWindow vs
CLAUDE_CODE_AUTO_COMPACT_WINDOW precedence was 'intentionally not
asserted' and cited fleetd.yaml's now-corrected comment as evidence the
question was open. Replace it with the measured answer from #618: the
env var wins, so autoCompactWindow is inert on a profile that sets both.
Kept the WARN-not-throw rationale paragraph above it untouched (#601)
and kept the ClaudeCodeArguments cross-reference, which now points to an
agreeing claim instead of a contradicting one. No behaviour change.
2026-09-22 10:50:59 +07:00
Dai Ha 8368a274a0 fleetd #618: state the measured auto-compact precedence, not 'unverified'
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 1m4s
CI / build (pull_request) Failing after 1m57s
ClaudeCodeArguments.withAutoCompactWindow's javadoc and FleetConfig's
warnConflictingAutoCompactWindows WARN text both used to say the
precedence between --autocompact and CLAUDE_CODE_AUTO_COMPACT_WINDOW was
not verified. fleetd #618 measured it: the env var wins, so the flag has
no effect when both are set. Update both texts to say so, name #618, and
warn that deleting the env var to resolve the conflict LOWERS the live
window rather than fixing anything. No behaviour change; the WARN still
fires on the same condition and stays a WARN (per #601).
2026-09-22 10:46:31 +07:00
ltms 17127efb88 Merge #601: pass auto-compact window to leads; warn instead of refusing on a conflict (CB-617)
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 1m16s
CI / build (push) Failing after 2m33s
2026-09-22 05:22:56 +02:00
ltms 203f034528 Merge #617: write FAILED instead of leaving a dead roll stuck at IN_PROGRESS (fleetd #615)
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 1m14s
CI / build (push) Failing after 1m47s
2026-09-22 05:21:21 +02:00
ltms 9ee16f5b85 Merge #616: report role-fallback gaps at boot, name contextHighNudge (fleetd #613)
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 1m15s
CI / build (push) Failing after 1m43s
2026-09-22 05:17:52 +02:00
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
Dai Ha be6c45ff78 CB-617 review: warn instead of refuse on conflicting autoCompactWindow
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Failing after 2m9s
rejectConflictingAutoCompactWindows threw and stopped fleetd from starting when a
Claude Code profile's autoCompactWindow flag and CLAUDE_CODE_AUTO_COMPACT_WINDOW
env var disagreed. Under launchd that is a restart loop, and the config that
would fix it (fleetd.yaml) is gitignored, so the cause is invisible on the host
where it bites (measured live: 4 profiles on this host trip it, including the
lead's own profile and the one every worker spawns on).

Renamed to warnConflictingAutoCompactWindows: it now logs a WARN naming each
offending profile with BOTH values (autoCompactWindow=... and
env.CLAUDE_CODE_AUTO_COMPACT_WINDOW=...) instead of throwing, so the daemon
starts and an operator can fix the config without reading the source. Equal
values still load silently.

Also reworded ClaudeCodeArguments' javadoc, which stated as fact that the env
var takes precedence over the flag. That was never measured, and this host's
own fleetd.yaml comment asserts the opposite — the javadoc no longer picks a
side.
2026-09-22 10:13:14 +07:00
Dai Ha e99cb70a8b CB-617: pass auto-compact window to leads 2026-09-22 10:12:56 +07:00
Dai Ha 987ccef4c7 fleetd #613: log role-fallback gaps at boot, name contextHighNudge in the heartbeat line
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 1m17s
CI / build (pull_request) Failing after 1m48s
- reportRoleFallbackGaps(cfg), called right after cfg.validateAll() in Fleetd.main, logs every
  MemberRole with no fleet.<role>s: pool (naming the profile count and the resolved
  defaultProfileFor(role) first choice) and, separately, every role with no
  fleet.charters.<role>: entry. Log only — the deliberate 'unconstrained' fallback in
  FleetConfig#candidateProfiles / CompositePeerLauncher#poolFor is unchanged, and a config with
  profiles: and no fleet: block still starts and still spawns.
- LeadHeartbeatLoop#start()'s boot line now also names contextHighNudge (fleetd #609), alongside
  the three settings it already logged.
- RoleFallbackGapReportTest (new) and two new LeadHeartbeatLoopTest cases pin both lines' content
  via a ListAppender, raising the dev.ltms.fleet logger past logback-test.xml's WARN override for
  the INFO-level lines.
2026-09-22 10:12:09 +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
17 changed files with 1722 additions and 87 deletions
+15 -6
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
@@ -200,13 +209,13 @@ herdrSocket: ~/.config/herdr/herdr.sock
# e.g. `env DISPLAY=:10.0 idea {dir}`. Best-effort: a failure is logged, never
# fails the spawn. Omit to open the member's module by hand. There is no close
# half yet — an opened module stays open until the operator closes it.
# autoCompactWindow → opt-in, default off. A bounded token window that forces a spawned member to
# compact its context instead of running on the backend's own default and dying
# mid-turn (losing its fleet_reply — the whole point of the turn — with it).
# autoCompactWindow → opt-in, default off. A bounded token window that forces a launched Claude Code
# lead or member to compact its context instead of running on the backend's own
# default. A member that runs out of context can die mid-turn and lose its fleet_reply.
# Validated at config load to [100000, 1000000] — the band Claude Code's own
# --autocompact flag accepts.
# CROSS-BACKEND SEMANTICS DIFFER: on claude-code this is a launch-time
# `--autocompact <tokens>` flag — the member compacts AT this window. opencode
# `--autocompact <tokens>` flag — the Claude Code session compacts AT this window. opencode
# has no equivalent flag (it only forces `compaction.auto: true`, unconditionally,
# already), so this is instead applied as the model's `limit.context` in the
# generated opencode.json — the member compacts WITHIN this window, not exactly
@@ -406,7 +415,7 @@ profiles:
# ideMcpUrl: http://127.0.0.1:29170/index-mcp/streamable-http # opt-in (CB-634): IDE code intelligence, pinned to the worktree
# ideProjectDir: fleetd # CB-634: module dir the IDE opens + the overlay pins (this repo's pom is in fleetd/)
# ideOpenCommand: env DISPLAY=:10.0 idea {dir} # CB-634 auto-open: opens {dir} in the IDE at spawn; omit to open by hand
# autoCompactWindow: 250000 # opt-in: bound member context; claude-code compacts AT this, opencode within it (model limit.context)
# autoCompactWindow: 250000 # opt-in: bound Claude Code lead/member context; claude-code compacts AT this, opencode within it (model limit.context)
gx11: # a second backend, so `placement: weighted` has a choice
baseUrl: http://gx01.gw:8000 # self-hosted; ccs handles the model + token
placement: tab
+121 -1
View File
@@ -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;
@@ -51,6 +53,7 @@ import dev.ltms.fleet.rest.FleetApp;
import dev.ltms.fleet.session.GitWorktrees;
import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.session.SessionReaper;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
@@ -174,6 +177,12 @@ public final class Fleetd {
// and the Fleetd-startup tests actually pin — see FleetConfig#validateAll's javadoc for
// why a name-by-name list here would have the same defect it replaces.
cfg.validateAll();
// fleetd #613: validateAll() (validateMembers() inside it) only refuses a slot that names a
// bad role or profile — it says nothing about a role that has NO pool or NO charter at all,
// because both are legitimate ("unconstrained") states, not errors. Report them here, right
// after validation passes, so an operator sees the gap once per restart instead of finding
// it later in a roster row (see reportRoleFallbackGaps' javadoc for the measured cause).
reportRoleFallbackGaps(cfg);
// fleetd #469, follow-up to #464: validateAll() (and validateCharters() inside it) only
// checks that a charter's KEY is a role wire name and its text is non-blank — it never
// looks at what the text actually names. This is the separate check that does: it asks
@@ -563,10 +572,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 +1649,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
@@ -2114,6 +2182,58 @@ public final class Fleetd {
}
}
/**
* fleetd #613: {@code FleetConfig.candidateProfiles(MemberRole)} (FleetConfig.java:1827) and
* {@code CompositePeerLauncher.poolFor} (CompositePeerLauncher.java:598-603) both fall back to
* <em>every</em> configured profile when a role has no {@code fleet.<role>s:} pool — a
* deliberate "unconstrained" behaviour, kept unchanged here, that lets a config with only
* {@code profiles:} and no {@code fleet:} block still spawn. That fallback is silent, and on
* the host that opened this ticket it widened an unqualified {@code hunter} spawn to all 8
* configured profiles and picked {@code local} as the resolved first choice — a profile every
* other pool on that same config gives weight 0 to. Report it once at boot instead, naming both
* how many profiles the gap opens onto and the exact first choice, since the first choice (not
* the pool size) is what actually surprised the operator.
*
* <p>A missing {@code fleet.charters.<role>:} entry is reported separately: the member still
* runs, but with only the launcher's own reply charter and no role contract. Unlike the pool
* gap this already has a per-spawn instrument ({@code HerdrPeerLauncher.logCharterReceipt},
* {@code SessionManager}'s {@code charterSource} roster field) — this boot line is the same
* information surfaced once, up front, rather than discovered per member later.
*
* <p>Never refuses to start over either gap — both are legitimate configurations, and this is a
* report, not a validation. Package-private so a test can call it directly and capture the log
* via a {@link ch.qos.logback.core.read.ListAppender}, the same pattern {@link
* #reportExhaustedPatternGap} and {@link #reportMemberCredentialsGap} already use.
*/
static void reportRoleFallbackGaps(FleetConfig cfg) {
List<String> poolGaps = new ArrayList<>();
List<String> charterGaps = new ArrayList<>();
int profileCount = cfg.profiles().size();
for (MemberRole role : MemberRole.values()) {
boolean hasPool = cfg.fleet() != null && !cfg.fleet().profilesFor(role).isEmpty();
if (!hasPool) {
String firstChoice = cfg.defaultProfileFor(role);
poolGaps.add(role.wireName() + " (may land on any of " + profileCount
+ " profile(s), first choice "
+ (firstChoice == null ? "none — no profiles configured" : "'" + firstChoice + "'")
+ ")");
}
String charter = cfg.fleet() == null ? null : cfg.fleet().charterFor(role);
if (charter == null || charter.isBlank()) {
charterGaps.add(role.wireName());
}
}
if (!poolGaps.isEmpty()) {
log.info("role fallback: no fleet.<role>s: pool for {} — an unqualified spawn of that "
+ "role falls back to every configured profile (deliberate; see "
+ "FleetConfig#candidateProfiles)", poolGaps);
}
if (!charterGaps.isEmpty()) {
log.info("role fallback: no fleet.charters: entry for {} — that role runs with only "
+ "the launcher's reply charter, no role contract", charterGaps);
}
}
/**
* fleetd #474: the one place both the startup call (right after {@code cfg.validateAll()} in
* {@link #main}) and the reload call (wired into {@code config}'s {@code extraValidation} above,
@@ -462,8 +462,8 @@ public record FleetConfig(
* profile that does not opt in. Read live off the current config, so it is
* HOT: a change takes effect on the next exhaustion classification / spawn,
* no restart needed.
* @param autoCompactWindow opt-in per-profile token window that forces a spawned member to
* auto-compact its context at (Claude Code) or within (opencode) a bound the
* @param autoCompactWindow opt-in per-profile token window that forces a launched Claude Code
* session to auto-compact its context at, or an opencode session within, a bound the
* operator chooses, instead of the backend's own default. {@code null} (the
* default) leaves today's behaviour exactly — opencode already forces
* {@code compaction.auto: true} unconditionally (CB-523) but has no absolute
@@ -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;
@@ -1864,6 +1875,7 @@ public record FleetConfig(
rejectDuplicateMemberSlots(yaml);
rejectNegativeMaxLoad(yaml);
rejectAutoCompactWindowOutOfRange(yaml);
warnConflictingAutoCompactWindows(yaml);
rejectMalformedProfilePatterns(yaml);
rejectUnknownKind(yaml);
rejectUnknownAuthMode(yaml);
@@ -2214,6 +2226,73 @@ public record FleetConfig(
}
}
/**
* Warn (never refuse to start) about a Claude Code profile whose auto-compaction flag and
* environment setting disagree.
*
* <p>Renamed from {@code rejectConflictingAutoCompactWindows} (fleetd #601 review, measured
* 2026-09-22): that method threw {@link IllegalStateException}, so {@link #load(Path)} refused
* to start on a config carrying this conflict. On this host, four profiles trip it, including
* the lead's own profile and the one every worker spawns on — so the throw is not a rare edge
* case. Under launchd, a throw inside {@code load()} is a restart loop, not an error an operator
* reads once, and the config that would fix it ({@code fleetd/fleetd.yaml}) is gitignored, so
* the cause is invisible on the host where it bites. A WARN gives the operator the same
* information — which profiles, and now both values, so they can fix it without reading the
* source — without ever taking the fleet down.
*
* <p>fleetd #618 measured which of the two inputs Claude Code actually follows when they
* disagree: the environment variable wins, so {@code autoCompactWindow} is inert on a profile
* that also sets the env var. This method only detects and reports the disagreement — it does
* not correct it — see {@link dev.ltms.fleet.launch.ClaudeCodeArguments} for the full measured
* precedence.
*
* <p>Equal values never warn: either input then produces the same session window, so there is
* nothing to reconcile.
*/
static void warnConflictingAutoCompactWindows(String yaml) {
Map<?, ?> raw;
try {
raw = YAML.readValue(yaml, Map.class);
} catch (IOException | IllegalArgumentException e) {
return;
}
if (raw == null || !(raw.get("profiles") instanceof Map<?, ?> profiles)) {
return;
}
List<String> names = new ArrayList<>();
List<String> detail = new ArrayList<>();
for (Map.Entry<?, ?> entry : profiles.entrySet()) {
if (!(entry.getValue() instanceof Map<?, ?> profile)
|| !(profile.get("autoCompactWindow") instanceof Number window)
|| !(profile.get("env") instanceof Map<?, ?> env)
|| !env.containsKey("CLAUDE_CODE_AUTO_COMPACT_WINDOW")) {
continue;
}
Object kind = profile.get("kind");
boolean claudeCode = kind == null || String.valueOf(kind).isBlank()
|| Profile.KIND_CLAUDE_CODE.equalsIgnoreCase(String.valueOf(kind));
Object envValue = env.get("CLAUDE_CODE_AUTO_COMPACT_WINDOW");
if (claudeCode && !String.valueOf(window).equals(String.valueOf(envValue))) {
String name = String.valueOf(entry.getKey());
names.add(name);
detail.add(name + " (autoCompactWindow=" + window
+ ", env.CLAUDE_CODE_AUTO_COMPACT_WINDOW=" + envValue + ")");
}
}
if (names.isEmpty()) {
return;
}
names.sort(String::compareTo);
detail.sort(String::compareTo);
log.warn("Claude Code profile(s) {} set disagreeing autoCompactWindow and env."
+ "CLAUDE_CODE_AUTO_COMPACT_WINDOW — the daemon starts anyway: {}. fleetd "
+ "#618 measured that CLAUDE_CODE_AUTO_COMPACT_WINDOW wins, so "
+ "autoCompactWindow is inert on these profiles. Set equal values on each "
+ "to resolve this — do not just delete the env var, since that LOWERS the "
+ "live window to autoCompactWindow's value rather than fixing anything.",
names, String.join(", ", detail));
}
/**
* Reject a profile whose {@code errorPattern} (fleetd #201 Unit 5) or {@code exhaustedPattern}
* (CB-578 stage A) is not a valid Java regex, naming the profile, the key, and the parser's own
@@ -0,0 +1,37 @@
package dev.ltms.fleet.launch;
import dev.ltms.fleet.config.FleetConfig;
import java.util.ArrayList;
import java.util.List;
/** Arguments shared by every fleetd path that starts Claude Code. */
public final class ClaudeCodeArguments {
private ClaudeCodeArguments() {
}
/**
* Append the configured Claude Code auto-compaction window when the profile opts in.
*
* <p>This flag and the environment variable {@code CLAUDE_CODE_AUTO_COMPACT_WINDOW} can
* disagree, and fleetd #618 measured which one Claude Code actually follows: the environment
* variable wins, ahead of this {@code --autocompact} flag, ahead of the settings file, ahead of
* clientdata, the experiment, and the model default. So when a profile sets both, the flag this
* method appends has NO effect — Claude Code reads {@code CLAUDE_CODE_AUTO_COMPACT_WINDOW}
* first and never consults the flag. {@link FleetConfig#load(java.nio.file.Path)} only WARNS
* when a Claude Code profile sets both to different values (see {@code
* FleetConfig.warnConflictingAutoCompactWindows}) — it does not stop the daemon from starting,
* and the launched session honours the env var, not this flag. Measured against Claude Code
* 2.1.278 (fleetd #618) — a later version could reorder this precedence.
*/
public static List<String> withAutoCompactWindow(List<String> argv, FleetConfig.Profile profile) {
if (profile.autoCompactWindow() == null) {
return argv;
}
List<String> withAutoCompact = new ArrayList<>(argv);
withAutoCompact.add("--autocompact");
withAutoCompact.add(String.valueOf(profile.autoCompactWindow()));
return withAutoCompact;
}
}
@@ -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);
}
}
@@ -8,6 +8,7 @@ import dev.ltms.fleet.herdr.PendingCloseMarker;
import dev.ltms.fleet.herdr.Tab;
import dev.ltms.fleet.herdr.Workspace;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.launch.ClaudeCodeArguments;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -358,7 +359,8 @@ public final class LeadLauncher {
}
/**
* The lead's argv: the profile's own command, the model pin, and the bridge MCP mount.
* The lead's argv: the profile's own command, the model and auto-compaction pins, and the bridge
* MCP mount.
*
* <p>No {@code --append-system-prompt}. That flag carries the worker reply charter, and a lead
* is not a worker — it reads its orchestration rules from the project's {@code CLAUDE.md} like
@@ -380,7 +382,7 @@ public final class LeadLauncher {
argv.add("--model");
argv.add(profile.model());
}
return argv;
return profile.isOpenCode() ? argv : ClaudeCodeArguments.withAutoCompactWindow(argv, profile);
}
/**
@@ -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());
@@ -8,6 +8,7 @@ import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.launch.ClaudeCodeArguments;
import dev.ltms.fleet.peer.Capability;
import dev.ltms.fleet.peer.PeerLauncher;
import org.slf4j.Logger;
@@ -298,7 +299,7 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
// has neither MCP nor a charter — session flags must be added into a list we own.
List<String> argv = mutableArgv(argvWithFleet(cfg, spec));
String agentSessionId = applySessionIdentity(argv, spec.sessionName(), spec.resumeSessionId());
return new Launch(workerEnv, argvWithAutoCompact(argvWithModel(argv, cfg), cfg), agentSessionId);
return new Launch(workerEnv, ClaudeCodeArguments.withAutoCompactWindow(argvWithModel(argv, cfg), cfg), agentSessionId);
}
/**
@@ -908,30 +909,6 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
return withModel;
}
/**
* Pin a bounded auto-compaction window on the command line via {@code --autocompact <tokens>},
* opt-in per profile (CB-634's sibling ticket: a member that runs out of context dies mid-turn
* and its {@code fleet_reply} — the whole point of the turn — is lost with it; opencode already
* forces {@code compaction.auto: true} unconditionally, CB-523, but Claude Code has no equivalent
* and runs at the backend's own default window).
*
* <p>Mirrors {@link #argvWithModel}: appended after it, so it survives the {@code ccs <profile>}
* wrapper the same way {@code --model} does, and outranks env/settings and the operator's own
* {@code argv}. Verified: {@code claude 2.1.241 --help} lists {@code --autocompact <auto|tokens>}
* (either the literal {@code auto}, or an integer 100k–1M) — {@link FleetConfig#load} rejects a
* configured value outside that band before this ever runs, so the flag Claude Code receives here
* is always in range.
*/
private static List<String> argvWithAutoCompact(List<String> argv, FleetConfig.Profile cfg) {
if (cfg.autoCompactWindow() == null) {
return argv;
}
List<String> withAutoCompact = mutableArgv(argv);
withAutoCompact.add("--autocompact");
withAutoCompact.add(String.valueOf(cfg.autoCompactWindow()));
return withAutoCompact;
}
// --- Agent-returning convenience spawns (used by callers/tests that want the herdr Agent) ---
/** Spawn a worker for the default profile in the resolved default cwd. */
@@ -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 ----------------------------------------------------------------------------------
@@ -203,62 +271,162 @@ public final class LeadHeartbeatLoop {
* does not evaluate the lead's idle state before the fleet has settled.
*/
public void start() {
log.info("idle-lead heartbeat: on — nudge lead after {}s idle (recheck {}ms, quiet cap {})",
TimeUnit.NANOSECONDS.toSeconds(idleAfterNanos), backoffMs, quietNudgeCap);
// fleetd #613: contextHighNudge added alongside the three settings already here — an
// operator otherwise cannot tell from the boot log whether the #609 handover notice is
// armed, and had to load the deployed jar's config to confirm it.
log.info("idle-lead heartbeat: on — nudge lead after {}s idle (recheck {}ms, quiet cap {}, "
+ "context-high nudge {})",
TimeUnit.NANOSECONDS.toSeconds(idleAfterNanos), backoffMs, quietNudgeCap,
contextHighNudge);
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());
}
}
@@ -0,0 +1,252 @@
package dev.ltms.fleet;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import dev.ltms.fleet.config.FleetConfig;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.slf4j.LoggerFactory;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #613: a {@code MemberRole} with no {@code fleet.<role>s:} pool falls back to
* <em>every</em> configured profile ({@code FleetConfig#candidateProfiles}), and one with no
* {@code fleet.charters.<role>:} entry runs with only the launcher's reply charter. Both are
* deliberate, legitimate states — neither is refused by {@code validateMembers()} — but both were
* silent at boot. On the host that opened this ticket, an unqualified {@code hunter} spawn silently
* widened to all 8 configured profiles and its resolved first choice was {@code local}, a profile
* every other pool on that same config gave weight 0 to.
*
* <p>{@link Fleetd#reportRoleFallbackGaps} must name every gapped role, and for a pool gap, the
* exact resolved first-choice profile — that number, not the pool size, is what actually surprised
* the operator. Mirrors {@link ExhaustedPatternGapReportTest}'s pattern: capture the real log via a
* {@link ListAppender} rather than asserting on the call site's source text.
*/
class RoleFallbackGapReportTest {
private static FleetConfig load(Path dir, String yaml) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml);
return FleetConfig.load(f);
}
/**
* The level this logger had before {@link #attach()} raised it, so {@link #detach} can put it
* back. {@code null} means "inherit from the parent" — the state this logger starts in.
*/
private static Level originalLevel;
/**
* {@code reportRoleFallbackGaps} logs at INFO, and {@code logback-test.xml} sets
* {@code dev.ltms.fleet} to WARN — so INFO events are dropped by the level check before any
* appender sees them. Raise the level for the duration of the test, exactly like {@code
* GitHostShapeReportTest#attach}.
*/
private static ListAppender<ILoggingEvent> attach() {
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
originalLevel = logger.getLevel();
logger.setLevel(Level.INFO);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
return appender;
}
private static void detach(ListAppender<ILoggingEvent> appender) {
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
logger.detachAppender(appender);
logger.setLevel(originalLevel);
}
private static List<String> infoMessages(ListAppender<ILoggingEvent> appender) {
return appender.list.stream()
.filter(e -> e.getLevel() == Level.INFO)
.map(ILoggingEvent::getFormattedMessage)
.toList();
}
/**
* Reproduces the shape measured in the ticket: {@code dev}, {@code reviewer} and
* {@code architect} each have a pool and a charter; {@code hunter} has neither. The pool-gap
* line must name {@code hunter}, the profile count (3), and the resolved first choice
* ({@code local}, the first profile in definition order) — and must not name the three healthy
* roles. The charter-gap line must separately name only {@code hunter}.
*/
@Test
void hunterWithNoPoolOrCharterIsNamedWithItsResolvedFirstChoice(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, """
profiles:
local:
baseUrl: http://gx00.gw:8000
sonnet:
baseUrl: https://llm.ltms.dev/v1
terra:
baseUrl: https://llm.ltms.dev/v1
fleet:
developers:
a:
profile: sonnet
reviewers:
b:
profile: terra
architects:
c:
profile: sonnet
charters:
dev: "dev charter text"
reviewer: "reviewer charter text"
architect: "architect charter text"
""");
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportRoleFallbackGaps(cfg);
} finally {
detach(appender);
}
List<String> infos = infoMessages(appender);
String poolLine = infos.stream()
.filter(m -> m.contains("no fleet.<role>s: pool"))
.findFirst()
.orElseThrow(() -> new AssertionError("expected a pool-gap INFO line: " + infos));
assertTrue(poolLine.contains("hunter"), poolLine);
assertTrue(poolLine.contains("3"), "must name the profile count the gap opens onto: " + poolLine);
assertTrue(poolLine.contains("'local'"),
"must name the resolved first-choice profile, the number that actually surprised "
+ "the operator: " + poolLine);
for (String healthy : List.of("dev", "reviewer", "architect")) {
assertFalse(poolLine.contains(healthy),
"pool-gap line must not name a role that has a pool: " + poolLine);
}
String charterLine = infos.stream()
.filter(m -> m.contains("no fleet.charters: entry"))
.findFirst()
.orElseThrow(() -> new AssertionError("expected a charter-gap INFO line: " + infos));
assertTrue(charterLine.contains("hunter"), charterLine);
for (String healthy : List.of("dev", "reviewer", "architect")) {
assertFalse(charterLine.contains(healthy),
"charter-gap line must not name a role that has a charter: " + charterLine);
}
}
/**
* The resolved first choice must be genuinely computed from definition order, not hardcoded —
* reordering {@code profiles:} so a different entry comes first changes the reported choice.
*/
@Test
void theResolvedFirstChoiceFollowsProfileDefinitionOrder(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, """
profiles:
sonnet:
baseUrl: https://llm.ltms.dev/v1
local:
baseUrl: http://gx00.gw:8000
fleet:
developers:
a:
profile: sonnet
charters:
dev: "dev charter text"
""");
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportRoleFallbackGaps(cfg);
} finally {
detach(appender);
}
String poolLine = infoMessages(appender).stream()
.filter(m -> m.contains("no fleet.<role>s: pool"))
.findFirst()
.orElseThrow();
// hunter, reviewer and architect all lack a pool here; each falls back to the full 2-profile
// set and the first-choice is 'sonnet' because it is first in profiles: definition order.
assertTrue(poolLine.contains("'sonnet'"), poolLine);
assertFalse(poolLine.contains("'local'"), poolLine);
}
/** A config with a pool and a charter for every role produces no role-fallback log at all. */
@Test
void everyRoleWithAPoolAndACharterProducesNoLogAtAll(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, """
profiles:
sonnet:
baseUrl: https://llm.ltms.dev/v1
fleet:
developers:
a:
profile: sonnet
reviewers:
b:
profile: sonnet
hunters:
c:
profile: sonnet
architects:
d:
profile: sonnet
charters:
dev: "dev charter text"
reviewer: "reviewer charter text"
hunter: "hunter charter text"
architect: "architect charter text"
""");
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportRoleFallbackGaps(cfg);
} finally {
detach(appender);
}
assertTrue(appender.list.isEmpty(),
"a config with no gaps must not print a per-role block: " + infoMessages(appender));
}
/**
* A config with only {@code profiles:} and no {@code fleet:} block at all must still be
* reported (every role is gapped, both pool and charter) rather than throwing — this is the
* exact shape {@code candidateProfiles}' fallback exists to keep starting.
*/
@Test
void aConfigWithNoFleetBlockAtAllReportsEveryRoleGapped(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, """
profiles:
sonnet:
baseUrl: https://llm.ltms.dev/v1
""");
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportRoleFallbackGaps(cfg);
} finally {
detach(appender);
}
List<String> infos = infoMessages(appender);
String poolLine = infos.stream()
.filter(m -> m.contains("no fleet.<role>s: pool"))
.findFirst()
.orElseThrow(() -> new AssertionError("expected a pool-gap INFO line: " + infos));
String charterLine = infos.stream()
.filter(m -> m.contains("no fleet.charters: entry"))
.findFirst()
.orElseThrow(() -> new AssertionError("expected a charter-gap INFO line: " + infos));
for (String role : List.of("dev", "hunter", "reviewer", "architect")) {
assertTrue(poolLine.contains(role), poolLine);
assertTrue(charterLine.contains(role), charterLine);
}
}
}
@@ -1,8 +1,11 @@
package dev.ltms.fleet.config;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.spi.ILoggingEvent;
import dev.ltms.fleet.auth.MemberRegistry;
import dev.ltms.fleet.msg.LeadMailbox;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.testing.CapturedLog;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
@@ -93,6 +96,72 @@ class FleetConfigTest {
"unset means off — today's behaviour, unchanged");
}
/**
* fleetd #601 review (measured 2026-09-22): this guard used to throw {@link
* IllegalStateException} and refuse to start. On a host with the conflict configured, that
* turned into a launchd restart loop with no readable cause, since {@code fleetd.yaml} is
* gitignored. It must now WARN and let the daemon start, and the warning must carry both values
* so an operator can fix the config without reading the source. Pinning the log line itself (via
* {@link CapturedLog}) rather than a snippet of production source text — the latter is the
* anti-pattern this repo avoids; the former is the actual observable behaviour a reader (or an
* alert on the log) depends on.
*/
@Test
void aClaudeProfileWithConflictingAutoCompactFlagAndEnvironmentWindowLoadsAndWarnsWithBothValues(
@TempDir Path dir) throws Exception {
Path f = dir.resolve("conflicting-auto-compact-window.yaml");
Files.writeString(f, """
profiles:
claude-profile:
autoCompactWindow: 250000
env:
CLAUDE_CODE_AUTO_COMPACT_WINDOW: "300000"
""");
FleetConfig cfg;
List<String> warnings;
try (CapturedLog log = CapturedLog.at(FleetConfig.class, Level.WARN)) {
cfg = FleetConfig.load(f);
warnings = log.events().stream().map(ILoggingEvent::getFormattedMessage).toList();
}
assertEquals(250_000, cfg.profiles().get("claude-profile").autoCompactWindow(),
"the disagreement is reported, not corrected — the flag value still loads as-is");
assertEquals(1, warnings.size(), "exactly one warning for the one conflicting profile: " + warnings);
String warning = warnings.get(0);
assertTrue(warning.contains("claude-profile"), "names the offending profile: " + warning);
assertTrue(warning.contains("autoCompactWindow=250000"), "names the flag value: " + warning);
assertTrue(warning.contains("CLAUDE_CODE_AUTO_COMPACT_WINDOW=300000"), "names the env value: " + warning);
}
/**
* The negative probe paired with the test above (per fleetd #601 review): a warning that fires
* on every load and a warning that never fires read the same from a single test, so both must be
* checked. No conflict here — the flag and the env value agree — so no warning should be logged.
*/
@Test
void aClaudeProfileWithEqualAutoCompactFlagAndEnvironmentWindowLoadsWithNoWarning(@TempDir Path dir)
throws Exception {
Path f = dir.resolve("equal-auto-compact-window.yaml");
Files.writeString(f, """
profiles:
claude-profile:
autoCompactWindow: 250000
env:
CLAUDE_CODE_AUTO_COMPACT_WINDOW: "250000"
""");
FleetConfig cfg;
List<ILoggingEvent> events;
try (CapturedLog log = CapturedLog.at(FleetConfig.class, Level.WARN)) {
cfg = FleetConfig.load(f);
events = log.events();
}
assertEquals(250_000, cfg.profiles().get("claude-profile").autoCompactWindow());
assertTrue(events.isEmpty(), "equal values must not warn: " + events);
}
// ── fleetd #201 Unit 5: errorPattern ────────────────────────────────────────────────────────
@Test
@@ -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());
@@ -32,13 +32,26 @@ class LeadLauncherTest {
Map.of("CLAUDE_CODE_AUTO_COMPACT_WINDOW", "300000"), null, null, true, null);
}
private static FleetConfig.Profile profileWithAutoCompactWindow(String kind) {
return new FleetConfig.Profile(
"opus", null, "claude-opus-5", null, "FLEETD_WORKER_TOKEN",
List.of("ccs", "ltms"), "tab", "fleet", null,
"http://127.0.0.1:8765/mcp", null, null,
null, null, kind, Map.of(), null, null, true, null, null,
null, null, null, 250_000);
}
private static FleetConfig configWith(FleetConfig.Leader lead) {
return configWith(lead, opusProfile());
}
private static FleetConfig configWith(FleetConfig.Leader lead, FleetConfig.Profile profile) {
Map<String, FleetConfig.Leader> leaders = new LinkedHashMap<>();
leaders.put("opus", lead);
FleetConfig.Fleet fleet =
new FleetConfig.Fleet(leaders, Map.of(), Map.of(), Map.of(), null);
return new FleetConfig(
null, null, Map.of("opus", opusProfile()), null, null, null, null, null,
null, null, Map.of("opus", profile), null, null, null, null, null,
null, null, fleet, null, "fixed", null).withDefaults();
}
@@ -330,6 +343,36 @@ class LeadLauncherTest {
"--model is appended last so it outranks the ccs wrapper (CB-533)");
}
@Test
void aClaudeLeadPassesItsConfiguredAutoCompactWindowToClaudeCode() {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile profile = profileWithAutoCompactWindow("claude-code");
launcher(herdr, configWith(lead("opus", "lead: opus", 1), profile)).ensureLeads();
List<String> args = startedArgs(herdr);
assertEquals("250000", args.get(args.indexOf("--autocompact") + 1));
}
@Test
void aClaudeLeadWithNoAutoCompactWindowGetsNoAutoCompactFlag() {
FakeHerdr herdr = new FakeHerdr();
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads();
assertFalse(startedArgs(herdr).contains("--autocompact"));
}
@Test
void aNonClaudeLeadDoesNotGetAnAutoCompactFlag() {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile profile = profileWithAutoCompactWindow("opencode");
launcher(herdr, configWith(lead("opus", "lead: opus", 1), profile)).ensureLeads();
assertFalse(startedArgs(herdr).contains("--autocompact"));
}
/** A lead runs on the operator's subscription. Nothing may move it off. */
@Test
void theLeadEnvCarriesNoAnthropicBinding() {
@@ -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,27 +1,46 @@
package dev.ltms.fleet.msg;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
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;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
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 +75,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 +94,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,17 +106,71 @@ 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");
}
// ── fleetd #613: the boot line names all four heartbeat settings ──────────────────────────
/**
* fleetd #613: {@link LeadHeartbeatLoop#start()}'s boot line named only 3 of the 4 constructor
* settings — {@code contextHighNudge} (fleetd #609) was missing, so an operator could not tell
* from the log whether the handover notice was armed. Captures the real log via a
* {@link ListAppender}, raising the logger's level past {@code logback-test.xml}'s
* {@code dev.ltms.fleet -> WARN} override for the duration of the call — the same seam {@code
* GitHostShapeReportTest#attach} uses for its own INFO-level boot line.
*/
private static String heartbeatBootLine(boolean contextHighNudge, ScheduledExecutorService scheduler) {
Logger logger = (Logger) LoggerFactory.getLogger(LeadHeartbeatLoop.class);
Level originalLevel = logger.getLevel();
logger.setLevel(Level.INFO);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
LeadHeartbeatLoop l = new LeadHeartbeatLoop(
new PrimaryRegistry("term_lead"), null /*agents*/, null /*inbox*/, List::of,
null /*pushLoop*/, scheduler, () -> 0L, IDLE_AFTER_NANOS, 1_000L, 3, null,
LeadHeartbeatLoop.LeadContextSource.none(), contextHighNudge);
l.start();
} finally {
logger.detachAppender(appender);
logger.setLevel(originalLevel);
}
return appender.list.stream()
.filter(e -> e.getLevel() == Level.INFO)
.map(ILoggingEvent::getFormattedMessage)
.filter(m -> m.startsWith("idle-lead heartbeat:"))
.findFirst()
.orElseThrow(() -> new AssertionError("expected the heartbeat boot line to be logged"));
}
@Test
void theBootLineNamesContextHighNudgeWhenArmed() {
String line = heartbeatBootLine(true, scheduler);
assertTrue(line.contains("300s idle"), line);
assertTrue(line.contains("recheck 1000ms"), line);
assertTrue(line.contains("quiet cap 3"), line);
assertTrue(line.contains("context-high nudge true"),
"the boot line must name the 4th setting, contextHighNudge, alongside the other "
+ "three: " + line);
}
@Test
void theBootLineNamesContextHighNudgeWhenOff() {
String line = heartbeatBootLine(false, scheduler);
assertTrue(line.contains("context-high nudge false"), line);
}
// ── (b) an idle lead within the quiet period is not yet injected ───────────────────────────
@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 +179,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 +190,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 +199,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 +220,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 +231,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 +243,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 +301,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() {
}
}
}