diff --git a/plugins/terminal/backend/src/com/intellij/terminal/backend/StateAwareTerminalSession.kt b/plugins/terminal/backend/src/com/intellij/terminal/backend/StateAwareTerminalSession.kt index 74024696fd28..01c1817de770 100644 --- a/plugins/terminal/backend/src/com/intellij/terminal/backend/StateAwareTerminalSession.kt +++ b/plugins/terminal/backend/src/com/intellij/terminal/backend/StateAwareTerminalSession.kt @@ -1,5 +1,6 @@ package com.intellij.terminal.backend +import com.intellij.openapi.diagnostic.logger import com.intellij.openapi.diagnostic.thisLogger import com.intellij.openapi.editor.impl.DocumentImpl import com.intellij.openapi.project.Project @@ -13,13 +14,17 @@ import com.intellij.util.asDisposable import kotlinx.coroutines.CoroutineName import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.ExperimentalCoroutinesApi +import kotlinx.coroutines.TimeoutCancellationException import kotlinx.coroutines.cancel import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.channels.SendChannel import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.buffer +import kotlinx.coroutines.flow.channelFlow import kotlinx.coroutines.flow.map import kotlinx.coroutines.flow.merge import kotlinx.coroutines.launch +import kotlinx.coroutines.withTimeout import org.jetbrains.plugins.terminal.block.reworked.TerminalSessionModel import org.jetbrains.plugins.terminal.block.reworked.TerminalSessionModelImpl import org.jetbrains.plugins.terminal.block.ui.TerminalUiUtils @@ -53,6 +58,7 @@ import org.jetbrains.plugins.terminal.view.impl.MutableTerminalOutputModelImpl import org.jetbrains.plugins.terminal.view.impl.updateContent import org.jetbrains.plugins.terminal.view.shellIntegration.impl.TerminalBlocksModelImpl import kotlin.coroutines.cancellation.CancellationException +import kotlin.time.Duration.Companion.seconds import kotlin.time.TimeSource /** @@ -172,13 +178,32 @@ internal class StateAwareTerminalSession( fun getHyperlinkFacade(isInAlternateBuffer: Boolean): BackendTerminalHyperlinkFacade? = if (isInAlternateBuffer) alternateBufferHyperlinkFacade else outputHyperlinkFacade - override suspend fun getOutputFlow(): Flow> = outputFlowProducer.getIncrementalUpdateFlow() + override suspend fun getOutputFlow(): Flow> { + return channelFlow { + try { + outputFlowProducer.getIncrementalUpdateFlow().collect { events -> + withTimeout(3.seconds) { + send(events) + } + } + } + catch (_: TimeoutCancellationException) { + // Downstream consumer is too slow, ending the flow to unblock the original session output flow processing. + // The collector should request a new flow and receive a state snapshot + LOG.info("Failed to emit output to the collector in 3 seconds, is there a connection problem? Terminating the output flow.") + } + }.buffer(Channel.RENDEZVOUS) + } override val isClosed: Boolean get() = delegate.isClosed override suspend fun hasRunningCommands(): Boolean = delegate.hasRunningCommands() + companion object { + private val LOG = logger() + } + private inner class State : MutableStateWithIncrementalUpdates> { override suspend fun applyUpdate(update: List): List { for (event in update) { diff --git a/plugins/terminal/frontend/src/com/intellij/terminal/frontend/view/impl/TerminalSessionController.kt b/plugins/terminal/frontend/src/com/intellij/terminal/frontend/view/impl/TerminalSessionController.kt index 6bf057146fd5..b231eb8817d4 100644 --- a/plugins/terminal/frontend/src/com/intellij/terminal/frontend/view/impl/TerminalSessionController.kt +++ b/plugins/terminal/frontend/src/com/intellij/terminal/frontend/view/impl/TerminalSessionController.kt @@ -10,7 +10,9 @@ import com.intellij.openapi.diagnostic.thisLogger import com.intellij.terminal.JBTerminalSystemSettingsProviderBase import com.intellij.terminal.frontend.view.hyperlinks.FrontendTerminalHyperlinkFacade import com.intellij.util.containers.DisposableWrapperList +import fleet.rpc.client.durable import kotlinx.coroutines.CompletableDeferred +import kotlinx.coroutines.CoroutineName import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.Runnable @@ -50,11 +52,17 @@ internal class TerminalSessionController( private val edtContext = Dispatchers.EDT + ModalityState.any().asContextElement() fun handleEvents(session: TerminalSession) { - coroutineScope.launch { - val outputFlow = session.getOutputFlow() - withContext(edtContext) { - outputFlow.collect { events -> - doHandleEvents(events) + coroutineScope.launch(Dispatchers.IO + CoroutineName("Output flow collection")) { + // Get output flow again even if it was terminated. + // It can happen in case of RemDev if there were any connection problems and backend decided to terminate the flow. + while (!session.isClosed) { + // Wrap the flow collection into `durable` call to retry in case of connection issues. + durable { + session.getOutputFlow().collect { events -> + withContext(edtContext) { + doHandleEvents(events) + } + } } } } diff --git a/plugins/terminal/src/org/jetbrains/plugins/terminal/session/impl/TerminalSession.kt b/plugins/terminal/src/org/jetbrains/plugins/terminal/session/impl/TerminalSession.kt index c65c9056787f..5cbe3c819444 100644 --- a/plugins/terminal/src/org/jetbrains/plugins/terminal/session/impl/TerminalSession.kt +++ b/plugins/terminal/src/org/jetbrains/plugins/terminal/session/impl/TerminalSession.kt @@ -20,6 +20,8 @@ interface TerminalSession { * Use this flow to handle the output events of the Terminal session. * * Underlying logic should continue reading the PTYs output stream only if there is some collector of this flow. + * If the flow collector is too slow (can't handle event in 3 seconds), the flow can be terminated, and + * you need to request a new flow and receive a state snapshot. */ suspend fun getOutputFlow(): Flow>