Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| e966cbadf9 | |||
| c87cc25aa6 | |||
| 3fb331145a |
@@ -80,7 +80,6 @@ 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;
|
||||
@@ -271,12 +270,15 @@ 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.
|
||||
HerdrAwaitOutcome herdrOutcome = awaitHerdr(herdr, System::nanoTime, Fleetd::sleepHerdrPoll);
|
||||
boolean herdrUp = logHerdrWaitOutcomeAndShouldReap(herdrOutcome);
|
||||
boolean herdrUp = awaitHerdr(herdr);
|
||||
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.
|
||||
@@ -1694,70 +1696,12 @@ public final class Fleetd {
|
||||
}
|
||||
|
||||
/**
|
||||
* 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.
|
||||
*/
|
||||
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).
|
||||
* Poll herdr's {@code ping} until it answers or {@link #HERDR_WAIT_SECONDS} elapses (CB-504).
|
||||
*
|
||||
* <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}
|
||||
* @return true if herdr answered, false if it never did
|
||||
*/
|
||||
static HerdrAwaitOutcome awaitHerdr(HerdrClient herdr, LongSupplier nanos, Runnable poller) {
|
||||
long start = nanos.getAsLong();
|
||||
long deadline = start + HERDR_WAIT_SECONDS * 1_000_000_000L;
|
||||
private static boolean awaitHerdr(HerdrClient herdr) {
|
||||
long deadline = System.nanoTime() + HERDR_WAIT_SECONDS * 1_000_000_000L;
|
||||
boolean waited = false;
|
||||
while (true) {
|
||||
try {
|
||||
@@ -1765,57 +1709,25 @@ public final class Fleetd {
|
||||
if (waited) {
|
||||
log.info("herdr is up");
|
||||
}
|
||||
return new HerdrAwaitOutcome(HerdrWaitResult.ANSWERED, nanos.getAsLong() - start);
|
||||
return true;
|
||||
} catch (HerdrException e) {
|
||||
if (nanos.getAsLong() >= deadline) {
|
||||
return new HerdrAwaitOutcome(HerdrWaitResult.DEADLINE_PASSED, nanos.getAsLong() - start);
|
||||
if (System.nanoTime() >= deadline) {
|
||||
return false;
|
||||
}
|
||||
if (!waited) {
|
||||
log.info("waiting up to {}s for the herdr socket…", HERDR_WAIT_SECONDS);
|
||||
waited = true;
|
||||
}
|
||||
poller.run();
|
||||
if (Thread.currentThread().isInterrupted()) {
|
||||
return new HerdrAwaitOutcome(HerdrWaitResult.INTERRUPTED, nanos.getAsLong() - start);
|
||||
try {
|
||||
Thread.sleep(HERDR_WAIT_POLL_MILLIS);
|
||||
} catch (InterruptedException ie) {
|
||||
Thread.currentThread().interrupt();
|
||||
return false;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 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() {
|
||||
}
|
||||
}
|
||||
|
||||
@@ -381,13 +381,18 @@ public final class LeadRollover {
|
||||
*/
|
||||
private void runRollover(PendingRollover p, FleetConfig.LeadRollover cfg) {
|
||||
String lead = p.leadTerminal();
|
||||
boolean turnSettled = waitUntilAtTurnBoundary(lead, cfg.turnSettleSeconds());
|
||||
if (!turnSettled) {
|
||||
log.warn("lead-rollover: pane {} did not reach a turn boundary (IDLE or DONE) 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());
|
||||
long rollStartMillis = nowMillis.getAsLong();
|
||||
TurnSettleResult turnResult = waitUntilAtTurnBoundary(lead, cfg.turnSettleSeconds());
|
||||
if (!turnResult.settled()) {
|
||||
// fleetd #494 follow-up: this line had the SAME defect as the /clear-timeout line below
|
||||
// — cfg.turnSettleSeconds() is the CONFIGURED budget, not how long this wait actually
|
||||
// ran. Print the measured elapsed time alongside it, labelled, exactly like the /clear
|
||||
// path already does.
|
||||
log.warn("lead-rollover: pane {} did not reach a turn boundary (IDLE or DONE) 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={}, configured={}s elapsed={}ms)",
|
||||
lead, p.token(), cfg.turnSettleSeconds(), turnResult.elapsedMillis());
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -395,15 +400,22 @@ 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 = 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={})",
|
||||
lead, cfg.clearSettleSeconds(), p.token());
|
||||
ClearSettleResult clearResult = waitForClearPickupAndSettle(lead, cfg.clearSettleSeconds());
|
||||
if (!clearResult.settled()) {
|
||||
// fleetd #494: cfg.clearSettleSeconds() is the CONFIGURED budget, not how long the wait
|
||||
// actually ran — an operator reading only that number wrongly believes it is a measured
|
||||
// duration. Print the measured elapsed time and nudge count alongside it, each labelled,
|
||||
// so the two can be compared at a glance.
|
||||
log.warn("lead-rollover: pane {} did not reach a turn boundary (IDLE or DONE) after "
|
||||
+ "/clear — NOT sending bootstrapText (token={}, configured={}s "
|
||||
+ "elapsed={}ms nudges={})",
|
||||
lead, p.token(), cfg.clearSettleSeconds(), clearResult.elapsedMillis(),
|
||||
clearResult.nudges());
|
||||
return;
|
||||
}
|
||||
agents.send(lead, cfg.bootstrapTextFor(p.handoverPath()));
|
||||
log.info("lead-rollover: rolled token={} lead={}", p.token(), lead);
|
||||
long rollElapsedMillis = nowMillis.getAsLong() - rollStartMillis;
|
||||
log.info("lead-rollover: rolled token={} lead={} elapsedMs={}", p.token(), lead, rollElapsedMillis);
|
||||
}
|
||||
|
||||
/** Drop a pending request without rolling. @return whether a pending request existed for {@code token} */
|
||||
@@ -475,9 +487,15 @@ public final class LeadRollover {
|
||||
* 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.)
|
||||
*
|
||||
* @return a {@link TurnSettleResult} whose {@code settled()} is {@code true} once a real
|
||||
* boundary was observed, {@code false} if {@code settleSeconds} elapses first.
|
||||
* {@code elapsedMillis()} is a MEASURED value from the injected {@link #nowMillis}
|
||||
* clock, never the configured {@code settleSeconds} budget (fleetd #494 follow-up).
|
||||
*/
|
||||
private boolean waitUntilAtTurnBoundary(String target, int settleSeconds) {
|
||||
long deadline = nowMillis.getAsLong() + TimeUnit.SECONDS.toMillis(settleSeconds);
|
||||
private TurnSettleResult waitUntilAtTurnBoundary(String target, int settleSeconds) {
|
||||
long startMillis = nowMillis.getAsLong();
|
||||
long deadline = startMillis + TimeUnit.SECONDS.toMillis(settleSeconds);
|
||||
while (nowMillis.getAsLong() < deadline) {
|
||||
AgentStatus status;
|
||||
try {
|
||||
@@ -488,13 +506,20 @@ public final class LeadRollover {
|
||||
status = null;
|
||||
}
|
||||
if (status == AgentStatus.IDLE || status == AgentStatus.DONE) {
|
||||
return true;
|
||||
return new TurnSettleResult(true, nowMillis.getAsLong() - startMillis);
|
||||
}
|
||||
settleSleeper.run();
|
||||
}
|
||||
return false;
|
||||
return new TurnSettleResult(false, nowMillis.getAsLong() - startMillis);
|
||||
}
|
||||
|
||||
/**
|
||||
* The measured outcome of {@link #waitUntilAtTurnBoundary} — fleetd #494 follow-up. The sibling
|
||||
* of {@link ClearSettleResult} for the FIRST wait, which never nudges, so it carries no nudge
|
||||
* count.
|
||||
*/
|
||||
private record TurnSettleResult(boolean settled, long elapsedMillis) {}
|
||||
|
||||
/**
|
||||
* 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>
|
||||
@@ -524,8 +549,9 @@ public final class LeadRollover {
|
||||
* 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>
|
||||
* choice {@code Injector} makes — and returns {@code settled() == true} anyway, logged at
|
||||
* {@code warn} with the measured elapsed time (fleetd #494) so an operator can see which
|
||||
* path ran and how long it actually took;</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>
|
||||
@@ -541,15 +567,20 @@ public final class LeadRollover {
|
||||
* 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
|
||||
* @return a {@link ClearSettleResult} whose {@code settled()} is {@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. {@code elapsedMillis()} and {@code nudges()} are MEASURED values (from the
|
||||
* injected {@link #nowMillis} clock and an actual nudge count), never the configured
|
||||
* {@code settleSeconds} budget (fleetd #494).
|
||||
*/
|
||||
private boolean waitForClearPickupAndSettle(String target, int settleSeconds) {
|
||||
long deadline = nowMillis.getAsLong() + TimeUnit.SECONDS.toMillis(settleSeconds);
|
||||
private ClearSettleResult waitForClearPickupAndSettle(String target, int settleSeconds) {
|
||||
long startMillis = nowMillis.getAsLong();
|
||||
long deadline = startMillis + TimeUnit.SECONDS.toMillis(settleSeconds);
|
||||
boolean pickedUp = false; // a WORKING sample has been observed since /clear was sent
|
||||
int idlePollsAwaitingPickup = 0;
|
||||
int nudges = 0;
|
||||
while (nowMillis.getAsLong() < deadline) {
|
||||
AgentStatus status;
|
||||
try {
|
||||
@@ -563,26 +594,56 @@ public final class LeadRollover {
|
||||
pickedUp = true;
|
||||
} else if (status == AgentStatus.IDLE || status == AgentStatus.DONE) {
|
||||
if (pickedUp) {
|
||||
return true; // a real WORKING -> IDLE/DONE completion boundary
|
||||
// a real WORKING -> IDLE/DONE completion boundary
|
||||
return new ClearSettleResult(true, nowMillis.getAsLong() - startMillis, nudges);
|
||||
}
|
||||
if (++idlePollsAwaitingPickup >= PICKUP_GRACE_POLLS) {
|
||||
log.info("lead-rollover: /clear on {} was never observed as WORKING after {} "
|
||||
long elapsedMillis = nowMillis.getAsLong() - startMillis;
|
||||
// fleetd #494: this release trades a possibly-unsubmitted /clear for progress
|
||||
// instead of wedging the roll — that trade is deliberate and stays. But it is
|
||||
// also exactly the case that reported false success in the real incident (the
|
||||
// whole roll "succeeded" after 438ms of a 20s budget), so raise it to WARN and
|
||||
// print the MEASURED elapsed time next to the target pane, not just the count.
|
||||
//
|
||||
// fleetd #494 follow-up (2nd pass): BOTH numbers in this line must come from
|
||||
// the loop's own counters, never from the PICKUP_GRACE_POLLS constant.
|
||||
// `idlePollsAwaitingPickup` and `nudges` each have exactly one write site in
|
||||
// this loop, on the same branch, so on this branch they cannot differ from
|
||||
// PICKUP_GRACE_POLLS / PICKUP_GRACE_POLLS - 1 today — no test can prove the
|
||||
// difference on this line, and printing the counters does not change that.
|
||||
// What it does buy: one source of truth instead of two, so a later change to
|
||||
// the loop (an early return, a second increment site, a different exit
|
||||
// condition) cannot leave this message reporting a number the loop no longer
|
||||
// produces. The place where `nudges` genuinely varies with the run — and is
|
||||
// covered by a test that can tell it apart from a constant — is the
|
||||
// /clear-timeout warn in runRollover, which prints clearResult.nudges().
|
||||
log.warn("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;
|
||||
+ "releasing rather than wedging the roll (elapsed={}ms)",
|
||||
target, idlePollsAwaitingPickup, nudges, elapsedMillis);
|
||||
return new ClearSettleResult(true, elapsedMillis, nudges);
|
||||
}
|
||||
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());
|
||||
} finally {
|
||||
nudges++; // an attempted nudge, whether or not the submit call itself threw
|
||||
}
|
||||
}
|
||||
// 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;
|
||||
return new ClearSettleResult(false, nowMillis.getAsLong() - startMillis, nudges);
|
||||
}
|
||||
|
||||
/**
|
||||
* The measured outcome of {@link #waitForClearPickupAndSettle} — fleetd #494. Carries the
|
||||
* MEASURED elapsed time (from the injected {@link #nowMillis} clock) and nudge count alongside
|
||||
* the settle/timeout decision, so callers can log them instead of the configured budget, which
|
||||
* is not how long the wait actually ran.
|
||||
*/
|
||||
private record ClearSettleResult(boolean settled, long elapsedMillis, int nudges) {}
|
||||
}
|
||||
|
||||
@@ -1,212 +0,0 @@
|
||||
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);
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,9 @@
|
||||
package dev.ltms.fleet.lead;
|
||||
|
||||
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.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
@@ -9,6 +13,7 @@ import dev.ltms.fleet.herdr.HerdrException;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.nio.file.Files;
|
||||
@@ -854,4 +859,207 @@ class LeadRolloverTest {
|
||||
assertFalse(bootstrapSent.contains("\"handover.md\""),
|
||||
"must not name the raw relative configured value in the text actually sent");
|
||||
}
|
||||
|
||||
// ---- fleetd #494: the log lines must print MEASURED values, never the configured budget --
|
||||
|
||||
private static ListAppender<ILoggingEvent> attachLog() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(LeadRollover.class);
|
||||
logger.setLevel(Level.DEBUG);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
return appender;
|
||||
}
|
||||
|
||||
private static void detachLog(ListAppender<ILoggingEvent> appender) {
|
||||
((Logger) LoggerFactory.getLogger(LeadRollover.class)).detachAppender(appender);
|
||||
}
|
||||
|
||||
private static ILoggingEvent lastEventContaining(ListAppender<ILoggingEvent> events, String substring) {
|
||||
return events.list.stream()
|
||||
.filter(e -> e.getFormattedMessage().contains(substring))
|
||||
.reduce((_, b) -> b)
|
||||
.orElseThrow(() -> new AssertionError("no log event contained \"" + substring
|
||||
+ "\"; got: " + events.list.stream().map(ILoggingEvent::getFormattedMessage).toList()));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[fleetd #494] the /clear-timeout warn line prints the MEASURED elapsed time and "
|
||||
+ "nudge count next to the configured budget, never the configured value alone")
|
||||
void clearTimeoutLogPrintsMeasuredElapsedAndNudgesNotJustConfigured() throws IOException {
|
||||
// Idle until /clear is sent, then permanently WORKING (a genuinely stuck /clear that never
|
||||
// reaches a completion boundary) — isolates the SECOND wait exactly like
|
||||
// clearThatNeverSettlesAfterwardsNeverSendsBootstrapText, but with a clock that ADVANCES on
|
||||
// every read so the measured elapsed time is a deterministic, non-zero value distinct from
|
||||
// the configured budget — pinning fleetd #494's fix, not just its absence of a hang.
|
||||
FakeHerdr fake = new FakeHerdr();
|
||||
HerdrClient flipsAfterClear = new HerdrClient() {
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) throws HerdrException {
|
||||
JsonNode result = fake.call(method, params);
|
||||
if ("agent.prompt".equals(method) && String.valueOf(params).contains("/clear")) {
|
||||
fake.agentStatus("working");
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
fake.close();
|
||||
}
|
||||
};
|
||||
Path handover = writeHandover("handover contents");
|
||||
FleetConfig.LeadRollover config =
|
||||
new FleetConfig.LeadRollover(handover.toString(), true, 3600, 20, 1 /*clearSettleSeconds*/, "boot text");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(flipsAfterClear, config, () -> clock.addAndGet(500));
|
||||
|
||||
ListAppender<ILoggingEvent> events = attachLog();
|
||||
try {
|
||||
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");
|
||||
|
||||
ILoggingEvent event = lastEventContaining(events, "NOT sending bootstrapText");
|
||||
assertEquals(Level.WARN, event.getLevel());
|
||||
String message = event.getFormattedMessage();
|
||||
assertTrue(message.contains("configured=1s"), "must label the configured budget: " + message);
|
||||
assertTrue(message.contains("elapsed=1500ms"), "must print the MEASURED elapsed time — "
|
||||
+ "with this fixture's advancing clock, the wait actually ran 1500ms against a "
|
||||
+ "1s(=1000ms) configured budget: " + message);
|
||||
assertTrue(message.contains("nudges=0"), "must print the measured nudge count (0 here — "
|
||||
+ "the pane was WORKING throughout, never IDLE/DONE, so no nudge was ever sent): "
|
||||
+ message);
|
||||
assertFalse(message.contains("within 1s"), "must not present the configured budget as if "
|
||||
+ "it were the measured wait duration: " + message);
|
||||
} finally {
|
||||
detachLog(events);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[fleetd #494] the /clear pickup-grace release line is WARN (was INFO) and prints "
|
||||
+ "the measured elapsed time next to the target pane")
|
||||
void clearGraceReleaseLogIsWarnWithMeasuredElapsed() throws IOException {
|
||||
// Default idle throughout — no WORKING sample is ever observed, so the pickup-grace wait
|
||||
// exhausts PICKUP_GRACE_POLLS and releases rather than wedging (see
|
||||
// clearPickupIsNudgedBeforeBootstrapTextWhenPaneStaysIdle for the un-logged half of this
|
||||
// scenario). A self-advancing clock makes the measured elapsed time deterministic and
|
||||
// provably distinct from a bare poll/nudge count.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Path handover = writeHandover("handover contents");
|
||||
FleetConfig.LeadRollover config = cfg(handover.toString()); // turnSettleSeconds=clearSettleSeconds=20
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, config, () -> clock.addAndGet(500));
|
||||
|
||||
ListAppender<ILoggingEvent> events = attachLog();
|
||||
try {
|
||||
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());
|
||||
|
||||
ILoggingEvent event = lastEventContaining(events, "releasing rather than wedging the roll");
|
||||
assertEquals(Level.WARN, event.getLevel(), "the grace-limit release must be WARN, not "
|
||||
+ "INFO — it is exactly the case that reported false success in the real incident "
|
||||
+ "this fix comes from (a roll that 'succeeded' after 438ms of a 20s budget)");
|
||||
String message = event.getFormattedMessage();
|
||||
assertTrue(message.contains(LEAD), "must name the target pane: " + message);
|
||||
assertTrue(message.contains("elapsed=4500ms"), "must print the MEASURED elapsed time — "
|
||||
+ "with this fixture's advancing clock, the wait ran 4500ms before releasing: "
|
||||
+ message);
|
||||
// fleetd #494 follow-up (2nd pass): both numbers here are DELIBERATE plain literals,
|
||||
// not derived from LeadRollover.PICKUP_GRACE_POLLS. A version of this assertion that
|
||||
// reads "(" + (LeadRollover.PICKUP_GRACE_POLLS - 1) + " of those were nudged)" builds
|
||||
// its expectation the same way the production code used to build the log line, so it
|
||||
// cannot tell a fixed constant apart from the measured counter — proved by reverting
|
||||
// the production fix and re-running: that mutant stayed green under the old assertion.
|
||||
// If PICKUP_GRACE_POLLS ever changes, THIS TEST MUST FAIL and a human must look at the
|
||||
// new message and update the literals below, not just re-derive them.
|
||||
assertTrue(message.contains("after 8 consecutive IDLE/DONE polls (7 of those were nudged)"),
|
||||
"must print the measured poll count and nudge count as plain numbers, not the "
|
||||
+ "PICKUP_GRACE_POLLS constant standing in for either: " + message);
|
||||
} finally {
|
||||
detachLog(events);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[fleetd #494] the success line prints the measured elapsed time for the whole roll")
|
||||
void successLogPrintsMeasuredElapsedForTheWholeRoll() throws IOException {
|
||||
// Same fixture as clearGraceReleaseLogIsWarnWithMeasuredElapsed: default idle throughout, so
|
||||
// the grace release fires and the roll still goes on to send bootstrapText and log success.
|
||||
// This is deliberately the SAME shape as the real incident (a roll that "succeeds" quickly)
|
||||
// — the missing signal was the elapsed time on this exact line.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Path handover = writeHandover("handover contents");
|
||||
FleetConfig.LeadRollover config = cfg(handover.toString());
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, config, () -> clock.addAndGet(500));
|
||||
|
||||
ListAppender<ILoggingEvent> events = attachLog();
|
||||
try {
|
||||
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());
|
||||
|
||||
ILoggingEvent event = lastEventContaining(events, "lead-rollover: rolled");
|
||||
assertEquals(Level.INFO, event.getLevel());
|
||||
String message = event.getFormattedMessage();
|
||||
// fleetd #494 follow-up: waitUntilAtTurnBoundary now also reads the injected clock one
|
||||
// extra time (to compute ITS OWN measured elapsed on the success path), so the shared
|
||||
// fixture clock advances by one more 500ms tick before the roll finishes than it did
|
||||
// before that follow-up — 7000ms, not 6500ms.
|
||||
assertTrue(message.contains("elapsedMs=7000"), "must print the MEASURED elapsed time for "
|
||||
+ "the whole roll — with this fixture's advancing clock, the full roll (turn-settle "
|
||||
+ "wait + /clear wait + bootstrapText) took 7000ms: " + message);
|
||||
} finally {
|
||||
detachLog(events);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[fleetd #494 follow-up] the turn-settle timeout warn line prints the MEASURED "
|
||||
+ "elapsed time next to the configured budget, never the configured value alone")
|
||||
void turnTimeoutLogPrintsMeasuredElapsedNotJustConfigured() throws IOException {
|
||||
// The brief that named the four items this ticket fixed left this exact sibling line out —
|
||||
// "did not reach a turn boundary (IDLE or DONE) within {}s after confirm()" — even though it
|
||||
// has the identical defect shape one method below. Same fixture shape as
|
||||
// clearTimeoutLogPrintsMeasuredElapsedAndNudgesNotJustConfigured, but for the FIRST wait: the
|
||||
// calling lead's own pane never goes idle, so waitUntilAtTurnBoundary times out.
|
||||
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");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, config, () -> clock.addAndGet(500));
|
||||
|
||||
ListAppender<ILoggingEvent> events = attachLog();
|
||||
try {
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
|
||||
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");
|
||||
|
||||
ILoggingEvent event = lastEventContaining(events, "refusing to send /clear at all");
|
||||
assertEquals(Level.WARN, event.getLevel());
|
||||
String message = event.getFormattedMessage();
|
||||
assertTrue(message.contains("configured=1s"), "must label the configured budget: " + message);
|
||||
assertTrue(message.contains("elapsed=1500ms"), "must print the MEASURED elapsed time — "
|
||||
+ "with this fixture's advancing clock, the wait actually ran 1500ms against a "
|
||||
+ "1s(=1000ms) configured budget: " + message);
|
||||
assertFalse(message.contains("within 1s"), "must not present the configured budget as if "
|
||||
+ "it were the measured wait duration: " + message);
|
||||
} finally {
|
||||
detachLog(events);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user