Compare commits
13 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| a42b12440c | |||
| a06426c33c | |||
| 28ea0de575 | |||
| 91792e11fc | |||
| 6539efe9fa | |||
| 68397f78d5 | |||
| af901ff1d2 | |||
| dac5f88812 | |||
| 2e663e5968 | |||
| 8c14ed2846 | |||
| 141ae3b04d | |||
| faefea14c4 | |||
| ea6896f2ef |
@@ -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");
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -159,7 +159,13 @@ done
|
|||||||
# The fix is structural, not another name to match: once a key line is masked, every following
|
# The fix is structural, not another name to match: once a key line is masked, every following
|
||||||
# line indented STRICTLY DEEPER than that key is masked too, by indentation alone, until the
|
# line indented STRICTLY DEEPER than that key is masked too, by indentation alone, until the
|
||||||
# indentation returns to the key's own level or shallower. This needs no knowledge of the key's
|
# indentation returns to the key's own level or shallower. This needs no knowledge of the key's
|
||||||
# name and so protects a block scalar under any masked key, present or future.
|
# name, so it covers a block scalar under any masked key — but ONLY while that key's own line is
|
||||||
|
# itself inside the hunk being printed. `diff -u` prints just three lines of context, so a block
|
||||||
|
# scalar's body often reaches this function with its key line left out; there is then nothing to
|
||||||
|
# anchor to, `masked` is never set, and the body prints in full. A blank line inside a block
|
||||||
|
# scalar loses the anchor the same way, because a blank diff line measures as indent 0. Both are
|
||||||
|
# measured and filed as fleetd #639 — do not read this paragraph as a guarantee that a masked
|
||||||
|
# key's value can never be printed.
|
||||||
#
|
#
|
||||||
# `redact` is always fed `diff -u` output, and every line of a unified diff starts with exactly
|
# `redact` is always fed `diff -u` output, and every line of a unified diff starts with exactly
|
||||||
# one of ' ', '+', '-' (the three body markers; '@'/'-'/'+' for the three header-line kinds too).
|
# one of ' ', '+', '-' (the three body markers; '@'/'-'/'+' for the three header-line kinds too).
|
||||||
|
|||||||
Reference in New Issue
Block a user