From 275ac0d251f0c6576ed14ec1becd60204b78e8f5 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 12 Sep 2026 19:26:37 +0700 Subject: [PATCH 1/2] fleetd #562: surface loop health --- CLAUDE.md | 2 +- .../src/main/java/dev/ltms/fleet/Fleetd.java | 6 +- .../java/dev/ltms/fleet/mcp/FleetMcp.java | 98 ++++++++++++++----- .../java/dev/ltms/fleet/rest/FleetApp.java | 54 +++++++--- .../java/dev/ltms/fleet/mcp/FleetMcpTest.java | 31 ++++++ .../dev/ltms/fleet/rest/FleetAppTest.java | 24 +++++ 6 files changed, 179 insertions(+), 36 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index 6eae178..e5f81e0 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -144,7 +144,7 @@ the merge — and merging on a reviewer's word is delegating it by proxy. | Confirm your own role | `fleet_whoami` | | See backends available | `fleet_profiles` | | Start a member | `fleet_spawn{role?, profile?, cwd?, worktree?, ticket?, sessionName?, resumeSessionId?}` → `sessionId` + `paneId` | -| See the fleet | `fleet_list` → `leads` (your peers) + `members` (each carries `agentSessionId` when its backend knows one) · one peer's state: `fleet_status{sessionId}` | +| See the fleet | `fleet_list` → `leads` (your peers) + `members` (each carries `agentSessionId` when its backend knows one) + `loopHealth` (`RUNNING`, `STALLED`, or `STOPPED` for `statusPoller` and `sessionReaper`) · one peer's state: `fleet_status{sessionId}` | | Delegate (blocking) | `fleet_send{sessionId, content}` | | Delegate (long task) | `fleet_send{sessionId, content, wait:false}` → ticket → `fleet_poll{ticket}` | | Answer a member's `fleet_ask` | `fleet_send{turnId, content}` — **not** `sessionId` | diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index 997c12f..90cdb20 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -21,6 +21,7 @@ import dev.ltms.fleet.inject.ExhaustedPatternLookup; import dev.ltms.fleet.inject.ExhaustionSink; import dev.ltms.fleet.inject.LiveExhaustedPatterns; import dev.ltms.fleet.inject.Injector; +import dev.ltms.fleet.inject.LoopWatchdog; import dev.ltms.fleet.inject.StatusPoller; import dev.ltms.fleet.inject.TurnListener; import dev.ltms.fleet.inject.MemberPresence; @@ -665,10 +666,13 @@ public final class Fleetd { return configured == null ? null : configured.effectiveCredentialId(); }, outagePolicy); + FleetMcp.LoopHealthSource loopHealth = new FleetMcp.LoopHealthSource(poller::health, + () -> reaper == null ? LoopWatchdog.State.STOPPED : reaper.health()); FleetMcp mcp = new FleetMcp(messages, workers, sessions, identity, presence, primaryRegistry, callers, FleetMcp.AuthorizationMode.ENFORCED, metrics, capacitySource(config, cfg, profile -> liveCountRef.get().apply(profile)), healthCoverageSource(config), + loopHealth, quarantineSource, leadMailbox, outageSource, @@ -760,7 +764,7 @@ public final class Fleetd { Javalin app = new FleetApp(herdr, memberHerdr, workers, sessions, messages, presence, mcp.servlet(), callers, metrics, deliverable, () -> MemberCredentialPolicyView.of(config.get().memberCredentials()), - quarantineSource, outageSource).build(); + quarantineSource, outageSource, loopHealth).build(); app.start(cfg.bind().host(), cfg.bind().port()); log.info("fleetd listening on {}:{}, herdr socket {}", cfg.bind().host(), cfg.bind().port(), socket); diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index 16af1ef..538b745 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -20,6 +20,7 @@ import dev.ltms.fleet.msg.Rendezvous; import dev.ltms.fleet.peer.PeerUnreachableException; import dev.ltms.fleet.placement.BackendOutagePolicy; import dev.ltms.fleet.placement.BackendQuarantine; +import dev.ltms.fleet.inject.LoopWatchdog; import dev.ltms.fleet.placement.PlacementException; import dev.ltms.fleet.session.SessionManager; import dev.ltms.fleet.session.MemberSession; @@ -108,6 +109,7 @@ public final class FleetMcp { private final Metrics metrics; // CB-502: null → auth failures not counted private final CapacitySource capacity; private final HealthCoverageSource healthCoverage; + private final LoopHealthSource loopHealth; private final QuarantineSource quarantine; /** fleetd #201 Unit 5: SEPARATE from {@link #quarantine} — see {@link OutageSource}'s doc. */ private final OutageSource outage; @@ -136,6 +138,15 @@ public final class FleetMcp { /** Coverage is supplied by the health wiring, not inferred from a missing dependency. */ public record HealthCoverageSource(Supplier value) { } + /** Progress states for fleetd's singleton background loops, read by {@code fleet_list} and {@code /healthz}. */ + public record LoopHealthSource(Supplier statusPoller, + Supplier sessionReaper) { + /** Inert source for callers that do not wire the background loops. */ + public static LoopHealthSource none() { + return new LoopHealthSource(() -> LoopWatchdog.State.STOPPED, () -> LoopWatchdog.State.STOPPED); + } + } + /** * CB-578 stage B quarantine facts used by {@code fleet_profiles}: a profile → credential id * lookup, plus the shared {@link BackendQuarantine} to read remaining cooldowns off. @@ -328,11 +339,22 @@ public final class FleetMcp { * instead of throwing. See {@link #handover}. */ public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions, - ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry, - CallerResolver callers, AuthorizationMode authorizationMode, Metrics metrics, - CapacitySource capacity, HealthCoverageSource healthCoverage, - QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage, - LeadSeatSource leadSeats, List peers, LeadRollover leadRollover) { + ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry, + CallerResolver callers, AuthorizationMode authorizationMode, Metrics metrics, + CapacitySource capacity, HealthCoverageSource healthCoverage, + QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage, + LeadSeatSource leadSeats, List peers, LeadRollover leadRollover) { + this(messages, workers, sessions, identity, presence, primaryRegistry, callers, authorizationMode, metrics, + capacity, healthCoverage, LoopHealthSource.none(), quarantine, leadChannel, outage, leadSeats, + peers, leadRollover); + } + + public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions, + ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry, + CallerResolver callers, AuthorizationMode authorizationMode, Metrics metrics, + CapacitySource capacity, HealthCoverageSource healthCoverage, LoopHealthSource loopHealth, + QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage, + LeadSeatSource leadSeats, List peers, LeadRollover leadRollover) { Objects.requireNonNull(callers, "callers"); this.authorizationEnforced = Objects.requireNonNull(authorizationMode, "authorizationMode") == AuthorizationMode.ENFORCED; @@ -343,6 +365,7 @@ public final class FleetMcp { this.outage = Objects.requireNonNull(outage, "outage"); this.leadSeats = Objects.requireNonNull(leadSeats, "leadSeats"); this.healthCoverage = healthCoverage; + this.loopHealth = Objects.requireNonNull(loopHealth, "loopHealth"); this.leadRollover = leadRollover; McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get(); this.transport = HttpServletStreamableServerTransportProvider.builder() @@ -474,7 +497,7 @@ public final class FleetMcp { (exchange, _) -> { McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_list", Map.of()), null); if (denied != null) return denied; - return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage, + return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage, leadSeats, callers.leads(), callerTerminal(exchange), new CoordinationSource(leadChannel, peers), @@ -1546,25 +1569,42 @@ public final class FleetMcp { * @param selfTerm the calling pane's terminal id, or blank for a caller with no pane */ static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, - Map leads, String selfTerm) { + Map leads, String selfTerm) { return listFleet(workers, sessions, null, CapacitySource.none(), new HealthCoverageSource(() -> "off"), - QuarantineSource.none(), leads, selfTerm); + LoopHealthSource.none(), QuarantineSource.none(), leads, selfTerm); + } + + static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages, + CapacitySource capacity, HealthCoverageSource healthCoverage, + QuarantineSource quarantine, Map leads, String selfTerm) { + return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, leads, selfTerm, + CoordinationSource.none()); } static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages, CapacitySource capacity, HealthCoverageSource healthCoverage, - QuarantineSource quarantine, Map leads, String selfTerm) { - return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, leads, selfTerm, + LoopHealthSource loopHealth, QuarantineSource quarantine, + Map leads, String selfTerm) { + return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, leads, selfTerm, CoordinationSource.none()); } + static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages, + CapacitySource capacity, HealthCoverageSource healthCoverage, + LoopHealthSource loopHealth, QuarantineSource quarantine, + Map leads, String selfTerm, + CoordinationSource coordination) { + return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, + OutageSource.none(), LeadSeatSource.none(), leads, selfTerm, coordination, false); + } + /** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */ static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages, CapacitySource capacity, HealthCoverageSource healthCoverage, QuarantineSource quarantine, OutageSource outage, Map leads, String selfTerm) { - return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage, - LeadSeatSource.none(), leads, selfTerm, CoordinationSource.none()); + return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage, + LeadSeatSource.none(), leads, selfTerm, CoordinationSource.none(), false); } /** @@ -1580,8 +1620,8 @@ public final class FleetMcp { CapacitySource capacity, HealthCoverageSource healthCoverage, QuarantineSource quarantine, Map leads, String selfTerm, CoordinationSource coordination) { - return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, OutageSource.none(), - LeadSeatSource.none(), leads, selfTerm, coordination); + return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, OutageSource.none(), + LeadSeatSource.none(), leads, selfTerm, coordination, false); } /** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */ @@ -1589,8 +1629,8 @@ public final class FleetMcp { CapacitySource capacity, HealthCoverageSource healthCoverage, QuarantineSource quarantine, OutageSource outage, Map leads, String selfTerm, CoordinationSource coordination) { - return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage, - LeadSeatSource.none(), leads, selfTerm, coordination); + return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage, + LeadSeatSource.none(), leads, selfTerm, coordination, false); } /** @@ -1611,7 +1651,7 @@ public final class FleetMcp { QuarantineSource quarantine, OutageSource outage, LeadSeatSource leadSeats, Map leads, String selfTerm, CoordinationSource coordination) { - return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage, + return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage, leadSeats, leads, selfTerm, coordination, false); } @@ -1631,8 +1671,18 @@ public final class FleetMcp { * an explicit {@code true} */ static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages, - CapacitySource capacity, HealthCoverageSource healthCoverage, - QuarantineSource quarantine, OutageSource outage, + CapacitySource capacity, HealthCoverageSource healthCoverage, + QuarantineSource quarantine, OutageSource outage, + LeadSeatSource leadSeats, Map leads, String selfTerm, + CoordinationSource coordination, boolean callerIsPrimary) { + return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, + outage, leadSeats, leads, selfTerm, coordination, callerIsPrimary); + } + + static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages, + CapacitySource capacity, HealthCoverageSource healthCoverage, + LoopHealthSource loopHealth, + QuarantineSource quarantine, OutageSource outage, LeadSeatSource leadSeats, Map leads, String selfTerm, CoordinationSource coordination, boolean callerIsPrimary) { try { @@ -1656,6 +1706,9 @@ public final class FleetMcp { Map result = new LinkedHashMap<>(); result.put("leads", leadRows); result.put("members", out); result.put("healthCoverage", healthCoverage.value().get()); + result.put("loopHealth", Map.of( + "statusPoller", loopHealth.statusPoller().get().name(), + "sessionReaper", loopHealth.sessionReaper().get().name())); // fleetd #439: coordinator/coordinatorView is lead-to-lead coordination state and must // never reach a worker or an architect -- gate BEFORE assembling it, not after, so the // key is absent rather than present-and-empty. @@ -2135,9 +2188,10 @@ public final class FleetMcp { + "cannot reliably re-identify: some backends (e.g. opencode) resolve it from the " + "member's working directory, which only uniquely identifies a member when it " + "was spawned into its own fleetd-provisioned worktree (worktree:true/); a " - + "member spawned without one shares its directory with others and never reports " - + "an id, however long it runs (fleetd #249). An empty 'members' " - + "means no members are spawned; it says nothing about peers. When capacity " + + "member spawned without one shares its directory with others and never reports " + + "an id, however long it runs (fleetd #249). An empty 'members' " + + "means no members are spawned; it says nothing about peers. 'loopHealth' reports " + + "the RUNNING, STALLED, or STOPPED state of statusPoller and sessionReaper. When capacity " + "facts are configured, a 'capacity' row per profile reports 'free' — the " + "slots a fresh fleet_spawn on that profile will actually be granted right " + "now (max(0, maxLoad - live)), the same check the spawn gate itself runs. A " diff --git a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java index 87acb9f..66aa167 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java @@ -94,6 +94,7 @@ public final class FleetApp { // .none() (the honest "feature not wired" view) for every constructor that does not pass one. private final FleetMcp.QuarantineSource quarantine; private final FleetMcp.OutageSource outage; + private final FleetMcp.LoopHealthSource loopHealth; private final ObjectMapper mapper = new ObjectMapper(); /** @@ -155,7 +156,8 @@ public final class FleetApp { HttpServlet mcpServlet, CallerResolver auth, Metrics metrics, Predicate deliverable, Supplier memberCredentials) { this(herdr, memberHerdr, workers, sessions, messages, presence, mcpServlet, auth, metrics, - deliverable, memberCredentials, FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none()); + deliverable, memberCredentials, FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), + FleetMcp.LoopHealthSource.none()); } /** @@ -169,8 +171,18 @@ public final class FleetApp { public FleetApp(HerdrClient herdr, HerdrClient memberHerdr, PeerLauncher workers, SessionManager sessions, MessageService messages, MemberPresence presence, HttpServlet mcpServlet, CallerResolver auth, Metrics metrics, - Predicate deliverable, Supplier memberCredentials, - FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage) { + Predicate deliverable, Supplier memberCredentials, + FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage) { + this(herdr, memberHerdr, workers, sessions, messages, presence, mcpServlet, auth, metrics, deliverable, + memberCredentials, quarantine, outage, FleetMcp.LoopHealthSource.none()); + } + + public FleetApp(HerdrClient herdr, HerdrClient memberHerdr, PeerLauncher workers, SessionManager sessions, + MessageService messages, MemberPresence presence, + HttpServlet mcpServlet, CallerResolver auth, Metrics metrics, + Predicate deliverable, Supplier memberCredentials, + FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage, + FleetMcp.LoopHealthSource loopHealth) { this.herdr = herdr; this.memberHerdr = memberHerdr != null ? memberHerdr : herdr; this.workers = workers; @@ -183,6 +195,7 @@ public final class FleetApp { this.memberCredentials = memberCredentials != null ? memberCredentials : MemberCredentialPolicyView::absent; this.quarantine = quarantine != null ? quarantine : FleetMcp.QuarantineSource.none(); this.outage = outage != null ? outage : FleetMcp.OutageSource.none(); + this.loopHealth = loopHealth != null ? loopHealth : FleetMcp.LoopHealthSource.none(); } /** Wire routes onto a fresh, unstarted Javalin instance. Caller starts it. */ @@ -291,15 +304,19 @@ public final class FleetApp { * spotted by comparing two numbers by eye. */ private void healthz(Context ctx) { + HealthzResponse response = healthzResponse(herdr, memberHerdr, loopHealth); + ctx.status(response.status()).json(response.body()); + } + + record HealthzResponse(int status, Map body) { } + + static HealthzResponse healthzResponse(HerdrClient herdr, HerdrClient memberHerdr, + FleetMcp.LoopHealthSource loopHealth) { JsonNode pong; try { pong = herdr.call("ping"); } catch (HerdrException e) { - ctx.status(503).json(Map.of( - "status", "degraded", - "herdr", "unreachable", - "detail", e.getMessage())); - return; + return degradedResponse("unreachable", e.getMessage(), loopHealth); } Map body = new LinkedHashMap<>(); body.put("status", "ok"); @@ -311,11 +328,11 @@ public final class FleetApp { try { memberPong = memberHerdr.call("ping"); } catch (HerdrException e) { - ctx.status(503).json(Map.of( + return new HealthzResponse(503, Map.of( "status", "degraded", "herdr", "member unreachable", - "detail", e.getMessage())); - return; + "detail", e.getMessage(), + "loopHealth", loopHealthView(loopHealth))); } int leadProtocol = pong.path("protocol").asInt(); int memberProtocol = memberPong.path("protocol").asInt(); @@ -326,7 +343,20 @@ public final class FleetApp { body.put("protocolMismatch", true); } } - ctx.status(200).json(body); + body.put("loopHealth", loopHealthView(loopHealth)); + return new HealthzResponse(200, body); + } + + private static Map loopHealthView(FleetMcp.LoopHealthSource loopHealth) { + return Map.of("statusPoller", loopHealth.statusPoller().get().name(), + "sessionReaper", loopHealth.sessionReaper().get().name()); + } + + private static HealthzResponse degradedResponse(String herdr, String detail, + FleetMcp.LoopHealthSource loopHealth) { + return new HealthzResponse(503, Map.of( + "status", "degraded", "herdr", herdr, "detail", detail, + "loopHealth", loopHealthView(loopHealth))); } /** diff --git a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java index 0664118..ee60653 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java @@ -8,6 +8,7 @@ import dev.ltms.fleet.config.FleetConfig; import dev.ltms.fleet.guard.SubscriptionGuard; import dev.ltms.fleet.herdr.AgentControl; import dev.ltms.fleet.herdr.FakeHerdr; +import dev.ltms.fleet.inject.LoopWatchdog; import dev.ltms.fleet.herdr.PaneLocator; import dev.ltms.fleet.herdr.WorkspaceControl; import dev.ltms.fleet.inject.Injector; @@ -1053,6 +1054,36 @@ class FleetMcpTest { assertFalse(out.contains("quarantinedForSeconds"), out); } + @Test + void loopHealthReportsStalledStatusPoller() { + String out = loopHealth(LoopWatchdog.State.STALLED, LoopWatchdog.State.RUNNING); + assertTrue(out.contains("\"statusPoller\":\"STALLED\""), + "fleet_list must report a stalled StatusPoller: " + out); + } + + @Test + void loopHealthReportsStoppedSessionReaperAsStopped() { + String out = loopHealth(LoopWatchdog.State.RUNNING, LoopWatchdog.State.STOPPED); + assertTrue(out.contains("\"sessionReaper\":\"STOPPED\""), + "fleet_list must report a deliberately stopped SessionReaper as STOPPED, not an alarm: " + out); + } + + @Test + void loopHealthReportsRunningStatusPoller() { + String out = loopHealth(LoopWatchdog.State.RUNNING, LoopWatchdog.State.STOPPED); + assertTrue(out.contains("\"statusPoller\":\"RUNNING\""), + "fleet_list must report a running StatusPoller: " + out); + } + + private static String loopHealth(LoopWatchdog.State statusPoller, LoopWatchdog.State sessionReaper) { + FakeHerdr herdr = new FakeHerdr(); + return textOf(FleetMcp.listFleet(workerService(herdr, "http://gx00.gw:8000", Set.of("gx00.gw")), + new SessionManager(workerService(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"))), null, + FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"), + new FleetMcp.LoopHealthSource(() -> statusPoller, () -> sessionReaper), + FleetMcp.QuarantineSource.none(), Map.of(), "")); + } + @Test void capacityIncludesConfiguredProfileWithoutMembers() { FakeHerdr h = new FakeHerdr(); diff --git a/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppTest.java b/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppTest.java index 7d014c1..5bb7a54 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppTest.java @@ -8,6 +8,7 @@ import dev.ltms.fleet.herdr.AgentControl; import dev.ltms.fleet.herdr.FakeHerdr; import dev.ltms.fleet.herdr.WorkspaceControl; import dev.ltms.fleet.inject.Injector; +import dev.ltms.fleet.inject.LoopWatchdog; import dev.ltms.fleet.inject.StatusPoller; import dev.ltms.fleet.inject.MemberPresence; import dev.ltms.fleet.msg.InMemoryReplyInbox; @@ -165,6 +166,29 @@ class FleetAppTest { assertEquals("degraded", mapper.readTree(res.body()).get("status").asText()); } + @Test + void healthzKeepsOkStatusAndReportsLoopHealthInItsBody() { + FleetMcp.LoopHealthSource loops = new FleetMcp.LoopHealthSource( + () -> LoopWatchdog.State.RUNNING, () -> LoopWatchdog.State.STOPPED); + + FleetApp.HealthzResponse ok = FleetApp.healthzResponse(new FakeHerdr(), new FakeHerdr(), loops); + assertEquals(200, ok.status(), "a healthy herdr must keep /healthz at 200 regardless of loop states"); + assertEquals(Map.of("statusPoller", "RUNNING", "sessionReaper", "STOPPED"), ok.body().get("loopHealth"), + "the /healthz body must report each loop state without making STOPPED an alarm"); + } + + @Test + void healthzKeepsDegradedStatusAndReportsLoopHealthInItsBody() { + FleetMcp.LoopHealthSource loops = new FleetMcp.LoopHealthSource( + () -> LoopWatchdog.State.RUNNING, () -> LoopWatchdog.State.STOPPED); + FleetApp.HealthzResponse degraded = FleetApp.healthzResponse(new FakeHerdr().healthy(false), + new FakeHerdr(), loops); + assertEquals(503, degraded.status(), + "an unreachable herdr must keep /healthz at 503 regardless of loop states"); + assertEquals(Map.of("statusPoller", "RUNNING", "sessionReaper", "STOPPED"), + degraded.body().get("loopHealth"), "the degraded /healthz body must retain loop states"); + } + @Test void sessionsMapsWorkspaceList() throws Exception { int port = startHealthy(); From 1513d4f260f6120811a589539d87f2ce7b5afc00 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 12 Sep 2026 20:16:23 +0700 Subject: [PATCH 2/2] fleetd #562 follow-up: extract loopHealthSource factory, pin its wiring PR #579's inline `new FleetMcp.LoopHealthSource(poller::health, ...)` in Fleetd.main had nothing a test could call directly. Measured: replacing poller::health with a constant () -> RUNNING compiled clean and left all 1771 tests green (see issue #562 comment "HOLD on PR #579"). Extracts the inline construction to a package-private Fleetd.loopHealthSource factory, the same style as the sibling capacitySource/healthCoverageSource factories, and adds FleetdLoopHealthSourceWiringTest with three separate assertions: the statusPoller half, the sessionReaper half, and the reaper == null branch (still STOPPED). --- .../src/main/java/dev/ltms/fleet/Fleetd.java | 28 ++- .../FleetdLoopHealthSourceWiringTest.java | 174 ++++++++++++++++++ 2 files changed, 200 insertions(+), 2 deletions(-) create mode 100644 fleetd/src/test/java/dev/ltms/fleet/FleetdLoopHealthSourceWiringTest.java diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index 90cdb20..a93637f 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -666,8 +666,7 @@ public final class Fleetd { return configured == null ? null : configured.effectiveCredentialId(); }, outagePolicy); - FleetMcp.LoopHealthSource loopHealth = new FleetMcp.LoopHealthSource(poller::health, - () -> reaper == null ? LoopWatchdog.State.STOPPED : reaper.health()); + FleetMcp.LoopHealthSource loopHealth = loopHealthSource(poller, reaper); FleetMcp mcp = new FleetMcp(messages, workers, sessions, identity, presence, primaryRegistry, callers, FleetMcp.AuthorizationMode.ENFORCED, metrics, capacitySource(config, cfg, profile -> liveCountRef.get().apply(profile)), @@ -1042,6 +1041,31 @@ public final class Fleetd { }); } + /** + * fleetd #562 follow-up: package-private factory for {@code fleet_list}'s and {@code + * /healthz}'s {@code loopHealth} source, extracted out of {@code main} for the same reason + * {@link #capacitySource} and {@link #healthCoverageSource} were. Before this ticket the + * {@link FleetMcp.LoopHealthSource} was built inline with a bare {@code new}, so there was + * nothing a test could call directly — measured: replacing {@code poller::health} with a + * constant {@code () -> LoopWatchdog.State.RUNNING} at the call site compiled clean and left + * the full suite green, meaning the daemon could report the {@link StatusPoller} as always + * {@code RUNNING} even while it was actually stalled. That is a false negative on the exact + * signal this ticket exists to surface, and is the mirror of a false positive muting a real + * monitoring component — worse, because there is no noise for anyone to notice and then + * silence. {@link FleetdLoopHealthSourceWiringTest} calls this factory directly and pins both + * halves separately, plus the {@code reaper == null} branch below. + * + *

