[rd][debugger] IJPL-177087 Trying to fix project leak 1

GitOrigin-RevId: 7f52f7f466f35a8574fa6ee0e657ed26c830ad0b
This commit is contained in:
Alexander Kuznetsov
2025-03-25 02:02:28 +00:00
committed by intellij-monorepo-bot
parent 3db6c38939
commit 22f6d6723d
2 changed files with 49 additions and 24 deletions
@@ -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<XDebuggerEvaluatorDto?> {
@@ -157,8 +156,8 @@ internal class BackendXDebugSessionApi : XDebugSessionApi {
override suspend fun currentExecutionStack(sessionId: XDebugSessionId): Flow<XExecutionStackDto?> {
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<XExecutionStacksEvent> {
val suspendContextEntity = debuggerEntity<XSuspendContextEntity>(suspendContextId.id) as? XSuspendContextEntity ?: return emptyFlow()
return channelFlow {
val channel = Channel<Deferred<XExecutionStacksEvent>?>(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<XExecutionStack>, 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)
})
}
})
}
@@ -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<XStackFramesEvent> {
val executionStackEntity = debuggerEntity<XExecutionStackEntity>(executionStackId.id) ?: return emptyFlow()
return channelFlow {
val channel = Channel<Deferred<XStackFramesEvent>?>(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<XStackFrame>, 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)
})
}
})
}