execute requests sequentially even in pooled thread (a lot of clients depend on it)

This commit is contained in:
Alexey Kudravtsev
2016-01-22 12:45:42 +03:00
parent 98735f0474
commit bd9ab218a0
3 changed files with 59 additions and 12 deletions
@@ -15,7 +15,6 @@
*/
package com.intellij.util;
import com.intellij.concurrency.JobScheduler;
import com.intellij.openapi.Disposable;
import com.intellij.openapi.application.Application;
import com.intellij.openapi.application.ApplicationActivationListener;
@@ -68,9 +67,7 @@ public class Alarm implements Disposable {
myDisposed = true;
cancelAllRequests();
if (myExecutorService != JobScheduler.getScheduler()) {
myExecutorService.shutdownNow();
}
myExecutorService.shutdownNow();
}
}
@@ -84,9 +81,9 @@ public class Alarm implements Disposable {
*/
SWING_THREAD,
@Deprecated
/**
* The action will be executed on a dedicated single shared thread, one per IDEA instance.
* The actions should be very fast to avoid blocking other Alarm instances that need the same thread.
* @deprecated Use {@link #POOLED_THREAD} instead
*/
SHARED_THREAD,
@@ -97,9 +94,9 @@ public class Alarm implements Disposable {
*/
POOLED_THREAD,
@Deprecated
/**
* A dedicated new thread is created for this Alarm instance to run the requests. No time limits are placed on the request execution time.
* In general it's advised to avoid this option because it may lead to too many OS resources being used.
* @deprecated Use {@link #POOLED_THREAD} instead
*/
OWN_THREAD
}
@@ -124,8 +121,7 @@ public class Alarm implements Disposable {
public Alarm(@NotNull ThreadToUse threadToUse, @Nullable Disposable parentDisposable) {
myThreadToUse = threadToUse;
myExecutorService = threadToUse == ThreadToUse.POOLED_THREAD ? JobScheduler.getScheduler() :
// have to restrict the number of running tasks because otherwise the (implicit) contract of
myExecutorService = // have to restrict the number of running tasks because otherwise the (implicit) contract of
// "addRequests with the same delay are executed in order" will be broken
AppExecutorUtil.createBoundedScheduledExecutorService(1);
@@ -272,7 +268,7 @@ public class Alarm implements Disposable {
}
@TestOnly
public void waitForAllExecuted(long timeout, @NotNull TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException {
void waitForAllExecuted(long timeout, @NotNull TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException {
List<Request> requests;
synchronized (LOCK) {
requests = new ArrayList<Request>(myRequests);
@@ -18,11 +18,32 @@ package com.intellij.util;
import com.intellij.testFramework.PlatformTestCase;
import com.intellij.util.ui.UIUtil;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
public class AlarmTest extends PlatformTestCase {
public void testTwoAddsWithZeroDelayMustExecuteSequentially() throws Exception {
Alarm alarm = new Alarm(getTestRootDisposable());
assertRequestsExecuteSequentially(alarm);
}
public void testAlarmRequestsShouldExecuteSequentiallyEvenInPooledThread() throws Exception {
Alarm alarm = new Alarm(Alarm.ThreadToUse.POOLED_THREAD, getTestRootDisposable());
assertRequestsExecuteSequentially(alarm);
}
public void testAlarmRequestsShouldExecuteSequentiallyEveryWhere() throws Exception {
Alarm alarm = new Alarm(Alarm.ThreadToUse.OWN_THREAD, getTestRootDisposable());
assertRequestsExecuteSequentially(alarm);
}
public void testAlarmRequestsShouldExecuteSequentiallyAbsolutelyEveryWhere() throws Exception {
Alarm alarm = new Alarm(Alarm.ThreadToUse.SHARED_THREAD, getTestRootDisposable());
assertRequestsExecuteSequentially(alarm);
}
private static void assertRequestsExecuteSequentially(Alarm alarm) throws InterruptedException, ExecutionException, TimeoutException {
int N = 100000;
StringBuffer log = new StringBuffer(N*4);
StringBuilder expected = new StringBuilder(N * 4);
@@ -17,6 +17,7 @@ package com.intellij.util.ui.update;
import com.intellij.concurrency.JobScheduler;
import com.intellij.testFramework.UsefulTestCase;
import com.intellij.util.Alarm;
import com.intellij.util.TimeoutUtil;
import com.intellij.util.WaitFor;
import com.intellij.util.containers.ContainerUtil;
@@ -385,7 +386,36 @@ public class MergingUpdateQueueTest extends UsefulTestCase {
waitForExecution(queue);
assertEquals(expected.toString(), actual.toString());
}
public void testAddRequestsInPooledThreadDoNotExecuteConcurrently() throws InterruptedException {
int delay = 10;
MergingUpdateQueue queue = new MergingUpdateQueue("x", delay, true, null, getTestRootDisposable(), null, Alarm.ThreadToUse.POOLED_THREAD);
queue.setPassThrough(false);
CountDownLatch startedExecuting1 = new CountDownLatch(1);
CountDownLatch canContinue = new CountDownLatch(1);
queue.queue(new Update("1") {
@Override
public void run() {
startedExecuting1.countDown();
try {
canContinue.await();
}
catch (InterruptedException e) {
throw new RuntimeException(e);
}
}
});
assertTrue(startedExecuting1.await(10, TimeUnit.SECONDS));
CountDownLatch startedExecuting2 = new CountDownLatch(1);
queue.queue(new Update("2") {
@Override
public void run() {
startedExecuting2.countDown();
}
});
TimeoutUtil.sleep(delay + 1000);
canContinue.countDown();
assertTrue(startedExecuting2.await(10, TimeUnit.SECONDS));
}
}