From 735b6af9765a34f1c5c7ce3d564c7b596c48a667 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 12 Sep 2026 15:48:07 +0700 Subject: [PATCH 1/2] fleetd #544: progress watchdog for StatusPoller and SessionReaper loops MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Each loop's virtual-thread runner (StatusPoller, SessionReaper) can die or get permanently parked in a herdr call with no read timeout, and nothing observed it: /healthz stayed green and Thread.isAlive() kept reporting true the whole time. Add LoopWatchdog (dev.ltms.fleet.inject — see its javadoc for why not dev.ltms.fleet.health, which would close a package cycle through session): each loop now records a monotonic last-completed-round timestamp (injectable LongSupplier clock, same pattern as Injector/SessionManager) and exposes it as a three-state health() fact — RUNNING, STALLED (dead or parked, indistinguishable from outside), STOPPED (stop() was called on purpose, never an alarm). This is the fleetd #512 shape: one flag cannot carry both "halted on purpose" and "halted unexpectedly", so stop() marks its own state explicitly instead of leaving state() to infer it from staleness. Staleness thresholds are derived from each loop's own poll interval with a documented multiplier: StatusPoller 40x (250ms -> 10s), SessionReaper 12x (5000ms -> 60s). Scope: observability only, per the ticket's own comment. No restart/recovery mechanism, no /healthz or REST/MCP wiring beyond the public health() API, no change to the per-item catch(Throwable) behavior (#543) or a process-wide uncaught-exception handler (ruled out on #538), and no deadline added to the herdr read itself (a separate, real ticket). 🤖 Generated with Claude Code Co-Authored-By: Claude --- .../dev/ltms/fleet/inject/LoopWatchdog.java | 91 +++++++++++ .../dev/ltms/fleet/inject/StatusPoller.java | 49 ++++++ .../dev/ltms/fleet/session/SessionReaper.java | 39 +++++ .../ltms/fleet/inject/LoopWatchdogTest.java | 61 +++++++ .../inject/StatusPollerWatchdogTest.java | 154 ++++++++++++++++++ .../session/SessionReaperWatchdogTest.java | 153 +++++++++++++++++ 6 files changed, 547 insertions(+) create mode 100644 fleetd/src/main/java/dev/ltms/fleet/inject/LoopWatchdog.java create mode 100644 fleetd/src/test/java/dev/ltms/fleet/inject/LoopWatchdogTest.java create mode 100644 fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerWatchdogTest.java create mode 100644 fleetd/src/test/java/dev/ltms/fleet/session/SessionReaperWatchdogTest.java diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/LoopWatchdog.java b/fleetd/src/main/java/dev/ltms/fleet/inject/LoopWatchdog.java new file mode 100644 index 0000000..4b97b39 --- /dev/null +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/LoopWatchdog.java @@ -0,0 +1,91 @@ +package dev.ltms.fleet.inject; + +import java.util.function.LongSupplier; + +/** + * A progress watchdog for a singleton background loop (fleetd #544): {@link StatusPoller} and + * {@code dev.ltms.fleet.session.SessionReaper} each hold one. It tracks the monotonic timestamp + * of the loop's last completed round and classifies health from that timestamp alone — + * never from thread liveness. That choice is deliberate: measured on {@code main} at + * {@code 7611b69}, {@code UnixSocketHerdrClient.call()} reads with no deadline anywhere in + * {@code herdr/}, so a herdrd that accepts a connection and never answers parks the calling + * loop's thread forever — the thread stays alive and {@code Thread.isAlive()} keeps reporting + * {@code true} the whole time. A last-round timestamp that stops advancing catches that parked + * case exactly the same way it catches a thread that died outright: either way, nothing marks a + * new round complete. + * + *

Lives in {@code inject} (rather than {@code health}, which would be the more obvious home + * for a health-reporting fact) because {@code health} already depends on {@code session} + * ({@code HealthSnapshot}/{@code FleetHealthMonitor} read {@code MemberSession}), and + * {@code session} already depends on {@code inject} ({@code SessionManager} implements + * {@code TurnListener} and extends {@code MemberPresence}). Putting this class in {@code health} + * would have made {@code SessionReaper} (in {@code session}) depend back on {@code health}, + * closing a cycle through {@code health -> session -> health} — {@code PackageCyclesTest} catches + * exactly this. {@code inject} has no dependency on {@code health}, so {@code session} depending + * on {@code inject} for this stays a one-way edge, same direction it already depends in. + * + *

