From f19e78eb7371828f31743daceb4e034082669aa8 Mon Sep 17 00:00:00 2001 From: Bartek Pacia Date: Wed, 3 Dec 2025 22:34:58 +0100 Subject: [PATCH] refactor [lsp]: add descriptive `CoroutineName`s to launched coroutines GitOrigin-RevId: c9aa7d304bca249bd503704b99b947cbf1829147 --- fleet/kernel/srcCommonMain/fleet/kernel/Transactor.kt | 6 ++++-- fleet/kernel/srcCommonMain/fleet/kernel/rete/Rete.kt | 2 +- .../com/jetbrains/lsp/implementation/protocolFraming.kt | 8 +++++--- .../com/jetbrains/lsp/implementation/withLsp.kt | 4 ++-- 4 files changed, 12 insertions(+), 8 deletions(-) diff --git a/fleet/kernel/srcCommonMain/fleet/kernel/Transactor.kt b/fleet/kernel/srcCommonMain/fleet/kernel/Transactor.kt index dfd16f3f2516..079358cc15b3 100644 --- a/fleet/kernel/srcCommonMain/fleet/kernel/Transactor.kt +++ b/fleet/kernel/srcCommonMain/fleet/kernel/Transactor.kt @@ -133,8 +133,10 @@ suspend fun Transactor.subscribe(capacity: Int = Channel.RENDEZVOUS, body: S val (send, receive) = channels(capacity) // trick: use channel in place of deferred, cause the latter one would hold the firstDB for the lifetime of the entire subscription val firstDB = Channel(1) - val job = launch(start = CoroutineStart.UNDISPATCHED, - context = Dispatchers.Unconfined) { + val job = launch( + start = CoroutineStart.UNDISPATCHED, + context = CoroutineName("transactor log collector") + Dispatchers.Unconfined, + ) { log.collect { e -> when (e) { is SubscriptionEvent.First -> { diff --git a/fleet/kernel/srcCommonMain/fleet/kernel/rete/Rete.kt b/fleet/kernel/srcCommonMain/fleet/kernel/rete/Rete.kt index 92499e5a75b5..a3f274848ae1 100644 --- a/fleet/kernel/srcCommonMain/fleet/kernel/rete/Rete.kt +++ b/fleet/kernel/srcCommonMain/fleet/kernel/rete/Rete.kt @@ -98,7 +98,7 @@ suspend fun withRete( kernel.subscribe(Channel.UNLIMITED) { db, changes -> val lastKnownDb = MutableStateFlow(ReteState.Db(db)) coroutineScope { - launch { + launch(CoroutineName("rete event loop")) { spannedScope("rete event loop") { // todo: implement a proper reconnect, this could still fail because of thread starvation changes.consumeAsFlow() diff --git a/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/implementation/protocolFraming.kt b/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/implementation/protocolFraming.kt index 2d5098259c2b..59c596c8b0f8 100644 --- a/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/implementation/protocolFraming.kt +++ b/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/implementation/protocolFraming.kt @@ -4,6 +4,7 @@ import com.jetbrains.lsp.protocol.LSP import fleet.util.decodeToStringUtf8 import fleet.util.encodeToByteArrayUtf8 import kotlinx.coroutines.CompletableDeferred +import kotlinx.coroutines.CoroutineName import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.channels.ReceiveChannel @@ -28,7 +29,7 @@ suspend fun withBaseProtocolFraming( coroutineScope { val (incomingSender, incomingReceiver) = channels() val (outgoingSender, outgoingReceiver) = channels(Channel.UNLIMITED) - val readJob = launch { + val readJob = launch(CoroutineName("frame reader")) { incomingSender.use { while (true) { val frame = reader.readFrame() @@ -40,7 +41,7 @@ suspend fun withBaseProtocolFraming( } } } - val writeJob = launch { + val writeJob = launch(CoroutineName("frame writer")) { outgoingReceiver.consumeEach { frame -> val success = writer.writeFrame(frame) if (!success) { @@ -76,7 +77,8 @@ private suspend fun ByteReader.readFrame(): JsonElement? { if (!readSomething) return null if (contentLength == -1) throw IllegalStateException("Content-Length header not found") readByteArray(contentLength) - } catch (e: Exception) { + } + catch (e: Exception) { when (e) { is IOException -> return null else -> throw e diff --git a/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/implementation/withLsp.kt b/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/implementation/withLsp.kt index c70a6d43f387..a739c1825f8d 100644 --- a/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/implementation/withLsp.kt +++ b/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/implementation/withLsp.kt @@ -102,7 +102,7 @@ suspend fun withLsp( val lspHandlerContext = LspHandlerContext(lspClient) - launch(createCoroutineContext(lspClient)) { + launch(CoroutineName("incoming requests accepter") + createCoroutineContext(lspClient)) { withSupervisor { supervisor -> val incomingRequestsJobs = MultiplatformConcurrentHashMap() incoming.consumeEach { jsonMessage -> @@ -113,7 +113,7 @@ suspend fun withLsp( isRequest(jsonMessage) -> { val request = LSP.json.decodeFromJsonElement(RequestMessage.serializer(), jsonMessage) - supervisor.launch(start = CoroutineStart.ATOMIC) { + supervisor.launch(context = CoroutineName("handler for ${request.method}"), start = CoroutineStart.ATOMIC) { val maybeHandler = handlers.requestHandler(request.method) ?.let { handler -> middleware.requestHandler(handler) } runCatching {