mirror of
https://gitflic.ru/project/openide/openide.git
synced 2026-09-27 10:03:11 +07:00
[threading] IJPL-206769: Introduce NonBlockingFlushQueue
This is a version of `FlushQueue` that allows non-blocking acquisition of write-intent lock with retained order of scheduled events GitOrigin-RevId: 9e37487a0378c7bad404175c59a5896713c0e217
This commit is contained in:
committed by
intellij-monorepo-bot
parent
225c3d483b
commit
640601b170
@@ -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<T> {
|
||||
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];
|
||||
|
||||
+430
@@ -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<RunnableInfo> = BulkArrayQueue()
|
||||
|
||||
/**
|
||||
* Runnables, which require write-intent lock, but that were skipped due to incompatible modality states.
|
||||
*/
|
||||
private val skippedWriteIntentQueue: Ref<ObjectArrayList<RunnableInfo>> = Ref(ObjectArrayList(100))
|
||||
|
||||
/**
|
||||
* The main queue for runnables not requiring write-intent lock.
|
||||
*/
|
||||
private val uiQueue: BulkArrayQueue<FlushQueueCommand> = BulkArrayQueue()
|
||||
|
||||
/**
|
||||
* Runnables, which do not require write-intent lock, but that were skipped due to incompatible modality states.
|
||||
*/
|
||||
private val skippedUiQueue: Ref<ObjectArrayList<RunnableInfo>> = 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<ObjectArrayList<RunnableInfo>>
|
||||
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<ObjectArrayList<RunnableInfo>>, mainQueue: BulkArrayQueue<in RunnableInfo>) {
|
||||
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<Unit, Throwable> {
|
||||
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
|
||||
}
|
||||
}
|
||||
@@ -48,6 +48,7 @@ import org.junit.platform.suite.api.Suite
|
||||
LockDowngradingTest::class,
|
||||
PlatformUtilitiesTest::class,
|
||||
WriteIntentReadActionTest::class,
|
||||
NonBlockingFlushQueueTest::class,
|
||||
|
||||
// propagation
|
||||
ThreadContextPropagationTest::class,
|
||||
|
||||
+533
@@ -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<String>()
|
||||
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<Boolean>()
|
||||
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<Boolean>()
|
||||
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<String>()
|
||||
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<Boolean>()
|
||||
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<String>()
|
||||
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<String>()
|
||||
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<Boolean>()
|
||||
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<String>()
|
||||
val modal = ModalityWrapper("Have non-modal after")
|
||||
|
||||
withContext(Dispatchers.UI) {
|
||||
enterModal(modal)
|
||||
pushNonModalUI { order.add("A") }
|
||||
pushNonModalUI { order.add("B") }
|
||||
}
|
||||
|
||||
spinQueue()
|
||||
assertEquals(emptyList<String>(), 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<String>()
|
||||
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<String>(), 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<String>()
|
||||
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<String>()
|
||||
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<String>()
|
||||
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)
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user