CB-643: join the message-layer evidence to the health monitor
CI / build (push) Successful in 1m8s
CI / contract (push) Successful in 1m9s

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.
This commit is contained in:
Dai Ha
2026-08-27 22:15:39 +07:00
parent a6095743f0
commit 7c4170ff6d
3 changed files with 146 additions and 7 deletions
+3 -3
View File
@@ -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
@@ -39,9 +39,26 @@ public final class FleetHealthMonitor {
private final BiConsumer<String, String> failTarget;
private final Map<String, HealthPrior> priors = new HashMap<>();
private final Map<String, HealthState> 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<String, Integer> 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;
@@ -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<String, String> 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<ILoggingEvent> 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<ILoggingEvent> 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<ILoggingEvent> 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);
}
}
}