From 85c90d440a856adc6779b016d8b589e7e039bd6e Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Fri, 4 Sep 2026 12:12:35 +0700 Subject: [PATCH] #297: map HerdrException on GET /agents and /members; GET /profiles reports quarantine + cool-off Two REST-only visibility gaps, both against the same shared instances FleetMcp reads (BackendQuarantine/BackendOutagePolicy), never recomputed: - GET /agents and GET /members let a HerdrException escape uncaught, outside the {error, detail} envelope every other failure path in FleetApp uses. Both now route through the existing herdrError() helper, matching healthz/ sessionStatus. GET /members is the endpoint's own comment names as the out-of-band path a lead falls back to when its MCP mount drops. - GET /profiles omitted the two outage states fleet_profiles already reports: quarantined (CB-578 stage B) and coolingOff (fleetd #201 Unit 5). FleetApp now takes the SAME FleetMcp.QuarantineSource/OutageSource instances Fleetd wires into FleetMcp (extracted to local vars in Fleetd.java so both doors share one object, not two independently-built copies of the same rule). FleetMcp itself is unchanged. Item 3 of the ticket (a capacity block on GET /members) is explicitly out of scope and was not added. --- .../src/main/java/dev/ltms/fleet/Fleetd.java | 28 ++-- .../java/dev/ltms/fleet/rest/FleetApp.java | 138 ++++++++++++++---- .../dev/ltms/fleet/rest/FleetAppTest.java | 92 +++++++++++- 3 files changed, 217 insertions(+), 41 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index f63d2d8..6058535 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -625,6 +625,19 @@ public final class Fleetd { log.info("auth: loopback-trust (any loopback non-worker caller is the primary)"); } + // fleetd #297: named once and reused verbatim below for FleetApp's GET /profiles, rather than + // built a second time — two independently-constructed sources reading the SAME BackendQuarantine + // / BackendOutagePolicy would still be able to drift (e.g. a future edit to the credentialIdFor + // closure in only one of the two places), exactly the shape #284 was. + FleetMcp.QuarantineSource quarantineSource = new FleetMcp.QuarantineSource(profile -> { + var configured = config.get().profiles().get(profile); + return configured == null ? null : configured.effectiveCredentialId(); + }, quarantine); + FleetMcp.OutageSource outageSource = new FleetMcp.OutageSource(profile -> { + var configured = config.get().profiles().get(profile); + return configured == null ? null : configured.effectiveCredentialId(); + }, outagePolicy); + FleetMcp mcp = new FleetMcp(messages, workers, sessions, identity, presence, primaryRegistry, callers, metrics, new FleetMcp.CapacitySource(profile -> liveCountRef.get().apply(profile), profile -> { @@ -636,15 +649,9 @@ public final class Fleetd { return FleetHealthMonitor.coverage(health != null && health.isEnabled(), health != null && health.notifications() != null && health.notifications().configured()); }), - new FleetMcp.QuarantineSource(profile -> { - var configured = config.get().profiles().get(profile); - return configured == null ? null : configured.effectiveCredentialId(); - }, quarantine), + quarantineSource, leadMailbox, - new FleetMcp.OutageSource(profile -> { - var configured = config.get().profiles().get(profile); - return configured == null ? null : configured.effectiveCredentialId(); - }, outagePolicy), + outageSource, new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads))); // CB-637: the receive half. Only constructed when a lead mailbox actually opened — with no @@ -715,9 +722,12 @@ public final class Fleetd { // GET /sessions must merge across both, or a down/unpolled member daemon is invisible. // fleetd #111: live (re-read-per-request) memberCredentials view for GET /member-credentials — // same hot-reload shape as the memberCredentials supplier passed to ClaudeCodeLauncher above. + // fleetd #297: quarantineSource/outageSource are the SAME instances passed to FleetMcp above — + // GET /profiles must report the identical quarantine/cool-off facts as fleet_profiles. Javalin app = new FleetApp(herdr, memberHerdr, workers, sessions, messages, presence, mcp.servlet(), callers, metrics, deliverable, - () -> MemberCredentialPolicyView.of(config.get().memberCredentials())).build(); + () -> MemberCredentialPolicyView.of(config.get().memberCredentials()), + quarantineSource, outageSource).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/rest/FleetApp.java b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java index 891e575..94c9420 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java @@ -9,6 +9,7 @@ import dev.ltms.fleet.auth.Principal; import dev.ltms.fleet.guard.GuardException; import dev.ltms.fleet.metrics.FleetMetrics; import dev.ltms.fleet.metrics.Metrics; +import dev.ltms.fleet.mcp.FleetMcp; import dev.ltms.fleet.herdr.Agent; import dev.ltms.fleet.herdr.HerdrClient; import dev.ltms.fleet.herdr.HerdrException; @@ -86,6 +87,12 @@ public final class FleetApp { // absent() (the honest "no policy configured" view) for every constructor that does not wire // a real one, so existing legacy call sites keep building without knowing this field exists. private final Supplier memberCredentials; + // fleetd #297: the SAME shared sources FleetMcp.profiles/fleet_profiles reads (BackendQuarantine + // and BackendOutagePolicy are each one instance for the whole daemon — see Fleetd wiring) so + // GET /profiles cannot drift from fleet_profiles about which profile is quarantined/cooling off. + // .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 ObjectMapper mapper = new ObjectMapper(); /** @@ -146,6 +153,23 @@ public final class FleetApp { MessageService messages, MemberPresence presence, 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()); + } + + /** + * @param quarantine the SAME {@link FleetMcp.QuarantineSource} instance passed to {@code + * FleetMcp} (fleetd #297), so {@code GET /profiles} reports the identical + * exhaustion-quarantine facts as {@code fleet_profiles} rather than a second, + * independently-computed copy + * @param outage the SAME {@link FleetMcp.OutageSource} instance passed to {@code FleetMcp} — + * see {@code quarantine}; a SEPARATE check from it, never merged in + */ + 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) { this.herdr = herdr; this.memberHerdr = memberHerdr != null ? memberHerdr : herdr; this.workers = workers; @@ -156,6 +180,8 @@ public final class FleetApp { this.auth = auth; this.metrics = metrics; this.memberCredentials = memberCredentials != null ? memberCredentials : MemberCredentialPolicyView::absent; + this.quarantine = quarantine != null ? quarantine : FleetMcp.QuarantineSource.none(); + this.outage = outage != null ? outage : FleetMcp.OutageSource.none(); } /** Wire routes onto a fresh, unstarted Javalin instance. Caller starts it. */ @@ -338,8 +364,15 @@ public final class FleetApp { if (!allow(ctx, routeAction("GET /agents"), null)) { return; } - ctx.status(200).json(Map.of("agents", - workers.list().stream().map(Agent.class::cast).map(FleetApp::view).toList())); + try { + ctx.status(200).json(Map.of("agents", + workers.list().stream().map(Agent.class::cast).map(FleetApp::view).toList())); + } catch (HerdrException e) { + // fleetd #297: workers.list() reaches herdr — a transport failure must land in the same + // {error, detail} envelope every other failure path here uses, not escape as a bare + // exception and leave Javalin's default handling to respond outside the JSON contract. + herdrError(ctx, e); + } } /** CB-304: bridge-owned roster merged with live herdr status by paneId. */ @@ -347,40 +380,85 @@ public final class FleetApp { if (!allow(ctx, routeAction("GET /members"), null)) { return; } - // CB-519: the registry key is a host-unique id, not the pane coordinate — join on terminal. - Map live = workers.list().stream() - .map(Agent.class::cast) - .filter(a -> a.terminalId() != null) - .collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b)); - // fleetd #209: this REST roster reports agentSessionId via SessionManager.rosterView, so it - // uses the resolving roster read (caller-driven, not a timer) rather than the plain one. - List> out = sessions.rosterResolved().stream() - .map(s -> SessionManager.rosterView(s, live.get(s.terminalId()))) - .toList(); - Map body = new LinkedHashMap<>(); - // fleetd #199: the endpoint became /members in the CB-634 rename but the body key stayed - // "workers", so a caller that read "members" saw an empty fleet and reported no members at - // all. "members" is the canonical key; "workers" stays as a deprecated alias so an existing - // REST consumer keeps working — the out-of-band path a lead falls back to when its MCP mount - // drops reads this endpoint. Drop the alias once nothing reads it. - body.put("members", out); - body.put("workers", out); - // CB-586: operator visibility for the refs/wip snapshot store without shelling into the - // repo — how many snapshot refs exist and roughly what they cost. Present only once a - // worktree session has established the repo, so a never-snapshotted fleet reports nothing. - sessions.wipRefs().ifPresent(st -> body.put("wipRefs", - Map.of("count", st.count(), "costBytes", st.costBytes()))); - ctx.status(200).json(body); + try { + // CB-519: the registry key is a host-unique id, not the pane coordinate — join on terminal. + Map live = workers.list().stream() + .map(Agent.class::cast) + .filter(a -> a.terminalId() != null) + .collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b)); + // fleetd #209: this REST roster reports agentSessionId via SessionManager.rosterView, so it + // uses the resolving roster read (caller-driven, not a timer) rather than the plain one. + List> out = sessions.rosterResolved().stream() + .map(s -> SessionManager.rosterView(s, live.get(s.terminalId()))) + .toList(); + Map body = new LinkedHashMap<>(); + // fleetd #199: the endpoint became /members in the CB-634 rename but the body key stayed + // "workers", so a caller that read "members" saw an empty fleet and reported no members at + // all. "members" is the canonical key; "workers" stays as a deprecated alias so an existing + // REST consumer keeps working — the out-of-band path a lead falls back to when its MCP mount + // drops reads this endpoint. Drop the alias once nothing reads it. + body.put("members", out); + body.put("workers", out); + // CB-586: operator visibility for the refs/wip snapshot store without shelling into the + // repo — how many snapshot refs exist and roughly what they cost. Present only once a + // worktree session has established the repo, so a never-snapshotted fleet reports nothing. + sessions.wipRefs().ifPresent(st -> body.put("wipRefs", + Map.of("count", st.count(), "costBytes", st.costBytes()))); + ctx.status(200).json(body); + } catch (HerdrException e) { + // fleetd #297: same reasoning as agents() above — this is the out-of-band roster a lead + // falls back to when its MCP mount drops, so it must stay inside the JSON error contract + // exactly when herdr is briefly unreachable, not escape as a bare exception. + herdrError(ctx, e); + } } - /** The configured worker profiles and which one a no-argument spawn uses. */ + /** + * The configured worker profiles, which one a no-argument spawn uses, and (fleetd #297) the two + * outage states {@code fleet_profiles} already reports: {@code quarantined} (CB-578 stage B — + * the backend reported it out of capacity) and {@code coolingOff} (fleetd #201 Unit 5 — the + * credential threw repeated non-exhaustion backend errors). Both are read from the SAME shared + * {@link FleetMcp.QuarantineSource}/{@link FleetMcp.OutageSource} instances {@code FleetMcp} + * reads, never recomputed, so the two doors cannot disagree about which profile is down and why. + * Independent checks, so a profile can appear in both maps at once; each map is present only + * when at least one profile is in that state. + */ private void profiles(Context ctx) { if (!allow(ctx, routeAction("GET /profiles"), null)) { return; } - ctx.status(200).json(Map.of( - "profiles", workers.profiles(), - "default", workers.defaultProfile() == null ? "" : workers.defaultProfile())); + Map body = new LinkedHashMap<>(); + body.put("profiles", workers.profiles()); + body.put("default", workers.defaultProfile() == null ? "" : workers.defaultProfile()); + Map quarantined = new LinkedHashMap<>(); + Map coolingOff = new LinkedHashMap<>(); + for (String profile : workers.profiles()) { + String credentialId = quarantine.credentialIdFor().apply(profile); + if (credentialId != null) { + quarantine.quarantine().remainingSeconds(credentialId).ifPresent(remaining -> { + Map row = new LinkedHashMap<>(); + row.put("credentialId", credentialId); + row.put("quarantinedForSeconds", remaining); + quarantined.put(profile, row); + }); + } + String outageCredentialId = outage.credentialIdFor().apply(profile); + if (outageCredentialId != null) { + outage.outagePolicy().remainingCoolOffSeconds(outageCredentialId).ifPresent(remaining -> { + Map row = new LinkedHashMap<>(); + row.put("credentialId", outageCredentialId); + row.put("coolingOffForSeconds", remaining); + coolingOff.put(profile, row); + }); + } + } + if (!quarantined.isEmpty()) { + body.put("quarantined", quarantined); + } + if (!coolingOff.isEmpty()) { + body.put("coolingOff", coolingOff); + } + ctx.status(200).json(body); } /** 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 a0a5768..dc9166b 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppTest.java @@ -19,6 +19,10 @@ import dev.ltms.fleet.session.SessionManager; import dev.ltms.fleet.session.Worktrees; import dev.ltms.fleet.member.ClaudeCodeLauncher; import dev.ltms.fleet.member.CompositePeerLauncher; +import dev.ltms.fleet.member.MemberCredentialPolicyView; +import dev.ltms.fleet.mcp.FleetMcp; +import dev.ltms.fleet.placement.BackendOutagePolicy; +import dev.ltms.fleet.placement.BackendQuarantine; import dev.ltms.fleet.placement.PlacementPolicies; import io.javalin.Javalin; import org.junit.jupiter.api.AfterEach; @@ -32,6 +36,7 @@ import java.util.List; import java.util.Map; import java.util.Set; import java.util.UUID; +import java.util.concurrent.TimeUnit; import java.util.function.Predicate; import static org.junit.jupiter.api.Assertions.*; @@ -70,6 +75,19 @@ class FleetAppTest { private int start(FakeHerdr herdr, String workerBaseUrl, Set allow, String placement, Worktrees worktrees, Predicate deliverable) { + return start(herdr, workerBaseUrl, allow, placement, worktrees, deliverable, + FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none()); + } + + /** + * fleetd #297: same wiring as above, plus the two SAME shared sources {@code GET /profiles} + * must read — lets a test prove the quarantined/coolingOff facts it reports come from a real + * {@link dev.ltms.fleet.placement.BackendQuarantine}/{@link + * dev.ltms.fleet.placement.BackendOutagePolicy}, exactly like {@code fleet_profiles}'s own tests. + */ + private int start(FakeHerdr herdr, String workerBaseUrl, Set allow, String placement, + Worktrees worktrees, Predicate deliverable, + FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage) { FleetConfig.Profile wcfg = new FleetConfig.Profile( "ltms-local", workerBaseUrl, "coder", null, "FLEETD_WORKER_TOKEN", null, placement, "fleet", "worker: {profile} #{n}", null, null, null); @@ -90,8 +108,9 @@ class FleetAppTest { // it directly so the inbox contract holds for those endpoints. inbox.own("term_a"); MessageService messages = new MessageService(agents, injector, rendezvous, inbox); - app = new FleetApp(herdr, workers, sessions, messages, this.presence, null, - null, null, id -> this.presence.isPresent(id) || deliverable.test(id)) + app = new FleetApp(herdr, herdr, workers, sessions, messages, this.presence, null, + null, null, id -> this.presence.isPresent(id) || deliverable.test(id), + MemberCredentialPolicyView::absent, quarantine, outage) .build().start("127.0.0.1", 0); return app.port(); } @@ -164,6 +183,22 @@ class FleetAppTest { assertEquals("idle", agents.get(0).get("status").asText()); } + /** + * fleetd #297 gap 1: {@code workers.list()} reaches herdr, and a transport failure there must + * land in the same {@code {error, detail}} envelope every other failure path in this file uses + * (see {@code herdrError}), not escape as a bare exception outside the JSON contract. + */ + @Test + void agentsMapsAHerdrFailureToTheJsonErrorEnvelope() throws Exception { + FakeHerdr down = new FakeHerdr().healthy(false); + int port = start(down, "http://gx00.gw:8000", Set.of("gx00.gw")); + HttpResponse res = req(port, "GET", "/agents"); + assertEquals(502, res.statusCode(), res.body()); + JsonNode body = mapper.readTree(res.body()); + assertEquals("herdr_error", body.get("error").asText()); + assertTrue(body.has("detail"), res.body()); + } + @Test void spawnWorkerLandsInOwnTabInWorkerSpaceAndInjectsBaseUrl() throws Exception { FakeHerdr herdr = new FakeHerdr(); @@ -203,6 +238,42 @@ class FleetAppTest { JsonNode body = mapper.readTree(req(port, "GET", "/profiles").body()); assertEquals("ltms-local", body.get("default").asText()); assertEquals("ltms-local", body.get("profiles").get(0).asText()); + assertFalse(body.has("quarantined"), "nothing is quarantined, so the key is omitted: " + body); + assertFalse(body.has("coolingOff"), "nothing is cooling off, so the key is omitted: " + body); + } + + /** + * fleetd #297 gap 2: {@code GET /profiles} must report the same two outage states {@code + * fleet_profiles} does — CB-578 stage B exhaustion quarantine and fleetd #201 Unit 5 cool-off — + * reading the SAME shared {@link BackendQuarantine}/{@link BackendOutagePolicy} instances rather + * than recomputing them. The two checks are independent, and this profile is deliberately put in + * both states at once, matching {@code FleetMcpTest}'s own coverage of that overlap. + */ + @Test + void profilesReportsQuarantineAndCoolingOffFromTheSameSharedSources() throws Exception { + FakeHerdr herdr = new FakeHerdr(); + BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30)); + quarantine.quarantine("shared-openai"); + FleetMcp.QuarantineSource quarantineSource = new FleetMcp.QuarantineSource( + profile -> "ltms-local".equals(profile) ? "shared-openai" : null, quarantine); + BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L); + outagePolicy.record("shared-openai", "t1", "API Error: rate limited"); + outagePolicy.record("shared-openai", "t2", "API Error: rate limited"); // 2nd distinct target starts the incident + FleetMcp.OutageSource outageSource = new FleetMcp.OutageSource( + profile -> "ltms-local".equals(profile) ? "shared-openai" : null, outagePolicy); + int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"), "tab", new GitWorktrees(), + ignored -> false, quarantineSource, outageSource); + + JsonNode body = mapper.readTree(req(port, "GET", "/profiles").body()); + assertTrue(body.has("quarantined"), body.toString()); + assertEquals("shared-openai", + body.get("quarantined").get("ltms-local").get("credentialId").asText()); + assertEquals(1800, + body.get("quarantined").get("ltms-local").get("quarantinedForSeconds").asLong()); + assertTrue(body.has("coolingOff"), body.toString()); + assertEquals("shared-openai", + body.get("coolingOff").get("ltms-local").get("credentialId").asText()); + assertEquals(60, body.get("coolingOff").get("ltms-local").get("coolingOffForSeconds").asLong()); } @Test @@ -239,6 +310,23 @@ class FleetAppTest { "liveStatus is unknown when herdr has no matching pane"); } + /** + * fleetd #297 gap 1: same reasoning as {@code agentsMapsAHerdrFailureToTheJsonErrorEnvelope} — + * {@code GET /members} is the endpoint's own comment names as "the out-of-band path a lead falls + * back to when its MCP mount drops", so it must stay inside the {@code {error, detail}} envelope + * exactly when herdr is briefly unreachable. + */ + @Test + void membersMapsAHerdrFailureToTheJsonErrorEnvelope() throws Exception { + FakeHerdr down = new FakeHerdr().healthy(false); + int port = start(down, "http://gx00.gw:8000", Set.of("gx00.gw")); + HttpResponse res = req(port, "GET", "/members"); + assertEquals(502, res.statusCode(), res.body()); + JsonNode body = mapper.readTree(res.body()); + assertEquals("herdr_error", body.get("error").asText()); + assertTrue(body.has("detail"), res.body()); + } + @Test void spawnWithACwdParamRootsTheWorkerThere() throws Exception { FakeHerdr herdr = new FakeHerdr();