Merge PR #644: fleetd #612 Shape A r12 — pin the assembled turn registrar
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 1m32s
CI / build (push) Failing after 1m57s

Pins FleetdAssembly.java:349 behaviourally. The test installs a throwing wrapper on the real assembled Injector's listener to construct the narrow between-delivery-and-completion window, then asserts runtime.completion() still resolves the registered waiter. No sleep: it advances a controllable nano clock. Teardown runs the captured shutdown hook and asserts herdr actually closed.

Lead verification, independent of the implementer's own proof, in a throwaway detached worktree at dac5f88:
- unmutated: Tests run: 1, Failures: 0 — BUILD SUCCESS
- :349 registrar -> (_, _) -> { }: Failures: 1 — "must wire the Injector registrar to this runtime's real CompletionResolver" expected: <true> but was: <false>
- mis-wire I built myself (differs from the implementer's): Fleetd.turnRegistrar(new CompletionResolver(agents, new Rendezvous(), ...)) — a fresh Rendezvous so the registry genuinely differs: Failures: 1, same assertion
- reverted, full suite: Tests run: 1894, Failures: 0, Errors: 0 — BUILD SUCCESS, 58s (1892 baseline + r5 + r12)

Test-only diff, no production code touched. Reflection is used only to install the throwing listener; that is how the failure window is constructed, not how the assertion is made. No reviewer fan-out: member capacity is committed to the remaining Shape A implementers.
This commit was merged in pull request #644.
This commit is contained in:
2026-10-02 04:09:29 +02:00
@@ -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.
*
* <p>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<String, String> 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<Rendezvous.Resolution> 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");
}
});
}
}