Merge branch 'master' into evm_networks_support
# Conflicts: # src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt
This commit is contained in:
@@ -51,6 +51,7 @@ cluster:
|
|||||||
min-peers: 2
|
min-peers: 2
|
||||||
upstreams:
|
upstreams:
|
||||||
- id: us-nodes
|
- id: us-nodes
|
||||||
|
node-id: 1
|
||||||
chain: auto
|
chain: auto
|
||||||
connection:
|
connection:
|
||||||
grpc:
|
grpc:
|
||||||
@@ -61,6 +62,7 @@ cluster:
|
|||||||
certificate: client.crt
|
certificate: client.crt
|
||||||
key: client.p8.key
|
key: client.p8.key
|
||||||
- id: infura-eth
|
- id: infura-eth
|
||||||
|
node-id: 2
|
||||||
chain: ethereum
|
chain: ethereum
|
||||||
role: fallback
|
role: fallback
|
||||||
labels:
|
labels:
|
||||||
@@ -80,6 +82,7 @@ cluster:
|
|||||||
username: ${INFURA_USER}
|
username: ${INFURA_USER}
|
||||||
password: ${INFURA_PASSWD}
|
password: ${INFURA_PASSWD}
|
||||||
- id: ethereum-pos
|
- id: ethereum-pos
|
||||||
|
node-id: 3
|
||||||
chain: ropsten
|
chain: ropsten
|
||||||
connection:
|
connection:
|
||||||
ethereum-pos:
|
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]"_
|
** 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
|
** 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
|
=== Roles and Fallback upstream
|
||||||
|
|
||||||
By default, the Dshackle connects to each upstream in a Round-Robin basis, i.e. sequentially one by one.
|
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_getUncleByBlockNumberAndIndex`
|
||||||
- `eth_feeHistory`
|
- `eth_feeHistory`
|
||||||
- `eth_getLogs`
|
- `eth_getLogs`
|
||||||
|
- `eth_getFilterChanges`
|
||||||
|
- `eth_newFilter`
|
||||||
|
- `eth_newBlockFilter`
|
||||||
|
- `eth_newPendingTransactionFilter`
|
||||||
|
|
||||||
.Plus following methods are answered directly by Dshackle
|
.Plus following methods are answered directly by Dshackle
|
||||||
- `net_version`
|
- `net_version`
|
||||||
|
|||||||
@@ -69,6 +69,7 @@ open class UpstreamsConfig {
|
|||||||
|
|
||||||
class Upstream<T : UpstreamConnection> {
|
class Upstream<T : UpstreamConnection> {
|
||||||
var id: String? = null
|
var id: String? = null
|
||||||
|
var nodeId: Int? = null
|
||||||
var chain: String? = null
|
var chain: String? = null
|
||||||
var options: Options? = null
|
var options: Options? = null
|
||||||
var isEnabled = true
|
var isEnabled = true
|
||||||
|
|||||||
@@ -31,6 +31,7 @@ class UpstreamsConfigReader(
|
|||||||
|
|
||||||
private val log = LoggerFactory.getLogger(UpstreamsConfigReader::class.java)
|
private val log = LoggerFactory.getLogger(UpstreamsConfigReader::class.java)
|
||||||
private val authConfigReader = AuthConfigReader()
|
private val authConfigReader = AuthConfigReader()
|
||||||
|
private val knownNodeIds: MutableSet<Int> = HashSet()
|
||||||
|
|
||||||
fun read(input: InputStream): UpstreamsConfig? {
|
fun read(input: InputStream): UpstreamsConfig? {
|
||||||
val configNode = readNode(input)
|
val configNode = readNode(input)
|
||||||
@@ -236,11 +237,22 @@ class UpstreamsConfigReader(
|
|||||||
log.warn("Invalid id: $id")
|
log.warn("Invalid id: $id")
|
||||||
return false
|
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<*>) {
|
internal fun readUpstreamCommon(upNode: MappingNode, upstream: UpstreamsConfig.Upstream<*>) {
|
||||||
upstream.id = getValueAsString(upNode, "id")
|
upstream.id = getValueAsString(upNode, "id")
|
||||||
|
upstream.nodeId = getValueAsInt(upNode, "node-id")
|
||||||
upstream.options = tryReadOptions(upNode)
|
upstream.options = tryReadOptions(upNode)
|
||||||
upstream.methods = tryReadMethods(upNode)
|
upstream.methods = tryReadMethods(upNode)
|
||||||
getValueAsBool(upNode, "enabled")?.let {
|
getValueAsBool(upNode, "enabled")?.let {
|
||||||
|
|||||||
@@ -21,6 +21,7 @@ import io.emeraldpay.dshackle.upstream.Upstream
|
|||||||
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcError
|
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcError
|
||||||
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcException
|
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcException
|
||||||
import io.emeraldpay.dshackle.upstream.signature.ResponseSigner
|
import io.emeraldpay.dshackle.upstream.signature.ResponseSigner
|
||||||
|
import java.util.concurrent.ConcurrentLinkedQueue
|
||||||
|
|
||||||
open class AlwaysQuorum : CallQuorum {
|
open class AlwaysQuorum : CallQuorum {
|
||||||
|
|
||||||
@@ -28,6 +29,7 @@ open class AlwaysQuorum : CallQuorum {
|
|||||||
private var result: ByteArray? = null
|
private var result: ByteArray? = null
|
||||||
private var rpcError: JsonRpcError? = null
|
private var rpcError: JsonRpcError? = null
|
||||||
private var sig: ResponseSigner.Signature? = null
|
private var sig: ResponseSigner.Signature? = null
|
||||||
|
private val resolvers: MutableCollection<Upstream> = ConcurrentLinkedQueue()
|
||||||
|
|
||||||
override fun init(head: Head) {
|
override fun init(head: Head) {
|
||||||
}
|
}
|
||||||
@@ -48,6 +50,7 @@ open class AlwaysQuorum : CallQuorum {
|
|||||||
result = response
|
result = response
|
||||||
resolved = true
|
resolved = true
|
||||||
sig = signature
|
sig = signature
|
||||||
|
resolvers.add(upstream)
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -64,6 +67,9 @@ open class AlwaysQuorum : CallQuorum {
|
|||||||
return rpcError
|
return rpcError
|
||||||
}
|
}
|
||||||
|
|
||||||
|
override fun getResolvedBy(): List<Upstream> =
|
||||||
|
resolvers.toList()
|
||||||
|
|
||||||
override fun toString(): String {
|
override fun toString(): String {
|
||||||
return "Quorum: Accept Any"
|
return "Quorum: Accept Any"
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -34,4 +34,5 @@ interface CallQuorum {
|
|||||||
fun getSignature(): ResponseSigner.Signature?
|
fun getSignature(): ResponseSigner.Signature?
|
||||||
fun getResult(): ByteArray?
|
fun getResult(): ByteArray?
|
||||||
fun getError(): JsonRpcError?
|
fun getError(): JsonRpcError?
|
||||||
|
fun getResolvedBy(): Collection<Upstream>
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -21,6 +21,7 @@ import io.emeraldpay.dshackle.upstream.Upstream
|
|||||||
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcError
|
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcError
|
||||||
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcException
|
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcException
|
||||||
import io.emeraldpay.dshackle.upstream.signature.ResponseSigner
|
import io.emeraldpay.dshackle.upstream.signature.ResponseSigner
|
||||||
|
import java.util.concurrent.ConcurrentLinkedQueue
|
||||||
import java.util.concurrent.atomic.AtomicReference
|
import java.util.concurrent.atomic.AtomicReference
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -34,6 +35,7 @@ class NotLaggingQuorum(val maxLag: Long = 0) : CallQuorum {
|
|||||||
private val failed = AtomicReference(false)
|
private val failed = AtomicReference(false)
|
||||||
private var rpcError: JsonRpcError? = null
|
private var rpcError: JsonRpcError? = null
|
||||||
private var sig: ResponseSigner.Signature? = null
|
private var sig: ResponseSigner.Signature? = null
|
||||||
|
private val resolvers: MutableCollection<Upstream> = ConcurrentLinkedQueue()
|
||||||
|
|
||||||
override fun init(head: Head) {
|
override fun init(head: Head) {
|
||||||
}
|
}
|
||||||
@@ -51,6 +53,7 @@ class NotLaggingQuorum(val maxLag: Long = 0) : CallQuorum {
|
|||||||
if (!lagging) {
|
if (!lagging) {
|
||||||
result.set(response)
|
result.set(response)
|
||||||
sig = signature
|
sig = signature
|
||||||
|
resolvers.add(upstream)
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
return false
|
return false
|
||||||
@@ -75,6 +78,8 @@ class NotLaggingQuorum(val maxLag: Long = 0) : CallQuorum {
|
|||||||
return rpcError
|
return rpcError
|
||||||
}
|
}
|
||||||
|
|
||||||
|
override fun getResolvedBy(): Collection<Upstream> =
|
||||||
|
resolvers.toList()
|
||||||
override fun toString(): String {
|
override fun toString(): String {
|
||||||
return "Quorum: late <= $maxLag blocks"
|
return "Quorum: late <= $maxLag blocks"
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -117,9 +117,9 @@ class QuorumRpcReader(
|
|||||||
return Function { quorumResult ->
|
return Function { quorumResult ->
|
||||||
quorumResult
|
quorumResult
|
||||||
.filter { it.isResolved() } // return nothing if not resolved
|
.filter { it.isResolved() } // return nothing if not resolved
|
||||||
.map {
|
.map { quorum ->
|
||||||
// TODO find actual quorum number
|
// 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)
|
.switchIfEmpty(defaultResult)
|
||||||
}
|
}
|
||||||
@@ -197,6 +197,7 @@ class QuorumRpcReader(
|
|||||||
class Result(
|
class Result(
|
||||||
val value: ByteArray,
|
val value: ByteArray,
|
||||||
val signature: ResponseSigner.Signature?,
|
val signature: ResponseSigner.Signature?,
|
||||||
val quorum: Int
|
val quorum: Int,
|
||||||
|
val resolvers: Collection<Byte>
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -23,6 +23,7 @@ import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcException
|
|||||||
import io.emeraldpay.dshackle.upstream.signature.ResponseSigner
|
import io.emeraldpay.dshackle.upstream.signature.ResponseSigner
|
||||||
import io.emeraldpay.etherjar.rpc.RpcException
|
import io.emeraldpay.etherjar.rpc.RpcException
|
||||||
import org.slf4j.LoggerFactory
|
import org.slf4j.LoggerFactory
|
||||||
|
import java.util.concurrent.ConcurrentLinkedQueue
|
||||||
|
|
||||||
abstract class ValueAwareQuorum<T>(
|
abstract class ValueAwareQuorum<T>(
|
||||||
val clazz: Class<T>
|
val clazz: Class<T>
|
||||||
@@ -30,6 +31,7 @@ abstract class ValueAwareQuorum<T>(
|
|||||||
|
|
||||||
private val log = LoggerFactory.getLogger(ValueAwareQuorum::class.java)
|
private val log = LoggerFactory.getLogger(ValueAwareQuorum::class.java)
|
||||||
private var rpcError: JsonRpcError? = null
|
private var rpcError: JsonRpcError? = null
|
||||||
|
private val resolvers: MutableCollection<Upstream> = ConcurrentLinkedQueue()
|
||||||
|
|
||||||
fun extractValue(response: ByteArray, clazz: Class<T>): T? {
|
fun extractValue(response: ByteArray, clazz: Class<T>): T? {
|
||||||
return Global.objectMapper.readValue(response.inputStream(), clazz)
|
return Global.objectMapper.readValue(response.inputStream(), clazz)
|
||||||
@@ -39,6 +41,7 @@ abstract class ValueAwareQuorum<T>(
|
|||||||
try {
|
try {
|
||||||
val value = extractValue(response, clazz)
|
val value = extractValue(response, clazz)
|
||||||
recordValue(response, value, signature, upstream)
|
recordValue(response, value, signature, upstream)
|
||||||
|
resolvers.add(upstream)
|
||||||
} catch (e: RpcException) {
|
} catch (e: RpcException) {
|
||||||
recordError(response, e.rpcMessage, signature, upstream)
|
recordError(response, e.rpcMessage, signature, upstream)
|
||||||
} catch (e: Exception) {
|
} catch (e: Exception) {
|
||||||
@@ -59,4 +62,7 @@ abstract class ValueAwareQuorum<T>(
|
|||||||
override fun getError(): JsonRpcError? {
|
override fun getError(): JsonRpcError? {
|
||||||
return rpcError
|
return rpcError
|
||||||
}
|
}
|
||||||
|
|
||||||
|
override fun getResolvedBy(): Collection<Upstream> =
|
||||||
|
resolvers.toList()
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -105,7 +105,8 @@ open class NativeCall(
|
|||||||
}
|
}
|
||||||
|
|
||||||
fun parseParams(it: ValidCallContext<RawCallDetails>): ValidCallContext<ParsedCallDetails> {
|
fun parseParams(it: ValidCallContext<RawCallDetails>): ValidCallContext<ParsedCallDetails> {
|
||||||
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))
|
return it.withPayload(ParsedCallDetails(it.payload.method, params))
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -237,17 +238,28 @@ open class NativeCall(
|
|||||||
matcher.withMatcher(heightMatcher)
|
matcher.withMatcher(heightMatcher)
|
||||||
}
|
}
|
||||||
val nonce = requestItem.nonce.let { if (it == 0L) null else it }
|
val nonce = requestItem.nonce.let { if (it == 0L) null else it }
|
||||||
|
val requestDecorator = getRequestDecorator(requestItem.method)
|
||||||
|
val resultDecorator = getResultDecorator(requestItem.method)
|
||||||
|
|
||||||
ValidCallContext(
|
ValidCallContext(
|
||||||
requestItem.id,
|
requestItem.id,
|
||||||
nonce,
|
nonce,
|
||||||
upstream,
|
upstream,
|
||||||
matcher.build(),
|
matcher.build(),
|
||||||
callQuorum,
|
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<ParsedCallDetails>): Mono<CallResult> {
|
fun fetch(ctx: ValidCallContext<ParsedCallDetails>): Mono<CallResult> {
|
||||||
return ctx.upstream.getRoutedApi(ctx.matcher)
|
return ctx.upstream.getRoutedApi(ctx.matcher)
|
||||||
.flatMap { api ->
|
.flatMap { api ->
|
||||||
@@ -286,7 +298,8 @@ open class NativeCall(
|
|||||||
return reader
|
return reader
|
||||||
.read(JsonRpcRequest(ctx.payload.method, ctx.payload.params, ctx.nonce))
|
.read(JsonRpcRequest(ctx.payload.method, ctx.payload.params, ctx.nonce))
|
||||||
.map {
|
.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 ->
|
.onErrorResume { t ->
|
||||||
val failure = when (t) {
|
val failure = when (t) {
|
||||||
@@ -353,14 +366,67 @@ open class NativeCall(
|
|||||||
fun getError(): CallError
|
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<Any>): List<Any>
|
||||||
|
}
|
||||||
|
|
||||||
|
open class NoneRequestDecorator : RequestDecorator {
|
||||||
|
override fun processRequest(request: List<Any>): List<Any> = request
|
||||||
|
}
|
||||||
|
|
||||||
|
open class GetFilterUpdatesDecorator : RequestDecorator {
|
||||||
|
override fun processRequest(request: List<Any>): List<Any> {
|
||||||
|
val filterId = request.first().toString()
|
||||||
|
val sanitized = filterId.substring(0, filterId.lastIndex - 1)
|
||||||
|
return listOf(sanitized)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
open class ValidCallContext<T>(
|
open class ValidCallContext<T>(
|
||||||
val id: Int,
|
val id: Int,
|
||||||
val nonce: Long?,
|
val nonce: Long?,
|
||||||
val upstream: Multistream,
|
val upstream: Multistream,
|
||||||
val matcher: Selector.Matcher,
|
val matcher: Selector.Matcher,
|
||||||
val callQuorum: CallQuorum,
|
val callQuorum: CallQuorum,
|
||||||
val payload: T
|
val payload: T,
|
||||||
|
val requestDecorator: RequestDecorator,
|
||||||
|
val resultDecorator: ResultDecorator
|
||||||
) : CallContext {
|
) : 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 {
|
override fun isValid(): Boolean {
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
@@ -374,7 +440,7 @@ open class NativeCall(
|
|||||||
}
|
}
|
||||||
|
|
||||||
fun <X> withPayload(payload: X): ValidCallContext<X> {
|
fun <X> withPayload(payload: X): ValidCallContext<X> {
|
||||||
return ValidCallContext(id, nonce, upstream, matcher, callQuorum, payload)
|
return ValidCallContext(id, nonce, upstream, matcher, callQuorum, payload, requestDecorator, resultDecorator)
|
||||||
}
|
}
|
||||||
|
|
||||||
fun getApis(): ApiSource {
|
fun getApis(): ApiSource {
|
||||||
|
|||||||
@@ -16,6 +16,7 @@
|
|||||||
*/
|
*/
|
||||||
package io.emeraldpay.dshackle.startup
|
package io.emeraldpay.dshackle.startup
|
||||||
|
|
||||||
|
import com.google.common.annotations.VisibleForTesting
|
||||||
import io.emeraldpay.dshackle.FileResolver
|
import io.emeraldpay.dshackle.FileResolver
|
||||||
import io.emeraldpay.dshackle.Global
|
import io.emeraldpay.dshackle.Global
|
||||||
import io.emeraldpay.dshackle.config.UpstreamsConfig
|
import io.emeraldpay.dshackle.config.UpstreamsConfig
|
||||||
@@ -53,7 +54,9 @@ import org.springframework.beans.factory.annotation.Autowired
|
|||||||
import org.springframework.stereotype.Repository
|
import org.springframework.stereotype.Repository
|
||||||
import java.net.URI
|
import java.net.URI
|
||||||
import java.util.concurrent.atomic.AtomicInteger
|
import java.util.concurrent.atomic.AtomicInteger
|
||||||
|
import java.util.function.Function
|
||||||
import javax.annotation.PostConstruct
|
import javax.annotation.PostConstruct
|
||||||
|
import kotlin.math.abs
|
||||||
|
|
||||||
@Repository
|
@Repository
|
||||||
open class ConfiguredUpstreams(
|
open class ConfiguredUpstreams(
|
||||||
@@ -65,6 +68,8 @@ open class ConfiguredUpstreams(
|
|||||||
private val log = LoggerFactory.getLogger(ConfiguredUpstreams::class.java)
|
private val log = LoggerFactory.getLogger(ConfiguredUpstreams::class.java)
|
||||||
private var seq = AtomicInteger(0)
|
private var seq = AtomicInteger(0)
|
||||||
|
|
||||||
|
private val hashes: MutableMap<Byte, Boolean> = HashMap()
|
||||||
|
|
||||||
@PostConstruct
|
@PostConstruct
|
||||||
fun start() {
|
fun start() {
|
||||||
log.debug("Starting upstreams")
|
log.debug("Starting upstreams")
|
||||||
@@ -77,7 +82,7 @@ open class ConfiguredUpstreams(
|
|||||||
log.debug("Start upstream ${up.id}")
|
log.debug("Start upstream ${up.id}")
|
||||||
if (up.connection is UpstreamsConfig.GrpcConnection) {
|
if (up.connection is UpstreamsConfig.GrpcConnection) {
|
||||||
val options = up.options ?: UpstreamsConfig.Options()
|
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 {
|
} else {
|
||||||
val chain = Global.chainById(up.chain)
|
val chain = Global.chainById(up.chain)
|
||||||
if (chain == Chain.UNSPECIFIED) {
|
if (chain == Chain.UNSPECIFIED) {
|
||||||
@@ -88,14 +93,22 @@ open class ConfiguredUpstreams(
|
|||||||
.merge(defaultOptions[chain] ?: UpstreamsConfig.Options.getDefaults())
|
.merge(defaultOptions[chain] ?: UpstreamsConfig.Options.getDefaults())
|
||||||
val upstream = when (BlockchainType.from(chain)) {
|
val upstream = when (BlockchainType.from(chain)) {
|
||||||
BlockchainType.EVM_POW -> {
|
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 -> {
|
BlockchainType.BITCOIN -> {
|
||||||
buildBitcoinUpstream(up.cast(UpstreamsConfig.BitcoinConnection::class.java), chain, options)
|
buildBitcoinUpstream(up.cast(UpstreamsConfig.BitcoinConnection::class.java), chain, options)
|
||||||
}
|
}
|
||||||
|
|
||||||
BlockchainType.EVM_POS -> {
|
BlockchainType.EVM_POS -> {
|
||||||
buildEthereumPosUpstream(up.cast(UpstreamsConfig.EthereumPosConnection::class.java), chain, options)
|
buildEthereumPosUpstream(
|
||||||
|
up.nodeId,
|
||||||
|
up.cast(UpstreamsConfig.EthereumPosConnection::class.java),
|
||||||
|
chain,
|
||||||
|
options
|
||||||
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
else -> {
|
else -> {
|
||||||
log.error("Chain is unsupported: ${up.chain}")
|
log.error("Chain is unsupported: ${up.chain}")
|
||||||
return@forEach
|
return@forEach
|
||||||
@@ -156,6 +169,7 @@ open class ConfiguredUpstreams(
|
|||||||
}
|
}
|
||||||
|
|
||||||
private fun buildEthereumPosUpstream(
|
private fun buildEthereumPosUpstream(
|
||||||
|
nodeId: Int?,
|
||||||
config: UpstreamsConfig.Upstream<UpstreamsConfig.EthereumPosConnection>,
|
config: UpstreamsConfig.Upstream<UpstreamsConfig.EthereumPosConnection>,
|
||||||
chain: Chain,
|
chain: Chain,
|
||||||
options: UpstreamsConfig.Options
|
options: UpstreamsConfig.Options
|
||||||
@@ -167,13 +181,26 @@ open class ConfiguredUpstreams(
|
|||||||
return null
|
return null
|
||||||
}
|
}
|
||||||
val urls = ArrayList<URI>()
|
val urls = ArrayList<URI>()
|
||||||
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)
|
val methods = buildMethods(config, chain)
|
||||||
if (connectorFactory == null) {
|
if (connectorFactory == null) {
|
||||||
return 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(
|
val upstream = EthereumPosRpcUpstream(
|
||||||
config.id!!,
|
config.id!!,
|
||||||
|
hash,
|
||||||
chain,
|
chain,
|
||||||
options, config.role,
|
options, config.role,
|
||||||
methods,
|
methods,
|
||||||
@@ -227,6 +254,7 @@ open class ConfiguredUpstreams(
|
|||||||
}
|
}
|
||||||
|
|
||||||
private fun buildEthereumUpstream(
|
private fun buildEthereumUpstream(
|
||||||
|
nodeId: Int?,
|
||||||
config: UpstreamsConfig.Upstream<UpstreamsConfig.EthereumConnection>,
|
config: UpstreamsConfig.Upstream<UpstreamsConfig.EthereumConnection>,
|
||||||
chain: Chain,
|
chain: Chain,
|
||||||
options: UpstreamsConfig.Options
|
options: UpstreamsConfig.Options
|
||||||
@@ -236,12 +264,22 @@ open class ConfiguredUpstreams(
|
|||||||
val urls = ArrayList<URI>()
|
val urls = ArrayList<URI>()
|
||||||
val methods = buildMethods(config, chain)
|
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) {
|
if (connectorFactory == null) {
|
||||||
return null
|
return null
|
||||||
}
|
}
|
||||||
|
|
||||||
|
val hashUrl = if (conn.preferHttp) conn.rpc?.url ?: conn.ws?.url else conn.ws?.url ?: conn.rpc?.url
|
||||||
val upstream = EthereumRpcUpstream(
|
val upstream = EthereumRpcUpstream(
|
||||||
config.id!!,
|
config.id!!,
|
||||||
|
getHash(nodeId, hashUrl!!),
|
||||||
chain,
|
chain,
|
||||||
options, config.role,
|
options, config.role,
|
||||||
methods,
|
methods,
|
||||||
@@ -253,12 +291,15 @@ open class ConfiguredUpstreams(
|
|||||||
}
|
}
|
||||||
|
|
||||||
private fun buildGrpcUpstream(
|
private fun buildGrpcUpstream(
|
||||||
|
nodeId: Int?,
|
||||||
config: UpstreamsConfig.Upstream<UpstreamsConfig.GrpcConnection>,
|
config: UpstreamsConfig.Upstream<UpstreamsConfig.GrpcConnection>,
|
||||||
options: UpstreamsConfig.Options
|
options: UpstreamsConfig.Options
|
||||||
) {
|
) {
|
||||||
val endpoint = config.connection!!
|
val endpoint = config.connection!!
|
||||||
|
val hash = getHash(nodeId, "${endpoint.host}:${endpoint.port}")
|
||||||
val ds = GrpcUpstreams(
|
val ds = GrpcUpstreams(
|
||||||
config.id!!,
|
config.id!!,
|
||||||
|
hash,
|
||||||
config.role,
|
config.role,
|
||||||
endpoint.host!!,
|
endpoint.host!!,
|
||||||
endpoint.port,
|
endpoint.port,
|
||||||
@@ -289,7 +330,12 @@ open class ConfiguredUpstreams(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private fun buildWsFactory(id: String, chain: Chain, conn: UpstreamsConfig.EthereumConnection, urls: ArrayList<URI>? = null): EthereumWsFactory? {
|
private fun buildWsFactory(
|
||||||
|
id: String,
|
||||||
|
chain: Chain,
|
||||||
|
conn: UpstreamsConfig.EthereumConnection,
|
||||||
|
urls: ArrayList<URI>? = null
|
||||||
|
): EthereumWsFactory? {
|
||||||
return conn.ws?.let { endpoint ->
|
return conn.ws?.let { endpoint ->
|
||||||
val wsApi = EthereumWsFactory(
|
val wsApi = EthereumWsFactory(
|
||||||
id, chain,
|
id, chain,
|
||||||
@@ -305,15 +351,45 @@ open class ConfiguredUpstreams(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private fun buildEthereumConnectorFactory(id: String, conn: UpstreamsConfig.EthereumConnection, chain: Chain, urls: ArrayList<URI>, forkChoice: ForkChoice, blockValidator: BlockValidator): EthereumConnectorFactory? {
|
private fun buildEthereumConnectorFactory(
|
||||||
|
id: String,
|
||||||
|
conn: UpstreamsConfig.EthereumConnection,
|
||||||
|
chain: Chain,
|
||||||
|
urls: ArrayList<URI>,
|
||||||
|
forkChoice: ForkChoice,
|
||||||
|
blockValidator: BlockValidator
|
||||||
|
): EthereumConnectorFactory? {
|
||||||
val wsFactoryApi = buildWsFactory(id, chain, conn, urls)
|
val wsFactoryApi = buildWsFactory(id, chain, conn, urls)
|
||||||
val httpFactory = buildHttpFactory(conn, urls)
|
val httpFactory = buildHttpFactory(conn, urls)
|
||||||
log.info("Using ${chain.chainName} upstream, at ${urls.joinToString()}")
|
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()) {
|
if (!connectorFactory.isValid()) {
|
||||||
log.warn("Upstream configuration is invalid (probably no http endpoint)")
|
log.warn("Upstream configuration is invalid (probably no http endpoint)")
|
||||||
return null
|
return null
|
||||||
}
|
}
|
||||||
return connectorFactory
|
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<Int, Int>>(
|
||||||
|
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()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -26,6 +26,7 @@ import java.util.concurrent.atomic.AtomicReference
|
|||||||
|
|
||||||
abstract class DefaultUpstream(
|
abstract class DefaultUpstream(
|
||||||
private val id: String,
|
private val id: String,
|
||||||
|
private val hash: Byte,
|
||||||
defaultLag: Long,
|
defaultLag: Long,
|
||||||
defaultAvail: UpstreamAvailability,
|
defaultAvail: UpstreamAvailability,
|
||||||
private val options: UpstreamsConfig.Options,
|
private val options: UpstreamsConfig.Options,
|
||||||
@@ -36,12 +37,14 @@ abstract class DefaultUpstream(
|
|||||||
|
|
||||||
constructor(
|
constructor(
|
||||||
id: String,
|
id: String,
|
||||||
|
hash: Byte,
|
||||||
options: UpstreamsConfig.Options,
|
options: UpstreamsConfig.Options,
|
||||||
role: UpstreamsConfig.UpstreamRole,
|
role: UpstreamsConfig.UpstreamRole,
|
||||||
targets: CallMethods?
|
targets: CallMethods?
|
||||||
) :
|
) :
|
||||||
this(
|
this(
|
||||||
id,
|
id,
|
||||||
|
hash,
|
||||||
Long.MAX_VALUE,
|
Long.MAX_VALUE,
|
||||||
UpstreamAvailability.UNAVAILABLE,
|
UpstreamAvailability.UNAVAILABLE,
|
||||||
options,
|
options,
|
||||||
@@ -52,12 +55,13 @@ abstract class DefaultUpstream(
|
|||||||
|
|
||||||
constructor(
|
constructor(
|
||||||
id: String,
|
id: String,
|
||||||
|
hash: Byte,
|
||||||
options: UpstreamsConfig.Options,
|
options: UpstreamsConfig.Options,
|
||||||
role: UpstreamsConfig.UpstreamRole,
|
role: UpstreamsConfig.UpstreamRole,
|
||||||
targets: CallMethods?,
|
targets: CallMethods?,
|
||||||
node: QuorumForLabels.QuorumItem?
|
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 status = AtomicReference(Status(defaultLag, defaultAvail, statusByLag(defaultLag, defaultAvail)))
|
||||||
private val statusStream = Sinks.many()
|
private val statusStream = Sinks.many()
|
||||||
@@ -139,6 +143,8 @@ abstract class DefaultUpstream(
|
|||||||
return targets ?: throw IllegalStateException("Methods are not set")
|
return targets ?: throw IllegalStateException("Methods are not set")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
override fun nodeId(): Byte = hash
|
||||||
|
|
||||||
private val quorumByLabel = node?.let { QuorumForLabels(it) }
|
private val quorumByLabel = node?.let { QuorumForLabels(it) }
|
||||||
?: QuorumForLabels(QuorumForLabels.QuorumItem.empty())
|
?: QuorumForLabels(QuorumForLabels.QuorumItem.empty())
|
||||||
|
|
||||||
|
|||||||
@@ -282,6 +282,8 @@ abstract class Multistream(
|
|||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
|
override fun nodeId(): Byte = 0
|
||||||
|
|
||||||
fun printStatus() {
|
fun printStatus() {
|
||||||
var height: Long? = null
|
var height: Long? = null
|
||||||
try {
|
try {
|
||||||
|
|||||||
@@ -396,4 +396,22 @@ class Selector {
|
|||||||
return "Matcher: ${describeInternal()}"
|
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
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -40,4 +40,6 @@ interface Upstream {
|
|||||||
fun isGrpc(): Boolean
|
fun isGrpc(): Boolean
|
||||||
|
|
||||||
fun <T : Upstream> cast(selfType: Class<T>): T
|
fun <T : Upstream> cast(selfType: Class<T>): T
|
||||||
|
|
||||||
|
fun nodeId(): Byte
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -31,7 +31,7 @@ abstract class BitcoinUpstream(
|
|||||||
callMethods: CallMethods,
|
callMethods: CallMethods,
|
||||||
node: QuorumForLabels.QuorumItem,
|
node: QuorumForLabels.QuorumItem,
|
||||||
val esploraClient: EsploraClient? = null
|
val esploraClient: EsploraClient? = null
|
||||||
) : DefaultUpstream(id, options, role, callMethods, node) {
|
) : DefaultUpstream(id, 0.toByte(), options, role, callMethods, node) {
|
||||||
|
|
||||||
constructor(
|
constructor(
|
||||||
id: String,
|
id: String,
|
||||||
|
|||||||
@@ -70,7 +70,14 @@ class DefaultEthereumMethods(
|
|||||||
"eth_feeHistory"
|
"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(
|
private val hardcodedMethods = listOf(
|
||||||
"net_version",
|
"net_version",
|
||||||
@@ -88,6 +95,7 @@ class DefaultEthereumMethods(
|
|||||||
|
|
||||||
override fun getQuorumFor(method: String): CallQuorum {
|
override fun getQuorumFor(method: String): CallQuorum {
|
||||||
return when {
|
return when {
|
||||||
|
filterMethods.contains(method) -> AlwaysQuorum()
|
||||||
hardcodedMethods.contains(method) -> AlwaysQuorum()
|
hardcodedMethods.contains(method) -> AlwaysQuorum()
|
||||||
firstValueMethods.contains(method) -> AlwaysQuorum()
|
firstValueMethods.contains(method) -> AlwaysQuorum()
|
||||||
anyResponseMethods.contains(method) -> NotLaggingQuorum(4)
|
anyResponseMethods.contains(method) -> NotLaggingQuorum(4)
|
||||||
|
|||||||
@@ -57,10 +57,26 @@ class EthereumCallSelector(
|
|||||||
return blockTagSelector(params, 1, head)
|
return blockTagSelector(params, 1, head)
|
||||||
} else if (method == "eth_getStorageAt") {
|
} else if (method == "eth_getStorageAt") {
|
||||||
return blockTagSelector(params, 2, head)
|
return blockTagSelector(params, 2, head)
|
||||||
|
} else if (method == "eth_getFilterChanges") {
|
||||||
|
return sameUpstreamMatcher(params)
|
||||||
}
|
}
|
||||||
return Mono.empty()
|
return Mono.empty()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private fun sameUpstreamMatcher(params: String): Mono<Selector.Matcher> {
|
||||||
|
val list = objectMapper.readerFor(Any::class.java).readValues<Any>(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<Selector.Matcher> {
|
private fun blockTagSelector(params: String, pos: Int, head: Head): Mono<Selector.Matcher> {
|
||||||
val list = objectMapper.readerFor(Any::class.java).readValues<Any>(params).readAll()
|
val list = objectMapper.readerFor(Any::class.java).readValues<Any>(params).readAll()
|
||||||
if (list.size < pos + 1) {
|
if (list.size < pos + 1) {
|
||||||
|
|||||||
@@ -36,13 +36,14 @@ import reactor.core.Disposable
|
|||||||
|
|
||||||
open class EthereumRpcUpstream(
|
open class EthereumRpcUpstream(
|
||||||
id: String,
|
id: String,
|
||||||
|
hash: Byte,
|
||||||
val chain: Chain,
|
val chain: Chain,
|
||||||
options: UpstreamsConfig.Options,
|
options: UpstreamsConfig.Options,
|
||||||
role: UpstreamsConfig.UpstreamRole,
|
role: UpstreamsConfig.UpstreamRole,
|
||||||
targets: CallMethods?,
|
targets: CallMethods?,
|
||||||
private val node: QuorumForLabels.QuorumItem?,
|
private val node: QuorumForLabels.QuorumItem?,
|
||||||
connectorFactory: ConnectorFactory
|
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 log = LoggerFactory.getLogger(EthereumRpcUpstream::class.java)
|
||||||
private val validator: EthereumUpstreamValidator = EthereumUpstreamValidator(this, getOptions())
|
private val validator: EthereumUpstreamValidator = EthereumUpstreamValidator(this, getOptions())
|
||||||
private val connector: EthereumConnector = connectorFactory.create(this, validator, chain)
|
private val connector: EthereumConnector = connectorFactory.create(this, validator, chain)
|
||||||
|
|||||||
@@ -24,11 +24,12 @@ import io.emeraldpay.dshackle.upstream.calls.CallMethods
|
|||||||
|
|
||||||
abstract class EthereumUpstream(
|
abstract class EthereumUpstream(
|
||||||
id: String,
|
id: String,
|
||||||
|
hash: Byte,
|
||||||
options: UpstreamsConfig.Options,
|
options: UpstreamsConfig.Options,
|
||||||
role: UpstreamsConfig.UpstreamRole,
|
role: UpstreamsConfig.UpstreamRole,
|
||||||
targets: CallMethods?,
|
targets: CallMethods?,
|
||||||
private val node: QuorumForLabels.QuorumItem?
|
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) {
|
private val capabilities = if (options.providesBalance != false) {
|
||||||
setOf(Capability.RPC, Capability.BALANCE)
|
setOf(Capability.RPC, Capability.BALANCE)
|
||||||
|
|||||||
@@ -36,13 +36,14 @@ import reactor.core.Disposable
|
|||||||
|
|
||||||
open class EthereumPosRpcUpstream(
|
open class EthereumPosRpcUpstream(
|
||||||
id: String,
|
id: String,
|
||||||
|
hash: Byte,
|
||||||
val chain: Chain,
|
val chain: Chain,
|
||||||
options: UpstreamsConfig.Options,
|
options: UpstreamsConfig.Options,
|
||||||
role: UpstreamsConfig.UpstreamRole,
|
role: UpstreamsConfig.UpstreamRole,
|
||||||
targets: CallMethods?,
|
targets: CallMethods?,
|
||||||
private val node: QuorumForLabels.QuorumItem?,
|
private val node: QuorumForLabels.QuorumItem?,
|
||||||
connectorFactory: ConnectorFactory
|
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 log = LoggerFactory.getLogger(EthereumPosRpcUpstream::class.java)
|
||||||
private val validator: EthereumUpstreamValidator = EthereumUpstreamValidator(this, getOptions())
|
private val validator: EthereumUpstreamValidator = EthereumUpstreamValidator(this, getOptions())
|
||||||
private val connector: EthereumConnector = connectorFactory.create(this, validator, chain)
|
private val connector: EthereumConnector = connectorFactory.create(this, validator, chain)
|
||||||
|
|||||||
@@ -24,11 +24,12 @@ import io.emeraldpay.dshackle.upstream.calls.CallMethods
|
|||||||
|
|
||||||
abstract class EthereumPosUpstream(
|
abstract class EthereumPosUpstream(
|
||||||
id: String,
|
id: String,
|
||||||
|
hash: Byte,
|
||||||
options: UpstreamsConfig.Options,
|
options: UpstreamsConfig.Options,
|
||||||
role: UpstreamsConfig.UpstreamRole,
|
role: UpstreamsConfig.UpstreamRole,
|
||||||
targets: CallMethods?,
|
targets: CallMethods?,
|
||||||
private val node: QuorumForLabels.QuorumItem?
|
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) {
|
private val capabilities = if (options.providesBalance != false) {
|
||||||
setOf(Capability.RPC, Capability.BALANCE)
|
setOf(Capability.RPC, Capability.BALANCE)
|
||||||
|
|||||||
@@ -51,6 +51,7 @@ import java.util.function.Function
|
|||||||
|
|
||||||
open class EthereumGrpcUpstream(
|
open class EthereumGrpcUpstream(
|
||||||
private val parentId: String,
|
private val parentId: String,
|
||||||
|
hash: Byte,
|
||||||
role: UpstreamsConfig.UpstreamRole,
|
role: UpstreamsConfig.UpstreamRole,
|
||||||
private val chain: Chain,
|
private val chain: Chain,
|
||||||
private val remote: ReactorBlockchainGrpc.ReactorBlockchainStub,
|
private val remote: ReactorBlockchainGrpc.ReactorBlockchainStub,
|
||||||
@@ -58,6 +59,7 @@ open class EthereumGrpcUpstream(
|
|||||||
overrideLabels: UpstreamsConfig.Labels?
|
overrideLabels: UpstreamsConfig.Labels?
|
||||||
) : EthereumUpstream(
|
) : EthereumUpstream(
|
||||||
"${parentId}_${chain.chainCode.lowercase(Locale.getDefault())}",
|
"${parentId}_${chain.chainCode.lowercase(Locale.getDefault())}",
|
||||||
|
hash,
|
||||||
UpstreamsConfig.Options.getDefaults(),
|
UpstreamsConfig.Options.getDefaults(),
|
||||||
role,
|
role,
|
||||||
null,
|
null,
|
||||||
|
|||||||
@@ -51,6 +51,7 @@ import java.util.function.Function
|
|||||||
|
|
||||||
open class EthereumPosGrpcUpstream(
|
open class EthereumPosGrpcUpstream(
|
||||||
private val parentId: String,
|
private val parentId: String,
|
||||||
|
hash: Byte,
|
||||||
role: UpstreamsConfig.UpstreamRole,
|
role: UpstreamsConfig.UpstreamRole,
|
||||||
private val chain: Chain,
|
private val chain: Chain,
|
||||||
private val remote: ReactorBlockchainGrpc.ReactorBlockchainStub,
|
private val remote: ReactorBlockchainGrpc.ReactorBlockchainStub,
|
||||||
@@ -59,6 +60,7 @@ open class EthereumPosGrpcUpstream(
|
|||||||
overrideLabels: UpstreamsConfig.Labels?
|
overrideLabels: UpstreamsConfig.Labels?
|
||||||
) : EthereumPosUpstream(
|
) : EthereumPosUpstream(
|
||||||
"${parentId}_${chain.chainCode.lowercase(Locale.getDefault())}",
|
"${parentId}_${chain.chainCode.lowercase(Locale.getDefault())}",
|
||||||
|
hash,
|
||||||
UpstreamsConfig.Options.getDefaults(),
|
UpstreamsConfig.Options.getDefaults(),
|
||||||
role,
|
role,
|
||||||
null, null
|
null, null
|
||||||
|
|||||||
@@ -52,6 +52,7 @@ import kotlin.concurrent.withLock
|
|||||||
|
|
||||||
class GrpcUpstreams(
|
class GrpcUpstreams(
|
||||||
private val id: String,
|
private val id: String,
|
||||||
|
private val hash: Byte,
|
||||||
private val role: UpstreamsConfig.UpstreamRole,
|
private val role: UpstreamsConfig.UpstreamRole,
|
||||||
private val host: String,
|
private val host: String,
|
||||||
private val port: Int,
|
private val port: Int,
|
||||||
@@ -206,7 +207,7 @@ class GrpcUpstreams(
|
|||||||
val current = known[chain]
|
val current = known[chain]
|
||||||
return if (current == null) {
|
return if (current == null) {
|
||||||
val rpcClient = JsonRpcGrpcClient(client!!, chain, metrics)
|
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
|
created.timeout = this.timeout
|
||||||
known[chain] = created
|
known[chain] = created
|
||||||
created.start()
|
created.start()
|
||||||
@@ -222,7 +223,7 @@ class GrpcUpstreams(
|
|||||||
val current = known[chain]
|
val current = known[chain]
|
||||||
return if (current == null) {
|
return if (current == null) {
|
||||||
val rpcClient = JsonRpcGrpcClient(client!!, chain, metrics)
|
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
|
created.timeout = this.timeout
|
||||||
known[chain] = created
|
known[chain] = created
|
||||||
created.start()
|
created.start()
|
||||||
|
|||||||
@@ -401,4 +401,18 @@ class UpstreamsConfigReaderSpec extends Specification {
|
|||||||
act.upstreams.get(0).role == UpstreamsConfig.UpstreamRole.PRIMARY
|
act.upstreams.get(0).role == UpstreamsConfig.UpstreamRole.PRIMARY
|
||||||
act.upstreams.get(1).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"
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -118,7 +118,7 @@ class NativeCallSpec extends Specification {
|
|||||||
def nativeCall = nativeCall()
|
def nativeCall = nativeCall()
|
||||||
nativeCall.quorumReaderFactory = Mock(QuorumReaderFactory) {
|
nativeCall.quorumReaderFactory = Mock(QuorumReaderFactory) {
|
||||||
1 * create(_, _, _) >> Mock(Reader) {
|
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,
|
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"() {
|
def "Parse empty params"() {
|
||||||
setup:
|
setup:
|
||||||
def nativeCall = nativeCall()
|
def nativeCall = nativeCall()
|
||||||
@@ -447,6 +509,64 @@ class NativeCallSpec extends Specification {
|
|||||||
act.payload.method == "eth_test"
|
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
|
@Ignore
|
||||||
//TODO
|
//TODO
|
||||||
def "Calls cache before remote"() {
|
def "Calls cache before remote"() {
|
||||||
|
|||||||
@@ -58,4 +58,37 @@ class ConfiguredUpstreamsSpec extends Specification {
|
|||||||
act instanceof ManagedCallMethods
|
act instanceof ManagedCallMethods
|
||||||
new String(act.executeHardcoded("foo_bar")) == "\"static_response\""
|
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
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -69,7 +69,7 @@ class EthereumPosRpcUpstreamMock extends EthereumPosRpcUpstream {
|
|||||||
}
|
}
|
||||||
|
|
||||||
EthereumPosRpcUpstreamMock(@NotNull String id, @NotNull Chain chain, @NotNull Reader<JsonRpcRequest, JsonRpcResponse> api, CallMethods methods, Map<String, String> labels) {
|
EthereumPosRpcUpstreamMock(@NotNull String id, @NotNull Chain chain, @NotNull Reader<JsonRpcRequest, JsonRpcResponse> api, CallMethods methods, Map<String, String> labels) {
|
||||||
super(id, chain,
|
super(id, (byte)id.hashCode(), chain,
|
||||||
UpstreamsConfig.Options.getDefaults(),
|
UpstreamsConfig.Options.getDefaults(),
|
||||||
UpstreamsConfig.UpstreamRole.PRIMARY,
|
UpstreamsConfig.UpstreamRole.PRIMARY,
|
||||||
methods,
|
methods,
|
||||||
|
|||||||
@@ -60,7 +60,7 @@ class EthereumRpcUpstreamMock extends EthereumRpcUpstream {
|
|||||||
}
|
}
|
||||||
|
|
||||||
EthereumRpcUpstreamMock(@NotNull String id, @NotNull Chain chain, @NotNull Reader<JsonRpcRequest, JsonRpcResponse> api, CallMethods methods) {
|
EthereumRpcUpstreamMock(@NotNull String id, @NotNull Chain chain, @NotNull Reader<JsonRpcRequest, JsonRpcResponse> api, CallMethods methods) {
|
||||||
super(id, chain,
|
super(id, id.hashCode().byteValue(), chain,
|
||||||
UpstreamsConfig.Options.getDefaults(),
|
UpstreamsConfig.Options.getDefaults(),
|
||||||
UpstreamsConfig.UpstreamRole.PRIMARY,
|
UpstreamsConfig.UpstreamRole.PRIMARY,
|
||||||
methods,
|
methods,
|
||||||
|
|||||||
@@ -51,6 +51,7 @@ class FilteredApisSpec extends Specification {
|
|||||||
def connectorFactory = new EthereumConnectorFactory(false, null, httpFactory, new MostWorkForkChoice(), BlockValidator.@Companion.ALWAYS_VALID)
|
def connectorFactory = new EthereumConnectorFactory(false, null, httpFactory, new MostWorkForkChoice(), BlockValidator.@Companion.ALWAYS_VALID)
|
||||||
new EthereumRpcUpstream(
|
new EthereumRpcUpstream(
|
||||||
"test",
|
"test",
|
||||||
|
(byte)123,
|
||||||
Chain.ETHEREUM,
|
Chain.ETHEREUM,
|
||||||
new UpstreamsConfig.Options(),
|
new UpstreamsConfig.Options(),
|
||||||
UpstreamsConfig.UpstreamRole.PRIMARY,
|
UpstreamsConfig.UpstreamRole.PRIMARY,
|
||||||
|
|||||||
@@ -349,4 +349,27 @@ class SelectorSpec extends Specification {
|
|||||||
!act
|
!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
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -179,4 +179,33 @@ class EthereumCallSelectorSpec extends Specification {
|
|||||||
then:
|
then:
|
||||||
act == new Selector.HeightMatcher(100)
|
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
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -30,6 +30,7 @@ class EthereumDirectReaderSpec extends Specification {
|
|||||||
|
|
||||||
String hash1 = "0x40d15edaff9acdabd2a1c96fd5f683b3300aad34e7015f34def3c56ba8a7ffb5"
|
String hash1 = "0x40d15edaff9acdabd2a1c96fd5f683b3300aad34e7015f34def3c56ba8a7ffb5"
|
||||||
String address1 = "0xe0aadb0a012dbcdc529c4c743d3e0385a0b54d3d"
|
String address1 = "0xe0aadb0a012dbcdc529c4c743d3e0385a0b54d3d"
|
||||||
|
List<Byte> resolvers = Collections.singletonList((byte)1)
|
||||||
|
|
||||||
def "Reads block by hash"() {
|
def "Reads block by hash"() {
|
||||||
setup:
|
setup:
|
||||||
@@ -53,8 +54,7 @@ class EthereumDirectReaderSpec extends Specification {
|
|||||||
1 * create(_, _, _) >> Mock(Reader) {
|
1 * create(_, _, _) >> Mock(Reader) {
|
||||||
1 * read(new JsonRpcRequest("eth_getBlockByHash", [hash1, false])) >> Mono.just(
|
1 * read(new JsonRpcRequest("eth_getBlockByHash", [hash1, false])) >> Mono.just(
|
||||||
new QuorumRpcReader.Result(
|
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 * create(_, _, _) >> Mock(Reader) {
|
||||||
1 * read(new JsonRpcRequest("eth_getBlockByHash", [hash1, false])) >> Mono.just(
|
1 * read(new JsonRpcRequest("eth_getBlockByHash", [hash1, false])) >> Mono.just(
|
||||||
new QuorumRpcReader.Result(
|
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 * create(_, _, _) >> Mock(Reader) {
|
||||||
1 * read(new JsonRpcRequest("eth_getBlockByNumber", ["0x64", false])) >> Mono.just(
|
1 * read(new JsonRpcRequest("eth_getBlockByNumber", ["0x64", false])) >> Mono.just(
|
||||||
new QuorumRpcReader.Result(
|
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 * create(_, _, _) >> Mock(Reader) {
|
||||||
1 * read(new JsonRpcRequest("eth_getTransactionByHash", [hash1])) >> Mono.just(
|
1 * read(new JsonRpcRequest("eth_getTransactionByHash", [hash1])) >> Mono.just(
|
||||||
new QuorumRpcReader.Result(
|
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 * create(_, _, _) >> Mock(Reader) {
|
||||||
1 * read(new JsonRpcRequest("eth_getTransactionByHash", [hash1])) >> Mono.just(
|
1 * read(new JsonRpcRequest("eth_getTransactionByHash", [hash1])) >> Mono.just(
|
||||||
new QuorumRpcReader.Result(
|
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 * create(_, _, _) >> Mock(Reader) {
|
||||||
1 * read(new JsonRpcRequest("eth_getBalance", [address1, "latest"])) >> Mono.just(
|
1 * read(new JsonRpcRequest("eth_getBalance", [address1, "latest"])) >> Mono.just(
|
||||||
new QuorumRpcReader.Result(
|
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 * create(_, _, _) >> Mock(Reader) {
|
||||||
1 * read(new JsonRpcRequest("eth_getBalance", [address1, "0xa8c9bb"])) >> Mono.just(
|
1 * read(new JsonRpcRequest("eth_getBalance", [address1, "0xa8c9bb"])) >> Mono.just(
|
||||||
new QuorumRpcReader.Result(
|
new QuorumRpcReader.Result(
|
||||||
Global.objectMapper.writeValueAsBytes("0x100"), null, 1
|
Global.objectMapper.writeValueAsBytes("0x100"), null, 1, resolvers
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -50,6 +50,8 @@ class EthereumGrpcUpstreamSpec extends Specification {
|
|||||||
Counter.builder("test2").register(TestingCommons.meterRegistry)
|
Counter.builder("test2").register(TestingCommons.meterRegistry)
|
||||||
)
|
)
|
||||||
|
|
||||||
|
def hash = (byte)123
|
||||||
|
|
||||||
def "Subscribe to head"() {
|
def "Subscribe to head"() {
|
||||||
setup:
|
setup:
|
||||||
def callData = [:]
|
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.setLag(0)
|
||||||
upstream.update(BlockchainOuterClass.DescribeChain.newBuilder()
|
upstream.update(BlockchainOuterClass.DescribeChain.newBuilder()
|
||||||
.setStatus(BlockchainOuterClass.ChainStatus.newBuilder().setQuorum(1).setAvailabilityValue(UpstreamAvailability.OK.grpcId))
|
.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.setLag(0)
|
||||||
upstream.update(BlockchainOuterClass.DescribeChain.newBuilder()
|
upstream.update(BlockchainOuterClass.DescribeChain.newBuilder()
|
||||||
.setStatus(BlockchainOuterClass.ChainStatus.newBuilder().setQuorum(1).setAvailabilityValue(UpstreamAvailability.OK.grpcId))
|
.setStatus(BlockchainOuterClass.ChainStatus.newBuilder().setQuorum(1).setAvailabilityValue(UpstreamAvailability.OK.grpcId))
|
||||||
@@ -201,7 +203,7 @@ class EthereumGrpcUpstreamSpec extends Specification {
|
|||||||
finished.complete(true)
|
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.setLag(0)
|
||||||
upstream.update(BlockchainOuterClass.DescribeChain.newBuilder()
|
upstream.update(BlockchainOuterClass.DescribeChain.newBuilder()
|
||||||
.setStatus(BlockchainOuterClass.ChainStatus.newBuilder().setQuorum(1).setAvailabilityValue(UpstreamAvailability.OK.grpcId))
|
.setStatus(BlockchainOuterClass.ChainStatus.newBuilder().setQuorum(1).setAvailabilityValue(UpstreamAvailability.OK.grpcId))
|
||||||
|
|||||||
34
src/test/resources/upstreams-node-id.yaml
Normal file
34
src/test/resources/upstreams-node-id.yaml
Normal file
@@ -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
|
||||||
Reference in New Issue
Block a user