Merge #502: awaitHerdr reports three outcomes with measured elapsed time, not one boolean (fleetd #498)
CI / contract (push) Successful in 58s
CI / build (push) Successful in 1m49s

Verified by the lead, not taken from the worker's report.

Trial-merged onto main (8f59019) locally and built the merge, because a clean auto-merge is not a
compiling merge:
  mvn clean install -> Tests run: 1687, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS
  FleetdAwaitHerdrTest -> Tests run: 6, Failures: 0
  Gitea CI on 274afaf: success

The worker mutated the three return values. I mutated the half it did not: the reap DECISION at the
call site. Changing logHerdrWaitOutcomeAndShouldReap to return true for INTERRUPTED gave 1 failure,
FleetdAwaitHerdrTest.interruptedLogsItsOwnMessageAndNeverClaimsTheBudgetElapsed:135 ("an interrupted
wait must not tell main to reap"). Two greps proved the mutant applied; restored, shasum
byte-identical, control 6/0 green.

Behaviour preserved: herdrUp is still true only for ANSWERED, and both its readers (the orphan reap
at :276 and the lead auto-launch gate at :375) see exactly what they saw before.

One correction to the worker's report, which changes nothing in the code. It wrote that a real
interrupt-detection regression "would also hang the daemon's startup thread forever". It would not:
in production `nanos` is System::nanoTime, so the deadline check still fires. The hang it hit was a
test-fixture property — a frozen injected clock with no iteration bound. That is fleetd #486's
shape, now seen in a second class, and it is recorded there.
This commit was merged in pull request #502.
This commit is contained in:
2026-09-12 05:01:02 +02:00
2 changed files with 318 additions and 18 deletions
+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() {
}
}
@@ -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);
}
}