From 91792e11fcd98888c81131aa1bb8338e850c5811 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Fri, 2 Oct 2026 04:20:53 +0200 Subject: [PATCH] fleetd #612 Shape A unit r4: pin quarantineSource + outageSource through both operator windows FleetdAssembly.java's quarantineSource (:471-472) and outageSource (:473-476) each feed two consumers: FleetMcp (fleet_profiles, :484/:486) and FleetApp (GET /profiles, :529). No existing test distinguished the two windows for either source. New test drives the real FleetdAssembly.assembleAndStart, classifies a real exhaustion/outage through the real CompletionResolver, and reads the result back through a real McpSyncClient (fleet_profiles) and a real HttpClient (GET /profiles), both authenticated via token-mode auth (sidesteps the in-JVM pid-resolution dead end). No source text is read; no production code changed. Verified with six mutation cycles (3 per site x 2 sites: FleetMcp starved, FleetApp starved, mis-wire with a disconnected collaborator), each run against the full unfiltered suite and reverted after confirming the expected test(s) alone went red. The quarantineSource FleetMcp-starve and mis-wire cycles also trip three pre-existing tests that read the live BackendQuarantine via FleetMcp#quarantineSource() for their own unrelated assertions - a pre-existing incidental coupling, not newly introduced here. Out of scope, noted per the ticket's dispatch comment: loopHealthSource (FleetdAssembly.java:478) shares this same two-consumer shape and is already assigned to a separate unit, r10. --- ...uarantineOutageDualWindowAssemblyTest.java | 368 ++++++++++++++++++ 1 file changed, 368 insertions(+) create mode 100644 fleetd/src/test/java/dev/ltms/fleet/FleetdQuarantineOutageDualWindowAssemblyTest.java diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdQuarantineOutageDualWindowAssemblyTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdQuarantineOutageDualWindowAssemblyTest.java new file mode 100644 index 0000000..43bcf34 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdQuarantineOutageDualWindowAssemblyTest.java @@ -0,0 +1,368 @@ +package dev.ltms.fleet; + +import dev.ltms.fleet.config.ConfigRef; +import dev.ltms.fleet.config.FleetConfig; +import dev.ltms.fleet.guard.SubscriptionGuard; +import dev.ltms.fleet.herdr.FakeHerdr; +import dev.ltms.fleet.herdr.HerdrClient; +import dev.ltms.fleet.inject.CompletionResolver; +import dev.ltms.fleet.msg.Rendezvous; +import dev.ltms.fleet.msg.TurnToken; +import dev.ltms.fleet.session.MemberSession; +import io.javalin.Javalin; +import io.modelcontextprotocol.client.McpClient; +import io.modelcontextprotocol.client.McpSyncClient; +import io.modelcontextprotocol.client.transport.HttpClientStreamableHttpTransport; +import io.modelcontextprotocol.spec.McpClientTransport; +import io.modelcontextprotocol.spec.McpSchema; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; +import org.junit.jupiter.api.io.TempDir; + +import java.net.URI; +import java.net.http.HttpClient; +import java.net.http.HttpRequest; +import java.net.http.HttpResponse; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; +import java.util.function.LongSupplier; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * fleetd #612 Shape A, unit r4 — {@link FleetdAssembly}'s {@code quarantineSource} (lines 471-472) + * and {@code outageSource} (lines 473-476 at {@code main} = {@code 141ae3b}), EACH of which feeds + * two separate consumers: {@code FleetMcp} ({@code fleet_profiles}) and {@code FleetApp} + * ({@code GET /profiles}), one call site ({@code :484}/{@code :486}) into the MCP constructor and + * the SAME shared local again ({@code :529}) into the REST constructor. + * + *

