Merge #559: progress watchdog for StatusPoller and SessionReaper loops (fleetd #544)
CI / shell-tests (push) Successful in 7s
CI / contract (push) Successful in 51s
CI / build (push) Successful in 2m19s

Verified by the lead on a merged tree (main 7a3b2bb + bfac141 = 5de807f):
mvn clean install exit 0, Tests run: 1744, Failures: 0, 130 surefire reports.

Three mutations run by the lead, all killed:
- StatusPoller.java:105 `watchdog.reset()` deleted (the surviving mutant from the
  first review) -> aRestartedLoopReportsRunningAgainNotStoppedForever:
  expected: <RUNNING> but was: <STOPPED>.
- LoopWatchdog.java:73 `lastRoundNanos = nowNanos.getAsLong()` deleted from reset()
  (not run by the worker) -> LoopWatchdogTest.resetClearsAPreviousStopAndTheStaleClock:
  expected: <RUNNING> but was: <STALLED>.
- LoopWatchdog.java:89 `>=` -> `>` boundary (not run by the worker) ->
  LoopWatchdogTest.reportsStalledOnceTheLastRoundAgesPastTheThreshold:
  expected: <STALLED> but was: <RUNNING>.

Each mutation counted the pristine full line 1 -> 0 by exact string equality
(awk '$0==p'), with a pristine control copy still reading its original count, and
was restored to a byte-identical file (shasum -a 256) before the next run.
Control build after all restores: exit 0, 1744 tests, 0 failures.
This commit was merged in pull request #559.
This commit is contained in:
2026-09-12 11:11:46 +02:00
6 changed files with 582 additions and 0 deletions
@@ -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 <em>completed</em> 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.
*
* <p>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.
*
* <p>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;
}
}
@@ -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();
}
}
@@ -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();
}
}
@@ -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");
}
}
@@ -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() {
}
}
}
@@ -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();
}
}
}