[rdct] IJPL-221336 fix 'fleetClient' occasionally not trying to reconnect

`fleet.rpc.core.DeferredSerializer` is using `rpcCoroutineScope` to run a long-running task.

Prior to the fix, the 'rpcScope' was set to be one from 'consumeAll' in the 'receiver'.
This blocked receiver coroutine from ever finishing and closing the 'requestsChannel'.
As a result, the 'receiver.await()' never finished either and the 'RpcClientDisconnectedException' was not rethrown,
triggering the reconnection cycle.

Fix regression after 338dcda101594c3a735b877cd345a6aabf1319bd where definition of 'this: CoroutineScope' was silently altered.

GitOrigin-RevId: 90ac428813735f2f026ce1e67eb988ed9b8e4f4e
This commit is contained in:
Aleksey Pivovarov
2025-12-04 19:10:26 +00:00
committed by intellij-monorepo-bot
parent d548e118d4
commit 66843efda8
@@ -146,6 +146,8 @@ private class RpcClient(
@OptIn(ExperimentalCoroutinesApi::class)
internal suspend fun work(abortOnError: Boolean) {
supervisorScope {
val rpcScope = this
val receiver = async(start = CoroutineStart.ATOMIC) {
consumeAll(transport.incoming, eventLoopChannel) {
val mergedIncomingAndTransport = flow {
@@ -165,7 +167,7 @@ private class RpcClient(
logger.trace { "Received ${event.message}" }
when (val message = event.message) {
is TransportMessage.Envelope -> {
acceptMessage(message.parseMessage(), message.origin)
acceptMessage(message.parseMessage(), message.origin, rpcScope)
}
is TransportMessage.RouteClosed -> {
grayList.putIfAbsent(message.address, CompletableDeferred())
@@ -293,7 +295,7 @@ private class RpcClient(
}
}
private fun CoroutineScope.acceptMessage(message: RpcMessage, senderRoute: UID) {
private fun acceptMessage(message: RpcMessage, senderRoute: UID, rpcScope: CoroutineScope) {
when (message) {
is RpcMessage.CallResult -> {
logger.trace { "Got CallResult: requestId = ${message.requestId}" }
@@ -324,7 +326,7 @@ private class RpcClient(
}
return@run resource to emptyList()
}
val (de, streamDescriptors) = withSerializationContext(rpc.call.displayName, rpc.token, this) {
val (de, streamDescriptors) = withSerializationContext(rpc.call.displayName, rpc.token, rpcScope) {
val kser = rpc.returnType.serializer(rpc.call.classMethodDisplayName())
val json = rpcJsonImplementationDetail()
json.decodeFromJsonElement(kser, message.result)
@@ -362,7 +364,7 @@ private class RpcClient(
if (stream != null) {
when (stream) {
is InternalStreamDescriptor.FromRemote -> {
val (element, streamDescriptors) = withSerializationContext(stream.displayName, stream.token, this) {
val (element, streamDescriptors) = withSerializationContext(stream.displayName, stream.token, rpcScope) {
rpcJsonImplementationDetail().decodeFromJsonElement(stream.elementSerializer, message.data)
}
for (internalDescriptor in registerStreams(streamDescriptors, stream.route, stream.prefetchStrategy)) {