fleetd #612: cover fleet list reporting sources
This commit is contained in:
@@ -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);
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user