mirror of
https://gitflic.ru/project/openide/openide.git
synced 2026-09-27 10:03:11 +07:00
refactor [lsp]: add descriptive CoroutineNames to launched coroutines
GitOrigin-RevId: c9aa7d304bca249bd503704b99b947cbf1829147
This commit is contained in:
committed by
intellij-monorepo-bot
parent
5c4989885a
commit
f19e78eb73
@@ -133,8 +133,10 @@ suspend fun <T> Transactor.subscribe(capacity: Int = Channel.RENDEZVOUS, body: S
|
||||
val (send, receive) = channels<Change>(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<DB>(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 -> {
|
||||
|
||||
@@ -98,7 +98,7 @@ suspend fun <T> withRete(
|
||||
kernel.subscribe(Channel.UNLIMITED) { db, changes ->
|
||||
val lastKnownDb = MutableStateFlow<ReteState>(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()
|
||||
|
||||
+5
-3
@@ -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<JsonElement>()
|
||||
val (outgoingSender, outgoingReceiver) = channels<JsonElement>(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
|
||||
|
||||
@@ -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<StringOrInt, Job>()
|
||||
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 {
|
||||
|
||||
Reference in New Issue
Block a user