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 extends Callable> tasks) {
+ throw unsupported();
+ }
+
+ @Override
+ public List> invokeAll(Collection extends Callable> tasks, long timeout, TimeUnit unit) {
+ throw unsupported();
+ }
+
+ @Override
+ public T invokeAny(Collection extends Callable> tasks) throws ExecutionException {
+ throw unsupported();
+ }
+
+ @Override
+ public T invokeAny(Collection extends Callable> 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");
}