fleetd #608: make anAlreadyCollectedTicketProducesNoNudge deterministic #611

Merged
ltms merged 1 commits from worker/fleetd-608-flaky-nudge-test-d0c2d1-3 into main 2026-09-20 12:26:16 +02:00
2 changed files with 234 additions and 6 deletions
@@ -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).
*
* <p>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.
*
* <p><strong>Deliberately narrow.</strong> {@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<Runnable> 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<Runnable> 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<Runnable> shutdownNow() {
shutdown = true;
List<Runnable> 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 <V> ScheduledFuture<V> schedule(Callable<V> 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 <T> Future<T> submit(Callable<T> task) {
throw unsupported();
}
@Override
public <T> Future<T> submit(Runnable task, T result) {
throw unsupported();
}
@Override
public Future<?> submit(Runnable task) {
throw unsupported();
}
@Override
public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks) {
throw unsupported();
}
@Override
public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) {
throw unsupported();
}
@Override
public <T> T invokeAny(Collection<? extends Callable<T>> tasks) throws ExecutionException {
throw unsupported();
}
@Override
public <T> T invokeAny(Collection<? extends Callable<T>> 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");
}
}
@@ -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");
}