fleetd #612: pin assembled turn registrar
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