From 769f28240878e6ee7e6d6277f8ae080090f9c255 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 10 Sep 2026 06:53:30 +0700 Subject: [PATCH] #386: give the stall detector a real-time clock, log the divergence FleetHealthMonitor.tick's stalled check compared two monotonic-clock readings (System.nanoTime(), which macOS freezes across a host sleep), so a member BUSY for 101 real minutes was never flagged. The monitor now also takes a wall-clock LongSupplier (realtimeClock), used only inside the stall check. Each tick measures how far the two clocks moved apart since the previous tick and folds any positive divergence into a running total; when a single tick's divergence exceeds one tick interval (the signature of a sleep, since a tick cannot run while the process itself is suspended) it logs one WARN naming how long the detector could not see. The correction is applied per member, keyed to when that member's current lastActivityAtNanos was first observed BUSY - not since monitor start - so a sleep that happened before a member went busy is never charged to it. Every other use of the monitor's clock (readiness grace, snapshot timestamp) is unchanged. Backend quarantine/cool-off, the lead tab scan, the completion resolver, the session reaper and the message service TTLs are untouched, per the ticket's decision. Existing FleetHealthMonitor/FleetHealth tests pass unmodified (none of them ticks a BUSY session more than once, so the drift path never engages for them). Two new tests: a frozen monotonic clock past the real-time threshold produces STALL_SUSPECTED, and a single sleep gap logs the divergence exactly once, not once per tick. --- .../src/main/java/dev/ltms/fleet/Fleetd.java | 6 +- .../ltms/fleet/health/FleetHealthMonitor.java | 112 +++++++++++++++++- .../fleet/health/FleetHealthMonitorTest.java | 81 +++++++++++++ 3 files changed, 197 insertions(+), 2 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index a52227e..b7e0829 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -584,8 +584,12 @@ public final class Fleetd { if (cfg.health() != null && cfg.health().isEnabled()) { // CB-580: a member found GONE/NEVER_READY must fail whatever ticket is waiting on it, // through the same idempotent target-wide operation CB-516 already uses on release. + // fleetd #386: System::nanoTime freezes across a macOS sleep, so the stall check also + // gets a wall-clock source to detect and correct for that freeze. Every other decision + // in FleetHealthMonitor stays on the monotonic clock, unchanged. healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler, - System::nanoTime, cfg.health().intervalOrDefault(), + System::nanoTime, () -> TimeUnit.MILLISECONDS.toNanos(System.currentTimeMillis()), + cfg.health().intervalOrDefault(), cfg.health().workingSuspectAfterOrDefault(), messages::abandon); String coverage = FleetHealthMonitor.coverage(true, cfg.health().notifications() != null && cfg.health().notifications().configured()); diff --git a/fleetd/src/main/java/dev/ltms/fleet/health/FleetHealthMonitor.java b/fleetd/src/main/java/dev/ltms/fleet/health/FleetHealthMonitor.java index cfa11b3..5482542 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/health/FleetHealthMonitor.java +++ b/fleetd/src/main/java/dev/ltms/fleet/health/FleetHealthMonitor.java @@ -38,15 +38,46 @@ public final class FleetHealthMonitor { */ static final long ASK_LAPSE_RECHECK_DELAY_SECONDS = 120; + /** + * fleetd #386: {@code System.nanoTime()} (or whatever {@link #clock} is) does not advance while + * macOS sleeps, so a raw {@code nowNanos - lastActivityAtNanos} comparison freezes with the + * host and can never cross {@link #workingSuspectAfterNanos}. This is a second, wall-clock + * source used ONLY inside the stall check ({@link #stallElapsedNanos}) to detect and correct + * for that freeze. Nothing else in this class reads it — every other decision (readiness grace, + * the fault classification itself) stays exactly on {@link #clock}, as the ticket requires. + */ + private static final LongSupplier DEFAULT_REALTIME_CLOCK = + () -> TimeUnit.MILLISECONDS.toNanos(System.currentTimeMillis()); + private final AgentControl agents; private final Supplier> roster; private final MessageService messages; private final ScheduledExecutorService scheduler; private final LongSupplier clock; + private final LongSupplier realtimeClock; private final long intervalSeconds; + private final long tickIntervalNanos; private final long workingSuspectAfterNanos; private final BiConsumer failTarget; private final Map priors = new HashMap<>(); + /** + * fleetd #386 clock-drift bookkeeping. {@code haveClockBaseline}/{@code lastTickMonoNanos}/ + * {@code lastTickRealNanos} track the previous tick's pair of readings so each new tick can + * measure how far the two clocks moved apart since then. {@code accumulatedDriftNanos} is the + * running total of every such divergence observed since this monitor started (never decreases — + * the monotonic clock can only lag real time, never lead it). {@code busyDriftBaselineNanos}/ + * {@code busyBaselineActivityNanos} record, per target, the value of {@code accumulatedDriftNanos} + * at the moment this monitor first saw that target's CURRENT {@code lastActivityAtNanos} while + * BUSY — so {@link #stallElapsedNanos} adds back only the drift observed DURING this BUSY span, + * never drift from a sleep that happened before the member went busy. All five fields are touched + * only from {@code tick()}, like {@link #priors}. + */ + private boolean haveClockBaseline = false; + private long lastTickMonoNanos; + private long lastTickRealNanos; + private long accumulatedDriftNanos = 0; + private final Map busyDriftBaselineNanos = new HashMap<>(); + private final Map busyBaselineActivityNanos = new HashMap<>(); /** * The live classification per member, and the only one of this class's three maps that more * than one scheduler task touches. {@code tick} writes it (and prunes it to the roster); @@ -88,12 +119,29 @@ public final class FleetHealthMonitor { public FleetHealthMonitor(AgentControl agents, Supplier> roster, MessageService messages, ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds, long workingSuspectAfterSeconds, BiConsumer failTarget) { + this(agents, roster, messages, scheduler, clock, DEFAULT_REALTIME_CLOCK, intervalSeconds, + workingSuspectAfterSeconds, failTarget); + } + + /** + * @param realtimeClock fleetd #386: a wall-clock nanosecond source (e.g. + * {@code System.currentTimeMillis()} converted to nanos) that keeps + * advancing while {@code clock} is frozen by a host sleep. Used only to + * correct the stall check — see the class-level javadoc on the + * clock-drift fields. + */ + public FleetHealthMonitor(AgentControl agents, Supplier> roster, MessageService messages, + ScheduledExecutorService scheduler, LongSupplier clock, LongSupplier realtimeClock, + long intervalSeconds, long workingSuspectAfterSeconds, + BiConsumer failTarget) { this.agents = agents; this.roster = roster; this.messages = messages; this.scheduler = scheduler; this.clock = clock; + this.realtimeClock = Objects.requireNonNull(realtimeClock, "realtimeClock"); this.intervalSeconds = intervalSeconds; + this.tickIntervalNanos = TimeUnit.SECONDS.toNanos(intervalSeconds); this.workingSuspectAfterNanos = TimeUnit.SECONDS.toNanos(workingSuspectAfterSeconds); this.failTarget = Objects.requireNonNull(failTarget, "failTarget"); } @@ -123,6 +171,7 @@ public final class FleetHealthMonitor { for (Agent agent : agentsNow) live.put(agent.terminalId(), agent); HashSet current = new HashSet<>(); long nowNanos = clock.getAsLong(); + long driftBeforeThisTick = observeClockDrift(nowNanos); for (MemberSession session : rosterNow) { current.add(session.terminalId()); Agent agent = live.get(session.terminalId()); @@ -133,7 +182,7 @@ public final class FleetHealthMonitor { && session.state() != MemberSession.State.SPAWNING; boolean readinessGraceElapsed = nowNanos - session.spawnedAtNanos() >= READINESS_GRACE_NANOS; boolean stalled = session.state() == MemberSession.State.BUSY - && nowNanos - session.lastActivityAtNanos() >= workingSuspectAfterNanos; + && stallElapsedNanos(session, nowNanos, driftBeforeThisTick) >= workingSuspectAfterNanos; // CB-643: the three message-layer facts CB-640 published. Read them here rather than // leaving them false — that constant is what made 8 of the 9 fault states dead. boolean queuedDelivery = messages.hasQueuedDelivery(session.terminalId()); @@ -151,6 +200,8 @@ public final class FleetHealthMonitor { priors.keySet().retainAll(current); states.keySet().retainAll(current); orphanStreaks.keySet().retainAll(current); + busyDriftBaselineNanos.keySet().retainAll(current); + busyBaselineActivityNanos.keySet().retainAll(current); } catch (Throwable error) { // Any unclassified collection failure must never kill the monitor's only scheduler task. log.warn("fleet health collection failed; will retry next tick", error); @@ -175,6 +226,65 @@ public final class FleetHealthMonitor { return streak >= ORPHAN_CONFIRM_TICKS; } + /** + * fleetd #386: compare this tick's monotonic and real-time readings against the previous + * tick's, and fold any positive divergence into {@link #accumulatedDriftNanos} (a ratchet — it + * never decreases, since the monotonic clock can only fall behind real time, never ahead of + * it). Logs once, at WARN, when that single tick's divergence exceeds one full tick interval — + * the signature of a host that slept between the two ticks (a tick literally cannot run while + * the process itself is suspended, so the whole sleep duration lands inside one tick's gap). + * + * @return {@link #accumulatedDriftNanos} as it stood BEFORE this tick's divergence was folded + * in — the baseline {@link #stallElapsedNanos} needs when a target is observed BUSY + * for the first time this tick, so a sleep that happened before this member went busy + * is not attributed to it. + */ + private long observeClockDrift(long nowNanos) { + long nowRealNanos = realtimeClock.getAsLong(); + long driftBeforeThisTick = accumulatedDriftNanos; + if (haveClockBaseline) { + long monoDelta = nowNanos - lastTickMonoNanos; + long realDelta = nowRealNanos - lastTickRealNanos; + long tickDrift = realDelta - monoDelta; + if (tickDrift > tickIntervalNanos) { + log.warn("fleet health: the monotonic clock did not advance for about {}s that the " + + "real clock did since the last tick (host likely slept); the stall " + + "detector could not see that time", TimeUnit.NANOSECONDS.toSeconds(tickDrift)); + } + if (tickDrift > 0) { + accumulatedDriftNanos = driftBeforeThisTick + tickDrift; + } + } + lastTickMonoNanos = nowNanos; + lastTickRealNanos = nowRealNanos; + haveClockBaseline = true; + return driftBeforeThisTick; + } + + /** + * fleetd #386: {@code nowNanos - lastActivityAtNanos} alone freezes across a host sleep, since + * both come from the monotonic {@link #clock}. This adds back the real-time drift observed + * since this BUSY span started — not the monitor's whole lifetime, so a sleep that happened + * before this member went busy never leaks into its stall reading (see the class-level javadoc + * on the drift fields). The baseline resets whenever {@code lastActivityAtNanos} changes (a new + * turn) or the member is not currently BUSY. + */ + private long stallElapsedNanos(MemberSession session, long nowNanos, long driftBeforeThisTick) { + String target = session.terminalId(); + if (session.state() != MemberSession.State.BUSY) { + busyDriftBaselineNanos.remove(target); + busyBaselineActivityNanos.remove(target); + return nowNanos - session.lastActivityAtNanos(); + } + Long baselineActivity = busyBaselineActivityNanos.get(target); + if (baselineActivity == null || baselineActivity != session.lastActivityAtNanos()) { + busyBaselineActivityNanos.put(target, session.lastActivityAtNanos()); + busyDriftBaselineNanos.put(target, driftBeforeThisTick); + } + long driftSinceBusyStart = accumulatedDriftNanos - busyDriftBaselineNanos.get(target); + return (nowNanos - session.lastActivityAtNanos()) + driftSinceBusyStart; + } + void reportTransition(String target, HealthState next) { HealthState previous = states.put(target, next); if (previous == next) return; diff --git a/fleetd/src/test/java/dev/ltms/fleet/health/FleetHealthMonitorTest.java b/fleetd/src/test/java/dev/ltms/fleet/health/FleetHealthMonitorTest.java index 6e83b8d..149422c 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/health/FleetHealthMonitorTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/health/FleetHealthMonitorTest.java @@ -206,6 +206,87 @@ class FleetHealthMonitorTest { .workingSuspectAfterOrDefault()); } + // --- fleetd #386: a stall detector whose only clock freezes with a sleeping host is worse + // than a silent one — it reports "quiet" for a member that was genuinely busy for hours. + + @Test void monotonicClockFrozenPastThresholdOnRealClockStillReportsStallSuspected() { + Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class); + ListAppender appender = new ListAppender<>(); + appender.start(); + logger.addAppender(appender); + try { + AtomicLong mono = new AtomicLong(0); + AtomicLong real = new AtomicLong(0); + FakeHerdr herdr = new FakeHerdr().withAgent("busy", "term_busy", "pane_busy", "tab_busy"); + var scheduler = Executors.newSingleThreadScheduledExecutor(); + FleetHealthMonitor monitor = monitorWithClocks(herdr, + List.of(member("term_busy", MemberSession.State.BUSY, 0, 0)), scheduler, + mono::get, real::get, 60, 600, (_, _) -> { }); + + monitor.tick(); // establishes the clock baseline; nothing has diverged yet + assertEquals(0, appender.list.stream().filter(event -> event.getFormattedMessage() + .contains("state=STALL_SUSPECTED")).count()); + + // The host "sleeps": the monotonic clock stands completely still while the real clock + // keeps moving, past the 600s stall threshold. + real.set(TimeUnit.SECONDS.toNanos(700)); + monitor.tick(); + monitor.stop(); + + assertTrue(appender.list.stream().anyMatch(event -> event.getFormattedMessage() + .contains("member=term_busy state=STALL_SUSPECTED")), + "the real clock crossed the stall threshold even though the monotonic clock never moved"); + } finally { + logger.detachAppender(appender); + } + } + + @Test void clockDivergenceIsLoggedOnceNotOncePerTick() { + Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class); + ListAppender appender = new ListAppender<>(); + appender.start(); + logger.addAppender(appender); + try { + AtomicLong mono = new AtomicLong(0); + AtomicLong real = new AtomicLong(0); + var scheduler = Executors.newSingleThreadScheduledExecutor(); + FleetHealthMonitor monitor = monitorWithClocks(new FakeHerdr(), List.of(), scheduler, + mono::get, real::get, 60, 600, (_, _) -> { }); + + monitor.tick(); // baseline: no divergence possible yet + + // One sleep gap: the monotonic clock is frozen while the real clock jumps far past one + // tick interval (60s). + real.set(TimeUnit.SECONDS.toNanos(700)); + monitor.tick(); + + // The host is awake again: both clocks advance together from here, so no more divergence. + mono.set(TimeUnit.SECONDS.toNanos(10)); + real.set(TimeUnit.SECONDS.toNanos(710)); + monitor.tick(); + mono.set(TimeUnit.SECONDS.toNanos(20)); + real.set(TimeUnit.SECONDS.toNanos(720)); + monitor.tick(); + monitor.stop(); + + assertEquals(1, appender.list.stream().filter(event -> event.getFormattedMessage() + .contains("the monotonic clock did not advance")) + .count(), "one sleep gap must produce exactly one divergence line, not one per tick"); + } finally { + logger.detachAppender(appender); + } + } + + private static FleetHealthMonitor monitorWithClocks(FakeHerdr herdr, List roster, + java.util.concurrent.ScheduledExecutorService scheduler, LongSupplier clock, + LongSupplier realtimeClock, long intervalSeconds, long workingSuspectAfterSeconds, + BiConsumer failTarget) { + AgentControl agents = new AgentControl(herdr); + return new FleetHealthMonitor(agents, () -> roster, + new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()), + scheduler, clock, realtimeClock, intervalSeconds, workingSuspectAfterSeconds, failTarget); + } + @Test void goneMemberRecoveryLogsOnceWithoutRefiringTargetFailure() { Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class); Level previousLevel = logger.getLevel();