From 28aba75be7d373f473aec315684c42be5fa3bcb3 Mon Sep 17 00:00:00 2001 From: Konstantin Nisht Date: Wed, 18 Jun 2025 09:51:50 +0200 Subject: [PATCH] [threading] IJPL-192074: Improve progress guarantees for `MergingUpdateQueue` GitOrigin-RevId: 03560101147fdd968189c479134eaab70bcaff2c --- .../util/ui/update/MergingUpdateQueue.kt | 10 ++++++--- .../MergingUpdateQueuePropagationTest.kt | 21 ++++++++++++++++--- 2 files changed, 25 insertions(+), 6 deletions(-) diff --git a/platform/ide-core/src/com/intellij/util/ui/update/MergingUpdateQueue.kt b/platform/ide-core/src/com/intellij/util/ui/update/MergingUpdateQueue.kt index 2f412edc7927..2a5d23ea79e3 100644 --- a/platform/ide-core/src/com/intellij/util/ui/update/MergingUpdateQueue.kt +++ b/platform/ide-core/src/com/intellij/util/ui/update/MergingUpdateQueue.kt @@ -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() } diff --git a/platform/platform-tests/testSrc/com/intellij/util/concurrency/MergingUpdateQueuePropagationTest.kt b/platform/platform-tests/testSrc/com/intellij/util/concurrency/MergingUpdateQueuePropagationTest.kt index 9683b5451fda..8bfc1cba919d 100644 --- a/platform/platform-tests/testSrc/com/intellij/util/concurrency/MergingUpdateQueuePropagationTest.kt +++ b/platform/platform-tests/testSrc/com/intellij/util/concurrency/MergingUpdateQueuePropagationTest.kt @@ -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) + } }