Merge worker/562-loop-health-wiring-test-99611c-5
This commit is contained in:
@@ -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` |
|
||||
|
||||
@@ -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,12 @@ public final class Fleetd {
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
}, outagePolicy);
|
||||
|
||||
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)),
|
||||
healthCoverageSource(config),
|
||||
loopHealth,
|
||||
quarantineSource,
|
||||
leadMailbox,
|
||||
outageSource,
|
||||
@@ -760,7 +763,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);
|
||||
@@ -1038,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.
|
||||
*
|
||||
* <p>{@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).
|
||||
|
||||
@@ -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<String> value) { }
|
||||
|
||||
/** Progress states for fleetd's singleton background loops, read by {@code fleet_list} and {@code /healthz}. */
|
||||
public record LoopHealthSource(Supplier<LoopWatchdog.State> statusPoller,
|
||||
Supplier<LoopWatchdog.State> 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<String> 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<String> 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<String> 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),
|
||||
@@ -1552,25 +1575,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<String, String> leads, String selfTerm) {
|
||||
Map<String, String> 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<String, String> 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<String, String> leads, String selfTerm) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, leads, selfTerm,
|
||||
LoopHealthSource loopHealth, QuarantineSource quarantine,
|
||||
Map<String, String> 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<String, String> 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<String, String> 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);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1586,8 +1626,8 @@ public final class FleetMcp {
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, Map<String, String> 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}). */
|
||||
@@ -1595,8 +1635,8 @@ public final class FleetMcp {
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
Map<String, String> 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);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1617,7 +1657,7 @@ public final class FleetMcp {
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
LeadSeatSource leadSeats, Map<String, String> 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);
|
||||
}
|
||||
|
||||
@@ -1637,8 +1677,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<String, String> 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<String, String> leads, String selfTerm,
|
||||
CoordinationSource coordination, boolean callerIsPrimary) {
|
||||
try {
|
||||
@@ -1662,6 +1712,9 @@ public final class FleetMcp {
|
||||
Map<String, Object> 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.
|
||||
@@ -2141,9 +2194,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/<slug>); 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 "
|
||||
|
||||
@@ -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<String> deliverable, Supplier<MemberCredentialPolicyView> 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<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials,
|
||||
FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage) {
|
||||
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> 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<String> deliverable, Supplier<MemberCredentialPolicyView> 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<String, Object> 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<String, Object> 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<String, String> 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)));
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -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.
|
||||
*
|
||||
* <p>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.
|
||||
*
|
||||
* <p>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<Capability> capabilities() {
|
||||
return Set.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> 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<String> 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<String> 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;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -9,6 +9,7 @@ import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
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;
|
||||
@@ -1084,6 +1085,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();
|
||||
|
||||
@@ -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();
|
||||
|
||||
Reference in New Issue
Block a user