diff --git a/platform/util/src/com/intellij/util/concurrency/AppDelayQueue.java b/platform/util/src/com/intellij/util/concurrency/AppDelayQueue.java new file mode 100644 index 000000000000..49e92e21997f --- /dev/null +++ b/platform/util/src/com/intellij/util/concurrency/AppDelayQueue.java @@ -0,0 +1,72 @@ +/* + * Copyright 2000-2016 JetBrains s.r.o. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package com.intellij.util.concurrency; + +import com.intellij.openapi.diagnostic.Logger; + +import java.util.concurrent.DelayQueue; +import java.util.concurrent.atomic.AtomicBoolean; + +/** + * This class implements the global delayed queue which is used by + * {@link AppScheduledExecutorService} and {@link BoundedScheduledExecutorService}. + * It starts the background thread which polls the queue for tasks ready to run and sends them to the appropriate executor. + * The {@link #shutdown()} must be called before disposal. + */ +class AppDelayQueue extends DelayQueue { + private static final Logger LOG = Logger.getInstance("#com.intellij.util.concurrency.AppDelayQueue"); + private final Thread scheduledToPooledTransferer; + private final AtomicBoolean shutdown = new AtomicBoolean(); + + AppDelayQueue() { + /** this thread takes the ready-to-execute scheduled tasks off the queue and passes them for immediate execution to {@link SchedulingWrapper#backendExecutorService} */ + scheduledToPooledTransferer = new Thread(new Runnable() { + @Override + public void run() { + while (!shutdown.get()) { + try { + final SchedulingWrapper.MyScheduledFutureTask task = take(); + if (LOG.isTraceEnabled()) { + LOG.trace("Took "+BoundedTaskExecutor.info(task)); + } + task.getBackendExecutorService().execute(task); + } + catch (InterruptedException e) { + if (!shutdown.get()) { + LOG.error(e); + } + } + } + LOG.debug("scheduledToPooledTransferer Stopped"); + } + }, "Periodic tasks thread"); + scheduledToPooledTransferer.start(); + } + + void shutdown() { + if (shutdown.getAndSet(true)) { + throw new IllegalStateException("Already shutdown"); + } + scheduledToPooledTransferer.interrupt(); + + try { + scheduledToPooledTransferer.join(); + } + catch (Exception e) { + throw new RuntimeException(e); + } + } +} diff --git a/platform/util/src/com/intellij/util/concurrency/AppScheduledExecutorService.java b/platform/util/src/com/intellij/util/concurrency/AppScheduledExecutorService.java index 9960427c2642..d5cbccfc4976 100644 --- a/platform/util/src/com/intellij/util/concurrency/AppScheduledExecutorService.java +++ b/platform/util/src/com/intellij/util/concurrency/AppScheduledExecutorService.java @@ -44,7 +44,7 @@ public class AppScheduledExecutorService extends SchedulingWrapper { } AppScheduledExecutorService() { - super(new BackendThreadPoolExecutor()); + super(new BackendThreadPoolExecutor(), new AppDelayQueue()); ((BackendThreadPoolExecutor)backendExecutorService).doSetThreadFactory(new ThreadFactory() { private final AtomicInteger counter = new AtomicInteger(); @NotNull @@ -89,14 +89,15 @@ public class AppScheduledExecutorService extends SchedulingWrapper { ((BackendThreadPoolExecutor)backendExecutorService).doShutdown(); } + @NotNull @Override List doShutdownNow() { return ContainerUtil.concat(super.doShutdownNow(), ((BackendThreadPoolExecutor)backendExecutorService).doShutdownNow()); } public void shutdownAppScheduledExecutorService() { + delayQueue.shutdown(); // shutdown delay queue first to avoid rejected execution exceptions in Alarm doShutdown(); - shutdownGlobalQueue(); } public int getBackendPoolExecutorSize() { diff --git a/platform/util/src/com/intellij/util/concurrency/BoundedScheduledExecutorService.java b/platform/util/src/com/intellij/util/concurrency/BoundedScheduledExecutorService.java index 2fe125d071da..b0ccd6008e7b 100644 --- a/platform/util/src/com/intellij/util/concurrency/BoundedScheduledExecutorService.java +++ b/platform/util/src/com/intellij/util/concurrency/BoundedScheduledExecutorService.java @@ -30,7 +30,8 @@ import java.util.concurrent.TimeUnit; */ class BoundedScheduledExecutorService extends SchedulingWrapper { BoundedScheduledExecutorService(@NotNull ExecutorService backendExecutor, int maxSimultaneousTasks) { - super(new BoundedTaskExecutor(backendExecutor, maxSimultaneousTasks)); + super(new BoundedTaskExecutor(backendExecutor, maxSimultaneousTasks), + ((AppScheduledExecutorService)AppExecutorUtil.getAppScheduledExecutorService()).delayQueue); assert !(backendExecutor instanceof ScheduledExecutorService) : "backendExecutor is already ScheduledExecutorService: " + backendExecutor; } diff --git a/platform/util/src/com/intellij/util/concurrency/SchedulingWrapper.java b/platform/util/src/com/intellij/util/concurrency/SchedulingWrapper.java index 2b016e92d172..605cdb94d4fc 100644 --- a/platform/util/src/com/intellij/util/concurrency/SchedulingWrapper.java +++ b/platform/util/src/com/intellij/util/concurrency/SchedulingWrapper.java @@ -17,7 +17,6 @@ package com.intellij.util.concurrency; import com.intellij.openapi.diagnostic.Logger; import com.intellij.openapi.util.Condition; -import com.intellij.openapi.util.EmptyRunnable; import com.intellij.util.IncorrectOperationException; import com.intellij.util.containers.ContainerUtil; import org.jetbrains.annotations.NotNull; @@ -37,47 +36,17 @@ import java.util.concurrent.atomic.AtomicLong; class SchedulingWrapper implements ScheduledExecutorService { private static final Logger LOG = Logger.getInstance("#com.intellij.util.concurrency.SchedulingWrapper"); private final AtomicBoolean shutdown = new AtomicBoolean(); - private static final MyScheduledFutureTask TOMB = new MyScheduledFutureTask(); - private static final Thread scheduledToPooledTransferer; @NotNull final ExecutorService backendExecutorService; + final AppDelayQueue delayQueue; - SchedulingWrapper(@NotNull final ExecutorService backendExecutorService) { + SchedulingWrapper(@NotNull final ExecutorService backendExecutorService, @NotNull AppDelayQueue delayQueue) { + this.delayQueue = delayQueue; if (backendExecutorService instanceof ScheduledExecutorService) { throw new IllegalArgumentException("backendExecutorService: "+backendExecutorService+" is already ScheduledExecutorService"); } this.backendExecutorService = backendExecutorService; } - static { - /** this thread takes the ready-to-execute scheduled tasks off the queue and passes them for immediate execution to {@link #backendExecutorService} */ - scheduledToPooledTransferer = new Thread(new Runnable() { - @Override - public void run() { - while (true) { - try { - final MyScheduledFutureTask task = delayQueue.take(); - if (task == TOMB) { - break; - } - if (LOG.isTraceEnabled()) { - LOG.trace("Took "+BoundedTaskExecutor.info(task)); - } - task.myBackendExecutorService.execute(task); - } - catch (InterruptedException e) { - throw new RuntimeException(e); - } - } - LOG.debug("scheduledToPooledTransferer Stopped"); - } - }, "Periodic tasks thread"); - scheduledToPooledTransferer.start(); - } - - private static class MyDelayQueue extends DelayQueue { - } - private static final MyDelayQueue delayQueue = new MyDelayQueue(); - @NotNull @Override public List shutdownNow() { @@ -98,12 +67,13 @@ class SchedulingWrapper implements ScheduledExecutorService { } } + @NotNull List doShutdownNow() { doShutdown(); // shutdown me first to avoid further delayQueue offers List result = ContainerUtil.filter(delayQueue, new Condition() { @Override public boolean value(MyScheduledFutureTask task) { - if (task.myBackendExecutorService == backendExecutorService) { + if (task.getBackendExecutorService() == backendExecutorService) { task.cancel(false); return true; } @@ -114,23 +84,10 @@ class SchedulingWrapper implements ScheduledExecutorService { if (LOG.isTraceEnabled()) { LOG.trace("Shutdown. Drained tasks: "+result); } + //noinspection unchecked return (List)result; } - static void shutdownGlobalQueue() { - if (!scheduledToPooledTransferer.isAlive()) { - throw new IllegalStateException("Already shutdown"); - } - - delayQueue.offer(TOMB); - try { - scheduledToPooledTransferer.join(); - } - catch (Exception e) { - throw new RuntimeException(e); - } - } - @Override public boolean isShutdown() { return shutdown.get(); @@ -146,12 +103,11 @@ class SchedulingWrapper implements ScheduledExecutorService { return isTerminated(); } - private static class MyScheduledFutureTask extends FutureTask implements RunnableScheduledFuture { + class MyScheduledFutureTask extends FutureTask implements RunnableScheduledFuture { /** * Sequence number to break ties FIFO */ private final long sequenceNumber; - private final ExecutorService myBackendExecutorService; /** * The time the task is enabled to execute in nanoTime units @@ -166,20 +122,11 @@ class SchedulingWrapper implements ScheduledExecutorService { */ private final long period; - // fake ctr for TOMB - private MyScheduledFutureTask() { - super(EmptyRunnable.getInstance(), null); - myBackendExecutorService = null; - period = 0; - sequenceNumber = 0; - } - /** * Creates a one-shot action with given nanoTime-based trigger time. */ - private MyScheduledFutureTask(@NotNull ExecutorService backendExecutorService, @NotNull Runnable r, V result, long ns) { + private MyScheduledFutureTask(@NotNull Runnable r, V result, long ns) { super(r, result); - myBackendExecutorService = backendExecutorService; time = ns; period = 0; sequenceNumber = sequencer.getAndIncrement(); @@ -188,9 +135,8 @@ class SchedulingWrapper implements ScheduledExecutorService { /** * Creates a periodic action with given nano time and period. */ - private MyScheduledFutureTask(@NotNull ExecutorService backendExecutorService, @NotNull Runnable r, V result, long ns, long period) { + private MyScheduledFutureTask(@NotNull Runnable r, V result, long ns, long period) { super(r, result); - myBackendExecutorService = backendExecutorService; time = ns; this.period = period; sequenceNumber = sequencer.getAndIncrement(); @@ -199,9 +145,8 @@ class SchedulingWrapper implements ScheduledExecutorService { /** * Creates a one-shot action with given nanoTime-based trigger time. */ - private MyScheduledFutureTask(@NotNull ExecutorService backendExecutorService, @NotNull Callable callable, long ns) { + private MyScheduledFutureTask(@NotNull Callable callable, long ns) { super(callable); - myBackendExecutorService = backendExecutorService; time = ns; period = 0; sequenceNumber = sequencer.getAndIncrement(); @@ -275,7 +220,7 @@ class SchedulingWrapper implements ScheduledExecutorService { LOG.trace("Executing " + BoundedTaskExecutor.info(this)); } boolean periodic = isPeriodic(); - if (myBackendExecutorService.isShutdown()) { + if (backendExecutorService.isShutdown()) { cancel(false); } else if (!periodic) { @@ -291,6 +236,11 @@ class SchedulingWrapper implements ScheduledExecutorService { public String toString() { return "Delay: " + getDelay(TimeUnit.MILLISECONDS) + "ms; " + BoundedTaskExecutor.info(this); } + + @NotNull + ExecutorService getBackendExecutorService() { + return backendExecutorService; + } } /** @@ -302,7 +252,7 @@ class SchedulingWrapper implements ScheduledExecutorService { /** * Returns the trigger time of a delayed action. */ - private static long triggerTime(@NotNull MyDelayQueue queue, long delay, TimeUnit unit) { + private static long triggerTime(@NotNull AppDelayQueue queue, long delay, TimeUnit unit) { return triggerTime(queue, unit.toNanos(delay < 0 ? 0 : delay)); } @@ -313,7 +263,7 @@ class SchedulingWrapper implements ScheduledExecutorService { /** * Returns the trigger time of a delayed action. */ - private static long triggerTime(@NotNull MyDelayQueue queue, long delay) { + private static long triggerTime(@NotNull AppDelayQueue queue, long delay) { return now() + (delay < Long.MAX_VALUE >> 1 ? delay : overflowFree(queue, delay)); } @@ -324,7 +274,7 @@ class SchedulingWrapper implements ScheduledExecutorService { * not yet been, while some other task is added with a delay of * Long.MAX_VALUE. */ - private static long overflowFree(@NotNull MyDelayQueue queue, long delay) { + private static long overflowFree(@NotNull AppDelayQueue queue, long delay) { Delayed head = queue.peek(); if (head != null) { long headDelay = head.getDelay(TimeUnit.NANOSECONDS); @@ -340,7 +290,7 @@ class SchedulingWrapper implements ScheduledExecutorService { public ScheduledFuture schedule(@NotNull Runnable command, long delay, @NotNull TimeUnit unit) { - MyScheduledFutureTask t = new MyScheduledFutureTask(backendExecutorService, command, null, triggerTime(delayQueue, delay, unit)); + MyScheduledFutureTask t = new MyScheduledFutureTask(command, null, triggerTime(delayQueue, delay, unit)); return delayedExecute(t); } @@ -361,7 +311,7 @@ class SchedulingWrapper implements ScheduledExecutorService { public ScheduledFuture schedule(@NotNull Callable callable, long delay, @NotNull TimeUnit unit) { - MyScheduledFutureTask t = new MyScheduledFutureTask(backendExecutorService, callable, triggerTime(delayQueue, delay, unit)); + MyScheduledFutureTask t = new MyScheduledFutureTask(callable, triggerTime(delayQueue, delay, unit)); return delayedExecute(t); } @@ -383,7 +333,7 @@ class SchedulingWrapper implements ScheduledExecutorService { if (delay <= 0) { throw new IllegalArgumentException("delay must be positive but got: "+delay); } - MyScheduledFutureTask sft = new MyScheduledFutureTask(backendExecutorService, command, + MyScheduledFutureTask sft = new MyScheduledFutureTask(command, null, triggerTime(delayQueue, initialDelay, unit), unit.toNanos(-delay));