fleetd #544: progress watchdog for StatusPoller and SessionReaper loops #559
@@ -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();
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user