Skip to content

Commit 9de9918

Browse files
committed
test(sse): cover bounded ordered dispatcher
1 parent efd8cfd commit 9de9918

1 file changed

Lines changed: 89 additions & 0 deletions

File tree

Lines changed: 89 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,89 @@
1+
package io.github.easy4j.hermes.api.sse;
2+
3+
import org.junit.jupiter.api.Test;
4+
5+
import java.util.ArrayList;
6+
import java.util.Collections;
7+
import java.util.List;
8+
import java.util.concurrent.CountDownLatch;
9+
import java.util.concurrent.TimeUnit;
10+
11+
import static org.junit.jupiter.api.Assertions.assertEquals;
12+
import static org.junit.jupiter.api.Assertions.assertThrows;
13+
import static org.junit.jupiter.api.Assertions.assertTrue;
14+
15+
class OrderedEventDispatcherTest {
16+
17+
@Test
18+
void dispatchesInOrderAndReportsQueuedTasks() throws Exception {
19+
OrderedEventDispatcher dispatcher = new OrderedEventDispatcher("run:test", 4);
20+
CountDownLatch firstStarted = new CountDownLatch(1);
21+
CountDownLatch releaseFirst = new CountDownLatch(1);
22+
CountDownLatch completed = new CountDownLatch(2);
23+
List<Integer> order = Collections.synchronizedList(new ArrayList<Integer>());
24+
try {
25+
dispatcher.dispatch(() -> {
26+
firstStarted.countDown();
27+
try {
28+
releaseFirst.await(2, TimeUnit.SECONDS);
29+
} catch (InterruptedException interrupted) {
30+
Thread.currentThread().interrupt();
31+
}
32+
order.add(1);
33+
completed.countDown();
34+
});
35+
assertTrue(firstStarted.await(2, TimeUnit.SECONDS));
36+
37+
dispatcher.dispatch(() -> {
38+
order.add(2);
39+
completed.countDown();
40+
});
41+
42+
assertEquals(1, dispatcher.queuedTaskCount());
43+
releaseFirst.countDown();
44+
assertTrue(completed.await(2, TimeUnit.SECONDS));
45+
assertEquals(java.util.Arrays.asList(1, 2), order);
46+
} finally {
47+
releaseFirst.countDown();
48+
dispatcher.close();
49+
}
50+
}
51+
52+
@Test
53+
void rejectsWhenBoundedQueueIsFull() throws Exception {
54+
OrderedEventDispatcher dispatcher = new OrderedEventDispatcher("overflow", 1);
55+
CountDownLatch workerStarted = new CountDownLatch(1);
56+
CountDownLatch releaseWorker = new CountDownLatch(1);
57+
try {
58+
dispatcher.dispatch(() -> {
59+
workerStarted.countDown();
60+
try {
61+
releaseWorker.await(2, TimeUnit.SECONDS);
62+
} catch (InterruptedException interrupted) {
63+
Thread.currentThread().interrupt();
64+
}
65+
});
66+
assertTrue(workerStarted.await(2, TimeUnit.SECONDS));
67+
dispatcher.dispatch(() -> { });
68+
69+
assertThrows(SseQueueOverflowException.class,
70+
() -> dispatcher.dispatch(() -> { }));
71+
} finally {
72+
releaseWorker.countDown();
73+
dispatcher.close();
74+
}
75+
}
76+
77+
@Test
78+
void closeIsIdempotentAndRejectsNewWork() {
79+
OrderedEventDispatcher dispatcher = new OrderedEventDispatcher(null, 1);
80+
dispatcher.close();
81+
dispatcher.close();
82+
83+
assertThrows(IllegalStateException.class,
84+
() -> dispatcher.dispatch(() -> { }));
85+
assertThrows(NullPointerException.class,
86+
() -> dispatcher.dispatch(null));
87+
assertEquals(0, dispatcher.queuedTaskCount());
88+
}
89+
}

0 commit comments

Comments
 (0)