problem: doesn't connect to Bitcoin networks over gRPC

This commit is contained in:
Igor Artamonov
2020-08-31 21:25:43 -04:00
parent 813f2ae4a2
commit ff3fc95711
22 changed files with 899 additions and 218 deletions

View File

@@ -23,5 +23,6 @@ class Defaults {
companion object {
val timeout: Duration = Duration.ofSeconds(60)
val timeoutInternal: Duration = timeout.dividedBy(4)
val retryConnection: Duration = Duration.ofSeconds(10)
}
}

View File

@@ -55,7 +55,7 @@ class BlockContainer(
}
@JvmStatic
fun from(raw: ByteArray): BlockContainer {
fun fromEthereumJson(raw: ByteArray): BlockContainer {
val block = Global.objectMapper.readValue(raw, BlockJson::class.java)
return from(block, raw)
}

View File

@@ -16,19 +16,16 @@
*/
package io.emeraldpay.dshackle.startup
import com.fasterxml.jackson.databind.ObjectMapper
import io.emeraldpay.dshackle.BlockchainType
import io.emeraldpay.dshackle.FileResolver
import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.cache.CachesFactory
import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.CurrentMultistreamHolder
import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinUpstream
import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinRpcUpstream
import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.dshackle.upstream.calls.ManagedCallMethods
import io.emeraldpay.dshackle.upstream.ethereum.EthereumRpcUpstream
import io.emeraldpay.dshackle.upstream.ethereum.EthereumUpstream
import io.emeraldpay.dshackle.upstream.ethereum.EthereumWsFactory
import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstreams
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcHttpClient
@@ -146,7 +143,7 @@ open class ConfiguredUpstreams(
}
val methods = buildMethods(config, chain)
val upstream = BitcoinUpstream(config.id
val upstream = BitcoinRpcUpstream(config.id
?: "bitcoin-${seq.getAndIncrement()}", chain, directApi,
options, config.role,
QuorumForLabels.QuorumItem(1, config.labels),

View File

@@ -58,6 +58,20 @@ class QuorumForLabels() {
return Collections.unmodifiableList(nodes)
}
override fun equals(other: Any?): Boolean {
if (this === other) return true
if (other !is QuorumForLabels) return false
if (nodes != other.nodes) return false
return true
}
override fun hashCode(): Int {
return nodes.hashCode()
}
/**
* Details for a single element (upstream, node or aggregation)
*/
@@ -67,6 +81,24 @@ class QuorumForLabels() {
return QuorumItem(0, UpstreamsConfig.Labels())
}
}
override fun equals(other: Any?): Boolean {
if (this === other) return true
if (other !is QuorumItem) return false
if (quorum != other.quorum) return false
if (labels != other.labels) return false
return true
}
override fun hashCode(): Int {
var result = quorum
result = 31 * result + labels.hashCode()
return result
}
}
}

View File

@@ -16,14 +16,12 @@
*/
package io.emeraldpay.dshackle.upstream
import com.fasterxml.jackson.databind.ObjectMapper
import io.emeraldpay.dshackle.BlockchainType
import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.cache.CachesEnabled
import io.emeraldpay.dshackle.cache.CachesFactory
import io.emeraldpay.dshackle.quorum.QuorumReaderFactory
import io.emeraldpay.dshackle.startup.UpstreamChange
import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinMultistream
import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinRpcUpstream
import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinUpstream
import io.emeraldpay.dshackle.upstream.calls.DefaultBitcoinMethods
import io.emeraldpay.dshackle.upstream.calls.CallMethods

View File

@@ -16,6 +16,7 @@
*/
package io.emeraldpay.dshackle.upstream
import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.startup.QuorumForLabels
import io.emeraldpay.dshackle.upstream.calls.CallMethods
@@ -46,6 +47,14 @@ abstract class DefaultUpstream(
return getStatus() == UpstreamAvailability.OK
}
fun onStatus(value: BlockchainOuterClass.ChainStatus) {
val available = value.availability
val quorum = value.quorum
setStatus(
if (available != null) UpstreamAvailability.fromGrpc(available.number) else UpstreamAvailability.UNAVAILABLE
)
}
override fun getStatus(): UpstreamAvailability {
return status.get().status
}

View File

