IJPL-238173 [rpc] Mark public and internal API in the fleet.rpc module to expose the necessary components for external plugin developers

Absolute minimum of APIs is exposed so far, namely:
- RemoteApiDescriptor
- RpcException
- RPC annotation
and their descendants


(cherry picked from commit 95b32b92665878ba2a35baad999458a5653eab05)

IJ-CR-195627

GitOrigin-RevId: 7e94238626f3ab69743de25653bc64a54857fa85
This commit is contained in:
Nikita Katkov
2026-03-11 17:30:54 +00:00
committed by intellij-monorepo-bot
parent 88b942cd9f
commit ab0fc99020
11 changed files with 100 additions and 3 deletions
+60
View File
@@ -0,0 +1,60 @@
fleet.rpc.RemoteApi
fleet.rpc.RemoteApiDescriptor
- a:call(fleet.rpc.RemoteApi,java.lang.String,java.lang.Object[],kotlin.coroutines.Continuation):java.lang.Object
- a:clientStub(kotlin.jvm.functions.Function3):fleet.rpc.RemoteApi
- a:getApiFqn():java.lang.String
- a:getSignature(java.lang.String):fleet.rpc.RpcSignature
@:fleet.rpc.Rpc
- java.lang.annotation.Annotation
f:fleet.rpc.client.DurableKt
- sf:durable(Z,kotlin.jvm.functions.Function2,kotlin.coroutines.Continuation):java.lang.Object
- bs:durable$default(Z,kotlin.jvm.functions.Function2,kotlin.coroutines.Continuation,I,java.lang.Object):java.lang.Object
f:fleet.rpc.client.FleetClientKt
- sf:fleetClient(fleet.rpc.client.ClientId,fleet.rpc.core.TransportFactory,Z,fleet.util.async.DelayStrategy,fleet.rpc.client.RpcInterceptor,java.lang.String):fleet.util.async.Resource
- bs:fleetClient$default(fleet.rpc.client.ClientId,fleet.rpc.core.TransportFactory,Z,fleet.util.async.DelayStrategy,fleet.rpc.client.RpcInterceptor,java.lang.String,I,java.lang.Object):fleet.util.async.Resource
- sf:proxy(fleet.rpc.client.FleetClient,fleet.rpc.RemoteApiDescriptor,fleet.util.UID,fleet.rpc.core.InstanceId):fleet.rpc.RemoteApi
f:fleet.rpc.client.RemoteIsCancelledException
- fleet.rpc.client.RpcClientException
- kotlinx.coroutines.CopyableThrowable
- <init>(java.lang.String,java.lang.Throwable):V
- createCopy():fleet.rpc.client.RemoteIsCancelledException
f:fleet.rpc.client.RouteClosedException
- fleet.rpc.client.RpcClientException
- kotlinx.coroutines.CopyableThrowable
- <init>(fleet.util.UID,java.lang.String,java.lang.Throwable):V
- b:<init>(fleet.util.UID,java.lang.String,java.lang.Throwable,I,kotlin.jvm.internal.DefaultConstructorMarker):V
- createCopy():fleet.rpc.client.RouteClosedException
- f:getRoute():fleet.util.UID
f:fleet.rpc.client.RpcCausalityTimeout
- fleet.rpc.client.RpcClientException
- kotlinx.coroutines.CopyableThrowable
- <init>(java.lang.String,java.lang.Throwable):V
- createCopy():fleet.rpc.client.RpcCausalityTimeout
f:fleet.rpc.client.RpcClientDisconnectedException
- fleet.rpc.client.RpcClientException
- kotlinx.coroutines.CopyableThrowable
- <init>(java.lang.String,java.lang.Throwable):V
- createCopy():fleet.rpc.client.RpcClientDisconnectedException
a:fleet.rpc.client.RpcClientException
- java.lang.RuntimeException
- <init>(java.lang.String,java.lang.Throwable):V
f:fleet.rpc.client.RpcServiceNotReady
- fleet.rpc.client.RpcClientException
- kotlinx.coroutines.CopyableThrowable
- <init>(fleet.rpc.core.RpcMessage$CallRequest,java.lang.Throwable):V
- b:<init>(fleet.rpc.core.RpcMessage$CallRequest,java.lang.Throwable,I,kotlin.jvm.internal.DefaultConstructorMarker):V
- createCopy():fleet.rpc.client.RpcServiceNotReady
f:fleet.rpc.client.RpcTimeoutException
- fleet.rpc.client.RpcClientException
- kotlinx.coroutines.CopyableThrowable
- <init>(java.lang.String,java.lang.Throwable):V
- b:<init>(java.lang.String,java.lang.Throwable,I,kotlin.jvm.internal.DefaultConstructorMarker):V
- createCopy():fleet.rpc.client.RpcTimeoutException
- f:getMsg():java.lang.String
f:fleet.rpc.client.UnresolvedServiceException
- fleet.rpc.client.RpcClientException
- kotlinx.coroutines.CopyableThrowable
- <init>(fleet.rpc.core.InstanceId,java.lang.Throwable):V
- b:<init>(fleet.rpc.core.InstanceId,java.lang.Throwable,I,kotlin.jvm.internal.DefaultConstructorMarker):V
- createCopy():fleet.rpc.client.UnresolvedServiceException
- f:getServiceId():fleet.rpc.core.InstanceId
@@ -1,5 +1,8 @@
package fleet.rpc
import org.jetbrains.annotations.ApiStatus
@ApiStatus.Internal
enum class EndpointKind {
Client,
Provider
@@ -9,6 +9,7 @@ import fleet.util.cast
import fleet.util.letIf
import kotlinx.serialization.KSerializer
import kotlinx.serialization.builtins.nullable
import org.jetbrains.annotations.ApiStatus
/**
* Base interface that must be implemented by every Fleet service.
@@ -26,6 +27,7 @@ interface RemoteApi<Metadata>
@Target(AnnotationTarget.CLASS)
annotation class Rpc
@ApiStatus.Internal
sealed interface RemoteKind {
data class Data(val serializer: KSerializer<*>) : RemoteKind
data class Flow(val elementKind: RemoteKind, val nullable: Boolean) : RemoteKind
@@ -36,6 +38,7 @@ sealed interface RemoteKind {
data class Resource(val descriptor: RemoteApiDescriptor<*>) : RemoteKind
}
@ApiStatus.Internal
fun RemoteKind.serializer(debugInfo: String): KSerializer<Any?> {
return when (this) {
is RemoteKind.Data -> serializer
@@ -47,8 +50,9 @@ fun RemoteKind.serializer(debugInfo: String): KSerializer<Any?> {
is RemoteKind.Resource -> error("Resource has no serializer")
}.cast()
}
@ApiStatus.Internal
data class ParameterDescriptor(val parameterName: String, val parameterKind: RemoteKind)
@ApiStatus.Internal
data class RpcSignature(val methodName: String, val parameters: Array<ParameterDescriptor>, val returnType: RemoteKind)
interface RemoteApiDescriptor<T : RemoteApi<*>> {
@@ -4,7 +4,9 @@ package fleet.rpc.client
import fleet.util.UID
import fleet.util.serialization.DataSerializer
import kotlinx.serialization.Serializable
import org.jetbrains.annotations.ApiStatus
@ApiStatus.Internal
@Serializable(with = ClientIdSerializer::class)
data class ClientId(val uid: UID)
@@ -25,10 +25,11 @@ import kotlinx.coroutines.CancellationException
import kotlinx.coroutines.CoroutineName
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import org.jetbrains.annotations.ApiStatus.Internal
import org.jetbrains.annotations.ApiStatus
import kotlin.concurrent.Volatile
import kotlin.coroutines.CoroutineContext
@ApiStatus.Internal
class FleetClient internal constructor(
val connectionStatus: StateFlow<ConnectionStatus<IRpcClient>>,
val stats: MutableStateFlow<TransportStats>,
@@ -41,7 +42,7 @@ class FleetClient internal constructor(
@Volatile
private var poison: CancellationException? = null
@Internal
@ApiStatus.Internal
val invocationHandlerFactory: InvocationHandlerFactory<ProxyClosure> =
reconnectingRpcClient(connectionStatus)
.asHandlerFactory()
@@ -17,10 +17,12 @@ import kotlinx.coroutines.flow.first
import kotlinx.coroutines.flow.mapNotNull
import kotlinx.coroutines.withContext
import kotlinx.coroutines.withTimeoutOrNull
import org.jetbrains.annotations.ApiStatus
import kotlin.concurrent.Volatile
import kotlin.coroutines.CoroutineContext
import kotlin.coroutines.coroutineContext
@ApiStatus.Internal
data class Call(val route: UID,
val service: InstanceId,
val signature: RpcSignature,
@@ -35,6 +37,7 @@ internal data class RpcStrategyContextElement(val awaitConnection: Boolean = tru
override val key: CoroutineContext.Key<*> get() = RpcStrategyContextElement
}
@ApiStatus.Internal
suspend fun <T> withoutAwaitingForReconnect(body: suspend CoroutineScope.() -> T): T {
val strategy = coroutineContext[RpcStrategyContextElement]?.copy(awaitConnection = false)
?: RpcStrategyContextElement(awaitConnection = false)
@@ -43,6 +46,7 @@ suspend fun <T> withoutAwaitingForReconnect(body: suspend CoroutineScope.() -> T
}
}
@ApiStatus.Internal
suspend fun <T> withPrefetchStrategy(prefetchStrategy: PrefetchStrategy, body: suspend CoroutineScope.() -> T): T {
val strategy = coroutineContext[RpcStrategyContextElement]?.copy(prefetchStrategy = prefetchStrategy)
?: RpcStrategyContextElement(prefetchStrategy = prefetchStrategy)
@@ -51,10 +55,12 @@ suspend fun <T> withPrefetchStrategy(prefetchStrategy: PrefetchStrategy, body: s
}
}
@ApiStatus.Internal
interface IRpcClient {
suspend fun call(call: Call, publish: (SuspendInvocationHandler.CallResult) -> Unit)
}
@ApiStatus.Internal
fun promisingRpcClient(promise: Deferred<IRpcClient>): IRpcClient {
return object : IRpcClient {
override suspend fun call(call: Call, publish: (SuspendInvocationHandler.CallResult) -> Unit) {
@@ -63,6 +69,7 @@ fun promisingRpcClient(promise: Deferred<IRpcClient>): IRpcClient {
}
}
@ApiStatus.Internal
fun IRpcClient.asHandlerFactory(): InvocationHandlerFactory<ProxyClosure> =
object : InvocationHandlerFactory<ProxyClosure> {
override fun handler(arg: ProxyClosure): SuspendInvocationHandler {
@@ -71,6 +71,7 @@ import kotlinx.coroutines.suspendCancellableCoroutine
import kotlinx.coroutines.withTimeoutOrNull
import kotlinx.serialization.builtins.serializer
import kotlinx.serialization.json.Json
import org.jetbrains.annotations.ApiStatus
import kotlin.coroutines.Continuation
import kotlin.coroutines.coroutineContext
import kotlin.coroutines.resumeWithException
@@ -90,6 +91,7 @@ private data class OutgoingRequest(
private data class OngoingRequest(val request: OutgoingRequest)
@ApiStatus.Internal
fun rpcClient(
transport: Transport,
origin: UID,
@@ -2,7 +2,9 @@
package fleet.rpc.client
import fleet.rpc.core.RpcMessage
import org.jetbrains.annotations.ApiStatus
@ApiStatus.Internal
interface RpcInterceptor {
companion object: RpcInterceptor {
override suspend fun interceptCallRequest(request: RpcMessage.CallRequest): RpcMessage.CallRequest {
@@ -18,6 +20,7 @@ interface RpcInterceptor {
suspend fun interceptCallResult(displayName: String, result: RpcMessage.CallResult) {}
}
@ApiStatus.Internal
operator fun RpcInterceptor.plus(another: RpcInterceptor): RpcInterceptor {
val one = this
return object : RpcInterceptor {
@@ -0,0 +1,5 @@
// Copyright 2000-2026 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license.
@Internal
package fleet.rpc.client.proxy;
import org.jetbrains.annotations.ApiStatus.Internal;
@@ -0,0 +1,5 @@
// Copyright 2000-2026 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license.
@Internal
package fleet.rpc.core;
import org.jetbrains.annotations.ApiStatus.Internal;
@@ -0,0 +1,5 @@
// Copyright 2000-2026 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license.
@Internal
package fleet.rpc.core.util;
import org.jetbrains.annotations.ApiStatus.Internal;