[plugin management] IJPL-248432 Fix flaky tests for PluginUpdateService

- Use MutableSharedFlow for plugin updates
PluginDto is mutated in-place when merging updates, so MutableStateFlow conflates all updates into one.
Fixing PluginDto is costly, so for now just make sure that every update is consumed.

 - Guard against exceptions in PluginUpdatesService
If external code throws an Exception, it won't cancel the coroutine

 - Dispatch callbacks in a dedicated coroutine
Hopping to UI dispatcher may park plugin update collector on EDT thread and lead to unwanted delays.

Merge-request: IJ-MR-210733
Merged-by: Konstantin Ripak <konstantin.ripak@jetbrains.com>

GitOrigin-RevId: 958d2c19913d16d4b1485d3cd7dd2cbbbd271f86
This commit is contained in:
Kostya Ripak
2026-08-11 14:33:03 +00:00
committed by intellij-monorepo-bot
parent e8e57c8082
commit e7318e3b1f
2 changed files with 42 additions and 22 deletions
@@ -7,6 +7,8 @@ import com.intellij.ide.plugins.PluginEnableStateChangedListener
import com.intellij.ide.plugins.PluginStateListener
import com.intellij.ide.plugins.PluginStateManager
import com.intellij.ide.plugins.api.PluginDto
import com.intellij.openapi.diagnostic.getOrHandleException
import com.intellij.openapi.diagnostic.logger
import com.intellij.openapi.extensions.PluginId
import com.intellij.openapi.updateSettings.impl.PluginUpdateHandler
import kotlinx.coroutines.CoroutineScope
@@ -23,6 +25,7 @@ import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext
import org.jetbrains.annotations.ApiStatus
import kotlin.coroutines.cancellation.CancellationException
import kotlin.time.Duration.Companion.milliseconds
@OptIn(FlowPreview::class)
@@ -33,6 +36,10 @@ class DefaultPluginUpdatesProvider(private val coroutineScope: CoroutineScope) :
private val updateRequestFlow = MutableSharedFlow<Unit>(replay = 1, extraBufferCapacity = 1, onBufferOverflow = BufferOverflow.DROP_OLDEST)
private var lastPluginUpdates: PluginUpdatesEvent? = null
companion object {
private val LOG = logger<DefaultPluginUpdatesProvider>()
}
init {
PluginStateManager.addStateListener(object : PluginStateListener {
override fun install(descriptor: IdeaPluginDescriptor) {
@@ -64,12 +71,14 @@ class DefaultPluginUpdatesProvider(private val coroutineScope: CoroutineScope) :
.debounce(300.milliseconds)
.collectLatest {
updateMutex.withLock {
val model = (PluginUpdateHandler.getInstance().loadAndStorePluginUpdates(null))
val pluginUpdates = PluginUpdatesEvent(model.pluginUpdates.markLocal(),
model.disabledPluginUpdates.markLocal(),
model.updatesFromCustomRepositories.markLocal())
lastPluginUpdates = pluginUpdates
emitUpdates(pluginUpdates)
runCatching {
val model = (PluginUpdateHandler.getInstance().loadAndStorePluginUpdates(null))
val pluginUpdates = PluginUpdatesEvent(model.pluginUpdates.markLocal(),
model.disabledPluginUpdates.markLocal(),
model.updatesFromCustomRepositories.markLocal())
lastPluginUpdates = pluginUpdates
emitUpdates(pluginUpdates)
}.getOrHandleException { e -> LOG.warn("Failed to load plugin updates:", e) }
}
}
}
@@ -7,10 +7,13 @@ import com.intellij.openapi.application.UI
import com.intellij.openapi.application.asContextElement
import com.intellij.openapi.components.Service
import com.intellij.openapi.components.service
import com.intellij.openapi.diagnostic.getOrHandleException
import com.intellij.openapi.diagnostic.logger
import com.intellij.openapi.diagnostic.rethrowControlFlowException
import com.intellij.openapi.extensions.PluginId
import com.intellij.openapi.progress.runBlockingMaybeCancellable
import com.intellij.util.concurrency.annotations.RequiresEdt
import kotlinx.coroutines.CancellationException
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.FlowPreview
@@ -23,15 +26,16 @@ import kotlinx.coroutines.flow.catch
import kotlinx.coroutines.flow.debounce
import kotlinx.coroutines.flow.filterNotNull
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.flow.firstOrNull
import kotlinx.coroutines.launch
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext
import org.jetbrains.annotations.ApiStatus
import org.jetbrains.annotations.VisibleForTesting
import java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.CopyOnWriteArrayList
import java.util.function.Consumer
import kotlin.concurrent.Volatile
import kotlin.concurrent.atomics.ExperimentalAtomicApi
import kotlin.time.Duration.Companion.milliseconds
@@ -48,7 +52,8 @@ typealias PluginUpdateCallback = Consumer<PluginUpdatesEvent>
class PluginUpdatesService(val coroutineScope: CoroutineScope) {
private val myCallbacks = CopyOnWriteArrayList<PluginUpdateCallback>()
private val pluginUpdateFlow = MutableStateFlow<PluginUpdatesEvent?>(null)
private val pluginUpdateFlow = MutableSharedFlow<PluginUpdatesEvent>(replay = 1, onBufferOverflow = BufferOverflow.DROP_OLDEST)
@Volatile private var lastUpdates: PluginUpdatesEvent? = null
private val updateIdsFlow = MutableStateFlow<Set<PluginId>?>(null)
private val updateRequestFlow = MutableSharedFlow<Unit>(replay = 1, extraBufferCapacity = 1, onBufferOverflow = BufferOverflow.DROP_OLDEST)
private val providerSnapshots = ConcurrentHashMap<PluginUpdatesProvider, PluginUpdatesEvent>()
@@ -74,10 +79,11 @@ class PluginUpdatesService(val coroutineScope: CoroutineScope) {
init {
startUpdateCollection()
startUpdateTrigger()
startDispatchingCallbacks()
}
private suspend fun ensureUpdatesStarted() {
if (pluginUpdateFlow.value == null) {
if (lastUpdates == null) {
triggerUpdates()
}
}
@@ -98,22 +104,24 @@ class PluginUpdatesService(val coroutineScope: CoroutineScope) {
}
private suspend fun onProviderUpdated() {
val merged = updateMutex.withLock {
updateMutex.withLock {
val merged = mergeUpdates(providerSnapshots.values)
pluginUpdateFlow.value = merged
lastUpdates = merged
pluginUpdateFlow.emit(merged)
updateIdsFlow.value = merged.all.mapTo(HashSet()) { plugin -> plugin.pluginId }
merged
}
withContext(Dispatchers.UI + ModalityState.any().asContextElement()) {
dispatchCallbacks(merged)
}
}
private fun startDispatchingCallbacks() = coroutineScope.launch(Dispatchers.UI + ModalityState.any().asContextElement()) {
pluginUpdateFlow
.collect { dispatchCallbacks(it) }
}
private fun startUpdateTrigger() = coroutineScope.launch {
updateRequestFlow
.debounce(300.milliseconds)
.collect {
PluginUpdatesProvider.getInstances().forEach { it.update() }
updateRequestFlow.debounce(300.milliseconds).collect {
PluginUpdatesProvider.getInstances().forEach {
runCatching { it.update() }.getOrHandleException { e -> LOG.warn("PluginUpdatesProvider.update() failed:", e) }
}
}
}
@@ -149,7 +157,7 @@ class PluginUpdatesService(val coroutineScope: CoroutineScope) {
suspend fun awaitUpdates(): Collection<PluginUiModel> {
ensureUpdatesStarted()
return pluginUpdateFlow.filterNotNull().first().all
return pluginUpdateFlow.first().all
}
@VisibleForTesting
@@ -169,11 +177,14 @@ class PluginUpdatesService(val coroutineScope: CoroutineScope) {
}
private fun getLastUpdates(): PluginUpdatesEvent? {
return pluginUpdateFlow.value
return lastUpdates
}
private fun dispatchCallbacks(updates: PluginUpdatesEvent) {
myCallbacks.forEach { it.accept(updates) }
myCallbacks.forEach {
runCatching { it.accept(updates) }
.getOrHandleException { e -> LOG.warn("PluginUpdatesEvent callback ${it.javaClass.name} failed:", e) }
}
}
@VisibleForTesting