From 82e5582669df34355cd7c86cd70015d262be0d02 Mon Sep 17 00:00:00 2001 From: Alexander Shparun Date: Tue, 18 Mar 2025 19:43:06 +0100 Subject: [PATCH] [fleet] more CoroutineNames GitOrigin-RevId: d8d32749d432dd51c7c74e700e5b0d1835d67cba --- fleet/kernel/src/fleet/kernel/Transactor.kt | 23 ++-- .../src/fleet/rpc/server/DirectRpcClient.kt | 3 +- .../rpc/server/ServerRequestDispatcher.kt | 115 +++++++++--------- fleet/rpc/src/fleet/rpc/client/RpcClient.kt | 4 +- .../core/src/fleet/util/async/WithLaunched.kt | 9 +- 5 files changed, 82 insertions(+), 72 deletions(-) diff --git a/fleet/kernel/src/fleet/kernel/Transactor.kt b/fleet/kernel/src/fleet/kernel/Transactor.kt index a10d6ee58792..d6479dfea767 100644 --- a/fleet/kernel/src/fleet/kernel/Transactor.kt +++ b/fleet/kernel/src/fleet/kernel/Transactor.kt @@ -87,7 +87,7 @@ interface Transactor : CoroutineContext.Element { internal suspend fun waitForDbSourceToCatchUpWithTimestamp(timestamp: Long) { val dbContext = DbContext.threadBound - if (dbContext.poison == null) { + if (dbContext.poison == null) { if (dbContext.impl.timestamp < timestamp) { val dbAfterTimestamp = currentCoroutineContext().dbSource.flow.first { db -> db.timestamp >= timestamp @@ -400,13 +400,13 @@ suspend fun withTransactor( val job = currentCoroutineContext().job job.ensureActive() val span = currentSpan.startChild( - SpanInfo( - name = "change", - job = job, - isScope = true, - startTimestampNano = null, - cause = null, - map = HashMap())) + SpanInfo( + name = "change", + job = job, + isScope = true, + startTimestampNano = null, + cause = null, + map = HashMap())) /** * DO NOT WRAP THIS BLOCK IN A SCOPE! * see `change suspend is atomic case 2` in [fleet.test.frontend.kernel.TransactorTest] @@ -471,8 +471,7 @@ suspend fun withTransactor( sharedFlow.emit(TransactorEvent.Init(timestamp = 0L, db = initialDb)) newSingleThreadCoroutineDispatcher("Kernel event loop thread ${kernelId}", DispatcherPriority.HIGH).use { coroutineDispatcher -> - launch(coroutineNameAppended("Changes processing job for $transactor") + coroutineDispatcher, - start = CoroutineStart.ATOMIC) { + launch(CoroutineName("Transactor loop $transactor") + coroutineDispatcher, start = CoroutineStart.ATOMIC) { spannedScope("kernel changes") { var ts = 1L consumeEach(priorityDispatchChannel, backgroundDispatchChannel) { changeTask -> @@ -537,7 +536,7 @@ suspend fun withTransactor( } }.use { try { - withContext(transactor + DbSource.ContextElement(FlowDbSource(transactor.dbState, debugName = "kernel $transactor")) + coroutineNameAppended("withKernel")) { + withContext(transactor + DbSource.ContextElement(FlowDbSource(transactor.dbState, debugName = "kernel $transactor"))) { body(transactor) } } @@ -555,7 +554,7 @@ private data class DbTimestamp(override val eid: EID) : Entity { } } -internal fun currentTimestamp(): Long = +internal fun currentTimestamp(): Long = DbTimestamp.single()[DbTimestamp.Timestamp] val Q.timestamp: Long diff --git a/fleet/rpc.server/src/fleet/rpc/server/DirectRpcClient.kt b/fleet/rpc.server/src/fleet/rpc/server/DirectRpcClient.kt index 72de2c2da4c9..73481692342c 100644 --- a/fleet/rpc.server/src/fleet/rpc/server/DirectRpcClient.kt +++ b/fleet/rpc.server/src/fleet/rpc/server/DirectRpcClient.kt @@ -10,6 +10,7 @@ import fleet.rpc.core.TransportMessage import fleet.util.UID import fleet.util.async.* import fleet.util.channels.channels +import kotlinx.coroutines.CoroutineName import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.launch @@ -36,7 +37,7 @@ fun RequestDispatcher.directRpcClient( cc(rpcClient) } } - }.span("directRpcClient") + }.span("directRpcClient").onContext(CoroutineName("directRpcClient")) suspend fun RequestDispatcher.withDirectRpcClient( interceptor: RpcInterceptor, diff --git a/fleet/rpc.server/src/fleet/rpc/server/ServerRequestDispatcher.kt b/fleet/rpc.server/src/fleet/rpc/server/ServerRequestDispatcher.kt index aa7fae1a3c0c..3d931e3b6eb0 100644 --- a/fleet/rpc.server/src/fleet/rpc/server/ServerRequestDispatcher.kt +++ b/fleet/rpc.server/src/fleet/rpc/server/ServerRequestDispatcher.kt @@ -35,70 +35,75 @@ class ServerRequestDispatcher(private val connectionListener: ConnectionListener return bannedEndpoints.asStateFlow() } - override suspend fun handleConnection(route: UID, - endpoint: EndpointKind, - presentableName: String?, - send: SendChannel, - receive: ReceiveChannel) { - val socketId = UID.random() - log.info { "handleConnection endpoint: $endpoint, route: $route, socket id: $socketId" } - receive.consume { - send.use { - bannedEndpoints.first { !it.contains(route) } - try { - val existing = connections.put(route, send) - if (existing != null) { - log.warn { "Replaced existing ${route}, will close previous socket" } - existing.close(RuntimeException("Replaced by other connection with same uid ${route}")) - } - log.info { "Notify $route is connected" } - broadcastSafely(TransportMessage.RouteOpened(route)) - connectionListener?.onConnect(endpoint, route, socketId, presentableName) - coroutineScope { - val connectionJob = launch { - receive.consumeEach { message -> - when (message) { - is TransportMessage.Envelope -> { - val destination = connections[message.destination] - if (destination != null) { - kotlin.runCatching { destination.send(message) } - .onFailure { ex -> - if (log.isTraceEnabled) { - log.trace(ex) { "Failed to send message from $route to ${message.destination}: $message" } - } else { - log.warn { "Failed to send message from $route to ${message.destination}" } + override suspend fun handleConnection( + route: UID, + endpoint: EndpointKind, + presentableName: String?, + send: SendChannel, + receive: ReceiveChannel, + ) { + withContext(CoroutineName("handleConnection $route")) { + val socketId = UID.random() + log.info { "handleConnection endpoint: $endpoint, route: $route, socket id: $socketId" } + receive.consume { + send.use { + bannedEndpoints.first { !it.contains(route) } + try { + val existing = connections.put(route, send) + if (existing != null) { + log.warn { "Replaced existing ${route}, will close previous socket" } + existing.close(RuntimeException("Replaced by other connection with same uid ${route}")) + } + log.info { "Notify $route is connected" } + broadcastSafely(TransportMessage.RouteOpened(route)) + connectionListener?.onConnect(endpoint, route, socketId, presentableName) + coroutineScope { + val connectionJob = launch { + receive.consumeEach { message -> + when (message) { + is TransportMessage.Envelope -> { + val destination = connections[message.destination] + if (destination != null) { + kotlin.runCatching { destination.send(message) } + .onFailure { ex -> + if (log.isTraceEnabled) { + log.trace(ex) { "Failed to send message from $route to ${message.destination}: $message" } + } + else { + log.warn { "Failed to send message from $route to ${message.destination}" } + } } - } + } + else { + val closed = TransportMessage.RouteClosed(message.destination) + log.trace { "Sending $closed to $route because route is not registered" } + send.send(closed) + } } - else { - val closed = TransportMessage.RouteClosed(message.destination) - log.trace { "Sending $closed to $route because route is not registered" } - send.send(closed) + else -> { + log.warn { "Good endpoints should send only TransportMessage.Envelope, but ${route} sends ${message}" } } } - else -> { - log.warn { "Good endpoints should send only TransportMessage.Envelope, but ${route} sends ${message}" } - } } } - } - val banned = async { bannedEndpoints.first { it.contains(route) } } - select { - connectionJob.onJoin { - banned.cancelAndJoin() - } - banned.onJoin { - connectionJob.cancelAndJoin() + val banned = async { bannedEndpoints.first { it.contains(route) } } + select { + connectionJob.onJoin { + banned.cancelAndJoin() + } + banned.onJoin { + connectionJob.cancelAndJoin() + } } } } - } - finally { - val removed = connections.remove(route, send) - connectionListener?.onDisconnect(endpoint, route, socketId) - if (removed) { - log.info { "Notify $route is disconnected" } - broadcastSafely(TransportMessage.RouteClosed(route)) + finally { + val removed = connections.remove(route, send) + connectionListener?.onDisconnect(endpoint, route, socketId) + if (removed) { + log.info { "Notify $route is disconnected" } + broadcastSafely(TransportMessage.RouteClosed(route)) + } } } } diff --git a/fleet/rpc/src/fleet/rpc/client/RpcClient.kt b/fleet/rpc/src/fleet/rpc/client/RpcClient.kt index eecbb56a208b..6dde2e9d282e 100644 --- a/fleet/rpc/src/fleet/rpc/client/RpcClient.kt +++ b/fleet/rpc/src/fleet/rpc/client/RpcClient.kt @@ -48,11 +48,11 @@ suspend fun rpcClient( ): T = newSingleThreadCoroutineDispatcher("rpc-client-$origin").use { dispatcher -> withSupervisor { supervisor -> - val client = RpcClient(coroutineScope = supervisor + supervisor.coroutineNameAppended("RpcClient"), + val client = RpcClient(coroutineScope = supervisor + CoroutineName("RpcScope"), transport = transport, origin = origin, requestInterceptor = requestInterceptor) - launch(start = CoroutineStart.ATOMIC, context = dispatcher) { client.work(abortOnError) } + launch(start = CoroutineStart.ATOMIC, context = dispatcher + CoroutineName("RpcClient")) { client.work(abortOnError) } .use { body(client) } diff --git a/fleet/util/core/src/fleet/util/async/WithLaunched.kt b/fleet/util/core/src/fleet/util/async/WithLaunched.kt index 6c04dacfb606..a48acace5949 100644 --- a/fleet/util/core/src/fleet/util/async/WithLaunched.kt +++ b/fleet/util/core/src/fleet/util/async/WithLaunched.kt @@ -2,6 +2,8 @@ package fleet.util.async import kotlinx.coroutines.* +import kotlin.coroutines.CoroutineContext +import kotlin.coroutines.EmptyCoroutineContext suspend fun T.use(body: suspend CoroutineScope.(T) -> R): R { return try { @@ -35,11 +37,14 @@ suspend fun withSupervisor(body: suspend CoroutineScope.(scope: CoroutineSco } } -suspend fun withCoroutineScope(body: suspend CoroutineScope.(scope: CoroutineScope) -> T): T { +suspend fun withCoroutineScope( + coroutineContext: CoroutineContext = EmptyCoroutineContext, + body: suspend CoroutineScope.(scope: CoroutineScope) -> T, +): T { val context = currentCoroutineContext() val job = Job(context.job) return try { - coroutineScope { body(CoroutineScope(context + job)) } + coroutineScope { body(CoroutineScope(context + coroutineContext + job)) } } finally { job.cancelAndJoin()