Merge CB-573: dormant fleet health monitor (M4 unit 1)
CI / build (push) Successful in 49s
CI / contract (push) Successful in 1m14s

An opt-in whole-fleet observer, separate from the 250ms delivery poller.
One AgentControl.list and one roster snapshot per tick, joined and fed to
the FleetHealth classifier, because a fault is a disagreement between the
two views at the same instant. Absent a health: block nothing is built
and no herdr call is made.

Adds bridge_list healthCoverage: off, detection-only, or full. Detection
is deliberately separate from notification, so a single-lead setup with
no webhook still gets detection and is told its coverage is partial
rather than being refused.

Two review fixes worth naming. tick() rescheduled itself as its last
statement with no try/catch, and a ScheduledExecutorService does not
re-run a task that threw — so the first agents.list failure would have
stopped health permanently and silently, which is exactly when the
control link is down. It now catches Throwable and reschedules in a
finally. And the snapshot fields this unit cannot supply are the named
constant NOT_YET_OBSERVED rather than bare false literals, because false
means no fault to this classifier.
This commit is contained in:
Dai Ha
2026-08-15 06:19:15 +02:00
8 changed files with 264 additions and 12 deletions
+10
View File
@@ -85,6 +85,16 @@ bind:
# backoffMs: 60000
# quietNudgeCap: 3
# Fleet health detection is dormant unless enabled. It reads one whole-fleet agent list per tick.
# It can run without a webhook; bridge_list then reports healthCoverage: detection-only.
# health:
# enabled: true
# intervalSeconds: 30 # minimum 15
# workingSuspectAfterSeconds: 600 # minimum 300
# paneProbeIntervalSeconds: 60 # minimum 60
# notifications:
# mode: disabled # disabled (default) or webhook
# herdr Unix socket. Omit to use the client default
# (${HERDR_SOCKET_PATH:-~/.config/herdr/herdr.sock}).
herdrSocket: ~/.config/herdr/herdr.sock
@@ -23,6 +23,7 @@ import dev.ltms.bridged.mcp.BridgeMcp;
import dev.ltms.bridged.mcp.ConnectionIdentity;
import dev.ltms.bridged.metrics.BridgedMetrics;
import dev.ltms.bridged.metrics.Metrics;
import dev.ltms.bridged.health.FleetHealthMonitor;
import dev.ltms.bridged.mcp.PrimaryRegistry;
import dev.ltms.bridged.mcp.LsofPeerPidLookup;
import dev.ltms.bridged.mcp.LsofProcessCwdLookup;
@@ -354,6 +355,26 @@ public final class Bridged {
MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox,
pushLoop, metrics);
// Health is a slow whole-fleet observer. Keep it separate from the 250ms delivery poller.
final FleetHealthMonitor healthMonitor;
var healthScheduler = Executors.newSingleThreadScheduledExecutor(r ->
Thread.ofVirtual().name("bridge-health-").unstarted(r));
if (cfg.health() != null && cfg.health().isEnabled()) {
healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler,
System::nanoTime, cfg.health().intervalOrDefault());
String coverage = FleetHealthMonitor.coverage(true,
cfg.health().notifications() != null && cfg.health().notifications().configured());
if ("detection-only".equals(coverage)) {
log.warn("fleet health: {} (no notification sink configured)", coverage);
} else {
log.info("fleet health: {}", coverage);
}
healthMonitor.start();
} else {
healthMonitor = null;
healthScheduler.shutdownNow();
}
// CB-520: the reply inbox only consumes for agents this gateway owns. own on acquire,
// release on teardown. Do this before CB-516 so the inbox is owned before any reply can land.
sessions.onAcquire(replyInbox::own);
@@ -393,7 +414,12 @@ public final class Bridged {
profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.maxLoad();
}, () -> config.get().profiles().keySet(), System::nanoTime));
}, () -> config.get().profiles().keySet(), System::nanoTime),
new BridgeMcp.HealthCoverageSource(() -> {
var health = config.get().health();
return FleetHealthMonitor.coverage(health != null && health.isEnabled(),
health != null && health.notifications() != null && health.notifications().configured());
}));
// CB-559: opt-in config reload. With no `configReload:` block nothing is constructed, so an
// upgraded daemon behaves exactly as before — the file is read once at boot and never again.
@@ -414,6 +440,7 @@ public final class Bridged {
messages.close();
pushLoop.close();
if (heartbeat != null) heartbeat.close(); // CB-551: stop the idle-lead heartbeat scheduler
if (healthMonitor != null) healthMonitor.stop();
if (configWatcher != null) configWatcher.stop(); // CB-559: stop polling the config file
mcp.close();
if (reaper != null) reaper.stop();
@@ -73,6 +73,7 @@ public record BridgedConfig(
Primary primary,
Fleet fleet,
LeadHeartbeat leadHeartbeat,
Health health,
String placement,
Auth auth,
ConfigReload configReload) {
@@ -83,7 +84,16 @@ public record BridgedConfig(
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, String placement, Auth auth) {
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, placement, auth, null);
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, null);
}
/** Back-compat form before the optional {@code health:} block was added. */
public BridgedConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, String placement, Auth auth, ConfigReload configReload) {
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, configReload);
}
/**
@@ -377,6 +387,17 @@ public record BridgedConfig(
boolean clearAfterTurn) {
}
/** Optional fleet detection. A missing block stays dormant. */
@JsonIgnoreProperties(ignoreUnknown = true)
public record Health(Boolean enabled, Integer intervalSeconds, Integer workingSuspectAfterSeconds,
Integer paneProbeIntervalSeconds, Notifications notifications) {
public boolean isEnabled() { return Boolean.TRUE.equals(enabled); }
public int intervalOrDefault() { return Math.max(15, intervalSeconds == null ? 30 : intervalSeconds); }
public record Notifications(String mode) {
public boolean configured() { return "webhook".equalsIgnoreCase(mode); }
}
}
/**
* External AMQP broker for durable, cross-restart reply delivery (CB-307 Stage 2). Its mere
* presence swaps the in-memory {@code ReplyInbox} for the AMQP-backed adapter; absent, bridged
@@ -812,7 +833,7 @@ public record BridgedConfig(
private static final Set<String> KNOWN_TOP_LEVEL_KEYS = Set.of(
"bind", "herdrSocket", "profiles", "guard", "worktreeRoot",
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
"leadHeartbeat", "placement", "auth", "configReload");
"leadHeartbeat", "health", "placement", "auth", "configReload");
/** Load and validate config from {@code path}. */
public static BridgedConfig load(Path path) {
@@ -1082,7 +1103,7 @@ public record BridgedConfig(
// defaults the fields of a block that IS present. Defaulting it here would start watching
// the file for every config that never asked to be watched.
return new BridgedConfig(b, herdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
broker, primary, f, leadHeartbeat, placementOrDefault, a, configReload);
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload);
}
/**
@@ -0,0 +1,107 @@
package dev.ltms.bridged.health;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.session.MemberSession;
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;
import java.util.concurrent.TimeUnit;
import java.util.function.LongSupplier;
import java.util.function.Supplier;
/** Slow whole-fleet evidence collection. It is deliberately separate from the delivery poller. */
public final class FleetHealthMonitor {
private static final Logger log = LoggerFactory.getLogger(FleetHealthMonitor.class);
private final AgentControl agents;
private final Supplier<List<MemberSession>> roster;
private final MessageService messages;
private final ScheduledExecutorService scheduler;
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) {
this.agents = agents;
this.roster = roster;
this.messages = messages;
this.scheduler = scheduler;
this.clock = clock;
this.intervalSeconds = intervalSeconds;
}
/** Pure per-member decision seam. */
static HealthDecision decide(HealthSnapshot snapshot, HealthPrior prior, long nowNanos) {
return FleetHealth.decide(snapshot, prior, nowNanos);
}
public void start() { scheduler.schedule(this::tick, intervalSeconds, TimeUnit.SECONDS); }
public void stop() { scheduler.shutdownNow(); }
// 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.
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);
}
}
}
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) {
return !enabled ? "off" : notificationConfigured ? "full" : "detection-only";
}
}
@@ -77,6 +77,7 @@ public final class BridgeMcp {
private final CallerResolver authz; // CB-501: null → authorization not enforced (legacy)
private final Metrics metrics; // CB-502: null → auth failures not counted
private final CapacitySource capacity;
private final HealthCoverageSource healthCoverage;
/** Capacity facts used by {@code bridge_list}; production must supply the placement live count. */
public record CapacitySource(Function<String, Integer> liveCount, Function<String, Integer> maxLoad,
@@ -86,6 +87,9 @@ public final class BridgeMcp {
boolean available() { return !configuredProfiles.get().isEmpty(); }
}
/** Coverage is supplied by the health wiring, not inferred from a missing dependency. */
public record HealthCoverageSource(Supplier<String> value) { }
/**
* @param callers resolves each call's {@link Principal}; {@code null} disables authorization.
* This surface needs its own enforcement: {@code /mcp} is a raw servlet on
@@ -95,8 +99,9 @@ public final class BridgeMcp {
*/
public BridgeMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
CallerResolver callers, Metrics metrics, CapacitySource capacity) {
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage) {
this.capacity = capacity;
this.healthCoverage = healthCoverage;
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
this.transport = HttpServletStreamableServerTransportProvider.builder()
.jsonMapper(json)
@@ -210,7 +215,7 @@ public final class BridgeMcp {
.toolCall(listTool(), (exchange, _) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied;
return listFleet(workers, sessions, messages, capacity,
return listFleet(workers, sessions, messages, capacity, healthCoverage,
callers == null ? Map.of() : callers.leads(),
callerTerminal(exchange));
})
@@ -724,11 +729,11 @@ public final class BridgeMcp {
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions,
Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, null, CapacitySource.none(), leads, selfTerm);
return listFleet(workers, sessions, null, CapacitySource.none(), new HealthCoverageSource(() -> "off"), leads, selfTerm);
}
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity,
CapacitySource capacity, HealthCoverageSource healthCoverage,
Map<String, String> leads, String selfTerm) {
try {
Map<String, Agent> live = workers.list().stream()
@@ -747,6 +752,7 @@ public final class BridgeMcp {
roster.stream().map(MemberSession::profile).forEach(profiles::add);
Map<String, Object> result = new LinkedHashMap<>();
result.put("leads", leadRows); result.put("members", out);
result.put("healthCoverage", healthCoverage.value().get());
if (capacity.available()) result.put("capacity", profiles.stream()
.map(profile -> capacityView(profile, capacity.liveCount(), capacity.maxLoad(), roster, messages,
capacity.clock().getAsLong())).toList());
@@ -0,0 +1,81 @@
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;
import dev.ltms.bridged.msg.InMemoryReplyInbox;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
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;
import java.util.concurrent.Executors;
import static org.junit.jupiter.api.Assertions.assertEquals;
class FleetHealthMonitorTest {
@Test void oneTickUsesOneFleetListForAnyRosterSize() {
FakeHerdr herdr = new FakeHerdr();
AgentControl agents = new AgentControl(herdr);
BridgedConfig.Profile profile = new BridgedConfig.Profile("test", "http://test:1", null,
null, null, null, null, null, null, null, null, null);
ClaudeCodeLauncher launcher = new ClaudeCodeLauncher(agents, new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("test")), Map.of("test", profile), "test", _ -> "token");
SessionManager sessions = new SessionManager(launcher);
sessions.acquire("test", null, null, null);
sessions.acquire("test", null, null, null);
herdr.calls.clear();
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);
monitor.tick();
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);
}
}
}
@@ -72,7 +72,7 @@ class BridgeMcpAuthzTest {
new PrimaryRegistry(null),
enforce ? CallerResolver.withLeadsAndMembers(identity, false, null,
Map::of, new MemberRegistry(null)) : null,
metrics, BridgeMcp.CapacitySource.none());
metrics, BridgeMcp.CapacitySource.none(), new BridgeMcp.HealthCoverageSource(() -> "off"));
return mcp;
}
@@ -474,7 +474,7 @@ class BridgeMcpTest {
sessions.acquire("ltms-local", null, null, null);
McpSchema.CallToolResult res = BridgeMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
sessions, null, new BridgeMcp.CapacitySource(profile -> 2, profile -> 2,
() -> Set.of("ltms-local"), () -> 0), Map.of(), "");
() -> Set.of("ltms-local"), () -> 0), new BridgeMcp.HealthCoverageSource(() -> "off"), Map.of(), "");
String out = textOf(res);
assertTrue(out.contains("\"maxLoad\":2"), out);
assertTrue(out.contains("\"live\":2"), out);
@@ -487,7 +487,7 @@ class BridgeMcpTest {
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
String out = textOf(BridgeMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
sessions, null, new BridgeMcp.CapacitySource(profile -> 0, profile -> 2,
() -> Set.of("terra"), () -> 0), Map.of(), ""));
() -> Set.of("terra"), () -> 0), new BridgeMcp.HealthCoverageSource(() -> "off"), Map.of(), ""));
assertTrue(out.contains("\"profile\":\"terra\""), out);
assertTrue(out.contains("\"live\":0"), out);
assertTrue(out.contains("\"free\":2"), out);
@@ -499,7 +499,7 @@ class BridgeMcpTest {
FakeHerdr h = new FakeHerdr();
String out = textOf(BridgeMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))), null,
BridgeMcp.CapacitySource.none(), Map.of(), ""));
BridgeMcp.CapacitySource.none(), new BridgeMcp.HealthCoverageSource(() -> "off"), Map.of(), ""));
assertFalse(out.contains("\"capacity\":"), out);
}