problem: confusing metric name and no metric to monitor RPC error responses
rel: #150
This commit is contained in:
@@ -73,13 +73,16 @@ abstract class BaseHandler(
|
|||||||
metricById(it.id)?.let { metrics ->
|
metricById(it.id)?.let { metrics ->
|
||||||
metrics.requestMetric.increment()
|
metrics.requestMetric.increment()
|
||||||
metrics.callMetric.record(System.currentTimeMillis() - startTime, TimeUnit.MILLISECONDS)
|
metrics.callMetric.record(System.currentTimeMillis() - startTime, TimeUnit.MILLISECONDS)
|
||||||
|
if (it.isError()) {
|
||||||
|
metrics.errorMetric.increment()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
handler.onResponse(it)
|
handler.onResponse(it)
|
||||||
}
|
}
|
||||||
.doOnError {
|
.doOnError {
|
||||||
// when error happened the whole flux is stopped and no result is produced, so we should mark all the requests as failed
|
// when error happened the whole flux is stopped and no result is produced, so we should mark all the requests as failed
|
||||||
items.forEach { item ->
|
items.forEach { item ->
|
||||||
requestMetrics.get(chain, item.method).errorMetric.increment()
|
requestMetrics.get(chain, item.method).failMetric.increment()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -173,29 +173,44 @@ class ProxyServer(
|
|||||||
interface RequestMetrics {
|
interface RequestMetrics {
|
||||||
val callMetric: Timer
|
val callMetric: Timer
|
||||||
val errorMetric: Counter
|
val errorMetric: Counter
|
||||||
|
val failMetric: Counter
|
||||||
val requestMetric: Counter
|
val requestMetric: Counter
|
||||||
}
|
}
|
||||||
|
|
||||||
class RequestMetricsBasic(chain: Chain) : RequestMetrics {
|
class RequestMetricsBasic(chain: Chain) : RequestMetrics {
|
||||||
override val callMetric = Timer.builder("request.jsonrpc.call")
|
override val callMetric = Timer.builder("request.jsonrpc.call")
|
||||||
|
.description("Time to process a call")
|
||||||
.tags("chain", chain.chainCode)
|
.tags("chain", chain.chainCode)
|
||||||
.register(Metrics.globalRegistry)
|
.register(Metrics.globalRegistry)
|
||||||
override val errorMetric = Counter.builder("request.jsonrpc.err")
|
override val errorMetric = Counter.builder("request.jsonrpc.err")
|
||||||
|
.description("Number of requests ended with an error response")
|
||||||
|
.tags("chain", chain.chainCode)
|
||||||
|
.register(Metrics.globalRegistry)
|
||||||
|
override val failMetric = Counter.builder("request.jsonrpc.fail")
|
||||||
|
.description("Number of requests failed to process")
|
||||||
.tags("chain", chain.chainCode)
|
.tags("chain", chain.chainCode)
|
||||||
.register(Metrics.globalRegistry)
|
.register(Metrics.globalRegistry)
|
||||||
override val requestMetric = Counter.builder("request.jsonrpc.request.total")
|
override val requestMetric = Counter.builder("request.jsonrpc.request.total")
|
||||||
|
.description("Number of requests")
|
||||||
.tags("chain", chain.chainCode)
|
.tags("chain", chain.chainCode)
|
||||||
.register(Metrics.globalRegistry)
|
.register(Metrics.globalRegistry)
|
||||||
}
|
}
|
||||||
|
|
||||||
class RequestMetricsWithMethod(chain: Chain, method: String) : RequestMetrics {
|
class RequestMetricsWithMethod(chain: Chain, method: String) : RequestMetrics {
|
||||||
override val callMetric = Timer.builder("request.jsonrpc.call")
|
override val callMetric = Timer.builder("request.jsonrpc.call")
|
||||||
|
.description("Time to process a call")
|
||||||
.tags("chain", chain.chainCode, "method", method)
|
.tags("chain", chain.chainCode, "method", method)
|
||||||
.register(Metrics.globalRegistry)
|
.register(Metrics.globalRegistry)
|
||||||
override val errorMetric = Counter.builder("request.jsonrpc.err")
|
override val errorMetric = Counter.builder("request.jsonrpc.err")
|
||||||
|
.description("Number of requests ended with an error response")
|
||||||
|
.tags("chain", chain.chainCode, "method", method)
|
||||||
|
.register(Metrics.globalRegistry)
|
||||||
|
override val failMetric = Counter.builder("request.jsonrpc.fail")
|
||||||
|
.description("Number of requests failed to process")
|
||||||
.tags("chain", chain.chainCode, "method", method)
|
.tags("chain", chain.chainCode, "method", method)
|
||||||
.register(Metrics.globalRegistry)
|
.register(Metrics.globalRegistry)
|
||||||
override val requestMetric = Counter.builder("request.jsonrpc.request.total")
|
override val requestMetric = Counter.builder("request.jsonrpc.request.total")
|
||||||
|
.description("Number of requests")
|
||||||
.tags("chain", chain.chainCode, "method", method)
|
.tags("chain", chain.chainCode, "method", method)
|
||||||
.register(Metrics.globalRegistry)
|
.register(Metrics.globalRegistry)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -57,7 +57,8 @@ class BlockchainRpc(
|
|||||||
.tag("type", "subscribeStatus")
|
.tag("type", "subscribeStatus")
|
||||||
.tag("chain", "NA")
|
.tag("chain", "NA")
|
||||||
.register(Metrics.globalRegistry)
|
.register(Metrics.globalRegistry)
|
||||||
private val errorMetric = Counter.builder("request.grpc.err")
|
private val failMetric = Counter.builder("request.grpc.fail")
|
||||||
|
.description("Number of requests failed to process")
|
||||||
.register(Metrics.globalRegistry)
|
.register(Metrics.globalRegistry)
|
||||||
private val chainMetrics = ChainValue { chain -> RequestMetrics(chain) }
|
private val chainMetrics = ChainValue { chain -> RequestMetrics(chain) }
|
||||||
|
|
||||||
@@ -71,9 +72,14 @@ class BlockchainRpc(
|
|||||||
metrics!!.nativeCallMetric.increment()
|
metrics!!.nativeCallMetric.increment()
|
||||||
startTime = System.currentTimeMillis()
|
startTime = System.currentTimeMillis()
|
||||||
}
|
}
|
||||||
).doOnNext {
|
).doOnNext { reply ->
|
||||||
metrics?.nativeCallRespMetric?.record(System.currentTimeMillis() - startTime, TimeUnit.MILLISECONDS)
|
metrics?.let { m ->
|
||||||
}.doOnError { errorMetric.increment() }
|
m.nativeCallRespMetric?.record(System.currentTimeMillis() - startTime, TimeUnit.MILLISECONDS)
|
||||||
|
if (!reply.succeed) {
|
||||||
|
m.nativeCallErrRespMetric.increment()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}.doOnError { failMetric.increment() }
|
||||||
}
|
}
|
||||||
|
|
||||||
override fun nativeSubscribe(request: Mono<BlockchainOuterClass.NativeSubscribeRequest>): Flux<BlockchainOuterClass.NativeSubscribeReplyItem> {
|
override fun nativeSubscribe(request: Mono<BlockchainOuterClass.NativeSubscribeRequest>): Flux<BlockchainOuterClass.NativeSubscribeReplyItem> {
|
||||||
@@ -86,14 +92,14 @@ class BlockchainRpc(
|
|||||||
}
|
}
|
||||||
).doOnNext {
|
).doOnNext {
|
||||||
metrics?.nativeSubscribeRespMetric?.increment()
|
metrics?.nativeSubscribeRespMetric?.increment()
|
||||||
}.doOnError { errorMetric.increment() }
|
}.doOnError { failMetric.increment() }
|
||||||
}
|
}
|
||||||
|
|
||||||
override fun subscribeHead(request: Mono<Common.Chain>): Flux<BlockchainOuterClass.ChainHead> {
|
override fun subscribeHead(request: Mono<Common.Chain>): Flux<BlockchainOuterClass.ChainHead> {
|
||||||
return streamHead.add(
|
return streamHead.add(
|
||||||
request
|
request
|
||||||
.doOnNext { chainMetrics.get(it.type).subscribeHeadMetric.increment() }
|
.doOnNext { chainMetrics.get(it.type).subscribeHeadMetric.increment() }
|
||||||
).doOnError { errorMetric.increment() }
|
).doOnError { failMetric.increment() }
|
||||||
}
|
}
|
||||||
|
|
||||||
override fun subscribeTxStatus(requestMono: Mono<BlockchainOuterClass.TxStatusRequest>): Flux<BlockchainOuterClass.TxStatus> {
|
override fun subscribeTxStatus(requestMono: Mono<BlockchainOuterClass.TxStatusRequest>): Flux<BlockchainOuterClass.TxStatus> {
|
||||||
@@ -105,11 +111,11 @@ class BlockchainRpc(
|
|||||||
trackTx.find { it.isSupported(chain) }?.let { track ->
|
trackTx.find { it.isSupported(chain) }?.let { track ->
|
||||||
track.subscribe(request)
|
track.subscribe(request)
|
||||||
.doOnNext { metrics.subscribeHeadRespMetric.increment() }
|
.doOnNext { metrics.subscribeHeadRespMetric.increment() }
|
||||||
.doOnError { errorMetric.increment() }
|
.doOnError { failMetric.increment() }
|
||||||
} ?: Flux.error(SilentException.UnsupportedBlockchain(chain))
|
} ?: Flux.error(SilentException.UnsupportedBlockchain(chain))
|
||||||
} catch (t: Throwable) {
|
} catch (t: Throwable) {
|
||||||
log.error("Internal error during Tx Subscription", t)
|
log.error("Internal error during Tx Subscription", t)
|
||||||
errorMetric.increment()
|
failMetric.increment()
|
||||||
Flux.error<BlockchainOuterClass.TxStatus>(IllegalStateException("Internal Error"))
|
Flux.error<BlockchainOuterClass.TxStatus>(IllegalStateException("Internal Error"))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -125,14 +131,14 @@ class BlockchainRpc(
|
|||||||
trackAddress.find { it.isSupported(chain, asset) }?.let { track ->
|
trackAddress.find { it.isSupported(chain, asset) }?.let { track ->
|
||||||
track.subscribe(request)
|
track.subscribe(request)
|
||||||
.doOnNext { metrics.subscribeBalanceRespMetric.increment() }
|
.doOnNext { metrics.subscribeBalanceRespMetric.increment() }
|
||||||
.doOnError { errorMetric.increment() }
|
.doOnError { failMetric.increment() }
|
||||||
} ?: Flux.error<BlockchainOuterClass.AddressBalance>(SilentException.UnsupportedBlockchain(chain))
|
} ?: Flux.error<BlockchainOuterClass.AddressBalance>(SilentException.UnsupportedBlockchain(chain))
|
||||||
.doOnSubscribe {
|
.doOnSubscribe {
|
||||||
log.error("Balance for $chain:$asset is not supported")
|
log.error("Balance for $chain:$asset is not supported")
|
||||||
}
|
}
|
||||||
} catch (t: Throwable) {
|
} catch (t: Throwable) {
|
||||||
log.error("Internal error during Balance Subscription", t)
|
log.error("Internal error during Balance Subscription", t)
|
||||||
errorMetric.increment()
|
failMetric.increment()
|
||||||
Flux.error<BlockchainOuterClass.AddressBalance>(IllegalStateException("Internal Error"))
|
Flux.error<BlockchainOuterClass.AddressBalance>(IllegalStateException("Internal Error"))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -160,7 +166,7 @@ class BlockchainRpc(
|
|||||||
}
|
}
|
||||||
} catch (t: Throwable) {
|
} catch (t: Throwable) {
|
||||||
log.error("Internal error during Balance Request", t)
|
log.error("Internal error during Balance Request", t)
|
||||||
errorMetric.increment()
|
failMetric.increment()
|
||||||
Flux.error<BlockchainOuterClass.AddressBalance>(IllegalStateException("Internal Error"))
|
Flux.error<BlockchainOuterClass.AddressBalance>(IllegalStateException("Internal Error"))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -182,20 +188,20 @@ class BlockchainRpc(
|
|||||||
}
|
}
|
||||||
.doOnError { t ->
|
.doOnError { t ->
|
||||||
log.error("Internal error during Fee Estimation", t)
|
log.error("Internal error during Fee Estimation", t)
|
||||||
errorMetric.increment()
|
failMetric.increment()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
override fun describe(request: Mono<BlockchainOuterClass.DescribeRequest>): Mono<BlockchainOuterClass.DescribeResponse> {
|
override fun describe(request: Mono<BlockchainOuterClass.DescribeRequest>): Mono<BlockchainOuterClass.DescribeResponse> {
|
||||||
describeMetric.increment()
|
describeMetric.increment()
|
||||||
return describe.describe(request)
|
return describe.describe(request)
|
||||||
.doOnError { errorMetric.increment() }
|
.doOnError { failMetric.increment() }
|
||||||
}
|
}
|
||||||
|
|
||||||
override fun subscribeStatus(request: Mono<BlockchainOuterClass.StatusRequest>): Flux<BlockchainOuterClass.ChainStatus> {
|
override fun subscribeStatus(request: Mono<BlockchainOuterClass.StatusRequest>): Flux<BlockchainOuterClass.ChainStatus> {
|
||||||
subscribeStatusMetric.increment()
|
subscribeStatusMetric.increment()
|
||||||
return subscribeStatus.subscribeStatus(request)
|
return subscribeStatus.subscribeStatus(request)
|
||||||
.doOnError { errorMetric.increment() }
|
.doOnError { failMetric.increment() }
|
||||||
}
|
}
|
||||||
|
|
||||||
class RequestMetrics(chain: Chain) {
|
class RequestMetrics(chain: Chain) {
|
||||||
@@ -208,6 +214,10 @@ class BlockchainRpc(
|
|||||||
.tag("chain", chain.chainCode)
|
.tag("chain", chain.chainCode)
|
||||||
.publishPercentileHistogram()
|
.publishPercentileHistogram()
|
||||||
.register(Metrics.globalRegistry)
|
.register(Metrics.globalRegistry)
|
||||||
|
val nativeCallErrRespMetric = Counter.builder("request.grpc.response.err")
|
||||||
|
.tag("type", "nativeCall")
|
||||||
|
.tag("chain", chain.chainCode)
|
||||||
|
.register(Metrics.globalRegistry)
|
||||||
val nativeSubscribeMetric = Counter.builder("request.grpc.request")
|
val nativeSubscribeMetric = Counter.builder("request.grpc.request")
|
||||||
.tag("type", "nativeSubscribe")
|
.tag("type", "nativeSubscribe")
|
||||||
.tag("chain", chain.chainCode)
|
.tag("chain", chain.chainCode)
|
||||||
|
|||||||
@@ -222,6 +222,7 @@ open class NativeCall(
|
|||||||
}
|
}
|
||||||
|
|
||||||
fun executeOnRemote(ctx: ValidCallContext<ParsedCallDetails>): Mono<CallResult> {
|
fun executeOnRemote(ctx: ValidCallContext<ParsedCallDetails>): Mono<CallResult> {
|
||||||
|
// check if method is allowed to be executed at all
|
||||||
if (!ctx.upstream.getMethods().isCallable(ctx.payload.method)) {
|
if (!ctx.upstream.getMethods().isCallable(ctx.payload.method)) {
|
||||||
return Mono.error(RpcException(RpcResponseError.CODE_METHOD_NOT_EXIST, "Unsupported method"))
|
return Mono.error(RpcException(RpcResponseError.CODE_METHOD_NOT_EXIST, "Unsupported method"))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -271,8 +271,8 @@ open class ConfiguredUpstreams(
|
|||||||
.tags(metricsTags)
|
.tags(metricsTags)
|
||||||
.publishPercentileHistogram()
|
.publishPercentileHistogram()
|
||||||
.register(Metrics.globalRegistry),
|
.register(Metrics.globalRegistry),
|
||||||
Counter.builder("upstream.rpc.err")
|
Counter.builder("upstream.rpc.fail")
|
||||||
.description("Errors received on request through HTTP JSON RPC connection")
|
.description("Number of failures of HTTP JSON RPC requests")
|
||||||
.tags(metricsTags)
|
.tags(metricsTags)
|
||||||
.register(Metrics.globalRegistry)
|
.register(Metrics.globalRegistry)
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -67,8 +67,8 @@ class EthereumWsUpstream(
|
|||||||
.tags(metricsTags)
|
.tags(metricsTags)
|
||||||
.publishPercentileHistogram()
|
.publishPercentileHistogram()
|
||||||
.register(Metrics.globalRegistry),
|
.register(Metrics.globalRegistry),
|
||||||
Counter.builder("upstream.ws.err")
|
Counter.builder("upstream.ws.fail")
|
||||||
.description("Errors received on request through WebSocket JSON RPC connection")
|
.description("Number of failures of WebSocket JSON RPC requests")
|
||||||
.tags(metricsTags)
|
.tags(metricsTags)
|
||||||
.register(Metrics.globalRegistry)
|
.register(Metrics.globalRegistry)
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -387,7 +387,7 @@ class WsConnection(
|
|||||||
.take(1)
|
.take(1)
|
||||||
.singleOrEmpty()
|
.singleOrEmpty()
|
||||||
.doOnNext { rpcMetrics?.timer?.record(System.nanoTime() - startTime, TimeUnit.NANOSECONDS) }
|
.doOnNext { rpcMetrics?.timer?.record(System.nanoTime() - startTime, TimeUnit.NANOSECONDS) }
|
||||||
.doOnError { rpcMetrics?.errors?.increment() }
|
.doOnError { rpcMetrics?.fails?.increment() }
|
||||||
.map { it.copyWithId(JsonRpcResponse.Id.from(originalId)) }
|
.map { it.copyWithId(JsonRpcResponse.Id.from(originalId)) }
|
||||||
.defaultIfEmpty(failResponse)
|
.defaultIfEmpty(failResponse)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -175,12 +175,12 @@ class GrpcUpstreams(
|
|||||||
|
|
||||||
val metrics = RpcMetrics(
|
val metrics = RpcMetrics(
|
||||||
Timer.builder("upstream.grpc.conn")
|
Timer.builder("upstream.grpc.conn")
|
||||||
.description("Request time through a gRPC connection")
|
.description("Request time through a Dshackle/gRPC connection")
|
||||||
.tags(metricsTags)
|
.tags(metricsTags)
|
||||||
.publishPercentileHistogram()
|
.publishPercentileHistogram()
|
||||||
.register(Metrics.globalRegistry),
|
.register(Metrics.globalRegistry),
|
||||||
Counter.builder("upstream.grpc.err")
|
Counter.builder("upstream.grpc.fail")
|
||||||
.description("Errors received on request through gRPC connection")
|
.description("Number of failures of Dshackle/gRPC requests")
|
||||||
.tags(metricsTags)
|
.tags(metricsTags)
|
||||||
.register(Metrics.globalRegistry)
|
.register(Metrics.globalRegistry)
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -79,7 +79,7 @@ class JsonRpcGrpcClient(
|
|||||||
val bytes = resp.payload.toByteArray()
|
val bytes = resp.payload.toByteArray()
|
||||||
Mono.just(JsonRpcResponse(bytes, null))
|
Mono.just(JsonRpcResponse(bytes, null))
|
||||||
} else {
|
} else {
|
||||||
metrics.errors.increment()
|
metrics.fails.increment()
|
||||||
Mono.error(
|
Mono.error(
|
||||||
RpcException(
|
RpcException(
|
||||||
RpcResponseError.CODE_UPSTREAM_CONNECTION_ERROR,
|
RpcResponseError.CODE_UPSTREAM_CONNECTION_ERROR,
|
||||||
|
|||||||
@@ -126,7 +126,7 @@ class JsonRpcHttpClient(
|
|||||||
is JsonRpcException -> JsonRpcResponse.error(t.error, JsonRpcResponse.NumberId(1))
|
is JsonRpcException -> JsonRpcResponse.error(t.error, JsonRpcResponse.NumberId(1))
|
||||||
else -> JsonRpcResponse.error(1, t.message ?: t.javaClass.name)
|
else -> JsonRpcResponse.error(1, t.message ?: t.javaClass.name)
|
||||||
}
|
}
|
||||||
metrics.errors.increment()
|
metrics.fails.increment()
|
||||||
Mono.just(err)
|
Mono.just(err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -20,5 +20,5 @@ import io.micrometer.core.instrument.Timer
|
|||||||
|
|
||||||
class RpcMetrics(
|
class RpcMetrics(
|
||||||
val timer: Timer,
|
val timer: Timer,
|
||||||
val errors: Counter
|
val fails: Counter
|
||||||
)
|
)
|
||||||
|
|||||||
Reference in New Issue
Block a user