diff --git a/platform/core-impl/src/com/intellij/openapi/application/impl/BulkArrayQueue.java b/platform/core-impl/src/com/intellij/openapi/application/impl/BulkArrayQueue.java index a9c613b1c743..54b7e73ff5a8 100644 --- a/platform/core-impl/src/com/intellij/openapi/application/impl/BulkArrayQueue.java +++ b/platform/core-impl/src/com/intellij/openapi/application/impl/BulkArrayQueue.java @@ -1,4 +1,4 @@ -// Copyright 2000-2024 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. +// Copyright 2000-2025 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. package com.intellij.openapi.application.impl; import it.unimi.dsi.fastutil.objects.ObjectArrayList; @@ -65,6 +65,13 @@ public final class BulkArrayQueue { return info; } + public T peekFirst() { + if (isEmpty()) return null; + //noinspection unchecked + return (T)myQueue[head]; + } + + private @NotNull T getAndNullize(int head) { //noinspection unchecked T t = (T)myQueue[head]; diff --git a/platform/core-impl/src/com/intellij/openapi/application/impl/NonBlockingFlushQueue.kt b/platform/core-impl/src/com/intellij/openapi/application/impl/NonBlockingFlushQueue.kt new file mode 100644 index 000000000000..5de241b8f76b --- /dev/null +++ b/platform/core-impl/src/com/intellij/openapi/application/impl/NonBlockingFlushQueue.kt @@ -0,0 +1,430 @@ +// Copyright 2000-2025 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. +package com.intellij.openapi.application.impl + +import com.intellij.concurrency.ContextAwareRunnable +import com.intellij.concurrency.resetThreadContext +import com.intellij.diagnostic.EventWatcher +import com.intellij.openapi.application.ApplicationManager +import com.intellij.openapi.application.ModalityState +import com.intellij.openapi.application.ThreadingSupport +import com.intellij.openapi.diagnostic.Logger +import com.intellij.openapi.diagnostic.ThrottledLogger +import com.intellij.openapi.progress.ProcessCanceledException +import com.intellij.openapi.progress.ProgressManager +import com.intellij.openapi.util.Condition +import com.intellij.openapi.util.Ref +import com.intellij.util.ExceptionUtil +import com.intellij.util.concurrency.ThreadingAssertions +import it.unimi.dsi.fastutil.objects.ObjectArrayList +import org.jetbrains.annotations.ApiStatus +import java.util.concurrent.TimeUnit +import java.util.concurrent.atomic.AtomicLong +import javax.swing.SwingUtilities + +/** + * Special UI Event Queue for runnables that interact with IJ Platform concepts, such as ModalityState and the RWI lock. + * + * ## Contract + * + * Normally, the execution on top of AWT EventQueue is fair (see the contracts of [java.awt.EventQueue]). + * IntelliJ Platform manipulates two entities that influence the ordering of events -- modality states and non-blocking write-intent lock. + * Hence, each event scheduled here has two additional values in its metadata: modality state and whether it runs under write-intent lock. + * An important part of this scheduler is that write-intent lock acquisition is non-blocking: if a runnable fails to acquire write-intent lock, + * it needs to be rescheduled to try later. + * + * The purpose of this class is to delay events in the way that retains ordering -- + * the main property of this scheduler is that + * **if there are two events `A` and `B` with the same metadata, and scheduling of `A` happens-before scheduling of `B`, then `A` will be executed before `B`** + * It is important that we do not enforce ordering guarantees between two runnables with different metadata: + * if event `A` is non-modal and event `B` is modal, it is undefined which one will execute first. + * Similarly, if event `A` requires lock and event `B` does not, it is undefined which one will execute first. + * + * ## Lock handling + * + * This class is an automaton with two states -- [WriteIntentLockMode.ALL] and [WriteIntentLockMode.UI_ONLY]. + * - In [WriteIntentLockMode.ALL] mode, the queue attempts to execute all runnables in order of their arrival. + * If a runnable does not require write-intent lock, it can be executed right away; + * if a runnable needs write-intent lock, [NonBlockingFlushQueue] acquires the lock in a non-blocking way. + * If a lock cannot be acquired, the queue transitions into [WriteIntentLockMode.UI_ONLY] + * and schedules a [WriteActionFinished] directive on termination of the existing write action. + * - In [WriteIntentLockMode.UI_ONLY] mode, the queue runs only runnables that do not require write-intent lock. + * When it reaches [WriteActionFinished], the queue transitions back to [WriteIntentLockMode.ALL]. + * + * ## Modality handling + * + * Modality states also influence the ordering of events. + * [NonBlockingFlushQueue] operates with two main queues: [writeIntentQueue] and [uiQueue]. + * When a polled runnable from each of those queues is deemed to have incorrect modality, + * the runnables gets to respective skipped queue ([skippedWriteIntentQueue] or [skippedUiQueue]). + * Once the modality state changes, skipped queues get appended back to the respective main queues. + * While this might not sound optimal, modality changes are relatively rare in the IntelliJ Platform, so the overhead of processing is negligible. + * + * ## Ordering guarantees + * + * The runnables are scheduled to the main queues in order of their arrival. + * The processing happens only on one thread, which ensures fairness of execution. + */ +@ApiStatus.Internal +class NonBlockingFlushQueue(private val threadingSupport: ThreadingSupport) { + + companion object { + private val LOG = Logger.getInstance(NonBlockingFlushQueue::class.java) + private val THROTTLED_LOG: ThrottledLogger = ThrottledLogger(LOG, TimeUnit.MINUTES.toMillis(1)) + } + + /** + * The main queue for runnables requiring write-intent lock. + */ + private val writeIntentQueue: BulkArrayQueue = BulkArrayQueue() + + /** + * Runnables, which require write-intent lock, but that were skipped due to incompatible modality states. + */ + private val skippedWriteIntentQueue: Ref> = Ref(ObjectArrayList(100)) + + /** + * The main queue for runnables not requiring write-intent lock. + */ + private val uiQueue: BulkArrayQueue = BulkArrayQueue() + + /** + * Runnables, which do not require write-intent lock, but that were skipped due to incompatible modality states. + */ + private val skippedUiQueue: Ref> = Ref(ObjectArrayList(100)) + + /** + * A monitor for interaction with [writeIntentQueue], [uiQueue] and [FLUSH_SCHEDULED]. + */ + private val lockObject = Any() + + /** + * We want to have as good latency guarantees for the executed events as possible. + * For this, we assign a monotonically increasing counter to each event and use it to determine the order of execution. + * If event A has a smaller associated counter than event B, and both A and B are allowed to be executed, then A will be executed first. + * This ensures that some late UI-only runnable does not run before all earlier WI runnables + */ + private val timeCounter = AtomicLong() + + /** + * Protection against too frequent scheduling requests. + */ + private var FLUSH_SCHEDULED: Boolean = false + + /** + * Requires lock on [lockObject] + */ + private fun setFlushScheduledGuard(value: Boolean) { + FLUSH_SCHEDULED = value + } + + /** + * Requires lock on [lockObject] + */ + private fun getFlushScheduledGuard(): Boolean { + return FLUSH_SCHEDULED + } + + /** + * States of [NonBlockingFlushQueue] + */ + private enum class WriteIntentLockMode { + ALL, + UI_ONLY, + } + + private sealed interface FlushQueueCommand { + val creationTime: Long + } + + private class RunnableInfo( + val runnable: Runnable, + val modalityState: ModalityState, + val isExpired: Condition<*>, + val needWriteIntent: Boolean, + override val creationTime: Long, + val queuedTimeNs: Long, + /** How many items were in queue at the moment this item was enqueued */ + val queueSize: Int, + // this field is not protected by intention. It gets mutated only on EDT + var wasInSkippedItems: Boolean, + ): FlushQueueCommand { + override fun toString(): String { + return "RunnableInfo[modalityState=$modalityState, needWi=$needWriteIntent, runnable=$runnable]" + } + } + + /** + * Special directive sent at the end of write action. + * When [com.intellij.openapi.application.impl.NonBlockingFlushQueue] encounters this directive, + * it transitions to the state which allows processing of skipped WI events. + */ + private class WriteActionFinished(override val creationTime: Long): FlushQueueCommand { + override fun toString(): String { + return "WriteActionFinished" + } + } + + private val FLUSH_NOW: ContextAwareRunnable = ContextAwareRunnable { flushNow() } + + /** + * Current state of [NonBlockingFlushQueue]. + * Can be changed only on EDT + */ + private var currentWriteIntentLockMode: WriteIntentLockMode = WriteIntentLockMode.ALL + + private fun pollNextEvent(): RunnableInfo? { + ThreadingAssertions.assertEventDispatchThread() + + val currentModality = LaterInvocator.getCurrentModalityState() + while (true) { + var skippedQueue: Ref> + var selectedQueue: BulkArrayQueue<*> + val incomingInfo = synchronized(lockObject) { + val topUiRunnable = uiQueue.peekFirst() + val topWiRunnable = if (currentWriteIntentLockMode == WriteIntentLockMode.ALL) writeIntentQueue.peekFirst() else null + if (topUiRunnable == null && topWiRunnable == null) { + // the queue is effectively empty + return null + } + else if (topUiRunnable != null && topWiRunnable == null) { + skippedQueue = skippedUiQueue + selectedQueue = uiQueue + topUiRunnable + } else if (topUiRunnable == null) { + // topWiRunnable != null + skippedQueue = skippedWriteIntentQueue + selectedQueue = writeIntentQueue + topWiRunnable + } else if (topWiRunnable == null) { + // topUiRunnable != null + skippedQueue = skippedUiQueue + selectedQueue = uiQueue + topUiRunnable + } else { + // both not null; select first arrived + if (topUiRunnable.creationTime < topWiRunnable.creationTime) { + skippedQueue = skippedUiQueue + selectedQueue = uiQueue + topUiRunnable + } else { + skippedQueue = skippedWriteIntentQueue + selectedQueue = writeIntentQueue + topWiRunnable + } + } + } ?: return null + + when (incomingInfo) { + is WriteActionFinished -> { + // there is a directive that write action has finished; it means that we can try to run skipped write actions again + require(currentWriteIntentLockMode == WriteIntentLockMode.UI_ONLY) { + "Write action finished, but FlushQueue is unexpectedly allowed to run all runnables" + } + if (threadingSupport.isWriteActionPending() || threadingSupport.isWriteActionInProgress()) { + threadingSupport.runWhenWriteActionIsCompleted { + requestFlush() + } + return null + } else { + uiQueue.pollFirst() // now we remove WriteActionFinished from the queue + currentWriteIntentLockMode = WriteIntentLockMode.ALL + } + } + is RunnableInfo -> { + if (!currentModality.accepts(incomingInfo.modalityState)) { + // Current modality is not acceptable; we must send such event to the Skipped Modality Queue + skippedQueue.get().add(incomingInfo) + incomingInfo.wasInSkippedItems = true + selectedQueue.pollFirst() // remove the skipped runnable + continue + } + if (incomingInfo.isExpired.value(null)) { + // this runnable is expired; let's continue searching for a suitable one + selectedQueue.pollFirst() // remove the skipped runnable + continue + } + requestFlush() // in case someone wrote "invokeLater { UIUtil.dispatchAllInvocationEvents(); }" + return incomingInfo + } + } + } + } + + private fun reincludeSkippedItems(list: Ref>, mainQueue: BulkArrayQueue) { + ThreadingAssertions.assertEventDispatchThread() + val size = list.get().size + if (size != 0) { + synchronized(lockObject) { + mainQueue.bulkEnqueueFirst(list.get()) + } + if (size < 100) { + list.get().clear() + } else { + list.set(ObjectArrayList(100)) + } + } + requestFlush() + } + + fun requestFlush() { + synchronized(lockObject) { + if (getFlushScheduledGuard()) { + return + } + if (uiQueue.isEmpty && (currentWriteIntentLockMode == WriteIntentLockMode.UI_ONLY || writeIntentQueue.isEmpty)) { + // so the queue is effectively empty + return + } + setFlushScheduledGuard(true) + } + + SwingUtilities.invokeLater(FLUSH_NOW) + } + + private fun flushNow() { + ThreadingAssertions.assertEventDispatchThread() + + synchronized(lockObject) { + setFlushScheduledGuard(false) + } + + val startTime = System.nanoTime() + while (true) { + val nextRunnable = pollNextEvent() + if (nextRunnable == null) { + // no more runnables to run... + break + } + + runNextEvent(nextRunnable) + + // we check if we want to abort only if a blocking event was executed + // other UI events are considered quick enough + if (InvocationUtil.priorityEventPending() || System.nanoTime() - startTime > 5_000_000) { + requestFlush() + break + } + } + } + + @Suppress("IncorrectCancellationExceptionHandling") + private fun runNextEvent(nextRunnable: RunnableInfo) { + val watcher = EventWatcher.getInstanceOrNull() + val waitingFinishedNs = System.nanoTime() + try { + if (nextRunnable.needWriteIntent) { + require(currentWriteIntentLockMode == WriteIntentLockMode.ALL) { + "Execution of Write-Intent-Locked runnables is not allowed in UI_ONLY mode" + } + // let's try to execute the runnable + val success = !(threadingSupport.isWriteActionPending() || threadingSupport.isWriteActionInProgress()) && threadingSupport.tryRunWriteIntentReadAction { + writeIntentQueue.pollFirst() // remove runnable from writeIntentQueue + resetThreadContext { + val progressManager = ProgressManager.getInstanceOrNull() + if (progressManager != null) { + progressManager.computePrioritized { + nextRunnable.runnable.run() + } + } else { + nextRunnable.runnable.run() + } + } + } + if (!success) { + // we failed, which means that there is a background write action; + // now we can to transition to the UI_ONLY state and execute only non-locking runnables + currentWriteIntentLockMode = WriteIntentLockMode.UI_ONLY + threadingSupport.runWhenWriteActionIsCompleted { + synchronized(lockObject) { + uiQueue.enqueue(WriteActionFinished(timeCounter.getAndIncrement())) + requestFlush() + } + } + } + } else { + uiQueue.pollFirst() // we are going to run this runnable; let's remove it from the queue + // no write-intent required; we can execute this directly + resetThreadContext { + nextRunnable.runnable.run() + } + } + } catch (_: ProcessCanceledException) { + // ignored + } catch (e: Throwable) { + if (ApplicationManager.getApplication().isUnitTestMode()) { + ExceptionUtil.rethrow(e) + } + if (!Logger.shouldRethrow(e)) { + LOG.error(e) + } + } finally { + reportStatistics(watcher, waitingFinishedNs, nextRunnable) + } + } + + private fun reportStatistics(watcher: EventWatcher?, waitingFinishedNs: Long, runnableInfo: RunnableInfo) { + if (watcher == null) return + val runnable: Runnable = runnableInfo.runnable + val executionFinishedNs = System.nanoTime() + val waitedInQueueNs: Long = waitingFinishedNs - runnableInfo.queuedTimeNs + val executionDurationNs = executionFinishedNs - waitingFinishedNs + + + //RC: ExceptionAnalyzer reports negative values here, but it is not clear there do they come from. + // The reasons I could think of now are: + // 1) oddities of .nanoTime() behavior under different CPU power-saving modes + // 2) changing .nanoTime() origin due to thead being relocated to another CPU + // 3) long overflow in (end-start) expression. + // those are straightforward reasons, but 1-2 was mostly solved (it seems to me) in a modern + // hardware/software, and 3 is hard to expect in our use-cases. Hence, negative values could be + // due to some other code bug I don't see right now. Safeguarding here prevents errors down the + // stack, but it also shifts value statistics + val waitedTimeInQueueNs_safe = if (waitedInQueueNs >= 0) waitedInQueueNs else 0 + val executionDurationNs_safe = if (executionDurationNs >= 0) executionDurationNs else 0 + + watcher.runnableTaskFinished(runnable, + waitedTimeInQueueNs_safe, + runnableInfo.queueSize, + executionDurationNs_safe, runnableInfo.wasInSkippedItems) + + if (waitedInQueueNs < 0 || executionDurationNs < 0) { + //maybe logs give us some hints about why the values are negative: + THROTTLED_LOG.info("waitedInQueueNs($waitedInQueueNs) | executionDurationNs($executionDurationNs) is negative -> unexpected state") + } + } + + /** + * Since the modality states are different now, we need to reevaluate the decisions about skipped runnables + */ + fun onModalityChanged() { + ThreadingAssertions.assertEventDispatchThread() + reincludeSkippedItems(skippedUiQueue, uiQueue) + reincludeSkippedItems(skippedWriteIntentQueue, writeIntentQueue) + } + + /** + * Adds [runnable] to the queue. + */ + fun push(modalityState: ModalityState, runnable: Runnable, isRunningUnderWriteIntent: Boolean, isExpired: Condition<*>) { + val creationTime = timeCounter.getAndIncrement() + val stamp = System.nanoTime() + synchronized(lockObject) { + val queueSize = writeIntentQueue.size() + uiQueue.size() // logically, this element goes to the end of the queue + val info = RunnableInfo(runnable, modalityState, isExpired, isRunningUnderWriteIntent, creationTime, stamp, queueSize, false) + if (isRunningUnderWriteIntent) { + writeIntentQueue.enqueue(info) + } else { + uiQueue.enqueue(info) + } + requestFlush() + } + } + + override fun toString(): String { + return "NonBlockingFlushQueue(currentWriteIntentLockMode=$currentWriteIntentLockMode, wiQueue=${writeIntentQueue.size()} elements, uiQueue=${uiQueue.size()}, skippedWriteIntentQueue=${skippedWriteIntentQueue.get().size} elements, skippedUiQueue=${skippedUiQueue.get().size} elements; flush scheduled: $FLUSH_SCHEDULED)" + } + + fun isFlushNow(runnable: Runnable): Boolean { + return runnable === FLUSH_NOW + } +} \ No newline at end of file diff --git a/platform/platform-tests/testSrc/com/intellij/concurrency/suites.kt b/platform/platform-tests/testSrc/com/intellij/concurrency/suites.kt index 950dc538590e..05f70576f3b6 100644 --- a/platform/platform-tests/testSrc/com/intellij/concurrency/suites.kt +++ b/platform/platform-tests/testSrc/com/intellij/concurrency/suites.kt @@ -48,6 +48,7 @@ import org.junit.platform.suite.api.Suite LockDowngradingTest::class, PlatformUtilitiesTest::class, WriteIntentReadActionTest::class, + NonBlockingFlushQueueTest::class, // propagation ThreadContextPropagationTest::class, diff --git a/platform/platform-tests/testSrc/com/intellij/openapi/application/impl/NonBlockingFlushQueueTest.kt b/platform/platform-tests/testSrc/com/intellij/openapi/application/impl/NonBlockingFlushQueueTest.kt new file mode 100644 index 000000000000..31452ee06d60 --- /dev/null +++ b/platform/platform-tests/testSrc/com/intellij/openapi/application/impl/NonBlockingFlushQueueTest.kt @@ -0,0 +1,533 @@ +// Copyright 2000-2025 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. +package com.intellij.openapi.application.impl + +import com.intellij.openapi.application.ModalityState +import com.intellij.openapi.application.UI +import com.intellij.openapi.application.asContextElement +import com.intellij.openapi.application.backgroundWriteAction +import com.intellij.platform.locking.impl.getGlobalThreadingSupport +import com.intellij.testFramework.common.timeoutRunBlocking +import com.intellij.testFramework.junit5.TestApplication +import com.intellij.util.ui.EDT +import kotlinx.coroutines.* +import kotlinx.coroutines.future.asCompletableFuture +import org.junit.jupiter.api.AfterEach +import org.junit.jupiter.api.Assertions.assertTrue +import org.junit.jupiter.api.BeforeEach +import org.junit.jupiter.api.Test +import java.util.concurrent.CompletableFuture +import java.util.concurrent.atomic.AtomicInteger +import javax.swing.SwingUtilities +import kotlin.test.assertEquals +import kotlin.test.assertFalse + +@TestApplication +class NonBlockingFlushQueueTest { + + lateinit var flushQueue: NonBlockingFlushQueue + lateinit var counter: AtomicInteger + + @BeforeEach + fun setUpFlushQueue() { + flushQueue = NonBlockingFlushQueue(getGlobalThreadingSupport()) + counter = AtomicInteger(0) + } + + @AfterEach + fun tearDownFlushQueue() { + } + + fun pushNonModalWI(runnable: Runnable) { + flushQueue.push(ModalityState.nonModal(), runnable, true) { false } + } + + fun pushNonModalUI(runnable: Runnable) { + flushQueue.push(ModalityState.nonModal(), runnable, false) { false } + } + + fun pushNonModalUIExpired(runnable: Runnable) { + flushQueue.push(ModalityState.nonModal(), runnable, false) { true } + } + + fun enterModal(modalEntity: Any) { + LaterInvocator.enterModal(modalEntity) + flushQueue.onModalityChanged() + } + + fun leaveModal(modalEntity: Any) { + flushQueue.onModalityChanged() + LaterInvocator.leaveModal(modalEntity) + } + + fun pushCurrentModalUI(runnable: Runnable) { + val modality = LaterInvocator.getCurrentModalityState() + flushQueue.push(modality, runnable, false) { false } + } + + fun pushCurrentModalWI(runnable: Runnable) { + val modality = LaterInvocator.getCurrentModalityState() + flushQueue.push(modality, runnable, true) { false } + } + + fun pushNonModalWIExpired(runnable: Runnable) { + flushQueue.push(ModalityState.nonModal(), runnable, true) { true } + } + + suspend fun spinQueue() { + withContext(Dispatchers.UI + ModalityState.any().asContextElement()) {} + suspendCancellableCoroutine { cont -> + SwingUtilities.invokeLater { + cont.resume(Unit, {_, _, _->}) + } + } + } + + @Test + fun `ordering is preserved for identical metadata - WI`(): Unit = timeoutRunBlocking { + withContext(Dispatchers.UI) { + pushNonModalWI { + assertTrue(EDT.isCurrentThreadEdt()) + assertEquals(0, counter.getAndIncrement()) + } + pushNonModalWI { + assertTrue(EDT.isCurrentThreadEdt()) + assertEquals(1, counter.getAndIncrement()) + } + } + spinQueue() + assertEquals(2, counter.get()) + } + + @Test + fun `ordering is preserved for identical metadata - UI only`(): Unit = timeoutRunBlocking { + withContext(Dispatchers.UI) { + pushNonModalUI { + assertTrue(EDT.isCurrentThreadEdt()) + assertEquals(0, counter.getAndIncrement()) + } + pushNonModalUI { + assertTrue(EDT.isCurrentThreadEdt()) + assertEquals(1, counter.getAndIncrement()) + } + } + spinQueue() + assertEquals(2, counter.get()) + } + + @Test + fun `UI event may overtake WI event when WI cannot acquire lock`(): Unit = timeoutRunBlocking { + val order = ArrayList() + val releaseBgWa = Job(coroutineContext.job) + val bgWaStarted = Job(coroutineContext.job) + + // Start a background write action that blocks WI execution on EDT + val bgJob = launch(Dispatchers.Default) { + backgroundWriteAction { + bgWaStarted.complete() + releaseBgWa.asCompletableFuture().join() + } + } + + withContext(Dispatchers.UI) { + bgWaStarted.join() + pushNonModalWI { + order.add("WI") + } + pushNonModalUI { + order.add("UI") + } + } + + // Pump the queue: UI should run while WI is delayed + spinQueue() + withContext(Dispatchers.UI) { + assertEquals(listOf("UI"), order) + } + + // Allow WI to proceed and ensure it eventually runs + releaseBgWa.complete() + bgJob.join() + + // Pump again to process the delayed WI + spinQueue() + withContext(Dispatchers.UI) { + assertEquals(listOf("UI", "WI"), order) + } + } + + @Test + fun `non-modal tasks are delayed while in modal state and resume after exit`(): Unit = timeoutRunBlocking { + val ran = CompletableFuture() + val modalEntity = Any() + + withContext(Dispatchers.UI) { + // Enter deeper modality + enterModal(modalEntity) + // Schedule a NON_MODAL runnable while current modality is modal + pushNonModalUI { + ran.complete(true) + } + } + + // Pump: task should NOT run yet due to modality mismatch + spinQueue() + assertFalse(ran.isDone) + + // Exit modality and notify our local queue about modality exit + withContext(Dispatchers.UI + ModalityState.any().asContextElement()) { + leaveModal(modalEntity) + } + + // Now the task should be re-included and executed + spinQueue() + assertTrue(ran.getNow(false)) + } + + @Test + fun `expired tasks are skipped`(): Unit = timeoutRunBlocking { + val ran = CompletableFuture() + withContext(Dispatchers.UI) { + pushNonModalUIExpired { + ran.complete(true) + } + } + spinQueue() + // The runnable must not be executed because it's expired + assertFalse(ran.isDone) + } + + @Test + fun `modal UI runs immediately while non-modal is deferred`(): Unit = timeoutRunBlocking { + val order = ArrayList() + val modalEntity = Any() + + withContext(Dispatchers.UI) { + enterModal(modalEntity) + pushCurrentModalUI { order.add("modal-ui") } + pushNonModalUI { order.add("non-modal-ui") } + } + + // Pump once: modal-ui should execute, non-modal should wait + spinQueue() + assertEquals(listOf("modal-ui"), order) + + withContext(Dispatchers.UI + ModalityState.any().asContextElement()) { + leaveModal(modalEntity) + } + + spinQueue() + assertEquals(listOf("modal-ui", "non-modal-ui"), order) + } + + data class ModalityWrapper(val name: String) + + @Test + fun `nested modalities keep non-modal deferred until full exit`(): Unit = timeoutRunBlocking { + val ran = CompletableFuture() + val outer = ModalityWrapper("Outer") + val inner = ModalityWrapper("Inner") + + withContext(Dispatchers.UI) { + enterModal(outer) + enterModal(inner) + pushNonModalUI { ran.complete(true) } + } + + spinQueue() + assertFalse(ran.isDone) + + // Leave only inner, still in outer modality -> still deferred + withContext(Dispatchers.UI + ModalityState.any().asContextElement()) { + leaveModal(inner) + } + spinQueue() + assertFalse(ran.isDone) + + // Leave outer -> now it should run + withContext(Dispatchers.UI + ModalityState.any().asContextElement()) { + leaveModal(outer) + } + spinQueue() + withContext(Dispatchers.UI) { + assertTrue(ran.getNow(false)) + } + } + + @Test + fun `WI ordering preserved across UI_ONLY then back to ALL`(): Unit = timeoutRunBlocking { + val order = ArrayList() + val releaseBgWa = Job(coroutineContext.job) + val bgWaStarted = Job(coroutineContext.job) + + val bgJob = launch(Dispatchers.Default) { + backgroundWriteAction { + bgWaStarted.complete() + releaseBgWa.asCompletableFuture().join() + } + } + + withContext(Dispatchers.UI) { + bgWaStarted.join() + // This WI will fail to acquire, switching to UI_ONLY and be skipped + pushNonModalWI { order.add("WI-A") } + // UI should run while WI is delayed + pushNonModalUI { order.add("UI-X") } + // Another WI enqueued while still UI_ONLY; must keep A before B on resume + pushNonModalWI { order.add("WI-B") } + } + + spinQueue() + // Only UI-X should have run so far + assertEquals(listOf("UI-X"), order) + + // Release background write action so WI can proceed; queue should switch back to ALL via WriteActionFinished + releaseBgWa.complete() + bgJob.join() + + spinQueue() + withContext(Dispatchers.UI) { + assertEquals(listOf("UI-X", "WI-A", "WI-B"), order) + } + } + + @Test + fun `entering modality during UI_ONLY delays skipped WI until modality exit`(): Unit = timeoutRunBlocking { + val order = ArrayList() + val modal = ModalityWrapper("Have WI before") + val releaseBgWa = Job(coroutineContext.job) + val bgWaStarted = Job(coroutineContext.job) + + val bgJob = launch(Dispatchers.Default) { + backgroundWriteAction { + bgWaStarted.complete() + releaseBgWa.asCompletableFuture().join() + } + } + + withContext(Dispatchers.UI) { + bgWaStarted.join() + // Fail WI acquisition -> UI_ONLY; this WI goes to skipped WI list + pushNonModalWI { order.add("WI-before-modal") } + // forcefully try to execute this event + yield() + // the event was delayed; we are now in UI_ONLY state + assertTrue(order.isEmpty()) + // Enter modality while in UI_ONLY; skipped WI must be moved under skipped modality and not run until exit + enterModal(modal) + // Even UI in current modal should run + pushCurrentModalUI { order.add("modal-UI") } + } + + spinQueue() + // Only modal UI should execute + assertEquals(listOf("modal-UI"), order) + + // Allow WI, but still inside modality -> WI must not run yet + releaseBgWa.complete() + bgJob.join() + spinQueue() + assertEquals(listOf("modal-UI"), order) + + // Exit modality -> WI should be re-included and executed + withContext(Dispatchers.UI + ModalityState.any().asContextElement()) { + leaveModal(modal) + } + spinQueue() + withContext(Dispatchers.UI) { + assertEquals(listOf("modal-UI", "WI-before-modal"), order) + } + } + + @Test + fun `expired WI tasks are skipped`(): Unit = timeoutRunBlocking { + val ran = CompletableFuture() + withContext(Dispatchers.UI) { + pushNonModalWIExpired { ran.complete(true) } + } + spinQueue() + assertFalse(ran.isDone) + } + + @Test + fun `deferred non-modal UI preserves ordering after modality exit`(): Unit = timeoutRunBlocking { + val order = ArrayList() + val modal = ModalityWrapper("Have non-modal after") + + withContext(Dispatchers.UI) { + enterModal(modal) + pushNonModalUI { order.add("A") } + pushNonModalUI { order.add("B") } + } + + spinQueue() + assertEquals(emptyList(), order) + + withContext(Dispatchers.UI + ModalityState.any().asContextElement()) { + leaveModal(modal) + } + + spinQueue() + assertEquals(listOf("A", "B"), order) + } + + @Test + fun `WI enqueued during UI_ONLY preserve relative order on resume`(): Unit = timeoutRunBlocking { + val order = ArrayList() + val releaseBgWa = Job(coroutineContext.job) + val bgWaStarted = Job(coroutineContext.job) + + val bgJob = launch(Dispatchers.Default) { + backgroundWriteAction { + bgWaStarted.complete() + releaseBgWa.asCompletableFuture().join() + } + } + + withContext(Dispatchers.UI) { + bgWaStarted.join() + pushNonModalWI { order.add("WI-1") } + pushNonModalWI { order.add("WI-2") } + } + + spinQueue() + // While UI_ONLY, neither WI should have run + assertEquals(emptyList(), order) + + releaseBgWa.complete() + bgJob.join() + + spinQueue() + withContext(Dispatchers.UI) { + assertEquals(listOf("WI-1", "WI-2"), order) + } + } + + @Test + fun `modal WI runs immediately and non-modal WI is deferred until modality exit`(): Unit = timeoutRunBlocking { + val order = ArrayList() + val modalEntity = ModalityWrapper("Have mixed inside") + + withContext(Dispatchers.UI) { + enterModal(modalEntity) + // enqueue modal WI, then non-modal WI, then one more modal WI + pushCurrentModalWI { order.add("modal-wi-1") } + pushNonModalWI { order.add("non-modal-wi-1") } + pushCurrentModalWI { order.add("modal-wi-2") } + } + + // While in modal state, only modal WI should execute and preserve their relative order + spinQueue() + assertEquals(listOf("modal-wi-1", "modal-wi-2"), order) + + // Exit modality, then the deferred non-modal WI should run + withContext(Dispatchers.UI + ModalityState.any().asContextElement()) { + leaveModal(modalEntity) + } + + spinQueue() + assertEquals(listOf("modal-wi-1", "modal-wi-2", "non-modal-wi-1"), order) + } + + @Test + fun `background write action inside modality - modal WI skipped until WA ends, ordered before non-modal UI on exit`(): Unit = timeoutRunBlocking { + val order = ArrayList() + val modal = ModalityWrapper("Has BG wa inside") + val releaseBgWa = Job(coroutineContext.job) + val bgWaStarted = Job(coroutineContext.job) + + val bgJob = launch(Dispatchers.Default) { + backgroundWriteAction { + bgWaStarted.complete() + releaseBgWa.asCompletableFuture().join() + } + } + + withContext(Dispatchers.UI + ModalityState.any().asContextElement()) { + enterModal(modal) + bgWaStarted.join() + // Enqueue a modal WI that will fail to acquire WI due to background WA + pushCurrentModalWI { order.add("M-WI-1") } + // This modal UI should still run + pushCurrentModalUI { order.add("M-UI-1") } + // Non-modal UI is not acceptable by current modality and must be deferred until modality exit + pushNonModalUI { order.add("NM-UI-1") } + // Another modal WI to check ordering among WI tasks + pushCurrentModalWI { order.add("M-WI-2") } + } + + // First pump: only modal UI should execute, WI are skipped, non-modal UI deferred by modality + spinQueue() + assertEquals(listOf("M-UI-1"), order) + + // Allow WI to proceed while still inside modality: skipped modal WI should now run in order + releaseBgWa.complete() + bgJob.join() + + spinQueue() + withContext(Dispatchers.UI + ModalityState.any().asContextElement()) { + // Non-modal UI still deferred due to modality; modal WI should have executed in FIFO order + assertEquals(listOf("M-UI-1", "M-WI-1", "M-WI-2"), order) + } + + // Exit modality: now the deferred non-modal UI should run + withContext(Dispatchers.UI + ModalityState.any().asContextElement()) { + leaveModal(modal) + } + + spinQueue() + withContext(Dispatchers.UI) { + assertEquals(listOf("M-UI-1", "M-WI-1", "M-WI-2", "NM-UI-1"), order) + } + } + + @Test + fun `modality exit re-includes skipped WI and skipped modality in correct order`(): Unit = timeoutRunBlocking { + val order = ArrayList() + val modal = ModalityWrapper("Complex exit") + val releaseBgWa = Job(coroutineContext.job) + val bgWaStarted = Job(coroutineContext.job) + + val bgJob = launch(Dispatchers.Default) { + backgroundWriteAction { + bgWaStarted.complete() + releaseBgWa.asCompletableFuture().join() + } + } + + withContext(Dispatchers.UI + ModalityState.any().asContextElement()) { + enterModal(modal) + bgWaStarted.join() + // Two modal WI that will fail to acquire WI and be skipped into skippedWriteIntentQueue + pushCurrentModalWI { order.add("M-WI-1") } + pushCurrentModalWI { order.add("M-WI-2") } + // Two non-modal UI that are incompatible with current modality and will be placed into skippedModalityQueue + pushNonModalUI { order.add("NM-UI-1") } + pushNonModalUI { order.add("NM-UI-2") } + // A modal UI that should execute immediately before we exit modality + pushCurrentModalUI { order.add("M-UI-1") } + } + + // Pump once to establish skipped queues: only modal UI should run + spinQueue() + assertEquals(listOf("M-UI-1"), order) + + // Exit modality while still in UI_ONLY due to background WA + withContext(Dispatchers.UI + ModalityState.any().asContextElement()) { + leaveModal(modal) + } + + // On modality exit, queue re-includes skipped WI first, then skipped modality; since still UI_ONLY, + // WI won’t run yet, but non-modal UI will now be acceptable and execute in FIFO order + spinQueue() + assertEquals(listOf("M-UI-1", "NM-UI-1", "NM-UI-2"), order) + + // Finish background write action; WI should now be re-included and execute in FIFO order after UI + releaseBgWa.complete() + bgJob.join() + + spinQueue() + withContext(Dispatchers.UI) { + assertEquals(listOf("M-UI-1", "NM-UI-1", "NM-UI-2", "M-WI-1", "M-WI-2"), order) + } + } +}