[threading] IJPL-206769: Allow using NonBlockingFlushQueue in LaterInvocator

GitOrigin-RevId: 4907d9bf2000879fe2269b92b793f1c012488232
This commit is contained in:
Konstantin Nisht
2025-10-06 09:27:44 +00:00
committed by intellij-monorepo-bot
parent 640601b170
commit 4c13774283
6 changed files with 105 additions and 54 deletions
@@ -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
@@ -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 {
@@ -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<ModalityStateListener> 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<Runnable> 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
@@ -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 {
@@ -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
@@ -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 {