[terminal] IJPL-234877 Try to fix stuck terminal after reconnect

There can be two issues:
1. We collect the flow on frontend out of `durable` block - if there is a connection failure, our logic won't start the collection again. Solved by wrapping both flow obtaining and collection into `durable` call. It is a widely used pattern.
2. If there is a connection problem, our backend will continue to emit events to the output flow. But will be blocked when RPC flow buffer fills up. It will lead to blocking reading from PTY and if it is blocked for a long time (10+ sec) it may break some applications, for example `btop`. Solved by adding a logic of flow termination if we fail to emit the event for 3 sec. The client will need to request the flow again and receive the state snapshot.

(cherry picked from commit 9838d2fc580b1574d0aa61fd5a01f8a25f9c6038)

IJ-CR-194211

GitOrigin-RevId: 1ee6253922fae7efa1de4fa1d20d506266d80b57
This commit is contained in:
Konstantin Hudyakov
2026-03-19 08:42:55 +00:00
committed by intellij-monorepo-bot
parent 7c88ec01aa
commit cc6351b890
3 changed files with 41 additions and 6 deletions
@@ -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<List<TerminalOutputEvent>> = outputFlowProducer.getIncrementalUpdateFlow()
override suspend fun getOutputFlow(): Flow<List<TerminalOutputEvent>> {
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<StateAwareTerminalSession>()
}
private inner class State : MutableStateWithIncrementalUpdates<List<TerminalOutputEvent>> {
override suspend fun applyUpdate(update: List<TerminalOutputEvent>): List<TerminalOutputEvent> {
for (event in update) {
@@ -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)
}
}
}
}
}
@@ -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<List<TerminalOutputEvent>>