#608 replace MessageService timing sleeps #623

Merged
ltms merged 1 commits from worker/608-sleeps-3a64ff-3 into main 2026-09-22 06:38:30 +02:00
@@ -1846,7 +1846,8 @@ class MessageServiceTest {
* but backed by {@link ManualScheduler} instead of a real timer (fleetd #608): its tick never * 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()}. * fires on its own — a test drives it explicitly via {@link ManualScheduler#runDueTasks()}.
*/ */
private record ManualPushWiring(MessageService service, FakeHerdr leadHerdr, ManualScheduler scheduler) private record ManualPushWiring(MessageService service, FakeHerdr leadHerdr, ManualScheduler scheduler,
ReplyPushLoop pushLoop)
implements AutoCloseable { implements AutoCloseable {
@Override @Override
public void close() { public void close() {
@@ -1861,14 +1862,20 @@ class MessageServiceTest {
* {@code anAlreadyCollectedTicketProducesNoNudge}, which sets it to 1 to prove exactly that. * {@code anAlreadyCollectedTicketProducesNoNudge}, which sets it to 1 to prove exactly that.
*/ */
private ManualPushWiring wireWithManualScheduler(int maxReminders, long backoffMs) { private ManualPushWiring wireWithManualScheduler(int maxReminders, long backoffMs) {
return wireWithManualScheduler(maxReminders, backoffMs, System::nanoTime);
}
/** As above, with an injectable clock for tests that exercise terminal-ticket pruning. */
private ManualPushWiring wireWithManualScheduler(int maxReminders, long backoffMs,
java.util.function.LongSupplier nowNanos) {
PrimaryRegistry registry = new PrimaryRegistry(null); PrimaryRegistry registry = new PrimaryRegistry(null);
registry.recordDelegation(T, LEAD); registry.recordDelegation(T, LEAD);
FakeHerdr leadHerdr = new FakeHerdr(); FakeHerdr leadHerdr = new FakeHerdr();
AgentControl leadAgents = new AgentControl(leadHerdr); AgentControl leadAgents = new AgentControl(leadHerdr);
ManualScheduler scheduler = new ManualScheduler(); ManualScheduler scheduler = new ManualScheduler();
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, leadAgents, inbox, scheduler, maxReminders, backoffMs); ReplyPushLoop pushLoop = new ReplyPushLoop(registry, leadAgents, inbox, scheduler, maxReminders, backoffMs);
MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop, null, System::nanoTime); MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop, null, nowNanos);
return new ManualPushWiring(service, leadHerdr, scheduler); return new ManualPushWiring(service, leadHerdr, scheduler, pushLoop);
} }
private void awaitNudge(FakeHerdr leadHerdr) throws InterruptedException { private void awaitNudge(FakeHerdr leadHerdr) throws InterruptedException {
@@ -1959,24 +1966,35 @@ class MessageServiceTest {
@Test @Test
void severalAsyncTicketsFinishingTogetherProduceOneCoalescedNudge() throws Exception { void severalAsyncTicketsFinishingTogetherProduceOneCoalescedNudge() throws Exception {
try (var wiring = wireWithPushLoop(1, 300)) { // wide backoff: both tickets land before the tick fires // A 1ms backoff is due immediately. ManualScheduler still cannot run it until this test
String first = wiring.service().sendAsync(T, "first task"); // explicitly calls runDueTasks(), so both terminal tickets join one scheduled tick.
awaitWaiting(); try (var wiring = wireWithManualScheduler(1, 1)) {
injectDelivery(); String first;
assertTrue(rendezvous.resolve(T, "first done")); String second;
// Settle without polling: poll() itself marks a ticket collected (that's the point of java.util.concurrent.CountDownLatch firstTerminalReached = new java.util.concurrent.CountDownLatch(1);
// anAlreadyCollectedTicketProducesNoNudge above) — using it here to detect completion wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(firstTerminalReached::countDown);
// would collect the ticket before the coalescing this test checks ever gets a chance. try {
Thread.sleep(100); first = wiring.service().sendAsync(T, "first task");
awaitWaiting();
injectDelivery();
assertTrue(rendezvous.resolve(T, "first done"));
assertTrue(firstTerminalReached.await(5, TimeUnit.SECONDS),
"the first ticket never reached its terminal phase");
String second = wiring.service().sendAsync(T, "second task"); java.util.concurrent.CountDownLatch secondTerminalReached = new java.util.concurrent.CountDownLatch(1);
awaitWaiting(); wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(secondTerminalReached::countDown);
injectDelivery(); second = wiring.service().sendAsync(T, "second task");
assertTrue(rendezvous.resolve(T, "second done")); awaitWaiting();
Thread.sleep(100); injectDelivery();
assertTrue(rendezvous.resolve(T, "second done"));
assertTrue(secondTerminalReached.await(5, TimeUnit.SECONDS),
"the second ticket never reached its terminal phase");
} finally {
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(null);
}
awaitNudge(wiring.leadHerdr()); assertEquals(1, wiring.scheduler().runDueTasks(),
Thread.sleep(200); // settle — nothing more should arrive beyond the one coalesced nudge "both terminal tickets must coalesce onto one scheduled tick");
long nudgeCount = wiring.leadHerdr().calls.stream() long nudgeCount = wiring.leadHerdr().calls.stream()
.filter(c -> c.method().equals("agent.prompt")).count(); .filter(c -> c.method().equals("agent.prompt")).count();
assertEquals(1, nudgeCount, "two tickets finishing together must produce ONE nudge, not two"); assertEquals(1, nudgeCount, "two tickets finishing together must produce ONE nudge, not two");
@@ -2055,7 +2073,8 @@ class MessageServiceTest {
@Test @Test
void answeringAQuestionStopsFurtherNudgesAboutIt() throws Exception { void answeringAQuestionStopsFurtherNudgesAboutIt() throws Exception {
try (var wiring = wireWithPushLoop(5, 50)) { // A 1ms backoff is due immediately, but ManualScheduler only ticks when this test asks it to.
try (var wiring = wireWithManualScheduler(5, 1)) {
String ticket = wiring.service().sendAsync(T, "task that asks"); String ticket = wiring.service().sendAsync(T, "task that asks");
awaitWaiting(); awaitWaiting();
injectDelivery(); injectDelivery();
@@ -2063,8 +2082,9 @@ class MessageServiceTest {
CompletableFuture<MessageService.AskResult> ask = CompletableFuture.supplyAsync( CompletableFuture<MessageService.AskResult> ask = CompletableFuture.supplyAsync(
() -> wiring.service().ask(T, "which config file?", 5000)); () -> wiring.service().ask(T, "which config file?", 5000));
MessageService.TaskView asking = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.ASKING); MessageService.TaskView asking = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.ASKING);
awaitQuestionPendingOn(wiring.pushLoop(), LEAD, asking.turnId());
awaitNudge(wiring.leadHerdr()); assertEquals(1, wiring.scheduler().runDueTasks(), "the open question must have one scheduled tick");
long callsBeforeAnswer = wiring.leadHerdr().calls.stream() long callsBeforeAnswer = wiring.leadHerdr().calls.stream()
.filter(c -> c.method().equals("agent.prompt")).count(); .filter(c -> c.method().equals("agent.prompt")).count();
@@ -2072,12 +2092,22 @@ class MessageServiceTest {
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000)); () -> wiring.service().answer(asking.turnId(), "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer()); assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting(); awaitWaiting();
java.util.concurrent.CountDownLatch terminalReached = new java.util.concurrent.CountDownLatch(1);
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(terminalReached::countDown);
assertTrue(rendezvous.resolve(T, "done")); assertTrue(rendezvous.resolve(T, "done"));
try {
assertTrue(terminalReached.await(5, TimeUnit.SECONDS),
"the answered ticket never reached its terminal phase");
} finally {
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(null);
}
answer.get(5, TimeUnit.SECONDS); answer.get(5, TimeUnit.SECONDS);
// Let several more ticks (and the ticket's own now-legitimate terminal nudge) fire — // Run the ticket's legitimate terminal nudge and every later scheduled tick through its
// none of them may still name the question's turnId, which is closed. // reminder cap. None may still name the closed question.
Thread.sleep(300); for (int tick = 0; tick < 6; tick++) {
assertEquals(1, wiring.scheduler().runDueTasks(), "expected one scheduled reminder tick");
}
boolean anyNamesClosedQuestion = wiring.leadHerdr().calls.stream() boolean anyNamesClosedQuestion = wiring.leadHerdr().calls.stream()
.filter(c -> c.method().equals("agent.prompt")) .filter(c -> c.method().equals("agent.prompt"))
.skip(callsBeforeAnswer) .skip(callsBeforeAnswer)
@@ -2160,35 +2190,45 @@ class MessageServiceTest {
// decideTickets hits the cap and STOPs — activeLeads drops the lead, but (before the fix) // decideTickets hits the cap and STOPs — activeLeads drops the lead, but (before the fix)
// pendingTickets never drops the ticket. That is the exact "cap already STOPped" branch of // pendingTickets never drops the ticket. That is the exact "cap already STOPped" branch of
// the bug report, reached deterministically rather than by timing it against a live tick. // the bug report, reached deterministically rather than by timing it against a live tick.
try (var wiring = wireWithPushLoop(1, 50, clock::get)) { // A 1ms backoff is due immediately, but ManualScheduler runs only the ticks below.
String stale = wiring.service().sendAsync(T, "first task"); try (var wiring = wireWithManualScheduler(1, 1, clock::get)) {
awaitWaiting(); String stale;
injectDelivery(); String fresh;
assertTrue(rendezvous.resolve(T, "stale result")); java.util.concurrent.CountDownLatch staleTerminalReached = new java.util.concurrent.CountDownLatch(1);
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(staleTerminalReached::countDown);
try {
stale = wiring.service().sendAsync(T, "first task");
awaitWaiting();
injectDelivery();
assertTrue(rendezvous.resolve(T, "stale result"));
assertTrue(staleTerminalReached.await(5, TimeUnit.SECONDS),
"the stale ticket never reached its terminal phase");
// Let the reminder loop fire its one nudge and hit the cap (STOP removes it from // Fire the stale ticket's one nudge, then its cap tick (STOP removes it from activeLeads;
// activeLeads; pendingTickets is untouched either way — that asymmetry is the bug). // pendingTickets is untouched either way — that asymmetry is the bug).
awaitNudge(wiring.leadHerdr()); assertEquals(1, wiring.scheduler().runDueTasks(), "the stale ticket must have one nudge tick");
Thread.sleep(300); assertEquals(1, wiring.scheduler().runDueTasks(), "the stale ticket must have one cap tick");
assertTrue(wiring.leadHerdr().lastCall("agent.prompt").params().toString().contains(stale), assertTrue(wiring.leadHerdr().lastCall("agent.prompt").params().toString().contains(stale),
"sanity: the stale ticket's own reminder must have fired first"); "sanity: the stale ticket's own reminder must have fired first");
// Cross the TTL — a real clock would need 10 minutes; the injected one does it instantly. // Cross the TTL — a real clock would need 10 minutes; the injected one does it instantly.
clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1)); clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1));
// A second, unrelated ticket to the same target/lead reaches sendAsync, which prunes. // A second, unrelated ticket to the same target/lead reaches sendAsync, which prunes.
String fresh = wiring.service().sendAsync(T, "second task"); java.util.concurrent.CountDownLatch freshTerminalReached = new java.util.concurrent.CountDownLatch(1);
awaitWaiting(); wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(freshTerminalReached::countDown);
injectDelivery(); fresh = wiring.service().sendAsync(T, "second task");
assertTrue(rendezvous.resolve(T, "fresh result")); awaitWaiting();
injectDelivery();
// The fresh ticket restarts the (now-dormant) reminder loop with its own nudge. assertTrue(rendezvous.resolve(T, "fresh result"));
long before = wiring.leadHerdr().calls.stream().filter(c -> c.method().equals("agent.prompt")).count(); assertTrue(freshTerminalReached.await(5, TimeUnit.SECONDS),
long deadline = System.currentTimeMillis() + 3000; "the fresh ticket never reached its terminal phase");
while (wiring.leadHerdr().calls.stream().filter(c -> c.method().equals("agent.prompt")).count() <= before } finally {
&& System.currentTimeMillis() < deadline) { wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(null);
Thread.sleep(10);
} }
// The fresh ticket restarts the now-dormant reminder loop with its own nudge.
assertEquals(1, wiring.scheduler().runDueTasks(), "the fresh ticket must have one nudge tick");
String latestNudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString(); String latestNudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString();
assertTrue(latestNudge.contains(fresh), "the fresh ticket's nudge must still arrive: " + latestNudge); assertTrue(latestNudge.contains(fresh), "the fresh ticket's nudge must still arrive: " + latestNudge);
assertFalse(latestNudge.contains(stale), assertFalse(latestNudge.contains(stale),