From af901ff1d2975e91d89d2a0af0e53c309d9c76fc Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Fri, 2 Oct 2026 04:08:02 +0200 Subject: [PATCH] fleetd #612: pin assembled loop health --- .../fleet/FleetdAssemblyLoopHealthTest.java | 180 ++++++++++++++++++ 1 file changed, 180 insertions(+) create mode 100644 fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyLoopHealthTest.java diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyLoopHealthTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyLoopHealthTest.java new file mode 100644 index 0000000..2cf457a --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyLoopHealthTest.java @@ -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 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 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"); + } +}