This is an automated email from the ASF dual-hosted git repository. qiaojialin pushed a commit to branch add_comment in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 5027a0b0b20cf3f426b154b843a54a4d6947a694 Author: qiaojialin <[email protected]> AuthorDate: Wed Apr 27 21:25:32 2022 +0800 spotless --- .../org/apache/iotdb/db/mpp/execution/Driver.java | 1 - .../org/apache/iotdb/db/mpp/execution/IDriver.java | 18 ++++----- .../db/mpp/schedule/AbstractDriverThread.java | 3 +- .../iotdb/db/mpp/schedule/DriverScheduler.java | 14 ++----- .../iotdb/db/mpp/schedule/IDriverScheduler.java | 3 +- .../iotdb/db/mpp/schedule/ITaskScheduler.java | 15 +++----- .../iotdb/db/mpp/schedule/task/DriverTask.java | 10 ++--- .../db/mpp/schedule/DefaultTaskSchedulerTest.java | 25 ++++-------- .../iotdb/db/mpp/schedule/DriverSchedulerTest.java | 12 ++---- .../DriverTaskTimeoutSentinelThreadTest.java | 45 ++++++++-------------- 10 files changed, 50 insertions(+), 96 deletions(-) diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/Driver.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/Driver.java index 27cb1d361f..ae220e4911 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/Driver.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/Driver.java @@ -47,7 +47,6 @@ import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static java.lang.Boolean.TRUE; import static org.apache.iotdb.db.mpp.operator.Operator.NOT_BLOCKED; - public abstract class Driver implements IDriver { protected static final Logger LOGGER = LoggerFactory.getLogger(Driver.class); diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/IDriver.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/IDriver.java index 4914760351..61e96dbe42 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/IDriver.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/IDriver.java @@ -25,8 +25,8 @@ import com.google.common.util.concurrent.ListenableFuture; import io.airlift.units.Duration; /** - * IDriver encapsulates some methods which are necessary for FragmentInstanceTaskExecutor to run a fragment - * instance + * IDriver encapsulates some methods which are necessary for FragmentInstanceTaskExecutor to run a + * fragment instance */ public interface IDriver { @@ -38,14 +38,14 @@ public interface IDriver { boolean isFinished(); /** - * run the IDriver for {@param duration} time slice, the time of this run is likely not - * to be equal to {@param duration}, the actual run time should be calculated by the caller + * run the IDriver for {@param duration} time slice, the time of this run is likely not to be + * equal to {@param duration}, the actual run time should be calculated by the caller * * @param duration how long should this IDriver run * @return the returned ListenableFuture<Void> is used to represent status of this processing if - * isDone() return true, meaning that this IDriver is not blocked and is ready for - * next processing. Otherwise, meaning that this IDriver is blocked and not ready for - * next processing. + * isDone() return true, meaning that this IDriver is not blocked and is ready for next + * processing. Otherwise, meaning that this IDriver is blocked and not ready for next + * processing. */ ListenableFuture<Void> processFor(Duration duration); @@ -66,8 +66,6 @@ public interface IDriver { */ void failed(Throwable t); - /** - * @return get SinkHandle of current IDriver - */ + /** @return get SinkHandle of current IDriver */ ISinkHandle getSinkHandle(); } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/schedule/AbstractDriverThread.java b/server/src/main/java/org/apache/iotdb/db/mpp/schedule/AbstractDriverThread.java index 13118452a8..c79304b0aa 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/schedule/AbstractDriverThread.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/schedule/AbstractDriverThread.java @@ -62,8 +62,7 @@ public abstract class AbstractDriverThread extends Thread implements Closeable { } /** Processing a task. */ - protected abstract void execute(DriverTask task) - throws InterruptedException, ExecutionException; + protected abstract void execute(DriverTask task) throws InterruptedException, ExecutionException; @Override public void close() throws IOException { diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/schedule/DriverScheduler.java b/server/src/main/java/org/apache/iotdb/db/mpp/schedule/DriverScheduler.java index 1089c3410c..63718faa6b 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/schedule/DriverScheduler.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/schedule/DriverScheduler.java @@ -73,12 +73,9 @@ public class DriverScheduler implements IDriverScheduler, IService { private DriverScheduler() { this.readyQueue = new L2PriorityQueue<>( - MAX_CAPACITY, - new DriverTask.SchedulePriorityComparator(), - new DriverTask()); + MAX_CAPACITY, new DriverTask.SchedulePriorityComparator(), new DriverTask()); this.timeoutQueue = - new L1PriorityQueue<>( - MAX_CAPACITY, new DriverTask.TimeoutComparator(), new DriverTask()); + new L1PriorityQueue<>(MAX_CAPACITY, new DriverTask.TimeoutComparator(), new DriverTask()); this.queryMap = new ConcurrentHashMap<>(); this.blockedTasks = Collections.synchronizedSet(new HashSet<>()); this.scheduler = new Scheduler(); @@ -91,8 +88,7 @@ public class DriverScheduler implements IDriverScheduler, IService { public void start() throws StartupException { for (int i = 0; i < WORKER_THREAD_NUM; i++) { AbstractDriverThread t = - new DriverTaskThread( - "Worker-Thread-" + i, workerGroups, readyQueue, scheduler); + new DriverTaskThread("Worker-Thread-" + i, workerGroups, readyQueue, scheduler); threads.add(t); t.start(); } @@ -124,9 +120,7 @@ public class DriverScheduler implements IDriverScheduler, IService { public void submitDrivers(QueryId queryId, List<IDriver> instances) { List<DriverTask> tasks = instances.stream() - .map( - v -> - new DriverTask(v, QUERY_TIMEOUT_MS, DriverTaskStatus.READY)) + .map(v -> new DriverTask(v, QUERY_TIMEOUT_MS, DriverTaskStatus.READY)) .collect(Collectors.toList()); queryMap .computeIfAbsent(queryId, v -> Collections.synchronizedSet(new HashSet<>())) diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/schedule/IDriverScheduler.java b/server/src/main/java/org/apache/iotdb/db/mpp/schedule/IDriverScheduler.java index b92f663eb2..30dc6aa38c 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/schedule/IDriverScheduler.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/schedule/IDriverScheduler.java @@ -44,7 +44,8 @@ public interface IDriverScheduler { void abortQuery(QueryId queryId); /** - * Abort all Drivers of the fragment instance. If the instance is not existed, nothing will happen. + * Abort all Drivers of the fragment instance. If the instance is not existed, nothing will + * happen. * * @param instanceId the id of the fragment instance to be aborted. */ diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/schedule/ITaskScheduler.java b/server/src/main/java/org/apache/iotdb/db/mpp/schedule/ITaskScheduler.java index 7f1b973c01..641aba8c7f 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/schedule/ITaskScheduler.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/schedule/ITaskScheduler.java @@ -25,16 +25,14 @@ import org.apache.iotdb.db.mpp.schedule.task.DriverTaskStatus; public interface ITaskScheduler { /** - * Switch a task from {@link DriverTaskStatus#BLOCKED} to {@link - * DriverTaskStatus#READY}. + * Switch a task from {@link DriverTaskStatus#BLOCKED} to {@link DriverTaskStatus#READY}. * * @param task the task to be switched. */ void blockedToReady(DriverTask task); /** - * Switch a task from {@link DriverTaskStatus#READY} to {@link - * DriverTaskStatus#RUNNING}. + * Switch a task from {@link DriverTaskStatus#READY} to {@link DriverTaskStatus#RUNNING}. * * @param task the task to be switched. * @return true if it's switched to the target status successfully, otherwise false. @@ -42,8 +40,7 @@ public interface ITaskScheduler { boolean readyToRunning(DriverTask task); /** - * Switch a task from {@link DriverTaskStatus#RUNNING} to {@link - * DriverTaskStatus#READY}. + * Switch a task from {@link DriverTaskStatus#RUNNING} to {@link DriverTaskStatus#READY}. * * @param task the task to be switched. * @param context the execution context of last running. @@ -51,8 +48,7 @@ public interface ITaskScheduler { void runningToReady(DriverTask task, ExecutionContext context); /** - * Switch a task from {@link DriverTaskStatus#RUNNING} to {@link - * DriverTaskStatus#BLOCKED}. + * Switch a task from {@link DriverTaskStatus#RUNNING} to {@link DriverTaskStatus#BLOCKED}. * * @param task the task to be switched. * @param context the execution context of last running. @@ -60,8 +56,7 @@ public interface ITaskScheduler { void runningToBlocked(DriverTask task, ExecutionContext context); /** - * Switch a task from {@link DriverTaskStatus#RUNNING} to {@link - * DriverTaskStatus#FINISHED}. + * Switch a task from {@link DriverTaskStatus#RUNNING} to {@link DriverTaskStatus#FINISHED}. * * @param task the task to be switched. * @param context the execution context of last running. diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/schedule/task/DriverTask.java b/server/src/main/java/org/apache/iotdb/db/mpp/schedule/task/DriverTask.java index 4844a0de89..e008f8df1c 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/schedule/task/DriverTask.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/schedule/task/DriverTask.java @@ -23,8 +23,8 @@ import org.apache.iotdb.db.mpp.common.FragmentInstanceId; import org.apache.iotdb.db.mpp.common.PlanFragmentId; import org.apache.iotdb.db.mpp.common.QueryId; import org.apache.iotdb.db.mpp.execution.IDriver; -import org.apache.iotdb.db.mpp.schedule.ExecutionContext; import org.apache.iotdb.db.mpp.schedule.DriverTaskThread; +import org.apache.iotdb.db.mpp.schedule.ExecutionContext; import org.apache.iotdb.db.mpp.schedule.queue.ID; import org.apache.iotdb.db.mpp.schedule.queue.IDIndexedAccessible; @@ -36,10 +36,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; -/** - * the scheduling element of {@link DriverTaskThread}. It wraps a single - * Driver. - */ +/** the scheduling element of {@link DriverTaskThread}. It wraps a single Driver. */ public class DriverTask implements IDIndexedAccessible { private DriverTaskID id; @@ -84,8 +81,7 @@ public class DriverTask implements IDIndexedAccessible { } public boolean isEndState() { - return status == DriverTaskStatus.ABORTED - || status == DriverTaskStatus.FINISHED; + return status == DriverTaskStatus.ABORTED || status == DriverTaskStatus.FINISHED; } public IDriver getFragmentInstance() { diff --git a/server/src/test/java/org/apache/iotdb/db/mpp/schedule/DefaultTaskSchedulerTest.java b/server/src/test/java/org/apache/iotdb/db/mpp/schedule/DefaultTaskSchedulerTest.java index 5689663133..213f742e69 100644 --- a/server/src/test/java/org/apache/iotdb/db/mpp/schedule/DefaultTaskSchedulerTest.java +++ b/server/src/test/java/org/apache/iotdb/db/mpp/schedule/DefaultTaskSchedulerTest.java @@ -80,8 +80,7 @@ public class DefaultTaskSchedulerTest { Assert.assertTrue(manager.getQueryMap().get(queryId).contains(testTask)); clear(); } - DriverTask testTask = - new DriverTask(mockDriver, 100L, DriverTaskStatus.BLOCKED); + DriverTask testTask = new DriverTask(mockDriver, 100L, DriverTaskStatus.BLOCKED); manager.getBlockedTasks().add(testTask); Set<DriverTask> taskSet = new HashSet<>(); taskSet.add(testTask); @@ -130,8 +129,7 @@ public class DefaultTaskSchedulerTest { Assert.assertTrue(manager.getQueryMap().get(queryId).contains(testTask)); clear(); } - DriverTask testTask = - new DriverTask(mockDriver, 100L, DriverTaskStatus.READY); + DriverTask testTask = new DriverTask(mockDriver, 100L, DriverTaskStatus.READY); Set<DriverTask> taskSet = new HashSet<>(); taskSet.add(testTask); manager.getQueryMap().put(queryId, taskSet); @@ -178,8 +176,7 @@ public class DefaultTaskSchedulerTest { Assert.assertTrue(manager.getQueryMap().get(queryId).contains(testTask)); clear(); } - DriverTask testTask = - new DriverTask(mockDriver, 100L, DriverTaskStatus.RUNNING); + DriverTask testTask = new DriverTask(mockDriver, 100L, DriverTaskStatus.RUNNING); Set<DriverTask> taskSet = new HashSet<>(); taskSet.add(testTask); manager.getQueryMap().put(queryId, taskSet); @@ -231,8 +228,7 @@ public class DefaultTaskSchedulerTest { Assert.assertTrue(manager.getQueryMap().get(queryId).contains(testTask)); clear(); } - DriverTask testTask = - new DriverTask(mockDriver, 100L, DriverTaskStatus.RUNNING); + DriverTask testTask = new DriverTask(mockDriver, 100L, DriverTaskStatus.RUNNING); Set<DriverTask> taskSet = new HashSet<>(); taskSet.add(testTask); manager.getQueryMap().put(queryId, taskSet); @@ -284,8 +280,7 @@ public class DefaultTaskSchedulerTest { Assert.assertTrue(manager.getQueryMap().get(queryId).contains(testTask)); clear(); } - DriverTask testTask = - new DriverTask(mockDriver, 100L, DriverTaskStatus.RUNNING); + DriverTask testTask = new DriverTask(mockDriver, 100L, DriverTaskStatus.RUNNING); Set<DriverTask> taskSet = new HashSet<>(); taskSet.add(testTask); manager.getQueryMap().put(queryId, taskSet); @@ -325,8 +320,7 @@ public class DefaultTaskSchedulerTest { }; for (DriverTaskStatus status : invalidStates) { DriverTask testTask1 = new DriverTask(mockDriver1, 100L, status); - DriverTask testTask2 = - new DriverTask(mockDriver2, 100L, DriverTaskStatus.BLOCKED); + DriverTask testTask2 = new DriverTask(mockDriver2, 100L, DriverTaskStatus.BLOCKED); Set<DriverTask> taskSet = new HashSet<>(); taskSet.add(testTask1); taskSet.add(testTask2); @@ -354,9 +348,7 @@ public class DefaultTaskSchedulerTest { } DriverTaskStatus[] validStates = new DriverTaskStatus[] { - DriverTaskStatus.RUNNING, - DriverTaskStatus.READY, - DriverTaskStatus.BLOCKED, + DriverTaskStatus.RUNNING, DriverTaskStatus.READY, DriverTaskStatus.BLOCKED, }; for (DriverTaskStatus status : validStates) { Mockito.reset(mockDriver1); @@ -366,8 +358,7 @@ public class DefaultTaskSchedulerTest { DriverTask testTask1 = new DriverTask(mockDriver1, 100L, status); - DriverTask testTask2 = - new DriverTask(mockDriver2, 100L, DriverTaskStatus.BLOCKED); + DriverTask testTask2 = new DriverTask(mockDriver2, 100L, DriverTaskStatus.BLOCKED); Set<DriverTask> taskSet = new HashSet<>(); taskSet.add(testTask1); taskSet.add(testTask2); diff --git a/server/src/test/java/org/apache/iotdb/db/mpp/schedule/DriverSchedulerTest.java b/server/src/test/java/org/apache/iotdb/db/mpp/schedule/DriverSchedulerTest.java index 07f3e67305..2dc325fb30 100644 --- a/server/src/test/java/org/apache/iotdb/db/mpp/schedule/DriverSchedulerTest.java +++ b/server/src/test/java/org/apache/iotdb/db/mpp/schedule/DriverSchedulerTest.java @@ -69,11 +69,9 @@ public class DriverSchedulerTest { Assert.assertEquals(2, manager.getQueryMap().get(queryId).size()); Assert.assertEquals(2, manager.getTimeoutQueue().size()); Assert.assertEquals(2, manager.getReadyQueue().size()); - DriverTask task1 = - manager.getTimeoutQueue().get(new DriverTaskID(instanceId1)); + DriverTask task1 = manager.getTimeoutQueue().get(new DriverTaskID(instanceId1)); Assert.assertNotNull(task1); - DriverTask task2 = - manager.getTimeoutQueue().get(new DriverTaskID(instanceId2)); + DriverTask task2 = manager.getTimeoutQueue().get(new DriverTaskID(instanceId2)); Assert.assertNotNull(task2); Assert.assertTrue(manager.getQueryMap().get(queryId).contains(task1)); Assert.assertTrue(manager.getQueryMap().get(queryId).contains(task2)); @@ -91,8 +89,7 @@ public class DriverSchedulerTest { Assert.assertEquals(3, manager.getQueryMap().get(queryId).size()); Assert.assertEquals(3, manager.getTimeoutQueue().size()); Assert.assertEquals(3, manager.getReadyQueue().size()); - DriverTask task3 = - manager.getTimeoutQueue().get(new DriverTaskID(instanceId3)); + DriverTask task3 = manager.getTimeoutQueue().get(new DriverTaskID(instanceId3)); Assert.assertNotNull(task3); Assert.assertTrue(manager.getQueryMap().get(queryId).contains(task3)); Assert.assertEquals(DriverTaskStatus.READY, task3.getStatus()); @@ -110,8 +107,7 @@ public class DriverSchedulerTest { Assert.assertEquals(1, manager.getQueryMap().get(queryId2).size()); Assert.assertEquals(4, manager.getTimeoutQueue().size()); Assert.assertEquals(4, manager.getReadyQueue().size()); - DriverTask task4 = - manager.getTimeoutQueue().get(new DriverTaskID(instanceId4)); + DriverTask task4 = manager.getTimeoutQueue().get(new DriverTaskID(instanceId4)); Assert.assertNotNull(task4); Assert.assertTrue(manager.getQueryMap().get(queryId2).contains(task4)); Assert.assertEquals(DriverTaskStatus.READY, task4.getStatus()); diff --git a/server/src/test/java/org/apache/iotdb/db/mpp/schedule/DriverTaskTimeoutSentinelThreadTest.java b/server/src/test/java/org/apache/iotdb/db/mpp/schedule/DriverTaskTimeoutSentinelThreadTest.java index 4ce669c1d2..9e6e7b3293 100644 --- a/server/src/test/java/org/apache/iotdb/db/mpp/schedule/DriverTaskTimeoutSentinelThreadTest.java +++ b/server/src/test/java/org/apache/iotdb/db/mpp/schedule/DriverTaskTimeoutSentinelThreadTest.java @@ -55,18 +55,15 @@ public class DriverTaskTimeoutSentinelThreadTest { PlanFragmentId fragmentId = new PlanFragmentId(queryId, 0); FragmentInstanceId instanceId = new FragmentInstanceId(fragmentId, "inst-0"); IndexedBlockingQueue<DriverTask> taskQueue = - new L1PriorityQueue<>( - 100, new DriverTask.TimeoutComparator(), new DriverTask()); + new L1PriorityQueue<>(100, new DriverTask.TimeoutComparator(), new DriverTask()); IDriver mockDriver = Mockito.mock(IDriver.class); Mockito.when(mockDriver.getInfo()).thenReturn(instanceId); AbstractDriverThread executor = - new DriverTaskThread( - "0", new ThreadGroup("timeout-test"), taskQueue, mockScheduler); + new DriverTaskThread("0", new ThreadGroup("timeout-test"), taskQueue, mockScheduler); // FINISHED status test - DriverTask testTask = - new DriverTask(mockDriver, 100L, DriverTaskStatus.FINISHED); + DriverTask testTask = new DriverTask(mockDriver, 100L, DriverTaskStatus.FINISHED); executor.execute(testTask); Assert.assertEquals(DriverTaskStatus.FINISHED, testTask.getStatus()); Mockito.verify(mockDriver, Mockito.never()).processFor(Mockito.any()); @@ -112,8 +109,7 @@ public class DriverTaskTimeoutSentinelThreadTest { return true; }); IndexedBlockingQueue<DriverTask> taskQueue = - new L1PriorityQueue<>( - 100, new DriverTask.TimeoutComparator(), new DriverTask()); + new L1PriorityQueue<>(100, new DriverTask.TimeoutComparator(), new DriverTask()); // Mock the instance with a cancelled future IDriver mockDriver = Mockito.mock(IDriver.class); @@ -125,10 +121,8 @@ public class DriverTaskTimeoutSentinelThreadTest { .thenReturn(Futures.immediateCancelledFuture()); AbstractDriverThread executor = - new DriverTaskThread( - "0", new ThreadGroup("timeout-test"), taskQueue, mockScheduler); - DriverTask testTask = - new DriverTask(mockDriver, 100L, DriverTaskStatus.READY); + new DriverTaskThread("0", new ThreadGroup("timeout-test"), taskQueue, mockScheduler); + DriverTask testTask = new DriverTask(mockDriver, 100L, DriverTaskStatus.READY); executor.execute(testTask); Mockito.verify(mockDriver, Mockito.times(1)).processFor(Mockito.any()); Assert.assertEquals( @@ -154,8 +148,7 @@ public class DriverTaskTimeoutSentinelThreadTest { return true; }); IndexedBlockingQueue<DriverTask> taskQueue = - new L1PriorityQueue<>( - 100, new DriverTask.TimeoutComparator(), new DriverTask()); + new L1PriorityQueue<>(100, new DriverTask.TimeoutComparator(), new DriverTask()); // Mock the instance with a cancelled future IDriver mockDriver = Mockito.mock(IDriver.class); @@ -166,10 +159,8 @@ public class DriverTaskTimeoutSentinelThreadTest { Mockito.when(mockDriver.processFor(Mockito.any())).thenReturn(Futures.immediateVoidFuture()); Mockito.when(mockDriver.isFinished()).thenReturn(true); AbstractDriverThread executor = - new DriverTaskThread( - "0", new ThreadGroup("timeout-test"), taskQueue, mockScheduler); - DriverTask testTask = - new DriverTask(mockDriver, 100L, DriverTaskStatus.READY); + new DriverTaskThread("0", new ThreadGroup("timeout-test"), taskQueue, mockScheduler); + DriverTask testTask = new DriverTask(mockDriver, 100L, DriverTaskStatus.READY); executor.execute(testTask); Mockito.verify(mockDriver, Mockito.times(1)).processFor(Mockito.any()); Assert.assertNull(testTask.getAbortCause()); @@ -194,8 +185,7 @@ public class DriverTaskTimeoutSentinelThreadTest { return true; }); IndexedBlockingQueue<DriverTask> taskQueue = - new L1PriorityQueue<>( - 100, new DriverTask.TimeoutComparator(), new DriverTask()); + new L1PriorityQueue<>(100, new DriverTask.TimeoutComparator(), new DriverTask()); // Mock the instance with a blocked future ListenableFuture<Void> mockFuture = Mockito.mock(ListenableFuture.class); @@ -217,10 +207,8 @@ public class DriverTaskTimeoutSentinelThreadTest { Mockito.when(mockDriver.processFor(Mockito.any())).thenReturn(mockFuture); Mockito.when(mockDriver.isFinished()).thenReturn(false); AbstractDriverThread executor = - new DriverTaskThread( - "0", new ThreadGroup("timeout-test"), taskQueue, mockScheduler); - DriverTask testTask = - new DriverTask(mockDriver, 100L, DriverTaskStatus.READY); + new DriverTaskThread("0", new ThreadGroup("timeout-test"), taskQueue, mockScheduler); + DriverTask testTask = new DriverTask(mockDriver, 100L, DriverTaskStatus.READY); executor.execute(testTask); Mockito.verify(mockDriver, Mockito.times(1)).processFor(Mockito.any()); Assert.assertNull(testTask.getAbortCause()); @@ -245,8 +233,7 @@ public class DriverTaskTimeoutSentinelThreadTest { return true; }); IndexedBlockingQueue<DriverTask> taskQueue = - new L1PriorityQueue<>( - 100, new DriverTask.TimeoutComparator(), new DriverTask()); + new L1PriorityQueue<>(100, new DriverTask.TimeoutComparator(), new DriverTask()); // Mock the instance with a ready future ListenableFuture<Void> mockFuture = Mockito.mock(ListenableFuture.class); @@ -268,10 +255,8 @@ public class DriverTaskTimeoutSentinelThreadTest { Mockito.when(mockDriver.processFor(Mockito.any())).thenReturn(mockFuture); Mockito.when(mockDriver.isFinished()).thenReturn(false); AbstractDriverThread executor = - new DriverTaskThread( - "0", new ThreadGroup("timeout-test"), taskQueue, mockScheduler); - DriverTask testTask = - new DriverTask(mockDriver, 100L, DriverTaskStatus.READY); + new DriverTaskThread("0", new ThreadGroup("timeout-test"), taskQueue, mockScheduler); + DriverTask testTask = new DriverTask(mockDriver, 100L, DriverTaskStatus.READY); executor.execute(testTask); Mockito.verify(mockDriver, Mockito.times(1)).processFor(Mockito.any()); Assert.assertNull(testTask.getAbortCause());
