From aa517ae0ec92f711871cbd9b0e97e06faa8bb69a Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 20 Sep 2026 17:18:45 +0700 Subject: [PATCH] fleetd #608: make anAlreadyCollectedTicketProducesNoNudge deterministic Replace the real ScheduledExecutorService backing ReplyPushLoop in this one test with ManualScheduler, a fake that only runs a tick when the test calls runDueTasks(). The old test bet a 300ms backoff was wide enough that collecting the ticket always won the race against the scheduler's own timer - true on an idle machine, false under a loaded full-suite run, which is exactly the flake reported. The rewritten test also waits on setAfterFinishAsyncTaskCompleteHookForTest (already used elsewhere in this file) instead of polling Phase.DONE, so it does not race CompletableFuture.complete()'s own publish-then-run-dependents gap (fleetd #399) while proving ReplyPushLoop.onTicketTerminal really ran before the ticket is collected. Verified: backoff=1 (the most hostile value) still passes; mutating ReplyPushLoop.ticketCollected to a no-op turns the test red with the same assertion message the original flake reported; three consecutive full-suite runs are green (1841/1841 each). --- .../dev/ltms/fleet/msg/ManualScheduler.java | 168 ++++++++++++++++++ .../ltms/fleet/msg/MessageServiceTest.java | 72 +++++++- 2 files changed, 234 insertions(+), 6 deletions(-) create mode 100644 fleetd/src/test/java/dev/ltms/fleet/msg/ManualScheduler.java diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/ManualScheduler.java b/fleetd/src/test/java/dev/ltms/fleet/msg/ManualScheduler.java new file mode 100644 index 0000000..08ae6c5 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/ManualScheduler.java @@ -0,0 +1,168 @@ +package dev.ltms.fleet.msg; + +import java.util.ArrayDeque; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Deque; +import java.util.List; +import java.util.concurrent.Callable; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Future; +import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; + +/** + * A {@link ScheduledExecutorService} for tests that never runs a task on a timer — it records what + * {@link ReplyPushLoop} schedules and only ever runs it when the test itself calls + * {@link #runDueTasks()} (fleetd #608). + * + *

The test this exists for ({@code anAlreadyCollectedTicketProducesNoNudge}) used to wire + * {@link ReplyPushLoop} to a real {@code Executors.newSingleThreadScheduledExecutor()} and bet a + * 300ms backoff was "wide" enough that the test's own work (collecting the ticket) always won the + * race against the scheduler's own timer firing the next tick. On an idle machine that held; under + * a loaded full-suite run it did not, and the test went red on perfectly correct code. Replacing + * the timer with this fake removes the race rather than widening it: nothing here ever runs on a + * schedule of its own, so no backoff value — however small — can make the tick fire before the test + * is ready for it. + * + *

