diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdListReportingSourcesAssemblyTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdListReportingSourcesAssemblyTest.java new file mode 100644 index 0000000..168c0b2 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdListReportingSourcesAssemblyTest.java @@ -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. + * + *

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 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 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 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); + } +}