[fleet, rpc] restore a guarantee that continuation will never be resumed twice

GitOrigin-RevId: 6c72b3758a3a5f79cef0e1f2dfd89f7db7242ce6
This commit is contained in:
Alexander Zolotov
2025-09-02 00:05:34 +00:00
committed by intellij-monorepo-bot
parent d627bca71c
commit fd3c8985dd
@@ -412,11 +412,12 @@ private class RpcClient(
private fun resumeAllOngoingCallsWithThrowable(throwable: Throwable) {
logger.debug(throwable) { "resumeAllOngoingCallsWithThrowable" }
val outgoingRpcIterator = outgoingRpc.iterator()
for ((key, value) in outgoingRpcIterator) {
outgoingRpcIterator.remove()
logger.trace { "resumeAllOngoingCallsWithThrowable: resume request $key" }
value.request.continuation.resumeWithException(throwable)
val outgoingKeys = outgoingRpc.keys.toList()
for (key in outgoingKeys) {
outgoingRpc.remove(key)?.let {
logger.trace { "resumeAllOngoingCallsWithThrowable: resume request $key" }
it.request.continuation.resumeWithException(throwable)
}
}
streams.values.removeAll {
it.closeStream(throwable)
@@ -427,11 +428,10 @@ private class RpcClient(
private fun resumeWithRouteClosed(route: UID) {
val message = "Route $route closed"
val outgoingRpcIterator = outgoingRpc.iterator()
for ((_, value) in outgoingRpcIterator) {
if (value.request.route == route) {
outgoingRpcIterator.remove()
value.request.continuation.resumeWithException(RouteClosedException(route, rpcCallFailureMessage(value.request.call, message)))
val outgoingKeys = outgoingRpc.filterValues { it.request.route == route }.keys
for (key in outgoingKeys) {
outgoingRpc.remove(key)?.let {
it.request.continuation.resumeWithException(RouteClosedException(route, rpcCallFailureMessage(it.request.call, message)))
}
}
streams.values.removeAll {