Merge #386: correct the stall check for a monotonic clock frozen by host sleep
This commit is contained in:
@@ -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());
|
||||
|
||||
@@ -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<List<MemberSession>> 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<String, String> failTarget;
|
||||
private final Map<String, HealthPrior> 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<String, Long> busyDriftBaselineNanos = new HashMap<>();
|
||||
private final Map<String, Long> 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<List<MemberSession>> roster, MessageService messages,
|
||||
ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds,
|
||||
long workingSuspectAfterSeconds, BiConsumer<String, String> 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<List<MemberSession>> roster, MessageService messages,
|
||||
ScheduledExecutorService scheduler, LongSupplier clock, LongSupplier realtimeClock,
|
||||
long intervalSeconds, long workingSuspectAfterSeconds,
|
||||
BiConsumer<String, String> 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<String> 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;
|
||||
|
||||
@@ -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<ILoggingEvent> 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<ILoggingEvent> 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<MemberSession> roster,
|
||||
java.util.concurrent.ScheduledExecutorService scheduler, LongSupplier clock,
|
||||
LongSupplier realtimeClock, long intervalSeconds, long workingSuspectAfterSeconds,
|
||||
BiConsumer<String, String> 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();
|
||||
|
||||
Reference in New Issue
Block a user