From 66e776d178082c711e53b359763834dea71f2913 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 27 Aug 2026 22:03:20 +0700 Subject: [PATCH] CB-641: wire herdr health evidence --- fleetd/fleetd.example.yaml | 11 +- .../src/main/java/dev/ltms/fleet/Fleetd.java | 3 +- .../dev/ltms/fleet/config/FleetConfig.java | 3 + .../ltms/fleet/health/FleetHealthMonitor.java | 32 +++- .../fleet/health/FleetHealthMonitorTest.java | 174 +++++++++++++++++- 5 files changed, 206 insertions(+), 17 deletions(-) diff --git a/fleetd/fleetd.example.yaml b/fleetd/fleetd.example.yaml index 52b964a..7f7cc6f 100644 --- a/fleetd/fleetd.example.yaml +++ b/fleetd/fleetd.example.yaml @@ -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 diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index 67bf198..9b66efc 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -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)) { diff --git a/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java b/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java index b486f60..a3183e9 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java +++ b/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java @@ -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); } } 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 7131c71..1c68f2c 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/health/FleetHealthMonitor.java +++ b/fleetd/src/main/java/dev/ltms/fleet/health/FleetHealthMonitor.java @@ -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> 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 failTarget; private final Map priors = new HashMap<>(); private final Map states = new HashMap<>(); @@ -48,13 +52,14 @@ public final class FleetHealthMonitor { */ public FleetHealthMonitor(AgentControl agents, Supplier> roster, MessageService messages, ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds, - BiConsumer failTarget) { + long workingSuspectAfterSeconds, BiConsumer 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 agentsNow = agents.list(); // Exactly one list call for this complete observation. List rosterNow = roster.get(); // One in-memory roster snapshot for this tick. + List 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 live = new HashMap<>(); for (Agent agent : agentsNow) live.put(agent.terminalId(), agent); HashSet 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()) { 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 cea9824..fb6a002 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/health/FleetHealthMonitorTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/health/FleetHealthMonitorTest.java @@ -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 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 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 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 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 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 roster, + java.util.concurrent.ScheduledExecutorService scheduler, + LongSupplier clock, 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, 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() {