Merge CB-641: wire herdr and time evidence into the fleet health monitor

This commit is contained in:
Dai Ha
2026-08-27 22:06:50 +07:00
5 changed files with 206 additions and 17 deletions
+5 -6
View File
@@ -93,12 +93,11 @@ bind:
# intervalSeconds → how often a tick runs (default 30). ENFORCED floor of 15: the code computes
# Math.max(15, intervalSeconds), so a lower value is silently raised, not
# rejected.
# workingSuspectAfterSeconds, paneProbeIntervalSeconds → accepted and parsed, but NOT YET READ by
# anything — the dormant monitor only consumes intervalSeconds today (CB-573
# shipped ahead of the evidence publishers these two knobs are for). Setting
# them changes nothing right now, and no minimum is enforced on either, because
# nothing reads them to enforce one. They exist so a later build can start
# honouring them without another config-shape change.
# workingSuspectAfterSeconds → age before a BUSY member is suspected of a stall (default 600).
# ENFORCED floor of 300: a lower value is silently raised.
# paneProbeIntervalSeconds → accepted and parsed, but NOT YET READ by anything. Setting it changes
# nothing right now. It exists so a later build can start honouring it without
# another config-shape change.
# notifications.mode → "webhook" flips what fleet_list REPORTS (healthCoverage: "full" instead
# of "detection-only") — it does NOT make fleetd send any webhook call; no
# delivery mechanism is implemented yet. Any other value, or omitting the
@@ -449,7 +449,8 @@ public final class Fleetd {
// 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.
healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler,
System::nanoTime, cfg.health().intervalOrDefault(), messages::abandon);
System::nanoTime, cfg.health().intervalOrDefault(),
cfg.health().workingSuspectAfterOrDefault(), messages::abandon);
String coverage = FleetHealthMonitor.coverage(true,
cfg.health().notifications() != null && cfg.health().notifications().configured());
if ("detection-only".equals(coverage)) {
@@ -642,6 +642,9 @@ public record FleetConfig(
Integer paneProbeIntervalSeconds, Notifications notifications) {
public boolean isEnabled() { return Boolean.TRUE.equals(enabled); }
public int intervalOrDefault() { return Math.max(15, intervalSeconds == null ? 30 : intervalSeconds); }
public int workingSuspectAfterOrDefault() {
return Math.max(300, workingSuspectAfterSeconds == null ? 600 : workingSuspectAfterSeconds);
}
public record Notifications(String mode) {
public boolean configured() { return "webhook".equalsIgnoreCase(mode); }
}
@@ -3,6 +3,7 @@ package dev.ltms.fleet.health;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.session.MemberSession;
import org.slf4j.Logger;
@@ -25,6 +26,8 @@ public final class FleetHealthMonitor {
/** Bounded attempts to run {@link #failTarget} for one transition. Never retried tick-to-tick (CB-580). */
static final int MAX_FAIL_TARGET_ATTEMPTS = 3;
// CB-641: Match the injector's 60s readiness gate so health allows a full first boot.
static final long READINESS_GRACE_NANOS = TimeUnit.SECONDS.toNanos(60);
private final AgentControl agents;
private final Supplier<List<MemberSession>> roster;
@@ -32,6 +35,7 @@ public final class FleetHealthMonitor {
private final ScheduledExecutorService scheduler;
private final LongSupplier clock;
private final long intervalSeconds;
private final long workingSuspectAfterNanos;
private final BiConsumer<String, String> failTarget;
private final Map<String, HealthPrior> priors = new HashMap<>();
private final Map<String, HealthState> states = new HashMap<>();
@@ -48,13 +52,14 @@ public final class FleetHealthMonitor {
*/
public FleetHealthMonitor(AgentControl agents, Supplier<List<MemberSession>> roster, MessageService messages,
ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds,
BiConsumer<String, String> failTarget) {
long workingSuspectAfterSeconds, BiConsumer<String, String> failTarget) {
this.agents = agents;
this.roster = roster;
this.messages = messages;
this.scheduler = scheduler;
this.clock = clock;
this.intervalSeconds = intervalSeconds;
this.workingSuspectAfterNanos = TimeUnit.SECONDS.toNanos(workingSuspectAfterSeconds);
this.failTarget = Objects.requireNonNull(failTarget, "failTarget");
}
@@ -69,28 +74,43 @@ public final class FleetHealthMonitor {
// Package-private so tests can run one tick without waiting.
void tick() {
try {
List<Agent> agentsNow = agents.list(); // Exactly one list call for this complete observation.
List<MemberSession> rosterNow = roster.get(); // One in-memory roster snapshot for this tick.
List<Agent> agentsNow;
boolean controlLinkDown = false;
try {
agentsNow = agents.list(); // Exactly one list call for this complete observation.
} catch (HerdrException error) {
agentsNow = List.of();
controlLinkDown = true;
log.warn("fleet health control link unavailable; classifying roster", error);
}
Map<String, Agent> live = new HashMap<>();
for (Agent agent : agentsNow) live.put(agent.terminalId(), agent);
HashSet<String> current = new HashSet<>();
long nowNanos = clock.getAsLong();
for (MemberSession session : rosterNow) {
current.add(session.terminalId());
Agent agent = live.get(session.terminalId());
AgentStatus status = agent == null ? AgentStatus.UNKNOWN : agent.status();
boolean accepted = messages.hasAcceptedDelivery(session.terminalId());
boolean present = agent != null;
boolean targetNotFound = !controlLinkDown && !present
&& session.state() != MemberSession.State.SPAWNING;
boolean readinessGraceElapsed = nowNanos - session.spawnedAtNanos() >= READINESS_GRACE_NANOS;
boolean stalled = session.state() == MemberSession.State.BUSY
&& nowNanos - session.lastActivityAtNanos() >= workingSuspectAfterNanos;
HealthSnapshot snapshot = new HealthSnapshot(session.state(), status, accepted, NOT_YET_OBSERVED,
messages.hasInboxMessage(session.terminalId()), agent != null, NOT_YET_OBSERVED,
NOT_YET_OBSERVED, NOT_YET_OBSERVED, NOT_YET_OBSERVED, NOT_YET_OBSERVED, NOT_YET_OBSERVED);
messages.hasInboxMessage(session.terminalId()), present, targetNotFound, controlLinkDown,
readinessGraceElapsed, NOT_YET_OBSERVED, NOT_YET_OBSERVED, stalled);
HealthDecision decision = decide(snapshot, priors.getOrDefault(session.terminalId(), HealthPrior.NONE),
clock.getAsLong());
nowNanos);
priors.put(session.terminalId(), decision.prior());
reportTransition(session.terminalId(), decision.state());
}
priors.keySet().retainAll(current);
states.keySet().retainAll(current);
} catch (Throwable error) {
// A list failure is health evidence, and must never kill the monitor's only scheduler task.
// Any unclassified collection failure must never kill the monitor's only scheduler task.
log.warn("fleet health collection failed; will retry next tick", error);
} finally {
if (!scheduler.isShutdown()) {
@@ -1,5 +1,6 @@
package dev.ltms.fleet.health;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
@@ -9,6 +10,8 @@ import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.msg.InMemoryReplyInbox;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.herdr.WorkspaceControl;
@@ -17,10 +20,15 @@ import dev.ltms.fleet.config.FleetConfig;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.BiConsumer;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -40,7 +48,7 @@ class FleetHealthMonitorTest {
MessageService messages = new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox());
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, sessions::roster, messages, scheduler, () -> 1, 60,
(_, _) -> { });
600, (_, _) -> { });
monitor.tick();
monitor.stop();
assertEquals(1, herdr.calls.stream().filter(call -> call.method().equals("agent.list")).count());
@@ -52,7 +60,7 @@ class FleetHealthMonitorTest {
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, java.util.List::of,
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
scheduler, () -> 1, 60, (_, _) -> { });
scheduler, () -> 1, 60, 600, (_, _) -> { });
monitor.tick();
herdr.healthy(true);
monitor.tick();
@@ -71,7 +79,7 @@ class FleetHealthMonitorTest {
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, java.util.List::of,
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
scheduler, () -> 1, 60, (_, _) -> { });
scheduler, () -> 1, 60, 600, (_, _) -> { });
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
monitor.stop();
@@ -82,6 +90,148 @@ class FleetHealthMonitorTest {
}
}
@Test void failedListMarksEveryMemberControlLinkDownAndReschedules() {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
FakeHerdr herdr = new FakeHerdr().healthy(false);
ScheduledThreadPoolExecutor scheduler = new ScheduledThreadPoolExecutor(1);
FleetHealthMonitor monitor = monitor(herdr, List.of(
member("term_one", MemberSession.State.READY, 0, 0),
member("term_two", MemberSession.State.BUSY, 0, 0)), scheduler, () -> 1, 600,
(_, _) -> { });
monitor.tick();
assertEquals(1, appender.list.stream().filter(event -> event.getFormattedMessage()
.contains("member=term_one state=CONTROL_LINK_DOWN")).count());
assertEquals(1, appender.list.stream().filter(event -> event.getFormattedMessage()
.contains("member=term_two state=CONTROL_LINK_DOWN")).count());
assertEquals(1, scheduler.getQueue().size());
monitor.stop();
} finally {
logger.detachAppender(appender);
}
}
@Test void missingRosterMemberGoesGoneAndFailsTargetOnceAcrossTicks() {
FakeHerdr herdr = new FakeHerdr();
RecordingFailTarget failTarget = new RecordingFailTarget();
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = monitor(herdr,
List.of(member("term_missing", MemberSession.State.READY, 0, 0)), scheduler,
() -> 1, 600, failTarget);
monitor.tick();
monitor.tick();
monitor.stop();
assertEquals(1, failTarget.calls.size());
assertEquals("term_missing", failTarget.calls.get(0).target());
assertTrue(failTarget.calls.get(0).reason().contains("GONE"));
}
@Test void spawningMemberBecomesNeverReadyOnlyAfterReadinessGrace() {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
AtomicLong clock = new AtomicLong(FleetHealthMonitor.READINESS_GRACE_NANOS - 1);
RecordingFailTarget failTarget = new RecordingFailTarget();
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = monitor(new FakeHerdr(),
List.of(member("term_starting", MemberSession.State.SPAWNING, 0, 0)), scheduler,
clock::get, 600, failTarget);
monitor.tick();
assertEquals(0, failTarget.calls.size());
clock.set(FleetHealthMonitor.READINESS_GRACE_NANOS);
monitor.tick();
monitor.stop();
assertEquals(1, failTarget.calls.size());
assertTrue(failTarget.calls.get(0).reason().contains("NEVER_READY"));
assertTrue(appender.list.stream().anyMatch(event -> event.getFormattedMessage()
.contains("state=NEVER_READY previous=STARTING")));
} finally {
logger.detachAppender(appender);
}
}
@Test void workingSuspectConfigChangesTheStallBoundary() {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
long nowNanos = TimeUnit.SECONDS.toNanos(600);
FakeHerdr herdr = new FakeHerdr()
.withAgent("short", "term_short", "pane_short", "tab_short")
.withAgent("long", "term_long", "pane_long", "tab_long");
FleetConfig.Health shortConfig = new FleetConfig.Health(true, 30, 300, null, null);
FleetConfig.Health longConfig = new FleetConfig.Health(true, 30, 601, null, null);
var shortScheduler = Executors.newSingleThreadScheduledExecutor();
var longScheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor shortMonitor = monitor(herdr,
List.of(member("term_short", MemberSession.State.BUSY, 0, 0)), shortScheduler,
() -> nowNanos, shortConfig.workingSuspectAfterOrDefault(), (_, _) -> { });
FleetHealthMonitor longMonitor = monitor(herdr,
List.of(member("term_long", MemberSession.State.BUSY, 0, 0)), longScheduler,
() -> nowNanos, longConfig.workingSuspectAfterOrDefault(), (_, _) -> { });
shortMonitor.tick();
longMonitor.tick();
shortMonitor.stop();
longMonitor.stop();
assertTrue(appender.list.stream().anyMatch(event -> event.getFormattedMessage()
.contains("member=term_short state=STALL_SUSPECTED")));
assertEquals(0, appender.list.stream().filter(event -> event.getFormattedMessage()
.contains("member=term_long state=STALL_SUSPECTED")).count());
} finally {
logger.detachAppender(appender);
}
}
@Test void workingSuspectConfigKeepsItsDefaultAndFloor() {
assertEquals(600, new FleetConfig.Health(true, null, null, null, null)
.workingSuspectAfterOrDefault());
assertEquals(300, new FleetConfig.Health(true, null, 1, null, null)
.workingSuspectAfterOrDefault());
}
@Test void goneMemberRecoveryLogsOnceWithoutRefiringTargetFailure() {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
Level previousLevel = logger.getLevel();
logger.setLevel(Level.INFO);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
FakeHerdr herdr = new FakeHerdr();
RecordingFailTarget failTarget = new RecordingFailTarget();
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = monitor(herdr,
List.of(member("term_recovered", MemberSession.State.READY, 0, 0)), scheduler,
() -> 1, 600, failTarget);
monitor.tick();
herdr.withAgent("recovered", "term_recovered", "pane_recovered", "tab_recovered");
monitor.tick();
monitor.stop();
assertEquals(1, failTarget.calls.size());
assertTrue(appender.list.stream().anyMatch(event -> event.getFormattedMessage()
.contains("member=term_recovered recovered state=IDLE previous=GONE")));
} finally {
logger.detachAppender(appender);
logger.setLevel(previousLevel);
}
}
// --- CB-580: a member that reaches GONE/NEVER_READY must fail its waiting tickets
private static FleetHealthMonitor monitorWith(BiConsumer<String, String> failTarget) {
@@ -90,7 +240,23 @@ class FleetHealthMonitorTest {
var scheduler = Executors.newSingleThreadScheduledExecutor();
return new FleetHealthMonitor(agents, java.util.List::of,
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
scheduler, () -> 1, 60, failTarget);
scheduler, () -> 1, 60, 600, failTarget);
}
private static FleetHealthMonitor monitor(FakeHerdr herdr, List<MemberSession> roster,
java.util.concurrent.ScheduledExecutorService scheduler,
LongSupplier clock, 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, 60, workingSuspectAfterSeconds, failTarget);
}
private static MemberSession member(String terminalId, MemberSession.State state,
long spawnedAtNanos, long lastActivityAtNanos) {
return new MemberSession("pane-" + terminalId, terminalId, "test", MemberRole.DEV, "/tmp", null,
spawnedAtNanos, lastActivityAtNanos, 0, state, null, null);
}
@Test void terminalTransitionFailsTheTargetOnce() {