From 4c1377428389db154de1f33d29ce27ff892bb6da Mon Sep 17 00:00:00 2001 From: Konstantin Nisht Date: Fri, 3 Oct 2025 16:53:57 +0200 Subject: [PATCH] [threading] IJPL-206769: Allow using `NonBlockingFlushQueue` in `LaterInvocator` GitOrigin-RevId: 4907d9bf2000879fe2269b92b793f1c012488232 --- .../application/ThreadingRuntimeFlags.kt | 6 +- .../impl/EdtCoroutineDispatcher.kt | 4 +- .../application/impl/LaterInvocator.java | 53 ++++++++++-- .../src/com/intellij/ide/IdeEventQueue.kt | 6 +- .../application/impl/ApplicationImpl.java | 84 +++++++++++-------- .../impl/EdtCoroutineDispatcherTest.kt | 6 +- 6 files changed, 105 insertions(+), 54 deletions(-) diff --git a/platform/core-api/src/com/intellij/openapi/application/ThreadingRuntimeFlags.kt b/platform/core-api/src/com/intellij/openapi/application/ThreadingRuntimeFlags.kt index a0ce2fa3e2a3..aaa5f79f366d 100644 --- a/platform/core-api/src/com/intellij/openapi/application/ThreadingRuntimeFlags.kt +++ b/platform/core-api/src/com/intellij/openapi/application/ThreadingRuntimeFlags.kt @@ -42,11 +42,11 @@ val installSuvorovProgress: Boolean = System.getProperty("ide.install.suvorov.pr val useDebouncedDrawingInSuvorovProgress: Boolean = System.getProperty("ide.suvorov.progress.debounced.drawing", "true").toBoolean() /** - * - `true` means that [kotlinx.coroutines.Dispatchers.EDT] will acquire write-intent lock in non-blocking way - * - `false` means that [kotlinx.coroutines.Dispatchers.EDT] will block on the acquisition of high-level write-intent lock. + * - `true` means that EDT runnables that require write-intent lock will acquire it in a non-blocking way + * - `false` means that the write-intent lock will be acquired in a blocking way */ @get:ApiStatus.Internal -val useNonBlockingIntentLockForEdtCoroutines: Boolean = System.getProperty("ide.non.blocking.write.intent.lock.for.edt.coroutines", "true").toBoolean() +val useNonBlockingFlushQueue: Boolean = System.getProperty("ide.use.non.blocking.flush.queue", "false").toBoolean() /** * Represents the deadline before blocking read lock acquisition starts compensating parallelism for coroutine worker threads diff --git a/platform/core-impl/src/com/intellij/openapi/application/impl/EdtCoroutineDispatcher.kt b/platform/core-impl/src/com/intellij/openapi/application/impl/EdtCoroutineDispatcher.kt index c5953ec038bb..ae67552f3fb1 100644 --- a/platform/core-impl/src/com/intellij/openapi/application/impl/EdtCoroutineDispatcher.kt +++ b/platform/core-impl/src/com/intellij/openapi/application/impl/EdtCoroutineDispatcher.kt @@ -44,7 +44,7 @@ internal sealed class EdtCoroutineDispatcher( else { DispatchedRunnable(context.job, lockingAwareBlock) } - val useWeakWriteIntent = useNonBlockingIntentLockForEdtCoroutines && type.lockBehavior == EdtDispatcherKind.LockBehavior.LOCKS_ALLOWED_MANDATORY_WRAPPING + val useWeakWriteIntent = useNonBlockingFlushQueue && type.lockBehavior == EdtDispatcherKind.LockBehavior.LOCKS_ALLOWED_MANDATORY_WRAPPING ApplicationManagerEx.getApplicationEx().dispatchCoroutineOnEDT(runnable, state, useWeakWriteIntent) } @@ -66,7 +66,7 @@ internal sealed class EdtCoroutineDispatcher( runnable } EdtDispatcherKind.LockBehavior.LOCKS_ALLOWED_MANDATORY_WRAPPING -> { - if (useNonBlockingIntentLockForEdtCoroutines) { + if (useNonBlockingFlushQueue) { runnable } else { diff --git a/platform/core-impl/src/com/intellij/openapi/application/impl/LaterInvocator.java b/platform/core-impl/src/com/intellij/openapi/application/impl/LaterInvocator.java index fba086c6a132..574c09f8abe8 100644 --- a/platform/core-impl/src/com/intellij/openapi/application/impl/LaterInvocator.java +++ b/platform/core-impl/src/com/intellij/openapi/application/impl/LaterInvocator.java @@ -20,6 +20,7 @@ import com.intellij.util.containers.CollectionFactory; import com.intellij.util.containers.ContainerUtil; import com.intellij.util.containers.Stack; import com.intellij.util.ui.EDT; +import kotlin.jvm.Volatile; import org.jetbrains.annotations.*; import javax.swing.*; @@ -46,6 +47,12 @@ public final class LaterInvocator { private static final EventDispatcher ourModalityStateMulticaster = EventDispatcher.create(ModalityStateListener.class); private static final FlushQueue ourEdtQueue = new FlushQueue(); + @Volatile + private static NonBlockingFlushQueue ourNonBlockingEdtQueue = null; + + public static void initializeNonBlockingFlushQueue(@NotNull ThreadingSupport threadingSupport) { + ourNonBlockingEdtQueue = new NonBlockingFlushQueue(threadingSupport); + } public static void addModalityStateListener(@NotNull ModalityStateListener listener, @NotNull Disposable parentDisposable) { if (!ourModalityStateMulticaster.getListeners().contains(listener)) { @@ -78,18 +85,36 @@ public final class LaterInvocator { @ApiStatus.Internal public static void invokeLater(@NotNull ModalityState modalityState, - @NotNull Condition expired, - @NotNull Runnable runnable) { + @NotNull Condition expired, + @NotNull Runnable runnable) { + invokeLater(modalityState, expired, true, runnable); + } + + @ApiStatus.Internal + public static void invokeLater(@NotNull ModalityState modalityState, + @NotNull Condition expired, + boolean needsWriteIntentLock, + @NotNull Runnable runnable) { SideEffectGuard.checkSideEffectAllowed(SideEffectGuard.EffectType.INVOKE_LATER); if (expired.value(null)) { return; } - ourEdtQueue.push(modalityState, expired, runnable); + if (useNonBlockingFlushQueue()) { + ourNonBlockingEdtQueue.push(modalityState, runnable, needsWriteIntentLock, expired); + } else { + ourEdtQueue.push(modalityState, expired, runnable); + } + } + + private static boolean useNonBlockingFlushQueue() { + // non-blocking flush queue can either be explicitly disabled, or be null, like in ServerApplication + // it is currently a TODO: use non-blocking flush queue in code server + return ThreadingRuntimeFlagsKt.getUseNonBlockingFlushQueue() && ourNonBlockingEdtQueue != null; } @RequiresBackgroundThread @ApiStatus.Internal - public static void invokeAndWait(@NotNull ModalityState modalityState, final @NotNull Runnable runnable) { + public static void invokeAndWait(@NotNull ModalityState modalityState, boolean wrapWithLocks, final @NotNull Runnable runnable) { final AtomicReference runnableRef = new AtomicReference<>(runnable); final Semaphore semaphore = new Semaphore(); semaphore.down(); @@ -118,7 +143,7 @@ public final class LaterInvocator { return "InvokeAndWait[" + (runnable == null ? "(cancelled)" : runnable.toString()) + "]"; } }; - invokeLater(modalityState, Conditions.alwaysFalse(), runnable1); + invokeLater(modalityState, Conditions.alwaysFalse(), wrapWithLocks, runnable1); try { while (!semaphore.waitFor(ConcurrencyUtil.DEFAULT_TIMEOUT_MS)) { ProgressManager.checkCanceled(); @@ -344,7 +369,11 @@ public final class LaterInvocator { } static boolean isFlushNow(@NotNull Runnable runnable) { - return ourEdtQueue.isFlushNow(runnable); + if (useNonBlockingFlushQueue()) { + return ourNonBlockingEdtQueue.isFlushNow(runnable); + } else { + return ourEdtQueue.isFlushNow(runnable); + } } public static void pollWriteThreadEventsOnce() { LOG.assertTrue(!SwingUtilities.isEventDispatchThread()); @@ -353,12 +382,20 @@ public final class LaterInvocator { @TestOnly public static @NotNull Object getLaterInvocatorEdtQueue() { - return ourEdtQueue.getQueue(); + if (useNonBlockingFlushQueue()) { + return ourNonBlockingEdtQueue; + } else { + return ourEdtQueue.getQueue(); + } } @RequiresEdt private static void reincludeSkippedItemsAndRequestFlush() { - ourEdtQueue.reincludeSkippedItems(); + if (useNonBlockingFlushQueue()) { + ourNonBlockingEdtQueue.onModalityChanged(); + } else { + ourEdtQueue.reincludeSkippedItems(); + } } @RequiresEdt diff --git a/platform/platform-impl/src/com/intellij/ide/IdeEventQueue.kt b/platform/platform-impl/src/com/intellij/ide/IdeEventQueue.kt index 08fb2b028413..38a8d12f024e 100644 --- a/platform/platform-impl/src/com/intellij/ide/IdeEventQueue.kt +++ b/platform/platform-impl/src/com/intellij/ide/IdeEventQueue.kt @@ -18,6 +18,7 @@ import com.intellij.openapi.Disposable import com.intellij.openapi.application.* import com.intellij.openapi.application.ex.ApplicationManagerEx import com.intellij.openapi.application.impl.InvocationUtil +import com.intellij.openapi.application.impl.LaterInvocator import com.intellij.openapi.components.serviceIfCreated import com.intellij.openapi.diagnostic.ControlFlowException import com.intellij.openapi.diagnostic.Logger @@ -140,6 +141,9 @@ class IdeEventQueue private constructor() : EventQueue() { assert(isDispatchThread()) { Thread.currentThread() } val systemEventQueue = Toolkit.getDefaultToolkit().systemEventQueue assert(systemEventQueue !is IdeEventQueue) { systemEventQueue } + if (useNonBlockingFlushQueue) { + LaterInvocator.initializeNonBlockingFlushQueue(threadingSupport) + } systemEventQueue.push(this) EDT.updateEdt() replaceDefaultKeyboardFocusManager() @@ -325,7 +329,7 @@ class IdeEventQueue private constructor() : EventQueue() { try { runCustomProcessors(finalEvent, preProcessors) performActivity(finalEvent, !nakedRunnable && isPureSwingEventWilEnabled && !threadingSupport.isInsideUnlockedWriteIntentLock()) { - if (progressManager == null) { + if (progressManager == null || (runnable != null && useNonBlockingFlushQueue && InvocationUtil.isFlushNow(runnable))) { _dispatchEvent(finalEvent) } else { 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 c21f136a9803..2b72df5f8c29 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 @@ -414,7 +414,7 @@ public final class ApplicationImpl extends ClientAwareComponentManager implement // Start from inner layer: transaction guard final var guarded = myTransactionGuard.wrapLaterInvocation(runnable, state); // Middle layer: lock and modality - final var locked = wrapWithRunIntendedWriteActionAndModality(guarded, ctxAware ? null : state); + final var locked = wrapWithRunIntendedWriteActionAndModality(guarded, true, ctxAware ? null : state); var finalRunnable = locked; // Outer layer, optional: context capture & reset if (propagateContext()) { @@ -422,19 +422,14 @@ public final class ApplicationImpl extends ClientAwareComponentManager implement finalRunnable = captured.getFirst(); expired = captured.getSecond(); } - LaterInvocator.invokeLater(state, expired, finalRunnable); + LaterInvocator.invokeLater(state, expired, true, finalRunnable); } @ApiStatus.Internal @Override - public void dispatchCoroutineOnEDT(Runnable runnable, ModalityState state, boolean acquireWriteIntentLockInNonBlockingWay) { + public void dispatchCoroutineOnEDT(Runnable runnable, ModalityState state, boolean needsWriteIntent) { var wrapped = myTransactionGuard.wrapCoroutineInvocation(runnable, state); - if (acquireWriteIntentLockInNonBlockingWay) { - scheduleWithWeakWriteIntentReadAction(wrapped, state); - } - else { - LaterInvocator.invokeLater(state, Conditions.alwaysFalse(), wrapped); - } + LaterInvocator.invokeLater(state, Conditions.alwaysFalse(), needsWriteIntent, wrapped); } private void scheduleWithWeakWriteIntentReadAction(@NotNull Runnable runnable, @NotNull ModalityState state) { @@ -587,41 +582,56 @@ public final class ApplicationImpl extends ClientAwareComponentManager implement // Start from inner layer: transaction guard final var guarded = myTransactionGuard.wrapLaterInvocation(runnable, state); // Middle layer: lock and modality - final var locked = wrapWithLocks ? wrapWithRunIntendedWriteActionAndModality(guarded, ctxAware ? null : state) : guarded; + boolean wrapWithLocksDeep = wrapWithLocks && !ThreadingRuntimeFlagsKt.getUseNonBlockingFlushQueue(); + final var locked = wrapWithRunIntendedWriteActionAndModality(guarded, wrapWithLocksDeep, ctxAware ? null : state); // Outer layer context capture & reset final var finalRunnable = AppImplKt.rethrowExceptions(AppScheduledExecutorService::captureContextCancellationForRunnableThatDoesNotOutliveContextScope, locked); - LaterInvocator.invokeAndWait(state, finalRunnable); + LaterInvocator.invokeAndWait(state, wrapWithLocks, finalRunnable); } - private @NotNull Runnable wrapWithRunIntendedWriteActionAndModality(@NotNull Runnable runnable, @Nullable ModalityState modalityState) { - return modalityState != null ? - new Runnable() { - @Override - public void run() { - ThreadContext.installThreadContext(ThreadContext.currentThreadContext().plus(asContextElement(modalityState)), true, () -> { - runIntendedWriteActionOnCurrentThread(runnable); - return Unit.INSTANCE; - }); - } + private @NotNull Runnable wrapWithRunIntendedWriteActionAndModality(@NotNull Runnable runnable, + boolean wrapWithLocks, + @Nullable ModalityState modalityState) { + if (modalityState == null && wrapWithLocks) { + return new Runnable() { + @Override + public void run() { + runIntendedWriteActionOnCurrentThread(runnable); + } - @Override - public String toString() { - return runnable.toString(); - } - } - : - new Runnable() { - @Override - public void run() { - runIntendedWriteActionOnCurrentThread(runnable); - } + @Override + public String toString() { + return runnable.toString(); + } + }; + } + else if (modalityState == null) { + // wrapWithLocks == false + return runnable; + } + else { + // modalityState != null + return new Runnable() { + @Override + public void run() { + ThreadContext.installThreadContext(ThreadContext.currentThreadContext().plus(asContextElement(modalityState)), true, () -> { + if (wrapWithLocks) { + runIntendedWriteActionOnCurrentThread(runnable); + } + else { + runnable.run(); + } + return Unit.INSTANCE; + }); + } - @Override - public String toString() { - return runnable.toString(); - } - }; + @Override + public String toString() { + return runnable.toString(); + } + }; + } } @Override diff --git a/platform/platform-tests/testSrc/com/intellij/openapi/application/impl/EdtCoroutineDispatcherTest.kt b/platform/platform-tests/testSrc/com/intellij/openapi/application/impl/EdtCoroutineDispatcherTest.kt index 3e91d4ab4d5e..81783a60b534 100644 --- a/platform/platform-tests/testSrc/com/intellij/openapi/application/impl/EdtCoroutineDispatcherTest.kt +++ b/platform/platform-tests/testSrc/com/intellij/openapi/application/impl/EdtCoroutineDispatcherTest.kt @@ -383,8 +383,8 @@ class EdtCoroutineDispatcherTest { } } finally { - assertThat(application.isReadAccessAllowed).isEqualTo(!useNonBlockingIntentLockForEdtCoroutines) - assertThat(application.isWriteIntentLockAcquired).isEqualTo(!useNonBlockingIntentLockForEdtCoroutines) + assertThat(application.isReadAccessAllowed).isEqualTo(!useNonBlockingFlushQueue) + assertThat(application.isWriteIntentLockAcquired).isEqualTo(!useNonBlockingFlushQueue) assertThat(application.isWriteAccessAllowed).isFalse } } @@ -530,7 +530,7 @@ class EdtCoroutineDispatcherTest { @Test fun `UI coroutine can be executed earlier then EDT coroutine`(): Unit = timeoutRunBlocking { - Assumptions.assumeTrue { useNonBlockingIntentLockForEdtCoroutines } + Assumptions.assumeTrue { useNonBlockingFlushQueue } val uiExecuted = AtomicBoolean() val edtExecuted = AtomicBoolean() backgroundWriteAction {