diff --git a/docs/04-upstream-config.adoc b/docs/04-upstream-config.adoc index eb6d3544..1a1279a9 100644 --- a/docs/04-upstream-config.adoc +++ b/docs/04-upstream-config.adoc @@ -51,6 +51,7 @@ cluster: min-peers: 2 upstreams: - id: us-nodes + node-id: 1 chain: auto connection: grpc: @@ -61,6 +62,7 @@ cluster: certificate: client.crt key: client.p8.key - id: infura-eth + node-id: 2 chain: ethereum role: fallback labels: @@ -80,6 +82,7 @@ cluster: username: ${INFURA_USER} password: ${INFURA_PASSWD} - id: ethereum-pos + node-id: 3 chain: ropsten connection: ethereum-pos: @@ -110,6 +113,12 @@ In the example above we have: ** label `[provider: infura]` is set for that particular upstream, which can be selected during a request.For example for some requests you may want to use nodes with that label only, i.e. _"send that tx to infura nodes only"_, or _"read only from archive node, with label [archive: true]"_ ** upstream validation (peers, sync status, etc) is disabled for that particular upstream +=== Nodes +`[node-id: 1]` is numeric node identifier defined in a range [1..255] and used to forward +`eth_getFilterChanges` request to the node where one of `eth_newFilter`, `eth_newBlockFilter` or `eth_newPendingTransactionFilter` methods was executed (because filter is s stateful method). + +*It's kindly recommended* to strictly associate _node-id_ parameter with a physical node and keep it during any configuration changes + === Roles and Fallback upstream By default, the Dshackle connects to each upstream in a Round-Robin basis, i.e. sequentially one by one. @@ -206,6 +215,10 @@ Dshackle currently supports - `eth_getUncleByBlockNumberAndIndex` - `eth_feeHistory` - `eth_getLogs` +- `eth_getFilterChanges` +- `eth_newFilter` +- `eth_newBlockFilter` +- `eth_newPendingTransactionFilter` .Plus following methods are answered directly by Dshackle - `net_version` diff --git a/src/main/kotlin/io/emeraldpay/dshackle/config/UpstreamsConfig.kt b/src/main/kotlin/io/emeraldpay/dshackle/config/UpstreamsConfig.kt index 168f2f30..c7ca9831 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/config/UpstreamsConfig.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/config/UpstreamsConfig.kt @@ -69,6 +69,7 @@ open class UpstreamsConfig { class Upstream { var id: String? = null + var nodeId: Int? = null var chain: String? = null var options: Options? = null var isEnabled = true diff --git a/src/main/kotlin/io/emeraldpay/dshackle/config/UpstreamsConfigReader.kt b/src/main/kotlin/io/emeraldpay/dshackle/config/UpstreamsConfigReader.kt index 91ea1d48..c5e5686f 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/config/UpstreamsConfigReader.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/config/UpstreamsConfigReader.kt @@ -31,6 +31,7 @@ class UpstreamsConfigReader( private val log = LoggerFactory.getLogger(UpstreamsConfigReader::class.java) private val authConfigReader = AuthConfigReader() + private val knownNodeIds: MutableSet = HashSet() fun read(input: InputStream): UpstreamsConfig? { val configNode = readNode(input) @@ -236,11 +237,22 @@ class UpstreamsConfigReader( log.warn("Invalid id: $id") return false } - return true + return upstream.nodeId?.let { + if (it !in 1..255) { + log.warn("Invalid node-id: $it. Must be in range [1, 255].") + false + } else if (!knownNodeIds.add(it)) { + log.warn("Duplicated node-id: $it. Must be in unique.") + false + } else { + true + } + } ?: true } internal fun readUpstreamCommon(upNode: MappingNode, upstream: UpstreamsConfig.Upstream<*>) { upstream.id = getValueAsString(upNode, "id") + upstream.nodeId = getValueAsInt(upNode, "node-id") upstream.options = tryReadOptions(upNode) upstream.methods = tryReadMethods(upNode) getValueAsBool(upNode, "enabled")?.let { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/quorum/AlwaysQuorum.kt b/src/main/kotlin/io/emeraldpay/dshackle/quorum/AlwaysQuorum.kt index a2b5d1dc..74ba4456 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/quorum/AlwaysQuorum.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/quorum/AlwaysQuorum.kt @@ -21,6 +21,7 @@ import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcError import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcException import io.emeraldpay.dshackle.upstream.signature.ResponseSigner +import java.util.concurrent.ConcurrentLinkedQueue open class AlwaysQuorum : CallQuorum { @@ -28,6 +29,7 @@ open class AlwaysQuorum : CallQuorum { private var result: ByteArray? = null private var rpcError: JsonRpcError? = null private var sig: ResponseSigner.Signature? = null + private val resolvers: MutableCollection = ConcurrentLinkedQueue() override fun init(head: Head) { } @@ -48,6 +50,7 @@ open class AlwaysQuorum : CallQuorum { result = response resolved = true sig = signature + resolvers.add(upstream) return true } @@ -64,6 +67,9 @@ open class AlwaysQuorum : CallQuorum { return rpcError } + override fun getResolvedBy(): List = + resolvers.toList() + override fun toString(): String { return "Quorum: Accept Any" } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/quorum/CallQuorum.kt b/src/main/kotlin/io/emeraldpay/dshackle/quorum/CallQuorum.kt index 58122631..e89f2920 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/quorum/CallQuorum.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/quorum/CallQuorum.kt @@ -34,4 +34,5 @@ interface CallQuorum { fun getSignature(): ResponseSigner.Signature? fun getResult(): ByteArray? fun getError(): JsonRpcError? + fun getResolvedBy(): Collection } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/quorum/NotLaggingQuorum.kt b/src/main/kotlin/io/emeraldpay/dshackle/quorum/NotLaggingQuorum.kt index 89b6bbf2..980f3c50 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/quorum/NotLaggingQuorum.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/quorum/NotLaggingQuorum.kt @@ -21,6 +21,7 @@ import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcError import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcException import io.emeraldpay.dshackle.upstream.signature.ResponseSigner +import java.util.concurrent.ConcurrentLinkedQueue import java.util.concurrent.atomic.AtomicReference /** @@ -34,6 +35,7 @@ class NotLaggingQuorum(val maxLag: Long = 0) : CallQuorum { private val failed = AtomicReference(false) private var rpcError: JsonRpcError? = null private var sig: ResponseSigner.Signature? = null + private val resolvers: MutableCollection = ConcurrentLinkedQueue() override fun init(head: Head) { } @@ -51,6 +53,7 @@ class NotLaggingQuorum(val maxLag: Long = 0) : CallQuorum { if (!lagging) { result.set(response) sig = signature + resolvers.add(upstream) return true } return false @@ -75,6 +78,8 @@ class NotLaggingQuorum(val maxLag: Long = 0) : CallQuorum { return rpcError } + override fun getResolvedBy(): Collection = + resolvers.toList() override fun toString(): String { return "Quorum: late <= $maxLag blocks" } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/quorum/QuorumRpcReader.kt b/src/main/kotlin/io/emeraldpay/dshackle/quorum/QuorumRpcReader.kt index 3b9c9694..cac1a03d 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/quorum/QuorumRpcReader.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/quorum/QuorumRpcReader.kt @@ -117,9 +117,9 @@ class QuorumRpcReader( return Function { quorumResult -> quorumResult .filter { it.isResolved() } // return nothing if not resolved - .map { + .map { quorum -> // TODO find actual quorum number - QuorumRpcReader.Result(it.getResult()!!, it.getSignature(), 1) + Result(quorum.getResult()!!, quorum.getSignature(), 1, quorum.getResolvedBy().map { it.nodeId() }) } .switchIfEmpty(defaultResult) } @@ -197,6 +197,7 @@ class QuorumRpcReader( class Result( val value: ByteArray, val signature: ResponseSigner.Signature?, - val quorum: Int + val quorum: Int, + val resolvers: Collection ) } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/quorum/ValueAwareQuorum.kt b/src/main/kotlin/io/emeraldpay/dshackle/quorum/ValueAwareQuorum.kt index 2e6233e6..bb8c5d3b 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/quorum/ValueAwareQuorum.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/quorum/ValueAwareQuorum.kt @@ -23,6 +23,7 @@ import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcException import io.emeraldpay.dshackle.upstream.signature.ResponseSigner import io.emeraldpay.etherjar.rpc.RpcException import org.slf4j.LoggerFactory +import java.util.concurrent.ConcurrentLinkedQueue abstract class ValueAwareQuorum( val clazz: Class @@ -30,6 +31,7 @@ abstract class ValueAwareQuorum( private val log = LoggerFactory.getLogger(ValueAwareQuorum::class.java) private var rpcError: JsonRpcError? = null + private val resolvers: MutableCollection = ConcurrentLinkedQueue() fun extractValue(response: ByteArray, clazz: Class): T? { return Global.objectMapper.readValue(response.inputStream(), clazz) @@ -39,6 +41,7 @@ abstract class ValueAwareQuorum( try { val value = extractValue(response, clazz) recordValue(response, value, signature, upstream) + resolvers.add(upstream) } catch (e: RpcException) { recordError(response, e.rpcMessage, signature, upstream) } catch (e: Exception) { @@ -59,4 +62,7 @@ abstract class ValueAwareQuorum( override fun getError(): JsonRpcError? { return rpcError } + + override fun getResolvedBy(): Collection = + resolvers.toList() } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt index 292bdbbb..78ec44ac 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt @@ -105,7 +105,8 @@ open class NativeCall( } fun parseParams(it: ValidCallContext): ValidCallContext { - val params = extractParams(it.payload.params) + val rawParams = extractParams(it.payload.params) + val params = it.requestDecorator.processRequest(rawParams) return it.withPayload(ParsedCallDetails(it.payload.method, params)) } @@ -237,17 +238,28 @@ open class NativeCall( matcher.withMatcher(heightMatcher) } val nonce = requestItem.nonce.let { if (it == 0L) null else it } + val requestDecorator = getRequestDecorator(requestItem.method) + val resultDecorator = getResultDecorator(requestItem.method) + ValidCallContext( requestItem.id, nonce, upstream, matcher.build(), callQuorum, - RawCallDetails(method, params) + RawCallDetails(method, params), + requestDecorator, + resultDecorator ) } } + private fun getRequestDecorator(method: String): RequestDecorator = + if (method == "eth_getFilterChanges") GetFilterUpdatesDecorator() else NoneRequestDecorator() + + private fun getResultDecorator(method: String): ResultDecorator = + if (CreateFilterDecorator.createFilterMethods.contains(method)) CreateFilterDecorator() else NoneResultDecorator() + fun fetch(ctx: ValidCallContext): Mono { return ctx.upstream.getRoutedApi(ctx.matcher) .flatMap { api -> @@ -286,7 +298,8 @@ open class NativeCall( return reader .read(JsonRpcRequest(ctx.payload.method, ctx.payload.params, ctx.nonce)) .map { - CallResult(ctx.id, ctx.nonce, it.value, null, it.signature) + val bytes = ctx.resultDecorator.processResult(it) + CallResult(ctx.id, ctx.nonce, bytes, null, it.signature) } .onErrorResume { t -> val failure = when (t) { @@ -353,14 +366,67 @@ open class NativeCall( fun getError(): CallError } + interface ResultDecorator { + fun processResult(result: QuorumRpcReader.Result): ByteArray + } + + open class NoneResultDecorator : ResultDecorator { + override fun processResult(result: QuorumRpcReader.Result): ByteArray = result.value + } + + open class CreateFilterDecorator : ResultDecorator { + + companion object { + const val quoteCode = '"'.code.toByte() + val createFilterMethods = listOf("eth_getFilterChanges", "eth_newFilter", "eth_newBlockFilter") + } + override fun processResult(result: QuorumRpcReader.Result): ByteArray { + val bytes = result.value + if (bytes.last() == quoteCode) { + val suffix = result.resolvers.first().toUByte().toString(16).padStart(2, padChar = '0').toByteArray() + bytes[bytes.lastIndex] = suffix.first() + return bytes + suffix.last() + quoteCode + } + return bytes + } + } + + interface RequestDecorator { + fun processRequest(request: List): List + } + + open class NoneRequestDecorator : RequestDecorator { + override fun processRequest(request: List): List = request + } + + open class GetFilterUpdatesDecorator : RequestDecorator { + override fun processRequest(request: List): List { + val filterId = request.first().toString() + val sanitized = filterId.substring(0, filterId.lastIndex - 1) + return listOf(sanitized) + } + } + open class ValidCallContext( val id: Int, val nonce: Long?, val upstream: Multistream, val matcher: Selector.Matcher, val callQuorum: CallQuorum, - val payload: T + val payload: T, + val requestDecorator: RequestDecorator, + val resultDecorator: ResultDecorator ) : CallContext { + + constructor( + id: Int, + nonce: Long?, + upstream: Multistream, + matcher: Selector.Matcher, + callQuorum: CallQuorum, + payload: T + ) : this(id, nonce, upstream, matcher, callQuorum, payload, NoneRequestDecorator(), NoneResultDecorator()) + override fun isValid(): Boolean { return true } @@ -374,7 +440,7 @@ open class NativeCall( } fun withPayload(payload: X): ValidCallContext { - return ValidCallContext(id, nonce, upstream, matcher, callQuorum, payload) + return ValidCallContext(id, nonce, upstream, matcher, callQuorum, payload, requestDecorator, resultDecorator) } fun getApis(): ApiSource { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt b/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt index bda626e2..3a86c6ac 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt @@ -16,6 +16,7 @@ */ package io.emeraldpay.dshackle.startup +import com.google.common.annotations.VisibleForTesting import io.emeraldpay.dshackle.FileResolver import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.config.UpstreamsConfig @@ -53,7 +54,9 @@ import org.springframework.beans.factory.annotation.Autowired import org.springframework.stereotype.Repository import java.net.URI import java.util.concurrent.atomic.AtomicInteger +import java.util.function.Function import javax.annotation.PostConstruct +import kotlin.math.abs @Repository open class ConfiguredUpstreams( @@ -65,6 +68,8 @@ open class ConfiguredUpstreams( private val log = LoggerFactory.getLogger(ConfiguredUpstreams::class.java) private var seq = AtomicInteger(0) + private val hashes: MutableMap = HashMap() + @PostConstruct fun start() { log.debug("Starting upstreams") @@ -77,7 +82,7 @@ open class ConfiguredUpstreams( log.debug("Start upstream ${up.id}") if (up.connection is UpstreamsConfig.GrpcConnection) { val options = up.options ?: UpstreamsConfig.Options() - buildGrpcUpstream(up.cast(UpstreamsConfig.GrpcConnection::class.java), options) + buildGrpcUpstream(up.nodeId, up.cast(UpstreamsConfig.GrpcConnection::class.java), options) } else { val chain = Global.chainById(up.chain) if (chain == Chain.UNSPECIFIED) { @@ -88,14 +93,22 @@ open class ConfiguredUpstreams( .merge(defaultOptions[chain] ?: UpstreamsConfig.Options.getDefaults()) val upstream = when (BlockchainType.from(chain)) { BlockchainType.EVM_POW -> { - buildEthereumUpstream(up.cast(UpstreamsConfig.EthereumConnection::class.java), chain, options) + buildEthereumUpstream(up.nodeId, up.cast(UpstreamsConfig.EthereumConnection::class.java), chain, options) } + BlockchainType.BITCOIN -> { buildBitcoinUpstream(up.cast(UpstreamsConfig.BitcoinConnection::class.java), chain, options) } + BlockchainType.EVM_POS -> { - buildEthereumPosUpstream(up.cast(UpstreamsConfig.EthereumPosConnection::class.java), chain, options) + buildEthereumPosUpstream( + up.nodeId, + up.cast(UpstreamsConfig.EthereumPosConnection::class.java), + chain, + options + ) } + else -> { log.error("Chain is unsupported: ${up.chain}") return@forEach @@ -156,6 +169,7 @@ open class ConfiguredUpstreams( } private fun buildEthereumPosUpstream( + nodeId: Int?, config: UpstreamsConfig.Upstream, chain: Chain, options: UpstreamsConfig.Options @@ -167,13 +181,26 @@ open class ConfiguredUpstreams( return null } val urls = ArrayList() - val connectorFactory = buildEthereumConnectorFactory(config.id!!, execution, chain, urls, NoChoiceWithPriorityForkChoice(conn.upstreamRating), BlockValidator.ALWAYS_VALID) + val connectorFactory = buildEthereumConnectorFactory( + config.id!!, + execution, + chain, + urls, + NoChoiceWithPriorityForkChoice(conn.upstreamRating), + BlockValidator.ALWAYS_VALID + ) val methods = buildMethods(config, chain) if (connectorFactory == null) { return null } + + val hashUrl = conn.execution!!.let { + if (it.preferHttp) it.rpc?.url ?: it.ws?.url else it.ws?.url ?: it.rpc?.url + } + val hash = getHash(nodeId, hashUrl!!) val upstream = EthereumPosRpcUpstream( config.id!!, + hash, chain, options, config.role, methods, @@ -227,6 +254,7 @@ open class ConfiguredUpstreams( } private fun buildEthereumUpstream( + nodeId: Int?, config: UpstreamsConfig.Upstream, chain: Chain, options: UpstreamsConfig.Options @@ -236,12 +264,22 @@ open class ConfiguredUpstreams( val urls = ArrayList() val methods = buildMethods(config, chain) - val connectorFactory = buildEthereumConnectorFactory(config.id!!, conn, chain, urls, MostWorkForkChoice(), EthereumBlockValidator()) + val connectorFactory = buildEthereumConnectorFactory( + config.id!!, + conn, + chain, + urls, + MostWorkForkChoice(), + EthereumBlockValidator() + ) if (connectorFactory == null) { return null } + + val hashUrl = if (conn.preferHttp) conn.rpc?.url ?: conn.ws?.url else conn.ws?.url ?: conn.rpc?.url val upstream = EthereumRpcUpstream( config.id!!, + getHash(nodeId, hashUrl!!), chain, options, config.role, methods, @@ -253,12 +291,15 @@ open class ConfiguredUpstreams( } private fun buildGrpcUpstream( + nodeId: Int?, config: UpstreamsConfig.Upstream, options: UpstreamsConfig.Options ) { val endpoint = config.connection!! + val hash = getHash(nodeId, "${endpoint.host}:${endpoint.port}") val ds = GrpcUpstreams( config.id!!, + hash, config.role, endpoint.host!!, endpoint.port, @@ -289,7 +330,12 @@ open class ConfiguredUpstreams( } } - private fun buildWsFactory(id: String, chain: Chain, conn: UpstreamsConfig.EthereumConnection, urls: ArrayList? = null): EthereumWsFactory? { + private fun buildWsFactory( + id: String, + chain: Chain, + conn: UpstreamsConfig.EthereumConnection, + urls: ArrayList? = null + ): EthereumWsFactory? { return conn.ws?.let { endpoint -> val wsApi = EthereumWsFactory( id, chain, @@ -305,15 +351,45 @@ open class ConfiguredUpstreams( } } - private fun buildEthereumConnectorFactory(id: String, conn: UpstreamsConfig.EthereumConnection, chain: Chain, urls: ArrayList, forkChoice: ForkChoice, blockValidator: BlockValidator): EthereumConnectorFactory? { + private fun buildEthereumConnectorFactory( + id: String, + conn: UpstreamsConfig.EthereumConnection, + chain: Chain, + urls: ArrayList, + forkChoice: ForkChoice, + blockValidator: BlockValidator + ): EthereumConnectorFactory? { val wsFactoryApi = buildWsFactory(id, chain, conn, urls) val httpFactory = buildHttpFactory(conn, urls) log.info("Using ${chain.chainName} upstream, at ${urls.joinToString()}") - val connectorFactory = EthereumConnectorFactory(conn.preferHttp, wsFactoryApi, httpFactory, forkChoice, blockValidator) + val connectorFactory = + EthereumConnectorFactory(conn.preferHttp, wsFactoryApi, httpFactory, forkChoice, blockValidator) if (!connectorFactory.isValid()) { log.warn("Upstream configuration is invalid (probably no http endpoint)") return null } return connectorFactory } + + @VisibleForTesting + private fun getHash(nodeId: Int?, obj: Any): Byte = + nodeId?.toByte() ?: (obj.hashCode() % 255).let { + if (it == 0) 1 else it + }.let { nonZeroHash -> + listOf>( + Function { i -> i }, + Function { i -> (-i) }, + Function { i -> 127 - abs(i) }, + Function { i -> abs(i) - 128 }, + ).map { + it.apply(nonZeroHash).toByte() + }.firstOrNull { + hashes[it] != true + }?.let { + hashes[it] = true + it + } ?: (Byte.MIN_VALUE..Byte.MAX_VALUE).first { + it != 0 && hashes[it.toByte()] != true + }.toByte() + } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/DefaultUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/DefaultUpstream.kt index 36470987..451a53b4 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/DefaultUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/DefaultUpstream.kt @@ -26,6 +26,7 @@ import java.util.concurrent.atomic.AtomicReference abstract class DefaultUpstream( private val id: String, + private val hash: Byte, defaultLag: Long, defaultAvail: UpstreamAvailability, private val options: UpstreamsConfig.Options, @@ -36,12 +37,14 @@ abstract class DefaultUpstream( constructor( id: String, + hash: Byte, options: UpstreamsConfig.Options, role: UpstreamsConfig.UpstreamRole, targets: CallMethods? ) : this( id, + hash, Long.MAX_VALUE, UpstreamAvailability.UNAVAILABLE, options, @@ -52,12 +55,13 @@ abstract class DefaultUpstream( constructor( id: String, + hash: Byte, options: UpstreamsConfig.Options, role: UpstreamsConfig.UpstreamRole, targets: CallMethods?, node: QuorumForLabels.QuorumItem? ) : - this(id, Long.MAX_VALUE, UpstreamAvailability.UNAVAILABLE, options, role, targets, node) + this(id, hash, Long.MAX_VALUE, UpstreamAvailability.UNAVAILABLE, options, role, targets, node) private val status = AtomicReference(Status(defaultLag, defaultAvail, statusByLag(defaultLag, defaultAvail))) private val statusStream = Sinks.many() @@ -139,6 +143,8 @@ abstract class DefaultUpstream( return targets ?: throw IllegalStateException("Methods are not set") } + override fun nodeId(): Byte = hash + private val quorumByLabel = node?.let { QuorumForLabels(it) } ?: QuorumForLabels(QuorumForLabels.QuorumItem.empty()) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt index 6319cede..c04dd9e0 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt @@ -282,6 +282,8 @@ abstract class Multistream( return false } + override fun nodeId(): Byte = 0 + fun printStatus() { var height: Long? = null try { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Selector.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Selector.kt index 4f57a7f3..6c105505 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Selector.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Selector.kt @@ -396,4 +396,22 @@ class Selector { return "Matcher: ${describeInternal()}" } } + + class SameNodeMatcher(private val upstreamHash: Byte) : Matcher { + override fun matches(up: Upstream): Boolean = + up.nodeId() == upstreamHash + + override fun describeInternal(): String = + "upstream node-id=${upstreamHash.toUByte()}" + + override fun toString(): String { + return "Matcher: ${describeInternal()}" + } + + override fun equals(other: Any?): Boolean { + if (other === this) return true + if (other !is SameNodeMatcher) return false + return other.upstreamHash == upstreamHash + } + } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Upstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Upstream.kt index d170acbc..e8de7910 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Upstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Upstream.kt @@ -40,4 +40,6 @@ interface Upstream { fun isGrpc(): Boolean fun cast(selfType: Class): T + + fun nodeId(): Byte } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinUpstream.kt index 6e02efff..b09709ff 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinUpstream.kt @@ -31,7 +31,7 @@ abstract class BitcoinUpstream( callMethods: CallMethods, node: QuorumForLabels.QuorumItem, val esploraClient: EsploraClient? = null -) : DefaultUpstream(id, options, role, callMethods, node) { +) : DefaultUpstream(id, 0.toByte(), options, role, callMethods, node) { constructor( id: String, diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/calls/DefaultEthereumMethods.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/calls/DefaultEthereumMethods.kt index f3a65f90..45067a48 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/calls/DefaultEthereumMethods.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/calls/DefaultEthereumMethods.kt @@ -70,7 +70,14 @@ class DefaultEthereumMethods( "eth_feeHistory" ) - private val allowedMethods = anyResponseMethods + firstValueMethods + specialMethods + headVerifiedMethods + private val filterMethods = listOf( + "eth_getFilterChanges", + "eth_newFilter", + "eth_newBlockFilter", + "eth_newPendingTransactionFilter", + ) + + private val allowedMethods = anyResponseMethods + firstValueMethods + specialMethods + headVerifiedMethods + filterMethods private val hardcodedMethods = listOf( "net_version", @@ -88,6 +95,7 @@ class DefaultEthereumMethods( override fun getQuorumFor(method: String): CallQuorum { return when { + filterMethods.contains(method) -> AlwaysQuorum() hardcodedMethods.contains(method) -> AlwaysQuorum() firstValueMethods.contains(method) -> AlwaysQuorum() anyResponseMethods.contains(method) -> NotLaggingQuorum(4) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/calls/EthereumCallSelector.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/calls/EthereumCallSelector.kt index 32b00dcb..79748acf 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/calls/EthereumCallSelector.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/calls/EthereumCallSelector.kt @@ -57,10 +57,26 @@ class EthereumCallSelector( return blockTagSelector(params, 1, head) } else if (method == "eth_getStorageAt") { return blockTagSelector(params, 2, head) + } else if (method == "eth_getFilterChanges") { + return sameUpstreamMatcher(params) } return Mono.empty() } + private fun sameUpstreamMatcher(params: String): Mono { + val list = objectMapper.readerFor(Any::class.java).readValues(params).readAll() + if (list.isEmpty()) { + return Mono.empty() + } + val filterId = list[0].toString() + if (filterId.length < 4) { + return Mono.just(Selector.SameNodeMatcher(0.toByte())) + } + val hashHex = filterId.substring(filterId.length - 2) + val nodeId = hashHex.toInt(16) + return Mono.just(Selector.SameNodeMatcher(nodeId.toByte())) + } + private fun blockTagSelector(params: String, pos: Int, head: Head): Mono { val list = objectMapper.readerFor(Any::class.java).readValues(params).readAll() if (list.size < pos + 1) { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumRpcUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumRpcUpstream.kt index b071af41..3bbe468f 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumRpcUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumRpcUpstream.kt @@ -36,13 +36,14 @@ import reactor.core.Disposable open class EthereumRpcUpstream( id: String, + hash: Byte, val chain: Chain, options: UpstreamsConfig.Options, role: UpstreamsConfig.UpstreamRole, targets: CallMethods?, private val node: QuorumForLabels.QuorumItem?, connectorFactory: ConnectorFactory -) : EthereumUpstream(id, options, role, targets, node), Lifecycle, Upstream, CachesEnabled { +) : EthereumUpstream(id, hash, options, role, targets, node), Lifecycle, Upstream, CachesEnabled { private val log = LoggerFactory.getLogger(EthereumRpcUpstream::class.java) private val validator: EthereumUpstreamValidator = EthereumUpstreamValidator(this, getOptions()) private val connector: EthereumConnector = connectorFactory.create(this, validator, chain) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumUpstream.kt index c14bd120..a06c7cb7 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumUpstream.kt @@ -24,11 +24,12 @@ import io.emeraldpay.dshackle.upstream.calls.CallMethods abstract class EthereumUpstream( id: String, + hash: Byte, options: UpstreamsConfig.Options, role: UpstreamsConfig.UpstreamRole, targets: CallMethods?, private val node: QuorumForLabels.QuorumItem? -) : DefaultUpstream(id, options, role, targets, node) { +) : DefaultUpstream(id, hash, options, role, targets, node) { private val capabilities = if (options.providesBalance != false) { setOf(Capability.RPC, Capability.BALANCE) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosRpcUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosRpcUpstream.kt index 7cd6cca6..50f201c2 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosRpcUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosRpcUpstream.kt @@ -36,13 +36,14 @@ import reactor.core.Disposable open class EthereumPosRpcUpstream( id: String, + hash: Byte, val chain: Chain, options: UpstreamsConfig.Options, role: UpstreamsConfig.UpstreamRole, targets: CallMethods?, private val node: QuorumForLabels.QuorumItem?, connectorFactory: ConnectorFactory -) : EthereumPosUpstream(id, options, role, targets, node), Lifecycle, Upstream, CachesEnabled { +) : EthereumPosUpstream(id, hash, options, role, targets, node), Lifecycle, Upstream, CachesEnabled { private val log = LoggerFactory.getLogger(EthereumPosRpcUpstream::class.java) private val validator: EthereumUpstreamValidator = EthereumUpstreamValidator(this, getOptions()) private val connector: EthereumConnector = connectorFactory.create(this, validator, chain) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosUpstream.kt index 4ce8253a..b33ac76c 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosUpstream.kt @@ -24,11 +24,12 @@ import io.emeraldpay.dshackle.upstream.calls.CallMethods abstract class EthereumPosUpstream( id: String, + hash: Byte, options: UpstreamsConfig.Options, role: UpstreamsConfig.UpstreamRole, targets: CallMethods?, private val node: QuorumForLabels.QuorumItem? -) : DefaultUpstream(id, options, role, targets, node) { +) : DefaultUpstream(id, hash, options, role, targets, node) { private val capabilities = if (options.providesBalance != false) { setOf(Capability.RPC, Capability.BALANCE) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcUpstream.kt index ab4084c5..d407d3b5 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcUpstream.kt @@ -51,6 +51,7 @@ import java.util.function.Function open class EthereumGrpcUpstream( private val parentId: String, + hash: Byte, role: UpstreamsConfig.UpstreamRole, private val chain: Chain, private val remote: ReactorBlockchainGrpc.ReactorBlockchainStub, @@ -58,6 +59,7 @@ open class EthereumGrpcUpstream( overrideLabels: UpstreamsConfig.Labels? ) : EthereumUpstream( "${parentId}_${chain.chainCode.lowercase(Locale.getDefault())}", + hash, UpstreamsConfig.Options.getDefaults(), role, null, diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumPosGrpcUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumPosGrpcUpstream.kt index 43acab2f..10a005cc 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumPosGrpcUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumPosGrpcUpstream.kt @@ -51,6 +51,7 @@ import java.util.function.Function open class EthereumPosGrpcUpstream( private val parentId: String, + hash: Byte, role: UpstreamsConfig.UpstreamRole, private val chain: Chain, private val remote: ReactorBlockchainGrpc.ReactorBlockchainStub, @@ -59,6 +60,7 @@ open class EthereumPosGrpcUpstream( overrideLabels: UpstreamsConfig.Labels? ) : EthereumPosUpstream( "${parentId}_${chain.chainCode.lowercase(Locale.getDefault())}", + hash, UpstreamsConfig.Options.getDefaults(), role, null, null diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreams.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreams.kt index db8774ad..c92fd624 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreams.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreams.kt @@ -52,6 +52,7 @@ import kotlin.concurrent.withLock class GrpcUpstreams( private val id: String, + private val hash: Byte, private val role: UpstreamsConfig.UpstreamRole, private val host: String, private val port: Int, @@ -206,7 +207,7 @@ class GrpcUpstreams( val current = known[chain] return if (current == null) { val rpcClient = JsonRpcGrpcClient(client!!, chain, metrics) - val created = EthereumGrpcUpstream(id, role, chain, client!!, rpcClient, labels) + val created = EthereumGrpcUpstream(id, hash, role, chain, client!!, rpcClient, labels) created.timeout = this.timeout known[chain] = created created.start() @@ -222,7 +223,7 @@ class GrpcUpstreams( val current = known[chain] return if (current == null) { val rpcClient = JsonRpcGrpcClient(client!!, chain, metrics) - val created = EthereumPosGrpcUpstream(id, role, chain, client!!, rpcClient, nodeRating, labels) + val created = EthereumPosGrpcUpstream(id, hash, role, chain, client!!, rpcClient, nodeRating, labels) created.timeout = this.timeout known[chain] = created created.start() diff --git a/src/test/groovy/io/emeraldpay/dshackle/config/UpstreamsConfigReaderSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/config/UpstreamsConfigReaderSpec.groovy index 2ffdb0fb..5b4589a8 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/config/UpstreamsConfigReaderSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/config/UpstreamsConfigReaderSpec.groovy @@ -401,4 +401,18 @@ class UpstreamsConfigReaderSpec extends Specification { act.upstreams.get(0).role == UpstreamsConfig.UpstreamRole.PRIMARY act.upstreams.get(1).role == UpstreamsConfig.UpstreamRole.PRIMARY } + + def "Parse node id"() { + setup: + def config = this.class.getClassLoader().getResourceAsStream("upstreams-node-id.yaml") + when: + def act = reader.read(config) + then: + act != null + act.upstreams.size() == 2 + act.upstreams[0].nodeId == 1 + act.upstreams[0].id == "has_node_id" + act.upstreams[1].nodeId == null + act.upstreams[1].id == "has_no_node_id" + } } diff --git a/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy index 8a7c9e1d..604e561f 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy @@ -118,7 +118,7 @@ class NativeCallSpec extends Specification { def nativeCall = nativeCall() nativeCall.quorumReaderFactory = Mock(QuorumReaderFactory) { 1 * create(_, _, _) >> Mock(Reader) { - 1 * read(_) >> Mono.just(new QuorumRpcReader.Result("\"foo\"".bytes, null, 1)) + 1 * read(_) >> Mono.just(new QuorumRpcReader.Result("\"foo\"".bytes, null, 1, Collections.singletonList((byte) 1))) } } def call = new NativeCall.ValidCallContext(1, 10, TestingCommons.multistream(TestingCommons.api()), Selector.empty, quorum, @@ -395,6 +395,68 @@ class NativeCallSpec extends Specification { } } + def "Prepare call adds decorator for eth_newFilter"() { + setup: + def methods = new ManagedCallMethods( + new DefaultEthereumMethods(Chain.ETHEREUM), + ["eth_newFilter"] as Set, [] as Set + ) + methods.setQuorum("eth_newFilter", "always") + def multistream = new MultistreamHolderMock.EthereumMultistreamMock(Chain.ETHEREUM, TestingCommons.upstream()) + multistream.customMethods = methods + multistream.customHead = Mock(Head) + def multistreamHolder = Mock(MultistreamHolder) { + _ * it.observeChains() >> Flux.empty() + } + def nativeCall = nativeCall(multistreamHolder) + + def req = BlockchainOuterClass.NativeCallRequest.newBuilder() + .setChain(Common.ChainRef.CHAIN_ETHEREUM) + .addItems( + BlockchainOuterClass.NativeCallItem.newBuilder() + .setId(1) + .setMethod("eth_newFilter") + ) + .build() + when: + def act = nativeCall.prepareCall(req, multistream) + .collectList().block(Duration.ofSeconds(1)).first() + then: + act instanceof NativeCall.ValidCallContext + act.resultDecorator instanceof NativeCall.CreateFilterDecorator + } + + def "Prepare call adds decorator for eth_getFilterChanges"() { + setup: + def methods = new ManagedCallMethods( + new DefaultEthereumMethods(Chain.ETHEREUM), + ["eth_getFilterChanges"] as Set, [] as Set + ) + methods.setQuorum("eth_getFilterChanges", "always") + def multistream = new MultistreamHolderMock.EthereumMultistreamMock(Chain.ETHEREUM, TestingCommons.upstream()) + multistream.customMethods = methods + multistream.customHead = Mock(Head) + def multistreamHolder = Mock(MultistreamHolder) { + _ * it.observeChains() >> Flux.empty() + } + def nativeCall = nativeCall(multistreamHolder) + + def req = BlockchainOuterClass.NativeCallRequest.newBuilder() + .setChain(Common.ChainRef.CHAIN_ETHEREUM) + .addItems( + BlockchainOuterClass.NativeCallItem.newBuilder() + .setId(1) + .setMethod("eth_getFilterChanges") + ) + .build() + when: + def act = nativeCall.prepareCall(req, multistream) + .collectList().block(Duration.ofSeconds(1)).first() + then: + act instanceof NativeCall.ValidCallContext + act.requestDecorator instanceof NativeCall.GetFilterUpdatesDecorator + } + def "Parse empty params"() { setup: def nativeCall = nativeCall() @@ -447,6 +509,64 @@ class NativeCallSpec extends Specification { act.payload.method == "eth_test" } + def "Decorate eth_getFilterUpdates params"() { + setup: + def nativeCall = nativeCall() + def ctx = new NativeCall.ValidCallContext(1, null, Stub(Multistream), Selector.empty, new AlwaysQuorum(), + new NativeCall.RawCallDetails("eth_getFilterUpdates", '["0xabcd"]'), + new NativeCall.GetFilterUpdatesDecorator(), new NativeCall.NoneResultDecorator()) + when: + def act = nativeCall.parseParams(ctx) + then: + act.id == 1 + act.payload.params == ["0xab"] + act.payload.method == "eth_getFilterUpdates" + } + + def "Decorate eth_newFilter result"() { + setup: + def quorum = new AlwaysQuorum() + + def nativeCall = nativeCall() + nativeCall.quorumReaderFactory = Mock(QuorumReaderFactory) { + 1 * create(_, _, _) >> Mock(Reader) { + 1 * read(_) >> Mono.just(new QuorumRpcReader.Result("\"0xab\"".bytes, null, 1, Collections.singletonList((byte)255))) + } + } + def call = new NativeCall.ValidCallContext(1, 10, TestingCommons.multistream(TestingCommons.api()), Selector.empty, quorum, + new NativeCall.ParsedCallDetails("eth_getFilterChanges", []), + new NativeCall.GetFilterUpdatesDecorator(), new NativeCall.CreateFilterDecorator()) + + when: + def resp = nativeCall.executeOnRemote(call).block(Duration.ofSeconds(1)) + def act = objectMapper.readValue(resp.result, Object) + then: + act == "0xabff" + resp.nonce == 10 + } + + def "Decorate eth_newFilter result with short nodeId"() { + setup: + def quorum = new AlwaysQuorum() + + def nativeCall = nativeCall() + nativeCall.quorumReaderFactory = Mock(QuorumReaderFactory) { + 1 * create(_, _, _) >> Mock(Reader) { + 1 * read(_) >> Mono.just(new QuorumRpcReader.Result("\"0xab\"".bytes, null, 1, Collections.singletonList((byte)1))) + } + } + def call = new NativeCall.ValidCallContext(1, 10, TestingCommons.multistream(TestingCommons.api()), Selector.empty, quorum, + new NativeCall.ParsedCallDetails("eth_getFilterChanges", []), + new NativeCall.GetFilterUpdatesDecorator(), new NativeCall.CreateFilterDecorator()) + + when: + def resp = nativeCall.executeOnRemote(call).block(Duration.ofSeconds(1)) + def act = objectMapper.readValue(resp.result, Object) + then: + act == "0xab01" + resp.nonce == 10 + } + @Ignore //TODO def "Calls cache before remote"() { diff --git a/src/test/groovy/io/emeraldpay/dshackle/startup/ConfiguredUpstreamsSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/startup/ConfiguredUpstreamsSpec.groovy index 106541c6..4e6f6eea 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/startup/ConfiguredUpstreamsSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/startup/ConfiguredUpstreamsSpec.groovy @@ -58,4 +58,37 @@ class ConfiguredUpstreamsSpec extends Specification { act instanceof ManagedCallMethods new String(act.executeHardcoded("foo_bar")) == "\"static_response\"" } + + def "Calculate node-id"() { + setup: + def configurer = new ConfiguredUpstreams(Stub(CurrentMultistreamHolder), Stub(FileResolver), Stub(UpstreamsConfig) + ) + expect: + configurer.getHash(node, src) == expected + + where: + node | src | expected + 1 | "" | 1 + 9 | "hohoho" | 9 + null | "hohoho" | 120 + } + + def "Calculate node-id conflicting results"() { + setup: + def configurer = new ConfiguredUpstreams(Stub(CurrentMultistreamHolder), Stub(FileResolver), Stub(UpstreamsConfig) + ) + when: + def h1 = configurer.getHash(null, "hohoho") + def h2 = configurer.getHash(null, "hohoho") + def h3 = configurer.getHash(null, "hohoho") + def h4 = configurer.getHash(null, "hohoho") + def h5 = configurer.getHash(null, "hohoho") + + then: + h1 == (byte)120 + h2 == (byte)-120 + h3 == (byte)-9 + h4 == (byte)8 + h5 == (byte)-128 + } } diff --git a/src/test/groovy/io/emeraldpay/dshackle/test/EthereumPosRpcUpstreamMock.groovy b/src/test/groovy/io/emeraldpay/dshackle/test/EthereumPosRpcUpstreamMock.groovy index 7c1f4954..e611b3d9 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/test/EthereumPosRpcUpstreamMock.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/test/EthereumPosRpcUpstreamMock.groovy @@ -69,7 +69,7 @@ class EthereumPosRpcUpstreamMock extends EthereumPosRpcUpstream { } EthereumPosRpcUpstreamMock(@NotNull String id, @NotNull Chain chain, @NotNull Reader api, CallMethods methods, Map labels) { - super(id, chain, + super(id, (byte)id.hashCode(), chain, UpstreamsConfig.Options.getDefaults(), UpstreamsConfig.UpstreamRole.PRIMARY, methods, diff --git a/src/test/groovy/io/emeraldpay/dshackle/test/EthereumRpcUpstreamMock.groovy b/src/test/groovy/io/emeraldpay/dshackle/test/EthereumRpcUpstreamMock.groovy index 89bfa684..369dc612 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/test/EthereumRpcUpstreamMock.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/test/EthereumRpcUpstreamMock.groovy @@ -60,7 +60,7 @@ class EthereumRpcUpstreamMock extends EthereumRpcUpstream { } EthereumRpcUpstreamMock(@NotNull String id, @NotNull Chain chain, @NotNull Reader api, CallMethods methods) { - super(id, chain, + super(id, id.hashCode().byteValue(), chain, UpstreamsConfig.Options.getDefaults(), UpstreamsConfig.UpstreamRole.PRIMARY, methods, diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/FilteredApisSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/FilteredApisSpec.groovy index d86b035d..ecec9e4b 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/FilteredApisSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/FilteredApisSpec.groovy @@ -51,6 +51,7 @@ class FilteredApisSpec extends Specification { def connectorFactory = new EthereumConnectorFactory(false, null, httpFactory, new MostWorkForkChoice(), BlockValidator.@Companion.ALWAYS_VALID) new EthereumRpcUpstream( "test", + (byte)123, Chain.ETHEREUM, new UpstreamsConfig.Options(), UpstreamsConfig.UpstreamRole.PRIMARY, diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/SelectorSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/SelectorSpec.groovy index ab21834f..fd8de016 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/SelectorSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/SelectorSpec.groovy @@ -349,4 +349,27 @@ class SelectorSpec extends Specification { !act } + def "Matches same nodeId"() { + setup: + def up = Mock(Upstream) { + nodeId() >> (byte)5 + } + def matcher = new Selector.SameNodeMatcher((byte)5) + when: + def act = matcher.matches(up) + then: + act + } + + def "Not matches nodeId"() { + setup: + def up = Mock(Upstream) { + nodeId() >> (byte)5 + } + def matcher = new Selector.SameNodeMatcher((byte)1) + when: + def act = matcher.matches(up) + then: + !act + } } diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/calls/EthereumCallSelectorSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/calls/EthereumCallSelectorSpec.groovy index ea5d15cb..8aded145 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/calls/EthereumCallSelectorSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/calls/EthereumCallSelectorSpec.groovy @@ -179,4 +179,33 @@ class EthereumCallSelectorSpec extends Specification { then: act == new Selector.HeightMatcher(100) } + + def "Get same matcher for getFilterChanges method"() { + setup: + def callSelector = new EthereumCallSelector(Mock(Reader)) + def head = Mock(Head) + + expect: + callSelector.getMatcher("eth_getFilterChanges", param, head).block() + == new Selector.SameNodeMatcher((byte)hash) + + where: + param | hash + '["0xff09"]' | 9 + '["0xff"]' | 255 + '[""]' | 0 + '["0x0"]' | 0 + } + + def "Get empty matcher for getFilterChanges method without params"() { + setup: + def callSelector = new EthereumCallSelector(Mock(Reader)) + def head = Mock(Head) + + when: + def act = callSelector.getMatcher("eth_getFilterChanges", "[]", head).block() + + then: + act == null + } } diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/EthereumDirectReaderSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/EthereumDirectReaderSpec.groovy index ac2781f5..78bcbe27 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/EthereumDirectReaderSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/EthereumDirectReaderSpec.groovy @@ -30,6 +30,7 @@ class EthereumDirectReaderSpec extends Specification { String hash1 = "0x40d15edaff9acdabd2a1c96fd5f683b3300aad34e7015f34def3c56ba8a7ffb5" String address1 = "0xe0aadb0a012dbcdc529c4c743d3e0385a0b54d3d" + List resolvers = Collections.singletonList((byte)1) def "Reads block by hash"() { setup: @@ -53,8 +54,7 @@ class EthereumDirectReaderSpec extends Specification { 1 * create(_, _, _) >> Mock(Reader) { 1 * read(new JsonRpcRequest("eth_getBlockByHash", [hash1, false])) >> Mono.just( new QuorumRpcReader.Result( - Global.objectMapper.writeValueAsBytes(json), null, 1 - ) + Global.objectMapper.writeValueAsBytes(json), null, 1, resolvers) ) } } @@ -84,7 +84,7 @@ class EthereumDirectReaderSpec extends Specification { 1 * create(_, _, _) >> Mock(Reader) { 1 * read(new JsonRpcRequest("eth_getBlockByHash", [hash1, false])) >> Mono.just( new QuorumRpcReader.Result( - Global.objectMapper.writeValueAsBytes(null), null, 1 + Global.objectMapper.writeValueAsBytes(null), null, 1, resolvers ) ) } @@ -119,7 +119,7 @@ class EthereumDirectReaderSpec extends Specification { 1 * create(_, _, _) >> Mock(Reader) { 1 * read(new JsonRpcRequest("eth_getBlockByNumber", ["0x64", false])) >> Mono.just( new QuorumRpcReader.Result( - Global.objectMapper.writeValueAsBytes(json), null, 1 + Global.objectMapper.writeValueAsBytes(json), null, 1, resolvers ) ) } @@ -155,7 +155,7 @@ class EthereumDirectReaderSpec extends Specification { 1 * create(_, _, _) >> Mock(Reader) { 1 * read(new JsonRpcRequest("eth_getTransactionByHash", [hash1])) >> Mono.just( new QuorumRpcReader.Result( - Global.objectMapper.writeValueAsBytes(json), null, 1 + Global.objectMapper.writeValueAsBytes(json), null, 1, resolvers ) ) } @@ -186,7 +186,7 @@ class EthereumDirectReaderSpec extends Specification { 1 * create(_, _, _) >> Mock(Reader) { 1 * read(new JsonRpcRequest("eth_getTransactionByHash", [hash1])) >> Mono.just( new QuorumRpcReader.Result( - Global.objectMapper.writeValueAsBytes(null), null, 1 + Global.objectMapper.writeValueAsBytes(null), null, 1, resolvers ) ) } @@ -217,7 +217,7 @@ class EthereumDirectReaderSpec extends Specification { 1 * create(_, _, _) >> Mock(Reader) { 1 * read(new JsonRpcRequest("eth_getBalance", [address1, "latest"])) >> Mono.just( new QuorumRpcReader.Result( - Global.objectMapper.writeValueAsBytes("0x100"), null, 1 + Global.objectMapper.writeValueAsBytes("0x100"), null, 1, resolvers ) ) } @@ -249,7 +249,7 @@ class EthereumDirectReaderSpec extends Specification { 1 * create(_, _, _) >> Mock(Reader) { 1 * read(new JsonRpcRequest("eth_getBalance", [address1, "0xa8c9bb"])) >> Mono.just( new QuorumRpcReader.Result( - Global.objectMapper.writeValueAsBytes("0x100"), null, 1 + Global.objectMapper.writeValueAsBytes("0x100"), null, 1, resolvers ) ) } diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcUpstreamSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcUpstreamSpec.groovy index 5418684b..af98978b 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcUpstreamSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcUpstreamSpec.groovy @@ -50,6 +50,8 @@ class EthereumGrpcUpstreamSpec extends Specification { Counter.builder("test2").register(TestingCommons.meterRegistry) ) + def hash = (byte)123 + def "Subscribe to head"() { setup: def callData = [:] @@ -81,7 +83,7 @@ class EthereumGrpcUpstreamSpec extends Specification { ) } }) - def upstream = new EthereumGrpcUpstream("test", UpstreamsConfig.UpstreamRole.PRIMARY, chain, client, new JsonRpcGrpcClient(client, chain, metrics), null) + def upstream = new EthereumGrpcUpstream("test", hash, UpstreamsConfig.UpstreamRole.PRIMARY, chain, client, new JsonRpcGrpcClient(client, chain, metrics), null) upstream.setLag(0) upstream.update(BlockchainOuterClass.DescribeChain.newBuilder() .setStatus(BlockchainOuterClass.ChainStatus.newBuilder().setQuorum(1).setAvailabilityValue(UpstreamAvailability.OK.grpcId)) @@ -139,7 +141,7 @@ class EthereumGrpcUpstreamSpec extends Specification { ) } }) - def upstream = new EthereumGrpcUpstream("test", UpstreamsConfig.UpstreamRole.PRIMARY, Chain.ETHEREUM, client, new JsonRpcGrpcClient(client, Chain.ETHEREUM, metrics), null) + def upstream = new EthereumGrpcUpstream("test", hash, UpstreamsConfig.UpstreamRole.PRIMARY, Chain.ETHEREUM, client, new JsonRpcGrpcClient(client, Chain.ETHEREUM, metrics), null) upstream.setLag(0) upstream.update(BlockchainOuterClass.DescribeChain.newBuilder() .setStatus(BlockchainOuterClass.ChainStatus.newBuilder().setQuorum(1).setAvailabilityValue(UpstreamAvailability.OK.grpcId)) @@ -201,7 +203,7 @@ class EthereumGrpcUpstreamSpec extends Specification { finished.complete(true) } }) - def upstream = new EthereumGrpcUpstream("test", UpstreamsConfig.UpstreamRole.PRIMARY, chain, client, new JsonRpcGrpcClient(client, chain, metrics), null) + def upstream = new EthereumGrpcUpstream("test", hash, UpstreamsConfig.UpstreamRole.PRIMARY, chain, client, new JsonRpcGrpcClient(client, chain, metrics), null) upstream.setLag(0) upstream.update(BlockchainOuterClass.DescribeChain.newBuilder() .setStatus(BlockchainOuterClass.ChainStatus.newBuilder().setQuorum(1).setAvailabilityValue(UpstreamAvailability.OK.grpcId)) diff --git a/src/test/resources/upstreams-node-id.yaml b/src/test/resources/upstreams-node-id.yaml new file mode 100644 index 00000000..cfea4978 --- /dev/null +++ b/src/test/resources/upstreams-node-id.yaml @@ -0,0 +1,34 @@ +version: v1 +upstreams: + - id: has_node_id + node-id: 1 + chain: ethereum + connection: + ethereum: + rpc: + url: "http://localhost:8545" + - id: has_no_node_id + chain: ethereum + connection: + ethereum: + rpc: + url: "http://localhost:8545" + ws: + url: "ws://localhost:8546" + - id: conflicted_node_id + node-id: 1 + chain: ethereum + connection: + prefer-http: true + ethereum: + rpc: + url: "http://localhost:9545" + ws: + url: "ws://localhost:9546" + - id: invalid_node_id + node-id: 256 + chain: ethereum + connection: + grpc: + host: "localhost" + port: 2449 \ No newline at end of file