From 274afafde6dca5817d70307560280d18b7273064 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 12 Sep 2026 09:54:59 +0700 Subject: [PATCH] fleetd #498: awaitHerdr distinguishes deadline-passed from interrupted, with measured elapsed time - 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. --- .../src/main/java/dev/ltms/fleet/Fleetd.java | 124 ++++++++-- .../dev/ltms/fleet/FleetdAwaitHerdrTest.java | 212 ++++++++++++++++++ 2 files changed, 318 insertions(+), 18 deletions(-) create mode 100644 fleetd/src/test/java/dev/ltms/fleet/FleetdAwaitHerdrTest.java diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index 429e1c7..dbfaf10 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -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). + * + *

{@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). + * + *

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() { } } diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdAwaitHerdrTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdAwaitHerdrTest.java new file mode 100644 index 0000000..9249c16 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdAwaitHerdrTest.java @@ -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: + *

+ * 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 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 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 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 attach() { + Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class); + logger.setLevel(Level.DEBUG); + ListAppender appender = new ListAppender<>(); + appender.start(); + logger.addAppender(appender); + return appender; + } + + private static void detach(ListAppender appender) { + ((Logger) LoggerFactory.getLogger(Fleetd.class)).detachAppender(appender); + } +} -- 2.52.0