diff --git a/server-common/src/test/java/org/a2aproject/sdk/server/events/MainEventBusProcessorPushNotificationOrderTest.java b/server-common/src/test/java/org/a2aproject/sdk/server/events/MainEventBusProcessorPushNotificationOrderTest.java index 67e6f6d32..655973e85 100644 --- a/server-common/src/test/java/org/a2aproject/sdk/server/events/MainEventBusProcessorPushNotificationOrderTest.java +++ b/server-common/src/test/java/org/a2aproject/sdk/server/events/MainEventBusProcessorPushNotificationOrderTest.java @@ -146,8 +146,17 @@ public void testSynchronousExecutorDoesNotLeakChainEntries() throws InterruptedE // One entry per DISTINCT task leaks, not repeated pushes for the same one -- // a single task's own map slot just gets overwritten each time. Only several // different task IDs reveal growth. + // + // Each task produces TWO push notifications: one for the initial Task event and one for the + // StatusUpdate event (both implement StreamingEventKind). With a synchronous executor, + // pushTask runs inside compute()'s lambda -- before the entry is stored in the map -- so + // the latch fires before the entry is even visible. To avoid a TOCTOU race between the + // while-loop exit (count==0 transiently) and the assertEquals re-reading count, we: + // 1. Set the latch to cover ALL push notifications (taskCount * 2). + // 2. Always sleep before the first count check so the last cleanup has time to run. + // 3. Assert on the captured count, not a fresh read. int taskCount = 10; - CountDownLatch latch = new CountDownLatch(taskCount); + CountDownLatch latch = new CountDownLatch(taskCount * 2); PushNotificationSender sender = (event, snapshot) -> latch.countDown(); mainEventBusProcessor = new MainEventBusProcessor(mainEventBus, taskStore, sender, queueManager); @@ -171,10 +180,16 @@ public void testSynchronousExecutorDoesNotLeakChainEntries() throws InterruptedE assertTrue(latch.await(5, TimeUnit.SECONDS), "All push notifications should complete"); long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5); - while (mainEventBusProcessor.pushNotificationChainCount() != 0 && System.nanoTime() < deadline) { + int count; + do { + // Always sleep before checking: with the sync executor, pushTask fires inside + // compute()'s lambda (before the entry is stored), so the latch may reach 0 + // slightly before the final whenComplete removal runs. A brief sleep lets the + // processor thread finish that last cleanup step. Thread.sleep(10); - } - assertEquals(0, mainEventBusProcessor.pushNotificationChainCount(), + count = mainEventBusProcessor.pushNotificationChainCount(); + } while (count != 0 && System.nanoTime() < deadline); + assertEquals(0, count, "A completed task's chain entry must not be left in the map"); }