Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand All @@ -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");
}

Expand Down
Loading