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:
@@ -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");
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user