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..fdfb7bb --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerWatchdogTest.java @@ -0,0 +1,172 @@ +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 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 + // 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..afaa622 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/session/SessionReaperWatchdogTest.java @@ -0,0 +1,170 @@ +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 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 + // 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(); + } + } +}