fleetd #480 correction round: defer the roll, and gate it on caller identity
CI / contract (pull_request) Successful in 1m14s
CI / build (pull_request) Successful in 2m3s

Two defects found after the fact, both from the original brief, both fixed here.

1. confirm() is called FROM the calling lead's own turn, so its pane is still
   WORKING and can never report injectable inside that same call. The old
   confirm() sent /clear before polling for that — the poll always timed out,
   but only after /clear had already fired and queued, destroying the lead's
   context with no fresh session ever started and a refusal return that lied
   about what had happened.

   Fix: confirm() now only validates and, if every gate passes, hands a
   one-shot continuation to a new continuationRunner (a real virtual thread in
   production, Runnable::run in tests) and returns RollDecision.approved()
   immediately - "scheduled", not "rolled". The continuation itself does the
   actual work, once the calling turn has ended: wait for the SAME pane to
   report injectable again (new turnSettleSeconds config key, default 20) -
   if this never happens, /clear is NEVER sent, at all - then /clear, then
   wait again (clearSettleSeconds, as before), then bootstrapText. The "no
   timer/scheduler, only confirm() can roll" invariant is restated precisely
   in LeadRollover's class javadoc: it is about initiative, not synchronicity
   - a single-shot continuation of an already-approved confirm() call still
   satisfies it; a recurring background loop would not.

2. confirm() resolved the pane to clear via PrimaryRegistry.primaryTerminal(),
   a single-slot lookup that is correct for a background loop with no caller
   but wrong here: on a daemon with more than one labelled lead tab, lead X's
   confirm() could clear lead Y's pane, violating the charter's "identity
   comes from the connection, never an argument" invariant.

   Fix: open() and confirm() now take the caller's terminal id as a parameter
   (resolved by the MCP layer from the connection - the later MCP-tool unit
   must pass it in, never accept it as a request field). confirm() refuses
   with a new NOT_YOUR_ROLLOVER reason unless it matches the terminal open()
   recorded. LeadRollover no longer depends on PrimaryRegistry at all.

Also: renamed RollResult to RollDecision (rolled -> accepted) to reflect the
new meaning - approved and scheduled, not necessarily cleared yet. Added
turnSettleSeconds to the leadRollover: config block (documented in
fleetd.example.yaml alongside the existing keys) and updated Fleetd.java's
leadRollover(...) factory to drop the primaryRegistry parameter, with
FleetdLeadRolloverWiringTest's source-text pin updated to match.

