skoppu22 commented on code in PR #360:
URL: https://github.com/apache/cassandra-sidecar/pull/360#discussion_r3436915257
##########
server/src/test/java/org/apache/cassandra/sidecar/job/OperationalJobManagerTest.java:
##########
@@ -166,4 +179,93 @@ protected Future<Void> executeInternal() throws
OperationalJobException
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
assertThat(tracker.get(jobId)).isNotNull();
}
+
+ @Test
+ void testCoordinatorCalledWhenJobRequiresCoordination() throws
InterruptedException
+ {
+ OperationalJobTracker tracker = new InMemoryOperationalJobTracker(4);
+ OperationalJobCoordinator coordinator =
mock(OperationalJobCoordinator.class);
+ when(coordinator.trySetActive(any(), any())).thenReturn(true);
+ OperationalJobManager manager = new OperationalJobManager(tracker,
coordinator, executorPool);
+ CountDownLatch latch = new CountDownLatch(1);
+
+ OperationalJob job = createCoordinatedJob(UUIDs.timeBased());
+ BiConsumer<OperationalJob, OperationalJobConflictException> onComplete
= (j, ex) -> {
+ assertThat(ex).isNull();
+ latch.countDown();
+ };
+
+ manager.trySubmitJob(job, onComplete, executorPool.service(),
SecondBoundConfiguration.parse("5s"));
+ assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
+ verify(coordinator).trySetActive(OperationType.MOVE, job.jobId());
+ }
+
+ @Test
+ void testConflictWhenCoordinatorReturnsFalse() throws InterruptedException
+ {
+ OperationalJobTracker tracker = new InMemoryOperationalJobTracker(4);
+ OperationalJobCoordinator coordinator =
mock(OperationalJobCoordinator.class);
+ when(coordinator.trySetActive(any(), any())).thenReturn(false);
+ OperationalJobManager manager = new OperationalJobManager(tracker,
coordinator, executorPool);
+ CountDownLatch latch = new CountDownLatch(1);
+
+ OperationalJob job = createCoordinatedJob(UUIDs.timeBased());
+ BiConsumer<OperationalJob, OperationalJobConflictException> onComplete
= (j, ex) -> {
+ assertThat(ex).isInstanceOf(OperationalJobConflictException.class);
+ assertThat(ex.getMessage()).contains("An active operation already
exists");
+ latch.countDown();
+ };
+
+ manager.trySubmitJob(job, onComplete, executorPool.service(),
SecondBoundConfiguration.parse("5s"));
+ assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
Review Comment:
Also verify that tracker doesn't have entry for this job after conflict
detected.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]