diff --git a/plugins/performanceTesting/event-bus/src/com/intellij/tools/ide/starter/bus/shared/SharedEventsFlow.kt b/plugins/performanceTesting/event-bus/src/com/intellij/tools/ide/starter/bus/shared/SharedEventsFlow.kt index a47521786130..c57c96320f14 100644 --- a/plugins/performanceTesting/event-bus/src/com/intellij/tools/ide/starter/bus/shared/SharedEventsFlow.kt +++ b/plugins/performanceTesting/event-bus/src/com/intellij/tools/ide/starter/bus/shared/SharedEventsFlow.kt @@ -1,6 +1,7 @@ package com.intellij.tools.ide.starter.bus.shared import com.fasterxml.jackson.module.kotlin.jacksonObjectMapper +import com.intellij.tools.ide.starter.bus.EventsBus import com.intellij.tools.ide.starter.bus.EventsFlow import com.intellij.tools.ide.starter.bus.events.Event import com.intellij.tools.ide.starter.bus.logger.EventBusLoggerFactory @@ -12,8 +13,10 @@ import kotlin.time.Duration private val LOG = EventBusLoggerFactory.getLogger(SharedEventsFlow::class.java) -class SharedEventsFlow(private val client: EventBusServerClient, - private val localEventsFlow: EventsFlow) : EventsFlow { +class SharedEventsFlow( + private val client: EventBusServerClient, + private val localEventsFlow: EventsFlow, +) : EventsFlow { private val objectMapper = jacksonObjectMapper() // Server returns all events. We need to save the count of processed events to avoid re-processing the event @@ -30,10 +33,12 @@ class SharedEventsFlow(private val client: EventBusServerClient, client.startServerProcess() } - override fun subscribe(eventClass: Class, - subscriber: Any, - timeout: Duration, - callback: suspend (event: EventType) -> Unit): Boolean { + override fun subscribe( + eventClass: Class, + subscriber: Any, + timeout: Duration, + callback: suspend (event: EventType) -> Unit, + ): Boolean { return localEventsFlow.subscribe(eventClass, subscriber, timeout, callback).also { if (it) client.newSubscriber(eventClass, timeout) } @@ -49,23 +54,25 @@ class SharedEventsFlow(private val client: EventBusServerClient, if (serverJob == null) { serverJob = CoroutineScope(Dispatchers.IO).launch { while (true) { - val allEvents = client.getEvents() - allEvents.entries.forEach { (eventName, events) -> - // Drop already processed events - events?.drop(processedEvents[eventName] ?: 0)?.forEach { - // Save the event as processed (here and not inside the launch) to avoid rerunning processing before the end of processing - processedEvents[eventName] = processedEvents.getOrDefault(eventName, 0) + 1 - // Processing events in not main flow to avoid blocking on nested events - launch { - try { - localEventsFlow.postAndWaitProcessing(it.second) - } - catch (e: Throwable) { - //can’t throw an exception into the main thread (even if CoroutineExceptionHandler used), so log it - LOG.info(e.stackTraceToString()) - } - finally { - client.processedEvent(it.first) + EventsBus.executeWithExceptionHandling(true) { + val allEvents = client.getEvents() + allEvents.entries.forEach { (eventName, events) -> + // Drop already processed events + events?.drop(processedEvents[eventName] ?: 0)?.forEach { + // Save the event as processed (here and not inside the launch) to avoid rerunning processing before the end of processing + processedEvents[eventName] = processedEvents.getOrDefault(eventName, 0) + 1 + // Processing events in not main flow to avoid blocking on nested events + launch { + try { + localEventsFlow.postAndWaitProcessing(it.second) + } + catch (e: Throwable) { + //can’t throw an exception into the main thread (even if CoroutineExceptionHandler used), so log it + LOG.info(e.stackTraceToString()) + } + finally { + client.processedEvent(it.first) + } } } }