fleetd #612 Shape A unit r4: pin quarantineSource + outageSource through both operator windows
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 1m29s
CI / build (pull_request) Failing after 1m56s

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.
This commit is contained in:
Dai Ha
2026-10-02 04:20:53 +02:00
parent 141ae3b04d
commit 91792e11fc
@@ -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.
*
* <p><strong>The property pinned here</strong> (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.
*
* <p><strong>No source-text assertion anywhere in this file.</strong> 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.
*
* <p><strong>How this gets past CB-185's own pid-resolution dead end.</strong> {@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.
*
* <p><strong>How the quarantined/cooling-off state is set up.</strong> 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<String, String> 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<Rendezvous.Resolution> 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<Rendezvous.Resolution> 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<Rendezvous.Resolution> 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<String> 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);
}
}