[threading] IJPL-192074: Improve progress guarantees for MergingUpdateQueue

GitOrigin-RevId: 03560101147fdd968189c479134eaab70bcaff2c
This commit is contained in:
Konstantin Nisht
2025-06-18 09:33:51 +00:00
committed by intellij-monorepo-bot
parent 8723132ebf
commit 28aba75be7
2 changed files with 25 additions and 6 deletions
@@ -24,6 +24,7 @@ import com.intellij.util.ui.EdtInvocationManager
import com.intellij.util.ui.update.UiNotifyConnector.Companion.installOn
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Job
import kotlinx.coroutines.currentCoroutineContext
import kotlinx.coroutines.ensureActive
import org.jetbrains.annotations.ApiStatus.Internal
import org.jetbrains.annotations.ApiStatus.Obsolete
@@ -342,15 +343,18 @@ open class MergingUpdateQueue @JvmOverloads constructor(
isFlushing = true
try {
val all = flushScheduledUpdates() ?: return@scheduleTask
try {
coroutineContext.ensureActive()
currentCoroutineContext().ensureActive()
} catch (e : CancellationException) {
all.forEachGuaranteed(Update::setRejected)
if (waiterForMerge.isDisposed) {
flushScheduledUpdates()?.forEachGuaranteed(Update::setRejected)
}
throw e
}
val all = flushScheduledUpdates() ?: return@scheduleTask
for (update in all) {
update.setProcessed()
}
@@ -21,11 +21,10 @@ import kotlinx.coroutines.Job
import kotlinx.coroutines.delay
import kotlinx.coroutines.future.asCompletableFuture
import org.assertj.core.api.Assertions.assertThat
import org.junit.jupiter.api.Assertions.assertEquals
import org.junit.jupiter.api.Assertions.assertFalse
import org.junit.jupiter.api.Assertions.assertTrue
import org.junit.jupiter.api.Assertions.*
import org.junit.jupiter.api.Test
import org.junit.jupiter.api.fail
import java.util.concurrent.atomic.AtomicInteger
import kotlin.coroutines.EmptyCoroutineContext
import kotlin.time.Duration.Companion.milliseconds
@@ -134,4 +133,20 @@ class MergingUpdateQueuePropagationTest {
delay(400.milliseconds)
assertTrue(update.isRejected)
}
@Test
fun `frequent flush of merging queue retains progress`(): Unit = timeoutRunBlocking(context = Dispatchers.Default) {
val queueProcessingJob = Job()
val queue = MergingUpdateQueue("test queue", 200, true, null, null, null, Alarm.ThreadToUse.POOLED_THREAD, coroutineScope = CoroutineScope(queueProcessingJob))
val counter = AtomicInteger()
val update = Update.create(1) {
counter.incrementAndGet()
}
queue.queue(update)
repeat(1000) {
queue.sendFlush()
}
delay(300)
assertThat(counter.get()).isEqualTo(1)
}
}