Compare commits

...

9 Commits

Author SHA1 Message Date
Dai Ha 274afafde6 fleetd #498: awaitHerdr distinguishes deadline-passed from interrupted, with measured elapsed time
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Successful in 1m51s
- awaitHerdr now returns a HerdrAwaitOutcome(HerdrWaitResult, elapsedNanos) instead of a bare
  boolean, so 'the wait budget genuinely ran out' and 'the waiting thread was interrupted' are
  two distinct, named states instead of the same false (fleetd #497's shape).
- awaitHerdr takes the clock (LongSupplier) and the per-poll sleep (Runnable) as required
  parameters, with no defaulted overload (fleetd #415), so a test can drive it.
- The startup call site is extracted into logHerdrWaitOutcomeAndShouldReap, since main() itself
  cannot be driven from a unit test; it logs a distinct message per outcome, always printing the
  measured elapsed time next to the configured budget, never the budget alone.
- Adds FleetdAwaitHerdrTest covering the seam (all three outcomes, plus the preserved interrupt
  flag) and the call site (the three distinct log messages), using ListAppender.
2026-09-12 09:54:59 +07:00
Dai Ha 60fa958da2 docs: a relative handoverPath lands in the LEAD's repo, not fleetd's (#491)
CI / contract (push) Successful in 1m29s
CI / build (push) Successful in 1m34s
On this host the lead's cwd IS the fleetd checkout, so fleetd/.gitignore
protects the handover file and the distinction is invisible. On fleet01 the
lead works in /home/ltms/LTMS/kb while fleetd sits in a different directory,
and that repo has no .handover rule — measured 2026-09-12.

Tell the lead to check its own workspace's .gitignore before writing, and to
report it rather than committing the file or editing someone's .gitignore.

Refs fleetd #491, #487, #480.
2026-09-12 08:46:17 +07:00
Dai Ha 9950361bc9 docs: #489 is fixed and deployed, but criterion 6 is still unmet
CI / contract (push) Successful in 1m12s
CI / build (push) Successful in 1m33s
The skill said the roll 'does nothing until #489 is merged and redeployed'.
Both have happened, so that sentence now reads as a green light. It is not
one: no roll has bootstrapped a fresh session end to end yet. Say that
plainly, and tell the lead to warn the operator before confirming.

Refs fleetd #489, #480.
2026-09-12 07:55:41 +07:00
ltms 9f74b6619a Merge #490: nudge the /clear submit keystroke before bootstrapText (fleetd #489)
CI / contract (push) Successful in 1m14s
CI / build (push) Successful in 1m30s
Fixes the live defect measured on 2026-09-12: the roll joined /clear and
bootstrapText into one line and Claude Code refused it as
"Unknown command: /clearFresh".

Verified by the lead before merge:
- mvn clean install: Tests run: 1677, Failures: 0, BUILD SUCCESS
- LeadRolloverTest baseline: 29 tests, 0 failures
- mutation PICKUP_GRACE_POLLS 8 -> 1: 1 failure (the regression test)
- mutation deleting the !clearSettled guard: 3 failures, 2 pre-existing tests
- mutation replacing waitForClearPickupAndSettle with "return true": 5 failures,
  including the strengthened pickup test
- positive control after each restore: 29 tests, 0 failures
- CI run 1738 on f687046: success

A separate mutation of the FIRST gate hung the suite instead of failing it.
That is fleetd #486, not a regression here; the reproduction is recorded there.
2026-09-12 02:52:52 +02:00
Dai Ha f687046450 fleetd #489 follow-up: fix stale class javadoc, off-by-one nudge count, weak test
CI / contract (pull_request) Successful in 1m22s
CI / build (pull_request) Successful in 1m33s
Three review corrections on top of the previous commit:

1. The class javadoc's four-step continuation list (lines 47-58) was stale.
   Step 1 said "report an injectable state", but waitUntilAtTurnBoundary's
   own javadoc excludes BLOCKED - fixed to say IDLE or DONE. Step 3 still
   described the old plain re-check ("the original, pre-correction wait...
   still here") - fixed to describe what waitForClearPickupAndSettle
   actually does: nudge while unpicked-up, then wait for a real WORKING ->
   IDLE/DONE boundary, releasing rather than wedging if WORKING never shows.

2. PICKUP_GRACE_POLLS=8 bounds the number of consecutive not-yet-picked-up
   polls, not the number of nudges - the 8th poll releases instead of
   nudging again, so 8 polls produce 7 nudges. The log.info in the release
   branch and two javadoc spots said "8 nudges"; fixed all three to state
   the poll count and the nudge count separately and correctly. Behavior
   and the constant are unchanged.

3. pickupSeenStopsNudgingAndBootstrapTextIsSent asserted only promptCallCount
   and sendKeysCallCount, both of which a return-true stub also satisfies.
   Added an assertion on the already-tracked postClearGetCalls counter
   (>= 2), which only a real post-/clear poll loop can produce - this is
   what makes the test fail against a return-true mutant.
2026-09-12 07:45:36 +07:00
Dai Ha c6058652be fleetd #489: nudge the /clear submit keystroke before bootstrapText
CI / contract (pull_request) Successful in 56s
CI / build (pull_request) Successful in 1m36s
LeadRollover.runRollover's second wait (after /clear) was a no-op: it polled
for IDLE/DONE, which /clear itself never leaves since it starts no real turn,
so it always returned true on the first poll. Combined with a direct
agents.send bypassing Injector (deliberate, to avoid wedging the pane), the
submit Enter that accompanies /clear could race the paste and leave it
unsubmitted — bootstrapText then landed concatenated onto the same input
line, exactly as measured live on 2026-09-12.

Replace that second wait with waitForClearPickupAndSettle, which copies the
pickup-nudge pattern Injector already ships for its own post-turn /clear
housekeeping (fleetd #306): nudge agents.submit while the pane hasn't
reported WORKING yet, release after PICKUP_GRACE_POLLS=8 nudges rather than
wedge, and require a real WORKING -> IDLE/DONE boundary once a pickup is
observed. BLOCKED stays excluded from both the nudge and the boundary check,
same as the (unchanged) first wait — a paused live turn is not settled, and
nudging Enter into an open prompt could wrongly answer it.

Adds four tests to LeadRolloverTest covering the paste-race regression
(nudge ordered between /clear and bootstrapText), a confirmed pickup, a
deadline expiry with no boundary ever reached, and a throwing submit().
2026-09-12 07:29:26 +07:00
Dai Ha 2a95da2ff4 docs: the handover skill must warn that the automatic roll is broken (#489)
CI / contract (push) Successful in 1m23s
CI / build (push) Successful in 1m53s
Section 11 said the bootstrap prompt landing was 'not yet proven
end-to-end'. It is now measured failing: the first real roll joined
/clear and bootstrapText into one line. Nothing was cleared, so the
failure is safe, but a lead that reads the old wording would reach for
fleet_handover expecting it to work.

Refs fleetd #489, #480.
2026-09-12 07:21:38 +07:00
Dai Ha 7f9137fcb6 fleetd #480: ignore the handover file, and note the absolute-path guarantee in the handover skill
CI / contract (push) Successful in 47s
CI / build (push) Successful in 1m52s
The lead rollover handover file now lives inside the workspace, at the relative
path fleetd.yaml's leadRollover.handoverPath names. It is a snapshot of one
moment's live state, so it must never enter git history.

The handover skill also now says the handoverPath fleetd hands back is always
absolute, even when the configured value is relative — a lead that resolves it
itself can pick a different file from the one the daemon checks.
2026-09-12 05:46:34 +07:00
Dai Ha 3bf3968bc7 Merge #487: resolve a relative leadRollover.handoverPath against the calling lead's workspace (fleetd #480 follow-up) 2026-09-12 05:43:12 +07:00
6 changed files with 652 additions and 42 deletions
+22 -4
View File
@@ -135,6 +135,20 @@ refused.**
1. **`fleet_handover{action: "open", reason: "<why now>"}`.** It returns a `token` and the
`handoverPath` you must write to. Nothing has happened to your pane yet.
**Write to exactly that path, and do not resolve it yourself.** It is always absolute, even when
the operator configured a relative `handoverPath`: fleetd resolves a relative one against your
own workspace before it hands it to you. The daemon and your pane can run in different
directories, so a path you resolve yourself can point at a different file from the one the daemon
will check.
**Check that the path is ignored by git before you write to it (#491).** A relative
`handoverPath` resolves inside YOUR workspace, which is usually a repository — and usually not
the `fleetd` one, so an ignore rule added to `fleetd` does not protect it. Run
`grep -n handover <your workspace>/.gitignore`. No output means the file you are about to write
will show up as untracked content in that repo. The file is a snapshot of live state and must
never be committed, so tell the operator rather than committing it or silently editing their
`.gitignore`.
2. **Write the handover file at that path**, following sections 1–10 above.
3. **Ask the operator, then `fleet_handover{action: "confirm", token, operatorConfirmed: true}`.**
@@ -158,10 +172,14 @@ fails.
session.
- **The roll can still refuse after `confirm` returns**, and by then there is no caller to tell.
Those outcomes are logged only, as `lead-rollover:` lines in the daemon log.
- **If the bootstrap prompt never lands, your context is gone and no fresh session starts.** This
has not yet been proven end-to-end (see fleetd #480). The recovery is the manual path: the file
is already written, so the operator starts a session and points it at the file. That is why you
write the file before you confirm, and never the other way round.
- **The bootstrap prompt has never yet landed, and the fix is unproven (fleetd #489).** The first
real rollover, on 2026-09-12, joined `/clear` and the bootstrap text into one line and Claude Code
refused it as `Unknown command: /clearFresh`. The pane was never cleared and no context was lost,
so the failure was safe — the roll simply did nothing. PR #490 fixed the cause and is deployed,
but no roll has bootstrapped a fresh session end to end yet. **Assume it may still fail, and tell
the operator so before you confirm.** The recovery is the same either way: the file is already
written, so the operator starts a session and points it at the file. That is why you write the
file before you confirm, and never the other way round.
## Writing style
+6
View File
@@ -19,3 +19,9 @@
fleetd.out
fleetd/fleetd.out
logs/
# fleetd #480: the lead rollover handover file. `leadRollover.handoverPath` points here, and the
# outgoing lead rewrites it on every rollover. It is a snapshot of one moment's live state —
# unpushed branches, running builds, open questions — so it is stale the moment it is written and
# has no business in git history.
.handover/
+106 -18
View File
@@ -80,6 +80,7 @@ import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Function;
import java.util.function.LongSupplier;
import java.util.function.Predicate;
import java.util.function.Supplier;
import java.util.regex.Pattern;
@@ -270,15 +271,12 @@ public final class Fleetd {
// the first thing that actually talks to herdr, so without this wait a boot-order race
// would crash the daemon into a restart loop. Wait, then degrade rather than die: serving
// with /healthz reporting "degraded" is strictly more useful than exiting.
boolean herdrUp = awaitHerdr(herdr);
HerdrAwaitOutcome herdrOutcome = awaitHerdr(herdr, System::nanoTime, Fleetd::sleepHerdrPoll);
boolean herdrUp = logHerdrWaitOutcomeAndShouldReap(herdrOutcome);
if (herdrUp) {
// CB-117: herdr keeps worker panes alive across a daemon restart, and their ids died
// with the previous process — reap those leaked orphans now, before we start serving.
workers.reapOrphanWorkers();
} else {
log.warn("herdr did not answer within {}s — starting anyway; /healthz will report "
+ "degraded until it comes up. Orphaned worker panes (if any) were NOT reaped.",
HERDR_WAIT_SECONDS);
}
// CB-301: authoritative session registry + lifecycle FSM on top of ClaudeCodeLauncher.
@@ -1696,12 +1694,70 @@ public final class Fleetd {
}
/**
* Poll herdr's {@code ping} until it answers or {@link #HERDR_WAIT_SECONDS} elapses (CB-504).
*
* @return true if herdr answered, false if it never did
* How {@link #awaitHerdr} ended (fleetd #498). The old code returned a bare {@code boolean},
* which collapsed two different facts onto the same {@code false}: the configured wait budget
* genuinely running out, and the waiting thread being interrupted possibly milliseconds in.
* Those need different operator messages — see {@link #logHerdrWaitOutcomeAndShouldReap} — so
* this is a third state, not a better number (the same shape fleetd #497 named). Never treat
* {@link #INTERRUPTED} as if it were {@link #DEADLINE_PASSED}: only the latter means herdr was
* actually given the full {@link #HERDR_WAIT_SECONDS} and still failed to answer.
*/
private static boolean awaitHerdr(HerdrClient herdr) {
long deadline = System.nanoTime() + HERDR_WAIT_SECONDS * 1_000_000_000L;
enum HerdrWaitResult {
/** herdr answered {@code ping} before the deadline. */
ANSWERED,
/** the configured {@link #HERDR_WAIT_SECONDS} budget elapsed with no answer. */
DEADLINE_PASSED,
/**
* the waiting thread was interrupted before the budget ran out — a different event from
* {@link #DEADLINE_PASSED} and must never be reported as "did not answer within Ns".
*/
INTERRUPTED
}
/**
* The outcome of one {@link #awaitHerdr} call, carrying the MEASURED elapsed wait time
* alongside {@link #result}. {@code elapsedNanos} is always measured against the {@code nanos}
* supplier passed to {@link #awaitHerdr} — never assume it equals the configured budget, the
* same defect fleetd #494 already fixed once in {@code LeadRollover}.
*/
record HerdrAwaitOutcome(HerdrWaitResult result, long elapsedNanos) {}
/**
* The real per-poll wait {@link #main} passes to {@link #awaitHerdr}: sleep
* {@link #HERDR_WAIT_POLL_MILLIS}, and on interruption re-set the thread's interrupt flag
* rather than throwing — {@link #awaitHerdr} detects an interruption by checking {@link
* Thread#isInterrupted()} right after this returns, so a poller that swallowed the flag
* instead of restoring it would make that check silently miss the interruption.
*/
private static void sleepHerdrPoll() {
try {
Thread.sleep(HERDR_WAIT_POLL_MILLIS);
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
}
}
/**
* Poll herdr's {@code ping} until it answers, the configured {@link #HERDR_WAIT_SECONDS}
* budget elapses, or the waiting thread is interrupted (CB-504, fleetd #498).
*
* <p>{@code nanos} and {@code poller} are required parameters with no defaulted overload
* (fleetd #415's shape: a defaulted overload is a silent survivor a green suite would vouch
* for) — the previous version read {@link System#nanoTime()} and called {@link Thread#sleep}
* directly, so nothing could drive it from a test. The one production call site in {@link
* #main} passes {@code System::nanoTime} and {@link #sleepHerdrPoll}.
*
* @param nanos a monotonic elapsed-time clock, e.g. {@code System::nanoTime} — never a
* wall-clock source, since only elapsed time (not a timestamp) is measured here
* @param poller called once per failed ping while the budget remains; must, on an
* {@link InterruptedException}, re-set the thread's interrupt flag rather than
* throw or swallow it — this method's interruption check reads that flag right
* after {@code poller.run()} returns
* @return the outcome and the measured elapsed wait time — see {@link HerdrAwaitOutcome}
*/
static HerdrAwaitOutcome awaitHerdr(HerdrClient herdr, LongSupplier nanos, Runnable poller) {
long start = nanos.getAsLong();
long deadline = start + HERDR_WAIT_SECONDS * 1_000_000_000L;
boolean waited = false;
while (true) {
try {
@@ -1709,25 +1765,57 @@ public final class Fleetd {
if (waited) {
log.info("herdr is up");
}
return true;
return new HerdrAwaitOutcome(HerdrWaitResult.ANSWERED, nanos.getAsLong() - start);
} catch (HerdrException e) {
if (System.nanoTime() >= deadline) {
return false;
if (nanos.getAsLong() >= deadline) {
return new HerdrAwaitOutcome(HerdrWaitResult.DEADLINE_PASSED, nanos.getAsLong() - start);
}
if (!waited) {
log.info("waiting up to {}s for the herdr socket…", HERDR_WAIT_SECONDS);
waited = true;
}
try {
Thread.sleep(HERDR_WAIT_POLL_MILLIS);
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
return false;
poller.run();
if (Thread.currentThread().isInterrupted()) {
return new HerdrAwaitOutcome(HerdrWaitResult.INTERRUPTED, nanos.getAsLong() - start);
}
}
}
}
/**
* Log the right message for {@code outcome} — never the configured {@link #HERDR_WAIT_SECONDS}
* budget alone, always the measured elapsed time next to it — and say whether {@link #main}
* should now reap orphan worker panes (fleetd #498).
*
* <p>Extracted out of {@link #main} so this decision is drivable from a test: {@link #main}
* boots the whole daemon and cannot itself be run in a unit test, but this is the exact,
* unmodified code {@link #main} calls for the decision, not a re-derivation of it.
*
* @return true only for {@link HerdrWaitResult#ANSWERED} — orphan workers are reaped only
* then, exactly as before this ticket
*/
static boolean logHerdrWaitOutcomeAndShouldReap(HerdrAwaitOutcome outcome) {
long elapsedMillis = TimeUnit.NANOSECONDS.toMillis(outcome.elapsedNanos());
if (outcome.result() == HerdrWaitResult.ANSWERED) {
return true;
}
if (outcome.result() == HerdrWaitResult.DEADLINE_PASSED) {
log.warn("herdr did not answer within the configured wait (configured={}s elapsed={}ms) "
+ "— starting anyway; /healthz will report degraded until it comes up. Orphaned "
+ "worker panes (if any) were NOT reaped.",
HERDR_WAIT_SECONDS, elapsedMillis);
return false;
}
// HerdrWaitResult.INTERRUPTED — a different fact from DEADLINE_PASSED (fleetd #498): the
// wait was cut short, not exhausted, and must never be reported as "did not answer within
// Ns" — that claim would be false and would send an operator to debug herdr for nothing.
log.warn("herdr wait was interrupted before the configured wait ran out (configured={}s "
+ "elapsed={}ms) — starting anyway; /healthz will report degraded until it comes "
+ "up. Orphaned worker panes (if any) were NOT reaped.",
HERDR_WAIT_SECONDS, elapsedMillis);
return false;
}
private Fleetd() {
}
}
@@ -44,15 +44,18 @@ import java.util.function.Supplier;
* {@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>wait for the lead's own pane to report a real turn boundary — {@code IDLE} or {@code
* DONE}, never merely {@code BLOCKED} — 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>wait for {@code /clear} to be picked up and settle, bounded by {@code clearSettleSeconds}
* (fleetd #489: no longer a plain re-check of the same boundary — {@code /clear} starts no
* turn of its own, so this instead nudges the submit keystroke while no pickup has been seen,
* then waits for a real {@code WORKING} → {@code IDLE}/{@code DONE} boundary once one has;
* see {@link #waitForClearPickupAndSettle})</li>
* <li>{@code agents.send(lead, cfg.bootstrapTextFor(p.handoverPath()))}</li>
* </ol>
* A {@link #confirm} that returns {@link RollDecision#approved()} therefore means <em>"every gate
@@ -111,6 +114,17 @@ 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;
/**
* How many consecutive not-yet-picked-up polls {@link #waitForClearPickupAndSettle} allows
* before releasing rather than wedging the roll — the same constant and the same
* release-not-wedge choice {@link dev.ltms.fleet.inject.Injector} already makes for its own
* post-turn {@code /clear} housekeeping (fleetd #306). <strong>This bounds the number of
* consecutive polls, not the number of nudges:</strong> the first {@code PICKUP_GRACE_POLLS - 1}
* of those polls each send a nudge, and the {@code PICKUP_GRACE_POLLS}th releases instead of
* nudging again — so 8 polls produce 7 nudges, not 8.
*/
static final int PICKUP_GRACE_POLLS = 8;
/**
* One request opened by {@link #open}, pending its {@link #confirm} (or {@link #cancel}).
*
@@ -381,7 +395,7 @@ public final class LeadRollover {
// /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 = waitUntilAtTurnBoundary(lead, cfg.clearSettleSeconds());
boolean clearSettled = waitForClearPickupAndSettle(lead, cfg.clearSettleSeconds());
if (!clearSettled) {
log.warn("lead-rollover: pane {} did not reach a turn boundary (IDLE or DONE) within {}s "
+ "after /clear — NOT sending bootstrapText (token={})",
@@ -440,11 +454,14 @@ public final class LeadRollover {
/**
* Poll {@link AgentControl#status} until {@code target} reports a real turn boundary — {@link
* AgentStatus#IDLE} or {@link AgentStatus#DONE} — 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
* AgentStatus#IDLE} or {@link AgentStatus#DONE} — bounded by {@code settleSeconds}. Used once by
* {@link #runRollover}, to wait for the CALLING turn's own pane to settle before {@code /clear}
* is ever sent at all — the {@code turnSettleSeconds} gate that makes this correction safe. The
* SECOND wait, after {@code /clear}, is {@link #waitForClearPickupAndSettle} instead (fleetd
* #489) — a plain boundary check is not enough there, because {@code /clear} starts no turn of
* its own, so this method would (wrongly) report "settled" on its very first poll whether or not
* {@code /clear} was actually picked up. 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.
*
* <p><strong>Deliberately not {@link AgentStatus#injectable()}.</strong> {@code injectable()}
@@ -452,12 +469,12 @@ public final class LeadRollover {
* turn" — and it accepts {@link AgentStatus#BLOCKED} for that purpose, because a pane paused on
* an approval prompt is safe to queue a message behind. This class asks a stricter question —
* "has the turn actually ended" — and {@code BLOCKED} answers no: it is a live turn that is
* merely paused, not one that has finished. Reusing {@code injectable()} here would let both
* waits fire into an open approval prompt mid-turn (the first wait would send {@code /clear}
* while the lead's own {@code confirm()}-calling turn is still live and paused on a prompt; the
* second would send {@code bootstrapText} the same way after {@code /clear}) — exactly the
* live-context-destroying failure the {@code turnSettleSeconds} gate exists to prevent. Do not
* "simplify" this back to {@code injectable()}.
* merely paused, not one that has finished. Reusing {@code injectable()} here would let this
* wait fire {@code /clear} while the lead's own {@code confirm()}-calling turn is still live and
* paused on a prompt — exactly the live-context-destroying failure the {@code turnSettleSeconds}
* gate exists to prevent. Do not "simplify" this back to {@code injectable()}. ({@link
* #waitForClearPickupAndSettle} keeps the same exclusion of {@code BLOCKED}, for the same
* reason, on the second wait.)
*/
private boolean waitUntilAtTurnBoundary(String target, int settleSeconds) {
long deadline = nowMillis.getAsLong() + TimeUnit.SECONDS.toMillis(settleSeconds);
@@ -477,4 +494,95 @@ public final class LeadRollover {
}
return false;
}
/**
* The SECOND wait in {@link #runRollover} — after {@code /clear} has been sent, waits for it to
* settle, bounded by {@code settleSeconds}. <strong>fleetd #489 — the paste-race fix.</strong>
* {@code /clear} does not start a real turn of its own, so a pane with no submit race simply
* stays {@link AgentStatus#IDLE} the whole time: {@link #waitUntilAtTurnBoundary} would (wrongly)
* call that "settled" on its very first poll, whether or not the {@code /clear} Enter actually
* landed. That was Fault 1, measured live on 2026-09-12 — the second gate was a no-op, so a
* {@code bootstrapText} send followed immediately, racing Fault 2: {@link AgentControl#submit}'s
* own javadoc already records that the submit accompanying a delivery "can race the paste —
* especially right as the worker's TUI becomes interactive — leaving the text unsubmitted"
* (CB-113). Because {@code runRollover} deliberately bypasses {@code Injector} for {@code
* /clear} (see this class's javadoc), it inherited none of {@code Injector}'s nudging — so the
* lost {@code /clear} Enter sat in the input box and {@code bootstrapText} was typed right after
* it, landing as one concatenated line.
*
* <p>This method copies the pickup-nudge pattern {@link dev.ltms.fleet.inject.Injector} already
* ships for exactly this, on its own post-turn {@code /clear} housekeeping (fleetd #306; see
* {@code Injector.java:288-340} and {@code Injector.java:437-442}):
* <ul>
* <li>an {@link AgentStatus#WORKING} sample means {@code /clear} was picked up as a real
* turn;</li>
* <li>until that happens, each poll that still reports {@link AgentStatus#IDLE} or {@link
* AgentStatus#DONE} re-sends the submit keystroke ({@link AgentControl#submit}) to nudge
* the raced Enter — for the first {@code PICKUP_GRACE_POLLS - 1} of {@link
* #PICKUP_GRACE_POLLS} consecutive such polls (i.e. {@code PICKUP_GRACE_POLLS - 1}
* nudges: 7, not 8, given {@code PICKUP_GRACE_POLLS = 8}). A second Enter on an empty
* Claude Code prompt is a no-op, so repeating it is safe;</li>
* <li>the {@code PICKUP_GRACE_POLLS}th consecutive such poll, with {@code WORKING} still never
* observed, releases rather than wedges the roll instead of nudging again — the same
* choice {@code Injector} makes — and returns {@code true} anyway, logged at {@code info}
* so an operator can see which path ran;</li>
* <li>once {@code WORKING} has been observed, nudging stops and this instead waits for a real
* {@code working → IDLE/DONE} completion boundary before returning {@code true}.</li>
* </ul>
*
* <p><strong>{@link AgentStatus#BLOCKED} is deliberately excluded from both the nudge and the
* boundary check</strong> — the same reasoning as {@link #waitUntilAtTurnBoundary}'s own
* javadoc: a paused live turn is not a settled one, and re-sending Enter into an open approval
* prompt could wrongly answer it. A {@code BLOCKED} sample (or an unreadable/{@link
* AgentStatus#UNKNOWN} one) simply keeps this polling, with no nudge and no release, until either
* a real boundary is reached or {@code settleSeconds} runs out.
*
* <p>{@link AgentControl#submit} can itself throw; a {@link RuntimeException} from it is
* swallowed and logged at {@code debug}, exactly like {@code Injector.java:437-442} — a failed
* nudge must not abort the roll.
*
* @return {@code true} once {@code /clear} has settled, or once the nudge budget was exhausted
* with no pickup ever observed (released rather than wedged); {@code false} if {@code
* settleSeconds} elapses first — the caller must NOT send {@code bootstrapText} in that
* case, exactly as before this fix
*/
private boolean waitForClearPickupAndSettle(String target, int settleSeconds) {
long deadline = nowMillis.getAsLong() + TimeUnit.SECONDS.toMillis(settleSeconds);
boolean pickedUp = false; // a WORKING sample has been observed since /clear was sent
int idlePollsAwaitingPickup = 0;
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 after "
+ "/clear: {}", target, e.toString());
status = null;
}
if (status == AgentStatus.WORKING) {
pickedUp = true;
} else if (status == AgentStatus.IDLE || status == AgentStatus.DONE) {
if (pickedUp) {
return true; // a real WORKING -> IDLE/DONE completion boundary
}
if (++idlePollsAwaitingPickup >= PICKUP_GRACE_POLLS) {
log.info("lead-rollover: /clear on {} was never observed as WORKING after {} "
+ "consecutive IDLE/DONE polls ({} of those were nudged) — "
+ "releasing rather than wedging the roll",
target, PICKUP_GRACE_POLLS, PICKUP_GRACE_POLLS - 1);
return true;
}
try {
agents.submit(target); // nudge a raced Enter (CB-113) so /clear actually submits
} catch (RuntimeException e) {
log.debug("lead-rollover: resubmit to {} failed (will retry next poll): {}",
target, e.getMessage());
}
}
// AgentStatus.BLOCKED or UNKNOWN (or an unreadable status, above): neither a pickup
// signal nor a boundary — keep polling without nudging or releasing.
settleSleeper.run();
}
return false;
}
}
@@ -0,0 +1,212 @@
package dev.ltms.fleet;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import com.fasterxml.jackson.databind.JsonNode;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.HerdrException;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.junit.jupiter.api.Assertions.fail;
/**
* fleetd #498: {@code Fleetd.awaitHerdr} used to return a bare {@code boolean}, collapsing "the
* configured wait budget genuinely ran out" and "the waiting thread was interrupted, possibly
* milliseconds in" onto the same {@code false} — and the caller's log line printed only the
* configured budget, never how long the wait actually ran. This class covers both halves of the
* fix:
* <ul>
* <li>the seam — {@link Fleetd#awaitHerdr} itself, driven with an injected clock and a stub
* {@link HerdrClient}, one test per {@link Fleetd.HerdrWaitResult};</li>
* <li>the call site — {@link Fleetd#logHerdrWaitOutcomeAndShouldReap}, the exact decision {@code
* main} calls (extracted here because {@code main} itself boots the whole daemon and cannot
* be driven from a unit test), pinning the three distinct log messages it emits.</li>
* </ul>
* Every expected message below is a plain literal, not built from {@code HERDR_WAIT_SECONDS} or
* any other production constant — a test that derives its expectation the way the code does
* cannot see a change to either (fleetd #496's identical trap).
*/
class FleetdAwaitHerdrTest {
// ---- the seam: Fleetd.awaitHerdr ----------------------------------------------------------
@Test
void answeredReturnsImmediatelyWithZeroElapsedAndNeverPolls() {
HerdrStub herdr = new HerdrStub(0); // succeeds on the very first call
LongSupplier clock = fixedClock(1_000L);
AtomicBoolean polled = new AtomicBoolean(false);
Runnable poller = () -> polled.set(true);
Fleetd.HerdrAwaitOutcome outcome = Fleetd.awaitHerdr(herdr, clock, poller);
assertEquals(Fleetd.HerdrWaitResult.ANSWERED, outcome.result());
assertEquals(0L, outcome.elapsedNanos(), "a fixed clock must measure zero elapsed time");
assertFalse(polled.get(), "herdr answering on the first try must never poll");
}
@Test
void deadlinePassedIsMeasuredNotAssumed() {
HerdrStub herdr = new HerdrStub(-1); // never succeeds
// call order inside awaitHerdr: start, then per failed attempt: deadline-check, elapsed-calc
ScriptedClock clock = new ScriptedClock(0L, 30_500_000_000L, 30_500_000_000L);
Runnable poller = () -> fail("the deadline was already exceeded on the first attempt — must not poll");
Fleetd.HerdrAwaitOutcome outcome = Fleetd.awaitHerdr(herdr, clock, poller);
assertEquals(Fleetd.HerdrWaitResult.DEADLINE_PASSED, outcome.result());
assertEquals(30_500_000_000L, outcome.elapsedNanos(),
"elapsed must be the MEASURED clock delta, not the configured budget");
}
@Test
void interruptedIsDistinctFromDeadlinePassedAndPreservesTheInterruptFlag() {
HerdrStub herdr = new HerdrStub(-1); // never succeeds
// start=0, deadline-check returns 500ms (well under the 30s budget) -> not deadline-passed,
// then the poller interrupts, and the elapsed-calc call returns 750ms.
ScriptedClock clock = new ScriptedClock(0L, 500_000_000L, 750_000_000L);
Runnable poller = () -> Thread.currentThread().interrupt();
try {
Fleetd.HerdrAwaitOutcome outcome = Fleetd.awaitHerdr(herdr, clock, poller);
assertEquals(Fleetd.HerdrWaitResult.INTERRUPTED, outcome.result());
assertEquals(750_000_000L, outcome.elapsedNanos(),
"elapsed must be measured even when the wait ends via interruption, not the deadline");
assertTrue(Thread.currentThread().isInterrupted(),
"the interrupt flag the old code re-set must still be set on return");
} finally {
Thread.interrupted(); // clear it so it cannot leak into another test on this thread
}
}
// ---- the call site: Fleetd.logHerdrWaitOutcomeAndShouldReap -------------------------------
@Test
void answeredLogsNothingAndSaysReap() {
ListAppender<ILoggingEvent> events = attach();
try {
boolean shouldReap = Fleetd.logHerdrWaitOutcomeAndShouldReap(
new Fleetd.HerdrAwaitOutcome(Fleetd.HerdrWaitResult.ANSWERED, 0L));
assertTrue(shouldReap, "only ANSWERED should tell main to reap orphan workers");
assertEquals(0, events.list.size(), "the answered path logs nothing itself");
} finally {
detach(events);
}
}
@Test
void deadlinePassedLogsConfiguredAndMeasuredElapsedTogether() {
ListAppender<ILoggingEvent> events = attach();
try {
boolean shouldReap = Fleetd.logHerdrWaitOutcomeAndShouldReap(
new Fleetd.HerdrAwaitOutcome(Fleetd.HerdrWaitResult.DEADLINE_PASSED, 30_500_000_000L));
assertFalse(shouldReap, "a deadline-passed wait must not tell main to reap");
assertEquals(1, events.list.size());
ILoggingEvent event = events.list.getFirst();
assertEquals(Level.WARN, event.getLevel());
assertEquals("herdr did not answer within the configured wait (configured=30s "
+ "elapsed=30500ms) — starting anyway; /healthz will report degraded until it "
+ "comes up. Orphaned worker panes (if any) were NOT reaped.",
event.getFormattedMessage());
} finally {
detach(events);
}
}
@Test
void interruptedLogsItsOwnMessageAndNeverClaimsTheBudgetElapsed() {
ListAppender<ILoggingEvent> events = attach();
try {
// 3ms: the ticket's own example of "a few milliseconds in", not the 30s budget.
boolean shouldReap = Fleetd.logHerdrWaitOutcomeAndShouldReap(
new Fleetd.HerdrAwaitOutcome(Fleetd.HerdrWaitResult.INTERRUPTED, 3_000_000L));
assertFalse(shouldReap, "an interrupted wait must not tell main to reap");
assertEquals(1, events.list.size());
ILoggingEvent event = events.list.getFirst();
assertEquals(Level.WARN, event.getLevel());
String message = event.getFormattedMessage();
assertEquals("herdr wait was interrupted before the configured wait ran out "
+ "(configured=30s elapsed=3ms) — starting anyway; /healthz will report "
+ "degraded until it comes up. Orphaned worker panes (if any) were NOT reaped.",
message);
assertFalse(message.contains("did not answer"),
"an interrupted wait must not be reported as if herdr failed to answer within the budget");
} finally {
detach(events);
}
}
// ---- fixtures --------------------------------------------------------------------------
/** Always returns the same value, i.e. a clock that measures zero elapsed time. */
private static LongSupplier fixedClock(long value) {
return () -> value;
}
/** Returns each value in order, then repeats the last one for any call beyond the list. */
private static final class ScriptedClock implements LongSupplier {
private final long[] values;
private int index;
ScriptedClock(long... values) {
this.values = values;
}
@Override
public long getAsLong() {
long v = values[Math.min(index, values.length - 1)];
if (index < values.length - 1) {
index++;
}
return v;
}
}
/** Fails {@code failuresBeforeSuccess} times, then succeeds forever; {@code -1} never succeeds. */
private static final class HerdrStub implements HerdrClient {
private final int failuresBeforeSuccess;
private int calls;
HerdrStub(int failuresBeforeSuccess) {
this.failuresBeforeSuccess = failuresBeforeSuccess;
}
@Override
public JsonNode call(String method, Object params) throws HerdrException {
calls++;
if (failuresBeforeSuccess < 0 || calls <= failuresBeforeSuccess) {
throw new HerdrException("herdr not up yet");
}
return null;
}
@Override
public void close() {
}
}
private static ListAppender<ILoggingEvent> attach() {
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
logger.setLevel(Level.DEBUG);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
return appender;
}
private static void detach(ListAppender<ILoggingEvent> appender) {
((Logger) LoggerFactory.getLogger(Fleetd.class)).detachAppender(appender);
}
}
@@ -458,6 +458,184 @@ class LeadRolloverTest {
assertEquals(0, bootstrapSends, "bootstrapText must never be sent when /clear did not settle");
}
// ---- fleetd #489: the second wait nudges the /clear pickup instead of being a no-op --------
/** Every {@code agent.send_keys} call {@code herdr} recorded — the submit-keystroke nudge. */
private static long sendKeysCallCount(FakeHerdr herdr) {
return herdr.calls.stream().filter(c -> "agent.send_keys".equals(c.method())).count();
}
@Test
@DisplayName("[fleetd #489] a pane that stays IDLE the whole time (the paste-race case, "
+ "measured live 2026-09-12) is nudged between /clear and bootstrapText, never lets "
+ "them concatenate into one line")
void clearPickupIsNudgedBeforeBootstrapTextWhenPaneStaysIdle() throws IOException {
FakeHerdr herdr = new FakeHerdr(); // default agentStatus is "idle" throughout — no WORKING sample ever
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(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason() + " / " + decision.detail());
int clearIdx = -1;
int bootstrapIdx = -1;
int firstNudgeIdx = -1;
for (int i = 0; i < herdr.calls.size(); i++) {
FakeHerdr.Call c = herdr.calls.get(i);
if ("agent.prompt".equals(c.method()) && String.valueOf(c.params()).contains("/clear") && clearIdx < 0) {
clearIdx = i;
} else if ("agent.prompt".equals(c.method()) && String.valueOf(c.params()).contains("read the handover file")) {
bootstrapIdx = i;
} else if ("agent.send_keys".equals(c.method()) && firstNudgeIdx < 0) {
firstNudgeIdx = i;
}
}
assertTrue(clearIdx >= 0, "/clear must have been sent");
assertTrue(bootstrapIdx >= 0, "bootstrapText must have been sent");
assertTrue(firstNudgeIdx >= 0, "at least one agent.send_keys nudge must go out — the pane "
+ "never reported WORKING, so the /clear Enter may have raced the paste, and only a "
+ "re-sent Enter proves the clear rather than concatenating bootstrapText onto "
+ "whatever sits unsubmitted in the input box");
assertTrue(firstNudgeIdx > clearIdx, "the nudge must happen AFTER /clear was sent, got call "
+ "order: " + herdr.calls);
assertTrue(firstNudgeIdx < bootstrapIdx, "the nudge must happen BEFORE bootstrapText is "
+ "sent — never concatenated onto the same input line, got call order: " + herdr.calls);
}
@Test
@DisplayName("[fleetd #489] once a WORKING sample confirms /clear was picked up, nudging stops "
+ "and bootstrapText is still sent after the pane returns to IDLE")
void pickupSeenStopsNudgingAndBootstrapTextIsSent() throws IOException {
FakeHerdr fake = new FakeHerdr();
// Scripts the SECOND wait only: idle (default) until /clear is sent, then the first status
// poll after /clear reports WORKING (a confirmed pickup), and every poll after that reports
// IDLE (the completion boundary). The first wait (turnSettleSeconds) never sees this
// sequence — it passes on its own first poll, before /clear is ever sent, on the default
// "idle" status.
AtomicLong postClearGetCalls = new AtomicLong(0);
HerdrClient scriptsPickupThenIdle = new HerdrClient() {
private volatile boolean clearSent = false;
@Override
public JsonNode call(String method, Object params) throws HerdrException {
if ("agent.prompt".equals(method) && String.valueOf(params).contains("/clear")) {
clearSent = true;
}
if (clearSent && "agent.get".equals(method)) {
long n = postClearGetCalls.incrementAndGet();
fake.agentStatus(n == 1 ? "working" : "idle");
}
return fake.call(method, params);
}
@Override
public void close() {
fake.close();
}
};
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(scriptsPickupThenIdle, cfg(handover.toString()), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason() + " / " + decision.detail());
assertEquals(2, promptCallCount(fake), "a confirmed WORKING pickup followed by IDLE must "
+ "still complete the full roll — /clear then bootstrapText");
assertEquals(0, sendKeysCallCount(fake), "once WORKING was observed, nudging must stop "
+ "immediately — no agent.send_keys call should ever have been needed or sent");
assertTrue(postClearGetCalls.get() >= 2, "the pane's status must have been polled AGAIN "
+ "after the WORKING sample, before bootstrapText was sent — this is what proves "
+ "the method actually waited for the WORKING -> IDLE completion boundary instead "
+ "of returning as soon as pickup was seen (or worse, without polling at all, as a "
+ "stub that just returns true would); got " + postClearGetCalls.get()
+ " agent.get call(s) after /clear");
}
@Test
@DisplayName("[fleetd #489] a pane that never reaches a turn boundary after /clear (stuck at "
+ "UNKNOWN, never WORKING either) lets clearSettleSeconds expire — bootstrapText is "
+ "never sent")
void clearPickupNeverSettlesWhenStatusNeverReachesABoundary() throws IOException {
FakeHerdr fake = new FakeHerdr();
HerdrClient stuckUnknownAfterClear = new HerdrClient() {
@Override
public JsonNode call(String method, Object params) throws HerdrException {
if ("agent.prompt".equals(method) && String.valueOf(params).contains("/clear")) {
// "wedged" maps to AgentStatus.UNKNOWN (see AgentStatus#fromWire) — neither a
// pickup signal (WORKING) nor a boundary (IDLE/DONE), and distinct from the
// already-covered BLOCKED case below.
fake.agentStatus("wedged");
}
return fake.call(method, params);
}
@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(stuckUnknownAfterClear, 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 the pane never "
+ "reaches a turn boundary after /clear, whether WORKING was ever observed or not");
assertEquals(0, sendKeysCallCount(fake), "an UNKNOWN status is neither a pickup signal nor "
+ "a boundary — it must never be nudged");
}
@Test
@DisplayName("[fleetd #489] a submit() nudge that throws does not abort the roll — /clear and "
+ "bootstrapText are both still sent")
void submitThatThrowsDoesNotAbortTheRoll() throws IOException {
FakeHerdr fake = new FakeHerdr(); // default idle throughout — nudging will be attempted
HerdrClient throwsOnSubmit = new HerdrClient() {
@Override
public JsonNode call(String method, Object params) throws HerdrException {
if ("agent.send_keys".equals(method)) {
throw new RuntimeException("simulated herdr transport failure on submit");
}
return fake.call(method, params);
}
@Override
public void close() {
fake.close();
}
};
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(throwsOnSubmit, cfg(handover.toString()), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason() + " / " + decision.detail());
var prompts = fake.calls.stream().filter(c -> "agent.prompt".equals(c.method())).toList();
assertEquals(2, prompts.size(), "a throwing submit() must be swallowed, not abort the roll "
+ "— /clear and bootstrapText must both still be sent");
assertTrue(prompts.get(0).params().toString().contains("/clear"));
assertTrue(prompts.get(1).params().toString().contains("read the handover file"));
}
@Test
@DisplayName("a full successful roll sends /clear then bootstrapText, in order, and consumes the token")
void successfulRollSendsClearThenBootstrapTextAndConsumesTheToken() throws IOException {