[vcs] Fix race in SingleTaskRunner

GitOrigin-RevId: df99f8eef9fa5a697549817cdcc75154842a2a41
This commit is contained in:
Ilia.Shulgin
2025-10-20 19:19:33 +00:00
committed by intellij-monorepo-bot
parent 36e0c8eccc
commit 92bef67ddd
4 changed files with 39 additions and 77 deletions
@@ -34,6 +34,8 @@ import org.jetbrains.annotations.CalledInAny
import java.lang.Runnable
import kotlin.time.Duration.Companion.milliseconds
private val REFRESH_DELAY = 100.milliseconds
@ApiStatus.Internal
class CommitChangesViewWithToolbarPanel(
changesView: ChangesListView,
@@ -42,7 +44,7 @@ class CommitChangesViewWithToolbarPanel(
val project: Project get() = changesView.project
private val settings get() = ChangesViewSettings.getInstance(project)
private val refresher = SingleTaskRunner(cs, 100.milliseconds) {
private val refresher = SingleTaskRunner(cs) {
refreshView()
}
@@ -57,8 +59,8 @@ class CommitChangesViewWithToolbarPanel(
init {
refresher.start()
cs.launch(Dispatchers.UI) {
refresher.getPendingTasksFlow().collect {
isBusy -> changesView.setPaintBusy(isBusy)
refresher.getIdleFlow().collect { idle ->
changesView.setPaintBusy(!idle)
}
}
}
@@ -119,13 +121,14 @@ class CommitChangesViewWithToolbarPanel(
@CalledInAny
private fun scheduleRefresh(withDelay: Boolean, @RequiresBackgroundThread callback: Runnable? = null) {
if (withDelay) {
if (!withDelay && callback == null) {
refresher.request()
return
}
else {
refresher.requestNow()
}
cs.launch {
if (withDelay) delay(REFRESH_DELAY)
refresher.request()
refresher.awaitNotBusy()
callback?.run()
}
@@ -6,11 +6,7 @@ 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.Channel
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.combine
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.flow.*
import org.jetbrains.annotations.ApiStatus
import kotlin.time.Duration
@@ -22,49 +18,39 @@ import kotlin.time.Duration
@ApiStatus.Internal
class SingleTaskRunner(
cs: CoroutineScope,
private val delay: Duration = Duration.Companion.ZERO,
private val task: suspend () -> Unit,
) {
private val requested = MutableStateFlow(false)
private val busy = MutableStateFlow(false)
/**
* Represents the state of the runner:
* 0 - idle,
* 1 - a single task is either queued or running,
* 2 - has one running task and one queued task.
*/
private val requestCounter = MutableStateFlow(0)
private val runNow = Channel<Unit>(capacity = Channel.CONFLATED)
fun getPendingTasksFlow(): Flow<Boolean> = requested.combine(busy) { requested, busy -> !requested && !busy }
fun getIdleFlow(): Flow<Boolean> = requestCounter.map { it == 0 }
private val processorJob = cs.launch(Dispatchers.Default, CoroutineStart.LAZY) {
try {
while (true) {
checkCanceled()
requested.first { it }
if (delay.isPositive()) {
withTimeoutOrNull(delay) { runNow.receive() }
}
busy.value = true
requested.value = false
requestCounter.first { it > 0 }
checkCanceled()
runCatching {
task()
}.getOrHandleException { LOG.error("Task failed", it) }
busy.value = false
requestCounter.update { it.dec() }
}
}
finally {
requested.value = false
busy.value = false
requestCounter.value = 0
LOG.debug { "Task processing finished" }
}
}
fun request() {
if (processorJob.isCancelled) return
requested.value = true
}
fun requestNow() {
if (processorJob.isCancelled) return
runNow.trySend(Unit)
requested.value = true
requestCounter.update { it.inc().coerceAtMost(2) }
}
fun start() {
@@ -76,10 +62,15 @@ class SingleTaskRunner(
* Await the state where the task is not executed and there are no requests to do so.
*/
suspend fun awaitNotBusy() {
getPendingTasksFlow().first { it }
getIdleFlow().first { it }
}
companion object {
private val LOG = logger<SingleTaskRunner>()
fun delayedTaskRunner(cs: CoroutineScope, delay: Duration, task: suspend () -> Unit): SingleTaskRunner = SingleTaskRunner(cs) {
delay(delay)
task()
}
}
}
@@ -5,18 +5,13 @@ import com.intellij.platform.util.coroutines.childScope
import com.intellij.platform.vcs.impl.shared.SingleTaskRunner
import com.intellij.testFramework.LoggedErrorProcessor
import com.intellij.testFramework.common.timeoutRunBlocking
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.TimeoutCancellationException
import kotlinx.coroutines.cancel
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.withTimeout
import org.junit.jupiter.api.Assertions
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
@@ -43,41 +38,13 @@ internal class SingleTaskRunnerTest {
}
}
@Test
fun `test instant execution`() = timeoutRunBlocking {
var counter = 0
withRunner({ counter++ }, delay = Duration.INFINITE) {
start()
request()
requestNow()
awaitNotBusy()
assertEquals(1, counter)
}
}
@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
val taskStarted = CompletableDeferred<Unit>()
val runAllowed = MutableStateFlow(false)
withRunner({
taskStarted.complete(Unit)
runAllowed.first { it }
counter++
}) {
@@ -85,11 +52,13 @@ internal class SingleTaskRunnerTest {
repeat(10) {
request()
taskStarted.await()
}
runAllowed.value = true
awaitNotBusy()
assertEquals(1, counter)
assertEquals(2, counter)
counter = 0
runAllowed.value = false
request()
assertThrows<TimeoutCancellationException> {
@@ -97,10 +66,10 @@ internal class SingleTaskRunnerTest {
awaitNotBusy()
}
}
assertEquals(1, counter)
assertEquals(0, counter)
runAllowed.value = true
awaitNotBusy()
assertEquals(2, counter)
assertEquals(1, counter)
}
}
@@ -131,11 +100,10 @@ internal class SingleTaskRunnerTest {
private inline fun CoroutineScope.withRunner(
noinline task: suspend () -> Unit,
delay: Duration = Companion.ZERO,
consumer: SingleTaskRunner.() -> Unit,
) {
val cs = childScope("BG runner")
SingleTaskRunner(cs, delay, task).also(consumer)
SingleTaskRunner(cs, task).also(consumer)
cs.cancel()
}
}
@@ -69,7 +69,7 @@ class GitUntrackedFilesHolder internal constructor(
get() = untrackedFiles.initialized
init {
updateRunner = SingleTaskRunner(cs, 500.milliseconds, ::update)
updateRunner = SingleTaskRunner.delayedTaskRunner(cs, 500.milliseconds, ::update)
cs.launch(start = CoroutineStart.UNDISPATCHED) {
try {
project.serviceAsync<InitialVfsRefreshService>().awaitInitialVfsRefreshFinished()