From 2e663e5968d0726492d50b05601ef2f7ff07577f Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Fri, 2 Oct 2026 04:04:54 +0200 Subject: [PATCH] fleetd #612: pin assembled turn registrar --- ...dAssemblyTurnRegistrarBehaviouralTest.java | 184 ++++++++++++++++++ 1 file changed, 184 insertions(+) create mode 100644 fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyTurnRegistrarBehaviouralTest.java diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyTurnRegistrarBehaviouralTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyTurnRegistrarBehaviouralTest.java new file mode 100644 index 0000000..3372ed4 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyTurnRegistrarBehaviouralTest.java @@ -0,0 +1,184 @@ +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.AgentStatus; +import dev.ltms.fleet.herdr.FakeHerdr; +import dev.ltms.fleet.herdr.HerdrClient; +import dev.ltms.fleet.inject.Injector; +import dev.ltms.fleet.inject.TurnListener; +import dev.ltms.fleet.msg.Rendezvous; +import dev.ltms.fleet.msg.TurnToken; +import io.javalin.Javalin; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.lang.reflect.Field; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; +import java.util.function.LongSupplier; + +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * fleetd #612 rank 12 — {@code FleetdAssembly} supplies the {@link Injector}'s {@code TurnRegistrar} + * with {@code Fleetd.turnRegistrar(completion)}. The normal delivery path cannot distinguish that + * registrar from {@code TurnRegistrar.NOOP}: {@link dev.ltms.fleet.inject.CompletionResolver#onDelivered} + * registers the same turn shortly afterwards. The distinction matters when a delivery listener throws + * after the pane received the message but before normal completion runs. The registrar must already have + * registered the waiter with the REAL assembled resolver, so a later completion can still resolve it. + * + *

The test installs a throwing wrapper around the real assembled listener. It does not replace the + * registrar or resolver. This constructs the narrow failure condition without a sleep, then observes the + * real resolver through {@link FleetdRuntime#completion()}. An inert registrar, or a registrar wired to a + * throwaway resolver, leaves this waiter's turn absent from the real resolver and makes the final assertion + * fail. + */ +class FleetdAssemblyTurnRegistrarBehaviouralTest { + + private static final String TARGET = "term_a"; + + private static final class ControllableResourcePorts implements ResourcePorts { + final FakeHerdr herdr = new FakeHerdr(); + final AtomicLong nowNanos = new AtomicLong(1_000_000_000L); + Runnable shutdownHook; + + void advanceSeconds(long seconds) { + nowNanos.addAndGet(TimeUnit.SECONDS.toNanos(seconds)); + } + + @Override + public Map environment() { + return Map.of(); + } + + @Override + public HerdrClient connectHerdr(Path socketPath) { + return herdr; + } + + @Override + public Fleetd.AmqpOpener replyInboxOpener() { + return (uri, prefetch) -> { + throw new UnsupportedOperationException("no broker: block is configured"); + }; + } + + @Override + public Fleetd.LeadMailboxOpener leadMailboxOpener() { + return (uri, selfCoordId, prefetch) -> { + throw new UnsupportedOperationException("no coordinator: block is configured"); + }; + } + + @Override + public LongSupplier nanoClock() { + return nowNanos::get; + } + + @Override + public LongSupplier wallClockNanos() { + return nowNanos::get; + } + + @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) { + // Do not bind a real port in this assembly test. + } + + @Override + public Runnable herdrPollWait() { + return () -> { + throw new UnsupportedOperationException("FakeHerdr is healthy; no poll wait is expected"); + }; + } + } + + 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 + idleSleepGuard: + enabled: false + fleet: + leaders: + primary: + tab: "lead: primary" + profile: sonnet + profiles: + sonnet: + subscription: true + argv: ["ccs", "sonnet"] + """); + return FleetConfig.load(file); + } + + @Test + void realAssembledResolverStillResolvesAfterDeliveredListenerThrows(@TempDir Path dir) throws Exception { + FleetConfig cfg = writeConfig(dir); + ControllableResourcePorts ports = new ControllableResourcePorts(); + ports.herdr.withTab("w2", "w2:t7", "lead: primary"); + FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, + new ConfigRef(dir.resolve("fleetd.yaml"), cfg), new SubscriptionGuard(cfg.guard().hostSet())), ports); + try { + Injector injector = runtime.injector(); + installThrowingDeliveredListener(injector); + + CompletableFuture waiter = new CompletableFuture<>(); + injector.enqueue(TARGET, "brief", new TurnToken(TARGET, waiter)); + + IllegalStateException thrown = assertThrows(IllegalStateException.class, + () -> injector.onStatus(TARGET, AgentStatus.IDLE), + "control: delivery must reach the installed listener and it must throw after delivery"); + assertTrue(thrown.getMessage().contains("listener failure"), thrown::getMessage); + + ports.herdr.readText("worker report after the listener failure"); + ports.advanceSeconds(3); + runtime.completion().resolveBeforePostAction(TARGET); + + assertTrue(waiter.isDone(), + "FleetdAssembly must wire the Injector registrar to this runtime's real CompletionResolver: " + + "after a delivered listener throws, resolveBeforePostAction must still find and " + + "resolve the registered waiter"); + } finally { + assertTrue(ports.shutdownHook != null, "control: assembly must capture its shutdown hook"); + ports.shutdownHook.run(); + assertTrue(ports.herdr.closed, "teardown control: the captured shutdown hook must close herdr"); + } + } + + private static void installThrowingDeliveredListener(Injector injector) throws Exception { + Field field = Injector.class.getDeclaredField("turnListener"); + field.setAccessible(true); + field.set(injector, new TurnListener() { + @Override + public void onTurnComplete(String target) { + } + + @Override + public void onDelivered(String target, TurnToken token) { + throw new IllegalStateException("listener failure after delivery"); + } + }); + } +}