fleetd #612: pin assembled loop health
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 1m1s
CI / build (pull_request) Failing after 2m23s

This commit is contained in:
Dai Ha
2026-10-02 04:08:02 +02:00
parent 141ae3b04d
commit af901ff1d2
@@ -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");
}
}