Merge PR #649: fleetd #612 Shape A r9+r11 — pin capacitySource, healthCoverageSource, coordinator peers via a real fleet_list round-trip
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 58s
CI / build (push) Failing after 1m52s

Test-only, 243 lines, one new file. ZERO production change — this is the point
of the PR, and it replaces PR #647, which added three public accessors to
FleetMcp on a false premise.

#647 argued a same-JVM test caller can never reach fleet_list's READ gate,
because FleetdAssembly hardcodes a real LsofPeerPidLookup, so the caller
resolves ANONYMOUS. The premise about lsof is true; the conclusion is not. In
CallerResolver.resolve, once c.terminal() == null the next branch is:

    if (tokenMode) {
        return presentedTokenMatches(authorizationHeader)
                ? Principal.primary(c.pid()) : Principal.anonymous();
    }

c.resolved() and c.scanComplete() guard only the LATER loopback-trust branch,
which token mode returns before reaching. So under auth.mode: token a bearer
token resolves to PRIMARY with no pid lookup on the path. This PR does exactly
that: real McpSyncClient callTool("fleet_list") over a real transport with an
Authorization: Bearer header, a faked leadMailboxOpener so no broker is touched,
and a real coordinator: block with a non-empty peers list. No reflection.

Lead verification, measured myself (not taken from the worker's report):

  git diff --stat origin/main -- fleetd/src/main/java/   -> EMPTY
  mutate :481 Fleetd.capacitySource(...) -> CapacitySource.none()
      -> Tests run: 3, Failures: 1
         fleetListReportsTheAssembledCapacitySource:222
         the other two tests stayed GREEN

The failure message carries the live fleet_list JSON body, showing
healthCoverage, loopHealth and coordinator.peers present with capacity absent —
so the round-trip, the PRIMARY resolution and the coordinator visibility are all
real, and the mutation removed exactly one thing.

Worker also reported, each with a grep -c anchor of 1 restored: health mutation
-> 1 failure named fleetListReportsTheAssembledHealthCoverageSource; peers
mutation -> 1 failure named fleetListReportsTheAssembledCoordinatorPeers; final
mvn clean install Tests run: 1899, Failures: 0 / BUILD SUCCESS. I reproduced the
capacity cycle only; the other two are the worker's measurement, not mine.

Carries forward #647's genuine find: the peers ternary at :489 is
inert-equals-absent, so the test configures a real coordinator: block with peers.

Detail on #612 and #647.
This commit is contained in:
Dai Ha
2026-10-02 04:53:56 +02:00
@@ -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);
}
}