give app scheduler a name for debug

This commit is contained in:
Alexey Kudravtsev
2016-07-01 15:54:04 +03:00
parent 4bca732a79
commit 4b1dd40778
2 changed files with 21 additions and 13 deletions
@@ -36,7 +36,7 @@ public class AppScheduledExecutorServiceTest extends TestCase {
}
public void testDelayedWorks() throws InterruptedException {
final AppScheduledExecutorService service = new AppScheduledExecutorService();
final AppScheduledExecutorService service = new AppScheduledExecutorService(getName());
final List<LogInfo> log = Collections.synchronizedList(new ArrayList<>());
assertFalse(service.isShutdown());
assertFalse(service.isTerminated());
@@ -96,7 +96,7 @@ public class AppScheduledExecutorServiceTest extends TestCase {
}
public void testMustNotBeAbleToShutdown() {
final AppScheduledExecutorService service = new AppScheduledExecutorService();
final AppScheduledExecutorService service = new AppScheduledExecutorService(getName());
try {
service.shutdown();
fail();
@@ -143,7 +143,7 @@ public class AppScheduledExecutorServiceTest extends TestCase {
}
public void testDelayedTasksReusePooledThreadIfExecuteAtDifferentTimes() throws InterruptedException, ExecutionException {
final AppScheduledExecutorService service = new AppScheduledExecutorService();
final AppScheduledExecutorService service = new AppScheduledExecutorService(getName());
final List<LogInfo> log = Collections.synchronizedList(new ArrayList<>());
// pre-start one thread
Future<?> future = service.submit(EmptyRunnable.getInstance());
@@ -181,13 +181,13 @@ public class AppScheduledExecutorServiceTest extends TestCase {
}
public void testDelayedTasksThatFiredAtTheSameTimeAreExecutedConcurrently() throws InterruptedException, ExecutionException {
final AppScheduledExecutorService service = new AppScheduledExecutorService();
final AppScheduledExecutorService service = new AppScheduledExecutorService(getName());
final List<LogInfo> log = Collections.synchronizedList(new ArrayList<>());
int delay = 500;
int N = 20;
List<? extends Future<?>> futures =
ContainerUtil.map(Collections.nCopies(N, ""), s -> service.schedule(()-> {
ContainerUtil.map(Collections.nCopies(N, ""), __ -> service.schedule(()-> {
log.add(new LogInfo(0));
TimeoutUtil.sleep(10 * 1000);
}
@@ -206,8 +206,8 @@ public class AppScheduledExecutorServiceTest extends TestCase {
assertTrue(service.awaitTermination(10, TimeUnit.SECONDS));
}
public void testAwaitTermination() throws InterruptedException, ExecutionException {
final AppScheduledExecutorService service = new AppScheduledExecutorService();
public void testAwaitTerminationMakesSureTasksTransferredToBackendExecutorAreFinished() throws InterruptedException, ExecutionException {
final AppScheduledExecutorService service = new AppScheduledExecutorService(getName());
final List<LogInfo> log = Collections.synchronizedList(new ArrayList<>());
int N = 20;
@@ -219,9 +219,11 @@ public class AppScheduledExecutorServiceTest extends TestCase {
}, delay, TimeUnit.MILLISECONDS
));
TimeoutUtil.sleep(delay);
while (!service.delayQueue.isEmpty()) {
long start = System.currentTimeMillis();
while (!service.delayQueue.isEmpty() && System.currentTimeMillis() < start + 20000) {
// wait till all tasks transferred to backend
}
assertTrue(service.delayQueue.toString(), service.delayQueue.isEmpty());
service.shutdownAppScheduledExecutorService();
assertTrue(service.awaitTermination(20, TimeUnit.SECONDS));
@@ -34,11 +34,12 @@ import java.util.concurrent.atomic.AtomicInteger;
public class AppScheduledExecutorService extends SchedulingWrapper {
private static final Logger LOG = Logger.getInstance("#org.jetbrains.ide.PooledThreadExecutor");
static final String POOLED_THREAD_PREFIX = "ApplicationImpl pooled thread ";
@NotNull private final String myName;
private Consumer<Thread> newThreadListener;
private final AtomicInteger counter = new AtomicInteger();
private static class Holder {
private static final AppScheduledExecutorService INSTANCE = new AppScheduledExecutorService();
private static final AppScheduledExecutorService INSTANCE = new AppScheduledExecutorService("Global instance");
}
@NotNull
@@ -46,8 +47,9 @@ public class AppScheduledExecutorService extends SchedulingWrapper {
return Holder.INSTANCE;
}
AppScheduledExecutorService() {
AppScheduledExecutorService(@NotNull final String name) {
super(new BackendThreadPoolExecutor(), new AppDelayQueue());
myName = name;
((BackendThreadPoolExecutor)backendExecutorService).doSetThreadFactory(new ThreadFactory() {
@NotNull
@Override
@@ -105,7 +107,7 @@ public class AppScheduledExecutorService extends SchedulingWrapper {
@NotNull
@TestOnly
public String statistics() {
return "app threads created counter = " + counter;
return myName + " threads created counter = " + counter;
}
public int getBackendPoolExecutorSize() {
@@ -120,18 +122,22 @@ public class AppScheduledExecutorService extends SchedulingWrapper {
private static class BackendThreadPoolExecutor extends ThreadPoolExecutor {
BackendThreadPoolExecutor() {
super(1, Integer.MAX_VALUE, 60, TimeUnit.SECONDS, new SynchronousQueue<Runnable>());
super(1, Integer.MAX_VALUE, 1, TimeUnit.MINUTES, new SynchronousQueue<Runnable>());
}
@Override
protected void beforeExecute(Thread t, Runnable r) {
if (LOG.isTraceEnabled()) {
LOG.trace("Running " + BoundedTaskExecutor.info(r) + " in " + t+" ("+System.identityHashCode(t)+")");
LOG.trace("beforeExecute " + BoundedTaskExecutor.info(r) + " in " + t);
}
}
@Override
protected void afterExecute(Runnable r, Throwable t) {
if (LOG.isTraceEnabled()) {
LOG.trace("afterExecute " + BoundedTaskExecutor.info(r) + " in " + Thread.currentThread());
}
if (t != null) {
LOG.error("Worker exited due to exception", t);
}