New tests: turnThatNeverSettlesSendsNoClearAtAll (the branch that matters
most - a turn that never ends means /clear is never sent) and
aDifferentLeadTerminalCannotConfirmAnotherLeadsRollover (NOT_YOUR_ROLLOVER),
plus a settle-after-clear timeout test and an open() input-validation test.
LeadRolloverTest: 11 -> 14 tests.
This commit is contained in:
Dai Ha
2026-09-11 06:44:58 +07:00
parent 5c12865c25
commit a94262271b
8 changed files with 425 additions and 201 deletions
+13 -4
View File
@@ -94,8 +94,12 @@ bind:
#
# 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 can ever roll a pane: there is no timer, no
# heartbeat and no timeout anywhere in this feature that fires one on its own.
# 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
@@ -103,15 +107,20 @@ bind:
# 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)
# # 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."
+23 -14
View File
@@ -593,10 +593,14 @@ public final class Fleetd {
}
// 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 scheduler of its own — nothing but an explicit
// confirm() call (wired to an MCP tool by a later ticket; nothing calls it yet) can ever
// roll a pane, so there is no background thread here to shut down.
LeadRollover leadRollover = leadRollover(cfg, primaryRegistry, router.leadAgents(), config);
// 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);
@@ -1075,22 +1079,27 @@ public final class Fleetd {
* 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.
*
* @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 primaryRegistry where the lead's terminal is looked up at roll time
* @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
* <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, PrimaryRegistry primaryRegistry,
AgentControl leadAgents, ConfigRef config) {
static LeadRollover leadRollover(FleetConfig cfg, AgentControl leadAgents, ConfigRef config) {
if (cfg.leadRollover() == null) {
return null;
}
return new LeadRollover(primaryRegistry, leadAgents, () -> config.get().leadRollover());
return new LeadRollover(leadAgents, () -> config.get().leadRollover());
}
/**
@@ -67,8 +67,10 @@ import java.util.function.Supplier;
* {@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 clearSettleSeconds}/{@code bootstrapText} fresh on every {@code open()}/{@code
* confirm()} call rather than capturing them into fields at construction — unlike its closest
* {@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
@@ -1346,7 +1346,7 @@ public record FleetConfig(
* <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
* five fields takes effect on the very next {@code open()}/{@code confirm()} call once the block
* 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
@@ -1355,6 +1355,16 @@ public record FleetConfig(
* 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
@@ -1366,6 +1376,9 @@ public record FleetConfig(
* @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}.
@@ -1375,10 +1388,12 @@ public record FleetConfig(
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record LeadRollover(String handoverPath, Boolean requireOperatorConfirm,
Integer maxDocAgeSeconds, Integer clearSettleSeconds, String bootstrapText) {
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
@@ -3,7 +3,6 @@ package dev.ltms.fleet.lead;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.mcp.PrimaryRegistry;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -12,49 +11,82 @@ import java.io.UncheckedIOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Map;
import java.util.Optional;
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
* the three handover-file checks and (by default) an explicit operator confirmation all pass —
* {@link #confirm} clears the lead's own pane and bootstraps a fresh session against that file.
* 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>Structural template:</strong> {@code dev.ltms.fleet.msg.LeadHeartbeatLoop} (see its
* class javadoc for why an opt-in config block matters and how a pane is safely injected). The two
* differ in shape on purpose:
* <ul>
* <li>{@code LeadHeartbeatLoop} is a scheduled loop that ticks — and can inject — on its own
* initiative. This class has NO scheduler, NO timer and NO background thread anywhere.
* {@link #open}, {@link #confirm} and {@link #cancel} are the entire public surface, and only
* an explicit call to {@link #confirm} can ever clear a pane — see that method's javadoc.</li>
* <li>{@code LeadHeartbeatLoop} captures its config numbers into {@code final} fields at
* construction ({@code idleAfterNanos}/{@code backoffMs}/{@code quietNudgeCap}), which is
* exactly what makes {@code leadHeartbeat:} DEFERRED in {@code ConfigRef}'s reload
* bookkeeping. This class instead holds a {@code Supplier<FleetConfig.LeadRollover>} and
* reads every field fresh on each call, which is what makes {@code leadRollover:} HOT
* instead — see {@code ConfigRef}'s class doc and {@link FleetConfig.LeadRollover}'s javadoc
* for the exact claim and its one restart-only caveat (the object itself is constructed only
* when the block is present at startup, in {@code Fleetd.java}, the same gate
* {@code leadHeartbeat:} uses).</li>
* </ul>
* <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><strong>{@code /clear} bypasses the Injector, on purpose.</strong> Like {@code
* dev.ltms.fleet.member.ClaudeCodeLauncher#clearContext}, {@link #confirm} calls {@code
* agents.send(lead, "/clear")} directly rather than routing it through {@code Injector}. A live
* probe (fleetd #480 fact-find) proved a {@code /clear} routed through {@code Injector} wedges that
* pane forever: fleetd read {@code state: busy} while herdr read {@code liveStatus: done}, and
* every later message to that pane queued behind a turn that could never complete — {@code /clear}
* produces no turn boundary, so the Injector's own turn never finishes.
* <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 {
@@ -63,15 +95,26 @@ public final class LeadRollover {
/** 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}). */
public record PendingRollover(String token, String handoverPath, long requestedAtMillis) {}
/**
* 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 rolled, or cancelled. */
/** {@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. */
@@ -82,56 +125,61 @@ public final class LeadRollover {
* The handover file's modified time is not after {@link #open}'s request timestamp, or is
* older than {@code maxDocAgeSeconds}.
*/
HANDOVER_STALE,
/** No lead terminal is known to clear. */
LEAD_UNKNOWN,
/**
* {@code /clear} was sent but the pane never reported an injectable state again within
* {@code clearSettleSeconds} — {@code bootstrapText} was NOT sent.
*/
CLEAR_DID_NOT_SETTLE
HANDOVER_STALE
}
/** The outcome of a {@link #confirm} call. */
public record RollResult(boolean rolled, RefusalReason reason, String detail) {
static RollResult success() {
return new RollResult(true, null, "rolled");
/**
* 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 RollResult refused(RefusalReason reason, String detail) {
return new RollResult(false, reason, detail);
static RollDecision refused(RefusalReason reason, String detail) {
return new RollDecision(false, reason, detail);
}
}
private final PrimaryRegistry primaryRegistry;
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. */
public LeadRollover(PrimaryRegistry primaryRegistry, AgentControl agents,
Supplier<FleetConfig.LeadRollover> configSupplier) {
this(primaryRegistry, agents, configSupplier, System::currentTimeMillis,
() -> sleepUninterruptibly(SETTLE_POLL_MS));
/** 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 and settle-poll sleeper, 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).
* 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(PrimaryRegistry primaryRegistry, AgentControl agents,
Supplier<FleetConfig.LeadRollover> configSupplier,
LongSupplier nowMillis, Runnable settleSleeper) {
this.primaryRegistry = primaryRegistry;
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) {
@@ -146,85 +194,115 @@ public final class LeadRollover {
/**
* The lead says it is ready to be replaced. Generates a token and records the resolved
* handover path and this moment's wall-clock timestamp — the baseline {@link #confirm} checks
* the handover file's modified time against.
* 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 reason free-text audit note (logged only; not otherwise interpreted or stored)
* @throws IllegalStateException if {@code leadRollover:} is not configured
* @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 reason) {
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, cfg.handoverPath(), requestedAt);
PendingRollover p = new PendingRollover(token, leadTerminal, cfg.handoverPath(), requestedAt);
pending.put(token, p);
log.info("lead-rollover: open token={} handoverPath={} reason={}", token, p.handoverPath(), reason);
log.info("lead-rollover: open token={} lead={} handoverPath={} reason={}",
token, leadTerminal, p.handoverPath(), reason);
return p;
}
/**
* Verify, then roll. <strong>Nothing else in this class — there is no timer, no heartbeat and
* no timeout anywhere in this feature — can ever call {@code agents.send}; this method is the
* only path that clears a pane.</strong>
* 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: the operator-confirmation gate, then the three handover-file checks (exists, not
* empty, fresh — see {@link #checkHandover}), then the roll itself: {@code /clear}, wait up to
* {@code clearSettleSeconds} for the pane to report an injectable state again, then send
* {@code bootstrapText}. The first failing check is returned. {@code token} stays pending on
* every refusal (so a caller can fix the problem — e.g. rewrite the handover file — and retry
* with the same token) and is consumed only once the roll actually succeeds.
* <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 RollResult confirm(String token, boolean operatorConfirmed) {
public RollDecision confirm(String callerTerminal, String token, boolean operatorConfirmed) {
FleetConfig.LeadRollover cfg = configSupplier.get();
if (cfg == null) {
return RollResult.refused(RefusalReason.NOT_CONFIGURED, "leadRollover: is not configured");
return RollDecision.refused(RefusalReason.NOT_CONFIGURED, "leadRollover: is not configured");
}
PendingRollover p = pending.get(token);
if (p == null) {
return RollResult.refused(RefusalReason.UNKNOWN_TOKEN,
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 RollResult.refused(RefusalReason.OPERATOR_NOT_CONFIRMED,
return RollDecision.refused(RefusalReason.OPERATOR_NOT_CONFIRMED,
"requireOperatorConfirm is true and operatorConfirmed was false");
}
RollResult docCheck = checkHandover(p, cfg);
RollDecision docCheck = checkHandover(p, cfg);
if (docCheck != null) {
return docCheck;
}
Optional<String> lead = primaryRegistry.primaryTerminal();
if (lead.isEmpty()) {
return RollResult.refused(RefusalReason.LEAD_UNKNOWN, "no lead terminal is known to clear");
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;
}
String leadTerminal = lead.get();
// 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(leadTerminal, "/clear");
boolean settled = waitUntilInjectable(leadTerminal, cfg.clearSettleSeconds());
if (!settled) {
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={})",
leadTerminal, cfg.clearSettleSeconds(), token);
return RollResult.refused(RefusalReason.CLEAR_DID_NOT_SETTLE,
"pane did not become injectable within " + cfg.clearSettleSeconds()
+ "s after /clear; bootstrapText was not sent");
lead, cfg.clearSettleSeconds(), p.token());
return;
}
agents.send(leadTerminal, cfg.bootstrapText());
pending.remove(token);
log.info("lead-rollover: rolled token={} lead={}", token, leadTerminal);
return RollResult.success();
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} */
@@ -238,10 +316,10 @@ public final class LeadRollover {
*
* @return the first failing check's refusal, or {@code null} when all three pass
*/
private RollResult checkHandover(PendingRollover p, FleetConfig.LeadRollover cfg) {
private RollDecision checkHandover(PendingRollover p, FleetConfig.LeadRollover cfg) {
Path path = Path.of(p.handoverPath());
if (!Files.exists(path)) {
return RollResult.refused(RefusalReason.HANDOVER_MISSING,
return RollDecision.refused(RefusalReason.HANDOVER_MISSING,
"handover file " + p.handoverPath() + " does not exist");
}
long size;
@@ -253,18 +331,18 @@ public final class LeadRollover {
throw new UncheckedIOException("failed to stat handover file " + p.handoverPath(), e);
}
if (size == 0) {
return RollResult.refused(RefusalReason.HANDOVER_EMPTY,
return RollDecision.refused(RefusalReason.HANDOVER_EMPTY,
"handover file " + p.handoverPath() + " is empty");
}
if (mtimeMillis <= p.requestedAtMillis()) {
return RollResult.refused(RefusalReason.HANDOVER_STALE,
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 RollResult.refused(RefusalReason.HANDOVER_STALE,
return RollDecision.refused(RefusalReason.HANDOVER_STALE,
"handover file " + p.handoverPath() + " is " + TimeUnit.MILLISECONDS.toSeconds(ageMillis)
+ "s old, older than maxDocAgeSeconds=" + cfg.maxDocAgeSeconds());
}
@@ -273,9 +351,12 @@ public final class LeadRollover {
/**
* Poll {@link AgentControl#status} until {@code target} reports an injectable state, bounded by
* {@code settleSeconds}. 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.
* {@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);
@@ -15,9 +15,9 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
*
* <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 its
* arguments with something that silently compiles, e.g. {@code primaryRegistry} swapped for
* {@code null}, or the whole assignment swapped for a bare {@code null} literal — leaves every
* 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.
*
@@ -49,12 +49,14 @@ class FleetdLeadRolloverWiringTest {
void mainStillCallsTheLeadRolloverFactory() throws Exception {
String source = fleetdSource();
assertTrue(source.contains(
"LeadRollover leadRollover = leadRollover(cfg, primaryRegistry, router.leadAgents(), config);"),
"LeadRollover leadRollover = leadRollover(cfg, router.leadAgents(), config);"),
"Fleetd.main must still assign `LeadRollover leadRollover = leadRollover(cfg, "
+ "primaryRegistry, router.leadAgents(), config);`. Dropping this call, or swapping "
+ "one of its arguments for something that still compiles (e.g. null in place of "
+ "primaryRegistry), leaves every behavioural test green — this source check is what "
+ "must go red instead.");
+ "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
@@ -103,7 +103,7 @@ class FleetConfigWithDefaultsPreservesEveryComponentTest {
// 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, "read the handover file"));
"/handover/guard.md", true, 3600, 20, 20, "read the handover file"));
assertNamesMatchComponents(v);
return v;
}
@@ -1,9 +1,11 @@
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.mcp.PrimaryRegistry;
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;
@@ -17,8 +19,9 @@ import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.*;
/**
* fleetd #480 Unit A — the lead-rollover executor. Six hard requirements, each pinned by a test
* named after the rule it protects (see class javadoc on {@link LeadRollover}):
* 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>
@@ -26,22 +29,38 @@ import static org.junit.jupiter.api.Assertions.*;
* {@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 (the {@code Fleetd.java} source-text pin) lives in
* 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, "read the handover file");
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,
return new FleetConfig.LeadRollover(handoverPath, requireOperatorConfirm, 3600, 20, 20,
"read the handover file");
}
@@ -49,11 +68,10 @@ class LeadRolloverTest {
return millis::get;
}
private static LeadRollover newRollover(FakeHerdr herdr, FleetConfig.LeadRollover config,
private static LeadRollover newRollover(HerdrClient herdr, FleetConfig.LeadRollover config,
LongSupplier nowMillis) {
AgentControl agents = new AgentControl(herdr);
PrimaryRegistry registry = new PrimaryRegistry("term_a");
return new LeadRollover(registry, agents, () -> config, nowMillis, () -> { });
return new LeadRollover(agents, () -> config, nowMillis, () -> { }, Runnable::run);
}
private Path writeHandover(String content) throws IOException {
@@ -62,6 +80,11 @@ class LeadRolloverTest {
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
@@ -74,7 +97,7 @@ class LeadRolloverTest {
assertNull(none, "sanity: an absent block really is null");
}
// ---- 2. only confirm() can roll --------------------------------------------------------
// ---- 2. only confirm() can (eventually) roll -------------------------------------------
@Test
@DisplayName("[HARD REQ 2] nothing but an explicit confirm() call can ever roll a pane")
@@ -84,12 +107,13 @@ class LeadRolloverTest {
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, cfg(handover.toString()), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open("context is full");
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
assertNotNull(pending.token());
assertTrue(herdr.calls.stream().noneMatch(c -> "agent.prompt".equals(c.method())),
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 confirm() may");
+ "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 ---------------------------------
@@ -102,12 +126,12 @@ class LeadRolloverTest {
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, cfg(missing.toString()), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open("context is full");
LeadRollover.RollResult result = rollover.confirm(pending.token(), true);
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertFalse(result.rolled());
assertEquals(LeadRollover.RefusalReason.HANDOVER_MISSING, result.reason());
assertTrue(herdr.calls.stream().noneMatch(c -> "agent.prompt".equals(c.method())));
assertFalse(decision.accepted());
assertEquals(LeadRollover.RefusalReason.HANDOVER_MISSING, decision.reason());
assertEquals(0, promptCallCount(herdr));
}
@Test
@@ -118,11 +142,11 @@ class LeadRolloverTest {
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, cfg(empty.toString()), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open("context is full");
LeadRollover.RollResult result = rollover.confirm(pending.token(), true);
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertFalse(result.rolled());
assertEquals(LeadRollover.RefusalReason.HANDOVER_EMPTY, result.reason());
assertFalse(decision.accepted());
assertEquals(LeadRollover.RefusalReason.HANDOVER_EMPTY, decision.reason());
}
@Test
@@ -135,16 +159,16 @@ class LeadRolloverTest {
// 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, "text");
new FleetConfig.LeadRollover(handover.toString(), true, 1 /*maxDocAgeSeconds*/, 20, 20, "text");
LeadRollover rollover = newRollover(herdr, config, fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open("context is full");
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.RollResult result = rollover.confirm(pending.token(), true);
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertFalse(result.rolled());
assertEquals(LeadRollover.RefusalReason.HANDOVER_STALE, result.reason());
assertFalse(decision.accepted());
assertEquals(LeadRollover.RefusalReason.HANDOVER_STALE, decision.reason());
}
// ---- 4. requireOperatorConfirm ---------------------------------------------------------
@@ -157,13 +181,13 @@ class LeadRolloverTest {
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, cfg(handover.toString(), true), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open("context is full");
LeadRollover.RollResult result = rollover.confirm(pending.token(), false);
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), false);
assertFalse(result.rolled());
assertEquals(LeadRollover.RefusalReason.OPERATOR_NOT_CONFIRMED, result.reason());
assertTrue(herdr.calls.stream().noneMatch(c -> "agent.prompt".equals(c.method())),
"must refuse before ever touching the handover file or sending /clear");
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 --------------------------------------------------------
@@ -171,7 +195,7 @@ class LeadRolloverTest {
@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();
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
@@ -181,17 +205,74 @@ class LeadRolloverTest {
FleetConfig.LeadRollover config = cfg(handover.toString());
LeadRollover rollover = newRollover(herdr, config, fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open("context is full");
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.RollResult result = rollover.confirm(pending.token(), true);
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(result.rolled(), "expected a successful roll; got refusal: " + result.reason()
+ " (" + result.detail() + ")");
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 ---------
@@ -204,13 +285,13 @@ class LeadRolloverTest {
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, cfg(handover.toString()), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open("context is full");
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.RollResult result = rollover.confirm(pending.token(), true);
assertFalse(result.rolled());
assertEquals(LeadRollover.RefusalReason.UNKNOWN_TOKEN, result.reason());
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertFalse(decision.accepted());
assertEquals(LeadRollover.RefusalReason.UNKNOWN_TOKEN, decision.reason());
}
@Test
@@ -218,35 +299,59 @@ class LeadRolloverTest {
void openThrowsWhenNotConfigured() {
FakeHerdr herdr = new FakeHerdr();
AgentControl agents = new AgentControl(herdr);
PrimaryRegistry registry = new PrimaryRegistry("term_a");
LeadRollover rollover = new LeadRollover(registry, agents, () -> null, () -> 1_000L, () -> { });
LeadRollover rollover = new LeadRollover(agents, () -> null, () -> 1_000L, () -> { }, Runnable::run);
assertThrows(IllegalStateException.class, () -> rollover.open("context is full"));
assertThrows(IllegalStateException.class, () -> rollover.open(LEAD, "context is full"));
}
@Test
@DisplayName("a pane that never reports injectable within clearSettleSeconds refuses with CLEAR_DID_NOT_SETTLE and never sends bootstrapText")
void clearThatNeverSettlesRefusesAndNeverSendsBootstrapText() throws IOException {
@DisplayName("open() throws IllegalArgumentException when leadTerminal is null or blank")
void openThrowsWhenLeadTerminalIsBlank() throws IOException {
FakeHerdr herdr = new FakeHerdr();
herdr.agentStatus("working"); // never becomes injectable
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, 1 /*clearSettleSeconds*/, "boot text");
// The clock must ADVANCE across waitUntilInjectable's poll loop, or a bounded loop against a
// frozen clock never reaches its own deadline. Each poll (open()'s stamp, the deadline
// computation, and every loop check) draws from the same supplier, so a fixed step per call
// reliably crosses the 1s settle bound within a few iterations without an infinite loop.
new FleetConfig.LeadRollover(handover.toString(), true, 3600, 20, 1 /*clearSettleSeconds*/, "boot text");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, config, () -> clock.addAndGet(500));
LeadRollover rollover = newRollover(flipsAfterClear, config, () -> clock.addAndGet(500));
LeadRollover.PendingRollover pending = rollover.open("context is full");
LeadRollover.RollResult result = rollover.confirm(pending.token(), true);
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertFalse(result.rolled());
assertEquals(LeadRollover.RefusalReason.CLEAR_DID_NOT_SETTLE, result.reason());
long bootstrapSends = herdr.calls.stream()
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 -> c.params().toString().contains("boot text"))
.filter(c -> String.valueOf(c.params()).contains("boot text"))
.count();
assertEquals(0, bootstrapSends, "bootstrapText must never be sent when /clear did not settle");
}
@@ -260,17 +365,18 @@ class LeadRolloverTest {
FleetConfig.LeadRollover config = cfg(handover.toString());
LeadRollover rollover = newRollover(herdr, config, fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open("context is full");
LeadRollover.RollResult result = rollover.confirm(pending.token(), true);
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(result.rolled(), "expected success; got: " + result.reason() + " / " + result.detail());
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 on success — a second confirm() with the same token is UNKNOWN_TOKEN
LeadRollover.RollResult again = rollover.confirm(pending.token(), true);
// 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());
}
}