[fleet] more CoroutineNames

GitOrigin-RevId: d8d32749d432dd51c7c74e700e5b0d1835d67cba
This commit is contained in:
Alexander Shparun
2025-03-18 21:20:42 +00:00
committed by intellij-monorepo-bot
parent b3aa0cfb93
commit 82e5582669
5 changed files with 82 additions and 72 deletions
+11 -12
View File
@@ -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 <T> 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 <T> 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 <T> 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
@@ -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,
@@ -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<TransportMessage>,
receive: ReceiveChannel<TransportMessage>) {
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<TransportMessage>,
receive: ReceiveChannel<TransportMessage>,
) {
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))
}
}
}
}
+2 -2
View File
@@ -48,11 +48,11 @@ suspend fun <T> 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)
}
@@ -2,6 +2,8 @@
package fleet.util.async
import kotlinx.coroutines.*
import kotlin.coroutines.CoroutineContext
import kotlin.coroutines.EmptyCoroutineContext
suspend fun <T : Job, R> T.use(body: suspend CoroutineScope.(T) -> R): R {
return try {
@@ -35,11 +37,14 @@ suspend fun <T> withSupervisor(body: suspend CoroutineScope.(scope: CoroutineSco
}
}
suspend fun<T> withCoroutineScope(body: suspend CoroutineScope.(scope: CoroutineScope) -> T): T {
suspend fun <T> 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()