invokeConcurrently() must execute multiple tasks concurrently even if one of them is wildly slow (which is often the case in inspections) to ensure progress

GitOrigin-RevId: 066dbab249819f0dfd56f7dc1b2f57936cf5e12e
This commit is contained in:
Alexey Kudravtsev
2019-10-03 11:33:07 +00:00
committed by intellij-monorepo-bot
parent d5ee1f1f53
commit 368dd7c859
2 changed files with 38 additions and 10 deletions
@@ -123,20 +123,14 @@ class ApplierCompleter<T> extends CountedCompleter<Void> {
@Nullable
private ApplierCompleter<T> execAndForkSubTasks() {
int hi = this.hi;
long start = System.currentTimeMillis();
ApplierCompleter<T> right = null;
Throwable throwable = null;
try {
for (int i = lo; i < hi; ++i) {
ProgressManager.checkCanceled();
if (!processor.process(array.get(i))) {
throw new ComputationAbortedException();
}
long finish = System.currentTimeMillis();
long elapsed = finish - start;
if (elapsed > 1 && hi - i >= 2) {
int availableParallelism = JobSchedulerImpl.getJobPoolParallelism() - getSurplusQueuedTaskCount();
if (hi - i >= 2) {
int availableParallelism = JobSchedulerImpl.getJobPoolParallelism() - Math.max(0,getSurplusQueuedTaskCount());
if (availableParallelism > 1) {
// fork off several sub-tasks at once to reduce rampup
for (int n=0; n<availableParallelism; n++) {
@@ -149,9 +143,11 @@ class ApplierCompleter<T> extends CountedCompleter<Void> {
right.fork();
hi = mid;
}
start = finish;
}
}
if (!processor.process(array.get(i))) {
throw new ComputationAbortedException();
}
}
// traverse the list looking for a task available for stealing
@@ -44,7 +44,8 @@ import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
import static com.intellij.util.TestTimeOut.*;
import static com.intellij.util.TestTimeOut.setTimeout;
public class JobUtilTest extends LightPlatformTestCase {
private static final AtomicInteger COUNT = new AtomicInteger();
@@ -538,4 +539,35 @@ public class JobUtilTest extends LightPlatformTestCase {
if (!stealHappened.get()) break; // tested that we wanted
}
}
public void testInvokeConcurrentlyMustExecuteMultipleTasksConcurrentlyEvenIfOneOfThemIsWildlySlow() {
int N = 8;
Integer[] times = new Integer[N];
for (int i=0; i<N; i++) {
Arrays.fill(times, 0);
times[i] = 1_000_000;
AtomicInteger executed = new AtomicInteger();
DaemonProgressIndicator progress = new DaemonProgressIndicator();
String enough = "enough is enough";
try {
JobLauncher.getInstance().invokeConcurrentlyUnderProgress(Arrays.asList(times.clone()), progress, time -> {
while ((time -= 100) >= 0) {
ProgressManager.checkCanceled();
TimeoutUtil.sleep(100);
}
if (executed.incrementAndGet() == times.length - 1) {
// executed all but the slowest one
throw new RuntimeException(enough);
}
return true;
});
fail();
}
catch (RuntimeException e) {
assertEquals(enough, e.getMessage());
}
}
}
}