From 65fc2bfe7ba399ae1acd58588da5dbf187c73c60 Mon Sep 17 00:00:00 2001 From: a10zn8 Date: Wed, 11 Jan 2023 15:39:35 +0400 Subject: [PATCH] missing responses through a WS connection --- .../upstream/ethereum/WsConnectionImpl.kt | 88 +++++++++++-------- 1 file changed, 50 insertions(+), 38 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionImpl.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionImpl.kt index 7d6315d1..227af6f0 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionImpl.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionImpl.kt @@ -53,6 +53,7 @@ import java.net.URI import java.time.Duration import java.time.Instant import java.util.Base64 +import java.util.concurrent.ConcurrentHashMap import java.util.concurrent.Executors import java.util.concurrent.ScheduledFuture import java.util.concurrent.TimeUnit @@ -100,10 +101,10 @@ open class WsConnectionImpl( private val resetBackoffExecutor = Executors.newScheduledThreadPool(2) private var resetBackoffTask: ScheduledFuture? = null - private val messages = Sinks + private val subscriptionResponses = Sinks .many() .multicast() - .directBestEffort() + .directBestEffort() private var rpcSend = Sinks .many() @@ -114,6 +115,9 @@ open class WsConnectionImpl( .many() .multicast() .directBestEffort() + + private val currentRequests = ConcurrentHashMap>() + private val sendIdSeq = AtomicInteger(IDS_START) private val sendExecutor = Executors.newSingleThreadExecutor() private var keepConnection = true @@ -272,41 +276,48 @@ open class WsConnectionImpl( fun onMessage(msg: ResponseWSParser.WsResponse): Mono { return Mono.fromCallable { - val status = messages.tryEmitNext(msg) - if (status.isFailure) { - if (status == Sinks.EmitResult.FAIL_ZERO_SUBSCRIBER) { - log.debug("No subscribers to WS response") - } else { - log.warn("Failed to proceed with a WS message: $status") - } + when (msg.type) { + ResponseWSParser.Type.RPC -> onMessageRpc(msg) + ResponseWSParser.Type.SUBSCRIPTION -> onMessageSubscription(msg) } }.then() } - open fun getRpcResponses(): Flux { - return Flux.from(messages.asFlux()) - .publishOn(Schedulers.boundedElastic()) - .filter { - it.type == ResponseWSParser.Type.RPC + fun onMessageRpc(msg: ResponseWSParser.WsResponse) { + val rpcResponse = JsonRpcResponse( + msg.value, msg.error, msg.id, null + ) + val sender = currentRequests.remove(msg.id.asNumber().toInt()) + if (sender == null) { + log.warn("Unknown response received for ${msg.id}") + } else { + try { + val emitResult = sender.tryEmitValue(rpcResponse) + if (emitResult.isFailure) { + log.debug("Response is ignored $emitResult") + } + } catch (t: Throwable) { + log.warn("Response is not processed ${t.message}", t) } - .map { msg -> - JsonRpcResponse( - msg.value, msg.error, msg.id, null - ) + } + } + + fun onMessageSubscription(msg: ResponseWSParser.WsResponse) { + val subscription = JsonRpcWsMessage( + msg.value, msg.error, msg.id.asString(), + ) + val status = subscriptionResponses.tryEmitNext(subscription) + if (status.isFailure) { + if (status == Sinks.EmitResult.FAIL_ZERO_SUBSCRIBER) { + log.debug("No subscribers to WS response") + } else { + log.warn("Failed to proceed with a WS message: $status") } + } } open fun getSubscribeResponses(): Flux { - return Flux.from(messages.asFlux()) - .publishOn(Schedulers.boundedElastic()) - .filter { - it.type == ResponseWSParser.Type.SUBSCRIPTION - } - .map { msg -> - JsonRpcWsMessage( - msg.value, msg.error, msg.id.asString(), - ) - } + return Flux.from(subscriptionResponses.asFlux()) } open fun callRpc(originalRequest: JsonRpcRequest): Mono { @@ -332,7 +343,9 @@ open class WsConnectionImpl( } fun waitForResponse(request: JsonRpcRequest, originalId: Int, startTime: Long): Mono { - val expectedId = request.id.toLong() + val internalId = request.id.toLong() + val onResponse = Sinks.one() + currentRequests[internalId.toInt()] = onResponse val noResponse = JsonRpcException( JsonRpcResponse.Id.from(originalId), JsonRpcError( @@ -342,13 +355,6 @@ open class WsConnectionImpl( false ) - val response = Flux.from(getRpcResponses()) - .doOnSubscribe { sendRpc(request) } - .filter { resp -> resp.id.asNumber() == expectedId } - .take(Defaults.timeout) - .take(1) - .singleOrEmpty() - val failOnDisconnect = Mono.from(disconnects.asFlux()) .flatMap { Mono.error( @@ -362,11 +368,16 @@ open class WsConnectionImpl( ) } - return response.or(failOnDisconnect) + return Mono.from(onResponse.asMono()).or(failOnDisconnect) + .doOnSubscribe { sendRpc(request) } + .take(Defaults.timeout) .doOnNext { rpcMetrics?.timer?.record(System.nanoTime() - startTime, TimeUnit.NANOSECONDS) } .doOnError { rpcMetrics?.fails?.increment() } .map { it.copyWithId(JsonRpcResponse.Id.from(originalId)) } - .switchIfEmpty(Mono.error(noResponse)) + .switchIfEmpty( + Mono.fromCallable { log.warn("No response for ${request.method} ${request.params}") }.then(Mono.error(noResponse)) + ) + .doFinally { currentRequests.remove(internalId.toInt()) } } override fun close() { @@ -374,5 +385,6 @@ open class WsConnectionImpl( keepConnection = false connection?.dispose() connection = null + currentRequests.clear() } }