From 66843efda89004f37c0081c35af7f052e9cfdf8a Mon Sep 17 00:00:00 2001 From: Aleksey Pivovarov Date: Thu, 4 Dec 2025 15:10:01 +0100 Subject: [PATCH] [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 --- fleet/rpc/srcCommonMain/fleet/rpc/client/RpcClient.kt | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/fleet/rpc/srcCommonMain/fleet/rpc/client/RpcClient.kt b/fleet/rpc/srcCommonMain/fleet/rpc/client/RpcClient.kt index 934f52287b9b..a2409677ee88 100644 --- a/fleet/rpc/srcCommonMain/fleet/rpc/client/RpcClient.kt +++ b/fleet/rpc/srcCommonMain/fleet/rpc/client/RpcClient.kt @@ -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)) {