Deliberately narrow. {@link ReplyPushLoop} calls exactly two methods on its + * {@link ScheduledExecutorService} field — {@link #schedule(Runnable, long, TimeUnit)} (every tick, + * including the very first one {@code startOrCoalesce} kicks off) and {@link #shutdownNow()} (on + * {@code ReplyPushLoop.stop()}/{@code close()}) — verified by reading every {@code scheduler.} call + * site in that class. Every other {@link ScheduledExecutorService} method throws + * {@link UnsupportedOperationException} rather than silently doing the wrong thing, so a future + * change to {@code ReplyPushLoop} that starts calling one of them fails this fake loudly, on the + * very first test that exercises it, instead of being quietly mishandled. + */ +final class ManualScheduler implements ScheduledExecutorService { + + private final Deque pending = new ArrayDeque<>(); + private boolean shutdown = false; + + /** + * Run every task pending as of the START of this call — a snapshot taken before anything runs. + * {@link ReplyPushLoop#tick} always reschedules its own next tick before returning (see + * {@code scheduleNext} in its {@code INJECT}/{@code WAIT_BUSY} branches), so without the + * snapshot a single call here would recurse forever. Taking it up front means one call is + * exactly one tick, deterministically, no matter what that tick itself goes on to schedule. + * + * @return how many tasks actually ran + */ + synchronized int runDueTasks() { + List due = new ArrayList<>(pending); + pending.clear(); + for (Runnable task : due) { + task.run(); + } + return due.size(); + } + + /** How many tasks are currently queued, without running any of them. */ + synchronized int pendingCount() { + return pending.size(); + } + + @Override + public synchronized ScheduledFuture schedule(Runnable command, long delay, TimeUnit unit) { + if (shutdown) { + throw new RejectedExecutionException("ManualScheduler is shut down"); + } + pending.add(command); + return null; // ReplyPushLoop discards the return value of every schedule() call it makes. + } + + @Override + public synchronized List shutdownNow() { + shutdown = true; + List left = new ArrayList<>(pending); + pending.clear(); + return left; + } + + @Override + public synchronized boolean isShutdown() { + return shutdown; + } + + // --- everything below: not called by ReplyPushLoop, and not supported by this fake (see the + // class javadoc's "deliberately narrow" note) ------------------------------------------------- + + @Override + public ScheduledFuture scheduleAtFixedRate(Runnable command, long initialDelay, long period, TimeUnit unit) { + throw unsupported(); + } + + @Override + public ScheduledFuture scheduleWithFixedDelay(Runnable command, long initialDelay, long delay, TimeUnit unit) { + throw unsupported(); + } + + @Override + public ScheduledFuture schedule(Callable callable, long delay, TimeUnit unit) { + throw unsupported(); + } + + @Override + public void shutdown() { + throw unsupported(); + } + + @Override + public boolean isTerminated() { + throw unsupported(); + } + + @Override + public boolean awaitTermination(long timeout, TimeUnit unit) { + throw unsupported(); + } + + @Override + public Future submit(Callable task) { + throw unsupported(); + } + + @Override + public Future submit(Runnable task, T result) { + throw unsupported(); + } + + @Override + public Future submit(Runnable task) { + throw unsupported(); + } + + @Override + public List> invokeAll(Collection> tasks) { + throw unsupported(); + } + + @Override + public List> invokeAll(Collection> tasks, long timeout, TimeUnit unit) { + throw unsupported(); + } + + @Override + public T invokeAny(Collection> tasks) throws ExecutionException { + throw unsupported(); + } + + @Override + public T invokeAny(Collection> tasks, long timeout, TimeUnit unit) + throws ExecutionException { + throw unsupported(); + } + + @Override + public void execute(Runnable command) { + throw unsupported(); + } + + private static UnsupportedOperationException unsupported() { + return new UnsupportedOperationException( + "ManualScheduler only supports schedule(Runnable, long, TimeUnit), shutdownNow() and isShutdown() " + + "— see the class javadoc"); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java index 0155b35..99f8f80 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -1841,6 +1841,36 @@ class MessageServiceTest { return new PushWiring(service, leadHerdr, scheduler); } + /** + * A {@link MessageService} wired to a real {@link ReplyPushLoop}, like {@link #wireWithPushLoop}, + * but backed by {@link ManualScheduler} instead of a real timer (fleetd #608): its tick never + * fires on its own — a test drives it explicitly via {@link ManualScheduler#runDueTasks()}. + */ + private record ManualPushWiring(MessageService service, FakeHerdr leadHerdr, ManualScheduler scheduler) + implements AutoCloseable { + @Override + public void close() { + scheduler.shutdownNow(); + } + } + + /** + * As {@link #wireWithPushLoop}, but the schedule never runs on its own. {@code backoffMs} is + * still threaded through to {@link ReplyPushLoop}'s constructor (it takes one), but nothing here + * ever waits it out, so its value cannot affect anything a test built on this observes — see + * {@code anAlreadyCollectedTicketProducesNoNudge}, which sets it to 1 to prove exactly that. + */ + private ManualPushWiring wireWithManualScheduler(int maxReminders, long backoffMs) { + PrimaryRegistry registry = new PrimaryRegistry(null); + registry.recordDelegation(T, LEAD); + FakeHerdr leadHerdr = new FakeHerdr(); + AgentControl leadAgents = new AgentControl(leadHerdr); + ManualScheduler scheduler = new ManualScheduler(); + ReplyPushLoop pushLoop = new ReplyPushLoop(registry, leadAgents, inbox, scheduler, maxReminders, backoffMs); + MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop, null, System::nanoTime); + return new ManualPushWiring(service, leadHerdr, scheduler); + } + private void awaitNudge(FakeHerdr leadHerdr) throws InterruptedException { long deadline = System.currentTimeMillis() + 3000; while (!leadHerdr.called("agent.prompt") && System.currentTimeMillis() < deadline) { @@ -1882,16 +1912,46 @@ class MessageServiceTest { @Test void anAlreadyCollectedTicketProducesNoNudge() throws Exception { - try (var wiring = wireWithPushLoop(1, 300)) { // wide backoff: poll before the first tick fires - String ticket = wiring.service().sendAsync(T, "long task"); - awaitWaiting(); - injectDelivery(); - assertTrue(rendezvous.resolve(T, "async result")); + // fleetd #608: backoff is 1ms — the most hostile value there is, the tick due immediately — + // and the test still must pass, because with ManualScheduler the tick never runs on a timer + // at all; it only runs when this test calls runDueTasks() below. The old version bet a 300ms + // backoff was "wide" enough that collecting the ticket always won the race against a real + // scheduler's own timer; that held on an idle machine and failed under a loaded full-suite + // run — a false red on correct code, since nothing here forced that ordering, it only made + // it likely. + try (var wiring = wireWithManualScheduler(1, 1)) { + java.util.concurrent.CountDownLatch terminalReached = new java.util.concurrent.CountDownLatch(1); + // finishAsyncTask's task.future.complete(...) runs every whenComplete registered on that + // same future — including sendAsync's own hook that calls ReplyPushLoop.onTicketTerminal + // — synchronously, before complete() returns (see finishAsyncTask's javadoc and fleetd + // #399). This test-only hook fires right after that complete() call, on the very same + // (async executor) thread, so waiting for it guarantees onTicketTerminal has already run + // and the ticket is really sitting in the push loop's pending set — unlike waiting for + // Phase.DONE via poll(), which CompletableFuture.complete() can make visible to another + // thread before every whenComplete dependent has actually finished running (the same + // publish-then-run-dependents gap awaitCompletionStamped exists to close elsewhere). + String ticket; + wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(terminalReached::countDown); + try { + ticket = wiring.service().sendAsync(T, "long task"); + awaitWaiting(); + injectDelivery(); + assertTrue(rendezvous.resolve(T, "async result")); + assertTrue(terminalReached.await(5, TimeUnit.SECONDS), + "the ticket never reached its terminal phase"); + } finally { + wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(null); + } + // Collect it — the exact action the nudge must never follow. MessageService.TaskView view = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.DONE); assertEquals(MessageService.Phase.DONE, view.phase()); - Thread.sleep(400); // let the scheduled tick run — it must find nothing pending + // Now run the pending tick explicitly. onTicketTerminal scheduled exactly one (the ticket + // is the only thing this lead has ever had pending); it must find nothing left pending — + // the collection above already removed it — and send no nudge. + assertEquals(1, wiring.scheduler().runDueTasks(), + "expected exactly the one tick onTicketTerminal scheduled"); assertFalse(wiring.leadHerdr().called("agent.prompt"), "a ticket the lead already polled must never be nudged"); }