fleetd #608: make anAlreadyCollectedTicketProducesNoNudge deterministic
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 1m21s
CI / build (pull_request) Failing after 1m34s

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).
This commit is contained in:
Dai Ha
2026-09-20 17:18:45 +07:00
parent 9a992d0f70
commit aa517ae0ec
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");
}