Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha 5fede82468 #382: give SpawnRequest a withProfile wither, guard it against the arity trap
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 1m49s
CompositePeerLauncher:372 rebuilt a routed SpawnRequest from a literal
new SpawnRequest(...) call listing six of the original request's own
accessors. That call is only correct because it happens to match the
canonical 6-arg constructor today; add a 7th component plus the
established back-compat constructor at the old (now-shorter) arity and
this call would silently rebind to it, dropping the new field on every
profile-routed spawn with no compile error — the same defect shape
already guarded on FleetConfig.withDefaults() (#357) and MemberSession
(#358).

Add SpawnRequest.withProfile(String), modeled on
MemberSession.withState/withActivity, and use it at the call site
instead. Add a guard test that resolves the true canonical constructor
by exact component types (never by argument count), gives every
component a distinctive value, and asserts every component but
profileName survives withProfile() unchanged.

Proved the guard against the real mechanism: temporarily dropped the
last (role) argument from withProfile()'s constructor call so it bound
to the 5-arg back-compat constructor — it still compiled, and the new
test failed, catching the silently-defaulted role. Restored the fix
and reconfirmed green.
2026-09-10 06:44:15 +07:00
6 changed files with 152 additions and 199 deletions
@@ -584,12 +584,8 @@ 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, () -> TimeUnit.MILLISECONDS.toNanos(System.currentTimeMillis()),
cfg.health().intervalOrDefault(),
System::nanoTime, cfg.health().intervalOrDefault(),
cfg.health().workingSuspectAfterOrDefault(), messages::abandon);
String coverage = FleetHealthMonitor.coverage(true,
cfg.health().notifications() != null && cfg.health().notifications().configured());
@@ -38,46 +38,15 @@ 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);
@@ -119,29 +88,12 @@ 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");
}
@@ -171,7 +123,6 @@ 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());
@@ -182,7 +133,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
&& stallElapsedNanos(session, nowNanos, driftBeforeThisTick) >= workingSuspectAfterNanos;
&& nowNanos - session.lastActivityAtNanos() >= 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());
@@ -200,8 +151,6 @@ 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);
@@ -226,65 +175,6 @@ 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;
@@ -369,8 +369,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
// CB-547a: route the chosen profile but keep the caller's session identity — dropping it
// here would silently sever the resume handle on every policy-routed spawn. CB-557: the
// role rides along for the same reason, or a routed spawn would be labelled as a dev.
SpawnRequest routedReq = new SpawnRequest(chosen.profile(), req.requestedCwd(), req.callerCwd(),
req.sessionName(), req.resumeSessionId(), req.role());
SpawnRequest routedReq = req.withProfile(chosen.profile());
try {
PeerHandle handle = d.spawn(routedReq);
spawnedBy.put(handle.id(), d);
@@ -29,4 +29,9 @@ public record SpawnRequest(String profileName, String requestedCwd, String calle
String sessionName, String resumeSessionId) {
this(profileName, requestedCwd, callerCwd, sessionName, resumeSessionId, null);
}
/** Return a copy of this request with {@code profileName} replaced by {@code profile}. */
public SpawnRequest withProfile(String profile) {
return new SpawnRequest(profile, requestedCwd, callerCwd, sessionName, resumeSessionId, role);
}
}
@@ -206,87 +206,6 @@ 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();
@@ -0,0 +1,144 @@
package dev.ltms.fleet.peer;
import org.junit.jupiter.api.Test;
import java.lang.reflect.Constructor;
import java.lang.reflect.RecordComponent;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.TreeSet;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* Fleetd #382, the same "defect factory" #357 and #358 guarded on {@code FleetConfig.withDefaults()}
* and {@link dev.ltms.fleet.session.MemberSession}'s rebuild sites — reproduced here on
* {@link SpawnRequest#withProfile(String)}.
*
* <p>{@code SpawnRequest} has back-compat constructors at arity 3 and 5 alongside its canonical
* arity-6 constructor. Before this ticket, {@code CompositePeerLauncher} routed a profile by
* building a fresh {@code SpawnRequest} from a literal {@code new SpawnRequest(...)} call listing
* six of the original request's own accessors. That call is only correct because it happens to
* name exactly six arguments today — add a 7th component and the established back-compat pattern
* (a new constructor at the old, now-shorter arity) and a call one argument short of the new
* canonical arity would silently rebind to that back-compat constructor, dropping the new
* component on every profile-routed spawn without any compile error. {@link #withProfile} replaces
* that literal call, so this test guards the ONE rebuild site instead of a call site scattered
* through a launcher.
*
* <p>The check below builds a {@link SpawnRequest} through the TRUE canonical constructor —
* resolved by the record's own component types via {@code getDeclaredConstructor}, never by
* argument count, so it can never itself land on a back-compat overload — with a real, distinctive,
* non-null value in every component, calls {@link SpawnRequest#withProfile(String)}, and asserts
* every component the method is not documented to change survives unchanged, while {@code
* profileName} comes back as the new value it was given. A component that comes back anything else
* was silently dropped — the shape of the defect this test exists to catch.
*
* <p>{@link #EXCLUDED_FROM_SURVIVAL_CHECK} is kept deliberately empty and size-pinned by
* {@link #exclusionListSizeIsPinned()}, for the same reason the other two guards pin theirs at
* zero: a checker whose escape hatch can grow to silence a failure is not a checker. Every one of
* {@link SpawnRequest}'s 6 current components has a real, non-null value here and none is excluded.
*/
class SpawnRequestWithProfilePreservesEveryComponentTest {
private static final RecordComponent[] COMPONENTS = SpawnRequest.class.getRecordComponents();
/** Deliberately empty today; grow it only with a matching justification, and re-pin the size. */
private static final Set<String> EXCLUDED_FROM_SURVIVAL_CHECK = Set.of();
/** One real, distinctive, non-null value per component — none of the 6 is excluded. */
private static Map<String, Object> baseValues() {
Map<String, Object> v = new LinkedHashMap<>();
v.put("profileName", "profile-guard");
v.put("requestedCwd", "/wt/requested-guard");
v.put("callerCwd", "/wt/caller-guard");
v.put("sessionName", "session-guard");
v.put("resumeSessionId", "resume-guard");
v.put("role", MemberRole.REVIEWER);
assertNamesMatchComponents(v);
return v;
}
/**
* Guards {@link #baseValues()} itself against drifting from the record's real shape — forgetting
* to add a new component here fails this assertion by name, rather than silently checking one
* component fewer than the record has.
*/
private static void assertNamesMatchComponents(Map<String, Object> values) {
Set<String> names = new TreeSet<>();
for (RecordComponent rc : COMPONENTS) {
names.add(rc.getName());
}
assertEquals(names, new TreeSet<>(values.keySet()),
"this test's value map has drifted from SpawnRequest's actual components — "
+ "update baseValues() alongside the record");
}
/**
* Builds a {@link SpawnRequest} through the TRUE canonical constructor — resolved by the
* record's own component types, not by argument count — so this never accidentally exercises a
* back-compat overload the way a literal {@code new SpawnRequest(...)} call risks doing.
*/
private static SpawnRequest requestOf(Map<String, Object> values) throws ReflectiveOperationException {
Class<?>[] types = Arrays.stream(COMPONENTS).map(RecordComponent::getType).toArray(Class<?>[]::new);
Object[] args = Arrays.stream(COMPONENTS).map(rc -> values.get(rc.getName())).toArray();
Constructor<SpawnRequest> ctor = SpawnRequest.class.getDeclaredConstructor(types);
return ctor.newInstance(args);
}
@Test
void exclusionListSizeIsPinned() {
assertEquals(0, EXCLUDED_FROM_SURVIVAL_CHECK.size(),
"EXCLUDED_FROM_SURVIVAL_CHECK grew from 0 — every entry needs a justification in "
+ "this test class's javadoc AND this assertion re-pinned to the new size; a "
+ "growing exclusion list that silences failures on its own is not a guard");
}
@Test
void withProfilePreservesEveryOtherComponent() throws ReflectiveOperationException {
Map<String, Object> base = baseValues();
SpawnRequest request = requestOf(base);
SpawnRequest result = request.withProfile("profile-updated");
Map<String, Object> expectedOverrides = Map.of("profileName", "profile-updated");
List<String> dropped = new ArrayList<>();
int checked = 0;
for (RecordComponent rc : COMPONENTS) {
String name = rc.getName();
if (EXCLUDED_FROM_SURVIVAL_CHECK.contains(name)) {
continue;
}
checked++;
Object expected = expectedOverrides.containsKey(name) ? expectedOverrides.get(name) : base.get(name);
Object actual;
try {
actual = rc.getAccessor().invoke(result);
} catch (ReflectiveOperationException e) {
throw new RuntimeException("failed to read SpawnRequest." + name + "()", e);
}
if (!Objects.equals(expected, actual)) {
dropped.add(String.format(Locale.ROOT,
"%s: withProfile() was expected to carry (%s) for '%s' but returned %s — a "
+ "component silently dropped by withProfile(), the shape of the defect "
+ "this test exists to catch (its final \"return new SpawnRequest(...)\" "
+ "call binding to a back-compat constructor instead of the true "
+ "canonical one)",
name, expected, name, actual));
}
}
System.out.printf(Locale.ROOT,
"SpawnRequest.withProfile() component-survival coverage — %d components, %d checked, "
+ "%d excluded, %d survived%n",
COMPONENTS.length, checked, EXCLUDED_FROM_SURVIVAL_CHECK.size(), checked - dropped.size());
assertEquals(List.of(), dropped,
"withProfile() silently dropped these components: " + dropped);
}
}