CB-573: add dormant fleet health monitor
This commit is contained in:
@@ -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,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());
|
||||
|
||||
Reference in New Issue
Block a user