{@code reaper} may be {@code null} — a {@link SessionReaper} is only constructed when + * {@code lifecycle.idleTtlSeconds} is configured (see the {@code reaper} local above) — and + * this factory preserves the existing behaviour of reporting {@link LoopWatchdog.State#STOPPED} + * in that case, rather than a {@code NullPointerException} on the first {@code fleet_list} or + * {@code /healthz} call. + */ + static FleetMcp.LoopHealthSource loopHealthSource(StatusPoller poller, SessionReaper reaper) { + return new FleetMcp.LoopHealthSource(poller::health, + () -> reaper == null ? LoopWatchdog.State.STOPPED : reaper.health()); + } + /** * fleetd #248: package-private factory for the member worktree/branch lookup {@link * CompletionResolver} uses to name a fallback report's worktree and branch (fleetd#241). diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdLoopHealthSourceWiringTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdLoopHealthSourceWiringTest.java new file mode 100644 index 0000000..32fa2da --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdLoopHealthSourceWiringTest.java @@ -0,0 +1,174 @@ +package dev.ltms.fleet; + +import dev.ltms.fleet.herdr.AgentControl; +import dev.ltms.fleet.herdr.FakeHerdr; +import dev.ltms.fleet.inject.Injector; +import dev.ltms.fleet.inject.LoopWatchdog; +import dev.ltms.fleet.inject.StatusPoller; +import dev.ltms.fleet.mcp.FleetMcp; +import dev.ltms.fleet.peer.Capability; +import dev.ltms.fleet.peer.PeerHandle; +import dev.ltms.fleet.peer.PeerLauncher; +import dev.ltms.fleet.peer.SpawnRequest; +import dev.ltms.fleet.placement.PlacementDecision; +import dev.ltms.fleet.session.SessionManager; +import dev.ltms.fleet.session.SessionReaper; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import java.util.List; +import java.util.Set; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +/** + * fleetd #562 follow-up (issue comment "HOLD on PR #579"): {@code Fleetd.main}'s {@code loopHealth} + * local used to be a bare {@code new FleetMcp.LoopHealthSource(poller::health, ...)} built inline, + * with nothing a test could call directly. Measured on that shape: replacing {@code + * poller::health} with a constant {@code () -> LoopWatchdog.State.RUNNING} at the call site + * compiled with 0 errors and left all 1771 existing tests green — the daemon could be changed to + * always report the {@link StatusPoller} as {@code RUNNING}, so the watchdog could never fire and + * a stalled poller would be invisible, while every test stayed green. That is exactly the false + * negative this ticket exists to prevent. + * + *

The five tests PR #579 added ({@code FleetMcpTest}, {@code FleetAppTest}) all build their own + * {@link FleetMcp.LoopHealthSource} directly with fixed lambdas — they prove the seam ({@code + * LoopHealthSource} reports what it is given) and nothing about what {@code Fleetd.main} actually + * gives it. This is the same hand-built-vs-config-wired shape as fleetd #561/#248/#426. + * + *

The fix extracts the inline {@code new} into {@link Fleetd#loopHealthSource}, a package-private + * factory in the same style as {@link Fleetd#capacitySource} and {@link Fleetd#healthCoverageSource} + * — which is exactly what makes it directly callable here. This test calls that factory with real + * {@link StatusPoller}/{@link SessionReaper} instances (never started, so no herdr or git I/O + * happens) and pins each half separately, plus the {@code reaper == null} branch: one invariant + * wired at three places needs three assertions, not one combined check whose non-zero total could + * hide a gap at any single place. + */ +class FleetdLoopHealthSourceWiringTest { + + @Test + @DisplayName("the statusPoller half reports the real poller's health, not a hardcoded state") + void statusPollerHalfReflectsThePollersRealHealth() { + // Stopped without ever being started — stop() still marks the watchdog STOPPED. A poller + // that has never reported RUNNING is the discriminating case: if Fleetd.loopHealthSource + // ever hardcoded RUNNING (the exact mutation this test exists to catch), this would fail. + StatusPoller stoppedPoller = freshPoller(); + stoppedPoller.stop(); + SessionReaper unusedReaper = freshReaper(); // present only to satisfy the signature + + FleetMcp.LoopHealthSource source = Fleetd.loopHealthSource(stoppedPoller, unusedReaper); + + assertEquals(LoopWatchdog.State.STOPPED, source.statusPoller().get(), + "the statusPoller supplier must delegate to the real poller's health() — " + + "replacing poller::health with a constant () -> RUNNING at the " + + "Fleetd.loopHealthSource call site must fail this assertion"); + } + + @Test + @DisplayName("the sessionReaper half reports the real reaper's health, not a hardcoded state") + void sessionReaperHalfReflectsTheReapersRealHealth() { + StatusPoller unusedPoller = freshPoller(); // present only to satisfy the signature + SessionReaper stoppedReaper = freshReaper(); + stoppedReaper.stop(); + + FleetMcp.LoopHealthSource source = Fleetd.loopHealthSource(unusedPoller, stoppedReaper); + + assertEquals(LoopWatchdog.State.STOPPED, source.sessionReaper().get(), + "the sessionReaper supplier must delegate to the real reaper's health() — " + + "replacing reaper.health() with a constant at the " + + "Fleetd.loopHealthSource call site must fail this assertion"); + } + + @Test + @DisplayName("a null reaper (idle ttl not configured) still reports STOPPED, not a crash") + void nullReaperStillReportsStopped() { + // SessionReaper is only constructed when lifecycle.idleTtlSeconds is configured (see the + // `reaper` local in Fleetd.main) — a real deployment routinely passes null here. That null + // check is real behaviour, not a simplification to delete: it must keep reporting STOPPED + // rather than throwing a NullPointerException on the first fleet_list/healthz call. + StatusPoller runningPoller = freshPoller(); + + FleetMcp.LoopHealthSource source = Fleetd.loopHealthSource(runningPoller, null); + + assertEquals(LoopWatchdog.State.STOPPED, source.sessionReaper().get(), + "reaper == null must still report STOPPED, exactly like an intentionally-stopped " + + "reaper would — do not delete this null check to simplify the wiring"); + } + + /** Never started, so no herdr call is ever made; freshly constructed reports RUNNING. */ + private static StatusPoller freshPoller() { + AgentControl agents = new AgentControl(new FakeHerdr()); + return new StatusPoller(agents, new Injector(agents), 1000); + } + + /** Never started, so no git/session I/O is ever made; freshly constructed reports RUNNING. */ + private static SessionReaper freshReaper() { + return new SessionReaper(new SessionManager(new NeverSpawnsLauncher()), 60, 1000); + } + + /** + * Same minimal shape as {@code FleetdBackendErrorSinkTest.NeverSpawnsLauncher} — every method + * throws or returns an empty/no-op value, since a {@link SessionReaper} that is only ever + * constructed and then stopped (never started) never calls any of them. + */ + private static final class NeverSpawnsLauncher implements PeerLauncher { + @Override + public Set capabilities() { + return Set.of(); + } + + @Override + public Set capabilitiesFor(String profileName) { + return Set.of(); + } + + @Override + public PeerHandle spawn(SpawnRequest req) { + throw new UnsupportedOperationException("not reachable — this test never acquires a session"); + } + + @Override + public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) { + throw new UnsupportedOperationException("not reachable — this test never acquires a session"); + } + + @Override + public Set profiles() { + return Set.of(); + } + + @Override + public String defaultProfile() { + return null; + } + + @Override + public String effectiveCwd(SpawnRequest req) { + throw new UnsupportedOperationException("not reachable — this test never acquires a session"); + } + + @Override + public List parityOverlay(String profileName) { + return List.of(); + } + + @Override + public List list() { + return List.of(); + } + + @Override + public int reapOrphanWorkers() { + return 0; + } + + @Override + public void stop(String id) { + } + + @Override + public boolean clearContext(String id) { + return false; + } + } +}