solution: better handle for minor upstream errors
This commit is contained in:
@@ -25,6 +25,7 @@ import io.emeraldpay.grpc.Chain
|
|||||||
import io.infinitape.etherjar.domain.BlockHash
|
import io.infinitape.etherjar.domain.BlockHash
|
||||||
import io.infinitape.etherjar.domain.TransactionId
|
import io.infinitape.etherjar.domain.TransactionId
|
||||||
import io.infinitape.etherjar.rpc.Commands
|
import io.infinitape.etherjar.rpc.Commands
|
||||||
|
import io.infinitape.etherjar.rpc.RpcException
|
||||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||||
import io.infinitape.etherjar.rpc.json.TransactionJson
|
import io.infinitape.etherjar.rpc.json.TransactionJson
|
||||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||||
@@ -252,6 +253,10 @@ class TrackTx(
|
|||||||
val execution = upstream.getApi(Selector.empty)
|
val execution = upstream.getApi(Selector.empty)
|
||||||
.flatMap { api -> api.executeAndConvert(Commands.eth().getTransaction(tx.txid)) }
|
.flatMap { api -> api.executeAndConvert(Commands.eth().getTransaction(tx.txid)) }
|
||||||
return execution
|
return execution
|
||||||
|
.onErrorResume(RpcException::class.java) { t ->
|
||||||
|
log.warn("Upstream error, ignoring. {}", t.rpcMessage)
|
||||||
|
Mono.empty<TransactionJson>()
|
||||||
|
}
|
||||||
.flatMap { updateFromBlock(upstream, tx, it) }
|
.flatMap { updateFromBlock(upstream, tx, it) }
|
||||||
.doOnError { t ->
|
.doOnError { t ->
|
||||||
log.error("Failed to load tx block", t)
|
log.error("Failed to load tx block", t)
|
||||||
|
|||||||
@@ -18,10 +18,9 @@ package io.emeraldpay.dshackle.upstream.ethereum
|
|||||||
import com.fasterxml.jackson.databind.ObjectMapper
|
import com.fasterxml.jackson.databind.ObjectMapper
|
||||||
import io.emeraldpay.dshackle.Defaults
|
import io.emeraldpay.dshackle.Defaults
|
||||||
import io.emeraldpay.dshackle.upstream.CallMethods
|
import io.emeraldpay.dshackle.upstream.CallMethods
|
||||||
import io.infinitape.etherjar.rpc.ReactorBatch
|
import io.grpc.Status
|
||||||
import io.infinitape.etherjar.rpc.ReactorRpcClient
|
import io.grpc.StatusRuntimeException
|
||||||
import io.infinitape.etherjar.rpc.RpcCall
|
import io.infinitape.etherjar.rpc.*
|
||||||
import io.infinitape.etherjar.rpc.RpcException
|
|
||||||
import io.infinitape.etherjar.rpc.json.ResponseJson
|
import io.infinitape.etherjar.rpc.json.ResponseJson
|
||||||
import org.slf4j.LoggerFactory
|
import org.slf4j.LoggerFactory
|
||||||
import reactor.core.publisher.Mono
|
import reactor.core.publisher.Mono
|
||||||
@@ -52,11 +51,18 @@ open class DirectEthereumApi(
|
|||||||
resp.result = it
|
resp.result = it
|
||||||
objectMapper.writer().writeValueAsBytes(resp)
|
objectMapper.writer().writeValueAsBytes(resp)
|
||||||
}
|
}
|
||||||
|
.onErrorResume(StatusRuntimeException::class.java) { t ->
|
||||||
|
if (t.status.code == Status.Code.CANCELLED) {
|
||||||
|
Mono.empty<ByteArray>()
|
||||||
|
} else {
|
||||||
|
Mono.error(RpcException(RpcResponseError.CODE_UPSTREAM_CONNECTION_ERROR, "gRPC error ${t.status}"))
|
||||||
|
}
|
||||||
|
}
|
||||||
.onErrorMap { t ->
|
.onErrorMap { t ->
|
||||||
if (RpcException::class.java.isAssignableFrom(t.javaClass)) {
|
if (RpcException::class.java.isAssignableFrom(t.javaClass)) {
|
||||||
t
|
t
|
||||||
} else {
|
} else {
|
||||||
log.warn("Convert to RPC error. Exception: ${t.message}")
|
log.warn("Convert to RPC error. Exception ${t.javaClass}:${t.message}", t)
|
||||||
RpcException(-32020, "Error reading from upstream", null, t)
|
RpcException(-32020, "Error reading from upstream", null, t)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user