Compare commits
10 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| a42b12440c | |||
| a06426c33c | |||
| 28ea0de575 | |||
| 91792e11fc | |||
| 6539efe9fa | |||
| 68397f78d5 | |||
| af901ff1d2 | |||
| dac5f88812 | |||
| 2e663e5968 | |||
| 8c14ed2846 |
@@ -837,41 +837,6 @@ public final class FleetMcp {
|
||||
return leadRollover;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #612 Shape A r9+r11 — {@code public} for the same cross-package reason as {@link
|
||||
* #quarantineSource()}, for the real {@link CapacitySource} this daemon was assembled with.
|
||||
* {@code fleet_list}'s {@code capacity} row is gated on {@link CapacitySource#available()},
|
||||
* which is {@code false} whenever {@code configuredProfiles} is empty — exactly what {@link
|
||||
* CapacitySource#none()} (the inert stand-in) always reports, so a test can tell a real,
|
||||
* populated source apart from one swapped for the inert variant at the {@code new FleetMcp(...)}
|
||||
* call site in {@code FleetdAssembly}.
|
||||
*/
|
||||
public CapacitySource capacitySource() {
|
||||
return capacity;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #612 Shape A r9+r11 — {@code public} for the same cross-package reason as {@link
|
||||
* #quarantineSource()}, for the real {@link HealthCoverageSource} this daemon was assembled
|
||||
* with. {@code fleet_list}'s {@code healthCoverage} field reads {@link
|
||||
* HealthCoverageSource#value()} straight off this exact instance.
|
||||
*/
|
||||
public HealthCoverageSource healthCoverageSource() {
|
||||
return healthCoverage;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #612 Shape A r9+r11 — {@code public} for the same cross-package reason as {@link
|
||||
* #quarantineSource()}, for the real {@code coordinator.peers} list this daemon was assembled
|
||||
* with (see {@link CoordinationSource#peers()}). {@code fleet_list}'s {@code coordinator.peers}
|
||||
* row reports exactly this list. Both an absent {@code coordinator:} block and a mis-wired call
|
||||
* site report {@code List.of()}, so a test must configure a real {@code coordinator:} block
|
||||
* with peers to tell a working wiring from the inert one.
|
||||
*/
|
||||
public List<String> coordinatorPeers() {
|
||||
return peers;
|
||||
}
|
||||
|
||||
// --- tool logic (thin adapters over the services; unit-testable) ---------------------------
|
||||
|
||||
/**
|
||||
|
||||
@@ -0,0 +1,180 @@
|
||||
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.LoopWatchdog;
|
||||
import dev.ltms.fleet.mcp.FleetMcp;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.lang.reflect.Field;
|
||||
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.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #612 Shape A, rank 10: {@link FleetdAssembly} creates one {@link FleetMcp.LoopHealthSource}
|
||||
* from the real started {@code StatusPoller} and {@code SessionReaper}, then gives it to two operator
|
||||
* windows. These tests reach the real assembled objects through {@link FleetdRuntime}, rather than
|
||||
* building a second source beside them. A hardcoded {@code RUNNING} source would pass a simple
|
||||
* "running" test, so the mutation proof also mis-wires the source to never-started loops: both
|
||||
* windows must then report {@code STOPPED} and these assertions go red.
|
||||
*/
|
||||
class FleetdAssemblyLoopHealthTest {
|
||||
|
||||
private static final class RecordingResourcePorts implements ResourcePorts {
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
Runnable shutdownHook;
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> new dev.ltms.fleet.msg.InMemoryReplyInbox();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException("no coordinator is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("healthy FakeHerdr must not be polled");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
// Bind runtime.app() to an ephemeral port only in the REST assertion below.
|
||||
}
|
||||
}
|
||||
|
||||
private final HttpClient http = HttpClient.newHttpClient();
|
||||
private RecordingResourcePorts ports;
|
||||
private FleetdRuntime runtime;
|
||||
private Javalin boundApp;
|
||||
|
||||
@AfterEach
|
||||
void tearDown() {
|
||||
if (boundApp != null) {
|
||||
boundApp.stop();
|
||||
}
|
||||
if (runtime != null) {
|
||||
runtime.close();
|
||||
}
|
||||
if (ports != null) {
|
||||
assertNotNull(ports.shutdownHook, "the assembly must register its shutdown hook");
|
||||
assertTrue(ports.herdr.closed, "FleetdRuntime.close must close the real assembled herdr client");
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
lifecycle:
|
||||
idleTtlSeconds: 600
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
health:
|
||||
enabled: false
|
||||
broker:
|
||||
uri: "amqp://fake-test-broker/vh"
|
||||
""");
|
||||
return FleetConfig.load(file);
|
||||
}
|
||||
|
||||
private void assemble(Path dir) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
ports = new RecordingResourcePorts();
|
||||
runtime = FleetdAssembly.assembleAndStart(
|
||||
new AssemblyInputs(cfg, new ConfigRef(dir.resolve("fleetd.yaml"), cfg),
|
||||
new SubscriptionGuard(cfg.guard().hostSet())), ports);
|
||||
}
|
||||
|
||||
@Test
|
||||
void fleetListUsesTheRunningLoopsInTheRealAssembledMcp(@TempDir Path dir) throws Exception {
|
||||
assemble(dir);
|
||||
|
||||
// The field is the exact source captured by FleetMcp's fleet_list handler. Reflection is
|
||||
// necessary because FleetMcp has no public source accessor; it is not a source-text check.
|
||||
FleetMcp.LoopHealthSource loopHealth = loopHealthOf(runtime.mcp());
|
||||
assertRunning(loopHealth, "FleetMcp's real fleet_list source");
|
||||
}
|
||||
|
||||
@Test
|
||||
void healthzUsesTheRunningLoopsInTheRealAssembledApp(@TempDir Path dir) throws Exception {
|
||||
assemble(dir);
|
||||
boundApp = runtime.app().start("127.0.0.1", 0);
|
||||
|
||||
HttpRequest request = HttpRequest.newBuilder(
|
||||
URI.create("http://127.0.0.1:" + boundApp.port() + "/healthz")).GET().build();
|
||||
HttpResponse<String> response = http.send(request, HttpResponse.BodyHandlers.ofString());
|
||||
assertEquals(200, response.statusCode(), response.body());
|
||||
assertTrue(response.body().contains("\"statusPoller\":\"RUNNING\""), response.body());
|
||||
assertTrue(response.body().contains("\"sessionReaper\":\"RUNNING\""), response.body());
|
||||
}
|
||||
|
||||
private static FleetMcp.LoopHealthSource loopHealthOf(FleetMcp mcp) throws Exception {
|
||||
Field field = FleetMcp.class.getDeclaredField("loopHealth");
|
||||
field.setAccessible(true);
|
||||
return (FleetMcp.LoopHealthSource) field.get(mcp);
|
||||
}
|
||||
|
||||
private static void assertRunning(FleetMcp.LoopHealthSource source, String consumer) {
|
||||
assertEquals(LoopWatchdog.State.RUNNING, source.statusPoller().get(),
|
||||
consumer + " must report the started real StatusPoller as RUNNING");
|
||||
assertEquals(LoopWatchdog.State.RUNNING, source.sessionReaper().get(),
|
||||
consumer + " must report the started real SessionReaper as RUNNING");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,184 @@
|
||||
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.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.inject.TurnListener;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.msg.TurnToken;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.lang.reflect.Field;
|
||||
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.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #612 rank 12 — {@code FleetdAssembly} supplies the {@link Injector}'s {@code TurnRegistrar}
|
||||
* with {@code Fleetd.turnRegistrar(completion)}. The normal delivery path cannot distinguish that
|
||||
* registrar from {@code TurnRegistrar.NOOP}: {@link dev.ltms.fleet.inject.CompletionResolver#onDelivered}
|
||||
* registers the same turn shortly afterwards. The distinction matters when a delivery listener throws
|
||||
* after the pane received the message but before normal completion runs. The registrar must already have
|
||||
* registered the waiter with the REAL assembled resolver, so a later completion can still resolve it.
|
||||
*
|
||||
* <p>The test installs a throwing wrapper around the real assembled listener. It does not replace the
|
||||
* registrar or resolver. This constructs the narrow failure condition without a sleep, then observes the
|
||||
* real resolver through {@link FleetdRuntime#completion()}. An inert registrar, or a registrar wired to a
|
||||
* throwaway resolver, leaves this waiter's turn absent from the real resolver and makes the final assertion
|
||||
* fail.
|
||||
*/
|
||||
class FleetdAssemblyTurnRegistrarBehaviouralTest {
|
||||
|
||||
private static final String TARGET = "term_a";
|
||||
|
||||
private static final class ControllableResourcePorts implements ResourcePorts {
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
final AtomicLong nowNanos = new AtomicLong(1_000_000_000L);
|
||||
Runnable shutdownHook;
|
||||
|
||||
void advanceSeconds(long seconds) {
|
||||
nowNanos.addAndGet(TimeUnit.SECONDS.toNanos(seconds));
|
||||
}
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> {
|
||||
throw new UnsupportedOperationException("no broker: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException("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) {
|
||||
shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
// Do not bind a real port in this assembly test.
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("FakeHerdr is healthy; no poll wait is expected");
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
fleet:
|
||||
leaders:
|
||||
primary:
|
||||
tab: "lead: primary"
|
||||
profile: sonnet
|
||||
profiles:
|
||||
sonnet:
|
||||
subscription: true
|
||||
argv: ["ccs", "sonnet"]
|
||||
""");
|
||||
return FleetConfig.load(file);
|
||||
}
|
||||
|
||||
@Test
|
||||
void realAssembledResolverStillResolvesAfterDeliveredListenerThrows(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
ControllableResourcePorts ports = new ControllableResourcePorts();
|
||||
ports.herdr.withTab("w2", "w2:t7", "lead: primary");
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg,
|
||||
new ConfigRef(dir.resolve("fleetd.yaml"), cfg), new SubscriptionGuard(cfg.guard().hostSet())), ports);
|
||||
try {
|
||||
Injector injector = runtime.injector();
|
||||
installThrowingDeliveredListener(injector);
|
||||
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = new CompletableFuture<>();
|
||||
injector.enqueue(TARGET, "brief", new TurnToken(TARGET, waiter));
|
||||
|
||||
IllegalStateException thrown = assertThrows(IllegalStateException.class,
|
||||
() -> injector.onStatus(TARGET, AgentStatus.IDLE),
|
||||
"control: delivery must reach the installed listener and it must throw after delivery");
|
||||
assertTrue(thrown.getMessage().contains("listener failure"), thrown::getMessage);
|
||||
|
||||
ports.herdr.readText("worker report after the listener failure");
|
||||
ports.advanceSeconds(3);
|
||||
runtime.completion().resolveBeforePostAction(TARGET);
|
||||
|
||||
assertTrue(waiter.isDone(),
|
||||
"FleetdAssembly must wire the Injector registrar to this runtime's real CompletionResolver: "
|
||||
+ "after a delivered listener throws, resolveBeforePostAction must still find and "
|
||||
+ "resolve the registered waiter");
|
||||
} finally {
|
||||
assertTrue(ports.shutdownHook != null, "control: assembly must capture its shutdown hook");
|
||||
ports.shutdownHook.run();
|
||||
assertTrue(ports.herdr.closed, "teardown control: the captured shutdown hook must close herdr");
|
||||
}
|
||||
}
|
||||
|
||||
private static void installThrowingDeliveredListener(Injector injector) throws Exception {
|
||||
Field field = Injector.class.getDeclaredField("turnListener");
|
||||
field.setAccessible(true);
|
||||
field.set(injector, new TurnListener() {
|
||||
@Override
|
||||
public void onTurnComplete(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onDelivered(String target, TurnToken token) {
|
||||
throw new IllegalStateException("listener failure after delivery");
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
-338
@@ -1,338 +0,0 @@
|
||||
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.mcp.FleetMcp;
|
||||
import dev.ltms.fleet.msg.LeadChannel;
|
||||
import dev.ltms.fleet.msg.LeadChannelHandle;
|
||||
import dev.ltms.fleet.msg.LeadMessage;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #612 Shape A, unit r9+r11 — replaces nothing (these three call sites were never covered
|
||||
* by the deleted source-text tests), and pins the three call sites {@code FleetdAssembly}'s
|
||||
* {@code new FleetMcp(...)} wires at {@code :481} ({@link Fleetd#capacitySource}), {@code :482}
|
||||
* ({@link Fleetd#healthCoverageSource}), and {@code :489} ({@code cfg.coordinator() == null ?
|
||||
* List.of() : cfg.coordinator().peers()}).
|
||||
*
|
||||
* <p>{@link FleetdCapacitySourceWiringTest} and {@link FleetdHealthCoverageSourceWiringTest}
|
||||
* already pin the two factory methods themselves in isolation — calling {@code Fleetd.capacitySource}
|
||||
* / {@code Fleetd.healthCoverageSource} directly with a hand-built {@link ConfigRef}/{@link
|
||||
* FleetConfig}. Neither proves that {@code FleetdAssembly}'s {@code new FleetMcp(...)} call
|
||||
* actually receives what those factories return, as opposed to e.g. {@link
|
||||
* FleetMcp.CapacitySource#none()} or a {@link ConfigRef}/{@link FleetConfig} that does not match
|
||||
* what the rest of the daemon was assembled with. The coordinator-peers call site has no isolated
|
||||
* factory at all — it is a plain ternary at the call site itself, so a test can only pin it at the
|
||||
* assembly level.
|
||||
*
|
||||
* <p>This class drives the real {@link FleetdAssembly#assembleAndStart}, then reads the resulting
|
||||
* sources straight off the exact {@link FleetMcp} instance {@link FleetdRuntime#mcp()} owns —
|
||||
* never a copy built alongside it. That needed three small, read-only, public accessors on {@link
|
||||
* FleetMcp} ({@code capacitySource()}, {@code healthCoverageSource()}, {@code coordinatorPeers()}),
|
||||
* added in this same change for the identical cross-package reason {@code quarantineSource()}/
|
||||
* {@code leadSeatSource()}/{@code leadRollover()} already exist (fleetd #612 B3): {@code fleet_list}
|
||||
* itself is reachable only through a real, authorized MCP/HTTP round trip, and driving the real
|
||||
* assembly's {@link dev.ltms.fleet.mcp.ConnectionIdentity} (built with a hardcoded real {@code
|
||||
* LsofPeerPidLookup}) always resolves a same-JVM test caller as {@code ANONYMOUS} — {@code
|
||||
* LsofPeerPidLookup} excludes its own pid, and a JUnit test's HTTP client shares the daemon's JVM
|
||||
* pid. {@code FleetdAssemblyFleetAppTest}'s own javadoc documents this exact dead end for {@code
|
||||
* GET /sessions}; {@code fleet_list} needs the identical {@code Authz.Action.READ} gate. See the
|
||||
* PR body for why this one change to {@code FleetMcp.java} was necessary rather than reshaping
|
||||
* {@code FleetdAssembly} itself.
|
||||
*/
|
||||
class FleetdCapacityHealthCoveragePeersAssemblyTest {
|
||||
|
||||
/** Minimal fake {@link ResourcePorts}: real {@link FakeHerdr}, a sentinel reply inbox, and an
|
||||
* optional fake lead channel for the one test that configures a {@code coordinator:} block. */
|
||||
private static final class RecordingResourcePorts implements ResourcePorts {
|
||||
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
final SentinelReplyInbox replyInbox = new SentinelReplyInbox();
|
||||
/** null unless a test wants the {@code coordinator:} path to actually open. */
|
||||
LeadChannelHandle leadChannel;
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> replyInbox;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
if (leadChannel == null) {
|
||||
throw new UnsupportedOperationException(
|
||||
"leadMailboxOpener must not be called — no coordinator: block is configured");
|
||||
}
|
||||
return leadChannel;
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
}
|
||||
|
||||
@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 static final class SentinelReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
@Override
|
||||
public void own(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void release(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void publish(String target, String msgId, String content) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<InboxMessage> peek(String target) {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean ack(String target, String msgId) {
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
|
||||
/** A fake {@link LeadChannelHandle}: never touches a broker. */
|
||||
private static final class FakeLeadChannel implements LeadChannelHandle {
|
||||
private final String selfCoordId;
|
||||
|
||||
FakeLeadChannel(String selfCoordId) {
|
||||
this.selfCoordId = selfCoordId;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void publish(String toCoordId, LeadMessage m) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<LeadMessage> peek() {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void ack(String msgId) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public String selfCoordId() {
|
||||
return selfCoordId;
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean heldDurable() {
|
||||
return true;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LeadChannel.MailboxState inspect(String coordId) {
|
||||
return LeadChannel.MailboxState.unknown(coordId);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir, String yaml) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml);
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
// --- rank 9: Fleetd.capacitySource, FleetdAssembly.java:481 ---------------------------------
|
||||
|
||||
/**
|
||||
* fleetd #612 rank 9. The inert stand-in {@link FleetMcp.CapacitySource#none()} always reports
|
||||
* an empty {@code configuredProfiles} set, which makes {@code fleet_list} report ZERO
|
||||
* configured profiles — loud, but cheap to pin, and this is the property named in the ticket.
|
||||
*/
|
||||
@Test
|
||||
@DisplayName("[rank 9, BEHAVIOURAL] the real assembled CapacitySource reports the configured "
|
||||
+ "profile with its real maxLoad, not an empty/inert source")
|
||||
void capacitySourceReportsTheConfiguredProfileWithItsRealCapacity(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
profiles:
|
||||
terra:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
model: terra
|
||||
maxLoad: 5
|
||||
guard:
|
||||
offSubscriptionHosts:
|
||||
- gx00.gw
|
||||
""");
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
|
||||
try (FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports)) {
|
||||
FleetMcp.CapacitySource capacity = runtime.mcp().capacitySource();
|
||||
|
||||
assertTrue(capacity.configuredProfiles().get().contains("terra"),
|
||||
"the real assembled CapacitySource must report the configured profile 'terra' — "
|
||||
+ "CapacitySource.none() (the inert stand-in a mis-wired call site could "
|
||||
+ "swap in) always reports an empty set, which would make fleet_list "
|
||||
+ "report ZERO configured profiles");
|
||||
assertEquals(5, capacity.maxLoad().apply("terra"),
|
||||
"the real assembled CapacitySource must report the configured profile's real "
|
||||
+ "maxLoad (5), not null (CapacitySource.none()'s maxLoad always answers "
|
||||
+ "null, whatever the profile)");
|
||||
|
||||
// Loud control, same instance: a profile nobody configured must NOT be reported, so this
|
||||
// assertion is not vacuously true against a source that reports everything.
|
||||
assertFalse(capacity.configuredProfiles().get().contains("no-such-profile"));
|
||||
}
|
||||
}
|
||||
|
||||
// --- rank 11: Fleetd.healthCoverageSource, FleetdAssembly.java:482 --------------------------
|
||||
|
||||
/**
|
||||
* fleetd #612 rank 11. The live daemon currently reports {@code healthCoverage: "detection-only"}
|
||||
* for exactly this shape of config (health enabled, no notifications configured) — a real,
|
||||
* non-default value this test can lean on, per the ticket's own note.
|
||||
*/
|
||||
@Test
|
||||
@DisplayName("[rank 11, BEHAVIOURAL] the real assembled HealthCoverageSource reports the "
|
||||
+ "configured coverage value")
|
||||
void healthCoverageSourceReportsTheConfiguredValue(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
health:
|
||||
enabled: true
|
||||
""");
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
|
||||
try (FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports)) {
|
||||
FleetMcp.HealthCoverageSource healthCoverage = runtime.mcp().healthCoverageSource();
|
||||
|
||||
assertEquals("detection-only", healthCoverage.value().get(),
|
||||
"the real assembled HealthCoverageSource must report 'detection-only' for "
|
||||
+ "health.enabled: true with no notifications configured — a mis-wired "
|
||||
+ "call site (e.g. a HealthCoverageSource built off an unrelated/empty "
|
||||
+ "ConfigRef) would report 'off' here instead, since an absent health: "
|
||||
+ "block also reports 'off'");
|
||||
}
|
||||
}
|
||||
|
||||
// --- rank 11: coordinator peers, FleetdAssembly.java:489 ------------------------------------
|
||||
|
||||
/**
|
||||
* fleetd #612 rank 11, the ternary trap the ticket calls out by name: {@code cfg.coordinator()
|
||||
* == null ? List.of() : cfg.coordinator().peers()}. An ABSENT {@code coordinator:} block and a
|
||||
* MIS-WIRED call site both report {@code List.of()} — indistinguishable unless a test actually
|
||||
* configures a {@code coordinator:} block with peers and checks the real, non-empty answer.
|
||||
*/
|
||||
@Test
|
||||
@DisplayName("[rank 11, BEHAVIOURAL] the real assembled FleetMcp reports coordinator.peers from "
|
||||
+ "a configured coordinator: block, not the no-config empty answer")
|
||||
void coordinatorPeersReportsTheConfiguredPeers(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
coordinator:
|
||||
uri: "amqp://fake-lead-broker/vh"
|
||||
selfId: "test-lead"
|
||||
peers:
|
||||
- "peer-one"
|
||||
- "peer-two"
|
||||
""");
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
ports.leadChannel = new FakeLeadChannel("test-lead");
|
||||
|
||||
try (FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports)) {
|
||||
List<String> peers = runtime.mcp().coordinatorPeers();
|
||||
|
||||
assertEquals(List.of("peer-one", "peer-two"), peers,
|
||||
"the real assembled FleetMcp must report exactly the peers declared under "
|
||||
+ "coordinator.peers, in order");
|
||||
|
||||
// Loud, named control: the trap is that List.of() is ALSO the correct answer with no
|
||||
// coordinator: block configured at all, so an assertion of emptiness here could never
|
||||
// tell a working wiring from the inert/no-config one. This expectation is non-empty,
|
||||
// which only the real wiring (not the ternary's other branch) can produce.
|
||||
assertFalse(peers.isEmpty(),
|
||||
"this test's own expected value must be non-empty, or it could not tell a real "
|
||||
+ "coordinator.peers wiring from cfg.coordinator() == null, which reports "
|
||||
+ "the identical List.of()");
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,213 @@
|
||||
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.mcp.FleetMcp;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.lang.reflect.Field;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
|
||||
/**
|
||||
* fleetd #612 Shape A, unit r5 — {@code FleetdAssembly.java:488} wires {@link
|
||||
* FleetMcp.LeadConfigDirSource} with {@code Fleetd.leadConfigDirSource(() -> config.get().profiles(),
|
||||
* leaders)}. {@link FleetdLeadConfigDirSourceWiringTest} already pins that the FACTORY itself
|
||||
* delegates to the real {@link Fleetd#leadConfigDirLookup} — but, by its own javadoc, it "does not
|
||||
* and structurally cannot cover" whether the real call site in {@code FleetdAssembly} still calls
|
||||
* that factory at all. Measured there: swapping that one-line call for a bare {@code
|
||||
* FleetMcp.LeadConfigDirSource.none()} compiles with 0 errors and leaves the full suite green.
|
||||
*
|
||||
* <p>This is the literal fleetd #602/#606 defect, one call site away from its own fix: {@code main}
|
||||
* (now {@code FleetdAssembly}) used to build {@code LeadConfigDirSource.none()} inline, the whole
|
||||
* suite passed, and the live daemon reported {@code "state":"unknown"} for every lead's context,
|
||||
* forever, with no test noticing. The fix extracted the factory; this test is the one that proves
|
||||
* {@code FleetdAssembly}'s own call site still reaches it.
|
||||
*
|
||||
* <p>This test drives the REAL {@link FleetMcp} the real {@link FleetdAssembly#assembleAndStart}
|
||||
* builds, reached through {@link FleetdRuntime#mcp()}, and reads the {@code leadConfigDirs} field it
|
||||
* was constructed with via reflection — {@code FleetMcp} exposes no public accessor for it (unlike
|
||||
* {@code quarantineSource()}/{@code leadSeatSource()}), so there is no non-reflective route to the
|
||||
* live instance. The assertion resolves a REAL lead name against a REAL configured {@code
|
||||
* configDir:}: {@link FleetMcp.LeadConfigDirSource#none()} (the historical defect, and the
|
||||
* mis-wire this test's mutation cycles reintroduce) always returns {@code null} regardless of the
|
||||
* input, so a non-null, config-matching answer is a property {@code none()} can never produce by
|
||||
* accident.
|
||||
*/
|
||||
class FleetdLeadConfigDirSourceAssemblyTest {
|
||||
|
||||
private static final String LEAD_NAME = "opus";
|
||||
private static final String LEAD_TAB = "lead: opus";
|
||||
private static final String LEAD_PROFILE = "sonnet";
|
||||
|
||||
private static final class RecordingResourcePorts implements ResourcePorts {
|
||||
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
final SentinelReplyInbox replyInbox = new SentinelReplyInbox();
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> replyInbox;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException(
|
||||
"leadMailboxOpener must not be called — no coordinator: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
}
|
||||
|
||||
@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 static final class SentinelReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
@Override
|
||||
public void own(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void release(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void publish(String target, String msgId, String content) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<InboxMessage> peek(String target) {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean ack(String target, String msgId) {
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir, String configDir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
broker:
|
||||
uri: "amqp://fake-test-broker/vh"
|
||||
fleet:
|
||||
leaders:
|
||||
%s:
|
||||
tab: "%s"
|
||||
profile: %s
|
||||
profiles:
|
||||
%s:
|
||||
subscription: true
|
||||
argv: ["ccs", "sonnet"]
|
||||
configDir: "%s"
|
||||
""".formatted(LEAD_NAME, LEAD_TAB, LEAD_PROFILE, LEAD_PROFILE, configDir));
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static FleetMcp.LeadConfigDirSource leadConfigDirSourceOf(FleetMcp mcp) throws Exception {
|
||||
Field field = FleetMcp.class.getDeclaredField("leadConfigDirs");
|
||||
field.setAccessible(true);
|
||||
return (FleetMcp.LeadConfigDirSource) field.get(mcp);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[BEHAVIOURAL] the real assembled LeadConfigDirSource resolves a lead's REAL "
|
||||
+ "configured configDir, not the none() stand-in's hardcoded null")
|
||||
void assembledLeadConfigDirSourceResolvesTheRealConfiguredConfigDir(@TempDir Path dir) throws Exception {
|
||||
String configuredConfigDir = "/mnt/fake-lead-configdir";
|
||||
FleetConfig cfg = writeConfig(dir, configuredConfigDir);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
// Label FakeHerdr's own default pane's tab (term_a / w2:p7 / w2:t7, already carrying a live
|
||||
// agent) to match fleet.leaders.opus.tab exactly, so LeadLauncher.ensureLeads() sees the
|
||||
// lead as already live and does not try to auto-launch a second one.
|
||||
ports.herdr.withTab("w2", "w2:t7", LEAD_TAB);
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
try {
|
||||
FleetMcp.LeadConfigDirSource source = leadConfigDirSourceOf(runtime.mcp());
|
||||
|
||||
assertEquals(configuredConfigDir, source.configDirFor().apply(LEAD_NAME),
|
||||
"fleet.leaders." + LEAD_NAME + ".profile (" + LEAD_PROFILE + ") configures "
|
||||
+ "configDir: " + configuredConfigDir + " — the real assembled source must "
|
||||
+ "resolve it. FleetMcp.LeadConfigDirSource.none() (the inert stand-in "
|
||||
+ "this test's mutation cycles swap the call site for, and the historical "
|
||||
+ "fleetd #602/#606 defect) always reports null here, whatever the input");
|
||||
|
||||
// A lead name the config does not recognise still resolves to null, not a crash — the
|
||||
// same source, applied to an input that must stay at the inert answer even on the real,
|
||||
// non-inert instance.
|
||||
assertNull(source.configDirFor().apply("no-such-lead"));
|
||||
} finally {
|
||||
// Surefire runs the whole suite in one JVM fork (fleetd/pom.xml sets no forkCount /
|
||||
// reuseForks), so the scheduler/loops this assembly starts must be torn down here, on the
|
||||
// failure path too — hence try/finally rather than a bare statement at the end.
|
||||
runtime.close();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,243 @@
|
||||
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.msg.LeadChannelHandle;
|
||||
import dev.ltms.fleet.msg.LeadMessage;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
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.http.HttpRequest;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #612 Shape A ranks 9 and 11: drives the real {@code fleet_list} MCP route through a
|
||||
* token-authenticated primary caller. The assertions read the response from the exact {@link
|
||||
* dev.ltms.fleet.mcp.FleetMcp} instance assembled by {@link FleetdAssembly}, rather than a source
|
||||
* scrape or a separately built reporting source.
|
||||
*
|
||||
* <p>The coordinator fixture contains a peer on purpose. Both a missing coordinator and an
|
||||
* incorrectly wired peers argument can render as an empty list, so the non-empty peer assertion
|
||||
* distinguishes the real call-site value from that inert result. Its fake mailbox never contacts a
|
||||
* broker.
|
||||
*/
|
||||
class FleetdListReportingSourcesAssemblyTest {
|
||||
|
||||
private static final String TOKEN = "fleetd-r9-r11-test-token";
|
||||
private static final String TOKEN_ENV = "FLEETD_R9_TEST_TOKEN";
|
||||
private static final String PROFILE = "capacity-profile";
|
||||
private static final String PEER = "peer-fleet";
|
||||
|
||||
private static final class FakeLeadChannel implements LeadChannelHandle {
|
||||
@Override
|
||||
public void publish(String toCoordId, LeadMessage message) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<LeadMessage> peek() {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void ack(String msgId) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public String selfCoordId() {
|
||||
return "test-lead";
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean heldDurable() {
|
||||
return true;
|
||||
}
|
||||
|
||||
@Override
|
||||
public MailboxState inspect(String coordId) {
|
||||
return MailboxState.unknown(coordId);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
|
||||
private static final class TestResourcePorts implements ResourcePorts {
|
||||
private final FakeHerdr herdr = new FakeHerdr();
|
||||
private Runnable shutdownHook;
|
||||
|
||||
@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() {
|
||||
return (uri, prefetch) -> new ReplyInbox() {
|
||||
@Override public void own(String target) { }
|
||||
@Override public void release(String target) { }
|
||||
@Override public void publish(String target, String msgId, String content) { }
|
||||
@Override public List<InboxMessage> peek(String target) { return List.of(); }
|
||||
@Override public boolean ack(String target, String msgId) { return false; }
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> new FakeLeadChannel();
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
// The test binds runtime.app() itself, below.
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("healthy FakeHerdr must not poll");
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
private FleetdRuntime runtime;
|
||||
private TestResourcePorts 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) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
auth:
|
||||
mode: token
|
||||
tokenEnv: %s
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
health:
|
||||
enabled: true
|
||||
notifications:
|
||||
mode: webhook
|
||||
profiles:
|
||||
%s:
|
||||
baseUrl: http://capacity.test:8000
|
||||
model: test-model
|
||||
maxLoad: 7
|
||||
coordinator:
|
||||
uri: amqp://fake-coordinator/vh
|
||||
selfId: test-lead
|
||||
peers:
|
||||
- %s
|
||||
""".formatted(TOKEN_ENV, PROFILE, PEER));
|
||||
return FleetConfig.load(file);
|
||||
}
|
||||
|
||||
private int assembleAndBind(Path dir) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
ports = new TestResourcePorts();
|
||||
runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config,
|
||||
new SubscriptionGuard(cfg.guard().hostSet())), ports);
|
||||
boundApp = runtime.app().start("127.0.0.1", 0);
|
||||
return boundApp.port();
|
||||
}
|
||||
|
||||
private static String fleetList(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();
|
||||
McpSchema.CallToolResult response = client.callTool(McpSchema.CallToolRequest.builder("fleet_list")
|
||||
.arguments(Map.of()).build());
|
||||
return ((McpSchema.TextContent) response.content().getFirst()).text();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
@Timeout(value = 15, unit = TimeUnit.SECONDS, threadMode = Timeout.ThreadMode.SEPARATE_THREAD)
|
||||
void fleetListReportsTheAssembledCapacitySource(@TempDir Path dir) throws Exception {
|
||||
String response = fleetList(assembleAndBind(dir));
|
||||
|
||||
assertTrue(response.contains("\"capacity\""), response);
|
||||
assertTrue(response.contains("\"" + PROFILE + "\""), response);
|
||||
assertTrue(response.contains("\"maxLoad\":7"), response);
|
||||
}
|
||||
|
||||
@Test
|
||||
@Timeout(value = 15, unit = TimeUnit.SECONDS, threadMode = Timeout.ThreadMode.SEPARATE_THREAD)
|
||||
void fleetListReportsTheAssembledHealthCoverageSource(@TempDir Path dir) throws Exception {
|
||||
String response = fleetList(assembleAndBind(dir));
|
||||
|
||||
assertTrue(response.contains("\"healthCoverage\":\"full\""), response);
|
||||
}
|
||||
|
||||
@Test
|
||||
@Timeout(value = 15, unit = TimeUnit.SECONDS, threadMode = Timeout.ThreadMode.SEPARATE_THREAD)
|
||||
void fleetListReportsTheAssembledCoordinatorPeers(@TempDir Path dir) throws Exception {
|
||||
String response = fleetList(assembleAndBind(dir));
|
||||
|
||||
assertTrue(response.contains("\"coordinator\""), response);
|
||||
assertTrue(response.contains("\"coordId\":\"" + PEER + "\""), response);
|
||||
}
|
||||
}
|
||||
+368
@@ -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);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user