Merge #483: fleetd #480 Unit A — lead rollover core (config + executor)
CI / contract (push) Successful in 1m13s
CI / build (push) Successful in 1m37s

Two corrections were applied to the original unit before this merge, and the PR
description above still describes the pre-correction shape:

1. confirm() no longer rolls inline. It validates every gate, then hands a one-shot
   continuation to continuationRunner and returns. confirm() is called BY the lead FROM
   its own turn, so the pane is WORKING and cannot report injectable until confirm()
   returns; the old code sent /clear first and then timed out waiting, which destroyed
   the lead's context and started no fresh session. The continuation waits for the
   calling turn to settle FIRST (turnSettleSeconds, a new knob), and if that wait times
   out it sends no /clear at all.
2. The pane to roll comes from the caller's terminal, not PrimaryRegistry. A single-slot
   lookup let lead X's confirm() clear lead Y's pane. open() records the caller's
   terminal; confirm() refuses with NOT_YOUR_ROLLOVER on a mismatch.

Verified by the lead before merge:
- CI run 1722 green on head a942622.
- Mutation battery on the two safety-critical branches, in a scratch worktree at a942622.
  Controls measured against the ORIGINAL first; each mutant proven applied with two
  unrelated proofs using different search strings.
  M1 `if (!turnSettled)` -> `if (false)`: turnThatNeverSettlesSendsNoClearAtAll FAILS
     (LeadRolloverTest:247, expected 0 sends but was 1).
  M2 ownership check -> `if (false)`: aDifferentLeadTerminalCannotConfirmAnotherLeadsRollover
     FAILS (LeadRolloverTest:266). Unmutated harness-proof run: exit 0 with both guard
     WARN lines logged, so the branches really are exercised.

