From bcee1f9581ddaf070fead5bda949ea20deb93bc5 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 06:10:29 +0200 Subject: [PATCH] CB-573: add dormant fleet health monitor --- bridged/bridged.example.yaml | 10 +++ .../main/java/dev/ltms/bridged/Bridged.java | 29 +++++++- .../ltms/bridged/config/BridgedConfig.java | 27 ++++++- .../bridged/health/FleetHealthMonitor.java | 70 +++++++++++++++++++ .../java/dev/ltms/bridged/mcp/BridgeMcp.java | 14 ++-- .../health/FleetHealthMonitorTest.java | 41 +++++++++++ .../ltms/bridged/mcp/BridgeMcpAuthzTest.java | 2 +- .../dev/ltms/bridged/mcp/BridgeMcpTest.java | 6 +- 8 files changed, 187 insertions(+), 12 deletions(-) create mode 100644 bridged/src/main/java/dev/ltms/bridged/health/FleetHealthMonitor.java create mode 100644 bridged/src/test/java/dev/ltms/bridged/health/FleetHealthMonitorTest.java diff --git a/bridged/bridged.example.yaml b/bridged/bridged.example.yaml index 6b68030..6bd6c0f 100644 --- a/bridged/bridged.example.yaml +++ b/bridged/bridged.example.yaml @@ -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 diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 85d057c..a43cb0e 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -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(); diff --git a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java index a8f7086..d072ff0 100644 --- a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java +++ b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java @@ -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 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 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); } /** diff --git a/bridged/src/main/java/dev/ltms/bridged/health/FleetHealthMonitor.java b/bridged/src/main/java/dev/ltms/bridged/health/FleetHealthMonitor.java new file mode 100644 index 0000000..c0d808c --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/health/FleetHealthMonitor.java @@ -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> roster; + private final MessageService messages; + private final ScheduledExecutorService scheduler; + private final LongSupplier clock; + private final long intervalSeconds; + private final Map priors = new HashMap<>(); + + public FleetHealthMonitor(AgentControl agents, Supplier> 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 agentsNow = agents.list(); // Exactly one list call for this complete observation. + Map 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"; + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java index a480aff..af11447 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java @@ -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 liveCount, Function 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 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 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 leads, String selfTerm) { try { Map live = workers.list().stream() @@ -747,6 +752,7 @@ public final class BridgeMcp { roster.stream().map(MemberSession::profile).forEach(profiles::add); Map 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()); diff --git a/bridged/src/test/java/dev/ltms/bridged/health/FleetHealthMonitorTest.java b/bridged/src/test/java/dev/ltms/bridged/health/FleetHealthMonitorTest.java new file mode 100644 index 0000000..e133ad5 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/health/FleetHealthMonitorTest.java @@ -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()); + } +} diff --git a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpAuthzTest.java b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpAuthzTest.java index fb2e067..7b1af8b 100644 --- a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpAuthzTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpAuthzTest.java @@ -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; } diff --git a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java index 48d9b77..43ba54d 100644 --- a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java @@ -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); }