diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/TaskManagerTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/TaskManagerTest.java index 81674557afa28..273b1ac35d2fc 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/TaskManagerTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/TaskManagerTest.java @@ -90,7 +90,6 @@ import java.util.stream.Collectors; import static java.util.Arrays.asList; -import static java.util.Collections.emptyList; import static java.util.Collections.emptyMap; import static java.util.Collections.emptySet; import static java.util.Collections.singleton; @@ -194,7 +193,7 @@ public class TaskManagerTest { private Consumer consumer; @org.mockito.Mock private ActiveTaskCreator activeTaskCreator; - @Mock(type = MockType.NICE) + @org.mockito.Mock private StandbyTaskCreator standbyTaskCreator; @org.mockito.Mock private Admin adminClient; @@ -295,17 +294,15 @@ public void shouldPrepareActiveTaskInStateUpdaterToBeRecycled() { final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); when(stateUpdater.getTasks()).thenReturn(mkSet(activeTaskToRecycle)); - expect(standbyTaskCreator.createTasks(Collections.emptyMap())).andReturn(emptySet()); - replay(standbyTaskCreator); taskManager.handleAssignment( Collections.emptyMap(), mkMap(mkEntry(activeTaskToRecycle.id(), activeTaskToRecycle.inputPartitions())) ); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).createTasks(consumer, Collections.emptyMap()); Mockito.verify(tasks).addPendingTaskToRecycle(activeTaskToRecycle.id(), activeTaskToRecycle.inputPartitions()); + Mockito.verify(standbyTaskCreator).createTasks(Collections.emptyMap()); } @Test @@ -316,18 +313,16 @@ public void shouldPrepareStandbyTaskInStateUpdaterToBeRecycled() { final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); when(stateUpdater.getTasks()).thenReturn(mkSet(standbyTaskToRecycle)); - expect(standbyTaskCreator.createTasks(Collections.emptyMap())).andReturn(emptySet()); - replay(standbyTaskCreator); taskManager.handleAssignment( mkMap(mkEntry(standbyTaskToRecycle.id(), standbyTaskToRecycle.inputPartitions())), Collections.emptyMap() ); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).createTasks(consumer, Collections.emptyMap()); Mockito.verify(stateUpdater).remove(standbyTaskToRecycle.id()); Mockito.verify(tasks).addPendingTaskToRecycle(standbyTaskToRecycle.id(), standbyTaskToRecycle.inputPartitions()); + Mockito.verify(standbyTaskCreator).createTasks(Collections.emptyMap()); } @Test @@ -338,15 +333,13 @@ public void shouldRemoveUnusedActiveTaskFromStateUpdater() { final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); when(stateUpdater.getTasks()).thenReturn(mkSet(activeTaskToClose)); - expect(standbyTaskCreator.createTasks(Collections.emptyMap())).andReturn(emptySet()); - replay(standbyTaskCreator); taskManager.handleAssignment(Collections.emptyMap(), Collections.emptyMap()); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).createTasks(consumer, Collections.emptyMap()); Mockito.verify(stateUpdater).remove(activeTaskToClose.id()); Mockito.verify(tasks).addPendingTaskToCloseClean(activeTaskToClose.id()); + Mockito.verify(standbyTaskCreator).createTasks(Collections.emptyMap()); } @Test @@ -357,15 +350,13 @@ public void shouldRemoveUnusedStandbyTaskFromStateUpdater() { final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); when(stateUpdater.getTasks()).thenReturn(mkSet(standbyTaskToClose)); - expect(standbyTaskCreator.createTasks(Collections.emptyMap())).andReturn(emptySet()); - replay(standbyTaskCreator); taskManager.handleAssignment(Collections.emptyMap(), Collections.emptyMap()); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).createTasks(consumer, Collections.emptyMap()); Mockito.verify(stateUpdater).remove(standbyTaskToClose.id()); Mockito.verify(tasks).addPendingTaskToCloseClean(standbyTaskToClose.id()); + Mockito.verify(standbyTaskCreator).createTasks(Collections.emptyMap()); } @Test @@ -377,18 +368,16 @@ public void shouldUpdateInputPartitionOfActiveTaskInStateUpdater() { final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); when(stateUpdater.getTasks()).thenReturn(mkSet(activeTaskToUpdateInputPartitions)); - expect(standbyTaskCreator.createTasks(Collections.emptyMap())).andReturn(emptySet()); - replay(standbyTaskCreator); taskManager.handleAssignment( mkMap(mkEntry(activeTaskToUpdateInputPartitions.id(), newInputPartitions)), Collections.emptyMap() ); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).createTasks(consumer, Collections.emptyMap()); Mockito.verify(stateUpdater).remove(activeTaskToUpdateInputPartitions.id()); Mockito.verify(tasks).addPendingTaskToUpdateInputPartitions(activeTaskToUpdateInputPartitions.id(), newInputPartitions); + Mockito.verify(standbyTaskCreator).createTasks(Collections.emptyMap()); } @Test @@ -399,16 +388,14 @@ public void shouldKeepReAssignedActiveTaskInStateUpdater() { final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); when(stateUpdater.getTasks()).thenReturn(mkSet(reassignedActiveTask)); - expect(standbyTaskCreator.createTasks(Collections.emptyMap())).andReturn(emptySet()); - replay(standbyTaskCreator); taskManager.handleAssignment( mkMap(mkEntry(reassignedActiveTask.id(), reassignedActiveTask.inputPartitions())), Collections.emptyMap() ); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).createTasks(consumer, Collections.emptyMap()); + Mockito.verify(standbyTaskCreator).createTasks(Collections.emptyMap()); } @Test @@ -419,17 +406,15 @@ public void shouldRemoveReAssignedRevokedActiveTaskInStateUpdaterFromPendingTask final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); when(stateUpdater.getTasks()).thenReturn(mkSet(reAssignedRevokedActiveTask)); - expect(standbyTaskCreator.createTasks(Collections.emptyMap())).andReturn(emptySet()); - replay(standbyTaskCreator); taskManager.handleAssignment( mkMap(mkEntry(reAssignedRevokedActiveTask.id(), reAssignedRevokedActiveTask.inputPartitions())), Collections.emptyMap() ); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).createTasks(consumer, Collections.emptyMap()); Mockito.verify(tasks).removePendingActiveTaskToSuspend(reAssignedRevokedActiveTask.id()); + Mockito.verify(standbyTaskCreator).createTasks(Collections.emptyMap()); } @Test @@ -441,19 +426,17 @@ public void shouldNeverUpdateInputPartitionsOfStandbyTaskInStateUpdater() { final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); when(stateUpdater.getTasks()).thenReturn(mkSet(standbyTaskToUpdateInputPartitions)); - expect(standbyTaskCreator.createTasks(Collections.emptyMap())).andReturn(emptySet()); - replay(standbyTaskCreator); taskManager.handleAssignment( Collections.emptyMap(), mkMap(mkEntry(standbyTaskToUpdateInputPartitions.id(), newInputPartitions)) ); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).createTasks(consumer, Collections.emptyMap()); Mockito.verify(stateUpdater, never()).remove(standbyTaskToUpdateInputPartitions.id()); Mockito.verify(tasks, never()) .addPendingTaskToUpdateInputPartitions(standbyTaskToUpdateInputPartitions.id(), newInputPartitions); + Mockito.verify(standbyTaskCreator).createTasks(Collections.emptyMap()); } @Test @@ -464,16 +447,14 @@ public void shouldKeepReAssignedStandbyTaskInStateUpdater() { final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); when(stateUpdater.getTasks()).thenReturn(mkSet(reAssignedStandbyTask)); - expect(standbyTaskCreator.createTasks(Collections.emptyMap())).andReturn(emptySet()); - replay(standbyTaskCreator); taskManager.handleAssignment( Collections.emptyMap(), mkMap(mkEntry(reAssignedStandbyTask.id(), reAssignedStandbyTask.inputPartitions())) ); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).createTasks(consumer, Collections.emptyMap()); + Mockito.verify(standbyTaskCreator).createTasks(Collections.emptyMap()); } @Test @@ -487,20 +468,18 @@ public void shouldAssignMultipleTasksInStateUpdater() { final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); when(stateUpdater.getTasks()).thenReturn(mkSet(activeTaskToClose, standbyTaskToRecycle)); - expect(standbyTaskCreator.createTasks(Collections.emptyMap())).andReturn(emptySet()); - replay(standbyTaskCreator); taskManager.handleAssignment( mkMap(mkEntry(standbyTaskToRecycle.id(), standbyTaskToRecycle.inputPartitions())), Collections.emptyMap() ); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).createTasks(consumer, Collections.emptyMap()); Mockito.verify(stateUpdater).remove(activeTaskToClose.id()); Mockito.verify(tasks).addPendingTaskToCloseClean(activeTaskToClose.id()); Mockito.verify(stateUpdater).remove(standbyTaskToRecycle.id()); Mockito.verify(tasks).addPendingTaskToRecycle(standbyTaskToRecycle.id(), standbyTaskToRecycle.inputPartitions()); + Mockito.verify(standbyTaskCreator).createTasks(Collections.emptyMap()); } @Test @@ -546,13 +525,11 @@ public void shouldCreateActiveTaskDuringAssignment() { final Map> tasksToBeCreated = mkMap( mkEntry(activeTaskToBeCreated.id(), activeTaskToBeCreated.inputPartitions())); when(activeTaskCreator.createTasks(consumer, tasksToBeCreated)).thenReturn(createdTasks); - expect(standbyTaskCreator.createTasks(Collections.emptyMap())).andReturn(emptySet()); - replay(standbyTaskCreator); taskManager.handleAssignment(tasksToBeCreated, Collections.emptyMap()); - verify(standbyTaskCreator); Mockito.verify(tasks).addPendingTaskToInit(createdTasks); + Mockito.verify(standbyTaskCreator).createTasks(Collections.emptyMap()); } @Test @@ -563,17 +540,15 @@ public void shouldCreateStandbyTaskDuringAssignment() { final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); final Set createdTasks = mkSet(standbyTaskToBeCreated); - expect(standbyTaskCreator.createTasks(mkMap( + when(standbyTaskCreator.createTasks(mkMap( mkEntry(standbyTaskToBeCreated.id(), standbyTaskToBeCreated.inputPartitions()))) - ).andReturn(createdTasks); - replay(standbyTaskCreator); + ).thenReturn(createdTasks); taskManager.handleAssignment( Collections.emptyMap(), mkMap(mkEntry(standbyTaskToBeCreated.id(), standbyTaskToBeCreated.inputPartitions())) ); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).createTasks(consumer, Collections.emptyMap()); Mockito.verify(tasks).addPendingTaskToInit(createdTasks); } @@ -589,21 +564,19 @@ public void shouldAssignActiveTaskInTasksRegistryToBeRecycledWithStateUpdaterEna final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); when(tasks.allTasks()).thenReturn(mkSet(activeTaskToRecycle)); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); - expect(standbyTaskCreator.createStandbyTaskFromActive(activeTaskToRecycle, activeTaskToRecycle.inputPartitions())) - .andReturn(recycledStandbyTask); - expect(standbyTaskCreator.createTasks(Collections.emptyMap())).andReturn(emptySet()); - replay(standbyTaskCreator); + when(standbyTaskCreator.createStandbyTaskFromActive(activeTaskToRecycle, activeTaskToRecycle.inputPartitions())) + .thenReturn(recycledStandbyTask); taskManager.handleAssignment( Collections.emptyMap(), mkMap(mkEntry(activeTaskToRecycle.id(), activeTaskToRecycle.inputPartitions())) ); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).closeAndRemoveTaskProducerIfNeeded(activeTaskToRecycle.id()); Mockito.verify(activeTaskCreator).createTasks(consumer, Collections.emptyMap()); Mockito.verify(activeTaskToRecycle).prepareCommit(); Mockito.verify(tasks).replaceActiveWithStandby(recycledStandbyTask); + Mockito.verify(standbyTaskCreator).createTasks(Collections.emptyMap()); } @Test @@ -614,7 +587,6 @@ public void shouldThrowDuringAssignmentIfStandbyTaskToRecycleIsFoundInTasksRegis final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); when(tasks.allTasks()).thenReturn(mkSet(standbyTaskToRecycle)); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); - replay(standbyTaskCreator); final IllegalStateException illegalStateException = assertThrows( IllegalStateException.class, @@ -625,7 +597,6 @@ public void shouldThrowDuringAssignmentIfStandbyTaskToRecycleIsFoundInTasksRegis ); assertEquals(illegalStateException.getMessage(), "Standby tasks should only be managed by the state updater"); - verify(standbyTaskCreator); Mockito.verifyNoInteractions(activeTaskCreator); } @@ -637,17 +608,15 @@ public void shouldAssignActiveTaskInTasksRegistryToBeClosedCleanlyWithStateUpdat final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); when(tasks.allTasks()).thenReturn(mkSet(activeTaskToClose)); - expect(standbyTaskCreator.createTasks(Collections.emptyMap())).andReturn(emptySet()); - replay(standbyTaskCreator); taskManager.handleAssignment(Collections.emptyMap(), Collections.emptyMap()); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).createTasks(consumer, Collections.emptyMap()); Mockito.verify(activeTaskCreator).closeAndRemoveTaskProducerIfNeeded(activeTaskToClose.id()); Mockito.verify(activeTaskToClose).prepareCommit(); Mockito.verify(activeTaskToClose).closeClean(); Mockito.verify(tasks).removeTask(activeTaskToClose); + Mockito.verify(standbyTaskCreator).createTasks(Collections.emptyMap()); } @Test @@ -658,7 +627,6 @@ public void shouldThrowDuringAssignmentIfStandbyTaskToCloseIsFoundInTasksRegistr final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); when(tasks.allTasks()).thenReturn(mkSet(standbyTaskToClose)); - replay(standbyTaskCreator); final IllegalStateException illegalStateException = assertThrows( IllegalStateException.class, @@ -666,7 +634,6 @@ public void shouldThrowDuringAssignmentIfStandbyTaskToCloseIsFoundInTasksRegistr ); assertEquals(illegalStateException.getMessage(), "Standby tasks should only be managed by the state updater"); - verify(standbyTaskCreator); Mockito.verifyNoInteractions(activeTaskCreator); } @@ -680,17 +647,15 @@ public void shouldAssignActiveTaskInTasksRegistryToUpdateInputPartitionsWithStat final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); when(tasks.allTasks()).thenReturn(mkSet(activeTaskToUpdateInputPartitions)); when(tasks.updateActiveTaskInputPartitions(activeTaskToUpdateInputPartitions, newInputPartitions)).thenReturn(true); - expect(standbyTaskCreator.createTasks(Collections.emptyMap())).andReturn(emptySet()); - replay(standbyTaskCreator); taskManager.handleAssignment( mkMap(mkEntry(activeTaskToUpdateInputPartitions.id(), newInputPartitions)), Collections.emptyMap() ); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).createTasks(consumer, Collections.emptyMap()); Mockito.verify(activeTaskToUpdateInputPartitions).updateInputPartitions(Mockito.eq(newInputPartitions), any()); + Mockito.verify(standbyTaskCreator).createTasks(Collections.emptyMap()); } @Test @@ -701,16 +666,14 @@ public void shouldResumeActiveRunningTaskInTasksRegistryWithStateUpdaterEnabled( final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); when(tasks.allTasks()).thenReturn(mkSet(activeTaskToResume)); - expect(standbyTaskCreator.createTasks(Collections.emptyMap())).andReturn(emptySet()); - replay(standbyTaskCreator); taskManager.handleAssignment( mkMap(mkEntry(activeTaskToResume.id(), activeTaskToResume.inputPartitions())), Collections.emptyMap() ); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).createTasks(consumer, Collections.emptyMap()); + Mockito.verify(standbyTaskCreator).createTasks(Collections.emptyMap()); } @Test @@ -721,19 +684,17 @@ public void shouldResumeActiveSuspendedTaskInTasksRegistryAndAddToStateUpdater() final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); when(tasks.allTasks()).thenReturn(mkSet(activeTaskToResume)); - expect(standbyTaskCreator.createTasks(Collections.emptyMap())).andReturn(emptySet()); - replay(standbyTaskCreator); taskManager.handleAssignment( mkMap(mkEntry(activeTaskToResume.id(), activeTaskToResume.inputPartitions())), Collections.emptyMap() ); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).createTasks(consumer, Collections.emptyMap()); Mockito.verify(activeTaskToResume).resume(); Mockito.verify(stateUpdater).add(activeTaskToResume); Mockito.verify(tasks).removeTask(activeTaskToResume); + Mockito.verify(standbyTaskCreator).createTasks(Collections.emptyMap()); } @Test @@ -745,7 +706,6 @@ public void shouldThrowDuringAssignmentIfStandbyTaskToUpdateInputPartitionsIsFou final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); when(tasks.allTasks()).thenReturn(mkSet(standbyTaskToUpdateInputPartitions)); - replay(standbyTaskCreator); final IllegalStateException illegalStateException = assertThrows( IllegalStateException.class, @@ -756,7 +716,6 @@ public void shouldThrowDuringAssignmentIfStandbyTaskToUpdateInputPartitionsIsFou ); assertEquals(illegalStateException.getMessage(), "Standby tasks should only be managed by the state updater"); - verify(standbyTaskCreator); Mockito.verifyNoInteractions(activeTaskCreator); } @@ -771,20 +730,18 @@ public void shouldAssignMultipleTasksInTasksRegistryWithStateUpdaterEnabled() { final TasksRegistry tasks = Mockito.mock(TasksRegistry.class); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); when(tasks.allTasks()).thenReturn(mkSet(activeTaskToClose)); - expect(standbyTaskCreator.createTasks(Collections.emptyMap())).andReturn(emptySet()); - replay(standbyTaskCreator); taskManager.handleAssignment( mkMap(mkEntry(activeTaskToCreate.id(), activeTaskToCreate.inputPartitions())), Collections.emptyMap() ); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).createTasks( consumer, mkMap(mkEntry(activeTaskToCreate.id(), activeTaskToCreate.inputPartitions())) ); Mockito.verify(activeTaskToClose).closeClean(); + Mockito.verify(standbyTaskCreator).createTasks(Collections.emptyMap()); } @Test @@ -850,13 +807,11 @@ public void shouldRecycleTasksRemovedFromStateUpdater() { taskManager = setUpTaskManager(StreamsConfigUtils.ProcessingMode.AT_LEAST_ONCE, tasks, true); when(activeTaskCreator.createActiveTaskFromStandby(task01, taskId01Partitions, consumer)).thenReturn(task01Converted); - expect(standbyTaskCreator.createStandbyTaskFromActive(eq(task00), eq(taskId00Partitions))) - .andStubReturn(task00Converted); - replay(standbyTaskCreator); + when(standbyTaskCreator.createStandbyTaskFromActive(task00, taskId00Partitions)) + .thenReturn(task00Converted); taskManager.checkStateUpdater(time.milliseconds(), noOpResetter); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).closeAndRemoveTaskProducerIfNeeded(any()); Mockito.verify(task00).suspend(); Mockito.verify(task01).suspend(); @@ -961,8 +916,8 @@ public void shouldHandleMultipleRemovedTasksFromStateUpdater() { when(stateUpdater.restoresActiveTasks()).thenReturn(true); when(activeTaskCreator.createActiveTaskFromStandby(taskToRecycle1, taskId01Partitions, consumer)) .thenReturn(convertedTask1); - expect(standbyTaskCreator.createStandbyTaskFromActive(eq(taskToRecycle0), eq(taskId00Partitions))) - .andStubReturn(convertedTask0); + when(standbyTaskCreator.createStandbyTaskFromActive(taskToRecycle0, taskId00Partitions)) + .thenReturn(convertedTask0); expect(consumer.assignment()).andReturn(emptySet()).anyTimes(); consumer.resume(anyObject()); expectLastCall().anyTimes(); @@ -977,11 +932,11 @@ public void shouldHandleMultipleRemovedTasksFromStateUpdater() { when(tasks.removePendingTaskToUpdateInputPartitions(taskToUpdateInputPartitions.id())).thenReturn(taskId04Partitions); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); taskManager.setMainConsumer(consumer); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.checkStateUpdater(time.milliseconds(), noOpResetter -> { }); - verify(standbyTaskCreator, consumer); + verify(consumer); Mockito.verify(activeTaskCreator, times(2)).closeAndRemoveTaskProducerIfNeeded(any()); Mockito.verify(convertedTask0).initializeIfNeeded(); Mockito.verify(convertedTask1).initializeIfNeeded(); @@ -1151,13 +1106,11 @@ public void shouldRecycleRestoredTask() { .inState(State.CREATED) .withInputPartitions(taskId00Partitions).build(); final TaskManager taskManager = setUpRecycleRestoredTask(statefulTask); - expect(standbyTaskCreator.createStandbyTaskFromActive(statefulTask, statefulTask.inputPartitions())) - .andStubReturn(standbyTask); - replay(standbyTaskCreator); + when(standbyTaskCreator.createStandbyTaskFromActive(statefulTask, statefulTask.inputPartitions())) + .thenReturn(standbyTask); taskManager.checkStateUpdater(time.milliseconds(), noOpResetter); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).closeAndRemoveTaskProducerIfNeeded(statefulTask.id()); Mockito.verify(statefulTask).suspend(); Mockito.verify(standbyTask).initializeIfNeeded(); @@ -1170,16 +1123,14 @@ public void shouldHandleExceptionThrownDuringConversionInRecycleRestoredTask() { .inState(State.RESTORING) .withInputPartitions(taskId00Partitions).build(); final TaskManager taskManager = setUpRecycleRestoredTask(statefulTask); - expect(standbyTaskCreator.createStandbyTaskFromActive(statefulTask, statefulTask.inputPartitions())) - .andThrow(new RuntimeException()); - replay(standbyTaskCreator); + when(standbyTaskCreator.createStandbyTaskFromActive(statefulTask, statefulTask.inputPartitions())) + .thenThrow(new RuntimeException()); assertThrows( StreamsException.class, () -> taskManager.checkStateUpdater(time.milliseconds(), noOpResetter) ); - verify(standbyTaskCreator); Mockito.verify(stateUpdater, never()).add(any()); Mockito.verify(statefulTask).closeDirty(); } @@ -1193,17 +1144,15 @@ public void shouldHandleExceptionThrownDuringTaskInitInRecycleRestoredTask() { .inState(State.CREATED) .withInputPartitions(taskId00Partitions).build(); final TaskManager taskManager = setUpRecycleRestoredTask(statefulTask); - expect(standbyTaskCreator.createStandbyTaskFromActive(statefulTask, statefulTask.inputPartitions())) - .andStubReturn(standbyTask); + when(standbyTaskCreator.createStandbyTaskFromActive(statefulTask, statefulTask.inputPartitions())) + .thenReturn(standbyTask); doThrow(StreamsException.class).when(standbyTask).initializeIfNeeded(); - replay(standbyTaskCreator); assertThrows( StreamsException.class, () -> taskManager.checkStateUpdater(time.milliseconds(), noOpResetter) ); - verify(standbyTaskCreator); Mockito.verify(stateUpdater, never()).add(any()); Mockito.verify(standbyTask).closeDirty(); } @@ -1371,8 +1320,8 @@ public void shouldHandleMultipleRestoredTasks() { .withInputPartitions(taskId04Partitions).build(); final TasksRegistry tasks = mock(TasksRegistry.class); final TaskManager taskManager = setUpTaskManager(ProcessingMode.AT_LEAST_ONCE, tasks, true); - expect(standbyTaskCreator.createStandbyTaskFromActive(taskToRecycle, taskToRecycle.inputPartitions())) - .andStubReturn(recycledStandbyTask); + when(standbyTaskCreator.createStandbyTaskFromActive(taskToRecycle, taskToRecycle.inputPartitions())) + .thenReturn(recycledStandbyTask); when(tasks.removePendingTaskToRecycle(taskToRecycle.id())).thenReturn(taskId01Partitions); when(tasks.removePendingTaskToRecycle( argThat(taskId -> !taskId.equals(taskToRecycle.id()))) @@ -1397,7 +1346,6 @@ public void shouldHandleMultipleRestoredTasks() { taskToCloseDirty, taskToUpdateInputPartitions )); - replay(standbyTaskCreator); taskManager.checkStateUpdater(time.milliseconds(), noOpResetter); @@ -1530,15 +1478,13 @@ public void shouldAddSubscribedTopicsFromAssignmentToTopologyMetadata() { mkEntry(taskId03, mkSet(t1p3)), mkEntry(taskId04, mkSet(t1p4)) ); - expect(standbyTaskCreator.createTasks(eq(standbyTasksAssignment))).andStubReturn(Collections.emptySet()); - replay(standbyTaskCreator); + when(standbyTaskCreator.createTasks(standbyTasksAssignment)).thenReturn(Collections.emptySet()); taskManager.handleAssignment(activeTasksAssignment, standbyTasksAssignment); Mockito.verify(topologyBuilder).addSubscribedTopicsFromAssignment(Mockito.eq(mkSet(t1p1, t1p2, t2p2)), Mockito.anyString()); Mockito.verify(topologyBuilder, never()).addSubscribedTopicsFromAssignment(Mockito.eq(mkSet(t1p3, t1p4)), Mockito.anyString()); Mockito.verify(activeTaskCreator).createTasks(any(), Mockito.eq(activeTasksAssignment)); - verify(standbyTaskCreator); } @Test @@ -1731,8 +1677,7 @@ public void shouldComputeOffsetSumFromCheckpointFileForUninitializedTask() throw taskManager.handleRebalanceStart(singleton("topic")); final StateMachineTask uninitializedTask = new StateMachineTask(taskId00, taskId00Partitions, true); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singleton(uninitializedTask)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); - replay(standbyTaskCreator); + taskManager.handleAssignment(taskId00Assignment, emptyMap()); assertThat(uninitializedTask.state(), is(State.CREATED)); @@ -1756,9 +1701,9 @@ public void shouldComputeOffsetSumFromCheckpointFileForClosedTask() throws Excep final StateMachineTask closedTask = new StateMachineTask(taskId00, taskId00Partitions, true); taskManager.handleRebalanceStart(singleton("topic")); + when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singleton(closedTask)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); - replay(standbyTaskCreator); + taskManager.handleAssignment(taskId00Assignment, emptyMap()); closedTask.suspend(); @@ -1820,7 +1765,6 @@ public void shouldCloseActiveUnassignedSuspendedTasksWhenClosingRevokedTasks() { expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singletonList(task00)); expectLastCall(); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(emptyList()); // `handleRevocation` consumer.commitSync(offsets); @@ -1830,7 +1774,7 @@ public void shouldCloseActiveUnassignedSuspendedTasksWhenClosingRevokedTasks() { consumer.commitSync(offsets); expectLastCall(); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -1859,9 +1803,8 @@ public void closeClean() { expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singletonList(task00)); expectLastCall(); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(emptyList()); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, emptyMap()); taskManager.handleRevocation(taskId00Partitions); @@ -1888,7 +1831,7 @@ public void shouldCloseActiveTasksWhenHandlingLostTasks() throws Exception { // `handleAssignment` expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(eq(taskId01Assignment))).andStubReturn(singletonList(task01)); + when(standbyTaskCreator.createTasks(taskId01Assignment)).thenReturn(singletonList(task01)); makeTaskFolders(taskId00.toString(), taskId01.toString()); expectLockObtainedFor(taskId00, taskId01); @@ -1901,7 +1844,7 @@ public void shouldCloseActiveTasksWhenHandlingLostTasks() throws Exception { taskManager.handleRebalanceStart(emptySet()); assertThat(taskManager.lockedTaskDirectories(), Matchers.is(mkSet(taskId00, taskId01))); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, taskId01Assignment); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -1943,14 +1886,13 @@ public void shouldThrowWhenHandlingClosingTasksOnProducerCloseError() { // `handleAssignment` expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(emptyList()); // `handleAssignment` consumer.commitSync(offsets); expectLastCall(); doThrow(new RuntimeException("KABOOM!")).when(activeTaskCreator).closeAndRemoveTaskProducerIfNeeded(taskId00); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -2021,9 +1963,8 @@ public void postCommit(final boolean enforceCheckpoint) { // `handleAssignment` expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); expect(consumer.assignment()).andReturn(taskId00Partitions); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), tp -> assertThat(tp, is(empty()))), is(true)); @@ -2059,9 +2000,8 @@ public void suspend() { expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); expect(consumer.assignment()).andReturn(taskId00Partitions); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), tp -> assertThat(tp, is(empty()))), is(true)); @@ -2094,14 +2034,12 @@ public void shouldCommitNonCorruptedTasksOnTaskCorruptedException() { // `handleAssignment` when(activeTaskCreator.createTasks(any(), Mockito.eq(assignment))) .thenReturn(asList(corruptedTask, nonCorruptedTask)); - expect(standbyTaskCreator.createTasks(anyObject())) - .andStubReturn(Collections.emptySet()); expectRestoreToBeCompleted(consumer); expect(consumer.assignment()).andReturn(taskId00Partitions); // check that we should not commit empty map either consumer.commitSync(eq(emptyMap())); expectLastCall().andStubThrow(new AssertionError("should not invoke commitSync when offset map is empty")); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), tp -> assertThat(tp, is(empty()))), is(true)); @@ -2136,10 +2074,8 @@ public void shouldNotCommitNonRunningNonCorruptedTasks() { // `handleAssignment` when(activeTaskCreator.createTasks(any(), Mockito.eq(assignment))) .thenReturn(asList(corruptedTask, nonRunningNonCorruptedTask)); - expect(standbyTaskCreator.createTasks(anyObject())) - .andStubReturn(Collections.emptySet()); expect(consumer.assignment()).andReturn(taskId00Partitions); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(assignment, emptyMap()); @@ -2169,13 +2105,13 @@ public Map prepareCommit() { }; // handleAssignment - expect(standbyTaskCreator.createTasks(eq(taskId00Assignment))).andStubReturn(singleton(corruptedStandby)); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId01Assignment))) .thenReturn(singleton(runningNonCorruptedActive)); + when(standbyTaskCreator.createTasks(taskId00Assignment)).thenReturn(singleton(corruptedStandby)); expectRestoreToBeCompleted(consumer); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId01Assignment, taskId00Assignment); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -2213,13 +2149,12 @@ public void shouldNotAttemptToCommitInHandleCorruptedDuringARebalance() { assignment.putAll(taskId01Assignment); when(activeTaskCreator.createTasks(any(), Mockito.eq(assignment))) .thenReturn(asList(corruptedActive, uncorruptedActive)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); expectRestoreToBeCompleted(consumer); expect(consumer.assignment()).andStubReturn(union(HashSet::new, taskId00Partitions, taskId01Partitions)); - replay(standbyTaskCreator, consumer, stateDirectory, stateManager); + replay(consumer, stateDirectory, stateManager); uncorruptedActive.setCommittableOffsetsAndMetadata(offsets); @@ -2266,7 +2201,6 @@ public void markChangelogAsCorrupted(final Collection partitions assignment.putAll(taskId01Assignment); when(activeTaskCreator.createTasks(any(), Mockito.eq(assignment))) .thenReturn(asList(corruptedActive, uncorruptedActive)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); expectRestoreToBeCompleted(consumer); @@ -2275,7 +2209,7 @@ public void markChangelogAsCorrupted(final Collection partitions expect(consumer.assignment()).andStubReturn(union(HashSet::new, taskId00Partitions, taskId01Partitions)); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -2343,7 +2277,6 @@ public void markChangelogAsCorrupted(final Collection partitions assignment.putAll(taskId01Assignment); when(activeTaskCreator.createTasks(any(), Mockito.eq(assignment))) .thenReturn(asList(corruptedActiveTask, uncorruptedActiveTask)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); expectRestoreToBeCompleted(consumer); @@ -2354,7 +2287,7 @@ public void markChangelogAsCorrupted(final Collection partitions expect(consumer.assignment()).andStubReturn(union(HashSet::new, taskId00Partitions, taskId01Partitions)); - replay(standbyTaskCreator, consumer, stateManager); + replay(consumer, stateManager); taskManager.handleAssignment(assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -2426,13 +2359,12 @@ public void markChangelogAsCorrupted(final Collection partitions when(activeTaskCreator.createTasks(any(), Mockito.eq(assignmentActive))) .thenReturn(asList(revokedActiveTask, unrevokedActiveTaskWithCommitNeeded, unrevokedActiveTaskWithoutCommitNeeded)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); expectLastCall(); consumer.commitSync(expectedCommittedOffsets); expectLastCall().andThrow(new TimeoutException()); expect(consumer.assignment()).andStubReturn(union(HashSet::new, taskId00Partitions, taskId01Partitions, taskId02Partitions)); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(assignmentActive, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -2490,7 +2422,6 @@ public void markChangelogAsCorrupted(final Collection partitions when(activeTaskCreator.createTasks(any(), Mockito.eq(assignmentActive))) .thenReturn(asList(revokedActiveTask, unrevokedActiveTask, unrevokedActiveTaskWithoutCommitNeeded)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); final ConsumerGroupMetadata groupMetadata = new ConsumerGroupMetadata("appId"); expect(consumer.groupMetadata()).andReturn(groupMetadata); @@ -2499,7 +2430,7 @@ public void markChangelogAsCorrupted(final Collection partitions expect(consumer.assignment()).andStubReturn(union(HashSet::new, taskId00Partitions, taskId01Partitions, taskId02Partitions)); - replay(standbyTaskCreator, consumer, stateManager); + replay(consumer, stateManager); taskManager.handleAssignment(assignmentActive, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -2525,11 +2456,10 @@ public void shouldCloseStandbyUnassignedTasksWhenCreatingNewTasks() { final Task task00 = new StateMachineTask(taskId00, taskId00Partitions, false); expectRestoreToBeCompleted(consumer); - expect(standbyTaskCreator.createTasks(eq(taskId00Assignment))).andStubReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(eq(Collections.emptyMap()))).andStubReturn(Collections.emptySet()); + when(standbyTaskCreator.createTasks(taskId00Assignment)).thenReturn(singletonList(task00)); consumer.commitSync(Collections.emptyMap()); expectLastCall(); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(emptyMap(), taskId00Assignment); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -2550,9 +2480,8 @@ public void shouldAddNonResumedSuspendedTasks() { // expect these calls twice (because we're going to tryToCompleteRestoration twice) expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(eq(taskId01Assignment))).andReturn(singletonList(task01)).anyTimes(); - expect(standbyTaskCreator.createTasks(eq(Collections.emptyMap()))).andReturn(Collections.emptySet()); - replay(standbyTaskCreator, consumer); + when(standbyTaskCreator.createTasks(taskId01Assignment)).thenReturn(singletonList(task01)); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, taskId01Assignment); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -2575,8 +2504,7 @@ public void shouldUpdateInputPartitionsAfterRebalance() { // expect these calls twice (because we're going to tryToCompleteRestoration twice) expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -2601,8 +2529,7 @@ public void shouldAddNewActiveTasks() { consumer.resume(eq(emptySet())); expectLastCall(); when(activeTaskCreator.createTasks(any(), Mockito.eq(assignment))).thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(eq(emptyMap()))).andStubReturn(emptyList()); - replay(consumer, standbyTaskCreator); + replay(consumer); taskManager.handleAssignment(assignment, emptyMap()); @@ -2641,8 +2568,7 @@ public void initializeIfNeeded() { consumer.resume(eq(emptySet())); expectLastCall(); when(activeTaskCreator.createTasks(any(), Mockito.eq(assignment))).thenReturn(asList(task00, task01)); - expect(standbyTaskCreator.createTasks(eq(emptyMap()))).andStubReturn(emptyList()); - replay(consumer, standbyTaskCreator); + replay(consumer); taskManager.handleAssignment(assignment, emptyMap()); @@ -2680,8 +2606,7 @@ public void completeRestoration(final java.util.function.Consumer changelogPartitions() { when(activeTaskCreator.createTasks(any(), Mockito.eq(assignment))).thenReturn(singletonList(task00)); doThrow(new RuntimeException("whatever")) .when(activeTaskCreator).closeAndRemoveTaskProducerIfNeeded(taskId00); - expect(standbyTaskCreator.createTasks(eq(emptyMap()))).andStubReturn(emptyList()); - replay(standbyTaskCreator); taskManager.handleAssignment(assignment, emptyMap()); @@ -3134,8 +3050,6 @@ public Set changelogPartitions() { when(activeTaskCreator.createTasks(any(), Mockito.eq(assignment))).thenReturn(singletonList(task00)); doThrow(new RuntimeException("whatever")).when(activeTaskCreator).closeThreadProducerIfNeeded(); - expect(standbyTaskCreator.createTasks(eq(emptyMap()))).andStubReturn(emptyList()); - replay(standbyTaskCreator); taskManager.handleAssignment(assignment, emptyMap()); @@ -3259,8 +3173,6 @@ public void suspend() { when(activeTaskCreator.createTasks(any(), Mockito.eq(assignment))).thenReturn(asList(task00, task01, task02)); doThrow(new RuntimeException("whatever")).when(activeTaskCreator).closeAndRemoveTaskProducerIfNeeded(Mockito.any()); doThrow(new RuntimeException("whatever all")).when(activeTaskCreator).closeThreadProducerIfNeeded(); - expect(standbyTaskCreator.createTasks(eq(emptyMap()))).andStubReturn(emptyList()); - replay(standbyTaskCreator); taskManager.handleAssignment(assignment, emptyMap()); @@ -3305,7 +3217,7 @@ public void shouldCloseStandbyTasksOnShutdown() { final Task task00 = new StateMachineTask(taskId00, taskId00Partitions, false); // `handleAssignment` - expect(standbyTaskCreator.createTasks(eq(assignment))).andStubReturn(singletonList(task00)); + when(standbyTaskCreator.createTasks(assignment)).thenReturn(singletonList(task00)); // `tryToCompleteRestoration` expect(consumer.assignment()).andReturn(emptySet()); @@ -3316,7 +3228,7 @@ public void shouldCloseStandbyTasksOnShutdown() { consumer.commitSync(Collections.emptyMap()); expectLastCall(); - replay(consumer, standbyTaskCreator); + replay(consumer); taskManager.handleAssignment(emptyMap(), assignment); assertThat(task00.state(), is(Task.State.CREATED)); @@ -3411,8 +3323,7 @@ public void shouldInitializeNewActiveTasks() { expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))) .thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -3425,13 +3336,13 @@ public void shouldInitializeNewActiveTasks() { } @Test - public void shouldInitializeNewStandbyTasks() { + public void shouldInitialiseNewStandbyTasks() { final StateMachineTask task01 = new StateMachineTask(taskId01, taskId01Partitions, false); expectRestoreToBeCompleted(consumer); - expect(standbyTaskCreator.createTasks(eq(taskId01Assignment))).andStubReturn(singletonList(task01)); + when(standbyTaskCreator.createTasks(taskId01Assignment)).thenReturn(singletonList(task01)); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(emptyMap(), taskId01Assignment); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -3466,12 +3377,12 @@ public void shouldCommitActiveAndStandbyTasks() { expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))) .thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(eq(taskId01Assignment))) - .andStubReturn(singletonList(task01)); + when(standbyTaskCreator.createTasks(taskId01Assignment)) + .thenReturn(singletonList(task01)); consumer.commitSync(offsets); expectLastCall(); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, taskId01Assignment); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -3510,12 +3421,12 @@ public void shouldCommitProvidedTasksIfNeeded() { expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(assignmentActive))) .thenReturn(Arrays.asList(task00, task01, task02)); - expect(standbyTaskCreator.createTasks(eq(assignmentStandby))) - .andStubReturn(Arrays.asList(task03, task04, task05)); + when(standbyTaskCreator.createTasks(assignmentStandby)) + .thenReturn(Arrays.asList(task03, task04, task05)); consumer.commitSync(eq(emptyMap())); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(assignmentActive, assignmentStandby); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -3542,10 +3453,9 @@ public void shouldNotCommitOffsetsIfOnlyStandbyTasksAssigned() { final StateMachineTask task00 = new StateMachineTask(taskId00, taskId00Partitions, false); expectRestoreToBeCompleted(consumer); - expect(standbyTaskCreator.createTasks(eq(taskId00Assignment))).andStubReturn(singletonList(task00)); - expectLastCall(); + when(standbyTaskCreator.createTasks(taskId00Assignment)).thenReturn(singletonList(task00)); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(Collections.emptyMap(), taskId00Assignment); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -3568,10 +3478,10 @@ public void shouldNotCommitActiveAndStandbyTasksWhileRebalanceInProgress() throw expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))) .thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(eq(taskId01Assignment))) - .andStubReturn(singletonList(task01)); + when(standbyTaskCreator.createTasks(taskId01Assignment)) + .thenReturn(singletonList(task01)); - replay(standbyTaskCreator, stateDirectory, consumer); + replay(stateDirectory, consumer); taskManager.handleAssignment(taskId00Assignment, taskId01Assignment); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -3681,8 +3591,7 @@ public Map prepareCommit() { expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -3706,9 +3615,9 @@ public Map prepareCommit() { }; expectRestoreToBeCompleted(consumer); - expect(standbyTaskCreator.createTasks(eq(taskId01Assignment))).andStubReturn(singletonList(task01)); + when(standbyTaskCreator.createTasks(taskId01Assignment)).thenReturn(singletonList(task01)); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(emptyMap(), taskId01Assignment); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -3741,9 +3650,8 @@ public Map purgeableOffsets() { expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -3777,8 +3685,7 @@ public Map purgeableOffsets() { expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -3851,12 +3758,12 @@ public void shouldMaybeCommitAllActiveTasksThatNeedCommit() { expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(assignmentActive))) .thenReturn(asList(task00, task01, task02, task03)); - expect(standbyTaskCreator.createTasks(eq(assignmentStandby))) - .andStubReturn(singletonList(task04)); + when(standbyTaskCreator.createTasks(assignmentStandby)) + .thenReturn(singletonList(task04)); consumer.commitSync(expectedCommittedOffsets); expectLastCall(); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(assignmentActive, assignmentStandby); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -3895,8 +3802,7 @@ public void shouldProcessActiveTasks() { expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(assignment))) .thenReturn(Arrays.asList(task00, task01)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -4008,8 +3914,7 @@ public boolean process(final long wallClockTime) { expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -4034,8 +3939,7 @@ public boolean process(final long wallClockTime) { expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))) .thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -4062,8 +3966,7 @@ public boolean maybePunctuateStreamTime() { expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -4084,8 +3987,7 @@ public boolean maybePunctuateStreamTime() { expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -4111,9 +4013,7 @@ public boolean maybePunctuateSystemTime() { expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(anyObject())) - .andStubReturn(Collections.emptySet()); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(true)); @@ -4134,8 +4034,7 @@ public Set changelogPartitions() { }; when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(taskId00Assignment, emptyMap()); assertThat(taskManager.tryToCompleteRestoration(time.milliseconds(), null), is(false)); @@ -4153,11 +4052,10 @@ public void shouldHaveRemainingPartitionsUncleared() { expectRestoreToBeCompleted(consumer); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singletonList(task00)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); consumer.commitSync(offsets); expectLastCall(); - replay(standbyTaskCreator, consumer); + replay(consumer); try (final LogCaptureAppender appender = LogCaptureAppender.createAndRegister(TaskManager.class)) { LogCaptureAppender.setClassLoggerToDebug(TaskManager.class); @@ -4310,11 +4208,11 @@ private Map handleAssignment(final Map allActiveTasks = new HashSet<>(runningTasks); allActiveTasks.addAll(restoringTasks); - expect(standbyTaskCreator.createTasks(eq(standbyAssignment))).andStubReturn(standbyTasks); + when(standbyTaskCreator.createTasks(standbyAssignment)).thenReturn(standbyTasks); when(activeTaskCreator.createTasks(any(), Mockito.eq(allActiveTasksAssignment))).thenReturn(allActiveTasks); expectRestoreToBeCompleted(consumer); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(allActiveTasksAssignment, standbyAssignment); taskManager.tryToCompleteRestoration(time.milliseconds(), null); @@ -4550,8 +4448,7 @@ public void suspend() { final Map> assignment = new HashMap<>(taskId00Assignment); assignment.putAll(taskId01Assignment); when(activeTaskCreator.createTasks(any(), Mockito.eq(assignment))).thenReturn(asList(task00, task01)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(assignment, Collections.emptyMap()); @@ -4576,19 +4473,18 @@ public void shouldConvertActiveTaskToStandbyTask() { expect(standbyTask.id()).andStubReturn(taskId00); when(activeTaskCreator.createTasks(any(), Mockito.eq(taskId00Assignment))).thenReturn(singletonList(activeTask)); - expect(standbyTaskCreator.createTasks(anyObject())).andStubReturn(Collections.emptySet()); activeTask.prepareRecycle(); expectLastCall().once(); - expect(standbyTaskCreator.createStandbyTaskFromActive(anyObject(), eq(taskId00Partitions))).andReturn(standbyTask); + when(standbyTaskCreator.createStandbyTaskFromActive(Mockito.any(), Mockito.eq(taskId00Partitions))).thenReturn(standbyTask); - replay(activeTask, standbyTask, standbyTaskCreator, consumer); + replay(activeTask, standbyTask, consumer); taskManager.handleAssignment(taskId00Assignment, Collections.emptyMap()); taskManager.handleAssignment(Collections.emptyMap(), taskId00Assignment); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator).closeAndRemoveTaskProducerIfNeeded(taskId00); Mockito.verify(activeTaskCreator).createTasks(any(), Mockito.eq(emptyMap())); + Mockito.verify(standbyTaskCreator, times(2)).createTasks(Collections.emptyMap()); } @Test @@ -4601,18 +4497,17 @@ public void shouldConvertStandbyTaskToActiveTask() { final StreamTask activeTask = mock(StreamTask.class); when(activeTask.id()).thenReturn(taskId00); when(activeTask.inputPartitions()).thenReturn(taskId00Partitions); - expect(standbyTaskCreator.createTasks(eq(taskId00Assignment))).andReturn(singletonList(standbyTask)); + when(standbyTaskCreator.createTasks(taskId00Assignment)).thenReturn(singletonList(standbyTask)); when(activeTaskCreator.createActiveTaskFromStandby(Mockito.eq(standbyTask), Mockito.eq(taskId00Partitions), any())) .thenReturn(activeTask); - expect(standbyTaskCreator.createTasks(eq(Collections.emptyMap()))).andReturn(Collections.emptySet()); - replay(standbyTaskCreator, consumer); + replay(consumer); taskManager.handleAssignment(Collections.emptyMap(), taskId00Assignment); taskManager.handleAssignment(taskId00Assignment, Collections.emptyMap()); - verify(standbyTaskCreator); Mockito.verify(activeTaskCreator, times(2)).createTasks(any(), Mockito.eq(emptyMap())); + Mockito.verify(standbyTaskCreator).createTasks(Collections.emptyMap()); } @Test