From 55018c676b27943a50005b25ee3b3bfdb7c6885b Mon Sep 17 00:00:00 2001 From: "Ilia.Shulgin" Date: Thu, 9 Oct 2025 15:48:54 +0200 Subject: [PATCH] [vcs] IJPL-173924 Fix value runNow stuck in SingleTaskRunner GitOrigin-RevId: 8131ac8743d128219096ef6ea547a47edaf79b42 --- .../vcs/impl/shared/SingleTaskRunner.kt | 12 ++++++---- .../intellij/vcsUtil/SingleTaskRunnerTest.kt | 23 +++++++++++++++++-- 2 files changed, 28 insertions(+), 7 deletions(-) diff --git a/platform/vcs-impl/shared/src/com/intellij/platform/vcs/impl/shared/SingleTaskRunner.kt b/platform/vcs-impl/shared/src/com/intellij/platform/vcs/impl/shared/SingleTaskRunner.kt index eadf17b86af2..dff47ebea0ed 100644 --- a/platform/vcs-impl/shared/src/com/intellij/platform/vcs/impl/shared/SingleTaskRunner.kt +++ b/platform/vcs-impl/shared/src/com/intellij/platform/vcs/impl/shared/SingleTaskRunner.kt @@ -6,8 +6,10 @@ import com.intellij.openapi.diagnostic.getOrHandleException import com.intellij.openapi.diagnostic.logger import com.intellij.openapi.progress.checkCanceled import kotlinx.coroutines.* -import kotlinx.coroutines.channels.BufferOverflow -import kotlinx.coroutines.flow.* +import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.flow.MutableStateFlow +import kotlinx.coroutines.flow.combine +import kotlinx.coroutines.flow.first import org.jetbrains.annotations.ApiStatus import kotlin.time.Duration @@ -25,7 +27,7 @@ class SingleTaskRunner( private val requested = MutableStateFlow(false) private val busy = MutableStateFlow(false) - private val runNow = MutableSharedFlow(replay = 1, onBufferOverflow = BufferOverflow.DROP_OLDEST) + private val runNow = Channel(capacity = Channel.CONFLATED) private val processorJob = cs.launch(Dispatchers.Default, CoroutineStart.LAZY) { try { @@ -33,7 +35,7 @@ class SingleTaskRunner( checkCanceled() requested.first { it } if (delay.isPositive()) { - withTimeoutOrNull(delay) { runNow.firstOrNull() } + withTimeoutOrNull(delay) { runNow.receive() } } busy.value = true requested.value = false @@ -58,7 +60,7 @@ class SingleTaskRunner( fun requestNow() { if (processorJob.isCancelled) return - runNow.tryEmit(Unit) + runNow.trySend(Unit) requested.value = true } diff --git a/platform/vcs-impl/testSrc/com/intellij/vcsUtil/SingleTaskRunnerTest.kt b/platform/vcs-impl/testSrc/com/intellij/vcsUtil/SingleTaskRunnerTest.kt index d8ae74d222b7..4ccd0c69c7d0 100644 --- a/platform/vcs-impl/testSrc/com/intellij/vcsUtil/SingleTaskRunnerTest.kt +++ b/platform/vcs-impl/testSrc/com/intellij/vcsUtil/SingleTaskRunnerTest.kt @@ -16,6 +16,7 @@ import org.junit.jupiter.api.Assertions.assertEquals import org.junit.jupiter.api.Test import org.junit.jupiter.api.assertThrows import kotlin.time.Duration +import kotlin.time.Duration.Companion internal class SingleTaskRunnerTest { @Test @@ -45,7 +46,7 @@ internal class SingleTaskRunnerTest { @Test fun `test instant execution`() = timeoutRunBlocking { var counter = 0 - withRunner({ counter++ }, delay = Duration.Companion.INFINITE) { + withRunner({ counter++ }, delay = Duration.INFINITE) { start() request() requestNow() @@ -54,6 +55,24 @@ internal class SingleTaskRunnerTest { } } + @Test + fun `test delayed execution after instant`() = timeoutRunBlocking { + var counter = 0 + withRunner({ counter++ }, delay = Duration.INFINITE) { + start() + request() + requestNow() + awaitNotBusy() + request() + assertThrows { + withTimeout(100) { + awaitNotBusy() + } + } + assertEquals(1, counter) + } + } + @Test fun `test execution debounced`() = timeoutRunBlocking { var counter = 0 @@ -112,7 +131,7 @@ internal class SingleTaskRunnerTest { private inline fun CoroutineScope.withRunner( noinline task: suspend () -> Unit, - delay: Duration = Duration.Companion.ZERO, + delay: Duration = Companion.ZERO, consumer: SingleTaskRunner.() -> Unit, ) { val cs = childScope("BG runner")