From cd856bf21c2cc3d76f53dd73df3b4eb47d43e610 Mon Sep 17 00:00:00 2001 From: Emmanuel Hugonnet Date: Mon, 28 Sep 2026 21:34:08 +0200 Subject: [PATCH] test(server): fix flaky testSynchronousExecutorDoesNotLeakChainEntries Two bugs caused intermittent failures: 1. Latch count was taskCount (10) but each task generates TWO push notifications -> one for the Task event and one for the StatusUpdate event, both implementing StreamingEventKind -> so the latch could fire with 10 notifications still in flight. Fixed to taskCount * 2. 2. TOCTOU race between the while-loop exit and assertEquals: with a synchronous executor, pushTask runs inside compute()'s lambda before the chain entry is stored, so the latch can reach 0 when count is transiently 0. The while exited on count==0 but assertEquals re-read count and could see 1 if the processor stored the entry in between. Fixed by switching to a do-while that always sleeps before the first check and asserts on the captured count, not a fresh read. Signed-off-by: Emmanuel Hugonnet --- ...BusProcessorPushNotificationOrderTest.java | 23 +++++++++++++++---- 1 file changed, 19 insertions(+), 4 deletions(-) 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"); }