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();