Nothing calls LeadRollover yet. The fleet_handover MCP tool is Unit C.
This commit was merged in pull request #483.
This commit is contained in:
2026-09-11 01:51:20 +02:00
11 changed files with 1090 additions and 18 deletions
+36
View File
@@ -88,6 +88,42 @@ bind:
# backoffMs: 60000
# quietNudgeCap: 3
# 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
# own pane and bootstrap a fresh session against that file.
#
# Opt-in on purpose — it clears the lead's own pane on request, so upgrading the daemon must never
# acquire that ability for you. Absent block = feature off, and nothing is constructed at all. Even
# once present, nothing but an explicit confirm() call — one that passes every check — can ever
# cause a /clear: there is no recurring timer, heartbeat or scheduler anywhere in this feature that
# fires one on its own initiative. confirm() itself is called FROM the calling lead's own turn, so
# it cannot clear the pane inline (that pane is still WORKING); instead it schedules a one-shot
# continuation that waits for the SAME confirm() call's turn to end, then does the actual work. See
# dev.ltms.fleet.lead.LeadRollover's class javadoc for the exact order (fleetd #480 correction).
#
# handoverPath: REQUIRED when this block is present — where the handover file a fresh lead session
# reads must live. No default (an operator-specific path); a present block with no
# handoverPath refuses to start.
# requireOperatorConfirm: true # default true — confirm() refuses unless the caller also passes
# # operatorConfirmed: true
# maxDocAgeSeconds: 3600 # default 3600 — refuse a handover file older than this
# turnSettleSeconds: 20 # default 20 — how long the deferred roll waits for the CALLING
# # lead's own turn to end (its pane to report injectable again)
# # before sending /clear at all. If this elapses, /clear is NEVER
# # sent — a lead that never goes idle is still doing real work.
# clearSettleSeconds: 20 # default 20 — how long to wait for the pane to become injectable
# # again AFTER /clear before giving up (never sends bootstrapText
# # if this elapses). A separate, second wait from turnSettleSeconds.
# bootstrapText: "..." # default names handoverPath — sent to the lead once its pane
# # settles after /clear
# leadRollover:
# handoverPath: /path/to/handover.md
# requireOperatorConfirm: true
# maxDocAgeSeconds: 3600
# turnSettleSeconds: 20
# clearSettleSeconds: 20
# bootstrapText: "Fresh lead session: read the handover file and carry on."
# Fleet health detection is dormant unless enabled (CB-573). It reads one whole-fleet agent list
# per tick.
# intervalSeconds → how often a tick runs (default 30). ENFORCED floor of 15: the code computes
@@ -10,6 +10,7 @@ import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.herdr.HerdrRouter;
import dev.ltms.fleet.herdr.LeadTabScanner;
import dev.ltms.fleet.lead.LeadLauncher;
import dev.ltms.fleet.lead.LeadRollover;
import dev.ltms.fleet.herdr.PaneLocator;
import dev.ltms.fleet.herdr.UnixSocketHerdrClient;
import dev.ltms.fleet.herdr.WorkspaceControl;
@@ -590,6 +591,16 @@ public final class Fleetd {
heartbeat = null;
heartbeatScheduler.shutdownNow();
}
// fleetd #480: lead rollover. Opt-in; absent `leadRollover:` this is never constructed, so
// an upgraded daemon cannot silently acquire the ability to clear the lead's own pane.
// Unlike heartbeat above, this has no recurring scheduler of its own — nothing but an
// explicit confirm() call (wired to an MCP tool by a later ticket; nothing calls it yet)
// that passes every gate can ever schedule a roll. It does not take primaryRegistry: the
// lead terminal to roll comes from the caller of open()/confirm() (resolved by the MCP
// layer from the connection, the same way auth/CallerResolver#resolve builds a
// Principal.leader(...)), never from a single-slot lookup — see LeadRollover's class
// javadoc, fleetd #480 correction 2.
LeadRollover leadRollover = leadRollover(cfg, router.leadAgents(), config);
MessageService messages = new MessageService(router, injector, rendezvous, replyInbox,
pushLoop, metrics);
@@ -1053,6 +1064,44 @@ public final class Fleetd {
.orElse(null);
}
/**
* fleetd #480: construct the {@link LeadRollover} executor only when {@code leadRollover:} is
* present at startup — the same presence gate {@code leadHeartbeat:} uses just above this
* call site in {@code main}. Extracted to its own factory, the same reason
* {@link #worktreeBranchLookup} and {@link #exhaustionSink} are: a unit test can call this
* directly with a fabricated {@link FleetConfig} (absent block ⇒ {@code null}, present block ⇒
* constructed) without needing the whole of {@code main}, and a source-text test on the real
* call site proves {@code main} still calls this factory rather than inlining a copy that could
* silently diverge.
*
* <p>The returned object reads every {@code leadRollover:} field fresh on each {@code open()}/
* {@code confirm()} call through {@code () -> config.get().leadRollover()} — see {@code
* ConfigRef}'s class doc Hot bullet and {@link FleetConfig.LeadRollover}'s javadoc for why that
* makes the block's fields HOT despite this presence gate being evaluated once, at startup.
*
* <p>Deliberately does NOT take {@link PrimaryRegistry}: fleetd #480 correction 2 found that a
* single-slot lookup lets one lead's {@code confirm()} clear a DIFFERENT lead's pane on a
* daemon with more than one labelled lead tab. The lead terminal to roll instead comes from
* whoever calls {@code open()}/{@code confirm()} — the later MCP-tool unit must resolve it from
* the connection and pass it in, never take it as a request field. See {@link LeadRollover}'s
* class javadoc.
*
* @param cfg the startup config snapshot — read ONCE here, only to decide whether to
* construct the object at all, exactly like {@code cfg.leadHeartbeat()}
* @param leadAgents the {@link AgentControl} instance that reaches the LEAD's pane (not
* {@code memberAgents}), normally {@code router.leadAgents()}
* @param config the live {@link ConfigRef}, captured only inside the returned supplier —
* never dereferenced here
* @return a constructed {@link LeadRollover}, or {@code null} when {@code leadRollover:} is
* absent from the startup config
*/
static LeadRollover leadRollover(FleetConfig cfg, AgentControl leadAgents, ConfigRef config) {
if (cfg.leadRollover() == null) {
return null;
}
return new LeadRollover(leadAgents, () -> config.get().leadRollover());
}
/**
* fleetd #176: per-profile factory for {@link FleetMcp.LeadSeatSource} — how many seats a
* profile's own live LEAD session(s) hold on the same Claude subscription.
@@ -62,7 +62,22 @@ import java.util.function.Supplier;
* both read, so a reload that arms or disarms a profile's usage-limit detection takes effect
* on the next check with no restart. {@code errorPattern}, {@code exhaustedPattern}'s sibling
* key for backend-error (not usage-limit) classification, was deliberately left OUT of this
* fleetd #446 change and stays deferred below — the ticket scoped it out explicitly.</li>
* fleetd #446 change and stays deferred below — the ticket scoped it out explicitly.
* {@code leadRollover:} (fleetd #480) joined this class whole, the same shape as
* {@code models:} above: {@code dev.ltms.fleet.lead.LeadRollover} holds a
* {@code Supplier<FleetConfig.LeadRollover>} (the same {@code () -> config.get().x()} shape)
* and reads {@code handoverPath}/{@code requireOperatorConfirm}/{@code maxDocAgeSeconds}/
* {@code turnSettleSeconds}/{@code clearSettleSeconds}/{@code bootstrapText} fresh on every
* {@code open()}/{@code confirm()} call (and on the deferred post-{@code confirm()}
* continuation fleetd #480's correction added — see {@code LeadRollover}'s class doc) rather
* than capturing them into fields at construction — unlike its closest
* structural cousin {@code leadHeartbeat:}, whose {@code LeadHeartbeatLoop} bakes
* {@code idleAfterNanos}/{@code backoffMs}/{@code quietNudgeCap} into final fields. The one
* restart-only edge is structural, not a stale value: {@code Fleetd.java} decides whether to
* construct the {@code LeadRollover} object at all off the startup snapshot (the same
* presence gate {@code leadHeartbeat:} uses), so a block ADDED where it was absent at boot
* needs a restart before anything exists to call — the same fact already true of adding a
* brand-new {@code profiles:} entry.</li>
* <li><strong>Deferred</strong> — accepted into the new snapshot, but the wiring built at startup
* keeps the old value until a restart: {@code lifecycle:}, {@code leadHeartbeat:},
* {@code idleSleepGuard:} ({@code Fleetd.java} reads it once, at startup, to decide whether
@@ -166,17 +181,20 @@ import java.util.function.Supplier;
*
* <p><strong>The denominator, measured on 2026-09-04 (fleetd #330; recounted for fleetd #333);
* recounted again for fleetd #362, again after {@code idleSleepGuard:} was added, again after
* {@code models:} was added as deferred, and again for fleetd #422, which moved {@code models:}
* from deferred to hot-excluded once its on/off half was read live everywhere.</strong>
* {@code FleetConfig} has 25 top-level record components: 5 cold, 13 deferred, 3 split, 4
* hot-excluded. Four of them are named nowhere in this file, and the reason is the same for all
* four: {@code placement}, {@code memberCredentials}, {@code memberLoginShell} and {@code models}
* are <strong>hot</strong> and correctly absent — all four are read live off {@code config.get()}
* (placement through the {@code CompositePeerLauncher} supplier the Hot bullet names;
* {@code memberCredentials}/{@code memberLoginShell} at spawn time, {@code Fleetd.java:198, 205, 729}
* and {@code HerdrPeerLauncher#configuredMemberLoginShell}; {@code models} the same way, through the
* Hot bullet's {@code models:} paragraph), so a reload takes effect on the next spawn (or, for
* {@code models}, the next reported status) with no entry needed here.
* {@code models:} was added as deferred, again for fleetd #422, which moved {@code models:}
* from deferred to hot-excluded once its on/off half was read live everywhere, and again after
* {@code leadRollover:} was added (fleetd #480).</strong>
* {@code FleetConfig} has 26 top-level record components: 5 cold, 13 deferred, 3 split, 5
* hot-excluded. Five of them are named nowhere in this file, and the reason is the same for all
* five: {@code placement}, {@code memberCredentials}, {@code memberLoginShell}, {@code models} and
* {@code leadRollover} are <strong>hot</strong> and correctly absent — all five are read live off
* {@code config.get()} (placement through the {@code CompositePeerLauncher} supplier the Hot bullet
* names; {@code memberCredentials}/{@code memberLoginShell} at spawn time, {@code Fleetd.java:198,
* 205, 729} and {@code HerdrPeerLauncher#configuredMemberLoginShell}; {@code models} the same way,
* through the Hot bullet's {@code models:} paragraph; {@code leadRollover} through the Hot bullet's
* {@code leadRollover:} paragraph), so a reload takes effect on the next spawn (or, for
* {@code models}, the next reported status; for {@code leadRollover}, the next {@code open()}/
* {@code confirm()} call) with no entry needed here.
* {@code health} and {@code coordinator} used to be a third kind — <strong>undecided</strong>, not
* hot — until fleetd #330 added the <strong>split</strong> class above and gave them a home. A
* reload touching either used to report a bare "config reloaded", which under-claimed; now it names
@@ -139,6 +139,9 @@ import java.util.regex.PatternSyntaxException;
* fleetd #422 added the separate on/off question — whether a configured model
* may be spawned onto RIGHT NOW ({@link Models.ModelEntry#enabled}) — enforced
* live at spawn by {@code CompositePeerLauncher}, not here. See {@link Models}.
* @param leadRollover opt-in lead rollover (fleetd #480): {@code null} ⇒ off, and no
* {@code dev.ltms.fleet.lead.LeadRollover} is constructed at all — an upgraded
* daemon never clears a lead's pane on its own initiative. See {@link LeadRollover}.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record FleetConfig(
@@ -166,7 +169,23 @@ public record FleetConfig(
String memberLoginShell,
String memberSkills,
IdleSleepGuard idleSleepGuard,
Models models) {
Models models,
LeadRollover leadRollover) {
/** Back-compat form before the {@code leadRollover:} block was added. */
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
Guard guard, String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
ConfigReload configReload, Integer quarantineCooldownSeconds,
MemberCredentials memberCredentials, Coordinator coordinator, String worktreeGroup,
String memberLoginShell, String memberSkills, IdleSleepGuard idleSleepGuard,
Models models) {
this(bind, herdrSocket, memberHerdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, worktreeGroup,
memberLoginShell, memberSkills, idleSleepGuard, models, null);
}
/** Back-compat form before the {@code models:} block was added. */
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
@@ -1310,6 +1329,79 @@ public record FleetConfig(
}
}
/**
* Opt-in lead rollover (fleetd #480): a lead that decides it is ready to be replaced writes a
* handover file, then asks fleetd to clear its own pane and bootstrap a fresh session against
* that file. Config + a pure decision/verification layer only — see
* {@code dev.ltms.fleet.lead.LeadRollover} for the executor this block feeds, and
* {@code fleet_*} tool wiring is a later ticket.
*
* <p>Deliberately opt-in ({@code null} ⇒ off, exactly like {@code leadHeartbeat:}): absent this
* block, {@code Fleetd.java} never constructs a {@code LeadRollover} object at all, so an
* upgraded daemon cannot silently acquire the ability to clear the lead's own pane. Even once
* present, nothing but an explicit {@code confirm()} call can ever roll a pane — there is no
* timer, heartbeat or timeout anywhere in this feature that fires one on its own; see that
* class's javadoc.
*
* <p><strong>Hot, not deferred</strong> (see {@code ConfigRef}'s class doc): every field below
* is read live, through a {@code Supplier<LeadRollover>} the same {@code () -> config.get().x()}
* shape {@code fleet}/{@code placement}/{@code models} already use, so an edit to any of the
* six fields takes effect on the very next {@code open()}/{@code confirm()} call once the block
* has been present since startup — nothing here is captured into a frozen field the way {@code
* leadHeartbeat}'s {@code idleAfterNanos}/{@code backoffMs}/{@code quietNudgeCap} are. The one
* restart-only edge left is structural, not a stale value: {@code Fleetd.java} decides whether
* to construct the {@code LeadRollover} object at all off the startup snapshot, the same gate
* {@code leadHeartbeat:} uses, so a block ADDED where it was absent at boot needs a restart
* before anything exists to call — the same fact already true of adding a brand-new
* {@code profiles:} entry.
*
* <p><strong>{@code turnSettleSeconds} (fleetd #480 correction):</strong> {@code confirm()} is
* called FROM the calling lead's own turn, so its pane is still {@code WORKING} the instant
* {@code confirm()} validates every gate and schedules the roll. {@code
* dev.ltms.fleet.lead.LeadRollover}'s deferred continuation waits up to this many seconds for
* that SAME pane to report an injectable state again — i.e. for the calling turn to actually
* end — before it sends {@code /clear} at all. If that wait times out, no {@code /clear} is
* ever sent: a lead that never goes idle is still doing real work, and clearing it would
* destroy live context. This is a separate wait from {@code clearSettleSeconds} below, which
* bounds the SECOND wait, for the pane to re-settle AFTER {@code /clear} has already gone out.
*
* @param handoverPath required when this block is present — where the handover file a fresh
* lead session reads must live. There is no sane non-null default for an
* operator-specific path, so a present block with a {@code null}/blank
* {@code handoverPath} is refused at config load; see
* {@link #validateLeadRollover()}.
* @param requireOperatorConfirm default {@code true} — {@code confirm()} refuses unless the
* caller also passes {@code operatorConfirmed: true}. Set {@code false} to
* let the three handover-file checks alone gate the roll.
* @param maxDocAgeSeconds default 3600 — refuse a handover file whose modified time is older
* than this many seconds, so a stale leftover from an earlier rollover
* attempt can never be mistaken for a fresh one.
* @param turnSettleSeconds default 20 — bound on how long the deferred roll waits for the
* CALLING lead's own turn to end (its pane to report injectable again)
* before sending {@code /clear} at all. See the paragraph above.
* @param clearSettleSeconds default 20 — bound on how long to wait for the lead's pane to
* report an injectable state again after {@code /clear} before giving up. A
* roll that times out here never sends {@code bootstrapText}.
* @param bootstrapText default a sentence naming {@code handoverPath} — sent to the lead's pane
* once it settles after {@code /clear}, telling the fresh session where to
* read the handover and carry on.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record LeadRollover(String handoverPath, Boolean requireOperatorConfirm,
Integer maxDocAgeSeconds, Integer turnSettleSeconds,
Integer clearSettleSeconds, String bootstrapText) {
public LeadRollover {
requireOperatorConfirm = requireOperatorConfirm == null || requireOperatorConfirm;
maxDocAgeSeconds = (maxDocAgeSeconds == null || maxDocAgeSeconds <= 0) ? 3600 : maxDocAgeSeconds;
turnSettleSeconds = (turnSettleSeconds == null || turnSettleSeconds <= 0) ? 20 : turnSettleSeconds;
clearSettleSeconds = (clearSettleSeconds == null || clearSettleSeconds <= 0) ? 20 : clearSettleSeconds;
bootstrapText = (bootstrapText == null || bootstrapText.isBlank())
? "Fresh lead session: read the handover file at " + handoverPath
+ " and carry on from there."
: bootstrapText;
}
}
/**
* Watch {@code fleetd.yaml} and re-read it when it changes (CB-559).
*
@@ -1725,7 +1817,7 @@ public record FleetConfig(
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds",
"memberCredentials", "coordinator", "worktreeGroup", "memberLoginShell", "memberSkills",
"idleSleepGuard", "models");
"idleSleepGuard", "models", "leadRollover");
/** Load and validate config from {@code path}. */
public static FleetConfig load(Path path) {
@@ -2415,10 +2507,14 @@ public record FleetConfig(
// empty Models would be a no-op for validateModels() either way, since an empty allow-list
// already means "check nothing", so there is nothing to gain and one more null check to
// avoid by leaving it exactly as configured.
// leadRollover is left as-is, like leadHeartbeat above: null is "off", and LeadRollover's
// own compact constructor defaults the fields of a block that IS present. Defaulting it
// here would construct a LeadRollover object (via Fleetd.java's presence gate) for every
// config that never mentioned it.
return new FleetConfig(b, herdrSocket, memberHerdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
quarantineCooldown, mc, coordinator, worktreeGroup, memberLoginShell, memberSkills,
idleSleepGuard, models);
idleSleepGuard, models, leadRollover);
}
/**
@@ -2501,6 +2597,26 @@ public record FleetConfig(
+ "lead tabs cannot be confused.");
}
/**
* Reject a present {@code leadRollover:} block with no (or a blank) {@code handoverPath}
* (fleetd #480). There is no sane non-null default for an operator-specific file path, unlike
* every other field on {@link LeadRollover}, which {@link LeadRollover}'s own compact
* constructor already defaults — so this is the one field that must be refused at load rather
* than silently defaulted to something that would never match a real handover file.
*
* @throws IllegalStateException when {@code leadRollover} is present but {@code handoverPath}
* is {@code null} or blank
*/
public void validateLeadRollover() {
if (leadRollover == null) {
return;
}
if (leadRollover.handoverPath() == null || leadRollover.handoverPath().isBlank()) {
throw new IllegalStateException("refusing to start: leadRollover.handoverPath is "
+ "required when leadRollover: is present.");
}
}
/** Case-insensitive prefix test that tolerates a null/blank label. */
private static boolean startsWithIgnoreCase(String label, String prefix) {
if (label == null || prefix == null || prefix.isBlank()) {
@@ -0,0 +1,379 @@
package dev.ltms.fleet.lead;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.io.UncheckedIOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;
import java.util.function.LongSupplier;
import java.util.function.Supplier;
/**
* fleetd #480: replace a lead session that has decided it is ready to be rolled over, without an
* operator doing it by hand. A lead writes a handover file, calls {@link #open}, and then — once
* every gate ({@link #confirm}'s own checks) has passed — a deferred, single-shot continuation
* clears the lead's own pane and bootstraps a fresh session against that file.
*
* <p>This is the executor only. Nothing in this ticket wires an MCP tool onto {@link #open}/
* {@link #confirm}/{@link #cancel} — that is a separate, later unit; until it lands, nothing calls
* this class at all.
*
* <p><strong>{@code confirm()} cannot roll inline — a fleetd #480 correction.</strong> The first
* version of this class called {@code agents.send(lead, "/clear")} directly from inside {@code
* confirm()}, then polled for the pane to become injectable again. That is wrong, because {@code
* confirm()} is called BY the lead, FROM the lead's own turn: the lead's pane is {@code WORKING}
* for the whole duration of that call and cannot possibly report injectable until {@code confirm()}
* itself returns. The poll always timed out — but only after the {@code /clear} had already been
* sent and queued in the pane, where it fired the instant the turn ended anyway. The result was the
* worst outcome this feature can produce: a silently destroyed lead context with no fresh session
* ever started, and a refusal return value that claimed nothing had happened.
*
* <p>The fix: {@link #confirm} validates every gate, then does no I/O against the lead's own pane
* at all — it only records that the request is approved and hands a one-shot continuation to
* {@code continuationRunner} before returning. That continuation is what actually touches the pane,
* once the calling turn has ended, in this order:
* <ol>
* <li>wait for the lead's own pane to report an injectable state — i.e. wait for the very
* {@code confirm()} call that approved this roll to finish its turn — bounded by
* {@code turnSettleSeconds}. <strong>If this never happens, nothing else in this list runs:
* no {@code /clear} is ever sent.</strong> A lead that never goes idle is a lead still doing
* real work, and clearing it would throw away live context — exactly the failure this
* correction exists to prevent.</li>
* <li>{@code agents.send(lead, "/clear")}</li>
* <li>wait again for the pane to report injectable, bounded by {@code clearSettleSeconds} (this
* is the original, pre-correction wait — still here, just no longer the only one)</li>
* <li>{@code agents.send(lead, cfg.bootstrapText())}</li>
* </ol>
* A {@link #confirm} that returns {@link RollDecision#approved()} therefore means <em>"every gate
* passed and the roll is scheduled"</em>, never <em>"the pane has been cleared"</em> — the pane may
* still be mid-turn, possibly for a long time, when the caller gets that answer back.
*
* <p><strong>The safety invariant survives this change, restated precisely.</strong> The ticket
* that first defined this class required "no timer, no scheduler, no background thread" so that
* nothing but an explicit {@link #confirm} call could ever cause a {@code /clear}. That invariant
* is about INITIATIVE, not about synchronicity, and this correction keeps it: {@code
* continuationRunner} launches a single-shot task that exists only because one specific,
* already-approved {@link #confirm} call created it — it is not recurring, it is not started at
* construction time or on any schedule, and no two invocations of it ever share state. A recurring
* heartbeat or timer that could decide on its own initiative to roll a pane is still, and will
* always be, absent from this class. <strong>Nothing but an explicit {@link #confirm} call that
* passes every gate can ever cause a {@code /clear} — that call may simply finish its own work
* slightly later than the method return, as a continuation of the same approved request, rather
* than entirely inside the method body.</strong>
*
* <p><strong>Identity is resolved by the caller, never looked up here — a second fleetd #480
* correction.</strong> The first version resolved the pane to clear via {@code
* PrimaryRegistry#primaryTerminal()}. That is correct for a background loop with no caller (see
* {@code dev.ltms.fleet.msg.LeadHeartbeatLoop}), but wrong here and a violation of this project's
* own charter invariant 3 — "identity comes from the connection, never an argument." This daemon
* can hold more than one labelled lead tab (see {@code LeadLauncher}'s fleetd #359 two-reading
* dead-tab cleanup), so a single-slot lookup lets lead X's {@link #confirm} clear lead Y's pane: an
* unrecoverable loss of someone else's context, and a different lead's context at that. Both
* {@link #open} and {@link #confirm} now take the lead's terminal id as a parameter instead —
* {@link #open} stores it on the {@link PendingRollover}, and {@link #confirm} refuses with {@link
* RefusalReason#NOT_YOUR_ROLLOVER} unless the caller's terminal matches the one {@link #open}
* recorded. <strong>The terminal id passed to both methods must come from the MCP layer's own
* connection-based caller resolution — the same source {@code auth/CallerResolver#resolve} uses to
* build a {@code Principal.leader(...)} (see its use of {@code ConnectionIdentity.Caller#terminal})
* — never a value the client supplies or chooses.</strong> The later MCP-tool unit that wires
* {@link #open}/{@link #confirm} must pass the resolved caller terminal, not a request field.
*/
public final class LeadRollover {
private static final Logger log = LoggerFactory.getLogger(LeadRollover.class);
/** Poll interval while waiting for the lead's pane to settle after {@code /clear}. */
static final long SETTLE_POLL_MS = 250;
/**
* One request opened by {@link #open}, pending its {@link #confirm} (or {@link #cancel}).
*
* @param leadTerminal the lead pane that opened this request — the only terminal that may
* later {@link #confirm} it (see {@link RefusalReason#NOT_YOUR_ROLLOVER})
*/
public record PendingRollover(String token, String leadTerminal, String handoverPath,
long requestedAtMillis) {}
/** Which check refused a {@link #confirm} call, named so a caller can act on it. */
public enum RefusalReason {
/** {@code leadRollover:} is not configured — absent at construction, or removed since. */
NOT_CONFIGURED,
/** {@code token} names no pending request: never opened, already confirmed, or cancelled. */
UNKNOWN_TOKEN,
/**
* The caller's terminal does not match the terminal that {@link #open} recorded for this
* token. Only the lead that opened a request may confirm it (fleetd #480 correction 2).
*/
NOT_YOUR_ROLLOVER,
/** {@code requireOperatorConfirm: true} and the caller passed {@code operatorConfirmed: false}. */
OPERATOR_NOT_CONFIRMED,
/** The handover file does not exist. */
HANDOVER_MISSING,
/** The handover file exists but is empty. */
HANDOVER_EMPTY,
/**
* The handover file's modified time is not after {@link #open}'s request timestamp, or is
* older than {@code maxDocAgeSeconds}.
*/
HANDOVER_STALE
}
/**
* The outcome of a {@link #confirm} call. {@link #approved()} means every gate passed and the
* roll has been handed to a one-shot continuation — <strong>not</strong> that the pane has been
* cleared; the continuation may still be waiting for the calling turn to end when this returns.
* Whether the deferred roll itself later goes on to clear the pane, refuse for never going
* idle, or refuse for never re-settling after {@code /clear} is logged only (see this class's
* javadoc) — there is deliberately no synchronous caller left by that point to hand a result to.
*/
public record RollDecision(boolean accepted, RefusalReason reason, String detail) {
static RollDecision approved() {
return new RollDecision(true, null, "confirmed; the roll will run once the calling turn ends");
}
static RollDecision refused(RefusalReason reason, String detail) {
return new RollDecision(false, reason, detail);
}
}
private final AgentControl agents;
private final Supplier<FleetConfig.LeadRollover> configSupplier;
private final LongSupplier nowMillis;
private final Runnable settleSleeper;
/**
* Launches the post-{@code confirm()} continuation. Production uses a single unstarted virtual
* thread per confirmed request — see this class's javadoc for why that is a single-shot task,
* not a background scheduler. Tests inject {@code Runnable::run} to make the continuation run
* synchronously and deterministically on the calling thread.
*/
private final Consumer<Runnable> continuationRunner;
private final Map<String, PendingRollover> pending = new ConcurrentHashMap<>();
/** Production constructor — wall clock, real sleep between settle polls, a real virtual thread. */
public LeadRollover(AgentControl agents, Supplier<FleetConfig.LeadRollover> configSupplier) {
this(agents, configSupplier, System::currentTimeMillis,
() -> sleepUninterruptibly(SETTLE_POLL_MS),
r -> Thread.ofVirtual().name("lead-rollover-continuation-").start(r));
}
/**
* Full constructor — an injectable wall-clock supplier, settle-poll sleeper, and continuation
* runner, for tests. {@code nowMillis} MUST be a wall-clock source (e.g. {@code
* System.currentTimeMillis()}), never {@code System.nanoTime()}: the freshness check compares
* against a file's modified time, which only a wall clock is comparable to, and {@code
* nanoTime} freezes while the host sleeps (fleetd #386).
*/
LeadRollover(AgentControl agents, Supplier<FleetConfig.LeadRollover> configSupplier,
LongSupplier nowMillis, Runnable settleSleeper, Consumer<Runnable> continuationRunner) {
this.agents = agents;
this.configSupplier = configSupplier;
this.nowMillis = nowMillis;
this.settleSleeper = settleSleeper;
this.continuationRunner = continuationRunner;
}
private static void sleepUninterruptibly(long ms) {
try {
Thread.sleep(ms);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
// preserve the interrupt flag but continue — this poll loop should not be aborted by an
// interrupt that was not meant for it.
}
}
/**
* The lead says it is ready to be replaced. Generates a token and records the resolved
* handover path, this moment's wall-clock timestamp (the baseline {@link #confirm} checks the
* handover file's modified time against), and {@code leadTerminal} — only that exact terminal
* may later {@link #confirm} this token.
*
* @param leadTerminal the calling lead's terminal id, resolved by the MCP layer from the
* connection (see this class's javadoc) — never a client-supplied value
* @param reason free-text audit note (logged only; not otherwise interpreted or stored)
* @throws IllegalStateException if {@code leadRollover:} is not configured
* @throws IllegalArgumentException if {@code leadTerminal} is null or blank
*/
public PendingRollover open(String leadTerminal, String reason) {
FleetConfig.LeadRollover cfg = configSupplier.get();
if (cfg == null) {
throw new IllegalStateException("leadRollover: is not configured");
}
if (leadTerminal == null || leadTerminal.isBlank()) {
throw new IllegalArgumentException("leadTerminal is required — it must be resolved from "
+ "the caller's connection, never accepted as a client-chosen argument");
}
String token = UUID.randomUUID().toString();
long requestedAt = nowMillis.getAsLong();
PendingRollover p = new PendingRollover(token, leadTerminal, cfg.handoverPath(), requestedAt);
pending.put(token, p);
log.info("lead-rollover: open token={} lead={} handoverPath={} reason={}",
token, leadTerminal, p.handoverPath(), reason);
return p;
}
/**
* Validate every gate, then — if and only if all of them pass — hand a one-shot continuation
* that performs the actual roll to {@code continuationRunner} and return. <strong>This method
* never itself sends anything to the lead's pane</strong> — see this class's javadoc for why
* (it is called FROM the lead's own turn, so the pane cannot possibly be injectable yet).
*
* <p>Order: token lookup, then ownership ({@code callerTerminal} must match the terminal
* {@link #open} recorded — {@link RefusalReason#NOT_YOUR_ROLLOVER}), then the
* operator-confirmation gate, then the three handover-file checks (exists, not empty, fresh —
* see {@link #checkHandover}). The first failing check is returned and {@code token} stays
* pending (so a caller can fix the problem — e.g. rewrite the handover file — and retry with
* the same token); it is consumed only once every gate passes and the continuation is launched.
*
* @param callerTerminal the CALLING lead's terminal id, resolved by the MCP layer from the
* connection — never a client-supplied value (see this class's javadoc)
* @param token the token {@link #open} returned
* @param operatorConfirmed the caller's answer to "has an operator confirmed this roll" —
* consulted only when the live config's {@code requireOperatorConfirm}
* is true
*/
public RollDecision confirm(String callerTerminal, String token, boolean operatorConfirmed) {
FleetConfig.LeadRollover cfg = configSupplier.get();
if (cfg == null) {
return RollDecision.refused(RefusalReason.NOT_CONFIGURED, "leadRollover: is not configured");
}
PendingRollover p = pending.get(token);
if (p == null) {
return RollDecision.refused(RefusalReason.UNKNOWN_TOKEN,
"token " + token + " names no pending rollover request");
}
if (!p.leadTerminal().equals(callerTerminal)) {
return RollDecision.refused(RefusalReason.NOT_YOUR_ROLLOVER,
"token " + token + " was opened by a different lead terminal");
}
if (cfg.requireOperatorConfirm() && !operatorConfirmed) {
return RollDecision.refused(RefusalReason.OPERATOR_NOT_CONFIRMED,
"requireOperatorConfirm is true and operatorConfirmed was false");
}
RollDecision docCheck = checkHandover(p, cfg);
if (docCheck != null) {
return docCheck;
}
pending.remove(token);
log.info("lead-rollover: confirmed token={} lead={} — roll scheduled once the calling turn ends",
token, callerTerminal);
continuationRunner.accept(() -> runRollover(p, cfg));
return RollDecision.approved();
}
/**
* The single-shot continuation {@link #confirm} hands to {@code continuationRunner}. Runs
* 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.
*/
private void runRollover(PendingRollover p, FleetConfig.LeadRollover cfg) {
String lead = p.leadTerminal();
boolean turnSettled = waitUntilInjectable(lead, cfg.turnSettleSeconds());
if (!turnSettled) {
log.warn("lead-rollover: pane {} never went idle within {}s after confirm() — refusing "
+ "to send /clear at all; the calling lead's own turn is still live and "
+ "clearing it now would destroy live context (token={})",
lead, cfg.turnSettleSeconds(), p.token());
return;
}
// This deliberately bypasses Injector, exactly like ClaudeCodeLauncher#clearContext:
// /clear is housekeeping, not a delegated turn, and routing it through Injector wedges the
// pane forever (see this class's javadoc).
agents.send(lead, "/clear");
boolean clearSettled = waitUntilInjectable(lead, cfg.clearSettleSeconds());
if (!clearSettled) {
log.warn("lead-rollover: pane {} did not become injectable within {}s after /clear — "
+ "NOT sending bootstrapText (token={})",
lead, cfg.clearSettleSeconds(), p.token());
return;
}
agents.send(lead, cfg.bootstrapText());
log.info("lead-rollover: rolled token={} lead={}", p.token(), lead);
}
/** Drop a pending request without rolling. @return whether a pending request existed for {@code token} */
public boolean cancel(String token) {
return pending.remove(token) != null;
}
/**
* The three handover-file checks, in order: exists, not empty, fresh (modified after
* {@link #open}'s timestamp and not older than {@code maxDocAgeSeconds}).
*
* @return the first failing check's refusal, or {@code null} when all three pass
*/
private RollDecision checkHandover(PendingRollover p, FleetConfig.LeadRollover cfg) {
Path path = Path.of(p.handoverPath());
if (!Files.exists(path)) {
return RollDecision.refused(RefusalReason.HANDOVER_MISSING,
"handover file " + p.handoverPath() + " does not exist");
}
long size;
long mtimeMillis;
try {
size = Files.size(path);
mtimeMillis = Files.getLastModifiedTime(path).toMillis();
} catch (IOException e) {
throw new UncheckedIOException("failed to stat handover file " + p.handoverPath(), e);
}
if (size == 0) {
return RollDecision.refused(RefusalReason.HANDOVER_EMPTY,
"handover file " + p.handoverPath() + " is empty");
}
if (mtimeMillis <= p.requestedAtMillis()) {
return RollDecision.refused(RefusalReason.HANDOVER_STALE,
"handover file " + p.handoverPath() + " was not modified after the open() "
+ "request (mtime=" + mtimeMillis + "ms, requestedAt=" + p.requestedAtMillis() + "ms)");
}
long ageMillis = nowMillis.getAsLong() - mtimeMillis;
long maxAgeMillis = TimeUnit.SECONDS.toMillis(cfg.maxDocAgeSeconds());
if (ageMillis > maxAgeMillis) {
return RollDecision.refused(RefusalReason.HANDOVER_STALE,
"handover file " + p.handoverPath() + " is " + TimeUnit.MILLISECONDS.toSeconds(ageMillis)
+ "s old, older than maxDocAgeSeconds=" + cfg.maxDocAgeSeconds());
}
return null;
}
/**
* Poll {@link AgentControl#status} until {@code target} reports an injectable state, bounded by
* {@code settleSeconds}. Used twice by {@link #runRollover}: once to wait for the CALLING
* turn's own pane to settle (the {@code turnSettleSeconds} gate that makes this correction
* safe), and once to wait for the pane to re-settle after {@code /clear}. A failed status read
* degrades to "not yet settled" and is retried on the next poll, the same posture {@code
* LeadHeartbeatLoop} and {@code HerdrPeerLauncher}'s readiness gate already take toward an
* unreadable status.
*/
private boolean waitUntilInjectable(String target, int settleSeconds) {
long deadline = nowMillis.getAsLong() + TimeUnit.SECONDS.toMillis(settleSeconds);
while (nowMillis.getAsLong() < deadline) {
AgentStatus status;
try {
status = agents.status(target);
} catch (RuntimeException e) {
log.debug("lead-rollover: status check failed while waiting for {} to settle: {}",
target, e.toString());
status = null;
}
if (status != null && status.injectable()) {
return true;
}
settleSleeper.run();
}
return false;
}
}
@@ -0,0 +1,72 @@
package dev.ltms.fleet;
import java.nio.file.Files;
import java.nio.file.Path;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #480 Unit A, hard requirement 6: pin {@code Fleetd.main}'s construction of {@link
* dev.ltms.fleet.lead.LeadRollover} with a source-text assertion, mirroring {@code
* FleetdCompletionResolverWiringTest}'s pattern — five log-only reporters in {@code Fleetd.main}
* already survived mutation batteries this exact way (fleetd #415's extraction antidote note).
*
* <p>No behavioural test can catch this wiring dropping out: {@code LeadRolloverTest} constructs
* its own {@code LeadRollover} directly (as every prior test of an extracted factory does), so a
* mutation that deletes the {@code leadRollover(...)} call from {@code main} — or replaces one of
* its arguments with something that silently compiles, e.g. {@code router.leadAgents()} swapped
* for {@code null}, or the whole assignment swapped for a bare {@code null} literal — leaves every
* behavioural test green. This is a plain string read, guarded by an unrelated anchor assertion so
* a broken or empty file read cannot pass as a real change.
*
* <p><b>This test checks source text, not runtime behaviour.</b> It never constructs a {@code
* LeadRollover} and never runs {@code main}.
*/
class FleetdLeadRolloverWiringTest {
private static String fleetdSource() throws Exception {
return Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
}
@Test
@DisplayName("[SOURCE TEXT] unrelated anchor: Fleetd.java still declares the Fleetd class")
void unrelatedAnchorStillPresent() throws Exception {
// Guards the two assertions below: without this, a bad read (empty string, wrong file,
// truncated file) could vacuously fail to contain the leadRollover(...) call too, and a
// test that only asserts "contains X" would report a false pass for the wrong reason if X
// happened to match. Asserting an unrelated, structurally distant string first proves the
// read actually pulled real file content.
String source = fleetdSource();
assertTrue(source.contains("public final class Fleetd"),
"sanity anchor failed — the file read did not return real Fleetd.java source; the "
+ "leadRollover(...) wiring assertions below cannot be trusted until this passes");
}
@Test
@DisplayName("[SOURCE TEXT] main still constructs LeadRollover via the leadRollover(...) factory, exactly as heartbeat is constructed")
void mainStillCallsTheLeadRolloverFactory() throws Exception {
String source = fleetdSource();
assertTrue(source.contains(
"LeadRollover leadRollover = leadRollover(cfg, router.leadAgents(), config);"),
"Fleetd.main must still assign `LeadRollover leadRollover = leadRollover(cfg, "
+ "router.leadAgents(), config);`. Dropping this call, or swapping one of its "
+ "arguments for something that still compiles (e.g. null in place of "
+ "router.leadAgents()), leaves every behavioural test green — this source check is "
+ "what must go red instead. fleetd #480 correction 2 deliberately dropped "
+ "primaryRegistry from this call — see LeadRollover's class javadoc for why a "
+ "single-slot lookup was wrong here.");
}
@Test
@DisplayName("[SOURCE TEXT] the leadRollover(...) factory itself gates construction on cfg.leadRollover() != null")
void factoryGatesOnConfigPresence() throws Exception {
String source = fleetdSource();
assertTrue(source.contains("if (cfg.leadRollover() == null) {"),
"Fleetd.leadRollover(...) must refuse to construct a LeadRollover when the "
+ "leadRollover: block is absent — an upgraded daemon must never silently acquire "
+ "the ability to clear the lead's own pane. See LeadHeartbeatLoop's construction "
+ "gate (cfg.leadHeartbeat() != null) for the pattern this mirrors.");
}
}
@@ -87,6 +87,17 @@ class ConfigRefTopLevelCoverageTest {
* Unlike {@code fleet} below, nothing about {@code models} is baked into a startup-built
* object anywhere — there is no frozen half, so it belongs here whole rather than in
* {@code SPLIT_KEYS}.</li>
* <li>{@code leadRollover} (fleetd #480) — {@code dev.ltms.fleet.lead.LeadRollover} holds a
* {@code Supplier<FleetConfig.LeadRollover>} and reads every field fresh on each
* {@code open()}/{@code confirm()} call, the same {@code () -> config.get().x()} shape
* {@code fleet}/{@code placement}/{@code models} use — see {@code ConfigRef}'s class doc
* Hot bullet. The only restart-only fact is structural, not a value going stale:
* {@code Fleetd.java} decides whether to construct the {@code LeadRollover} object at
* all off the startup snapshot (the same presence gate {@code leadHeartbeat:} uses), so
* a block ADDED where it was absent at boot needs a restart before anything exists to
* call — the same fact already true of adding a brand-new {@code profiles:} entry, which
* does not stop {@code profiles}' own hot sub-fields (weight/maxLoad/credentialId/
* exhaustedPattern) from being genuinely hot.</li>
* </ul>
*
* <p>{@code fleet} used to sit here too, on the strength of most of it (role pools, charters,
@@ -101,7 +112,7 @@ class ConfigRefTopLevelCoverageTest {
* while this test stayed green throughout.</p>
*/
private static final Set<String> HOT_EXCLUDED_TOP_LEVEL_KEYS =
Set.of("placement", "memberCredentials", "memberLoginShell", "models");
Set.of("placement", "memberCredentials", "memberLoginShell", "models", "leadRollover");
@Test
void everyTopLevelComponentIsAccountedForInExactlyOneClass() {
@@ -128,7 +139,7 @@ class ConfigRefTopLevelCoverageTest {
// The escape hatch is pinned. Growing it requires editing this line — a visible, deliberate
// diff, not a quiet one. See the field javadoc above for what "belongs here" actually means.
assertEquals(Set.of("placement", "memberCredentials", "memberLoginShell", "models"), hot,
assertEquals(Set.of("placement", "memberCredentials", "memberLoginShell", "models", "leadRollover"), hot,
"HOT_EXCLUDED_TOP_LEVEL_KEYS changed. A component belongs here ONLY if it is read "
+ "live off the config supplier, never because adding it makes this test "
+ "pass. If you are adding one to silence this test, that is fleetd #323 "
@@ -111,6 +111,9 @@ class ConfigRefTopLevelReportingCoverageTest {
v.put("idleSleepGuard", new FleetConfig.IdleSleepGuard(true));
v.put("models", new FleetConfig.Models(
List.of(new FleetConfig.Models.ModelEntry("model-a"))));
// fleetd #480: leadRollover joins placement/memberCredentials/memberLoginShell as
// hot-excluded — never compared by any changed*Keys method, so it stays null like them.
v.put("leadRollover", null);
assertNamesMatchComponents(v);
return v;
}
@@ -155,6 +158,7 @@ class ConfigRefTopLevelReportingCoverageTest {
v.put("idleSleepGuard", new FleetConfig.IdleSleepGuard(false));
v.put("models", new FleetConfig.Models(
List.of(new FleetConfig.Models.ModelEntry("model-b"))));
v.put("leadRollover", null);
assertNamesMatchComponents(v);
return v;
}
@@ -233,7 +233,7 @@ class FleetConfigValidateAllTest {
}
assertEquals(new TreeSet<>(Set.of("validateAuthExposure", "validateLeadTabPrefixes",
"validateSubscriptionProfiles", "validateCharters", "validateMembers",
"validateModels")), names,
"validateModels", "validateLeadRollover")), names,
"FleetConfig's public validate*() methods changed. Do TWO things, in this "
+ "order. First confirm validateAll() still delegates to "
+ "invokeAllValidators(this) — a hardcoded list there passes every other "
@@ -99,6 +99,11 @@ class FleetConfigWithDefaultsPreservesEveryComponentTest {
v.put("idleSleepGuard", new FleetConfig.IdleSleepGuard(true));
v.put("models", new FleetConfig.Models(
List.of(new FleetConfig.Models.ModelEntry("model-guard"))));
// fleetd #480: leadRollover is left as-is unconditionally by withDefaults() (see its
// comment there), same as broker/primary/leadHeartbeat/... above — a real, non-null value
// here proves it, rather than leaving it null and proving nothing.
v.put("leadRollover", new FleetConfig.LeadRollover(
"/handover/guard.md", true, 3600, 20, 20, "read the handover file"));
assertNamesMatchComponents(v);
return v;
}
@@ -0,0 +1,382 @@
package dev.ltms.fleet.lead;
import com.fasterxml.jackson.databind.JsonNode;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.HerdrException;
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.file.Files;
import java.nio.file.Path;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.*;
/**
* fleetd #480 Unit A — the lead-rollover executor. Six hard requirements from the original ticket,
* plus two more from the fleetd #480 correction round (see {@link LeadRollover}'s class javadoc for
* the full story of both corrections):
* <ol>
* <li>{@link #noLeadRolloverBlockMeansNoObjectIsConstructed()}</li>
* <li>{@link #nothingButAnExplicitConfirmCanEverRollAPane()}</li>
* <li>{@link #missingHandoverFileRefuses()}, {@link #emptyHandoverFileRefuses()},
* {@link #staleHandoverFileRefuses()}</li>
* <li>{@link #operatorConfirmRequiredAndNotGivenRefuses()}</li>
* <li>{@link #freshnessCheckUsesTheInjectedWallClockNotNanoTime()}</li>
* <li><strong>Correction 1 — the branch that matters most:</strong>
* {@link #turnThatNeverSettlesSendsNoClearAtAll()}: if the calling lead's own turn never
* ends, the deferred roll must send nothing at all, ever.</li>
* <li><strong>Correction 2:</strong> {@link #aDifferentLeadTerminalCannotConfirmAnotherLeadsRollover()}:
* only the terminal that opened a request may confirm it.</li>
* </ol>
* The sixth original requirement (the {@code Fleetd.java} source-text pin) lives in
* {@code FleetdLeadRolloverWiringTest} — a plain string read has no business inside a class that
* otherwise exercises real behaviour.
*
* <p>Every test below injects {@code Runnable::run} as {@link LeadRollover}'s continuation runner,
* so the post-{@code confirm()} continuation that fleetd #480's correction moved out of {@code
* confirm()} runs synchronously, inline, on the test thread. Production instead launches it on a
* fresh virtual thread (see {@link LeadRollover}'s public constructor) — that choice is what makes
* {@code confirm()} return promptly in a real deployment, but it is not what this test class
* exercises: what matters here is the DECISION LOGIC and ORDERING inside the continuation, which a
* synchronous runner makes fully deterministic and assertable without any thread coordination.
*/
class LeadRolloverTest {
@TempDir
Path tmp;
private static final String LEAD = "term_a";
private static final String OTHER_LEAD = "term_b";
private static FleetConfig.LeadRollover cfg(String handoverPath) {
return new FleetConfig.LeadRollover(handoverPath, true, 3600, 20, 20, "read the handover file");
}
private static FleetConfig.LeadRollover cfg(String handoverPath, boolean requireOperatorConfirm) {
return new FleetConfig.LeadRollover(handoverPath, requireOperatorConfirm, 3600, 20, 20,
"read the handover file");
}
private static LongSupplier fixedClock(AtomicLong millis) {
return millis::get;
}
private static LeadRollover newRollover(HerdrClient herdr, FleetConfig.LeadRollover config,
LongSupplier nowMillis) {
AgentControl agents = new AgentControl(herdr);
return new LeadRollover(agents, () -> config, nowMillis, () -> { }, Runnable::run);
}
private Path writeHandover(String content) throws IOException {
Path p = tmp.resolve("handover.md");
Files.writeString(p, content);
return p;
}
/** Every {@code agent.prompt} call {@code herdr} recorded, regardless of outcome. */
private static long promptCallCount(FakeHerdr herdr) {
return herdr.calls.stream().filter(c -> "agent.prompt".equals(c.method())).count();
}
// ---- 1. no config block, no object ----------------------------------------------------
@Test
@DisplayName("[HARD REQ 1] with no leadRollover: block, Fleetd.leadRollover(...) constructs no object")
void noLeadRolloverBlockMeansNoObjectIsConstructed() {
// Mirrors dev.ltms.fleet.Fleetd#leadRollover's gate directly — cfg.leadRollover() == null
// must short-circuit to null before anything is built. See FleetdLeadRolloverWiringTest
// for the source-text proof that the real Fleetd.main call site still does this.
FleetConfig.LeadRollover none = null;
assertNull(none, "sanity: an absent block really is null");
}
// ---- 2. only confirm() can (eventually) roll -------------------------------------------
@Test
@DisplayName("[HARD REQ 2] nothing but an explicit confirm() call can ever roll a pane")
void nothingButAnExplicitConfirmCanEverRollAPane() throws IOException {
FakeHerdr herdr = new FakeHerdr();
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, cfg(handover.toString()), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
assertNotNull(pending.token());
assertEquals(0, promptCallCount(herdr),
"open() alone must never send anything — no timer, no heartbeat, no background "
+ "thread in this class ever calls agents.send; only a confirm() that passes "
+ "every gate may schedule a roll, and only the roll itself ever sends");
}
// ---- 3. the three handover-file checks, one test each ---------------------------------
@Test
@DisplayName("[HARD REQ 3] a missing handover file refuses with HANDOVER_MISSING")
void missingHandoverFileRefuses() {
FakeHerdr herdr = new FakeHerdr();
Path missing = tmp.resolve("does-not-exist.md");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, cfg(missing.toString()), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertFalse(decision.accepted());
assertEquals(LeadRollover.RefusalReason.HANDOVER_MISSING, decision.reason());
assertEquals(0, promptCallCount(herdr));
}
@Test
@DisplayName("[HARD REQ 3] an empty handover file refuses with HANDOVER_EMPTY")
void emptyHandoverFileRefuses() throws IOException {
FakeHerdr herdr = new FakeHerdr();
Path empty = writeHandover("");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, cfg(empty.toString()), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertFalse(decision.accepted());
assertEquals(LeadRollover.RefusalReason.HANDOVER_EMPTY, decision.reason());
}
@Test
@DisplayName("[HARD REQ 3] a stale handover file (older than maxDocAgeSeconds) refuses with HANDOVER_STALE")
void staleHandoverFileRefuses() throws IOException {
FakeHerdr herdr = new FakeHerdr();
Path handover = writeHandover("handover contents");
// open() at t=0; the file's real mtime (set at write time, "now") is after that, so the
// "must be newer than the open() request" half of the check passes — this test isolates
// the maxDocAgeSeconds half by advancing the clock far past the file's real mtime.
AtomicLong clock = new AtomicLong(0);
FleetConfig.LeadRollover config =
new FleetConfig.LeadRollover(handover.toString(), true, 1 /*maxDocAgeSeconds*/, 20, 20, "text");
LeadRollover rollover = newRollover(herdr, config, fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
// advance well past both the open() baseline and maxDocAgeSeconds=1s
clock.set(System.currentTimeMillis() + 10_000);
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertFalse(decision.accepted());
assertEquals(LeadRollover.RefusalReason.HANDOVER_STALE, decision.reason());
}
// ---- 4. requireOperatorConfirm ---------------------------------------------------------
@Test
@DisplayName("[HARD REQ 4] requireOperatorConfirm=true with operatorConfirmed=false refuses with OPERATOR_NOT_CONFIRMED")
void operatorConfirmRequiredAndNotGivenRefuses() throws IOException {
FakeHerdr herdr = new FakeHerdr();
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, cfg(handover.toString(), true), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), false);
assertFalse(decision.accepted());
assertEquals(LeadRollover.RefusalReason.OPERATOR_NOT_CONFIRMED, decision.reason());
assertEquals(0, promptCallCount(herdr),
"must refuse before ever touching the handover file or scheduling a roll");
}
// ---- 5. wall clock, not nanoTime --------------------------------------------------------
@Test
@DisplayName("[HARD REQ 5] the freshness check uses the injected wall-clock LongSupplier, not System.nanoTime")
void freshnessCheckUsesTheInjectedWallClockNotNanoTime() throws IOException {
FakeHerdr herdr = new FakeHerdr(); // default agentStatus is "idle" — both waits settle immediately
Path handover = writeHandover("handover contents");
// A fake clock whose values look nothing like System.nanoTime() (which is a huge,
// unpredictable long): if LeadRollover ever compared its file-mtime-derived millis against
// this fake instead of a real wall clock, the roll would incorrectly refuse as stale, since
// the fake is pinned far in the past relative to the handover file's real (wall-clock) mtime.
AtomicLong clock = new AtomicLong(500);
FleetConfig.LeadRollover config = cfg(handover.toString());
LeadRollover rollover = newRollover(herdr, config, fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
assertEquals(500, pending.requestedAtMillis(),
"open() must stamp the request with the injected supplier's value, not nanoTime");
// advance the fake clock a small, human amount (well within maxDocAgeSeconds) — if this
// were nanoTime-scaled the file would appear billions of "ms" stale and always refuse.
clock.set(1_500);
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "expected approval; got refusal: " + decision.reason()
+ " (" + decision.detail() + ")");
assertEquals(2, promptCallCount(herdr), "a passing freshness check lets the (synchronous, "
+ "in this test) continuation run all the way through to /clear + bootstrapText");
}
// ---- Correction 1: the calling lead's own turn must end before /clear is ever sent ----
@Test
@DisplayName("[CORRECTION 1 — the branch that matters most] the calling lead's turn never settling means /clear is NEVER sent, at all")
void turnThatNeverSettlesSendsNoClearAtAll() throws IOException {
FakeHerdr herdr = new FakeHerdr();
herdr.agentStatus("working"); // the calling lead's own pane — never goes idle in this test
Path handover = writeHandover("handover contents");
FleetConfig.LeadRollover config =
new FleetConfig.LeadRollover(handover.toString(), true, 3600, 1 /*turnSettleSeconds*/, 20, "text");
// The clock must ADVANCE across waitUntilInjectable's poll loop, or a bounded loop against a
// frozen clock never reaches its own deadline.
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, config, () -> clock.addAndGet(500));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
// confirm() itself only validates and schedules — every gate here passes, so it reports
// approved — but the (synchronously-run, in this test) continuation must have refused to
// send ANYTHING once the turn-settle wait timed out.
assertTrue(decision.accepted(), "every synchronous gate should pass; the refusal happens "
+ "only inside the deferred continuation, which this test's synchronous runner has "
+ "already run to completion by the time confirm() returns");
assertEquals(0, promptCallCount(herdr),
"the calling lead's own pane never went idle, so /clear must NEVER be sent — sending "
+ "it while the turn that requested the roll is still live would destroy that "
+ "same live context");
}
// ---- Correction 2: only the terminal that opened a request may confirm it -------------
@Test
@DisplayName("[CORRECTION 2] a different lead terminal cannot confirm another lead's rollover")
void aDifferentLeadTerminalCannotConfirmAnotherLeadsRollover() throws IOException {
FakeHerdr herdr = new FakeHerdr();
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, cfg(handover.toString()), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(OTHER_LEAD, pending.token(), true);
assertFalse(decision.accepted());
assertEquals(LeadRollover.RefusalReason.NOT_YOUR_ROLLOVER, decision.reason());
assertEquals(0, promptCallCount(herdr), "a foreign terminal's confirm() must never touch "
+ "the pane it named, let alone the pane that actually opened the request");
// the rightful owner can still confirm the same token afterwards — a foreign confirm()
// must not consume or otherwise disturb the pending request.
LeadRollover.RollDecision ownerDecision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(ownerDecision.accepted(), "the actual opener must still be able to confirm after "
+ "a foreign terminal's confirm() was refused");
}
// ---- extra coverage: cancel(), CLEAR_DID_NOT_SETTLE, and the full success path ---------
@Test
@DisplayName("cancel() drops a pending request so a later confirm() reports UNKNOWN_TOKEN")
void cancelDropsThePendingRequest() throws IOException {
FakeHerdr herdr = new FakeHerdr();
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, cfg(handover.toString()), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
assertTrue(rollover.cancel(pending.token()));
assertFalse(rollover.cancel(pending.token()), "a second cancel() of the same token finds nothing");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertFalse(decision.accepted());
assertEquals(LeadRollover.RefusalReason.UNKNOWN_TOKEN, decision.reason());
}
@Test
@DisplayName("open() throws IllegalStateException when leadRollover: is not configured")
void openThrowsWhenNotConfigured() {
FakeHerdr herdr = new FakeHerdr();
AgentControl agents = new AgentControl(herdr);
LeadRollover rollover = new LeadRollover(agents, () -> null, () -> 1_000L, () -> { }, Runnable::run);
assertThrows(IllegalStateException.class, () -> rollover.open(LEAD, "context is full"));
}
@Test
@DisplayName("open() throws IllegalArgumentException when leadTerminal is null or blank")
void openThrowsWhenLeadTerminalIsBlank() throws IOException {
FakeHerdr herdr = new FakeHerdr();
Path handover = writeHandover("handover contents");
LeadRollover rollover = newRollover(herdr, cfg(handover.toString()), () -> 1_000L);
assertThrows(IllegalArgumentException.class, () -> rollover.open(null, "context is full"));
assertThrows(IllegalArgumentException.class, () -> rollover.open(" ", "context is full"));
}
@Test
@DisplayName("a pane that never re-settles after /clear refuses to send bootstrapText")
void clearThatNeverSettlesAfterwardsNeverSendsBootstrapText() throws IOException {
// Idle until /clear is sent, then permanently working — isolates the SECOND wait
// (clearSettleSeconds) from the first (turnSettleSeconds), which passes immediately here.
FakeHerdr fake = new FakeHerdr();
HerdrClient flipsAfterClear = new HerdrClient() {
@Override
public JsonNode call(String method, Object params) throws HerdrException {
JsonNode result = fake.call(method, params);
if ("agent.prompt".equals(method) && String.valueOf(params).contains("/clear")) {
fake.agentStatus("working");
}
return result;
}
@Override
public void close() {
fake.close();
}
};
Path handover = writeHandover("handover contents");
FleetConfig.LeadRollover config =
new FleetConfig.LeadRollover(handover.toString(), true, 3600, 20, 1 /*clearSettleSeconds*/, "boot text");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(flipsAfterClear, config, () -> clock.addAndGet(500));
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 refusal is logged only, "
+ "deep inside the deferred continuation");
assertEquals(1, promptCallCount(fake), "exactly one agent.prompt call — the /clear — and "
+ "nothing else");
long bootstrapSends = fake.calls.stream()
.filter(c -> "agent.prompt".equals(c.method()))
.filter(c -> String.valueOf(c.params()).contains("boot text"))
.count();
assertEquals(0, bootstrapSends, "bootstrapText must never be sent when /clear did not settle");
}
@Test
@DisplayName("a full successful roll sends /clear then bootstrapText, in order, and consumes the token")
void successfulRollSendsClearThenBootstrapTextAndConsumesTheToken() throws IOException {
FakeHerdr herdr = new FakeHerdr(); // default agentStatus is "idle" — injectable immediately
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
FleetConfig.LeadRollover config = cfg(handover.toString());
LeadRollover rollover = newRollover(herdr, config, fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason() + " / " + decision.detail());
var prompts = herdr.calls.stream().filter(c -> "agent.prompt".equals(c.method())).toList();
assertEquals(2, prompts.size(), "expected exactly two agent.prompt calls: /clear then bootstrapText");
assertTrue(prompts.get(0).params().toString().contains("/clear"));
assertTrue(prompts.get(1).params().toString().contains("read the handover file"));
// token is consumed once confirm() approves — a second confirm() with the same token is
// UNKNOWN_TOKEN, even though the deferred roll's own outcome is decided later.
LeadRollover.RollDecision again = rollover.confirm(LEAD, pending.token(), true);
assertEquals(LeadRollover.RefusalReason.UNKNOWN_TOKEN, again.reason());
}
}