nonBlockingReadAction: make Submission.myExpirationDisposables thread-safe (part of IJPL-148971)

GitOrigin-RevId: 2e484a50f823f772566dc03c623892f45690bd8f
This commit is contained in:
Alexey Kudravtsev
2025-11-12 18:42:01 +00:00
committed by intellij-monorepo-bot
parent 3839bc027a
commit 648435b6a5
2 changed files with 66 additions and 26 deletions
@@ -23,7 +23,6 @@ import com.intellij.openapi.progress.util.ProgressIndicatorUtils;
import com.intellij.openapi.project.Project;
import com.intellij.openapi.project.ex.ProjectEx;
import com.intellij.openapi.project.impl.ProjectImpl;
import com.intellij.openapi.util.CheckedDisposable;
import com.intellij.openapi.util.Disposer;
import com.intellij.openapi.util.Ref;
import com.intellij.openapi.vfs.VirtualFile;
@@ -54,6 +53,7 @@ import java.util.*;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReferenceArray;
import java.util.function.BooleanSupplier;
import java.util.function.Consumer;
@@ -75,7 +75,7 @@ public final class NonBlockingReadActionImpl<T> implements NonBlockingReadAction
private final Consumer<? super T> myUiThreadAction;
private final ContextConstraint @NotNull [] myConstraints;
private final BooleanSupplier @NotNull [] myCancellationConditions;
private final Set<? extends Disposable> myDisposables;
private final @Unmodifiable Disposable[] myDisposables;
private final @Nullable @Unmodifiable ListWithFixedHashCode myCoalesceEquality;
private final @Nullable ProgressIndicator myProgressIndicator;
/** Original computation passed in */
@@ -132,8 +132,9 @@ public final class NonBlockingReadActionImpl<T> implements NonBlockingReadAction
private static final ContextConstraint[] EMPTY_CONSTRAINTS = new ContextConstraint[0];
private static final BooleanSupplier[] EMPTY_CONDITIONS = new BooleanSupplier[0];
private static final Disposable[] EMPTY_DISPOSABLE_ARRAY = new Disposable[0];
NonBlockingReadActionImpl(@NotNull Callable<? extends T> computation) {
this(computation, null, null, EMPTY_CONSTRAINTS, EMPTY_CONDITIONS, Collections.emptySet(), null, null);
this(computation, null, null, EMPTY_CONSTRAINTS, EMPTY_CONDITIONS, EMPTY_DISPOSABLE_ARRAY, null, null);
}
private NonBlockingReadActionImpl(@NotNull Callable<? extends T> computation,
@@ -141,7 +142,7 @@ public final class NonBlockingReadActionImpl<T> implements NonBlockingReadAction
@Nullable Consumer<? super T> uiThreadAction,
ContextConstraint @NotNull [] constraints,
BooleanSupplier @NotNull [] cancellationConditions,
@NotNull Set<? extends Disposable> disposables,
@NotNull Disposable @NotNull [] disposables,
@Unmodifiable @Nullable ListWithFixedHashCode coalesceEquality,
@Nullable ProgressIndicator progressIndicator) {
myOriginalComputation = computation;
@@ -196,10 +197,9 @@ public final class NonBlockingReadActionImpl<T> implements NonBlockingReadAction
@Override
public @NotNull NonBlockingReadAction<T> expireWith(@NotNull Disposable parentDisposable) {
Set<Disposable> disposables = new HashSet<>(myDisposables);
disposables.add(parentDisposable);
Disposable[] newDisposables = ArrayUtil.indexOf(myDisposables, parentDisposable) == -1 ? ArrayUtil.append(myDisposables, parentDisposable) : myDisposables;
return new NonBlockingReadActionImpl<>(myOriginalComputation, myModalityState, myUiThreadAction, myConstraints, myCancellationConditions,
disposables,
newDisposables,
myCoalesceEquality, myProgressIndicator);
}
@@ -307,7 +307,8 @@ public final class NonBlockingReadActionImpl<T> implements NonBlockingReadAction
}
}
private final List<Disposable> myExpirationDisposables = new ArrayList<>();
private final @NotNull AtomicReferenceArray<@Nullable Disposable> myExpirationDisposables;
private static final @NotNull AtomicReferenceArray<@Nullable Disposable> EMPTY_ARRAY = new AtomicReferenceArray<>(0);
Submission(@NotNull NonBlockingReadActionImpl<T> builder,
@NotNull Executor backgroundThreadExecutor,
@@ -329,27 +330,28 @@ public final class NonBlockingReadActionImpl<T> implements NonBlockingReadAction
if (shouldTrackInTests()) {
ourTasksForTestMode.add(this);
}
if (!builder.myDisposables.isEmpty()) {
expireWithDisposables(this.builder.myDisposables);
Disposable[] disposables = this.builder.myDisposables;
if (disposables.length == 0) {
myExpirationDisposables = EMPTY_ARRAY;
}
else {
myExpirationDisposables = new AtomicReferenceArray<>(disposables.length);
expireWithDisposables(disposables);
}
}
private void expireWithDisposables(@NotNull Set<? extends Disposable> disposables) {
for (Disposable parent : disposables) {
private void expireWithDisposables(Disposable @NotNull [] disposables) {
for (int i = 0; i < disposables.length; i++) {
Disposable parent = disposables[i];
if (parent instanceof Project ? ((Project)parent).isDisposed() : Disposer.isDisposed(parent)) {
cancel();
break;
}
Disposable child = new CheckedDisposable() {
private volatile boolean disposed;
@Override
public boolean isDisposed() {
return disposed;
}
// need separate child instance for each parent, to be able to register in Disposer for them all
//noinspection Anonymous2MethodRef,Convert2Lambda
Disposable child = new Disposable() {
@Override
public void dispose() {
disposed = true;
// NB: We call here `super.cancel()` directly instead of `cancel()`
// The reason is that `Job` is needed to cover the scheduling of `myUiThreadAction`,
// so its lifetime is bigger than the lifetime of computation in NBRA, hence `Job` should not be cancelled in `dispose`.
@@ -363,7 +365,7 @@ public final class NonBlockingReadActionImpl<T> implements NonBlockingReadAction
cancel();
break;
}
myExpirationDisposables.add(child);
myExpirationDisposables.set(i, child);
}
}
@@ -414,6 +416,7 @@ public final class NonBlockingReadActionImpl<T> implements NonBlockingReadAction
if (cleanedHandle.compareAndSet(this, false, true)) {
cleanup();
}
disposeExpirationDisposables();
}
private void cleanup() {
@@ -427,9 +430,6 @@ public final class NonBlockingReadActionImpl<T> implements NonBlockingReadAction
if (builder.myCoalesceEquality != null) {
release();
}
for (Disposable disposable : myExpirationDisposables) {
Disposer.dispose(disposable);
}
if (hasUnboundedExecutor()) {
ourUnboundedSubmissionTracker.unregisterSubmission(myStartTrace);
}
@@ -438,6 +438,15 @@ public final class NonBlockingReadActionImpl<T> implements NonBlockingReadAction
}
}
private void disposeExpirationDisposables() {
for (int i = 0; i < myExpirationDisposables.length(); i++) {
Disposable disposable = myExpirationDisposables.getAndSet(i, null);
if (disposable != null) {
Disposer.dispose(disposable);
}
}
}
private void acquire() {
assert builder.myCoalesceEquality != null;
synchronized (ourTasksByEquality) {
@@ -12,6 +12,7 @@ import com.intellij.openapi.progress.ProcessCanceledException;
import com.intellij.openapi.progress.ProgressIndicator;
import com.intellij.openapi.progress.ProgressManager;
import com.intellij.openapi.progress.util.ProgressIndicatorBase;
import com.intellij.openapi.util.CheckedDisposable;
import com.intellij.openapi.util.Disposer;
import com.intellij.openapi.util.Pair;
import com.intellij.openapi.util.text.StringUtil;
@@ -21,6 +22,7 @@ import com.intellij.psi.impl.PsiDocumentManagerImpl;
import com.intellij.psi.util.PsiUtilCore;
import com.intellij.testFramework.*;
import com.intellij.tools.ide.metrics.benchmark.Benchmark;
import com.intellij.util.ConcurrencyUtil;
import com.intellij.util.ExceptionUtil;
import com.intellij.util.TimeoutUtil;
import com.intellij.util.concurrency.AppExecutorUtil;
@@ -539,7 +541,7 @@ public class NonBlockingReadActionTest extends LightPlatformTestCase {
}
}
public void test_submit_doesNot_fail_without_readAction_when_parent_isDisposed() {
public void test_submit_doesNot_fail_without_readAction_when_parent_isDisposed() throws Exception {
ExecutorService executor = AppExecutorUtil.createBoundedApplicationPoolExecutor(StringUtil.capitalize(getName()), 10);
for (int i = 0; i < 50; i++) {
@@ -557,7 +559,36 @@ public class NonBlockingReadActionTest extends LightPlatformTestCase {
}
parents.forEach(Disposer::dispose);
futures.forEach(f -> PlatformTestUtil.waitForFuture(f, 50_000));
ConcurrencyUtil.getAll(50, TimeUnit.SECONDS, futures);
}
}
public void testSubmitDoesNotFailWhenParentsAreDisposedConcurrently() throws Exception {
ExecutorService executor = AppExecutorUtil.createBoundedApplicationPoolExecutor(StringUtil.capitalize(getName()), 10);
for (int i = 0; i < 50; i++) {
List<CheckedDisposable> parents = new ArrayList<>(200);
List<Future<?>> futures = new ArrayList<>();
for (int j = 0; j < 100; j++) {
CheckedDisposable parent = Disposer.newCheckedDisposable();
CheckedDisposable parent2 = Disposer.newCheckedDisposable();
parents.add(parent);
parents.add(parent2);
futures.add(executor.submit(() -> ReadAction.nonBlocking(() -> {}).expireWith(parent).expireWith(parent2).submit(executor).get()));
futures.add(executor.submit(() -> {
try {
ReadAction.nonBlocking(() -> {}).expireWith(parent).expireWith(parent2).executeSynchronously();
}
catch (ProcessCanceledException ignore) {
}
}));
}
parents.forEach(Disposer::dispose);
ConcurrencyUtil.getAll(50, TimeUnit.SECONDS, futures);
for (CheckedDisposable parent : parents) {
assertTrue(parent.isDisposed());
}
}
}