CB-573: add dormant fleet health monitor
This commit is contained in:
@@ -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;
|
||||
@@ -348,6 +349,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);
|
||||
@@ -387,7 +408,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.
|
||||
@@ -408,6 +434,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,70 @@
|
||||
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.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<>();
|
||||
|
||||
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() {
|
||||
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());
|
||||
}
|
||||
scheduler.schedule(this::tick, intervalSeconds, TimeUnit.SECONDS);
|
||||
}
|
||||
|
||||
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,41 @@
|
||||
package dev.ltms.bridged.health;
|
||||
|
||||
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 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());
|
||||
}
|
||||
}
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user