mirror of
https://gitflic.ru/project/openide/openide.git
synced 2026-09-27 10:03:11 +07:00
[vcs] IJPL-173924 Fix value runNow stuck in SingleTaskRunner
GitOrigin-RevId: 8131ac8743d128219096ef6ea547a47edaf79b42
This commit is contained in:
committed by
intellij-monorepo-bot
parent
44cf8223ef
commit
55018c676b
+7
-5
@@ -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<Unit>(replay = 1, onBufferOverflow = BufferOverflow.DROP_OLDEST)
|
||||
private val runNow = Channel<Unit>(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
|
||||
}
|
||||
|
||||
|
||||
@@ -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<TimeoutCancellationException> {
|
||||
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")
|
||||
|
||||
Reference in New Issue
Block a user