|
17 | 17 |
|
18 | 18 | package org.apache.dolphinscheduler.server.master.runner; |
19 | 19 |
|
| 20 | +import static org.junit.jupiter.api.Assertions.assertEquals; |
| 21 | +import static org.junit.jupiter.api.Assertions.assertFalse; |
| 22 | +import static org.junit.jupiter.api.Assertions.assertTrue; |
20 | 23 | import static org.mockito.Mockito.mock; |
21 | | -import static org.mockito.Mockito.verify; |
22 | 24 |
|
| 25 | +import org.apache.dolphinscheduler.dao.entity.WorkerGroup; |
23 | 26 | import org.apache.dolphinscheduler.server.master.engine.task.runnable.ITaskExecutionRunnable; |
24 | 27 |
|
25 | | -import org.junit.jupiter.api.Assertions; |
26 | | -import org.junit.jupiter.api.BeforeEach; |
| 28 | +import java.time.Duration; |
| 29 | +import java.util.Arrays; |
| 30 | +import java.util.List; |
| 31 | + |
| 32 | +import org.awaitility.Awaitility; |
27 | 33 | import org.junit.jupiter.api.Test; |
28 | 34 | import org.junit.jupiter.api.extension.ExtendWith; |
29 | 35 | import org.mockito.InjectMocks; |
30 | | -import org.mockito.MockitoAnnotations; |
31 | 36 | import org.mockito.junit.jupiter.MockitoExtension; |
32 | 37 |
|
| 38 | +import com.google.common.truth.Truth; |
| 39 | + |
33 | 40 | @ExtendWith(MockitoExtension.class) |
34 | 41 | public class WorkerGroupTaskDispatcherManagerTest { |
35 | 42 |
|
36 | 43 | @InjectMocks |
37 | | - private WorkerGroupTaskDispatcherManager workerGroupTaskDispatcherManager; |
38 | | - |
39 | | - @BeforeEach |
40 | | - public void setUp() { |
41 | | - MockitoAnnotations.openMocks(this); |
42 | | - } |
| 44 | + private WorkerGroupTaskDispatcherManager manager; |
43 | 45 |
|
44 | 46 | @Test |
45 | | - public void testAddTaskToExistingWorkerGroup() { |
46 | | - String workerGroup = "testWorkerGroup"; |
47 | | - ITaskExecutionRunnable task = mock(ITaskExecutionRunnable.class); |
48 | | - long delay = 1000L; |
| 47 | + public void testAddTaskToExistingWorkerGroup_ShouldReturnTrue() { |
| 48 | + String workerGroupName = "testGroup"; |
| 49 | + manager.addWorkerGroup(workerGroupName); |
| 50 | + ITaskExecutionRunnable mockTask = mock(ITaskExecutionRunnable.class); |
49 | 51 |
|
50 | | - workerGroupTaskDispatcherManager.add(workerGroup, task, delay); |
| 52 | + boolean result = manager.add(workerGroupName, mockTask, 0L); |
51 | 53 |
|
52 | | - // not have workerGroup queue,cannot add |
53 | | - Assertions.assertEquals(0, |
54 | | - workerGroupTaskDispatcherManager.getDispatchWorkerMap().size()); |
| 54 | + assertTrue(result); |
| 55 | + } |
| 56 | + |
| 57 | + @Test |
| 58 | + public void testAddTaskToNonExistingWorkerGroup_ShouldReturnFalse() { |
| 59 | + String workerGroupName = "nonExistingGroup"; |
| 60 | + ITaskExecutionRunnable mockTask = mock(ITaskExecutionRunnable.class); |
| 61 | + boolean result = manager.add(workerGroupName, mockTask, 0L); |
| 62 | + assertFalse(result); |
55 | 63 | } |
56 | 64 |
|
57 | 65 | @Test |
58 | | - public void testAddTaskToNonExistingWorkerGroup() { |
59 | | - String workerGroup = "nonExistingWorkerGroup"; |
60 | | - ITaskExecutionRunnable task = mock(ITaskExecutionRunnable.class); |
61 | | - long delay = 1000L; |
62 | | - workerGroupTaskDispatcherManager.addWorkerGroup(workerGroup); |
63 | | - workerGroupTaskDispatcherManager.add(workerGroup, task, delay); |
64 | | - |
65 | | - Assertions.assertTrue( |
66 | | - workerGroupTaskDispatcherManager.getWorkerGroupPriorityDelayQueueMap().containsKey(workerGroup)); |
| 66 | + public void testDeleteExistingWorkerGroup_ShouldRemoveGroup() throws Exception { |
| 67 | + String workerGroupName = "testGroup"; |
| 68 | + manager.addWorkerGroup(workerGroupName); |
| 69 | + |
| 70 | + manager.deleteWorkerGroup(workerGroupName); |
| 71 | + |
| 72 | + Awaitility.await() |
| 73 | + .atMost(Duration.ofSeconds(5)) |
| 74 | + .untilAsserted(() -> { |
| 75 | + Truth.assertThat(manager.getDispatchWorkerMap().isEmpty()).isTrue(); |
| 76 | + }); |
67 | 77 | } |
68 | 78 |
|
69 | 79 | @Test |
70 | | - public void testStopWorkerGroup() { |
71 | | - String workerGroup = "testWorkerGroup"; |
72 | | - ITaskExecutionRunnable task = mock(ITaskExecutionRunnable.class); |
73 | | - long delay = 1000L; |
74 | | - |
75 | | - workerGroupTaskDispatcherManager.addWorkerGroup(workerGroup); |
76 | | - workerGroupTaskDispatcherManager.add(workerGroup, task, delay); |
77 | | - Assertions.assertTrue( |
78 | | - workerGroupTaskDispatcherManager.getWorkerGroupPriorityDelayQueueMap().get(workerGroup).size() > 0); |
| 80 | + public void testAddNewWorkerGroup_ShouldAddGroup() { |
| 81 | + String workerGroupName = "newGroup"; |
| 82 | + manager.addWorkerGroup(workerGroupName); |
| 83 | + assertFalse(manager.getDispatchWorkerMap().isEmpty()); |
79 | 84 | } |
80 | 85 |
|
81 | 86 | @Test |
82 | | - public void testAddWorkerGroup() { |
83 | | - String workerGroup = "newWorkerGroup"; |
| 87 | + public void testClose_ShouldShutdownScheduler() throws Exception { |
| 88 | + manager.addWorkerGroup("testGroup"); |
| 89 | + manager.add("testGroup", mock(ITaskExecutionRunnable.class), 0); |
| 90 | + WorkerGroupTaskDispatcher dispatcher = manager.getDispatchWorkerMap().get("testGroup"); |
84 | 91 |
|
85 | | - workerGroupTaskDispatcherManager.addWorkerGroup(workerGroup); |
| 92 | + manager.deleteWorkerGroup("testGroup"); |
| 93 | + |
| 94 | + Awaitility.await() |
| 95 | + .atMost(Duration.ofSeconds(5)) |
| 96 | + .untilAsserted(() -> { |
| 97 | + Truth.assertThat(dispatcher.getStatus() == DispatchWorkerStatus.DELETE_SUCCESS).isTrue(); |
| 98 | + }); |
86 | 99 |
|
87 | | - Assertions.assertTrue( |
88 | | - workerGroupTaskDispatcherManager.getWorkerGroupPriorityDelayQueueMap().containsKey(workerGroup)); |
89 | | - Assertions.assertTrue(workerGroupTaskDispatcherManager.getWorkerGroupPriorityDelayQueueMap() |
90 | | - .containsKey(workerGroup)); |
91 | 100 | } |
92 | 101 |
|
93 | 102 | @Test |
94 | | - public void testClose() throws Exception { |
95 | | - String workerGroup = "testWorkerGroup"; |
96 | | - WorkerGroupTaskDispatcher looper = mock(WorkerGroupTaskDispatcher.class); |
| 103 | + public void testOnWorkerGroupAdd_ShouldAddWorkerGroups() { |
| 104 | + WorkerGroup group1 = new WorkerGroup(); |
| 105 | + WorkerGroup group2 = new WorkerGroup(); |
| 106 | + group1.setName("testGroup1"); |
| 107 | + group2.setName("testGroup2"); |
| 108 | + List<WorkerGroup> workerGroups = Arrays.asList(group1, group2); |
| 109 | + manager.onWorkerGroupAdd(workerGroups); |
| 110 | + assertEquals(2, manager.getDispatchWorkerMap().size()); |
| 111 | + } |
97 | 112 |
|
98 | | - workerGroupTaskDispatcherManager.getDispatchWorkerMap().put(workerGroup, looper); |
| 113 | + @Test |
| 114 | + public void testOnWorkerGroupDelete_ShouldDeleteWorkerGroups() { |
| 115 | + WorkerGroup group1 = new WorkerGroup(); |
| 116 | + WorkerGroup group2 = new WorkerGroup(); |
| 117 | + group1.setName("testGroup1"); |
| 118 | + group2.setName("testGroup2"); |
| 119 | + List<WorkerGroup> workerGroups = Arrays.asList(group1, group2); |
| 120 | + workerGroups.forEach(workerGroup -> manager.addWorkerGroup(workerGroup.getName())); |
| 121 | + |
| 122 | + manager.onWorkerGroupDelete(workerGroups); |
| 123 | + |
| 124 | + Awaitility.await() |
| 125 | + .atMost(Duration.ofSeconds(5)) |
| 126 | + .untilAsserted(() -> { |
| 127 | + Truth.assertThat(manager.getDispatchWorkerMap().isEmpty()).isTrue(); |
| 128 | + }); |
99 | 129 |
|
100 | | - workerGroupTaskDispatcherManager.deleteWorkerGroup(workerGroup); |
| 130 | + } |
101 | 131 |
|
102 | | - verify(looper).close(); |
| 132 | + @Test |
| 133 | + public void testOnCloseWorkerGroupTaskDispatcherManager() throws Exception { |
| 134 | + WorkerGroup group1 = new WorkerGroup(); |
| 135 | + WorkerGroup group2 = new WorkerGroup(); |
| 136 | + group1.setName("testGroup1"); |
| 137 | + group2.setName("testGroup2"); |
| 138 | + List<WorkerGroup> workerGroups = Arrays.asList(group1, group2); |
| 139 | + workerGroups.forEach(workerGroup -> manager.addWorkerGroup(workerGroup.getName())); |
| 140 | + |
| 141 | + manager.close(); |
| 142 | + Awaitility.await() |
| 143 | + .atMost(Duration.ofSeconds(5)) |
| 144 | + .untilAsserted(() -> { |
| 145 | + Truth.assertThat(manager.getDispatchWorkerMap().isEmpty()).isTrue(); |
| 146 | + }); |
103 | 147 | } |
104 | 148 | } |
0 commit comments