diff --git a/src/main/kotlin/io/emeraldpay/dshackle/quorum/QuorumRpcReader.kt b/src/main/kotlin/io/emeraldpay/dshackle/quorum/QuorumRpcReader.kt index 4a8911b9..214b95bc 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/quorum/QuorumRpcReader.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/quorum/QuorumRpcReader.kt @@ -31,6 +31,7 @@ import reactor.util.function.Tuple2 import reactor.util.function.Tuple3 import reactor.util.function.Tuples import java.util.Optional +import java.util.concurrent.atomic.AtomicInteger import java.util.function.BiFunction import java.util.function.Function @@ -190,6 +191,9 @@ class QuorumRpcReader( } } + fun getValidAttemptsCount(): AtomicInteger = + apiControl.attempts() + class Result( val value: ByteArray, val signature: ResponseSigner.Signature?, diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt index a6d99c3e..3f006391 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt @@ -24,6 +24,7 @@ import io.emeraldpay.dshackle.SilentException import io.emeraldpay.dshackle.quorum.CallQuorum import io.emeraldpay.dshackle.quorum.NotLaggingQuorum import io.emeraldpay.dshackle.quorum.QuorumReaderFactory +import io.emeraldpay.dshackle.quorum.QuorumRpcReader import io.emeraldpay.dshackle.startup.ConfiguredUpstreams import io.emeraldpay.dshackle.upstream.ApiSource import io.emeraldpay.dshackle.upstream.Multistream @@ -40,6 +41,7 @@ import io.emeraldpay.etherjar.rpc.RpcException import io.emeraldpay.etherjar.rpc.RpcResponseError import io.emeraldpay.grpc.BlockchainType import io.emeraldpay.grpc.Chain +import io.micrometer.core.instrument.Metrics import org.apache.commons.lang3.StringUtils import org.slf4j.LoggerFactory import org.springframework.beans.factory.annotation.Autowired @@ -48,6 +50,7 @@ import reactor.core.publisher.Flux import reactor.core.publisher.Mono import reactor.kotlin.core.publisher.toMono import java.util.EnumMap +import java.util.concurrent.atomic.AtomicInteger @Service open class NativeCall( @@ -64,7 +67,9 @@ open class NativeCall( init { multistreamHolder.observeChains().subscribe { chain -> - if ((BlockchainType.from(chain) == BlockchainType.ETHEREUM_POS || BlockchainType.from(chain) == BlockchainType.ETHEREUM) && !ethereumCallSelectors.containsKey(chain) + if ((BlockchainType.from(chain) == BlockchainType.ETHEREUM_POS || BlockchainType.from(chain) == BlockchainType.ETHEREUM) && !ethereumCallSelectors.containsKey( + chain + ) ) { multistreamHolder.getUpstream(chain)?.let { up -> val reader = up.cast(EthereumPosMultiStream::class.java).getReader() @@ -118,7 +123,10 @@ open class NativeCall( return result.build() } - fun buildSignature(nonce: Long, signature: ResponseSigner.Signature): BlockchainOuterClass.NativeCallReplySignature { + fun buildSignature( + nonce: Long, + signature: ResponseSigner.Signature + ): BlockchainOuterClass.NativeCallReplySignature { val msg = BlockchainOuterClass.NativeCallReplySignature.newBuilder() msg.signature = ByteString.copyFrom(signature.value) msg.keyId = signature.keyId @@ -154,6 +162,11 @@ open class NativeCall( val matcher = Selector.convertToMatcher(request.selector) if (!configuredUpstreams.hasMatchingUpstream(chain, matcher)) { + if (Global.metricsExtended) { + Metrics.globalRegistry + .counter("no_matching_upstream", "chain", chain.chainCode, "matcher", matcher.describeInternal()) + .increment() + } return Flux.error(CallFailure(0, SilentException.NoMatchingUpstream(matcher))) } @@ -188,7 +201,11 @@ open class NativeCall( val errorMessage = "The method $method does not exist/is not available" return Mono.just( InvalidCallContext( - CallError(requestItem.id, errorMessage, JsonRpcError(RpcResponseError.CODE_METHOD_NOT_EXIST, errorMessage)) + CallError( + requestItem.id, + errorMessage, + JsonRpcError(RpcResponseError.CODE_METHOD_NOT_EXIST, errorMessage) + ) ) ) } @@ -216,7 +233,14 @@ open class NativeCall( matcher.withMatcher(heightMatcher) } val nonce = requestItem.nonce.let { if (it == 0L) null else it } - ValidCallContext(requestItem.id, nonce, upstream, matcher.build(), callQuorum, RawCallDetails(method, params)) + ValidCallContext( + requestItem.id, + nonce, + upstream, + matcher.build(), + callQuorum, + RawCallDetails(method, params) + ) } } @@ -250,6 +274,11 @@ open class NativeCall( return Mono.error(RpcException(RpcResponseError.CODE_METHOD_NOT_EXIST, "Unsupported method")) } val reader = quorumReaderFactory.create(ctx.getApis(), ctx.callQuorum, signer) + val counter = if (reader is QuorumRpcReader) { + reader.getValidAttemptsCount() + } else { + AtomicInteger(-1) + } return reader .read(JsonRpcRequest(ctx.payload.method, ctx.payload.params, ctx.nonce)) .map { @@ -264,10 +293,47 @@ open class NativeCall( Mono.just(failure) } .switchIfEmpty( - Mono.just(CallResult.fail(ctx.id, ctx.nonce, 1, "No response or no available upstream for ${ctx.payload.method}")) + Mono.fromSupplier { + counter.get().let { attempts -> + CallResult.fail( + ctx.id, + ctx.nonce, + 1, + errorMessage(attempts, ctx.payload.method) + ).also { + countFailure(attempts, ctx) + } + } + } ) } + private fun errorMessage(attempts: Int, method: String): String = + when (attempts) { + -1 -> "No response or no available upstream for $method" + 0 -> "No available upstream for $method" + else -> "No response for $method" + } + + private fun countFailure(counter: Int, ctx: ValidCallContext) = + when (counter) { + -1 -> "UNDEFINED" + 0 -> "NO_AVAIL_UPSTREAM" + else -> "NO_RESPONSE" + }.let { reason -> + if (Global.metricsExtended) { + Metrics.globalRegistry.counter( + "native_call_failure", + "upstream", + ctx.upstream.getId(), + "reason", + reason, + "chain", + ctx.upstream.chain.chainCode, + ).increment() + } + } + @Suppress("UNCHECKED_CAST") private fun extractParams(jsonParams: String): List { if (StringUtils.isEmpty(jsonParams)) { @@ -346,7 +412,13 @@ open class NativeCall( } } - open class CallResult(val id: Int, val nonce: Long?, val result: ByteArray?, val error: CallError?, val signature: ResponseSigner.Signature?) { + open class CallResult( + val id: Int, + val nonce: Long?, + val result: ByteArray?, + val error: CallError?, + val signature: ResponseSigner.Signature? + ) { companion object { fun ok(id: Int, nonce: Long?, result: ByteArray, signature: ResponseSigner.Signature?): CallResult { return CallResult(id, nonce, result, null, signature) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ApiSource.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ApiSource.kt index c9f747fe..66a977d4 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ApiSource.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ApiSource.kt @@ -17,6 +17,7 @@ package io.emeraldpay.dshackle.upstream import org.reactivestreams.Publisher +import java.util.concurrent.atomic.AtomicInteger interface ApiSource : Publisher { @@ -26,4 +27,6 @@ interface ApiSource : Publisher { * Must be called before actual use, it spins off control flow of the API Source */ fun request(tries: Int) + + fun attempts(): AtomicInteger } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/FilteredApis.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/FilteredApis.kt index 5f72e89b..36ee6615 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/FilteredApis.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/FilteredApis.kt @@ -28,6 +28,7 @@ import reactor.core.publisher.Flux import reactor.core.publisher.Sinks import java.time.Duration import java.util.EnumMap +import java.util.concurrent.atomic.AtomicInteger import java.util.concurrent.locks.Lock import java.util.concurrent.locks.ReentrantLock import kotlin.concurrent.withLock @@ -89,6 +90,8 @@ class FilteredApis( private val secondaryUpstreams: List private val standardWithFallback: List + private val counter: AtomicInteger = AtomicInteger(0) + private var started = false private val control = Sinks.many().unicast().onBackpressureBuffer() @@ -177,7 +180,13 @@ class FilteredApis( .doFinally { metrics[chain]?.tried?.record(count.toDouble()) } } - result.filter { up -> up.isAvailable() && matcher.matches(up) } + result.filter { up -> + (up.isAvailable() && matcher.matches(up)).also { + if (it) { + counter.incrementAndGet() + } + } + } .zipWith(control.asFlux()) .map { it.t1 } .doOnSubscribe { @@ -201,6 +210,9 @@ class FilteredApis( } } + override fun attempts(): AtomicInteger = + counter + override fun toString(): String { return "Filter API: ${allUpstreams.size} upstreams with $matcher" }