solution: use reactor-based rpc client

This commit is contained in:
Igor Artamonov
2019-10-11 22:54:12 -04:00
parent fd95aa4868
commit 77685c314a
16 changed files with 131 additions and 121 deletions

View File

@@ -24,8 +24,7 @@ import io.emeraldpay.dshackle.upstream.ethereum.EthereumUpstream
import io.emeraldpay.dshackle.upstream.ethereum.EthereumWs
import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstreams
import io.emeraldpay.grpc.Chain
import io.infinitape.etherjar.rpc.DefaultRpcClient
import io.infinitape.etherjar.rpc.transport.DefaultRpcTransport
import io.infinitape.etherjar.rpc.http.ReactorHttpRpcClient
import org.apache.commons.lang3.StringUtils
import org.slf4j.LoggerFactory
import org.springframework.beans.factory.annotation.Autowired
@@ -132,18 +131,18 @@ open class ConfiguredUpstreams(
currentUpstreams.getDefaultMethods(chain)
}
conn.rpc?.let { endpoint ->
val rpcTransport = DefaultRpcTransport(endpoint.url)
val rpcClient = ReactorHttpRpcClient.newBuilder()
.setTarget(endpoint.url)
conn.rpc?.basicAuth?.let { auth ->
rpcTransport.setBasicAuth(auth.username, auth.password)
rpcClient.setBasicAuth(auth.username, auth.password)
}
conn.rpc?.tls?.let { tls ->
tls.ca?.let { ca ->
fileResolver.resolve(ca).inputStream().use { cert -> rpcTransport.setTrustedCertificate(cert) }
fileResolver.resolve(ca).inputStream().use { cert -> rpcClient.setTrustedCertificate(cert) }
}
}
val rpcClient = DefaultRpcClient(rpcTransport)
rpcApi = DirectEthereumApi(
rpcClient,
rpcClient.build(),
objectMapper,
methods
).apply {

View File

@@ -20,33 +20,50 @@ import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.upstream.ethereum.EthereumUpstream
import io.infinitape.etherjar.rpc.Batch
import io.infinitape.etherjar.rpc.Commands
import io.infinitape.etherjar.rpc.ReactorBatch
import org.slf4j.LoggerFactory
import org.springframework.scheduling.concurrent.CustomizableThreadFactory
import reactor.core.publisher.Flux
import reactor.core.publisher.Mono
import reactor.core.scheduler.Schedulers
import java.time.Duration
import java.util.concurrent.Executors
import java.util.concurrent.TimeUnit
class UpstreamValidator(
private val ethereumUpstream: EthereumUpstream,
private val options: UpstreamsConfig.Options
) {
companion object {
private val log = LoggerFactory.getLogger(UpstreamValidator::class.java)
val scheduler = Schedulers.fromExecutor(Executors.newCachedThreadPool(CustomizableThreadFactory("validator")))
}
fun validate(): Mono<UpstreamAvailability> {
val batch = Batch()
val peerCount = batch.add(Commands.net().peerCount())
val syncing = batch.add(Commands.eth().syncing())
val batch = ReactorBatch()
val peerCount = batch.add(Commands.net().peerCount()).result
val syncing = batch.add(Commands.eth().syncing()).result
return ethereumUpstream.getApi(Selector.empty)
.map { api -> api.rpcClient.execute(batch) }
.flatMap { Mono.fromCompletionStage(it) }
.timeout(Defaults.timeout)
.map {
if (syncing.get().isSyncing) {
UpstreamAvailability.SYNCING
} else if (options.minPeers != null && peerCount.get() < options.minPeers!!) {
UpstreamAvailability.IMMATURE
.subscribeOn(scheduler)
.flatMapMany { api -> api.rpcClient.execute(batch) }
.timeout(Defaults.timeout, Mono.error(Exception("Validation timeout")))
.then(syncing)
.flatMap { value ->
if (value.isSyncing) {
Mono.just(UpstreamAvailability.SYNCING)
} else {
UpstreamAvailability.OK
peerCount.map { count ->
val minPeers = options.minPeers ?: 1
if (count < minPeers) {
UpstreamAvailability.IMMATURE
} else {
UpstreamAvailability.OK
}
}
}
}.onErrorContinue { _, _ -> UpstreamAvailability.UNAVAILABLE }
}
.doOnError { err -> log.warn("Failed to validate upstream", err)}
.onErrorReturn(UpstreamAvailability.UNAVAILABLE)
}
fun start(): Flux<UpstreamAvailability> {

View File

@@ -18,8 +18,9 @@ package io.emeraldpay.dshackle.upstream.ethereum
import com.fasterxml.jackson.databind.ObjectMapper
import io.emeraldpay.dshackle.Defaults
import io.emeraldpay.dshackle.upstream.CallMethods
import io.infinitape.etherjar.rpc.ReactorBatch
import io.infinitape.etherjar.rpc.ReactorRpcClient
import io.infinitape.etherjar.rpc.RpcCall
import io.infinitape.etherjar.rpc.RpcClient
import io.infinitape.etherjar.rpc.RpcException
import io.infinitape.etherjar.rpc.json.ResponseJson
import org.slf4j.LoggerFactory
@@ -27,7 +28,7 @@ import reactor.core.publisher.Mono
import java.time.Duration
open class DirectEthereumApi(
val rpcClient: RpcClient,
val rpcClient: ReactorRpcClient,
private val objectMapper: ObjectMapper,
val targets: CallMethods
): EthereumApi(objectMapper) {
@@ -68,8 +69,7 @@ open class DirectEthereumApi(
}
private fun callUpstream(method: String, params: List<Any>): Mono<out Any> {
return Mono.fromCompletionStage(
rpcClient.execute(RpcCall.create(method, Any::class.java, params))
).timeout(timeout, Mono.error(RpcException(-32603, "Upstream timeout")))
return rpcClient.execute(RpcCall.create(method, Any::class.java, params))
.timeout(timeout, Mono.error(RpcException(-32603, "Upstream timeout")))
}
}

View File

@@ -18,6 +18,7 @@ package io.emeraldpay.dshackle.upstream.ethereum
import com.fasterxml.jackson.databind.ObjectMapper
import io.emeraldpay.dshackle.upstream.Upstream
import io.infinitape.etherjar.rpc.*
import org.slf4j.LoggerFactory
import reactor.core.publisher.Mono
import java.io.InputStream
@@ -25,6 +26,10 @@ abstract class EthereumApi(
objectMapper: ObjectMapper
) {
companion object {
private val log = LoggerFactory.getLogger(EthereumApi::class.java)
}
private val jacksonRpcConverter = JacksonRpcConverter(objectMapper)
var upstream: Upstream? = null
@@ -40,5 +45,6 @@ abstract class EthereumApi(
return execute(0, rpcCall.method, rpcCall.params as List<Any>)
.flatMap(convertToJS)
.map(rpcCall.converter::apply)
.doOnError { err -> log.debug("Failed to read from upstream", err) }
}
}

View File

@@ -18,35 +18,44 @@ package io.emeraldpay.dshackle.upstream.ethereum
import io.emeraldpay.dshackle.Defaults
import io.infinitape.etherjar.rpc.Batch
import io.infinitape.etherjar.rpc.Commands
import io.infinitape.etherjar.rpc.ReactorBatch
import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import org.springframework.scheduling.concurrent.CustomizableThreadFactory
import reactor.core.Disposable
import reactor.core.publisher.Flux
import reactor.core.publisher.Mono
import reactor.core.scheduler.Schedulers
import java.time.Duration
import java.util.concurrent.Executors
class EthereumRpcHead(
private val api: DirectEthereumApi,
private val interval: Duration = Duration.ofSeconds(10)
): DefaultEthereumHead(), Lifecycle {
companion object {
val scheduler = Schedulers.fromExecutor(Executors.newCachedThreadPool(CustomizableThreadFactory("ethereum-rpc-head")))
}
private val log = LoggerFactory.getLogger(EthereumRpcHead::class.java)
private var refreshSubscription: Disposable? = null
override fun start() {
val base = Flux.interval(interval)
.publishOn(scheduler)
.flatMap {
val batch = Batch()
val f = batch.add(Commands.eth().blockNumber)
api.rpcClient.execute(batch)
Mono.fromCompletionStage(f).timeout(Defaults.timeout, Mono.empty())
api.rpcClient
.execute(Commands.eth().blockNumber)
.subscribeOn(scheduler)
.timeout(Defaults.timeout, Mono.error(Exception("Block number not received")))
}
.flatMap {
val batch = Batch()
val f = batch.add(Commands.eth().getBlock(it))
api.rpcClient.execute(batch)
Mono.fromCompletionStage(f).timeout(Defaults.timeout, Mono.empty())
api.rpcClient
.execute(Commands.eth().getBlock(it))
.subscribeOn(scheduler)
.timeout(Defaults.timeout, Mono.error(Exception("Block data not received")))
}
.onErrorContinue { err, _ ->
log.debug("RPC error ${err.message}")

View File

@@ -30,7 +30,7 @@ import io.emeraldpay.grpc.Chain
import io.infinitape.etherjar.domain.BlockHash
import io.infinitape.etherjar.domain.TransactionId
import io.infinitape.etherjar.rpc.*
import io.infinitape.etherjar.rpc.emerald.EmeraldGrpcTransport
import io.infinitape.etherjar.rpc.emerald.ReactorEmeraldClient
import io.infinitape.etherjar.rpc.json.BlockJson
import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
@@ -38,7 +38,6 @@ import reactor.core.Disposable
import reactor.core.publisher.Flux
import reactor.core.publisher.Mono
import reactor.core.publisher.toMono
import java.lang.Exception
import java.math.BigInteger
import java.time.Duration
import java.util.*
@@ -50,9 +49,9 @@ import kotlin.collections.ArrayList
open class GrpcUpstream(
private val parentId: String,
private val chain: Chain,
private val client: ReactorBlockchainGrpc.ReactorBlockchainStub,
private val blockchainStub: ReactorBlockchainGrpc.ReactorBlockchainStub,
private val objectMapper: ObjectMapper,
private val grpcTransport: EmeraldGrpcTransport
private val rpcClient: ReactorEmeraldClient
): DefaultUpstream(), Lifecycle {
private var allLabels: Collection<UpstreamsConfig.Labels> = ArrayList<UpstreamsConfig.Labels>()
@@ -68,11 +67,10 @@ open class GrpcUpstream(
open fun createApi(matcher: Selector.Matcher): DirectEthereumApi {
val targets = this.getMethods()
val transport = Selector.extractLabels(matcher)?.let { selector ->
grpcTransport.copyWithSelector(selector.asProto())
} ?: grpcTransport
val rpcClient = DefaultRpcClient(transport)
return DirectEthereumApi(rpcClient, objectMapper, targets).let {
val client = Selector.extractLabels(matcher)?.let { selector ->
rpcClient.copyWithSelector(selector.asProto())
} ?: rpcClient
return DirectEthereumApi(client, objectMapper, targets).let {
it.upstream = this
it
}
@@ -91,10 +89,10 @@ open class GrpcUpstream(
val retry: Function<Flux<BlockchainOuterClass.ChainHead>, Flux<BlockchainOuterClass.ChainHead>> = Function {
setStatus(UpstreamAvailability.UNAVAILABLE)
client.subscribeHead(chainRef)
blockchainStub.subscribeHead(chainRef)
}
val flux = client.subscribeHead(chainRef)
val flux = blockchainStub.subscribeHead(chainRef)
.compose(GrpcRetry.ManyToMany.retryAfter(retry, Duration.ofSeconds(5)))
observeHead(flux)
}
@@ -136,7 +134,7 @@ open class GrpcUpstream(
}
}
}.onErrorContinue { err, _ ->
log.error("Head subscription error: ${err.message}")
log.error("Head subscription error. ${err.javaClass.name}:${err.message}", err)
}.doOnNext {
setStatus(UpstreamAvailability.OK)
}

View File

@@ -26,7 +26,7 @@ import io.emeraldpay.dshackle.upstream.UpstreamChange
import io.emeraldpay.grpc.Chain
import io.grpc.ManagedChannelBuilder
import io.grpc.netty.NettyChannelBuilder
import io.infinitape.etherjar.rpc.emerald.EmeraldGrpcTransport
import io.infinitape.etherjar.rpc.emerald.ReactorEmeraldClient
import io.netty.handler.ssl.*
import org.apache.commons.lang3.StringUtils
import org.apache.commons.lang3.exception.ExceptionUtils
@@ -56,7 +56,7 @@ class GrpcUpstreams(
private var client: ReactorBlockchainGrpc.ReactorBlockchainStub? = null
private val known = HashMap<Chain, GrpcUpstream>()
private val lock = ReentrantLock()
private var grpcTransport: EmeraldGrpcTransport? = null
private var grpcTransport: ReactorEmeraldClient? = null
fun start(): Flux<UpstreamChange> {
val channel: ManagedChannelBuilder<*> = if (auth != null && StringUtils.isNotEmpty(auth.ca)) {
@@ -73,12 +73,9 @@ class GrpcUpstreams(
val client = ReactorBlockchainGrpc.newReactorStub(channel.build())
this.client = client
var i = 0
val grpcExecutor = Executors.newCachedThreadPool { r -> Thread(r, "grpc-up-$id-${i++}") };
this.grpcTransport = EmeraldGrpcTransport.newBuilder()
this.grpcTransport = ReactorEmeraldClient.newBuilder()
.forChannel(client.channel)
.setObjectMapper(objectMapper)
.setExecutorService(grpcExecutor)
.build()
val statusSubscription = AtomicReference<Disposable>()
@@ -109,8 +106,6 @@ class GrpcUpstreams(
prev?.dispose()
subscription
}
}.doFinally {
grpcExecutor.shutdown()
}
return updates