Pins three call sites the assembly owns and nothing observed: - healthFailTarget (FleetdAssembly:429) — inert, a dead member's waiting ticket sits PENDING for the full 30-minute async timeout instead of failing immediately. - releaseCleanup (:447) — inert, every teardown leaks three things: a stuck rendezvous waiter, an unreleased reply-inbox consumer, and a stale lead binding. - requireOperatorConfirm (:402/:409, fleetd #630) — dropping the 14th constructor argument selects #621's 13-arg overload, which hardcodes true, silently reverting the operator's fix on a host that set requireOperatorConfirm: false. Pinned in both directions, plus an assertion that the two notice strings differ, so no constant satisfies both. Verified by the lead beyond the worker's proof: its releaseCleanup mutation killed all three cleanups at once and so proved only the first assertion had teeth. Starving them one at a time — abandon kept, release starved; then abandon and release kept, forget starved — each fails its own named assertion. All three leaks are pinned independently. A review pass found these three tests assembled real schedulers and never tore them down, by any route: no close(), no shutdownHook, no @AfterEach, no finally. Surefire runs one JVM fork for the whole suite, so those loops outlived their tests. Fixed by capturing the hook and running it in a finally, with an assertion on FakeHerdr.closed so the teardown itself is pinned rather than assumed. MERGE RESOLUTION BY THE LEAD, the same collision as #633. These 3 test files each add a ResourcePorts fake, and #633 landed first making herdrPollWait() abstract with no default. Git reported a clean merge that did not compile. Added the override to all 3, matching the established convention for an always-healthy fake — a Runnable that throws, verified first that none of the three uses healthy(false), so the tripwire can only fire if the test's herdr behaviour changes. Full suite on the resolved merge: 1892 tests, 0 failures (1889 + this branch's 3).
This commit is contained in:
+198
@@ -0,0 +1,198 @@
|
||||
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.health.FleetHealthMonitor;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
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.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.BiConsumer;
|
||||
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 rank 7 — {@code FleetdAssembly.java:429} wires {@link FleetHealthMonitor}'s {@code
|
||||
* failTarget} callback with {@code Fleetd.healthFailTarget(messages)}. {@link
|
||||
* FleetdHealthFailTargetWiringTest} already pins that the FACTORY itself delegates to {@code
|
||||
* messages::abandon}, but it calls {@code Fleetd.healthFailTarget} directly — it never drives {@code
|
||||
* FleetdAssembly.assembleAndStart} and so cannot see whether the real call site at {@code :429}
|
||||
* still passes it the real, assembled {@link MessageService}. Swapping that argument for a no-op
|
||||
* {@code (a, b) -> {}} compiles clean and leaves the whole suite — including the factory-level test
|
||||
* — green: a dead member's waiting ticket then sits {@code PENDING} for the full 30-minute async
|
||||
* timeout instead of failing immediately.
|
||||
*
|
||||
* <p>This test assembles the real daemon with {@code health.enabled: true}, pulls the REAL {@code
|
||||
* failTarget} {@link BiConsumer} out of the REAL, assembled {@link FleetHealthMonitor} (via
|
||||
* reflection — the field is package-private to {@code dev.ltms.fleet.health}, and nothing public
|
||||
* exposes it; {@code StatusPollerResilienceTest} already uses the same technique in this suite), and
|
||||
* invokes it directly against the REAL {@link MessageService} {@link FleetdRuntime#messages()}
|
||||
* returns. A no-op lambda swapped in at the call site leaves the ticket {@code PENDING} forever,
|
||||
* which this test catches; the real one fails it.
|
||||
*/
|
||||
class FleetdAssemblyHealthFailTargetBehaviouralTest {
|
||||
|
||||
private static final String TARGET = "term_a";
|
||||
|
||||
private static final class RecordingResourcePorts implements ResourcePorts {
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
Runnable shutdownHook;
|
||||
|
||||
@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 Runnable herdrPollWait() {
|
||||
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
|
||||
};
|
||||
}
|
||||
@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) {
|
||||
this.shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
health:
|
||||
enabled: true
|
||||
intervalSeconds: 30
|
||||
""");
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@Test
|
||||
@DisplayName("[BEHAVIOURAL] the real assembled FleetHealthMonitor's failTarget reaches the real "
|
||||
+ "MessageService.abandon, not a no-op")
|
||||
void assembledHealthFailTargetReachesRealMessages(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
assertNotNull(ports.shutdownHook, "FleetdAssembly must have registered a shutdown hook");
|
||||
|
||||
// Surefire runs the whole suite in one JVM fork, so the scheduler/loops this assembly starts
|
||||
// (SessionReaper, StatusPoller, the health monitor) must be torn down here, on the failure
|
||||
// path too — hence the try/finally, not just a statement at the end of the happy path.
|
||||
try {
|
||||
FleetHealthMonitor healthMonitor = runtime.healthMonitor();
|
||||
assertNotNull(healthMonitor, "health.enabled: true in this test's config, so "
|
||||
+ "FleetdAssembly.assembleAndStart must have built a real FleetHealthMonitor");
|
||||
|
||||
Field field = FleetHealthMonitor.class.getDeclaredField("failTarget");
|
||||
field.setAccessible(true);
|
||||
BiConsumer<String, String> failTarget = (BiConsumer<String, String>) field.get(healthMonitor);
|
||||
assertNotNull(failTarget, "FleetHealthMonitor's failTarget must never be null — the "
|
||||
+ "constructor itself requires it");
|
||||
|
||||
MessageService messages = runtime.messages();
|
||||
|
||||
// --- loud control: prove the assembled MessageService is actually wired up and a ticket is
|
||||
// genuinely PENDING before failTarget ever runs. If this fails, the test below would pass
|
||||
// vacuously on a MessageService that never got a ticket in the first place. TARGET has no
|
||||
// live agent behind it (no session was ever acquired), so nothing resolves this ticket on
|
||||
// its own — it stays PENDING until failTarget (or a timeout) ends it.
|
||||
String ticket = messages.sendAsync(TARGET, "long task");
|
||||
MessageService.TaskView before = messages.poll(ticket);
|
||||
assertEquals(MessageService.Phase.PENDING, before.phase(),
|
||||
"control: the async ticket must be PENDING before failTarget runs");
|
||||
|
||||
failTarget.accept(TARGET, "member unreachable (health monitor)");
|
||||
|
||||
MessageService.TaskView after = awaitTerminal(messages, ticket);
|
||||
assertEquals(MessageService.Phase.FAILED, after.phase(),
|
||||
"FleetdAssembly.java:429 must pass Fleetd.healthFailTarget(messages) built from the "
|
||||
+ "SAME assembled MessageService — a no-op BiConsumer at that call site leaves "
|
||||
+ "this ticket PENDING for the full 30-minute async timeout instead of failing it");
|
||||
assertTrue(after.detail() != null && after.detail().contains("member unreachable"),
|
||||
"the failure reason passed to failTarget.accept must reach MessageService.abandon and "
|
||||
+ "end up in the ticket's detail");
|
||||
} finally {
|
||||
// Proof the teardown actually ran, not just an assurance that a finally was added: the
|
||||
// captured shutdown hook's close order (FleetdAssemblyLifecycleTest) closes the herdr
|
||||
// client last, so ports.herdr.closed flips to true only if this hook really executed.
|
||||
ports.shutdownHook.run();
|
||||
assertTrue(ports.herdr.closed, "the captured shutdown hook must have run and closed herdr — "
|
||||
+ "proof this test's assembled background loops/scheduler were torn down");
|
||||
}
|
||||
}
|
||||
|
||||
private static MessageService.TaskView awaitTerminal(MessageService messages, String ticket)
|
||||
throws InterruptedException {
|
||||
long deadline = System.currentTimeMillis() + 5000;
|
||||
MessageService.TaskView view = messages.poll(ticket);
|
||||
while (view.phase() == MessageService.Phase.PENDING && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(10);
|
||||
view = messages.poll(ticket);
|
||||
}
|
||||
return view;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,232 @@
|
||||
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.mcp.PrimaryRegistry;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import dev.ltms.fleet.msg.ReplyPushLoop;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
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.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.Consumer;
|
||||
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 rank 6 — {@code FleetdAssembly.java:447} wires {@code
|
||||
* sessions.onRelease(Fleetd.releaseCleanup(messages, replyInbox, primaryRegistry))}. {@link
|
||||
* FleetdReleaseCleanupWiringTest} already pins that the FACTORY {@code Fleetd.releaseCleanup}
|
||||
* itself reaches all three collaborators — but it calls the factory directly, never {@code
|
||||
* FleetdAssembly.assembleAndStart}, so it cannot see whether the real call site at {@code :447}
|
||||
* still registers it (as opposed to a no-op {@code detail -> { }}) or still passes it the REAL,
|
||||
* assembled {@code messages}/{@code replyInbox}/{@code primaryRegistry}. Swapping the registered
|
||||
* listener for a no-op at that call site compiles clean and leaves the whole suite — including the
|
||||
* factory-level test — green: EVERY teardown then leaks a stuck rendezvous waiter, an unreleased
|
||||
* reply-inbox consumer, and a stale lead binding, all three at once.
|
||||
*
|
||||
* <p>This test assembles the real daemon with {@code idleSleepGuard.enabled: false} — the ONLY
|
||||
* other {@code onRelease} registration in {@code FleetdAssembly} (see {@code
|
||||
* dev.ltms.fleet.power.IdleSleepGuard}'s own wiring at {@code FleetdAssembly.java:228}) — so the
|
||||
* real {@link SessionManager}'s release-listener list holds exactly the one listener this call site
|
||||
* registers. It pulls that REAL listener out via reflection (the list itself is private, like
|
||||
* {@code StatusPollerResilienceTest}'s use of the same technique elsewhere in this suite), invokes
|
||||
* it directly, and asserts all three collaborator effects against the REAL, assembled {@link
|
||||
* MessageService} ({@link FleetdRuntime#messages()}), the REAL {@link ReplyInbox} ({@link
|
||||
* FleetdRuntime#replyInbox()}), and the REAL {@link PrimaryRegistry} — reached through {@link
|
||||
* FleetdRuntime#pushLoop()}, the only other accessor that was handed the same {@code
|
||||
* primaryRegistry} instance ({@code FleetdAssembly.java:382}), since {@code FleetMcp} never exposes
|
||||
* it.
|
||||
*/
|
||||
class FleetdAssemblyReleaseCleanupBehaviouralTest {
|
||||
|
||||
private static final String TARGET = "term_a";
|
||||
|
||||
private static final class RecordingResourcePorts implements ResourcePorts {
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
Runnable shutdownHook;
|
||||
|
||||
@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 Runnable herdrPollWait() {
|
||||
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
|
||||
};
|
||||
}
|
||||
@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) {
|
||||
this.shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
""");
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@Test
|
||||
@DisplayName("[BEHAVIOURAL] the real assembled release listener reaches messages.abandon, "
|
||||
+ "replyInbox.release, AND primaryRegistry.forgetDelegation — all three leaks at once")
|
||||
void assembledReleaseListenerReachesAllThreeCollaborators(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
assertNotNull(ports.shutdownHook, "FleetdAssembly must have registered a shutdown hook");
|
||||
|
||||
// Surefire runs the whole suite in one JVM fork, so the scheduler/loops this assembly starts
|
||||
// must be torn down here, on the failure path too — hence the try/finally, not just a
|
||||
// statement at the end of the happy path.
|
||||
try {
|
||||
// --- reach into SessionManager's private release-listener list. idleSleepGuard.enabled:
|
||||
// false above means FleetdAssembly.java:228 never registers, so this list must hold EXACTLY
|
||||
// the one listener :447 registers.
|
||||
Field listenersField = SessionManager.class.getDeclaredField("releaseListeners");
|
||||
listenersField.setAccessible(true);
|
||||
List<Consumer<SessionManager.ReleaseDetail>> releaseListeners =
|
||||
(List<Consumer<SessionManager.ReleaseDetail>>) listenersField.get(runtime.sessions());
|
||||
assertEquals(1, releaseListeners.size(), "control: with idleSleepGuard.enabled: false, "
|
||||
+ "FleetdAssembly.java:447 must be the ONLY onRelease registration — a different "
|
||||
+ "count means this test is no longer isolating the call site it claims to pin");
|
||||
Consumer<SessionManager.ReleaseDetail> releaseListener = releaseListeners.get(0);
|
||||
|
||||
MessageService messages = runtime.messages();
|
||||
ReplyInbox replyInbox = runtime.replyInbox();
|
||||
|
||||
// primaryRegistry is never exposed by FleetdRuntime directly — ReplyPushLoop is the other
|
||||
// collaborator FleetdAssembly.java:382 hands the SAME instance to, so reach it from there.
|
||||
Field primaryRegistryField = ReplyPushLoop.class.getDeclaredField("primaryRegistry");
|
||||
primaryRegistryField.setAccessible(true);
|
||||
PrimaryRegistry primaryRegistry = (PrimaryRegistry) primaryRegistryField.get(runtime.pushLoop());
|
||||
assertNotNull(primaryRegistry, "control: the assembled ReplyPushLoop must hold a real "
|
||||
+ "PrimaryRegistry instance");
|
||||
|
||||
// --- loud controls: set up the "before" state each collaborator's effect is measured
|
||||
// against, against the REAL assembled objects. If any of these three fails, the test below
|
||||
// would pass vacuously because the subject it claims to observe never existed in the first
|
||||
// place.
|
||||
String ticket = messages.sendAsync(TARGET, "long task");
|
||||
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase(),
|
||||
"control: the async ticket must be PENDING before the release listener runs");
|
||||
|
||||
replyInbox.own(TARGET);
|
||||
replyInbox.publish(TARGET, "msg-1", "hello");
|
||||
assertEquals(1, replyInbox.peek(TARGET).size(),
|
||||
"control: the reply inbox must own TARGET and hold one message before the release "
|
||||
+ "listener runs");
|
||||
|
||||
primaryRegistry.recordDelegation(TARGET, "lead-1");
|
||||
assertEquals("lead-1", primaryRegistry.nudgeTargetFor(TARGET).orElse(null),
|
||||
"control: the delegation must be recorded before the release listener runs");
|
||||
|
||||
// --- the one call under test: invoke the REAL, assembled release listener directly, the
|
||||
// same way SessionManager.release(...) would on a real teardown.
|
||||
releaseListener.accept(new SessionManager.ReleaseDetail(TARGET, null, null, null, null));
|
||||
|
||||
MessageService.TaskView after = awaitTerminal(messages, ticket);
|
||||
assertEquals(MessageService.Phase.FAILED, after.phase(),
|
||||
"FleetdAssembly.java:447 must register a listener that calls messages.abandon(...) "
|
||||
+ "on the SAME assembled MessageService — an inert listener leaves this "
|
||||
+ "ticket PENDING for the full 30-minute async timeout");
|
||||
assertTrue(after.detail() != null && after.detail().contains("released"),
|
||||
"the abandon reason must say the worker session was released");
|
||||
|
||||
assertTrue(replyInbox.peek(TARGET).isEmpty(),
|
||||
"FleetdAssembly.java:447 must register a listener that calls replyInbox.release(...) "
|
||||
+ "— an inert listener leaves the inbox still owning TARGET with its message");
|
||||
|
||||
assertTrue(primaryRegistry.nudgeTargetFor(TARGET).isEmpty(),
|
||||
"FleetdAssembly.java:447 must register a listener that calls "
|
||||
+ "primaryRegistry.forgetDelegation(...) — an inert listener leaves the stale "
|
||||
+ "delegation in place");
|
||||
} finally {
|
||||
// Proof the teardown actually ran, not just an assurance that a finally was added: the
|
||||
// captured shutdown hook's close order (FleetdAssemblyLifecycleTest) closes the herdr
|
||||
// client last, so ports.herdr.closed flips to true only if this hook really executed.
|
||||
ports.shutdownHook.run();
|
||||
assertTrue(ports.herdr.closed, "the captured shutdown hook must have run and closed herdr — "
|
||||
+ "proof this test's assembled background loops/scheduler were torn down");
|
||||
}
|
||||
}
|
||||
|
||||
private static MessageService.TaskView awaitTerminal(MessageService messages, String ticket)
|
||||
throws InterruptedException {
|
||||
long deadline = System.currentTimeMillis() + 5000;
|
||||
MessageService.TaskView view = messages.poll(ticket);
|
||||
while (view.phase() == MessageService.Phase.PENDING && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(10);
|
||||
view = messages.poll(ticket);
|
||||
}
|
||||
return view;
|
||||
}
|
||||
}
|
||||
+241
@@ -0,0 +1,241 @@
|
||||
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.lead.LeadContextGauge;
|
||||
import dev.ltms.fleet.msg.LeadHeartbeatLoop;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.lang.reflect.Field;
|
||||
import java.lang.reflect.Method;
|
||||
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.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #630, and fleetd #612's ranks for this call site. {@code FleetdAssembly.java:402}
|
||||
* computes {@code requireOperatorConfirm} from the effective {@code leadRollover.requireOperatorConfirm}
|
||||
* config, and {@code :409} threads it as the 14th argument into the full {@link LeadHeartbeatLoop}
|
||||
* constructor. Measured on 26f1986: dropping that one argument so the 13-argument overload is
|
||||
* selected instead (it delegates with {@code true} hardcoded — see that overload's own javadoc,
|
||||
* fleetd #621) compiles with 0 errors and leaves all 1883 tests green, both with and without the
|
||||
* argument. In production this means the daemon keeps starting and keeps nudging, but the
|
||||
* context-high notice silently goes back to telling EVERY lead to ask the operator before a
|
||||
* context roll — on a host that set {@code requireOperatorConfirm: false} specifically so it would
|
||||
* not have to. That is the operator's own fix silently reverting, with a fully green suite.
|
||||
*
|
||||
* <p>{@code LeadHeartbeatLoopTest} already proves {@link LeadHeartbeatLoop}'s package-private
|
||||
* {@code contextNotice(boolean, LeadContextGauge.Reading, boolean, boolean)} branches correctly on
|
||||
* its own {@code requireOperatorConfirm} argument — that the METHOD works. It says nothing about
|
||||
* which value {@code FleetdAssembly} actually passes into the constructed loop, so it is not
|
||||
* reused here as coverage for the call site.
|
||||
*
|
||||
* <p>This test assembles the real daemon TWICE — once with {@code leadRollover.requireOperatorConfirm:
|
||||
* false}, once with {@code true} — pulls the REAL {@code requireOperatorConfirm} field out of the
|
||||
* REAL, assembled {@link LeadHeartbeatLoop} each time (reflection: the field, and {@code
|
||||
* contextNotice} itself, are package-private to {@code dev.ltms.fleet.msg}, and nothing public
|
||||
* exposes either — the same technique {@code StatusPollerResilienceTest} already uses in this
|
||||
* suite), and calls the REAL {@code contextNotice} method with that field's value to produce the
|
||||
* actual notice text the assembled loop would append to a nudge. Both directions are asserted: a
|
||||
* one-directional test here would pass on a constant.
|
||||
*/
|
||||
class FleetdAssemblyRequireOperatorConfirmBehaviouralTest {
|
||||
|
||||
private static final class RecordingResourcePorts implements ResourcePorts {
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
Runnable shutdownHook;
|
||||
|
||||
@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 Runnable herdrPollWait() {
|
||||
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
|
||||
};
|
||||
}
|
||||
@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) {
|
||||
this.shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir, boolean requireOperatorConfirm) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
leadHeartbeat:
|
||||
idleAfterSeconds: 600
|
||||
backoffMs: 15000
|
||||
quietNudgeCap: 5
|
||||
leadRollover:
|
||||
handoverPath: handover.md
|
||||
requireOperatorConfirm: %s
|
||||
""".formatted(requireOperatorConfirm));
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
/** Carries both the assembled loop under test AND its {@link RecordingResourcePorts}, so the
|
||||
* caller can tear the assembly down (this test assembles the real daemon TWICE — see the class
|
||||
* javadoc — and each assembly needs its own teardown, not just the last one). */
|
||||
private record Assembled(LeadHeartbeatLoop heartbeat, RecordingResourcePorts ports) {
|
||||
}
|
||||
|
||||
private static Assembled assembleHeartbeat(Path dir, boolean requireOperatorConfirm) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir, requireOperatorConfirm);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
assertNotNull(ports.shutdownHook, "FleetdAssembly must have registered a shutdown hook");
|
||||
|
||||
LeadHeartbeatLoop heartbeat = runtime.heartbeat();
|
||||
assertNotNull(heartbeat, "control: leadHeartbeat: is configured, so FleetdAssembly.assembleAndStart "
|
||||
+ "must have built a real LeadHeartbeatLoop");
|
||||
return new Assembled(heartbeat, ports);
|
||||
}
|
||||
|
||||
/** Pulls the REAL {@code requireOperatorConfirm} field off the REAL, assembled loop. */
|
||||
private static boolean assembledRequireOperatorConfirm(LeadHeartbeatLoop heartbeat) throws Exception {
|
||||
Field field = LeadHeartbeatLoop.class.getDeclaredField("requireOperatorConfirm");
|
||||
field.setAccessible(true);
|
||||
return field.getBoolean(heartbeat);
|
||||
}
|
||||
|
||||
/** Calls the REAL, package-private {@code contextNotice(boolean, Reading, boolean, boolean)} via reflection. */
|
||||
private static String contextNotice(boolean enabled, LeadContextGauge.Reading reading, boolean alreadyNotified,
|
||||
boolean requireOperatorConfirm) throws Exception {
|
||||
Method method = LeadHeartbeatLoop.class.getDeclaredMethod("contextNotice", boolean.class,
|
||||
LeadContextGauge.Reading.class, boolean.class, boolean.class);
|
||||
method.setAccessible(true);
|
||||
return (String) method.invoke(null, enabled, reading, alreadyNotified, requireOperatorConfirm);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[BEHAVIOURAL] the real assembled LeadHeartbeatLoop's context-high notice tracks "
|
||||
+ "leadRollover.requireOperatorConfirm — BOTH directions")
|
||||
void assembledRequireOperatorConfirmControlsNoticeWording(@TempDir Path dir) throws Exception {
|
||||
LeadContextGauge.Reading highReading = new LeadContextGauge.Reading(LeadContextGauge.State.HIGH,
|
||||
250_000L, 2);
|
||||
|
||||
String noticeFalse;
|
||||
String noticeTrue;
|
||||
|
||||
// --- direction 1: requireOperatorConfirm: false -----------------------------------------
|
||||
Path falseDir = dir.resolve("false");
|
||||
Files.createDirectories(falseDir);
|
||||
Assembled assembledFalse = assembleHeartbeat(falseDir, false);
|
||||
// Surefire runs the whole suite in one JVM fork, so each assembly's scheduler/loops must be
|
||||
// torn down here, on the failure path too — hence try/finally per assembly (this test
|
||||
// assembles TWICE, so both need their own teardown, not just the last one).
|
||||
try {
|
||||
boolean fieldFalse = assembledRequireOperatorConfirm(assembledFalse.heartbeat());
|
||||
assertFalse(fieldFalse, "FleetdAssembly.java:402/:409 must thread leadRollover."
|
||||
+ "requireOperatorConfirm: false into the assembled LeadHeartbeatLoop's own field — "
|
||||
+ "dropping the 14th constructor argument selects the 13-argument overload, which "
|
||||
+ "hardcodes true regardless of config (fleetd #621), and this would read true instead");
|
||||
|
||||
noticeFalse = contextNotice(true, highReading, false, fieldFalse);
|
||||
assertTrue(noticeFalse.contains("Decide for yourself when to confirm"),
|
||||
"with requireOperatorConfirm: false, the assembled loop's own notice must tell the "
|
||||
+ "lead it can decide for itself — got: " + noticeFalse);
|
||||
assertFalse(noticeFalse.contains("ask the operator") || noticeFalse.contains("Only the operator"),
|
||||
"with requireOperatorConfirm: false, the assembled loop's own notice must NOT ask the "
|
||||
+ "operator — got: " + noticeFalse);
|
||||
} finally {
|
||||
// Proof the teardown actually ran, not just an assurance that a finally was added: the
|
||||
// captured shutdown hook's close order (FleetdAssemblyLifecycleTest) closes the herdr
|
||||
// client last, so ports.herdr.closed flips to true only if this hook really executed.
|
||||
assembledFalse.ports().shutdownHook.run();
|
||||
assertTrue(assembledFalse.ports().herdr.closed, "the captured shutdown hook must have run "
|
||||
+ "and closed herdr — proof this assembly's background loops/scheduler were torn down");
|
||||
}
|
||||
|
||||
// --- direction 2: requireOperatorConfirm: true -------------------------------------------
|
||||
Path trueDir = dir.resolve("true");
|
||||
Files.createDirectories(trueDir);
|
||||
Assembled assembledTrue = assembleHeartbeat(trueDir, true);
|
||||
try {
|
||||
boolean fieldTrue = assembledRequireOperatorConfirm(assembledTrue.heartbeat());
|
||||
assertTrue(fieldTrue, "FleetdAssembly.java:402/:409 must thread leadRollover."
|
||||
+ "requireOperatorConfirm: true into the assembled LeadHeartbeatLoop's own field");
|
||||
|
||||
noticeTrue = contextNotice(true, highReading, false, fieldTrue);
|
||||
assertTrue(noticeTrue.contains("ask the operator") && noticeTrue.contains("Only the operator can approve the roll"),
|
||||
"with requireOperatorConfirm: true, the assembled loop's own notice must ask the "
|
||||
+ "operator — got: " + noticeTrue);
|
||||
assertFalse(noticeTrue.contains("Decide for yourself when to confirm"),
|
||||
"with requireOperatorConfirm: true, the assembled loop's own notice must NOT tell "
|
||||
+ "the lead it can decide for itself — got: " + noticeTrue);
|
||||
} finally {
|
||||
assembledTrue.ports().shutdownHook.run();
|
||||
assertTrue(assembledTrue.ports().herdr.closed, "the captured shutdown hook must have run "
|
||||
+ "and closed herdr — proof this assembly's background loops/scheduler were torn down");
|
||||
}
|
||||
|
||||
// --- the two directions must actually differ: a constant return would pass both assertion
|
||||
// blocks above vacuously if they happened to share wording, so compare them directly too.
|
||||
assertTrue(!noticeFalse.equals(noticeTrue),
|
||||
"the two directions must produce genuinely different notice text — got the same "
|
||||
+ "text for both: " + noticeFalse);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user