From 7c4170ff6ddc48908ec41245cc3a102ad4c5f418 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 27 Aug 2026 22:15:39 +0700 Subject: [PATCH] CB-643: join the message-layer evidence to the health monitor CB-640 published the three message-layer facts and CB-641 wired the herdr and time ones. This joins them, so every HealthSnapshot field now carries real evidence and the NOT_YET_OBSERVED placeholder is gone. That constant is what made 8 of the 9 fault states unreachable, GONE and NEVER_READY included, which is why CB-580's failTarget never fired. hasOrphanedDelegation is a true snapshot, but it can read true for one tick during an ordinary race: an async ticket exists before its virtual thread reaches rendezvous.open, so for that instant nothing is accepted or queued behind it. decide maps the field straight to DELEGATION_ORPHANED with no smoothing, so one racy read would log a fault that clears on the next tick. The monitor now requires two consecutive observations. That costs one interval on a real orphan and removes the false positive. Two tests drive real ticks against a genuinely orphaned ticket (an unanswered fleet_ask that lapsed back to PENDING), not the seam: one tick reports nothing, two report once, and a single clean tick in between resets the streak. Also correct two config comments. paneProbeIntervalSeconds is parsed and read by nothing, so its "minimum 60" note promised a floor that does not exist. 970 tests green. --- fleetd/fleetd.example.yaml | 6 +- .../ltms/fleet/health/FleetHealthMonitor.java | 46 +++++++- .../fleet/health/FleetHealthMonitorTest.java | 101 ++++++++++++++++++ 3 files changed, 146 insertions(+), 7 deletions(-) diff --git a/fleetd/fleetd.example.yaml b/fleetd/fleetd.example.yaml index 7f7cc6f..f878fed 100644 --- a/fleetd/fleetd.example.yaml +++ b/fleetd/fleetd.example.yaml @@ -104,9 +104,9 @@ bind: # block, reports "detection-only". # health: # enabled: true -# intervalSeconds: 30 -# workingSuspectAfterSeconds: 600 -# paneProbeIntervalSeconds: 60 +# intervalSeconds: 30 # floor 15 +# workingSuspectAfterSeconds: 600 # floor 300 — how long BUSY with no activity means STALL_SUSPECTED +# paneProbeIntervalSeconds: 60 # parsed, but nothing reads it yet — changing it changes nothing # notifications: # mode: disabled 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 1c68f2c..ee61559 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/health/FleetHealthMonitor.java +++ b/fleetd/src/main/java/dev/ltms/fleet/health/FleetHealthMonitor.java @@ -39,9 +39,26 @@ public final class FleetHealthMonitor { private final BiConsumer failTarget; private final Map priors = new HashMap<>(); private final Map states = new HashMap<>(); + /** + * CB-643: consecutive ticks on which a target looked like an orphaned delegation. The fact + * {@link MessageService#hasOrphanedDelegation} reports is a true snapshot, but it can read true + * for one tick during an ordinary race — an async ticket exists before its virtual thread has + * reached {@code rendezvous.open()}, so for that instant nothing is accepted or queued behind + * it. {@code decide} maps the field straight to {@code DELEGATION_ORPHANED} with no cross-tick + * smoothing of its own, so a single racy read would log a fault that clears on the next tick. + * Requiring two consecutive observations costs one interval of latency on a real orphan and + * removes that false positive entirely. + */ + private final Map orphanStreaks = new HashMap<>(); - // These facts need the evidence publishers introduced by later M4 units. They are not negatives. - private static final boolean NOT_YET_OBSERVED = false; + /** How many consecutive ticks a target must look orphaned before health reports it (CB-643). */ + static final int ORPHAN_CONFIRM_TICKS = 2; + + // CB-643: every HealthSnapshot field now carries real evidence. The NOT_YET_OBSERVED placeholder + // that stood in for 7 of the 12 is gone, and with it the reason 8 of the 9 fault states were + // unreachable — GONE and NEVER_READY included, which is what kept CB-580's failTarget from ever + // firing. Do not reintroduce a constant here: a field with no publisher is a dead state, and the + // tests pass either way, so nothing else will tell you. /** * @param failTarget CB-568's idempotent target-wide failure operation (e.g. {@code messages::abandon}), @@ -99,9 +116,15 @@ public final class FleetHealthMonitor { 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, + // 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()); + boolean replyStranded = messages.hasStrandedReply(session.terminalId()); + boolean orphanedDelegation = confirmOrphan(session.terminalId(), + messages.hasOrphanedDelegation(session.terminalId())); + HealthSnapshot snapshot = new HealthSnapshot(session.state(), status, accepted, queuedDelivery, messages.hasInboxMessage(session.terminalId()), present, targetNotFound, controlLinkDown, - readinessGraceElapsed, NOT_YET_OBSERVED, NOT_YET_OBSERVED, stalled); + readinessGraceElapsed, orphanedDelegation, replyStranded, stalled); HealthDecision decision = decide(snapshot, priors.getOrDefault(session.terminalId(), HealthPrior.NONE), nowNanos); priors.put(session.terminalId(), decision.prior()); @@ -109,6 +132,7 @@ public final class FleetHealthMonitor { } priors.keySet().retainAll(current); states.keySet().retainAll(current); + orphanStreaks.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); @@ -119,6 +143,20 @@ public final class FleetHealthMonitor { } } + /** + * Debounce {@link MessageService#hasOrphanedDelegation} across ticks (CB-643). Returns true only + * once {@code observed} has held for {@link #ORPHAN_CONFIRM_TICKS} consecutive ticks; a single + * false reading resets the streak, so a transient race never reaches the classifier. + */ + private boolean confirmOrphan(String target, boolean observed) { + if (!observed) { + orphanStreaks.remove(target); + return false; + } + int streak = orphanStreaks.merge(target, 1, Integer::sum); + return streak >= ORPHAN_CONFIRM_TICKS; + } + 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 fb6a002..08360a7 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/health/FleetHealthMonitorTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/health/FleetHealthMonitorTest.java @@ -5,6 +5,7 @@ import ch.qos.logback.classic.Logger; import ch.qos.logback.classic.spi.ILoggingEvent; import ch.qos.logback.core.read.ListAppender; import dev.ltms.fleet.herdr.AgentControl; +import dev.ltms.fleet.herdr.AgentStatus; import dev.ltms.fleet.herdr.FakeHerdr; import dev.ltms.fleet.inject.Injector; import dev.ltms.fleet.msg.InMemoryReplyInbox; @@ -333,4 +334,104 @@ class FleetHealthMonitorTest { throw new RuntimeException("boom"); } } + + // --- CB-643: the two-tick gate on an orphaned delegation ------------------------------------ + + private static final String ORPHAN_WARNING = "member=term_a state=DELEGATION_ORPHANED"; + + /** + * Everything one orphan test needs: a live target whose async ticket really is orphaned. The + * recipe is the one MessageServiceTest proves for {@code hasOrphanedDelegation} — an unanswered + * {@code fleet_ask} lapses, so the ticket returns to PENDING while its forward waiter is already + * closed. {@code rendezvous} is exposed so a test can turn the fact off and on again. + */ + private record OrphanFleet(FakeHerdr herdr, Rendezvous rendezvous, MessageService messages, + FleetHealthMonitor monitor) { + } + + private static OrphanFleet orphanedTarget(BiConsumer failTarget) throws Exception { + FakeHerdr herdr = new FakeHerdr().withAgent("worker", "term_a", "pane-term_a", "tab_a") + .readText("$ prompt"); + AgentControl agents = new AgentControl(herdr); + Rendezvous rendezvous = new Rendezvous(); + Injector injector = new Injector(agents); + InMemoryReplyInbox inbox = new InMemoryReplyInbox(); + inbox.own("term_a"); + MessageService messages = new MessageService(agents, injector, rendezvous, inbox); + + messages.sendAsync("term_a", "task that asks"); + long deadline = System.currentTimeMillis() + 2000; + while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + } + assertTrue(rendezvous.isWaiting("term_a"), "the async send should have opened its waiter"); + injector.onStatus("term_a", AgentStatus.IDLE); // deliver the task + injector.onStatus("term_a", AgentStatus.WORKING); // the worker picks it up + assertEquals(MessageService.AskOutcome.TIMED_OUT, + messages.ask("term_a", "which config?", 200).outcome()); + assertTrue(messages.hasOrphanedDelegation("term_a"), + "the lapsed ask should leave a PENDING ticket with nothing in flight"); + + FleetHealthMonitor monitor = new FleetHealthMonitor(agents, () -> List.of( + member("term_a", MemberSession.State.READY, 0, 0)), + messages, new ScheduledThreadPoolExecutor(1), () -> 1, 60, 600, failTarget); + return new OrphanFleet(herdr, rendezvous, messages, monitor); + } + + private static long orphanWarnings(ListAppender appender) { + return appender.list.stream() + .filter(event -> event.getFormattedMessage().contains(ORPHAN_WARNING)).count(); + } + + @Test void oneOrphanObservationIsNotYetReportedButTwoAre() throws Exception { + Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class); + ListAppender appender = new ListAppender<>(); + appender.start(); + logger.addAppender(appender); + OrphanFleet fleet = orphanedTarget((_, _) -> { }); + try { + fleet.monitor().tick(); + assertEquals(0, orphanWarnings(appender), + "one observation can be an ordinary race, so health must not report it yet"); + + fleet.monitor().tick(); + assertEquals(1, orphanWarnings(appender), + "a second consecutive observation confirms the orphan and is reported once"); + } finally { + fleet.messages().abandon("term_a", "test over"); + fleet.monitor().stop(); + logger.detachAppender(appender); + } + } + + @Test void aSingleCleanTickResetsTheOrphanStreak() throws Exception { + Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class); + ListAppender appender = new ListAppender<>(); + appender.start(); + logger.addAppender(appender); + OrphanFleet fleet = orphanedTarget((_, _) -> { }); + try { + fleet.monitor().tick(); // first observation: streak 1 + + // Something is in flight for the target again, so the orphan fact reads false. + var waiter = fleet.rendezvous().open("term_a"); + fleet.monitor().tick(); + assertEquals(0, orphanWarnings(appender), "a clean tick must clear the streak"); + + // Deregister that waiter — resolving it is not enough, the sender's close() is what + // removes it — so the ticket is orphaned again, from a streak of zero. + fleet.rendezvous().close("term_a", waiter); + assertTrue(fleet.messages().hasOrphanedDelegation("term_a")); + fleet.monitor().tick(); + assertEquals(0, orphanWarnings(appender), + "the streak restarted, so this first observation is not reported either"); + + fleet.monitor().tick(); + assertEquals(1, orphanWarnings(appender), "two consecutive observations report once"); + } finally { + fleet.messages().abandon("term_a", "test over"); + fleet.monitor().stop(); + logger.detachAppender(appender); + } + } }