@@ -0,0 +1,108 @@
/**
* Copyright (c) 2020 EmeraldPay, Inc
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package io.emeraldpay.dshackle.upstream.bitcoin
import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.startup.QuorumForLabels
import io.emeraldpay.dshackle.upstream.*
import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import io.emeraldpay.grpc.Chain
import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import reactor.core.Disposable
open class BitcoinRpcUpstream(
id: String,
chain: Chain,
private val directApi: Reader<JsonRpcRequest, JsonRpcResponse>,
options: UpstreamsConfig.Options,
role: UpstreamsConfig.UpstreamRole,
node: QuorumForLabels.QuorumItem,
callMethods: CallMethods
) : BitcoinUpstream(id, chain, options, role, callMethods, node), Lifecycle {
companion object {
private val log = LoggerFactory.getLogger(BitcoinRpcUpstream::class.java)
}
private val head: Head = createHead()
private var validatorSubscription: Disposable? = null
private fun createHead(): Head {
return BitcoinRpcHead(
directApi,
ExtractBlock()
)
}
override fun getHead(): Head {
return head
}
override fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse> {
return directApi
}
override fun getLabels(): Collection<UpstreamsConfig.Labels> {
return listOf(UpstreamsConfig.Labels())
}
override fun <T : Upstream> cast(selfType: Class<T>): T {
if (!selfType.isAssignableFrom(this.javaClass)) {
throw ClassCastException("Cannot cast ${this.javaClass} to $selfType")
}
return this as T
}
override fun isRunning(): Boolean {
var runningAny = validatorSubscription != null
if (head is Lifecycle) {
runningAny = runningAny || head.isRunning
}
return runningAny
}
override fun start() {
log.info("Configured for ${chain.chainName}")
if (head is Lifecycle) {
if (!head.isRunning) {
head.start()
}
}
validatorSubscription?.dispose()
if (getOptions().disableValidation != null && getOptions().disableValidation!!) {
this.setLag(0)
this.setStatus(UpstreamAvailability.OK)
} else {
val validator = BitcoinUpstreamValidator(directApi, getOptions())
validatorSubscription = validator.start()
.subscribe(this::setStatus)
}
}
override fun stop() {
if (head is Lifecycle) {
head.stop()
}
validatorSubscription?.dispose()
}
}

View File

@@ -15,97 +15,30 @@
*/
package io.emeraldpay.dshackle.upstream.bitcoin
import com.fasterxml.jackson.databind.ObjectMapper
import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.startup.QuorumForLabels
import io.emeraldpay.dshackle.upstream.*
import io.emeraldpay.dshackle.upstream.DefaultUpstream
import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import io.emeraldpay.dshackle.upstream.calls.DefaultBitcoinMethods
import io.emeraldpay.grpc.Chain
import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import reactor.core.Disposable
import reactor.core.publisher.Mono
open class BitcoinUpstream(
abstract class BitcoinUpstream(
id: String,
val chain: Chain,
private val directApi: Reader<JsonRpcRequest, JsonRpcResponse>,
options: UpstreamsConfig.Options,
role: UpstreamsConfig.UpstreamRole,
node: QuorumForLabels.QuorumItem,
callMethods: CallMethods
) : DefaultUpstream(id, options, role, callMethods, node), Lifecycle {
callMethods: CallMethods,
node: QuorumForLabels.QuorumItem
) : DefaultUpstream(id, options, role, callMethods, node) {
constructor(id: String,
chain: Chain,
options: UpstreamsConfig.Options,
role: UpstreamsConfig.UpstreamRole) : this(id, chain, options, role, DefaultBitcoinMethods(), QuorumForLabels.QuorumItem.empty())
companion object {
private val log = LoggerFactory.getLogger(BitcoinUpstream::class.java)
}
private val head: Head = createHead()
private var validatorSubscription: Disposable? = null
private fun createHead(): Head {
return BitcoinRpcHead(
directApi,
ExtractBlock()
)
}
override fun getHead(): Head {
return head
}
override fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse> {
return directApi
}
override fun getLabels(): Collection<UpstreamsConfig.Labels> {
return listOf(UpstreamsConfig.Labels())
}
override fun <T : Upstream> cast(selfType: Class<T>): T {
if (!selfType.isAssignableFrom(this.javaClass)) {
throw ClassCastException("Cannot cast ${this.javaClass} to $selfType")
}
return this as T
}
override fun isRunning(): Boolean {
var runningAny = validatorSubscription != null
if (head is Lifecycle) {
runningAny = runningAny || head.isRunning
}
runningAny = runningAny
return runningAny
}
override fun start() {
log.info("Configured for ${chain.chainName}")
if (head is Lifecycle) {
if (!head.isRunning) {
head.start()
}
}
validatorSubscription?.dispose()
if (getOptions().disableValidation != null && getOptions().disableValidation!!) {
this.setLag(0)
this.setStatus(UpstreamAvailability.OK)
} else {
val validator = BitcoinUpstreamValidator(directApi, getOptions())
validatorSubscription = validator.start()
.subscribe(this::setStatus)
}
}
override fun stop() {
if (head is Lifecycle) {
head.stop()
}
validatorSubscription?.dispose()
}
}

View File

@@ -16,14 +16,12 @@
*/
package io.emeraldpay.dshackle.upstream.ethereum
import com.fasterxml.jackson.databind.ObjectMapper
import io.emeraldpay.dshackle.Defaults
import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import io.infinitape.etherjar.hex.HexQuantity
import io.infinitape.etherjar.rpc.Commands
import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import org.springframework.scheduling.concurrent.CustomizableThreadFactory
@@ -71,7 +69,7 @@ class EthereumRpcHead(
.timeout(Defaults.timeout, Mono.error(Exception("Block data not received")))
}
.map {
BlockContainer.from(it.getResult())
BlockContainer.fromEthereumJson(it.getResult())
}
.onErrorContinue { err, _ ->
log.debug("RPC error ${err.message}")

View File

@@ -92,7 +92,7 @@ class EthereumWsFactory(
}
}
.flatMap(JsonRpcResponse::requireResult)
.map { BlockContainer.from(it) }
.map { BlockContainer.fromEthereumJson(it) }
}.repeatWhenEmpty { n ->
Repeat.times<Any>(5)
.exponentialBackoff(Duration.ofMillis(50), Duration.ofMillis(500))

View File

@@ -0,0 +1,132 @@
/**
* Copyright (c) 2020 EmeraldPay, Inc
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package io.emeraldpay.dshackle.upstream.grpc
import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.api.proto.ReactorBlockchainGrpc
import io.emeraldpay.dshackle.Defaults
import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.data.BlockId
import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.UpstreamAvailability
import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinUpstream
import io.emeraldpay.dshackle.upstream.bitcoin.ExtractBlock
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcGrpcClient
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import io.emeraldpay.grpc.Chain
import io.infinitape.etherjar.rpc.RpcException
import org.reactivestreams.Publisher
import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import reactor.core.publisher.Mono
import java.math.BigInteger
import java.time.Instant
import java.util.concurrent.TimeoutException
import java.util.function.Function
class BitcoinGrpcUpstream(
private val parentId: String,
chain: Chain,
private val remote: ReactorBlockchainGrpc.ReactorBlockchainStub,
private val client: JsonRpcGrpcClient
) : BitcoinUpstream(
"$parentId/${chain.chainCode}",
chain,
UpstreamsConfig.Options.getDefaults(),
UpstreamsConfig.UpstreamRole.STANDARD
), GrpcUpstream, Lifecycle {
companion object {
private val log = LoggerFactory.getLogger(BitcoinGrpcUpstream::class.java)
}
private val extractBlock = ExtractBlock()
private val defaultReader: Reader<JsonRpcRequest, JsonRpcResponse> = client.forSelector(Selector.empty)
private val blockConverter: Function<BlockchainOuterClass.ChainHead, BlockContainer> = Function { value ->
val block = BlockContainer(
value.height,
BlockId.from(value.blockId),
BigInteger(1, value.weight.toByteArray()),
Instant.ofEpochMilli(value.timestamp),
false,
null,
null
)
block
}
private val reloadBlock: Function<BlockContainer, Publisher<BlockContainer>> = Function { existingBlock ->
// head comes without transaction data
// need to download transactions for the block
defaultReader.read(JsonRpcRequest("getblock", listOf(existingBlock.hash.toHex())))
.flatMap(JsonRpcResponse::requireResult)
.map(extractBlock::extract)
.timeout(timeout, Mono.error(TimeoutException("Timeout from upstream")))
.doOnError { t ->
setStatus(UpstreamAvailability.UNAVAILABLE)
val msg = "Failed to download block data for chain $chain on $parentId"
if (t is RpcException || t is TimeoutException) {
log.warn("$msg. Message: ${t.message}")
} else {
log.error(msg, t)
}
}
}
private val upstreamStatus = GrpcUpstreamStatus()
private val grpcHead = GrpcHead(chain, this, blockConverter, reloadBlock)
var timeout = Defaults.timeout
override fun getHead(): Head {
return grpcHead
}
override fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse> {
return defaultReader
}
override fun getLabels(): Collection<UpstreamsConfig.Labels> {
return upstreamStatus.getLabels()
}
override fun <T : Upstream> cast(selfType: Class<T>): T {
if (!selfType.isAssignableFrom(this.javaClass)) {
throw ClassCastException("Cannot cast ${this.javaClass} to $selfType")
}
return this as T
}
override fun isRunning(): Boolean {
return grpcHead.isRunning
}
override fun start() {
grpcHead.start(remote)
}
override fun stop() {
grpcHead.stop()
}
override fun update(conf: BlockchainOuterClass.DescribeChain) {
upstreamStatus.update(conf)
conf.status?.let { status -> onStatus(status) }
}
}

View File

@@ -16,9 +16,7 @@
*/
package io.emeraldpay.dshackle.upstream.grpc
import com.salesforce.reactorgrpc.GrpcRetry
import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.api.proto.Common
import io.emeraldpay.api.proto.ReactorBlockchainGrpc
import io.emeraldpay.dshackle.Defaults
import io.emeraldpay.dshackle.config.UpstreamsConfig
@@ -28,7 +26,6 @@ import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.startup.QuorumForLabels
import io.emeraldpay.dshackle.upstream.*
import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.dshackle.upstream.calls.DirectCallMethods
import io.emeraldpay.dshackle.upstream.ethereum.*
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcGrpcClient
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
@@ -36,167 +33,107 @@ import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import io.emeraldpay.grpc.Chain
import io.infinitape.etherjar.domain.BlockHash
import io.infinitape.etherjar.rpc.*
import org.reactivestreams.Publisher
import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import reactor.core.Disposable
import reactor.core.publisher.Flux
import reactor.core.publisher.Mono
import reactor.core.publisher.toMono
import java.math.BigInteger
import java.time.Duration
import java.time.Instant
import java.util.*
import java.util.concurrent.TimeoutException
import java.util.concurrent.atomic.AtomicReference
import java.util.function.Function
import kotlin.collections.ArrayList
open class EthereumGrpcUpstream(
private val parentId: String,
private val chain: Chain,
private val blockchainStub: ReactorBlockchainGrpc.ReactorBlockchainStub,
private val remote: ReactorBlockchainGrpc.ReactorBlockchainStub,
private val client: JsonRpcGrpcClient
) : EthereumUpstream(
"$parentId/${chain.chainCode}",
UpstreamsConfig.Options.getDefaults(),
UpstreamsConfig.UpstreamRole.STANDARD,
null, null
), Lifecycle {
), GrpcUpstream, Lifecycle {
private val blockConverter: Function<BlockchainOuterClass.ChainHead, BlockContainer> = Function { value ->
val block = BlockContainer(
value.height,
BlockId.from(BlockHash.from("0x" + value.blockId)),
BigInteger(1, value.weight.toByteArray()),
Instant.ofEpochMilli(value.timestamp),
false,
null,
null
)
block
}
private val reloadBlock: Function<BlockContainer, Publisher<BlockContainer>> = Function { existingBlock ->
// head comes without transaction data
// need to download transactions for the block
defaultReader.read(JsonRpcRequest("eth_getBlockByHash", listOf(existingBlock.hash.toHexWithPrefix(), false)))
.flatMap(JsonRpcResponse::requireResult)
.map {
BlockContainer.fromEthereumJson(it)
}
.timeout(timeout, Mono.error(TimeoutException("Timeout from upstream")))
.doOnError { t ->
setStatus(UpstreamAvailability.UNAVAILABLE)
val msg = "Failed to download block data for chain $chain on $parentId"
if (t is RpcException || t is TimeoutException) {
log.warn("$msg. Message: ${t.message}")
} else {
log.error(msg, t)
}
}
}
private var allLabels: Collection<UpstreamsConfig.Labels> = ArrayList<UpstreamsConfig.Labels>()
private val log = LoggerFactory.getLogger(EthereumGrpcUpstream::class.java)
private val nodes = AtomicReference<QuorumForLabels>(QuorumForLabels())
private val head = DefaultEthereumHead()
private var targets: CallMethods? = null
private var headSubscription: Disposable? = null
var timeout = Defaults.timeout
private val upstreamStatus = GrpcUpstreamStatus()
private val grpcHead = GrpcHead(chain, this, blockConverter, reloadBlock)
private val defaultReader: Reader<JsonRpcRequest, JsonRpcResponse> = client.forSelector(Selector.empty)
var timeout = Defaults.timeout
override fun start() {
if (this.isRunning) return
val chainRef = Common.Chain.newBuilder()
.setTypeValue(chain.id)
.build()
.toMono()
val retry: Function<Flux<BlockchainOuterClass.ChainHead>, Flux<BlockchainOuterClass.ChainHead>> = Function {
setStatus(UpstreamAvailability.UNAVAILABLE)
blockchainStub.subscribeHead(chainRef)
}
val flux = blockchainStub.subscribeHead(chainRef)
.compose(GrpcRetry.ManyToMany.retryAfter(retry, Duration.ofSeconds(5)))
observeHead(flux)
grpcHead.start(remote)
}
override fun isRunning(): Boolean {
return headSubscription != null
return grpcHead.isRunning
}
override fun stop() {
headSubscription?.dispose()
headSubscription = null
grpcHead.stop()
}
internal fun observeHead(flux: Flux<BlockchainOuterClass.ChainHead>) {
val base = flux.map { value ->
val block = BlockContainer(
value.height,
BlockId.from(BlockHash.from("0x" + value.blockId)),
BigInteger(1, value.weight.toByteArray()),
Instant.ofEpochMilli(value.timestamp),
false,
null,
null
)
block
}.distinctUntilChanged {
it.hash
}.filter { block ->
val curr = head.getCurrent()
curr == null || curr.difficulty < block.difficulty
}.flatMap {
defaultReader.read(JsonRpcRequest("eth_getBlockByHash", listOf(it.hash.toHexWithPrefix(), false)))
.flatMap(JsonRpcResponse::requireResult)
.map {
BlockContainer.from(it)
}
.timeout(timeout, Mono.error(TimeoutException("Timeout from upstream")))
.doOnError { t ->
setStatus(UpstreamAvailability.UNAVAILABLE)
val msg = "Failed to download block data for chain $chain on $parentId"
if (t is RpcException || t is TimeoutException) {
log.warn("$msg. Message: ${t.message}")
} else {
log.error(msg, t)
}
}
}.onErrorContinue { err, _ ->
log.error("Head subscription error. ${err.javaClass.name}:${err.message}", err)
}.doOnNext {
setStatus(UpstreamAvailability.OK)
}
headSubscription = head.follow(base)
}
fun init(conf: BlockchainOuterClass.DescribeChain) {
targets = DirectCallMethods(conf.supportedMethodsList.toSet())
val nodes = QuorumForLabels()
val allLabels = ArrayList<UpstreamsConfig.Labels>()
conf.nodesList.forEach { remoteNode ->
val node = QuorumForLabels.QuorumItem(remoteNode.quorum,
remoteNode.labelsList.let { provided ->
val labels = UpstreamsConfig.Labels()
provided.forEach {
labels[it.name] = it.value
}
allLabels.add(labels)
labels
}
)
nodes.add(node)
}
this.nodes.set(nodes)
this.allLabels = Collections.unmodifiableCollection(allLabels)
override fun update(conf: BlockchainOuterClass.DescribeChain) {
upstreamStatus.update(conf)
conf.status?.let { status -> onStatus(status) }
}
fun onStatus(value: BlockchainOuterClass.ChainStatus) {
val available = value.availability
val quorum = value.quorum
setStatus(
if (available != null) UpstreamAvailability.fromGrpc(available.number) else UpstreamAvailability.UNAVAILABLE
)
}
override fun getQuorumByLabel(): QuorumForLabels {
return nodes.get()
return upstreamStatus.getNodes()
}
// ------------------------------------------------------------------------------------------
override fun getLabels(): Collection<UpstreamsConfig.Labels> {
return allLabels
return upstreamStatus.getLabels()
}
override fun getMethods(): CallMethods {
return targets ?: throw IllegalStateException("Upstream is not initialized yet")
return upstreamStatus.getCallMethods()
}
override fun isAvailable(): Boolean {
return super.isAvailable() && head.getCurrent() != null && nodes.get().getAll().any {
return super.isAvailable() && grpcHead.getCurrent() != null && getQuorumByLabel().getAll().any {
it.quorum > 0
}
}
override fun getHead(): Head {
return head
return grpcHead
}
override fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse> {

View File

@@ -0,0 +1,121 @@
/**
* Copyright (c) 2020 EmeraldPay, Inc
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package io.emeraldpay.dshackle.upstream.grpc
import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.api.proto.Common
import io.emeraldpay.api.proto.ReactorBlockchainGrpc
import io.emeraldpay.dshackle.Defaults
import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.upstream.AbstractHead
import io.emeraldpay.dshackle.upstream.DefaultUpstream
import io.emeraldpay.dshackle.upstream.UpstreamAvailability
import io.emeraldpay.grpc.Chain
import org.reactivestreams.Publisher
import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import reactor.core.Disposable
import reactor.core.publisher.Flux
import reactor.core.publisher.Mono
import reactor.util.retry.Retry
import java.time.Duration
import java.util.function.Function
class GrpcHead(
private val chain: Chain,
private val parent: DefaultUpstream,
/**
* Converted from remote head details to the block container, which could be partial at this point
*/
private val converter: Function<BlockchainOuterClass.ChainHead, BlockContainer>,
/**
* Populate block data with all missing details, of any
*/
private val enhancer: Function<BlockContainer, Publisher<BlockContainer>>?
) : AbstractHead(), Lifecycle {
companion object {
private val log = LoggerFactory.getLogger(GrpcHead::class.java)
}
private var headSubscription: Disposable? = null
/**
* Initiate a new head subscription with connection to the remote
*/
fun start(remote: ReactorBlockchainGrpc.ReactorBlockchainStub) {
if (this.isRunning) {
stop()
}
val source = Flux.concat(
// first connect immediately
Flux.just(remote),
// following requests do with delay, give it a time to recover
Flux.just(remote).repeat().delayElements(Defaults.retryConnection)
).flatMap(this::subscribeHead)
start(source)
}
fun subscribeHead(client: ReactorBlockchainGrpc.ReactorBlockchainStub): Publisher<BlockchainOuterClass.ChainHead> {
val chainRef = Common.Chain.newBuilder()
.setTypeValue(chain.id)
.build()
return client.subscribeHead(chainRef)
// simple retry on failure, if eventually failed then it supposed to resubscribe later from outer method
.retryWhen(Retry.backoff(4, Duration.ofSeconds(1)))
.onErrorContinue { err, _ ->
log.warn("Disconnected $chain from ${parent.getId()}: ${err.message}")
parent.setStatus(UpstreamAvailability.UNAVAILABLE)
Mono.empty<BlockchainOuterClass.ChainHead>()
}
}
/**
* Initiate a new head from provided source of head details
*/
fun start(source: Flux<BlockchainOuterClass.ChainHead>) {
var blocks = source.map(converter)
.distinctUntilChanged {
it.hash
}.filter { block ->
val curr = this.getCurrent()
curr == null || curr.difficulty < block.difficulty
}
if (enhancer != null) {
blocks = blocks.flatMap(enhancer)
}
blocks = blocks.onErrorContinue { err, _ ->
log.error("Head subscription error. ${err.javaClass.name}:${err.message}", err)
}
headSubscription = super.follow(blocks)
}
override fun isRunning(): Boolean {
return !(headSubscription?.isDisposed ?: true)
}
override fun start() {
log.error("Use start with provides source")
}
override fun stop() {
headSubscription?.dispose()
}
}

View File

@@ -0,0 +1,28 @@
/**
* Copyright (c) 2020 EmeraldPay, Inc
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package io.emeraldpay.dshackle.upstream.grpc
import io.emeraldpay.api.proto.BlockchainOuterClass
interface GrpcUpstream {
/**
* Update the configuration of the upstream with the new data.
* Called on the first creation, and each time a new state received from upstream
*/
fun update(conf: BlockchainOuterClass.DescribeChain)
}

View File

@@ -0,0 +1,72 @@
/**
* Copyright (c) 2020 EmeraldPay, Inc
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package io.emeraldpay.dshackle.upstream.grpc
import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.startup.QuorumForLabels
import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.dshackle.upstream.calls.DirectCallMethods
import org.slf4j.LoggerFactory
import java.util.*
import java.util.concurrent.atomic.AtomicReference
import kotlin.collections.ArrayList
class GrpcUpstreamStatus {
companion object {
private val log = LoggerFactory.getLogger(GrpcUpstreamStatus::class.java)
}
private val allLabels: AtomicReference<Collection<UpstreamsConfig.Labels>> = AtomicReference(emptyList())
private val nodes = AtomicReference<QuorumForLabels>(QuorumForLabels())
private var targets: CallMethods? = null
fun update(conf: BlockchainOuterClass.DescribeChain) {
val updateLabels = ArrayList<UpstreamsConfig.Labels>()
val updateNodes = QuorumForLabels()
conf.nodesList.forEach { remoteNode ->
val node = QuorumForLabels.QuorumItem(remoteNode.quorum,
remoteNode.labelsList.let { provided ->
val labels = UpstreamsConfig.Labels()
provided.forEach {
labels[it.name] = it.value
}
updateLabels.add(labels)
labels
}
)
updateNodes.add(node)
}
this.nodes.set(updateNodes)
this.allLabels.set(Collections.unmodifiableCollection(updateLabels))
this.targets = DirectCallMethods(conf.supportedMethodsList.toSet())
}
fun getLabels(): Collection<UpstreamsConfig.Labels> {
return allLabels.get()
}
fun getNodes(): QuorumForLabels {
return nodes.get()
}
fun getCallMethods(): CallMethods {
return targets ?: throw IllegalStateException("Upstream is not initialized yet")
}
}

View File

@@ -16,7 +16,6 @@
*/
package io.emeraldpay.dshackle.upstream.grpc
import com.fasterxml.jackson.databind.ObjectMapper
import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.api.proto.ReactorBlockchainGrpc
import io.emeraldpay.dshackle.BlockchainType
@@ -25,6 +24,7 @@ import io.emeraldpay.dshackle.FileResolver
import io.emeraldpay.dshackle.config.AuthConfig
import io.emeraldpay.dshackle.upstream.UpstreamAvailability
import io.emeraldpay.dshackle.startup.UpstreamChange
import io.emeraldpay.dshackle.upstream.DefaultUpstream
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcGrpcClient
import io.emeraldpay.grpc.Chain
import io.grpc.ManagedChannelBuilder
@@ -54,7 +54,7 @@ class GrpcUpstreams(
var timeout = Defaults.timeout
private var client: ReactorBlockchainGrpc.ReactorBlockchainStub? = null
private val known = HashMap<Chain, EthereumGrpcUpstream>()
private val known = HashMap<Chain, DefaultUpstream>()
private val lock = ReentrantLock()
fun start(): Flux<UpstreamChange> {
@@ -115,7 +115,7 @@ class GrpcUpstreams(
try {
val chain = Chain.byId(chainDetails.chain.number)
val up = getOrCreate(chain)
(up.upstream as EthereumGrpcUpstream).init(chainDetails)
(up.upstream as GrpcUpstream).update(chainDetails)
up
} catch (e: Throwable) {
log.warn("Skip unsupported upstream ${chainDetails.chain} on $id: ${e.message}")
@@ -157,9 +157,17 @@ class GrpcUpstreams(
}
fun getOrCreate(chain: Chain): UpstreamChange {
if (BlockchainType.fromBlockchain(chain) != BlockchainType.ETHEREUM) {
val blockchainType = BlockchainType.fromBlockchain(chain)
if (blockchainType == BlockchainType.ETHEREUM) {
return getOrCreateEthereum(chain)
} else if (blockchainType == BlockchainType.BITCOIN) {
return getOrCreateBitcoin(chain)
} else {
throw IllegalArgumentException("Unsupported blockchain: $chain")
}
}
fun getOrCreateEthereum(chain: Chain): UpstreamChange {
lock.withLock {
val current = known[chain]
return if (current == null) {
@@ -175,7 +183,23 @@ class GrpcUpstreams(
}
}
fun get(chain: Chain): EthereumGrpcUpstream {
fun getOrCreateBitcoin(chain: Chain): UpstreamChange {
lock.withLock {
val current = known[chain]
return if (current == null) {
val rpcClient = JsonRpcGrpcClient(client!!, chain)
val created = BitcoinGrpcUpstream(id, chain, client!!, rpcClient)
created.timeout = this.timeout
known[chain] = created
created.start()
UpstreamChange(chain, created, UpstreamChange.ChangeType.ADDED)
} else {
UpstreamChange(chain, current, UpstreamChange.ChangeType.REVALIDATED)
}
}
}
fun get(chain: Chain): DefaultUpstream {
return known[chain]!!
}
}