The property pinned here (from the ticket): when the real assembly has built + * a real {@link dev.ltms.fleet.placement.BackendQuarantine} holding a quarantined credential, and a + * real {@link dev.ltms.fleet.placement.BackendOutagePolicy} holding a cooling-off credential, BOTH + * operator windows must report that state — the real assembled {@code FleetMcp} (reached through + * {@link FleetdRuntime#mcp()}) and the real assembled REST surface (reached through {@link + * FleetdRuntime#app()}). A mutation that starves one consumer while leaving the other wired must + * make only that consumer's assertion go red. + * + *

No source-text assertion anywhere in this file. Both windows are read off the + * REAL running objects: {@code fleet_profiles} is called through a real MCP client over a real + * HTTP connection to the servlet {@link FleetdAssembly} actually mounted, and {@code GET /profiles} + * is called through a real {@link java.net.http.HttpClient} against the real bound {@link + * FleetdRuntime#app()}. Neither is a copy built alongside the assembly for this test's benefit. + * + *

How this gets past CB-185's own pid-resolution dead end. {@code + * FleetdAssemblyFleetAppTest}'s class javadoc explains that {@code GET /sessions} cannot be driven + * over real HTTP here because {@code LsofPeerPidLookup} excludes its own pid and an in-process test + * client/server share one JVM pid — every such request resolves {@code ANONYMOUS} and is refused + * before the handler runs. {@code GET /profiles} and {@code fleet_profiles} sit behind the exact + * same {@code Authz.Action.READ} gate. This test sidesteps the dead end instead of hitting it: + * {@code auth.mode: token} (see {@link dev.ltms.fleet.auth.CallerResolver#resolve}) resolves a + * caller to {@code PRIMARY} from a valid {@code Authorization: Bearer} header ALONE, with no pid + * resolution involved at all — the same technique {@code FleetMcpContextExtractorTest} already uses + * to drive a real {@code fleet_whoami} call through the real transport. + * + *

How the quarantined/cooling-off state is set up. Both {@code BackendQuarantine} + * and {@code BackendOutagePolicy} are private to the collaborators the assembly wires them into, and + * (unlike {@code quarantineSource()}) neither {@code FleetMcp} nor {@code FleetApp} exposes a public + * accessor for the live {@code BackendOutagePolicy} instance. Rather than add one (acceptance + * criterion 1: no production change), this test drives the REAL production classification path — + * exactly the recipe {@code FleetdExhaustedPatternAssemblyTest} (quarantine) and {@code + * FleetdCompletionResolverAssemblyTest} (cool-off) already proved works end to end against this same + * {@link FleetdAssembly#assembleAndStart}: acquire a real {@link MemberSession}, feed the real {@link + * CompletionResolver} a pane scrape matching the profile's configured {@code exhaustedPattern} / + * {@code errorPattern}, and let the real {@code exhaustionSink}/{@code backendErrorSink} write into + * the real, shared tracker. Each setup step asserts its own {@link Rendezvous.Kind} as a CONTROL — + * if the resolver were never actually exercised, the setup itself fails loudly before either window + * is ever read. + */ +class FleetdQuarantineOutageDualWindowAssemblyTest { + + private static final String TOKEN = "s3cret-r4-token"; + private static final String TOKEN_ENV = "FLEETD_R4_TEST_TOKEN"; + + private static final class ControllableResourcePorts implements ResourcePorts { + + final FakeHerdr herdr; + final AtomicLong nowNanos = new AtomicLong(1_000_000_000L); // arbitrary non-zero start + Runnable shutdownHook; + + ControllableResourcePorts(FakeHerdr herdr) { + this.herdr = herdr; + } + + void advanceSeconds(long seconds) { + nowNanos.addAndGet(TimeUnit.SECONDS.toNanos(seconds)); + } + + @Override + public Map environment() { + return Map.of(TOKEN_ENV, TOKEN); + } + + @Override + public HerdrClient connectHerdr(Path socketPath) { + return herdr; + } + + @Override + public Fleetd.AmqpOpener replyInboxOpener() { + // Never invoked: this test's config has no `broker:` block. + return (uri, prefetch) -> { + throw new UnsupportedOperationException("replyInboxOpener must not be called — no broker: block"); + }; + } + + @Override + public Fleetd.LeadMailboxOpener leadMailboxOpener() { + // Never invoked: this test's config has no `coordinator:` block. + return (uri, selfCoordId, prefetch) -> { + throw new UnsupportedOperationException( + "leadMailboxOpener must not be called — no coordinator: block is configured"); + }; + } + + @Override + public LongSupplier nanoClock() { + return nowNanos::get; + } + + @Override + public LongSupplier wallClockNanos() { + return nowNanos::get; + } + + @Override + public ScheduledExecutorService newScheduler(String purpose) { + return Executors.newSingleThreadScheduledExecutor(); + } + + @Override + public void addShutdownHook(Runnable hook) { + this.shutdownHook = hook; + } + + @Override + public void startHttp(Javalin app, String host, int port) { + // Deliberately never bind here — this test binds runtime.app() itself, for real, below. + } + + @Override + public Runnable herdrPollWait() { + // Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls. + return () -> { + throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy"); + }; + } + } + + private FleetdRuntime runtime; + private ControllableResourcePorts ports; + private Javalin boundApp; + + @AfterEach + void tearDown() { + if (boundApp != null) { + boundApp.stop(); + } + if (ports != null && ports.shutdownHook != null) { + ports.shutdownHook.run(); + } + } + + private static FleetConfig writeConfig(Path dir, int cooldownSeconds) throws Exception { + Path f = dir.resolve("fleetd.yaml"); + Files.writeString(f, """ + bind: + host: 127.0.0.1 + port: 8765 + auth: + mode: token + tokenEnv: %s + idleSleepGuard: + enabled: false + quarantineCooldownSeconds: %d + profiles: + exhaustprofile: + baseUrl: http://exhausthost.local:8000 + model: sonnet + exhaustedPattern: "usage limit reached" + coolprofile: + baseUrl: http://coolhost.local:8000 + model: sonnet + errorPattern: "credential outage" + guard: + offSubscriptionHosts: + - exhausthost.local + - coolhost.local + """.formatted(TOKEN_ENV, cooldownSeconds)); + return FleetConfig.load(f); + } + + /** Assembles the real graph, then binds the real {@code Javalin app} to an ephemeral port. */ + private int assembleAndBind(Path dir) throws Exception { + FleetConfig cfg = writeConfig(dir, 120); + ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg); + SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet()); + ports = new ControllableResourcePorts(new FakeHerdr()); + runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports); + boundApp = runtime.app().start("127.0.0.1", 0); + return boundApp.port(); + } + + /** + * Drives the real assembled {@link CompletionResolver} through a scrape matching {@code + * exhaustprofile}'s configured {@code exhaustedPattern}, exactly {@code + * FleetdExhaustedPatternAssemblyTest}'s own recipe, so the real {@code exhaustionSink} quarantines + * the credential ({@code effectiveCredentialId() == "exhaustprofile"}, no explicit credentialId + * configured). + */ + private void quarantineExhaustProfile(Path dir) { + MemberSession session = runtime.sessions().acquire("exhaustprofile", null, dir.toString(), null); + String target = session.terminalId(); + + CompletionResolver completion = runtime.completion(); + CompletableFuture waiter = new CompletableFuture<>(); + + ports.herdr.readText("idle, nothing yet"); + completion.onDelivered(target, new TurnToken(target, waiter, null)); + ports.herdr.readText("usage limit reached: try again in a few hours"); + ports.advanceSeconds(3); // clear CompletionResolver.MIN_TURN_NANOS (2s), no real sleep + completion.resolveBeforePostAction(target); + + Rendezvous.Resolution resolution = waiter.getNow(null); + assertTrue(resolution != null && resolution.kind() == Rendezvous.Kind.BACKEND_EXHAUSTED, + "CONTROL: setup must classify as BACKEND_EXHAUSTED before either window is read — " + + "if this fails, the assembled resolver was never actually exercised: " + resolution); + } + + /** + * Drives the real assembled {@link CompletionResolver} with TWO distinct targets on {@code + * coolprofile}, each matching its configured {@code errorPattern}, exactly {@code + * FleetdCompletionResolverAssemblyTest}'s own recipe, so the real {@code backendErrorSink} cools + * the credential off ({@code effectiveCredentialId() == "coolprofile"}). + */ + private void coolOffCoolProfile(Path dir) { + MemberSession s1 = runtime.sessions().acquire("coolprofile", null, dir.toString(), null); + MemberSession s2 = runtime.sessions().acquire("coolprofile", null, dir.toString(), null); + CompletionResolver completion = runtime.completion(); + + String t1 = s1.terminalId(); + CompletableFuture w1 = new CompletableFuture<>(); + ports.herdr.readText("idle 1"); + completion.onDelivered(t1, new TurnToken(t1, w1, null)); + ports.herdr.readText("credential outage: upstream 503"); + ports.advanceSeconds(3); + completion.resolveBeforePostAction(t1); + Rendezvous.Resolution r1 = w1.getNow(null); + assertTrue(r1 != null && r1.kind() == Rendezvous.Kind.FAILED, + "CONTROL: target1's setup must classify FAILED (backend error): " + r1); + + String t2 = s2.terminalId(); + CompletableFuture w2 = new CompletableFuture<>(); + ports.herdr.readText("idle 2"); + completion.onDelivered(t2, new TurnToken(t2, w2, null)); + ports.herdr.readText("credential outage: upstream 503 again"); + ports.advanceSeconds(3); + completion.resolveBeforePostAction(t2); + Rendezvous.Resolution r2 = w2.getNow(null); + assertTrue(r2 != null && r2.kind() == Rendezvous.Kind.FAILED, + "CONTROL: target2's setup must classify FAILED (backend error) — two distinct " + + "targets are required to cross BackendOutagePolicy.THRESHOLD: " + r2); + } + + /** Calls the real {@code fleet_profiles} tool over a real MCP client, token-authenticated as PRIMARY. */ + private static McpSchema.CallToolResult callProfilesViaMcp(int port) { + HttpRequest.Builder requestTemplate = HttpRequest.newBuilder().header("Authorization", "Bearer " + TOKEN); + McpClientTransport transport = HttpClientStreamableHttpTransport.builder("http://127.0.0.1:" + port) + .endpoint("/mcp") + .requestBuilder(requestTemplate) + .build(); + try (McpSyncClient client = McpClient.sync(transport).build()) { + client.initialize(); + return client.callTool(McpSchema.CallToolRequest.builder("fleet_profiles").arguments(Map.of()).build()); + } + } + + private static String textOf(McpSchema.CallToolResult r) { + return ((McpSchema.TextContent) r.content().getFirst()).text(); + } + + /** Calls the real {@code GET /profiles} route over a real {@link HttpClient}, same token. */ + private static String getProfilesViaRest(int port) throws Exception { + HttpClient http = HttpClient.newHttpClient(); + HttpRequest req = HttpRequest.newBuilder(URI.create("http://127.0.0.1:" + port + "/profiles")) + .header("Authorization", "Bearer " + TOKEN) + .GET().build(); + HttpResponse res = http.send(req, HttpResponse.BodyHandlers.ofString()); + assertEquals(200, res.statusCode(), "GET /profiles must succeed with the real token: " + res.body()); + return res.body(); + } + + // --- quarantineSource (FleetdAssembly.java :471-472) -------------------------------------- + + @Test + @Timeout(value = 15, unit = TimeUnit.SECONDS, threadMode = Timeout.ThreadMode.SEPARATE_THREAD) + void quarantinedCredentialIsReportedByTheRealAssembledFleetMcp(@TempDir Path dir) throws Exception { + int port = assembleAndBind(dir); + quarantineExhaustProfile(dir); + + String out = textOf(callProfilesViaMcp(port)); + + assertTrue(out.contains("\"quarantined\""), "fleet_profiles must report a quarantined " + + "section once the real BackendQuarantine holds a quarantined credential: " + out); + assertTrue(out.contains("\"exhaustprofile\""), out); + assertTrue(out.contains("\"quarantinedForSeconds\""), out); + } + + @Test + @Timeout(value = 15, unit = TimeUnit.SECONDS, threadMode = Timeout.ThreadMode.SEPARATE_THREAD) + void quarantinedCredentialIsReportedByTheRealAssembledFleetApp(@TempDir Path dir) throws Exception { + int port = assembleAndBind(dir); + quarantineExhaustProfile(dir); + + String out = getProfilesViaRest(port); + + assertTrue(out.contains("\"quarantined\""), "GET /profiles must report a quarantined " + + "section once the real BackendQuarantine holds a quarantined credential: " + out); + assertTrue(out.contains("\"exhaustprofile\""), out); + assertTrue(out.contains("\"quarantinedForSeconds\""), out); + } + + // --- outageSource (FleetdAssembly.java :473-476) ------------------------------------------- + + @Test + @Timeout(value = 15, unit = TimeUnit.SECONDS, threadMode = Timeout.ThreadMode.SEPARATE_THREAD) + void coolingOffCredentialIsReportedByTheRealAssembledFleetMcp(@TempDir Path dir) throws Exception { + int port = assembleAndBind(dir); + coolOffCoolProfile(dir); + + String out = textOf(callProfilesViaMcp(port)); + + assertTrue(out.contains("\"coolingOff\""), "fleet_profiles must report a coolingOff " + + "section once the real BackendOutagePolicy holds a cooling-off credential: " + out); + assertTrue(out.contains("\"coolprofile\""), out); + assertTrue(out.contains("\"coolingOffForSeconds\""), out); + } + + @Test + @Timeout(value = 15, unit = TimeUnit.SECONDS, threadMode = Timeout.ThreadMode.SEPARATE_THREAD) + void coolingOffCredentialIsReportedByTheRealAssembledFleetApp(@TempDir Path dir) throws Exception { + int port = assembleAndBind(dir); + coolOffCoolProfile(dir); + + String out = getProfilesViaRest(port); + + assertTrue(out.contains("\"coolingOff\""), "GET /profiles must report a coolingOff " + + "section once the real BackendOutagePolicy holds a cooling-off credential: " + out); + assertTrue(out.contains("\"coolprofile\""), out); + assertTrue(out.contains("\"coolingOffForSeconds\""), out); + } +}