mirror of
https://gitflic.ru/project/openide/openide.git
synced 2026-09-27 10:03:11 +07:00
run highlighting and find-in-project in "impatient" mode (which cancels whenever pending write action is detected) to fix thread starvation issues like IDEA-162320 IDEA stuck after calling Find in path during indexing
This commit is contained in:
@@ -31,15 +31,13 @@ public class ApplicationUtil {
|
||||
// throws exception if can't grab read action right now
|
||||
public static <T> T tryRunReadAction(@NotNull final Computable<T> computable) throws CannotRunReadActionException {
|
||||
final Ref<T> result = new Ref<T>();
|
||||
if (((ApplicationEx)ApplicationManager.getApplication()).tryRunReadAction(new Runnable() {
|
||||
tryRunReadAction(new Runnable() {
|
||||
@Override
|
||||
public void run() {
|
||||
result.set(computable.compute());
|
||||
}
|
||||
})) {
|
||||
return result.get();
|
||||
}
|
||||
throw new CannotRunReadActionException();
|
||||
});
|
||||
return result.get();
|
||||
}
|
||||
|
||||
public static void tryRunReadAction(@NotNull final Runnable computable) throws CannotRunReadActionException {
|
||||
@@ -48,27 +46,6 @@ public class ApplicationUtil {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Allows to interrupt a process which does not performs checkCancelled() calls by itself.
|
||||
* Note that the process may continue to run in background indefinitely - so <b>avoid using this method unless absolutely needed</b>.
|
||||
*/
|
||||
public static <T> T runWithCheckCanceled(@NotNull final Computable<T> computable, @NotNull ProgressIndicator indicator) {
|
||||
try {
|
||||
return runWithCheckCanceled(new Callable<T>() {
|
||||
@Override
|
||||
public T call() throws Exception {
|
||||
return computable.compute();
|
||||
}
|
||||
}, indicator);
|
||||
}
|
||||
catch (RuntimeException e) {
|
||||
throw e;
|
||||
}
|
||||
catch (Exception e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Allows to interrupt a process which does not performs checkCancelled() calls by itself.
|
||||
* Note that the process may continue to run in background indefinitely - so <b>avoid using this method unless absolutely needed</b>.
|
||||
@@ -112,11 +89,6 @@ public class ApplicationUtil {
|
||||
}
|
||||
}
|
||||
|
||||
public static class CannotRunReadActionException extends RuntimeException {
|
||||
@SuppressWarnings({"NullableProblems", "NonSynchronizedMethodOverridesSynchronizedMethod"})
|
||||
@Override
|
||||
public Throwable fillInStackTrace() {
|
||||
return this;
|
||||
}
|
||||
public static class CannotRunReadActionException extends ProcessCanceledException {
|
||||
}
|
||||
}
|
||||
@@ -315,49 +315,66 @@ public class PsiSearchHelperImpl implements PsiSearchHelper {
|
||||
final AtomicInteger counter = new AtomicInteger(alreadyProcessedFiles);
|
||||
final AtomicBoolean canceled = new AtomicBoolean(false);
|
||||
|
||||
boolean completed = true;
|
||||
while (true) {
|
||||
List<VirtualFile> failedList = new SmartList<>();
|
||||
final List<VirtualFile> failedFiles = Collections.synchronizedList(failedList);
|
||||
final Processor<VirtualFile> processor = vfile -> {
|
||||
try {
|
||||
TooManyUsagesStatus.getFrom(progress).pauseProcessingIfTooManyUsages();
|
||||
processVirtualFile(vfile, progress, localProcessor, canceled, counter, totalSize);
|
||||
}
|
||||
catch (ApplicationUtil.CannotRunReadActionException action) {
|
||||
failedFiles.add(vfile);
|
||||
}
|
||||
return !canceled.get();
|
||||
};
|
||||
if (ApplicationManager.getApplication().isWriteAccessAllowed() || ((ApplicationEx)ApplicationManager.getApplication()).isWriteActionPending()) {
|
||||
// no point in processing in separate threads - they are doomed to fail to obtain read action anyway
|
||||
completed &= ContainerUtil.process(files, processor);
|
||||
return processFilesConcurrentlyDespiteWriteActions(myManager.getProject(), files, progress, vfile -> {
|
||||
TooManyUsagesStatus.getFrom(progress).pauseProcessingIfTooManyUsages();
|
||||
processVirtualFile(vfile, progress, localProcessor, canceled);
|
||||
if (progress.isRunning()) {
|
||||
double fraction = (double)counter.incrementAndGet() / totalSize;
|
||||
progress.setFraction(fraction);
|
||||
}
|
||||
else {
|
||||
completed &= JobLauncher.getInstance().invokeConcurrentlyUnderProgress(files, progress, false, false, processor);
|
||||
}
|
||||
|
||||
if (failedFiles.isEmpty()) {
|
||||
break;
|
||||
}
|
||||
// we failed to run read action in job launcher thread
|
||||
// run read action in our thread instead to wait for a write action to complete and resume parallel processing
|
||||
DumbService.getInstance(myManager.getProject()).runReadActionInSmartMode(EmptyRunnable.getInstance());
|
||||
files = failedList;
|
||||
}
|
||||
return completed;
|
||||
return !canceled.get();
|
||||
});
|
||||
}
|
||||
finally {
|
||||
myManager.finishBatchFilesProcessingMode();
|
||||
}
|
||||
}
|
||||
|
||||
// Tries to run {@code localProcessor} for each file in {@code files} concurrently on ForkJoinPool.
|
||||
// When encounters write action request, stops all threads, waits for write action to finish and re-starts all threads again.
|
||||
// {@localProcessor} must be as idempotent as possible.
|
||||
public static boolean processFilesConcurrentlyDespiteWriteActions(@NotNull Project project,
|
||||
@NotNull List<VirtualFile> files,
|
||||
@NotNull final ProgressIndicator progress,
|
||||
@NotNull final Processor<VirtualFile> localProcessor) {
|
||||
final AtomicBoolean canceled = new AtomicBoolean(false);
|
||||
|
||||
boolean completed = true;
|
||||
while (true) {
|
||||
List<VirtualFile> failedList = new SmartList<>();
|
||||
final List<VirtualFile> failedFiles = Collections.synchronizedList(failedList);
|
||||
final Processor<VirtualFile> processor = vfile -> {
|
||||
try {
|
||||
return localProcessor.process(vfile);
|
||||
}
|
||||
catch (ApplicationUtil.CannotRunReadActionException action) {
|
||||
failedFiles.add(vfile);
|
||||
}
|
||||
return !canceled.get();
|
||||
};
|
||||
if (ApplicationManager.getApplication().isWriteAccessAllowed() || ((ApplicationEx)ApplicationManager.getApplication()).isWriteActionPending()) {
|
||||
// no point in processing in separate threads - they are doomed to fail to obtain read action anyway
|
||||
completed &= ContainerUtil.process(files, processor);
|
||||
}
|
||||
else {
|
||||
completed &= JobLauncher.getInstance().invokeConcurrentlyUnderProgress(files, progress, false, true, processor);
|
||||
}
|
||||
|
||||
if (failedFiles.isEmpty()) {
|
||||
break;
|
||||
}
|
||||
// we failed to run read action in job launcher thread
|
||||
// run read action in our thread instead to wait for a write action to complete and resume parallel processing
|
||||
DumbService.getInstance(project).runReadActionInSmartMode(EmptyRunnable.getInstance());
|
||||
files = failedList;
|
||||
}
|
||||
return completed;
|
||||
}
|
||||
|
||||
private void processVirtualFile(@NotNull final VirtualFile vfile,
|
||||
@NotNull final ProgressIndicator progress,
|
||||
@NotNull final Processor<? super PsiFile> localProcessor,
|
||||
@NotNull final AtomicBoolean canceled,
|
||||
@NotNull AtomicInteger counter,
|
||||
int totalSize) throws ApplicationUtil.CannotRunReadActionException {
|
||||
@NotNull final AtomicBoolean canceled) throws ApplicationUtil.CannotRunReadActionException {
|
||||
final PsiFile file = ApplicationUtil.tryRunReadAction(() -> vfile.isValid() ? myManager.findFile(vfile) : null);
|
||||
if (file != null && !(file instanceof PsiBinaryFile)) {
|
||||
// load contents outside read action
|
||||
@@ -392,10 +409,6 @@ public class PsiSearchHelperImpl implements PsiSearchHelper {
|
||||
}
|
||||
});
|
||||
}
|
||||
if (progress.isRunning()) {
|
||||
double fraction = (double)counter.incrementAndGet() / totalSize;
|
||||
progress.setFraction(fraction);
|
||||
}
|
||||
}
|
||||
|
||||
private void getFilesWithText(@NotNull GlobalSearchScope scope,
|
||||
|
||||
+4
-3
@@ -33,6 +33,7 @@ import com.intellij.lang.annotation.HighlightSeverity;
|
||||
import com.intellij.lang.injection.InjectedLanguageManager;
|
||||
import com.intellij.openapi.actionSystem.IdeActions;
|
||||
import com.intellij.openapi.application.ApplicationManager;
|
||||
import com.intellij.openapi.application.impl.ApplicationImpl;
|
||||
import com.intellij.openapi.diagnostic.Logger;
|
||||
import com.intellij.openapi.editor.Document;
|
||||
import com.intellij.openapi.editor.RangeMarker;
|
||||
@@ -246,14 +247,15 @@ public class LocalInspectionsPass extends ProgressableTextEditorHighlightingPass
|
||||
pair -> {
|
||||
LocalInspectionToolWrapper toolWrapper = pair.getKey();
|
||||
Set<String> dialectIdsSpecifiedForTool = pair.getValue();
|
||||
return runToolOnElements(toolWrapper, dialectIdsSpecifiedForTool, iManager, isOnTheFly, indicator, elements, session, init, elementDialectIds);
|
||||
((ApplicationImpl)ApplicationManager.getApplication()).executeByImpatientReader(()->runToolOnElements(toolWrapper, dialectIdsSpecifiedForTool, iManager, isOnTheFly, indicator, elements, session, init, elementDialectIds));
|
||||
return true;
|
||||
};
|
||||
boolean result = JobLauncher.getInstance().invokeConcurrentlyUnderProgress(entries, indicator, myFailFastOnAcquireReadAction, processor);
|
||||
if (!result) throw new ProcessCanceledException();
|
||||
return init;
|
||||
}
|
||||
|
||||
private boolean runToolOnElements(@NotNull final LocalInspectionToolWrapper toolWrapper,
|
||||
private void runToolOnElements(@NotNull final LocalInspectionToolWrapper toolWrapper,
|
||||
Set<String> dialectIdsSpecifiedForTool,
|
||||
@NotNull final InspectionManager iManager,
|
||||
final boolean isOnTheFly,
|
||||
@@ -289,7 +291,6 @@ public class LocalInspectionsPass extends ProgressableTextEditorHighlightingPass
|
||||
appendDescriptors(getFile(), holder.getResults(), toolWrapper);
|
||||
}
|
||||
applyIncrementally[0] = false; // do not apply incrementally outside visible range
|
||||
return true;
|
||||
}
|
||||
|
||||
private void visitRestElementsAndCleanup(@NotNull final ProgressIndicator indicator,
|
||||
|
||||
+14
-7
@@ -26,6 +26,8 @@ import com.intellij.openapi.Disposable;
|
||||
import com.intellij.openapi.application.ApplicationManager;
|
||||
import com.intellij.openapi.application.ModalityState;
|
||||
import com.intellij.openapi.application.ex.ApplicationManagerEx;
|
||||
import com.intellij.openapi.application.ex.ApplicationUtil;
|
||||
import com.intellij.openapi.application.impl.ApplicationImpl;
|
||||
import com.intellij.openapi.application.impl.ApplicationInfoImpl;
|
||||
import com.intellij.openapi.diagnostic.Logger;
|
||||
import com.intellij.openapi.editor.Document;
|
||||
@@ -401,13 +403,18 @@ class PassExecutorService implements Disposable {
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
try {
|
||||
doRun();
|
||||
}
|
||||
catch (RuntimeException | Error e) {
|
||||
saveException(e,myUpdateProgress);
|
||||
throw e;
|
||||
}
|
||||
((ApplicationImpl)ApplicationManager.getApplication()).executeByImpatientReader(() -> {
|
||||
try {
|
||||
doRun();
|
||||
}
|
||||
catch (ApplicationUtil.CannotRunReadActionException e) {
|
||||
myUpdateProgress.cancel();
|
||||
}
|
||||
catch (RuntimeException | Error e) {
|
||||
saveException(e, myUpdateProgress);
|
||||
throw e;
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
private void doRun() {
|
||||
|
||||
@@ -17,7 +17,6 @@ package com.intellij.find.impl;
|
||||
|
||||
import com.google.common.collect.HashMultiset;
|
||||
import com.google.common.collect.Multiset;
|
||||
import com.intellij.concurrency.JobSchedulerImpl;
|
||||
import com.intellij.find.FindBundle;
|
||||
import com.intellij.find.FindModel;
|
||||
import com.intellij.find.findInProject.FindInProjectManager;
|
||||
@@ -30,8 +29,10 @@ import com.intellij.openapi.fileTypes.FileType;
|
||||
import com.intellij.openapi.fileTypes.FileTypeManager;
|
||||
import com.intellij.openapi.module.Module;
|
||||
import com.intellij.openapi.module.ModuleManager;
|
||||
import com.intellij.openapi.progress.*;
|
||||
import com.intellij.openapi.progress.util.ProgressWrapper;
|
||||
import com.intellij.openapi.progress.EmptyProgressIndicator;
|
||||
import com.intellij.openapi.progress.ProcessCanceledException;
|
||||
import com.intellij.openapi.progress.ProgressIndicator;
|
||||
import com.intellij.openapi.progress.ProgressManager;
|
||||
import com.intellij.openapi.project.DumbService;
|
||||
import com.intellij.openapi.project.Project;
|
||||
import com.intellij.openapi.project.ProjectUtil;
|
||||
@@ -58,8 +59,8 @@ import com.intellij.usageView.UsageInfo;
|
||||
import com.intellij.usages.FindUsagesProcessPresentation;
|
||||
import com.intellij.usages.UsageLimitUtil;
|
||||
import com.intellij.usages.impl.UsageViewManagerImpl;
|
||||
import com.intellij.util.*;
|
||||
import com.intellij.util.concurrency.AppExecutorUtil;
|
||||
import com.intellij.util.Processor;
|
||||
import com.intellij.util.Processors;
|
||||
import com.intellij.util.containers.ContainerUtil;
|
||||
import com.intellij.util.indexing.FileBasedIndex;
|
||||
import com.intellij.util.indexing.FileBasedIndexImpl;
|
||||
@@ -67,10 +68,6 @@ import org.jetbrains.annotations.NotNull;
|
||||
import org.jetbrains.annotations.Nullable;
|
||||
|
||||
import java.util.*;
|
||||
import java.util.concurrent.CancellationException;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Future;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
@@ -79,7 +76,6 @@ import java.util.concurrent.atomic.AtomicLong;
|
||||
* @author peter
|
||||
*/
|
||||
class FindInProjectTask {
|
||||
private static final ExecutorService ourExecutor = AppExecutorUtil.createBoundedApplicationPoolExecutor("find in path", JobSchedulerImpl.CORES_COUNT);
|
||||
private static final Logger LOG = Logger.getInstance("#com.intellij.find.impl.FindInProjectTask");
|
||||
private static final int FILES_SIZE_LIMIT = 70 * 1024 * 1024; // megabytes.
|
||||
private static final int SINGLE_FILE_SIZE_LIMIT = 5 * 1024 * 1024; // megabytes.
|
||||
@@ -148,7 +144,7 @@ class FindInProjectTask {
|
||||
if (LOG.isDebugEnabled()) {
|
||||
LOG.debug("Searching for " + myFindModel.getStringToFind() + " in " + otherFiles.size() + " non-indexed files");
|
||||
}
|
||||
|
||||
myProgress.checkCanceled();
|
||||
long start = System.currentTimeMillis();
|
||||
searchInFiles(otherFiles, processPresentation, consumer);
|
||||
if (canRelyOnIndices && otherFiles.size() > 1000) {
|
||||
@@ -196,33 +192,35 @@ class FindInProjectTask {
|
||||
private void searchInFiles(@NotNull Collection<VirtualFile> virtualFiles,
|
||||
@NotNull FindUsagesProcessPresentation processPresentation,
|
||||
@NotNull final Processor<UsageInfo> consumer) {
|
||||
AtomicInteger i = new AtomicInteger();
|
||||
AtomicInteger count = new AtomicInteger();
|
||||
AtomicInteger occurrenceCount = new AtomicInteger();
|
||||
AtomicInteger processedFileCount = new AtomicInteger();
|
||||
|
||||
Consumer<VirtualFile> searchInFile = virtualFile -> {
|
||||
final int index = i.incrementAndGet();
|
||||
if (!virtualFile.isValid()) return;
|
||||
Processor<VirtualFile> processor = virtualFile -> {
|
||||
if (!virtualFile.isValid()) return true;
|
||||
|
||||
long fileLength = UsageViewManagerImpl.getFileLength(virtualFile);
|
||||
if (fileLength == -1) return; // Binary or invalid
|
||||
if (fileLength == -1) return true; // Binary or invalid
|
||||
|
||||
final boolean skipProjectFile = ProjectUtil.isProjectOrWorkspaceFile(virtualFile) && !myFindModel.isSearchInProjectFiles();
|
||||
if (skipProjectFile && !Registry.is("find.search.in.project.files")) return;
|
||||
if (skipProjectFile && !Registry.is("find.search.in.project.files")) return true;
|
||||
|
||||
if (fileLength > SINGLE_FILE_SIZE_LIMIT) {
|
||||
myLargeFiles.add(virtualFile);
|
||||
return;
|
||||
return true;
|
||||
}
|
||||
|
||||
myProgress.checkCanceled();
|
||||
myProgress.setFraction((double)index / virtualFiles.size());
|
||||
if (myProgress.isRunning()) {
|
||||
double fraction = (double)processedFileCount.incrementAndGet() / virtualFiles.size();
|
||||
myProgress.setFraction(fraction);
|
||||
}
|
||||
String text = FindBundle.message("find.searching.for.string.in.file.progress",
|
||||
myFindModel.getStringToFind(), virtualFile.getPresentableUrl());
|
||||
myProgress.setText(text);
|
||||
myProgress.setText2(FindBundle.message("find.searching.for.string.in.file.occurrences.progress", count));
|
||||
myProgress.setText2(FindBundle.message("find.searching.for.string.in.file.occurrences.progress", occurrenceCount));
|
||||
|
||||
Pair.NonNull<PsiFile, VirtualFile> pair = ReadAction.compute(() -> findFile(virtualFile));
|
||||
if (pair == null) return;
|
||||
if (pair == null) return true;
|
||||
PsiFile psiFile = pair.first;
|
||||
VirtualFile sourceVirtualFile = pair.second;
|
||||
int countInFile = FindInProjectUtil.processUsagesInFile(psiFile, sourceVirtualFile, myFindModel, info -> skipProjectFile || consumer.process(info));
|
||||
@@ -233,10 +231,10 @@ class FindInProjectTask {
|
||||
model.setSearchInProjectFiles(true);
|
||||
FindInProjectManager.getInstance(myProject).startFindInProject(model);
|
||||
});
|
||||
return;
|
||||
return true;
|
||||
}
|
||||
|
||||
count.addAndGet(countInFile);
|
||||
occurrenceCount.addAndGet(countInFile);
|
||||
if (countInFile > 0) {
|
||||
if (myTotalFilesSize.addAndGet(fileLength) > FILES_SIZE_LIMIT && myWarningShown.compareAndSet(false, true)) {
|
||||
String message = FindBundle.message("find.excessive.total.size.prompt",
|
||||
@@ -245,41 +243,9 @@ class FindInProjectTask {
|
||||
UsageLimitUtil.showAndCancelIfAborted(myProject, message, processPresentation.getUsageViewPresentation());
|
||||
}
|
||||
}
|
||||
return true;
|
||||
};
|
||||
ProgressIndicator indicator = ProgressIndicatorProvider.getGlobalProgressIndicator();
|
||||
forkJoin(virtualFiles, searchInFile, indicator != null ? indicator : new EmptyProgressIndicator());
|
||||
}
|
||||
|
||||
private static void forkJoin(@NotNull Collection<VirtualFile> virtualFiles, Consumer<VirtualFile> searchInFile, ProgressIndicator indicator) {
|
||||
List<Future<?>> futures = Collections.synchronizedList(new ArrayList<>());
|
||||
for (VirtualFile file : virtualFiles) {
|
||||
futures.add(ourExecutor.submit(() -> ProgressManager.getInstance().runProcess(() -> {
|
||||
try {
|
||||
searchInFile.consume(file);
|
||||
}
|
||||
catch (ProcessCanceledException e) {
|
||||
futures.forEach(future -> future.cancel(false));
|
||||
}
|
||||
}, ProgressWrapper.wrap(indicator))));
|
||||
}
|
||||
waitForFutures(futures);
|
||||
}
|
||||
|
||||
private static void waitForFutures(List<Future<?>> futures) {
|
||||
for (Future<?> future : futures) {
|
||||
try {
|
||||
future.get();
|
||||
}
|
||||
catch (InterruptedException e) {
|
||||
throw new ProcessCanceledException();
|
||||
}
|
||||
catch (CancellationException e) {
|
||||
throw new ProcessCanceledException();
|
||||
}
|
||||
catch (ExecutionException e) {
|
||||
ExceptionUtil.rethrowAllAsUnchecked(e.getCause());
|
||||
}
|
||||
}
|
||||
PsiSearchHelperImpl.processFilesConcurrentlyDespiteWriteActions(myProject, new ArrayList<>(virtualFiles), myProgress, processor);
|
||||
}
|
||||
|
||||
// must return non-binary files
|
||||
|
||||
@@ -17,6 +17,7 @@ package com.intellij.concurrency;
|
||||
|
||||
import com.intellij.openapi.application.ApplicationManager;
|
||||
import com.intellij.openapi.application.ex.ApplicationManagerEx;
|
||||
import com.intellij.openapi.application.impl.ApplicationImpl;
|
||||
import com.intellij.openapi.progress.ProcessCanceledException;
|
||||
import com.intellij.openapi.progress.ProgressIndicator;
|
||||
import com.intellij.openapi.progress.ProgressManager;
|
||||
@@ -44,6 +45,7 @@ import java.util.concurrent.CountedCompleter;
|
||||
*/
|
||||
class ApplierCompleter<T> extends CountedCompleter<Void> {
|
||||
private final boolean runInReadAction;
|
||||
private final boolean failFastOnAcquireReadAction;
|
||||
private final ProgressIndicator progressIndicator;
|
||||
@NotNull
|
||||
private final List<T> array;
|
||||
@@ -68,6 +70,7 @@ class ApplierCompleter<T> extends CountedCompleter<Void> {
|
||||
|
||||
ApplierCompleter(ApplierCompleter<T> parent,
|
||||
boolean runInReadAction,
|
||||
boolean failFastOnAcquireReadAction,
|
||||
@NotNull ProgressIndicator progressIndicator,
|
||||
@NotNull List<T> array,
|
||||
@NotNull Processor<? super T> processor,
|
||||
@@ -77,6 +80,7 @@ class ApplierCompleter<T> extends CountedCompleter<Void> {
|
||||
ApplierCompleter<T> next) {
|
||||
super(parent);
|
||||
this.runInReadAction = runInReadAction;
|
||||
this.failFastOnAcquireReadAction = failFastOnAcquireReadAction;
|
||||
this.progressIndicator = progressIndicator;
|
||||
this.array = array;
|
||||
this.processor = processor;
|
||||
@@ -88,7 +92,12 @@ class ApplierCompleter<T> extends CountedCompleter<Void> {
|
||||
|
||||
@Override
|
||||
public void compute() {
|
||||
wrapInReadActionAndIndicator(this::execAndForkSubTasks);
|
||||
if (failFastOnAcquireReadAction) {
|
||||
((ApplicationImpl)ApplicationManager.getApplication()).executeByImpatientReader(()-> wrapInReadActionAndIndicator(this::execAndForkSubTasks));
|
||||
}
|
||||
else {
|
||||
wrapInReadActionAndIndicator(this::execAndForkSubTasks);
|
||||
}
|
||||
}
|
||||
|
||||
private void wrapInReadActionAndIndicator(@NotNull final Runnable process) {
|
||||
@@ -100,6 +109,7 @@ class ApplierCompleter<T> extends CountedCompleter<Void> {
|
||||
} : process;
|
||||
ProgressIndicator existing = ProgressManager.getInstance().getProgressIndicator();
|
||||
if (existing == progressIndicator) {
|
||||
// we are already wrapped in an indicator - most probably because we came here from helper which steals children tasks
|
||||
toRun.run();
|
||||
}
|
||||
else {
|
||||
@@ -127,7 +137,7 @@ class ApplierCompleter<T> extends CountedCompleter<Void> {
|
||||
long elapsed = finish - start;
|
||||
if (elapsed > 5 && hi - i >= 2 && getSurplusQueuedTaskCount() <= JobSchedulerImpl.CORES_COUNT) {
|
||||
int mid = i + hi >>> 1;
|
||||
right = new ApplierCompleter<>(this, runInReadAction, progressIndicator, array, processor, mid, hi, failedSubTasks, right);
|
||||
right = new ApplierCompleter<>(this, runInReadAction, failFastOnAcquireReadAction, progressIndicator, array, processor, mid, hi, failedSubTasks, right);
|
||||
//children.add(right);
|
||||
addToPendingCount(1);
|
||||
right.fork();
|
||||
@@ -221,14 +231,15 @@ class ApplierCompleter<T> extends CountedCompleter<Void> {
|
||||
final boolean[] result = {true};
|
||||
// these tasks could not be executed in the other thread; do them here
|
||||
for (final ApplierCompleter<T> task : failedSubTasks) {
|
||||
ApplicationManager.getApplication().runReadAction(() -> task.wrapInReadActionAndIndicator(() -> {
|
||||
for (int i = task.lo; i < task.hi; ++i) {
|
||||
if (!task.processor.process(task.array.get(i))) {
|
||||
result[0] = false;
|
||||
break;
|
||||
ApplicationManager.getApplication().runReadAction(() ->
|
||||
task.wrapInReadActionAndIndicator(() -> {
|
||||
for (int i = task.lo; i < task.hi; ++i) {
|
||||
if (!task.processor.process(task.array.get(i))) {
|
||||
result[0] = false;
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}));
|
||||
}));
|
||||
}
|
||||
return result[0];
|
||||
}
|
||||
|
||||
@@ -15,7 +15,9 @@
|
||||
*/
|
||||
package com.intellij.concurrency;
|
||||
|
||||
import com.intellij.openapi.application.ApplicationManager;
|
||||
import com.intellij.openapi.application.ex.ApplicationManagerEx;
|
||||
import com.intellij.openapi.application.ex.ApplicationUtil;
|
||||
import com.intellij.openapi.progress.ProcessCanceledException;
|
||||
import com.intellij.openapi.progress.ProgressIndicator;
|
||||
import com.intellij.openapi.progress.ProgressManager;
|
||||
@@ -56,14 +58,19 @@ public class JobLauncherImpl extends JobLauncher {
|
||||
HeavyProcessLatch.INSTANCE.stopThreadPrioritizing();
|
||||
|
||||
List<ApplierCompleter<T>> failedSubTasks = Collections.synchronizedList(new ArrayList<>());
|
||||
ApplierCompleter<T> applier = new ApplierCompleter<>(null, runInReadAction, wrapper, things, thingProcessor, 0, things.size(), failedSubTasks, null);
|
||||
ApplierCompleter<T> applier = new ApplierCompleter<>(null, runInReadAction, failFastOnAcquireReadAction, wrapper, things, thingProcessor, 0, things.size(), failedSubTasks, null);
|
||||
try {
|
||||
ForkJoinPool.commonPool().invoke(applier);
|
||||
if (applier.throwable != null) throw applier.throwable;
|
||||
}
|
||||
catch (ApplierCompleter.ComputationAbortedException e) {
|
||||
// one of the processors returned false
|
||||
return false;
|
||||
}
|
||||
catch (ApplicationUtil.CannotRunReadActionException e) {
|
||||
// failFastOnAcquireReadAction==true and one of the processors called runReadAction() during the pending write action
|
||||
throw e;
|
||||
}
|
||||
catch (ProcessCanceledException e) {
|
||||
// task1.processor returns false and the task cancels the indicator
|
||||
// then task2 calls checkCancel() and get here
|
||||
@@ -91,7 +98,10 @@ public class JobLauncherImpl extends JobLauncher {
|
||||
//}
|
||||
if (things.isEmpty()) return true;
|
||||
|
||||
if (things.size() <= 1 || JobSchedulerImpl.CORES_COUNT <= CORES_FORK_THRESHOLD) {
|
||||
if (things.size() <= 1 ||
|
||||
JobSchedulerImpl.CORES_COUNT <= CORES_FORK_THRESHOLD ||
|
||||
runInReadAction && ApplicationManager.getApplication().isWriteAccessAllowed()
|
||||
) {
|
||||
final AtomicBoolean result = new AtomicBoolean(true);
|
||||
Runnable runnable = () -> ProgressManager.getInstance().executeProcessUnderProgress(() -> {
|
||||
//noinspection ForLoopReplaceableByForEach
|
||||
|
||||
@@ -33,6 +33,7 @@ import com.intellij.openapi.Disposable;
|
||||
import com.intellij.openapi.actionSystem.ex.ActionUtil;
|
||||
import com.intellij.openapi.application.*;
|
||||
import com.intellij.openapi.application.ex.ApplicationEx;
|
||||
import com.intellij.openapi.application.ex.ApplicationUtil;
|
||||
import com.intellij.openapi.command.CommandProcessor;
|
||||
import com.intellij.openapi.components.ServiceKt;
|
||||
import com.intellij.openapi.components.impl.PlatformComponentManagerImpl;
|
||||
@@ -218,6 +219,21 @@ public class ApplicationImpl extends PlatformComponentManagerImpl implements App
|
||||
myLock = new ReadMostlyRWLock(edt);
|
||||
}
|
||||
|
||||
/**
|
||||
* Executes a {@code runnable} in an "impatient" mode.
|
||||
* In this mode any attempt to call {@link #runReadAction(Runnable)}
|
||||
* would fail (i.e. throw {@link ApplicationUtil.CannotRunReadActionException})
|
||||
* if there is a pending write action.
|
||||
*/
|
||||
public void executeByImpatientReader(@NotNull Runnable runnable) throws ApplicationUtil.CannotRunReadActionException {
|
||||
if (isDispatchThread()) {
|
||||
runnable.run();
|
||||
}
|
||||
else {
|
||||
myLock.executeByImpatientReader(runnable);
|
||||
}
|
||||
}
|
||||
|
||||
private boolean disposeSelf(final boolean checkCanCloseProject) {
|
||||
final ProjectManagerImpl manager = (ProjectManagerImpl)ProjectManagerEx.getInstanceEx();
|
||||
if (manager != null) {
|
||||
|
||||
+24
-1
@@ -17,6 +17,7 @@ package com.intellij.openapi.application.impl;
|
||||
|
||||
import com.intellij.diagnostic.ThreadDumper;
|
||||
import com.intellij.openapi.application.AccessToken;
|
||||
import com.intellij.openapi.application.ex.ApplicationUtil;
|
||||
import com.intellij.openapi.diagnostic.Attachment;
|
||||
import com.intellij.openapi.diagnostic.Logger;
|
||||
import com.intellij.openapi.progress.ProgressManager;
|
||||
@@ -65,7 +66,7 @@ class ReadMostlyRWLock {
|
||||
@NotNull private final Thread thread; // its thread
|
||||
private volatile boolean readRequested; // this reader is requesting or obtained read access. Written by reader thread only, read by writer.
|
||||
private volatile boolean blocked; // this reader is blocked waiting for the writer thread to release write lock. Written by reader thread only, read by writer.
|
||||
|
||||
private boolean impatientReads; // true if should throw PCE on contented read lock
|
||||
Reader(@NotNull Thread readerThread) {
|
||||
thread = readerThread;
|
||||
}
|
||||
@@ -106,6 +107,9 @@ class ReadMostlyRWLock {
|
||||
if (iteration > SPIN_TO_WAIT_FOR_LOCK) {
|
||||
status.blocked = true;
|
||||
try {
|
||||
if (status.impatientReads) {
|
||||
throw new ApplicationUtil.CannotRunReadActionException();
|
||||
}
|
||||
LockSupport.parkNanos(this, 1000000); // unparked by writeUnlock
|
||||
}
|
||||
finally {
|
||||
@@ -117,6 +121,25 @@ class ReadMostlyRWLock {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Executes a {@code runnable} in an "impatient" mode.
|
||||
* In this mode any attempt to grab read lock
|
||||
* will fail (i.e. throw {@link ApplicationUtil.CannotRunReadActionException})
|
||||
* if there is a pending write lock request.
|
||||
*/
|
||||
void executeByImpatientReader(@NotNull Runnable runnable) throws ApplicationUtil.CannotRunReadActionException {
|
||||
checkReadThreadAccess();
|
||||
Reader status = R.get();
|
||||
boolean old = status.impatientReads;
|
||||
try {
|
||||
status.impatientReads = true;
|
||||
runnable.run();
|
||||
}
|
||||
finally {
|
||||
status.impatientReads = old;
|
||||
}
|
||||
}
|
||||
|
||||
void readUnlock() {
|
||||
checkReadThreadAccess();
|
||||
Reader status = R.get();
|
||||
|
||||
+44
-3
@@ -19,6 +19,7 @@ import com.intellij.concurrency.JobSchedulerImpl;
|
||||
import com.intellij.openapi.Disposable;
|
||||
import com.intellij.openapi.application.*;
|
||||
import com.intellij.openapi.application.ex.ApplicationEx;
|
||||
import com.intellij.openapi.application.ex.ApplicationUtil;
|
||||
import com.intellij.openapi.progress.ProgressIndicator;
|
||||
import com.intellij.openapi.progress.ProgressManager;
|
||||
import com.intellij.openapi.progress.Task;
|
||||
@@ -44,6 +45,7 @@ import java.util.ArrayList;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.Callable;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.Future;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
@@ -123,7 +125,7 @@ public class ApplicationImplTest extends LightPlatformTestCase {
|
||||
runReadWrites(readIterations, writeIterations, 10000);
|
||||
}
|
||||
|
||||
private static void runReadWrites(final int readIterations, final int writeIterations, int expectedMs) throws InterruptedException {
|
||||
private static void runReadWrites(final int readIterations, final int writeIterations, int expectedMs) {
|
||||
final ApplicationImpl application = (ApplicationImpl)ApplicationManager.getApplication();
|
||||
Disposable disposable = Disposer.newDisposable();
|
||||
application.disableEventsUntil(disposable);
|
||||
@@ -331,6 +333,7 @@ public class ApplicationImplTest extends LightPlatformTestCase {
|
||||
private volatile boolean tryingToStartWriteAction;
|
||||
private volatile boolean readStarted;
|
||||
private volatile List<Thread> readThreads;
|
||||
@SuppressWarnings("StringConcatenationInsideStringBufferAppend") // to prevent tearing
|
||||
public void testReadWontStartWhenWriteIsPending() throws Throwable {
|
||||
int N = 5;
|
||||
final AtomicBoolean[] anotherThreadStarted = new AtomicBoolean[N];
|
||||
@@ -554,14 +557,13 @@ public class ApplicationImplTest extends LightPlatformTestCase {
|
||||
}
|
||||
|
||||
public void testSuspendWriteActionDelaysForeignReadActions() throws Exception {
|
||||
List<String> log = Collections.synchronizedList(new ArrayList<>());
|
||||
|
||||
Semaphore mayStartForeignRead = new Semaphore();
|
||||
mayStartForeignRead.down();
|
||||
|
||||
List<Future> futures = new ArrayList<>();
|
||||
|
||||
ApplicationImpl app = (ApplicationImpl)ApplicationManager.getApplication();
|
||||
List<String> log = Collections.synchronizedList(new ArrayList<>());
|
||||
futures.add(app.executeOnPooledThread(() -> {
|
||||
assertTrue(mayStartForeignRead.waitFor(1000));
|
||||
ReadAction.run(() -> log.add("foreign read"));
|
||||
@@ -673,4 +675,43 @@ public class ApplicationImplTest extends LightPlatformTestCase {
|
||||
private static boolean isEscapingThreadAssertion(AssertionError e) {
|
||||
return e.getMessage().contains("should have been terminated");
|
||||
}
|
||||
|
||||
public void testReadActionInImpatientModeShouldThrowWhenThereIsAPendingWrite() throws ExecutionException, InterruptedException {
|
||||
AtomicBoolean stopRead = new AtomicBoolean();
|
||||
AtomicBoolean readAcquired = new AtomicBoolean();
|
||||
ApplicationImpl app = (ApplicationImpl)ApplicationManager.getApplication();
|
||||
Future<?> readAction1 = app.executeOnPooledThread(() ->
|
||||
app.runReadAction(() -> {
|
||||
readAcquired.set(true);
|
||||
try {
|
||||
while (!stopRead.get()) ;
|
||||
}
|
||||
finally {
|
||||
readAcquired.set(false);
|
||||
}
|
||||
})
|
||||
);
|
||||
while (!readAcquired.get());
|
||||
Future<?> readAction2 = app.executeOnPooledThread(() -> {
|
||||
// wait for write action attempt to start
|
||||
while (!app.isWriteActionPending());
|
||||
app.executeByImpatientReader(() -> {
|
||||
try {
|
||||
app.runReadAction(EmptyRunnable.getInstance());
|
||||
fail("Must have been failed");
|
||||
}
|
||||
catch (ApplicationUtil.CannotRunReadActionException ignored) {
|
||||
|
||||
}
|
||||
finally {
|
||||
stopRead.set(true);
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
app.runWriteAction(EmptyRunnable.getInstance());
|
||||
|
||||
readAction2.get();
|
||||
readAction1.get();
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user