CB-573: keep health monitoring after failures
This commit is contained in:
@@ -9,6 +9,7 @@ import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
@@ -26,6 +27,10 @@ public final class FleetHealthMonitor {
|
||||
private final LongSupplier clock;
|
||||
private final long intervalSeconds;
|
||||
private final Map<String, HealthPrior> priors = new HashMap<>();
|
||||
private final Map<String, HealthState> states = 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;
|
||||
|
||||
public FleetHealthMonitor(AgentControl agents, Supplier<List<MemberSession>> roster, MessageService messages,
|
||||
ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds) {
|
||||
@@ -47,21 +52,53 @@ public final class FleetHealthMonitor {
|
||||
|
||||
// Package-private so tests can run one tick without waiting.
|
||||
void tick() {
|
||||
List<Agent> agentsNow = agents.list(); // Exactly one list call for this complete observation.
|
||||
Map<String, Agent> live = new HashMap<>();
|
||||
for (Agent agent : agentsNow) live.put(agent.terminalId(), agent);
|
||||
for (MemberSession session : roster.get()) { // One in-memory roster snapshot for the same tick.
|
||||
Agent agent = live.get(session.terminalId());
|
||||
AgentStatus status = agent == null ? AgentStatus.UNKNOWN : agent.status();
|
||||
boolean accepted = messages.hasAcceptedDelivery(session.terminalId());
|
||||
HealthSnapshot snapshot = new HealthSnapshot(session.state(), status, accepted, false,
|
||||
messages.hasInboxMessage(session.terminalId()), agent != null, false, false,
|
||||
false, false, false, false);
|
||||
HealthDecision decision = decide(snapshot, priors.getOrDefault(session.terminalId(), HealthPrior.NONE),
|
||||
clock.getAsLong());
|
||||
priors.put(session.terminalId(), decision.prior());
|
||||
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.
|
||||
Map<String, Agent> live = new HashMap<>();
|
||||
for (Agent agent : agentsNow) live.put(agent.terminalId(), agent);
|
||||
HashSet<String> current = new HashSet<>();
|
||||
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());
|
||||
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);
|
||||
HealthDecision decision = decide(snapshot, priors.getOrDefault(session.terminalId(), HealthPrior.NONE),
|
||||
clock.getAsLong());
|
||||
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.
|
||||
log.warn("fleet health collection failed; will retry next tick", error);
|
||||
} finally {
|
||||
if (!scheduler.isShutdown()) {
|
||||
scheduler.schedule(this::tick, intervalSeconds, TimeUnit.SECONDS);
|
||||
}
|
||||
}
|
||||
scheduler.schedule(this::tick, intervalSeconds, TimeUnit.SECONDS);
|
||||
}
|
||||
|
||||
void reportTransition(String target, HealthState next) {
|
||||
HealthState previous = states.put(target, next);
|
||||
if (previous == next) return;
|
||||
if (fault(next)) {
|
||||
log.warn("fleet health member={} state={} previous={}", target, next, previous);
|
||||
} else if (previous != null && fault(previous)) {
|
||||
log.info("fleet health member={} recovered state={} previous={}", target, next, previous);
|
||||
}
|
||||
}
|
||||
|
||||
private static boolean fault(HealthState state) {
|
||||
return switch (state) {
|
||||
case NEVER_READY, GONE, TURN_BOUNDARY_LOST, ERROR_ON_SCREEN, STALL_SUSPECTED,
|
||||
MUTE, REPLY_STRANDED, DELEGATION_ORPHANED, CONTROL_LINK_DOWN -> true;
|
||||
default -> false;
|
||||
};
|
||||
}
|
||||
|
||||
public static String coverage(boolean enabled, boolean notificationConfigured) {
|
||||
|
||||
@@ -1,5 +1,8 @@
|
||||
package dev.ltms.bridged.health;
|
||||
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.bridged.herdr.AgentControl;
|
||||
import dev.ltms.bridged.herdr.FakeHerdr;
|
||||
import dev.ltms.bridged.inject.Injector;
|
||||
@@ -12,6 +15,7 @@ import dev.ltms.bridged.herdr.WorkspaceControl;
|
||||
import dev.ltms.bridged.guard.SubscriptionGuard;
|
||||
import dev.ltms.bridged.config.BridgedConfig;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
@@ -38,4 +42,40 @@ class FleetHealthMonitorTest {
|
||||
monitor.stop();
|
||||
assertEquals(1, herdr.calls.stream().filter(call -> call.method().equals("agent.list")).count());
|
||||
}
|
||||
|
||||
@Test void failedTickDoesNotStopTheNextTick() {
|
||||
FakeHerdr herdr = new FakeHerdr().healthy(false);
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
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);
|
||||
monitor.tick();
|
||||
herdr.healthy(true);
|
||||
monitor.tick();
|
||||
monitor.stop();
|
||||
assertEquals(2, herdr.calls.stream().filter(call -> call.method().equals("agent.list")).count());
|
||||
}
|
||||
|
||||
@Test void faultTransitionLogsOnlyOnceUntilItChanges() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
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);
|
||||
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
|
||||
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
|
||||
monitor.stop();
|
||||
assertEquals(1, appender.list.stream().filter(event -> event.getFormattedMessage()
|
||||
.contains("member=term_a state=TURN_BOUNDARY_LOST")).count());
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user