diff --git a/platform/xdebugger-impl/backend/src/com/intellij/platform/debugger/impl/backend/BackendXDebugSessionApi.kt b/platform/xdebugger-impl/backend/src/com/intellij/platform/debugger/impl/backend/BackendXDebugSessionApi.kt index c07b467f870d..38808c9f34b5 100644 --- a/platform/xdebugger-impl/backend/src/com/intellij/platform/debugger/impl/backend/BackendXDebugSessionApi.kt +++ b/platform/xdebugger-impl/backend/src/com/intellij/platform/debugger/impl/backend/BackendXDebugSessionApi.kt @@ -24,11 +24,10 @@ import fleet.kernel.rete.collect import fleet.kernel.rete.query import fleet.kernel.withEntities import fleet.rpc.core.toRpc -import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.* +import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.channels.awaitClose import kotlinx.coroutines.flow.* -import kotlinx.coroutines.launch -import kotlinx.coroutines.withContext internal class BackendXDebugSessionApi : XDebugSessionApi { override suspend fun currentEvaluator(sessionId: XDebugSessionId): Flow { @@ -157,8 +156,8 @@ internal class BackendXDebugSessionApi : XDebugSessionApi { override suspend fun currentExecutionStack(sessionId: XDebugSessionId): Flow { val sessionEntity = entity(XDebugSessionEntity.SessionId, sessionId) ?: return emptyFlow() val session = sessionEntity.session as? XDebugSessionImpl ?: return emptyFlow() - return withEntities(sessionEntity) { - channelFlow { + return channelFlow { + withEntities(sessionEntity) { session.getCurrentExecutionStackFlow().asEntityFlow { executionStack -> XExecutionStackEntity.new(this, executionStack) }.collect { stackAndId -> @@ -205,11 +204,23 @@ internal class BackendXDebugSessionApi : XDebugSessionApi { override suspend fun computeExecutionStacks(suspendContextId: XSuspendContextId): Flow { val suspendContextEntity = debuggerEntity(suspendContextId.id) as? XSuspendContextEntity ?: return emptyFlow() return channelFlow { + val channel = Channel?>(capacity = Channel.UNLIMITED) + + launch { + for (event in channel) { + if (event == null) { + channel.close() + this@channelFlow.close() + break + } + send(event.await()) + } + } + withEntities(suspendContextEntity) { suspendContextEntity.obj.computeExecutionStacks(object : XSuspendContext.XExecutionStackContainer { override fun addExecutionStack(executionStacks: List, last: Boolean) { - this@channelFlow.launch { - // TODO[IJPL-177087] delete entities!!! + channel.trySend(this@channelFlow.async { withEntities(suspendContextEntity) { val stackEntities = executionStacks.map { stack -> change { @@ -226,19 +237,18 @@ internal class BackendXDebugSessionApi : XDebugSessionApi { XExecutionStackDto(XExecutionStackId(stack.id), stack.obj.displayName, stack.obj.icon?.rpcId()) } } - send(XExecutionStacksEvent.NewExecutionStacks(stacks, last)) - if (last) { - this@channelFlow.close() - } + XExecutionStacksEvent.NewExecutionStacks(stacks, last) } + }) + if (last) { + channel.trySend(null) } } override fun errorOccurred(errorMessage: @NlsContexts.DialogMessage String) { - this@channelFlow.launch { - send(XExecutionStacksEvent.ErrorOccurred(errorMessage)) - this@channelFlow.close() - } + channel.trySend(this@channelFlow.async { + XExecutionStacksEvent.ErrorOccurred(errorMessage) + }) } }) } diff --git a/platform/xdebugger-impl/backend/src/com/intellij/platform/debugger/impl/backend/BackendXExecutionStackApi.kt b/platform/xdebugger-impl/backend/src/com/intellij/platform/debugger/impl/backend/BackendXExecutionStackApi.kt index 3cf5d10af835..2557a0eb1a9a 100644 --- a/platform/xdebugger-impl/backend/src/com/intellij/platform/debugger/impl/backend/BackendXExecutionStackApi.kt +++ b/platform/xdebugger-impl/backend/src/com/intellij/platform/debugger/impl/backend/BackendXExecutionStackApi.kt @@ -11,6 +11,9 @@ import com.intellij.xdebugger.impl.rhizome.XStackFrameEntity import com.intellij.xdebugger.impl.rpc.* import fleet.kernel.change import fleet.kernel.withEntities +import kotlinx.coroutines.Deferred +import kotlinx.coroutines.async +import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.channels.awaitClose import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.channelFlow @@ -31,11 +34,24 @@ internal class BackendXExecutionStackApi : XExecutionStackApi { override suspend fun computeStackFrames(executionStackId: XExecutionStackId, firstFrameIndex: Int): Flow { val executionStackEntity = debuggerEntity(executionStackId.id) ?: return emptyFlow() return channelFlow { + val channel = Channel?>(capacity = Channel.UNLIMITED) + + launch { + for (event in channel) { + if (event == null) { + channel.close() + this@channelFlow.close() + break + } + this@channelFlow.send(event.await()) + } + } + withEntities(executionStackEntity) { val executionStack = executionStackEntity.obj executionStack.computeStackFrames(firstFrameIndex, object : XExecutionStack.XStackFrameContainer { override fun addStackFrames(stackFrames: List, last: Boolean) { - this@channelFlow.launch { + channel.trySend(this@channelFlow.async { withEntities(executionStackEntity) { val frameEntities = stackFrames.map { frame -> change { @@ -50,19 +66,18 @@ internal class BackendXExecutionStackApi : XExecutionStackApi { val stacks = frameEntities.map { frame -> XStackFrameDto(XStackFrameId(frame.id), frame.obj.sourcePosition?.toRpc()) } - send(XStackFramesEvent.XNewStackFrames(stacks, last)) - if (last) { - this@channelFlow.close() - } + XStackFramesEvent.XNewStackFrames(stacks, last) } + }) + if (last) { + channel.trySend(null) } } override fun errorOccurred(errorMessage: @NlsContexts.DialogMessage String) { - this@channelFlow.launch { - send(XStackFramesEvent.ErrorOccurred(errorMessage)) - this@channelFlow.close() - } + channel.trySend(this@channelFlow.async { + XStackFramesEvent.ErrorOccurred(errorMessage) + }) } }) }