Three states, not two (the fleetd #512 shape — one flag carrying two conditions that need + * opposite handling). {@code stop()} also makes the loop go quiet, exactly like a crash or a + * hang does, so a caller must say which one happened: {@link #markStoppedByCaller()} records + * "this halt was on purpose," and {@link #state()} reports {@link State#STOPPED} for it + * regardless of how stale the last round looks — an intentional stop is never an alarm. + */ +public final class LoopWatchdog { + + /** + * {@code RUNNING} — the loop is going and its last round finished recently. + * {@code STALLED} — the loop was never stopped on purpose, and its last-completed-round + * timestamp is older than the threshold. Covers both a dead loop and a parked one. + * {@code STOPPED} — {@link #markStoppedByCaller()} was called. Not an alarm. + */ + public enum State { RUNNING, STALLED, STOPPED } + + private final LongSupplier nowNanos; + private final long staleAfterNanos; + private volatile long lastRoundNanos; + private volatile boolean stoppedByCaller; + + /** + * @param nowNanos a monotonic elapsed-time clock, e.g. {@code System::nanoTime} live, a + * controllable stub in tests. + * @param staleAfterNanos how long a last-completed-round timestamp may age before {@link #state()} + * reports {@link State#STALLED}. Derive this from the loop's own poll + * interval with generous headroom — see the caller's own derivation comment. + */ + public LoopWatchdog(LongSupplier nowNanos, long staleAfterNanos) { + this.nowNanos = nowNanos; + this.staleAfterNanos = staleAfterNanos; + this.lastRoundNanos = nowNanos.getAsLong(); + } + + /** Call once per completed round, from inside the loop. */ + public void recordRoundComplete() { + lastRoundNanos = nowNanos.getAsLong(); + } + + /** + * Call from {@code start()} (and any future restart): clears a previous stop's mark and + * resets the clock, so the loop gets a fresh grace period before its first round completes + * rather than being judged against a timestamp left over from before it was (re)started. + */ + public void reset() { + stoppedByCaller = false; + lastRoundNanos = nowNanos.getAsLong(); + } + + /** + * Call from {@code stop()}: marks the halt as intentional, so {@link #state()} reports + * {@link State#STOPPED} rather than {@link State#STALLED} no matter how stale the last round is. + */ + public void markStoppedByCaller() { + stoppedByCaller = true; + } + + /** The current fact about this loop, per the three-state contract above. */ + public State state() { + if (stoppedByCaller) { + return State.STOPPED; + } + return (nowNanos.getAsLong() - lastRoundNanos >= staleAfterNanos) ? State.STALLED : State.RUNNING; + } +} diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/StatusPoller.java b/fleetd/src/main/java/dev/ltms/fleet/inject/StatusPoller.java index 3c47cf9..0aa64c8 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/inject/StatusPoller.java +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/StatusPoller.java @@ -8,6 +8,8 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.util.Set; +import java.util.concurrent.TimeUnit; +import java.util.function.LongSupplier; /** * Drives the {@link Injector} by sampling each active worker's {@code agent_status} and @@ -22,11 +24,25 @@ public final class StatusPoller { private static final Logger log = LoggerFactory.getLogger(StatusPoller.class); + /** + * How many multiples of {@code intervalMillis} the last-completed round may age before + * {@link #health()} reports {@link LoopWatchdog.State#STALLED} (fleetd #544). A round this + * loop runs is normally a small fraction of one interval — it only takes longer when a herdr + * call is genuinely wedged (herdr's {@code UnixSocketHerdrClient} has no read deadline, so a + * herdrd that accepts and never answers parks the thread forever) — so a wide multiplier + * avoids false alarms from an ordinarily slow round. At the production interval + * ({@link Injector#POLL_INTERVAL_MILLIS} = 250ms) this puts the threshold at 10s: long enough + * that a transient slow poll never trips it, short enough that a stuck poller is visible well + * before anything else downstream (delivery, `/healthz` consumers) would show a symptom. + */ + private static final long STALE_AFTER_INTERVAL_MULTIPLIER = 40; + private final AgentControl agents; private final HerdrRouter router; private final Injector injector; private final StatusRefiner refiner; private final long intervalMillis; + private final LoopWatchdog watchdog; private volatile boolean running; private Thread thread; @@ -36,14 +52,26 @@ public final class StatusPoller { public StatusPoller(AgentControl agents, Injector injector, StatusRefiner refiner, long intervalMillis) { + this(agents, injector, refiner, intervalMillis, System::nanoTime); + } + + /** Full constructor — for tests: an injectable monotonic clock (fleetd #544). */ + StatusPoller(AgentControl agents, Injector injector, StatusRefiner refiner, + long intervalMillis, LongSupplier nowNanos) { this.agents = agents; this.router = null; this.injector = injector; this.refiner = refiner; this.intervalMillis = intervalMillis; + this.watchdog = new LoopWatchdog(nowNanos, staleAfterNanos(intervalMillis)); } public StatusPoller(HerdrRouter router, Injector injector, long intervalMillis) { + this(router, injector, intervalMillis, System::nanoTime); + } + + /** Full constructor — for tests: an injectable monotonic clock (fleetd #544). */ + StatusPoller(HerdrRouter router, Injector injector, long intervalMillis, LongSupplier nowNanos) { this.agents = null; this.router = router; this.injector = injector; @@ -53,12 +81,28 @@ public final class StatusPoller { // against the LEAD daemon even though this field points at the member one. this.refiner = new StatusRefiner(router.memberAgents()); this.intervalMillis = intervalMillis; + this.watchdog = new LoopWatchdog(nowNanos, staleAfterNanos(intervalMillis)); + } + + private static long staleAfterNanos(long intervalMillis) { + return TimeUnit.MILLISECONDS.toNanos(Math.max(intervalMillis, 1)) * STALE_AFTER_INTERVAL_MULTIPLIER; + } + + /** + * This loop's progress fact (fleetd #544): {@code RUNNING}, {@code STALLED} (dead or parked — + * indistinguishable from outside, and never observed on purpose), or {@code STOPPED} + * ({@link #stop()} was called). See {@link LoopWatchdog} for why staleness — not thread + * liveness — is the signal. + */ + public LoopWatchdog.State health() { + return watchdog.state(); } /** Start the polling loop on a virtual thread. Idempotent. */ public synchronized void start() { if (running) return; running = true; + watchdog.reset(); // fleetd #544: a fresh grace period, not last run's stale mark/timestamp. thread = Thread.ofVirtual().name("status-poller").start(this::loop); log.info("status poller started (interval {}ms)", intervalMillis); } @@ -90,6 +134,10 @@ public final class StatusPoller { log.error("unexpected failure polling {}; skipping this round", target, e); } } + // fleetd #544: a round is "complete" only once every active target has been polled — + // a target parked mid-for-loop (control.status(target) never returning) means this + // line is never reached, so the watchdog goes stale exactly like a dead loop would. + watchdog.recordRoundComplete(); sleep(); } } finally { @@ -112,6 +160,7 @@ public final class StatusPoller { /** Stop the polling loop. Idempotent. */ public synchronized void stop() { running = false; + watchdog.markStoppedByCaller(); // fleetd #544: this halt is on purpose — health() must say STOPPED, not STALLED. if (thread != null) thread.interrupt(); } } diff --git a/fleetd/src/main/java/dev/ltms/fleet/session/SessionReaper.java b/fleetd/src/main/java/dev/ltms/fleet/session/SessionReaper.java index 4a85342..4ecbd8c 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/session/SessionReaper.java +++ b/fleetd/src/main/java/dev/ltms/fleet/session/SessionReaper.java @@ -1,9 +1,11 @@ package dev.ltms.fleet.session; +import dev.ltms.fleet.inject.LoopWatchdog; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.util.concurrent.TimeUnit; +import java.util.function.LongSupplier; /** * Periodic virtual-thread reaper that tears down {@code READY}/{@code DONE} sessions which have @@ -14,6 +16,16 @@ public final class SessionReaper { private static final Logger log = LoggerFactory.getLogger(SessionReaper.class); private static final long DEFAULT_INTERVAL_MILLIS = 5000; + /** + * How many multiples of {@code intervalMillis} the last-completed round may age before + * {@link #health()} reports {@link LoopWatchdog.State#STALLED} (fleetd #544). One round here + * is {@code sessions.reapIdle} plus the (rare, best-effort) WIP sweep — both normally finish + * in a small fraction of one interval. At the default 5s interval this puts the threshold at + * 60s: generous enough that an occasional slow git call in the WIP sweep never trips it, short + * enough that a genuinely wedged reap round (mirroring the herdr-read-with-no-deadline hang + * {@code StatusPoller} guards against — see fleetd #544) is caught inside about a minute. + */ + private static final long STALE_AFTER_INTERVAL_MULTIPLIER = 12; /** CB-586: the refs/wip age floor — never sweep a snapshot younger than 24h (the CB-586 rule). */ private static final long WIP_MIN_AGE_MILLIS = TimeUnit.HOURS.toMillis(24); /** @@ -25,6 +37,7 @@ public final class SessionReaper { private final SessionManager sessions; private final long idleTtlNanos; private final long intervalMillis; + private final LoopWatchdog watchdog; private volatile boolean running; private Thread thread; /** @@ -45,15 +58,36 @@ public final class SessionReaper { /** Construct a reaper with an explicit polling interval (useful for tests). */ public SessionReaper(SessionManager sessions, long idleTtlSeconds, long intervalMillis) { + this(sessions, idleTtlSeconds, intervalMillis, System::nanoTime); + } + + /** Full constructor — for tests: an injectable monotonic clock (fleetd #544). */ + SessionReaper(SessionManager sessions, long idleTtlSeconds, long intervalMillis, LongSupplier nowNanos) { this.sessions = sessions; this.idleTtlNanos = TimeUnit.SECONDS.toNanos(idleTtlSeconds); this.intervalMillis = intervalMillis; + this.watchdog = new LoopWatchdog(nowNanos, staleAfterNanos(intervalMillis)); + } + + private static long staleAfterNanos(long intervalMillis) { + return TimeUnit.MILLISECONDS.toNanos(Math.max(intervalMillis, 1)) * STALE_AFTER_INTERVAL_MULTIPLIER; + } + + /** + * This loop's progress fact (fleetd #544): {@code RUNNING}, {@code STALLED} (dead or parked — + * indistinguishable from outside, and never observed on purpose), or {@code STOPPED} + * ({@link #stop()} was called). See {@link LoopWatchdog} for why staleness — not thread + * liveness — is the signal. + */ + public LoopWatchdog.State health() { + return watchdog.state(); } /** Start the reaper loop on a virtual thread. Idempotent. */ public synchronized void start() { if (running) return; running = true; + watchdog.reset(); // fleetd #544: a fresh grace period, not last run's stale mark/timestamp. thread = Thread.ofVirtual().name("session-reaper").start(this::loop); log.info("session reaper started (idle ttl {}s, interval {}ms)", TimeUnit.NANOSECONDS.toSeconds(idleTtlNanos), intervalMillis); @@ -68,6 +102,10 @@ public final class SessionReaper { log.error("session reaper iteration failed; continuing", e); } maybeSweepWipRefs(); + // fleetd #544: a round is "complete" only after reapIdle AND the WIP-sweep gate have + // both returned — a hang in either (e.g. a wedged git call) means this line is never + // reached, so the watchdog goes stale exactly like a dead loop would. + watchdog.recordRoundComplete(); sleep(); } } finally { @@ -118,6 +156,7 @@ public final class SessionReaper { /** Stop the reaper loop. Idempotent. */ public synchronized void stop() { running = false; + watchdog.markStoppedByCaller(); // fleetd #544: this halt is on purpose — health() must say STOPPED, not STALLED. if (thread != null) thread.interrupt(); } } diff --git a/fleetd/src/test/java/dev/ltms/fleet/inject/LoopWatchdogTest.java b/fleetd/src/test/java/dev/ltms/fleet/inject/LoopWatchdogTest.java new file mode 100644 index 0000000..71042e4 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/inject/LoopWatchdogTest.java @@ -0,0 +1,61 @@ +package dev.ltms.fleet.inject; + +import org.junit.jupiter.api.Test; + +import java.util.concurrent.atomic.AtomicLong; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +/** fleetd #544: the three-state contract in isolation from either loop that composes it. */ +class LoopWatchdogTest { + + private static final long STALE_AFTER_NANOS = 1000; + + @Test + void reportsRunningBeforeTheThresholdElapses() { + AtomicLong clock = new AtomicLong(0); + LoopWatchdog watchdog = new LoopWatchdog(clock::get, STALE_AFTER_NANOS); + clock.set(STALE_AFTER_NANOS - 1); + assertEquals(LoopWatchdog.State.RUNNING, watchdog.state()); + } + + @Test + void reportsStalledOnceTheLastRoundAgesPastTheThreshold() { + AtomicLong clock = new AtomicLong(0); + LoopWatchdog watchdog = new LoopWatchdog(clock::get, STALE_AFTER_NANOS); + clock.set(STALE_AFTER_NANOS); + assertEquals(LoopWatchdog.State.STALLED, watchdog.state()); + } + + @Test + void recordRoundCompleteResetsTheStaleClock() { + AtomicLong clock = new AtomicLong(0); + LoopWatchdog watchdog = new LoopWatchdog(clock::get, STALE_AFTER_NANOS); + clock.set(STALE_AFTER_NANOS - 1); + watchdog.recordRoundComplete(); // last round is now "now" again + clock.set(2 * STALE_AFTER_NANOS - 2); + assertEquals(LoopWatchdog.State.RUNNING, watchdog.state(), + "a fresh round completion must push the staleness deadline forward"); + } + + @Test + void markStoppedByCallerReportsStoppedRegardlessOfStaleness() { + AtomicLong clock = new AtomicLong(0); + LoopWatchdog watchdog = new LoopWatchdog(clock::get, STALE_AFTER_NANOS); + watchdog.markStoppedByCaller(); + clock.set(STALE_AFTER_NANOS * 1000); + assertEquals(LoopWatchdog.State.STOPPED, watchdog.state(), + "an intentional stop must never be reported as STALLED"); + } + + @Test + void resetClearsAPreviousStopAndTheStaleClock() { + AtomicLong clock = new AtomicLong(0); + LoopWatchdog watchdog = new LoopWatchdog(clock::get, STALE_AFTER_NANOS); + watchdog.markStoppedByCaller(); + clock.set(STALE_AFTER_NANOS * 1000); + watchdog.reset(); + assertEquals(LoopWatchdog.State.RUNNING, watchdog.state(), + "reset() must clear both the stop mark and the stale timestamp"); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerWatchdogTest.java b/fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerWatchdogTest.java new file mode 100644 index 0000000..e575305 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerWatchdogTest.java @@ -0,0 +1,154 @@ +package dev.ltms.fleet.inject; + +import com.fasterxml.jackson.databind.JsonNode; +import dev.ltms.fleet.herdr.AgentControl; +import dev.ltms.fleet.herdr.HerdrClient; +import dev.ltms.fleet.msg.TestTurnTokens; +import org.junit.jupiter.api.Test; + +import java.lang.reflect.Field; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * fleetd #544: {@link StatusPoller#health()} — the progress watchdog, not a liveness check. Every + * test drives {@link LoopWatchdog}'s clock explicitly (never real elapsed time), so the threshold + * crossing is deterministic rather than a race against the poll interval. + */ +class StatusPollerWatchdogTest { + + private static final long INTERVAL_MILLIS = 50; + // Mirrors StatusPoller.staleAfterNanos(50): 50ms * the 40x multiplier = 2000ms. + private static final long STALE_AFTER_NANOS = TimeUnit.MILLISECONDS.toNanos(INTERVAL_MILLIS) * 40; + + @Test + void aLoopParkedInAHerdrCallThatNeverReturnsReportsStalled() throws Exception { + CountDownLatch entered = new CountDownLatch(1); + AgentControl agents = new AgentControl(new BlocksOnAgentGet(entered)); + Injector injector = new Injector(agents); + AtomicLong clock = new AtomicLong(0); + StatusPoller poller = new StatusPoller(agents, injector, new StatusRefiner(agents), + INTERVAL_MILLIS, clock::get); + poller.start(); + try { + injector.enqueue("w1:p1", "hello", TestTurnTokens.inert("w1:p1")); + assertTrue(entered.await(2, TimeUnit.SECONDS), + "the poller must have entered the blocking herdr call"); + + assertEquals(LoopWatchdog.State.RUNNING, poller.health(), + "freshly parked, before the threshold elapses, must still read RUNNING"); + + clock.set(STALE_AFTER_NANOS + 1); + assertEquals(LoopWatchdog.State.STALLED, poller.health(), + "a loop parked mid-round past the staleness threshold must report STALLED"); + } finally { + poller.stop(); + } + } + + @Test + void aLoopThatDiedReportsStalled() throws Exception { + AtomicLong clock = new AtomicLong(0); + // intervalMillis=-1: Thread.sleep(-1) throws IllegalArgumentException, uncaught by loop(), + // which kills the carrier thread — the same "died" shape StatusPollerResilienceTest uses. + StatusPoller poller = new StatusPoller(new AgentControl(new IdleHerdr()), new Injector(new AgentControl(new IdleHerdr())), + new StatusRefiner(new AgentControl(new IdleHerdr())), -1, clock::get); + poller.start(); + Thread thread = threadOf(poller); + thread.join(2000); + assertFalse(thread.isAlive(), "the loop must have died from Thread.sleep(-1)"); + + // staleAfterNanos(-1) = TimeUnit.MILLISECONDS.toNanos(max(-1,1)) * 40 = 40ms. + clock.set(TimeUnit.MILLISECONDS.toNanos(1) * 40 + 1); + assertEquals(LoopWatchdog.State.STALLED, poller.health(), + "a dead loop's last-completed round goes stale and must report STALLED"); + } + + @Test + void aLoopStoppedOnPurposeReportsStoppedNotStalled() throws Exception { + AtomicLong clock = new AtomicLong(0); + StatusPoller poller = new StatusPoller(new AgentControl(new IdleHerdr()), new Injector(new AgentControl(new IdleHerdr())), + new StatusRefiner(new AgentControl(new IdleHerdr())), INTERVAL_MILLIS, clock::get); + poller.start(); + poller.stop(); + + // Advance the clock far past any staleness threshold — an intentional stop must never + // be reported as STALLED no matter how old the last round looks (fleetd #512 shape). + clock.set(STALE_AFTER_NANOS * 100); + assertEquals(LoopWatchdog.State.STOPPED, poller.health(), + "stop() must report STOPPED, never STALLED, however stale the last round looks"); + } + + @Test + void aHealthyLoopReportsRunning() throws Exception { + // Real elapsed time on purpose, unlike the other tests here: a clock frozen at 0 would read + // RUNNING even if recordRoundComplete() were never wired into loop() at all, since the + // constructor's own initial timestamp already satisfies "not stale yet". Running for several + // multiples of this instance's own threshold (5ms interval * 40 = 200ms) and still reading + // RUNNING instead proves the loop's repeated round completions are the thing keeping the + // staleness deadline pushed forward. + long fastIntervalMillis = 5; + StatusPoller poller = new StatusPoller(new AgentControl(new IdleHerdr()), new Injector(new AgentControl(new IdleHerdr())), + new StatusRefiner(new AgentControl(new IdleHerdr())), fastIntervalMillis, System::nanoTime); + poller.start(); + try { + Thread.sleep(800); + assertEquals(LoopWatchdog.State.RUNNING, poller.health(), + "a healthy loop's repeated round completions must keep pushing the staleness deadline forward"); + } finally { + poller.stop(); + } + } + + private static Thread threadOf(StatusPoller poller) throws ReflectiveOperationException { + Field field = StatusPoller.class.getDeclaredField("thread"); + field.setAccessible(true); + return (Thread) field.get(poller); + } + + /** Never has any active target, so the loop just idles and completes empty rounds. */ + private static final class IdleHerdr implements HerdrClient { + @Override + public JsonNode call(String method, Object params) { + throw new AssertionError("unexpected herdr call: " + method); + } + + @Override + public void close() { + } + } + + private static final class BlocksOnAgentGet implements HerdrClient { + private final CountDownLatch entered; + + private BlocksOnAgentGet(CountDownLatch entered) { + this.entered = entered; + } + + @Override + public JsonNode call(String method, Object params) { + if ("agent.get".equals(method)) { + entered.countDown(); + try { + // Blocks forever, mirroring UnixSocketHerdrClient.readLine()'s deadline-free + // read (fleetd #544) — a test double standing in for a parked socket, not a + // real one. + new CountDownLatch(1).await(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + throw new AssertionError("unreachable: only stop()'s interrupt reaches here"); + } + throw new AssertionError("unexpected herdr call: " + method); + } + + @Override + public void close() { + } + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/session/SessionReaperWatchdogTest.java b/fleetd/src/test/java/dev/ltms/fleet/session/SessionReaperWatchdogTest.java new file mode 100644 index 0000000..9993db1 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/session/SessionReaperWatchdogTest.java @@ -0,0 +1,153 @@ +package dev.ltms.fleet.session; + +import dev.ltms.fleet.config.FleetConfig; +import dev.ltms.fleet.guard.SubscriptionGuard; +import dev.ltms.fleet.herdr.AgentControl; +import dev.ltms.fleet.inject.LoopWatchdog; +import dev.ltms.fleet.herdr.FakeHerdr; +import dev.ltms.fleet.herdr.WorkspaceControl; +import dev.ltms.fleet.member.ClaudeCodeLauncher; +import org.junit.jupiter.api.Test; + +import java.lang.reflect.Field; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; +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; + +/** + * fleetd #544: {@link SessionReaper#health()} — the progress watchdog, not a liveness check. Every + * test drives {@link LoopWatchdog}'s clock explicitly (never real elapsed time), so the threshold + * crossing is deterministic rather than a race against the poll interval. + */ +class SessionReaperWatchdogTest { + + private static final long INTERVAL_MILLIS = 50; + // Mirrors SessionReaper.staleAfterNanos(50): 50ms * the 12x multiplier = 600ms. + private static final long STALE_AFTER_NANOS = TimeUnit.MILLISECONDS.toNanos(INTERVAL_MILLIS) * 12; + + @Test + void aLoopParkedInAReapCallThatNeverReturnsReportsStalled() throws Exception { + CountDownLatch entered = new CountDownLatch(1); + // reapIdle's very first statement reads sessions' own injected clock — blocking there + // parks the reaper thread mid-round, exactly as a wedged call would (fleetd #544). + SessionManager sessions = sessionManager(new BlocksOnFirstCall(entered)); + AtomicLong watchdogClock = new AtomicLong(0); + SessionReaper reaper = new SessionReaper(sessions, 60, INTERVAL_MILLIS, watchdogClock::get); + reaper.start(); + try { + assertTrue(entered.await(2, TimeUnit.SECONDS), + "the reaper must have entered the blocking reapIdle call"); + + assertEquals(LoopWatchdog.State.RUNNING, reaper.health(), + "freshly parked, before the threshold elapses, must still read RUNNING"); + + watchdogClock.set(STALE_AFTER_NANOS + 1); + assertEquals(LoopWatchdog.State.STALLED, reaper.health(), + "a loop parked mid-round past the staleness threshold must report STALLED"); + } finally { + reaper.stop(); + } + } + + @Test + void aLoopThatDiedReportsStalled() throws Exception { + AtomicLong watchdogClock = new AtomicLong(0); + // intervalMillis=-1: Thread.sleep(-1) throws IllegalArgumentException, uncaught by loop(), + // which kills the carrier thread — the same "died" shape SessionReaperResilienceTest uses. + SessionReaper reaper = new SessionReaper(sessionManager(System::nanoTime), 60, -1, watchdogClock::get); + reaper.start(); + Thread thread = threadOf(reaper); + thread.join(2000); + assertFalse(thread.isAlive(), "the loop must have died from Thread.sleep(-1)"); + + // staleAfterNanos(-1) = TimeUnit.MILLISECONDS.toNanos(max(-1,1)) * 12 = 12ms. + watchdogClock.set(TimeUnit.MILLISECONDS.toNanos(1) * 12 + 1); + assertEquals(LoopWatchdog.State.STALLED, reaper.health(), + "a dead loop's last-completed round goes stale and must report STALLED"); + } + + @Test + void aLoopStoppedOnPurposeReportsStoppedNotStalled() throws Exception { + AtomicLong watchdogClock = new AtomicLong(0); + SessionReaper reaper = new SessionReaper(sessionManager(System::nanoTime), 60, INTERVAL_MILLIS, watchdogClock::get); + reaper.start(); + reaper.stop(); + + // Advance the clock far past any staleness threshold — an intentional stop must never + // be reported as STALLED no matter how old the last round looks (fleetd #512 shape). + watchdogClock.set(STALE_AFTER_NANOS * 100); + assertEquals(LoopWatchdog.State.STOPPED, reaper.health(), + "stop() must report STOPPED, never STALLED, however stale the last round looks"); + } + + @Test + void aHealthyLoopReportsRunning() throws Exception { + // Real elapsed time on purpose, unlike the other tests here: a clock frozen at 0 would read + // RUNNING even if recordRoundComplete() were never wired into loop() at all, since the + // constructor's own initial timestamp already satisfies "not stale yet". Running for several + // multiples of this instance's own threshold (5ms interval * 12 = 60ms) and still reading + // RUNNING instead proves the loop's repeated round completions are the thing keeping the + // staleness deadline pushed forward. + long fastIntervalMillis = 5; + SessionReaper reaper = new SessionReaper(sessionManager(System::nanoTime), 60, fastIntervalMillis, System::nanoTime); + reaper.start(); + try { + Thread.sleep(400); + assertEquals(LoopWatchdog.State.RUNNING, reaper.health(), + "a healthy loop's repeated round completions must keep pushing the staleness deadline forward"); + } finally { + reaper.stop(); + } + } + + private static SessionManager sessionManager(LongSupplier clock) { + FakeHerdr herdr = new FakeHerdr(); + FleetConfig.Profile cfg = new FleetConfig.Profile( + "ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", + List.of("ccs", "ltms-local"), "tab", "fleetd-workers", + "worker: {profile} #{n}", null, null, null); + ClaudeCodeLauncher launcher = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr), + new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null); + return new SessionManager(launcher, new FakeWorktrees(), clock); + } + + private static Thread threadOf(SessionReaper reaper) throws ReflectiveOperationException { + Field field = SessionReaper.class.getDeclaredField("thread"); + field.setAccessible(true); + return (Thread) field.get(reaper); + } + + /** Blocks forever on its first call, then falls back to real time if ever unblocked. */ + private static final class BlocksOnFirstCall implements LongSupplier { + private final CountDownLatch entered; + private volatile boolean first = true; + + private BlocksOnFirstCall(CountDownLatch entered) { + this.entered = entered; + } + + @Override + public synchronized long getAsLong() { + if (first) { + first = false; + entered.countDown(); + try { + // Blocks forever, mirroring a genuinely wedged call (fleetd #544) — a test + // double standing in for a stuck dependency, not a real socket. + new CountDownLatch(1).await(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + return System.nanoTime(); + } + } +} -- 2.52.0 From bfac14108f19fb668762051f5be7732fc5cfa694 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 12 Sep 2026 16:00:58 +0700 Subject: [PATCH 2/2] fleetd #544: pin the sticky-STOPPED-across-restart invariant MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review of PR #559 (issue comment #16944) found a surviving mutant: removing watchdog.reset() from StatusPoller.start() (and the identical line in SessionReaper.start()) passed the entire suite. stoppedByCaller is sticky and reset() — called only from start() — is the only thing that clears it. Both loops document start() as idempotent and loop()'s own error log says "it can be restarted", so stop() followed by start() is an anticipated path. Without reset() wired into start(), health() would report STOPPED forever after a restart even though the loop is genuinely running again. Add aRestartedLoopReportsRunningAgainNotStoppedForever to both StatusPollerWatchdogTest and SessionReaperWatchdogTest, pinning "an intentional stop must not outlive the restart that follows it". Verified via the standard mutation cycle: exact-line anchor (not regex, to avoid the \Q-style false match the reviewer flagged) counted pristine 1 -> mutated 0, test goes red with its own message, restored, shasum -a 256 byte-identical, green again. mvn clean install: exit 0, BUILD SUCCESS, Tests run: 1734, Failures: 0, Errors: 0, Skipped: 0 (cross-checked against 130 surefire report files). No production code changed — the reset() call under test was already correct; it simply had nothing pinning it. 🤖 Generated with Claude Code Co-Authored-By: Claude --- .../fleet/inject/StatusPollerWatchdogTest.java | 18 ++++++++++++++++++ .../session/SessionReaperWatchdogTest.java | 17 +++++++++++++++++ 2 files changed, 35 insertions(+) diff --git a/fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerWatchdogTest.java b/fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerWatchdogTest.java index e575305..fdfb7bb 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerWatchdogTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerWatchdogTest.java @@ -84,6 +84,24 @@ class StatusPollerWatchdogTest { "stop() must report STOPPED, never STALLED, however stale the last round looks"); } + @Test + void aRestartedLoopReportsRunningAgainNotStoppedForever() throws Exception { + // fleetd #544 review (issue comment #16944): stoppedByCaller is sticky, and reset() — + // called only from start() — is the sole thing that clears it. start() is documented + // idempotent and loop()'s own error line says "it can be restarted", so stop() followed + // by start() is an anticipated path. Without the reset() call in start(), health() would + // report STOPPED forever after a restart even though the loop is genuinely running again. + AtomicLong clock = new AtomicLong(0); + StatusPoller poller = new StatusPoller(new AgentControl(new IdleHerdr()), new Injector(new AgentControl(new IdleHerdr())), + new StatusRefiner(new AgentControl(new IdleHerdr())), INTERVAL_MILLIS, clock::get); + poller.start(); + poller.stop(); + poller.start(); + + assertEquals(LoopWatchdog.State.RUNNING, poller.health(), + "an intentional stop must not outlive the restart that follows it"); + } + @Test void aHealthyLoopReportsRunning() throws Exception { // Real elapsed time on purpose, unlike the other tests here: a clock frozen at 0 would read diff --git a/fleetd/src/test/java/dev/ltms/fleet/session/SessionReaperWatchdogTest.java b/fleetd/src/test/java/dev/ltms/fleet/session/SessionReaperWatchdogTest.java index 9993db1..afaa622 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/session/SessionReaperWatchdogTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/session/SessionReaperWatchdogTest.java @@ -88,6 +88,23 @@ class SessionReaperWatchdogTest { "stop() must report STOPPED, never STALLED, however stale the last round looks"); } + @Test + void aRestartedLoopReportsRunningAgainNotStoppedForever() throws Exception { + // fleetd #544 review (issue comment #16944): stoppedByCaller is sticky, and reset() — + // called only from start() — is the sole thing that clears it. start() is documented + // idempotent and loop()'s own error line says "it can be restarted", so stop() followed + // by start() is an anticipated path. Without the reset() call in start(), health() would + // report STOPPED forever after a restart even though the loop is genuinely running again. + AtomicLong watchdogClock = new AtomicLong(0); + SessionReaper reaper = new SessionReaper(sessionManager(System::nanoTime), 60, INTERVAL_MILLIS, watchdogClock::get); + reaper.start(); + reaper.stop(); + reaper.start(); + + assertEquals(LoopWatchdog.State.RUNNING, reaper.health(), + "an intentional stop must not outlive the restart that follows it"); + } + @Test void aHealthyLoopReportsRunning() throws Exception { // Real elapsed time on purpose, unlike the other tests here: a clock frozen at 0 would read -- 2.52.0