diff --git a/platform/core-impl/src/com/intellij/openapi/application/ex/ApplicationUtil.java b/platform/core-impl/src/com/intellij/openapi/application/ex/ApplicationUtil.java index 2b2784d5004a..2bddb9addbf8 100644 --- a/platform/core-impl/src/com/intellij/openapi/application/ex/ApplicationUtil.java +++ b/platform/core-impl/src/com/intellij/openapi/application/ex/ApplicationUtil.java @@ -31,15 +31,13 @@ public class ApplicationUtil { // throws exception if can't grab read action right now public static T tryRunReadAction(@NotNull final Computable computable) throws CannotRunReadActionException { final Ref result = new Ref(); - 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 avoid using this method unless absolutely needed. - */ - public static T runWithCheckCanceled(@NotNull final Computable computable, @NotNull ProgressIndicator indicator) { - try { - return runWithCheckCanceled(new Callable() { - @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 avoid using this method unless absolutely needed. @@ -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 { } } \ No newline at end of file diff --git a/platform/indexing-impl/src/com/intellij/psi/impl/search/PsiSearchHelperImpl.java b/platform/indexing-impl/src/com/intellij/psi/impl/search/PsiSearchHelperImpl.java index 2ea7cad9da13..c83ff3c17986 100644 --- a/platform/indexing-impl/src/com/intellij/psi/impl/search/PsiSearchHelperImpl.java +++ b/platform/indexing-impl/src/com/intellij/psi/impl/search/PsiSearchHelperImpl.java @@ -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 failedList = new SmartList<>(); - final List failedFiles = Collections.synchronizedList(failedList); - final Processor 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 files, + @NotNull final ProgressIndicator progress, + @NotNull final Processor localProcessor) { + final AtomicBoolean canceled = new AtomicBoolean(false); + + boolean completed = true; + while (true) { + List failedList = new SmartList<>(); + final List failedFiles = Collections.synchronizedList(failedList); + final Processor 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 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, diff --git a/platform/lang-impl/src/com/intellij/codeInsight/daemon/impl/LocalInspectionsPass.java b/platform/lang-impl/src/com/intellij/codeInsight/daemon/impl/LocalInspectionsPass.java index c067adc24da3..ad6877a4e6dc 100644 --- a/platform/lang-impl/src/com/intellij/codeInsight/daemon/impl/LocalInspectionsPass.java +++ b/platform/lang-impl/src/com/intellij/codeInsight/daemon/impl/LocalInspectionsPass.java @@ -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 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 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, diff --git a/platform/lang-impl/src/com/intellij/codeInsight/daemon/impl/PassExecutorService.java b/platform/lang-impl/src/com/intellij/codeInsight/daemon/impl/PassExecutorService.java index 57db3cd614a4..d6464839a125 100644 --- a/platform/lang-impl/src/com/intellij/codeInsight/daemon/impl/PassExecutorService.java +++ b/platform/lang-impl/src/com/intellij/codeInsight/daemon/impl/PassExecutorService.java @@ -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() { diff --git a/platform/lang-impl/src/com/intellij/find/impl/FindInProjectTask.java b/platform/lang-impl/src/com/intellij/find/impl/FindInProjectTask.java index 639f8debab54..b453a5570ed1 100644 --- a/platform/lang-impl/src/com/intellij/find/impl/FindInProjectTask.java +++ b/platform/lang-impl/src/com/intellij/find/impl/FindInProjectTask.java @@ -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 virtualFiles, @NotNull FindUsagesProcessPresentation processPresentation, @NotNull final Processor consumer) { - AtomicInteger i = new AtomicInteger(); - AtomicInteger count = new AtomicInteger(); + AtomicInteger occurrenceCount = new AtomicInteger(); + AtomicInteger processedFileCount = new AtomicInteger(); - Consumer searchInFile = virtualFile -> { - final int index = i.incrementAndGet(); - if (!virtualFile.isValid()) return; + Processor 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 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 virtualFiles, Consumer searchInFile, ProgressIndicator indicator) { - List> 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> 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 diff --git a/platform/platform-impl/src/com/intellij/concurrency/ApplierCompleter.java b/platform/platform-impl/src/com/intellij/concurrency/ApplierCompleter.java index e62efbd1426d..64d3da1f24e2 100644 --- a/platform/platform-impl/src/com/intellij/concurrency/ApplierCompleter.java +++ b/platform/platform-impl/src/com/intellij/concurrency/ApplierCompleter.java @@ -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 extends CountedCompleter { private final boolean runInReadAction; + private final boolean failFastOnAcquireReadAction; private final ProgressIndicator progressIndicator; @NotNull private final List array; @@ -68,6 +70,7 @@ class ApplierCompleter extends CountedCompleter { ApplierCompleter(ApplierCompleter parent, boolean runInReadAction, + boolean failFastOnAcquireReadAction, @NotNull ProgressIndicator progressIndicator, @NotNull List array, @NotNull Processor processor, @@ -77,6 +80,7 @@ class ApplierCompleter extends CountedCompleter { ApplierCompleter next) { super(parent); this.runInReadAction = runInReadAction; + this.failFastOnAcquireReadAction = failFastOnAcquireReadAction; this.progressIndicator = progressIndicator; this.array = array; this.processor = processor; @@ -88,7 +92,12 @@ class ApplierCompleter extends CountedCompleter { @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 extends CountedCompleter { } : 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 extends CountedCompleter { 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 extends CountedCompleter { final boolean[] result = {true}; // these tasks could not be executed in the other thread; do them here for (final ApplierCompleter 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]; } diff --git a/platform/platform-impl/src/com/intellij/concurrency/JobLauncherImpl.java b/platform/platform-impl/src/com/intellij/concurrency/JobLauncherImpl.java index 1f2169c49c2d..bdd7d56c25f4 100644 --- a/platform/platform-impl/src/com/intellij/concurrency/JobLauncherImpl.java +++ b/platform/platform-impl/src/com/intellij/concurrency/JobLauncherImpl.java @@ -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> failedSubTasks = Collections.synchronizedList(new ArrayList<>()); - ApplierCompleter applier = new ApplierCompleter<>(null, runInReadAction, wrapper, things, thingProcessor, 0, things.size(), failedSubTasks, null); + ApplierCompleter 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 diff --git a/platform/platform-impl/src/com/intellij/openapi/application/impl/ApplicationImpl.java b/platform/platform-impl/src/com/intellij/openapi/application/impl/ApplicationImpl.java index bf3580c55a03..82a71f46953b 100644 --- a/platform/platform-impl/src/com/intellij/openapi/application/impl/ApplicationImpl.java +++ b/platform/platform-impl/src/com/intellij/openapi/application/impl/ApplicationImpl.java @@ -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) { diff --git a/platform/platform-impl/src/com/intellij/openapi/application/impl/ReadMostlyRWLock.java b/platform/platform-impl/src/com/intellij/openapi/application/impl/ReadMostlyRWLock.java index 3d0f90c3a8ac..d80ab9524095 100644 --- a/platform/platform-impl/src/com/intellij/openapi/application/impl/ReadMostlyRWLock.java +++ b/platform/platform-impl/src/com/intellij/openapi/application/impl/ReadMostlyRWLock.java @@ -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(); diff --git a/platform/platform-tests/testSrc/com/intellij/openapi/application/impl/ApplicationImplTest.java b/platform/platform-tests/testSrc/com/intellij/openapi/application/impl/ApplicationImplTest.java index 641d2e350e91..2c59780f1923 100644 --- a/platform/platform-tests/testSrc/com/intellij/openapi/application/impl/ApplicationImplTest.java +++ b/platform/platform-tests/testSrc/com/intellij/openapi/application/impl/ApplicationImplTest.java @@ -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 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 log = Collections.synchronizedList(new ArrayList<>()); - Semaphore mayStartForeignRead = new Semaphore(); mayStartForeignRead.down(); List futures = new ArrayList<>(); ApplicationImpl app = (ApplicationImpl)ApplicationManager.getApplication(); + List 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(); + } }