Merge pull request #101 from p2p-org/subscriptions-refactoring

small renaming refactoring and supported subscription feature
This commit is contained in:
a10zn8
2022-12-28 15:42:11 +04:00
committed by GitHub
60 changed files with 325 additions and 261 deletions

View File

@@ -126,7 +126,7 @@ class QuorumRpcReader(
} }
fun callApi(api: Upstream, key: JsonRpcRequest): Mono<Tuple3<ByteArray, Optional<ResponseSigner.Signature>, Upstream>> { fun callApi(api: Upstream, key: JsonRpcRequest): Mono<Tuple3<ByteArray, Optional<ResponseSigner.Signature>, Upstream>> {
return api.getApi() return api.getIngressReader()
.read(key) .read(key)
.flatMap { response -> .flatMap { response ->
response.requireResult() response.requireResult()

View File

@@ -0,0 +1,6 @@
package io.emeraldpay.dshackle.reader
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
typealias JsonRpcReader = Reader<JsonRpcRequest, JsonRpcResponse>

View File

@@ -36,13 +36,14 @@ class Describe(
return requestMono.map { _ -> return requestMono.map { _ ->
val resp = BlockchainOuterClass.DescribeResponse.newBuilder() val resp = BlockchainOuterClass.DescribeResponse.newBuilder()
multistreamHolder.getAvailable().forEach { chain -> multistreamHolder.getAvailable().forEach { chain ->
multistreamHolder.getUpstream(chain)?.let { chainUpstreams -> multistreamHolder.getUpstream(chain).let { chainUpstreams ->
val status = subscribeStatus.chainStatus(chain, chainUpstreams.getStatus(), chainUpstreams) val status = subscribeStatus.chainStatus(chain, chainUpstreams.getStatus(), chainUpstreams)
val targets = chainUpstreams.getMethods().getSupportedMethods() val targets = chainUpstreams.getMethods().getSupportedMethods()
val capabilities: MutableSet<Capability> = mutableSetOf() val capabilities: MutableSet<Capability> = mutableSetOf()
val chainDescription = BlockchainOuterClass.DescribeChain.newBuilder() val chainDescription = BlockchainOuterClass.DescribeChain.newBuilder()
.setChain(Common.ChainRef.forNumber(chain.id)) .setChain(Common.ChainRef.forNumber(chain.id))
.addAllSupportedMethods(targets) .addAllSupportedMethods(targets)
.addAllSupportedSubscriptions(chainUpstreams.getEgressSubscription().getAvailableTopics())
.setStatus(status) .setStatus(status)
.setCurrentHeight(chainUpstreams.getHead().getCurrentHeight() ?: 0) .setCurrentHeight(chainUpstreams.getHead().getCurrentHeight() ?: 0)
chainUpstreams.getAll().let { ups -> chainUpstreams.getAll().let { ups ->

View File

@@ -274,7 +274,7 @@ open class NativeCall(
if (method in DefaultEthereumMethods.newFilterMethods) CreateFilterDecorator() else NoneResultDecorator() if (method in DefaultEthereumMethods.newFilterMethods) CreateFilterDecorator() else NoneResultDecorator()
fun fetch(ctx: ValidCallContext<ParsedCallDetails>): Mono<CallResult> { fun fetch(ctx: ValidCallContext<ParsedCallDetails>): Mono<CallResult> {
return ctx.upstream.getRoutedApi(localRouterEnabled) return ctx.upstream.getLocalReader(localRouterEnabled)
.flatMap { api -> .flatMap { api ->
api.read(JsonRpcRequest(ctx.payload.method, ctx.payload.params, ctx.nonce, ctx.forwardedSelector)) api.read(JsonRpcRequest(ctx.payload.method, ctx.payload.params, ctx.nonce, ctx.forwardedSelector))
.flatMap(JsonRpcResponse::requireResult) .flatMap(JsonRpcResponse::requireResult)

View File

@@ -97,7 +97,7 @@ open class NativeSubscribe(
} }
open fun subscribe(chain: Chain, method: String, params: Any?, matcher: Selector.Matcher): Flux<out Any> = open fun subscribe(chain: Chain, method: String, params: Any?, matcher: Selector.Matcher): Flux<out Any> =
getUpstream(chain).getSubscriptionApi().subscribe(method, params, matcher) getUpstream(chain).getEgressSubscription().subscribe(method, params, matcher)
private fun getUpstream(chain: Chain): EthereumLikeMultistream = private fun getUpstream(chain: Chain): EthereumLikeMultistream =
multistreamHolder.getUpstream(chain).let { it as EthereumLikeMultistream } multistreamHolder.getUpstream(chain).let { it as EthereumLikeMultistream }

View File

@@ -88,7 +88,7 @@ class TrackERC20Address(
val asset = request.asset.code.lowercase(Locale.getDefault()) val asset = request.asset.code.lowercase(Locale.getDefault())
val tokenDefinition = tokens[TokenId(chain, asset)] ?: return Flux.empty() val tokenDefinition = tokens[TokenId(chain, asset)] ?: return Flux.empty()
val logs = getUpstream(chain) val logs = getUpstream(chain)
.getSubscriptionApi().logs .getEgressSubscription().logs
.create( .create(
listOf(tokenDefinition.token.contract), listOf(tokenDefinition.token.contract),
listOf(EventId.fromSignature("Transfer", "address", "address", "uint256")), listOf(EventId.fromSignature("Transfer", "address", "address", "uint256")),

View File

@@ -22,7 +22,7 @@ import io.emeraldpay.dshackle.Chain
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
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.* import io.emeraldpay.dshackle.upstream.*
import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinRpcHead import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinRpcHead
import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinRpcUpstream import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinRpcUpstream
@@ -41,8 +41,6 @@ import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice
import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice
import io.emeraldpay.dshackle.upstream.forkchoice.NoChoiceWithPriorityForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.NoChoiceWithPriorityForkChoice
import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstreams import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstreams
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.boot.ApplicationArguments import org.springframework.boot.ApplicationArguments
import org.springframework.boot.ApplicationRunner import org.springframework.boot.ApplicationRunner
@@ -226,7 +224,7 @@ open class ConfiguredUpstreams(
log.warn("Upstream doesn't have API configuration") log.warn("Upstream doesn't have API configuration")
return null return null
} }
val directApi: Reader<JsonRpcRequest, JsonRpcResponse> = httpFactory.create(config.id, chain) val directApi: JsonRpcReader = httpFactory.create(config.id, chain)
val esplora = conn.esplora?.let { endpoint -> val esplora = conn.esplora?.let { endpoint ->
val tls = endpoint.tls?.let { tls -> val tls = endpoint.tls?.let { tls ->
tls.ca?.let { ca -> tls.ca?.let { ca ->

View File

@@ -0,0 +1,33 @@
/**
* Copyright (c) 2022 EmeraldPay, Inc
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package io.emeraldpay.dshackle.upstream
import reactor.core.publisher.Flux
interface EgressSubscription {
fun getAvailableTopics(): List<String>
fun subscribe(topic: String, params: Any?, matcher: Selector.Matcher): Flux<out Any>
}
class EmptyEgressSubscription : EgressSubscription {
override fun getAvailableTopics(): List<String> {
return emptyList()
}
override fun subscribe(topic: String, params: Any?, matcher: Selector.Matcher): Flux<out Any> {
return Flux.empty()
}
}

View File

@@ -15,13 +15,6 @@
*/ */
package io.emeraldpay.dshackle.upstream package io.emeraldpay.dshackle.upstream
open class NoUpstreamSubscriptions : UpstreamSubscriptions { interface HasEgressSubscription {
fun getEgressSubscription(): EgressSubscription
companion object {
val DEFAULT = NoUpstreamSubscriptions()
}
override fun <T> get(method: String): SubscriptionConnect<T>? {
return null
}
} }

View File

@@ -1,10 +1,8 @@
package io.emeraldpay.dshackle.upstream package io.emeraldpay.dshackle.upstream
import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
interface HttpFactory { interface HttpFactory {
fun create(id: String?, chain: Chain): Reader<JsonRpcRequest, JsonRpcResponse> fun create(id: String?, chain: Chain): JsonRpcReader
} }

View File

@@ -2,10 +2,8 @@ package io.emeraldpay.dshackle.upstream
import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.config.AuthConfig import io.emeraldpay.dshackle.config.AuthConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcHttpClient import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcHttpClient
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import io.emeraldpay.dshackle.upstream.rpcclient.RpcMetrics import io.emeraldpay.dshackle.upstream.rpcclient.RpcMetrics
import io.micrometer.core.instrument.Counter import io.micrometer.core.instrument.Counter
import io.micrometer.core.instrument.Metrics import io.micrometer.core.instrument.Metrics
@@ -17,7 +15,7 @@ open class HttpRpcFactory(
private val basicAuth: AuthConfig.ClientBasicAuth?, private val basicAuth: AuthConfig.ClientBasicAuth?,
private val tls: ByteArray? private val tls: ByteArray?
) : HttpFactory { ) : HttpFactory {
override fun create(id: String?, chain: Chain): Reader<JsonRpcRequest, JsonRpcResponse> { override fun create(id: String?, chain: Chain): JsonRpcReader {
val metricsTags = listOf( val metricsTags = listOf(
// "unknown" is not supposed to happen // "unknown" is not supposed to happen
Tag.of("upstream", id ?: "unknown"), Tag.of("upstream", id ?: "unknown"),

View File

@@ -18,7 +18,8 @@ package io.emeraldpay.dshackle.upstream
/** /**
* Subscriptions available on the current upstream * Subscriptions available on the current upstream
*/ */
interface UpstreamSubscriptions { interface IngressSubscription {
fun <T> get(method: String): SubscriptionConnect<T>? fun getAvailableTopics(): List<String>
fun <T> get(topic: String): SubscriptionConnect<T>?
} }

View File

@@ -20,12 +20,10 @@ import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.cache.CachesEnabled import io.emeraldpay.dshackle.cache.CachesEnabled
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.startup.UpstreamChangeEvent import io.emeraldpay.dshackle.startup.UpstreamChangeEvent
import io.emeraldpay.dshackle.upstream.calls.AggregatedCallMethods import io.emeraldpay.dshackle.upstream.calls.AggregatedCallMethods
import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import io.micrometer.core.instrument.Gauge import io.micrometer.core.instrument.Gauge
import io.micrometer.core.instrument.Meter import io.micrometer.core.instrument.Meter
import io.micrometer.core.instrument.Metrics import io.micrometer.core.instrument.Metrics
@@ -55,7 +53,7 @@ abstract class Multistream(
val chain: Chain, val chain: Chain,
private val upstreams: MutableList<Upstream>, private val upstreams: MutableList<Upstream>,
val caches: Caches, val caches: Caches,
) : Upstream, Lifecycle { ) : Upstream, Lifecycle, HasEgressSubscription {
companion object { companion object {
private val log = LoggerFactory.getLogger(Multistream::class.java) private val log = LoggerFactory.getLogger(Multistream::class.java)
@@ -175,9 +173,9 @@ abstract class Multistream(
/** /**
* Finds an API that leverages caches and other optimizations/transformations of the request. * Finds an API that leverages caches and other optimizations/transformations of the request.
*/ */
abstract fun getRoutedApi(localEnabled: Boolean): Mono<Reader<JsonRpcRequest, JsonRpcResponse>> abstract fun getLocalReader(localEnabled: Boolean): Mono<JsonRpcReader>
override fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse> { override fun getIngressReader(): JsonRpcReader {
throw NotImplementedError("Immediate direct API is not implemented for Aggregated Upstream") throw NotImplementedError("Immediate direct API is not implemented for Aggregated Upstream")
} }

View File

@@ -0,0 +1,31 @@
/**
* Copyright (c) 2022 EmeraldPay, Inc
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package io.emeraldpay.dshackle.upstream
open class NoIngressSubscription : IngressSubscription {
companion object {
val DEFAULT = NoIngressSubscription()
}
override fun getAvailableTopics(): List<String> {
return listOf()
}
override fun <T> get(topic: String): SubscriptionConnect<T>? {
return null
}
}

View File

@@ -17,10 +17,8 @@
package io.emeraldpay.dshackle.upstream package io.emeraldpay.dshackle.upstream
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
interface Upstream { interface Upstream {
@@ -28,7 +26,12 @@ interface Upstream {
fun getStatus(): UpstreamAvailability fun getStatus(): UpstreamAvailability
fun observeStatus(): Flux<UpstreamAvailability> fun observeStatus(): Flux<UpstreamAvailability>
fun getHead(): Head fun getHead(): Head
fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse>
/**
* Get an actual reader that access the current upstream
*/
fun getIngressReader(): JsonRpcReader
fun getOptions(): UpstreamsConfig.Options fun getOptions(): UpstreamsConfig.Options
fun getRole(): UpstreamsConfig.UpstreamRole fun getRole(): UpstreamsConfig.UpstreamRole
fun setLag(lag: Long) fun setLag(lag: Long)

View File

@@ -18,13 +18,11 @@ package io.emeraldpay.dshackle.upstream.bitcoin
import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.* import io.emeraldpay.dshackle.upstream.*
import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.calls.DefaultBitcoinMethods import io.emeraldpay.dshackle.upstream.calls.DefaultBitcoinMethods
import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
@@ -98,15 +96,15 @@ open class BitcoinMultistream(
/** /**
* Finds an API that executed directly on a remote. * Finds an API that executed directly on a remote.
*/ */
open fun getDirectApi(matcher: Selector.Matcher): Mono<Reader<JsonRpcRequest, JsonRpcResponse>> { open fun getDirectApi(matcher: Selector.Matcher): Mono<JsonRpcReader> {
val apis = getApiSource(matcher) val apis = getApiSource(matcher)
apis.request(1) apis.request(1)
return Mono.from(apis) return Mono.from(apis)
.map(Upstream::getApi) .map(Upstream::getIngressReader)
.switchIfEmpty(Mono.error(Exception("No API available for $chain"))) .switchIfEmpty(Mono.error(Exception("No API available for $chain")))
} }
override fun getRoutedApi(localEnabled: Boolean): Mono<Reader<JsonRpcRequest, JsonRpcResponse>> { override fun getLocalReader(localEnabled: Boolean): Mono<JsonRpcReader> {
return Mono.just(callRouter) return Mono.just(callRouter)
} }
@@ -144,6 +142,10 @@ open class BitcoinMultistream(
return this as T return this as T
} }
override fun getEgressSubscription(): EgressSubscription {
return EmptyEgressSubscription()
}
override fun isRunning(): Boolean { override fun isRunning(): Boolean {
return super.isRunning() || reader.isRunning() return super.isRunning() || reader.isRunning()
} }

View File

@@ -16,7 +16,7 @@
package io.emeraldpay.dshackle.upstream.bitcoin package io.emeraldpay.dshackle.upstream.bitcoin
import io.emeraldpay.dshackle.Defaults import io.emeraldpay.dshackle.Defaults
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.AbstractHead import io.emeraldpay.dshackle.upstream.AbstractHead
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Lifecycle
@@ -33,7 +33,7 @@ import java.time.Duration
import java.util.concurrent.Executors import java.util.concurrent.Executors
class BitcoinRpcHead( class BitcoinRpcHead(
private val api: Reader<JsonRpcRequest, JsonRpcResponse>, private val api: JsonRpcReader,
private val extractBlock: ExtractBlock, private val extractBlock: ExtractBlock,
private val interval: Duration = Duration.ofSeconds(15) private val interval: Duration = Duration.ofSeconds(15)
) : Head, AbstractHead(MostWorkForkChoice(), awaitHeadTimeoutMs = 1200_000), Lifecycle { ) : Head, AbstractHead(MostWorkForkChoice(), awaitHeadTimeoutMs = 1200_000), Lifecycle {

View File

@@ -17,7 +17,7 @@ package io.emeraldpay.dshackle.upstream.bitcoin
import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.startup.QuorumForLabels import io.emeraldpay.dshackle.startup.QuorumForLabels
import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.Capability
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
@@ -25,15 +25,13 @@ import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.UpstreamAvailability
import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import reactor.core.Disposable import reactor.core.Disposable
open class BitcoinRpcUpstream( open class BitcoinRpcUpstream(
id: String, id: String,
chain: Chain, chain: Chain,
private val directApi: Reader<JsonRpcRequest, JsonRpcResponse>, private val directApi: JsonRpcReader,
private val head: Head, private val head: Head,
options: UpstreamsConfig.Options, options: UpstreamsConfig.Options,
role: UpstreamsConfig.UpstreamRole, role: UpstreamsConfig.UpstreamRole,
@@ -65,7 +63,7 @@ open class BitcoinRpcUpstream(
return head return head
} }
override fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse> { override fun getIngressReader(): JsonRpcReader {
return directApi return directApi
} }

View File

@@ -16,7 +16,7 @@
package io.emeraldpay.dshackle.upstream.bitcoin package io.emeraldpay.dshackle.upstream.bitcoin
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.UpstreamAvailability
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
@@ -29,7 +29,7 @@ import java.time.Duration
import java.util.concurrent.Executors import java.util.concurrent.Executors
class BitcoinUpstreamValidator( class BitcoinUpstreamValidator(
private val api: Reader<JsonRpcRequest, JsonRpcResponse>, private val api: JsonRpcReader,
private val options: UpstreamsConfig.Options private val options: UpstreamsConfig.Options
) { ) {

View File

@@ -2,7 +2,7 @@ package io.emeraldpay.dshackle.upstream.bitcoin
import io.emeraldpay.dshackle.Defaults import io.emeraldpay.dshackle.Defaults
import io.emeraldpay.dshackle.data.BlockContainer import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.AbstractHead import io.emeraldpay.dshackle.upstream.AbstractHead
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Lifecycle
@@ -19,7 +19,7 @@ import java.time.Duration
class BitcoinZMQHead( class BitcoinZMQHead(
private val server: ZMQServer, private val server: ZMQServer,
private val api: Reader<JsonRpcRequest, JsonRpcResponse>, private val api: JsonRpcReader,
private val extractBlock: ExtractBlock, private val extractBlock: ExtractBlock,
) : Head, AbstractHead(MostWorkForkChoice(), awaitHeadTimeoutMs = 1200_000), Lifecycle { ) : Head, AbstractHead(MostWorkForkChoice(), awaitHeadTimeoutMs = 1200_000), Lifecycle {

View File

@@ -17,7 +17,7 @@ package io.emeraldpay.dshackle.upstream.bitcoin
import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.SilentException import io.emeraldpay.dshackle.SilentException
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.bitcoin.data.RpcUnspent import io.emeraldpay.dshackle.upstream.bitcoin.data.RpcUnspent
import io.emeraldpay.dshackle.upstream.bitcoin.data.SimpleUnspent import io.emeraldpay.dshackle.upstream.bitcoin.data.SimpleUnspent
import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.dshackle.upstream.calls.CallMethods
@@ -39,7 +39,7 @@ import reactor.core.publisher.Mono
class LocalCallRouter( class LocalCallRouter(
private val methods: CallMethods, private val methods: CallMethods,
private val reader: BitcoinReader, private val reader: BitcoinReader,
) : Reader<JsonRpcRequest, JsonRpcResponse> { ) : JsonRpcReader {
companion object { companion object {
private val log = LoggerFactory.getLogger(LocalCallRouter::class.java) private val log = LoggerFactory.getLogger(LocalCallRouter::class.java)

View File

@@ -17,13 +17,12 @@ package io.emeraldpay.dshackle.upstream.ethereum
import io.emeraldpay.dshackle.Defaults import io.emeraldpay.dshackle.Defaults
import io.emeraldpay.dshackle.data.BlockContainer import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.AbstractHead import io.emeraldpay.dshackle.upstream.AbstractHead
import io.emeraldpay.dshackle.upstream.BlockValidator import io.emeraldpay.dshackle.upstream.BlockValidator
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import io.emeraldpay.etherjar.hex.HexQuantity import io.emeraldpay.etherjar.hex.HexQuantity
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
@@ -38,7 +37,7 @@ open class DefaultEthereumHead(
private val log = LoggerFactory.getLogger(DefaultEthereumHead::class.java) private val log = LoggerFactory.getLogger(DefaultEthereumHead::class.java)
} }
fun getLatestBlock(api: Reader<JsonRpcRequest, JsonRpcResponse>): Mono<BlockContainer> { fun getLatestBlock(api: JsonRpcReader): Mono<BlockContainer> {
return api.read(JsonRpcRequest("eth_blockNumber", emptyList())) return api.read(JsonRpcRequest("eth_blockNumber", emptyList()))
.subscribeOn(EthereumRpcHead.scheduler) .subscribeOn(EthereumRpcHead.scheduler)
.timeout(Defaults.timeout, Mono.error(Exception("Block number not received"))) .timeout(Defaults.timeout, Mono.error(Exception("Block number not received")))

View File

@@ -59,7 +59,7 @@ open class ERC20Balance {
open fun getBalance(upstream: EthereumPosRpcUpstream, token: ERC20Token, address: Address): Mono<BigInteger> { open fun getBalance(upstream: EthereumPosRpcUpstream, token: ERC20Token, address: Address): Mono<BigInteger> {
return upstream return upstream
.getApi() .getIngressReader()
.read(prepareEthCall(token, address, upstream.getHead())) .read(prepareEthCall(token, address, upstream.getHead()))
.flatMap(JsonRpcResponse::requireStringResult) .flatMap(JsonRpcResponse::requireStringResult)
.map { Hex32.from(it).asQuantity().value } .map { Hex32.from(it).asQuantity().value }

View File

@@ -1,5 +1,6 @@
package io.emeraldpay.dshackle.upstream.ethereum package io.emeraldpay.dshackle.upstream.ethereum
import io.emeraldpay.dshackle.upstream.EgressSubscription
import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.ConnectLogs import io.emeraldpay.dshackle.upstream.ethereum.subscribe.ConnectLogs
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.ConnectNewHeads import io.emeraldpay.dshackle.upstream.ethereum.subscribe.ConnectNewHeads
@@ -10,13 +11,13 @@ import io.emeraldpay.etherjar.hex.Hex32
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
open class EthereumSubscriptionApi( open class EthereumEgressSubscription(
val upstream: EthereumLikeMultistream, val upstream: EthereumLikeMultistream,
val pendingTxesSource: PendingTxesSource val pendingTxesSource: PendingTxesSource?
) { ) : EgressSubscription {
companion object { companion object {
private val log = LoggerFactory.getLogger(EthereumSubscriptionApi::class.java) private val log = LoggerFactory.getLogger(EthereumEgressSubscription::class.java)
const val METHOD_NEW_HEADS = "newHeads" const val METHOD_NEW_HEADS = "newHeads"
const val METHOD_LOGS = "logs" const val METHOD_LOGS = "logs"
@@ -24,16 +25,30 @@ open class EthereumSubscriptionApi(
const val METHOD_PENDING_TXES = "newPendingTransactions" const val METHOD_PENDING_TXES = "newPendingTransactions"
} }
private val availableTopics = listOf(
METHOD_NEW_HEADS,
METHOD_LOGS,
METHOD_SYNCING,
).let {
if (pendingTxesSource != null) {
it + METHOD_PENDING_TXES
} else {
it
}
}
private val newHeads = ConnectNewHeads(upstream) private val newHeads = ConnectNewHeads(upstream)
open val logs = ConnectLogs(upstream) open val logs = ConnectLogs(upstream)
private val syncing = ConnectSyncing(upstream) private val syncing = ConnectSyncing(upstream)
override fun getAvailableTopics() = availableTopics
@Suppress("UNCHECKED_CAST") @Suppress("UNCHECKED_CAST")
open fun subscribe(method: String, params: Any?, matcher: Selector.Matcher): Flux<out Any> { override fun subscribe(topic: String, params: Any?, matcher: Selector.Matcher): Flux<out Any> {
if (method == METHOD_NEW_HEADS) { if (topic == METHOD_NEW_HEADS) {
return newHeads.connect(matcher) return newHeads.connect(matcher)
} }
if (method == METHOD_LOGS) { if (topic == METHOD_LOGS) {
val paramsMap = try { val paramsMap = try {
if (params != null && Map::class.java.isAssignableFrom(params.javaClass)) { if (params != null && Map::class.java.isAssignableFrom(params.javaClass)) {
readLogsRequest(params as Map<String, Any?>) readLogsRequest(params as Map<String, Any?>)
@@ -41,17 +56,17 @@ open class EthereumSubscriptionApi(
LogsRequest(emptyList(), emptyList()) LogsRequest(emptyList(), emptyList())
} }
} catch (t: Throwable) { } catch (t: Throwable) {
return Flux.error(UnsupportedOperationException("Invalid parameter for $method. Error: ${t.message}")) return Flux.error(UnsupportedOperationException("Invalid parameter for $topic. Error: ${t.message}"))
} }
return logs.create(paramsMap.address, paramsMap.topics).connect(matcher) return logs.create(paramsMap.address, paramsMap.topics).connect(matcher)
} }
if (method == METHOD_SYNCING) { if (topic == METHOD_SYNCING) {
return syncing.connect(matcher) return syncing.connect(matcher)
} }
if (method == METHOD_PENDING_TXES) { if (topic == METHOD_PENDING_TXES) {
return pendingTxesSource.connect(matcher) return pendingTxesSource?.connect(matcher) ?: Flux.empty()
} }
return Flux.error(UnsupportedOperationException("Method $method is not supported")) return Flux.error(UnsupportedOperationException("Method $topic is not supported"))
} }
data class LogsRequest( data class LogsRequest(

View File

@@ -15,10 +15,10 @@
*/ */
package io.emeraldpay.dshackle.upstream.ethereum package io.emeraldpay.dshackle.upstream.ethereum
import io.emeraldpay.dshackle.upstream.UpstreamSubscriptions import io.emeraldpay.dshackle.upstream.IngressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.PendingTxesSource import io.emeraldpay.dshackle.upstream.ethereum.subscribe.PendingTxesSource
interface EthereumUpstreamSubscriptions : UpstreamSubscriptions { interface EthereumIngressSubscription : IngressSubscription {
fun getPendingTxes(): PendingTxesSource? fun getPendingTxes(): PendingTxesSource?
} }

View File

@@ -1,14 +1,14 @@
package io.emeraldpay.dshackle.upstream.ethereum package io.emeraldpay.dshackle.upstream.ethereum
import io.emeraldpay.api.proto.BlockchainOuterClass import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.dshackle.upstream.HasEgressSubscription
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstream
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
interface EthereumLikeMultistream : Upstream { interface EthereumLikeMultistream : Upstream, HasEgressSubscription {
fun getReader(): EthereumCachingReader fun getReader(): EthereumCachingReader
fun getSubscriptionApi(): EthereumSubscriptionApi
fun getHead(mather: Selector.Matcher): Head fun getHead(mather: Selector.Matcher): Head

View File

@@ -17,7 +17,7 @@ package io.emeraldpay.dshackle.upstream.ethereum
import io.emeraldpay.dshackle.data.BlockId import io.emeraldpay.dshackle.data.BlockId
import io.emeraldpay.dshackle.data.TxId import io.emeraldpay.dshackle.data.TxId
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
@@ -36,15 +36,15 @@ import java.math.BigInteger
* *
* @see EthereumCachingReader * @see EthereumCachingReader
*/ */
class LocalCallRouter( class EthereumLocalReader(
private val reader: EthereumCachingReader, private val reader: EthereumCachingReader,
private val methods: CallMethods, private val methods: CallMethods,
private val head: Head, private val head: Head,
private val localEnabled: Boolean private val localEnabled: Boolean
) : Reader<JsonRpcRequest, JsonRpcResponse> { ) : JsonRpcReader {
companion object { companion object {
private val log = LoggerFactory.getLogger(LocalCallRouter::class.java) private val log = LoggerFactory.getLogger(EthereumLocalReader::class.java)
} }
override fun read(key: JsonRpcRequest): Mono<JsonRpcResponse> { override fun read(key: JsonRpcRequest): Mono<JsonRpcResponse> {

View File

@@ -20,7 +20,7 @@ import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.* import io.emeraldpay.dshackle.upstream.*
import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.AggregatedPendingTxes import io.emeraldpay.dshackle.upstream.ethereum.subscribe.AggregatedPendingTxes
@@ -28,8 +28,6 @@ import io.emeraldpay.dshackle.upstream.ethereum.subscribe.NoPendingTxes
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.PendingTxesSource import io.emeraldpay.dshackle.upstream.ethereum.subscribe.PendingTxesSource
import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice
import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstream import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstream
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.util.ConcurrentReferenceHashMap import org.springframework.util.ConcurrentReferenceHashMap
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
@@ -51,7 +49,7 @@ open class EthereumMultistream(
ConcurrentReferenceHashMap(16, ConcurrentReferenceHashMap.ReferenceType.WEAK) ConcurrentReferenceHashMap(16, ConcurrentReferenceHashMap.ReferenceType.WEAK)
private val reader: EthereumCachingReader = EthereumCachingReader(this, this.caches, getMethodsFactory()) private val reader: EthereumCachingReader = EthereumCachingReader(this, this.caches, getMethodsFactory())
private var subscribe = EthereumSubscriptionApi(this, NoPendingTxes()) private var subscribe = EthereumEgressSubscription(this, NoPendingTxes())
private val supportsEIP1559 = when (chain) { private val supportsEIP1559 = when (chain) {
Chain.ETHEREUM, Chain.TESTNET_ROPSTEN, Chain.TESTNET_GOERLI, Chain.TESTNET_RINKEBY -> true Chain.ETHEREUM, Chain.TESTNET_ROPSTEN, Chain.TESTNET_GOERLI, Chain.TESTNET_RINKEBY -> true
@@ -76,7 +74,7 @@ open class EthereumMultistream(
val pendingTxes: PendingTxesSource = upstreams val pendingTxes: PendingTxesSource = upstreams
.mapNotNull { .mapNotNull {
it.getUpstreamSubscriptions().getPendingTxes() it.getIngressSubscription().getPendingTxes()
}.let { }.let {
if (it.isEmpty()) { if (it.isEmpty()) {
NoPendingTxes() NoPendingTxes()
@@ -86,7 +84,7 @@ open class EthereumMultistream(
AggregatedPendingTxes(it) AggregatedPendingTxes(it)
} }
} }
subscribe = EthereumSubscriptionApi(this, pendingTxes) subscribe = EthereumEgressSubscription(this, pendingTxes)
} }
override fun start() { override fun start() {
@@ -175,12 +173,12 @@ open class EthereumMultistream(
return this as T return this as T
} }
override fun getRoutedApi(localEnabled: Boolean): Mono<Reader<JsonRpcRequest, JsonRpcResponse>> { override fun getEgressSubscription(): EgressSubscription {
return Mono.just(LocalCallRouter(reader, getMethods(), getHead(), localEnabled)) return subscribe
} }
override fun getSubscriptionApi(): EthereumSubscriptionApi { override fun getLocalReader(localEnabled: Boolean): Mono<JsonRpcReader> {
return subscribe return Mono.just(EthereumLocalReader(reader, getMethods(), getHead(), localEnabled))
} }
override fun getHead(mather: Selector.Matcher): Head = override fun getHead(mather: Selector.Matcher): Head =

View File

@@ -16,12 +16,10 @@
*/ */
package io.emeraldpay.dshackle.upstream.ethereum package io.emeraldpay.dshackle.upstream.ethereum
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.BlockValidator import io.emeraldpay.dshackle.upstream.BlockValidator
import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.scheduling.concurrent.CustomizableThreadFactory import org.springframework.scheduling.concurrent.CustomizableThreadFactory
import reactor.core.Disposable import reactor.core.Disposable
@@ -31,7 +29,7 @@ import java.time.Duration
import java.util.concurrent.Executors import java.util.concurrent.Executors
class EthereumRpcHead( class EthereumRpcHead(
private val api: Reader<JsonRpcRequest, JsonRpcResponse>, private val api: JsonRpcReader,
forkChoice: ForkChoice, forkChoice: ForkChoice,
upstreamId: String, upstreamId: String,
blockValidator: BlockValidator, blockValidator: BlockValidator,

View File

@@ -20,7 +20,7 @@ import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.cache.CachesEnabled import io.emeraldpay.dshackle.cache.CachesEnabled
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.startup.QuorumForLabels import io.emeraldpay.dshackle.startup.QuorumForLabels
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstream
@@ -28,8 +28,6 @@ import io.emeraldpay.dshackle.upstream.UpstreamAvailability
import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.dshackle.upstream.ethereum.connectors.ConnectorFactory import io.emeraldpay.dshackle.upstream.ethereum.connectors.ConnectorFactory
import io.emeraldpay.dshackle.upstream.ethereum.connectors.EthereumConnector import io.emeraldpay.dshackle.upstream.ethereum.connectors.EthereumConnector
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle import org.springframework.context.Lifecycle
import reactor.core.Disposable import reactor.core.Disposable
@@ -70,8 +68,8 @@ open class EthereumRpcUpstream(
} }
} }
override fun getUpstreamSubscriptions(): EthereumUpstreamSubscriptions { override fun getIngressSubscription(): EthereumIngressSubscription {
return connector.getUpstreamSubscriptions() return connector.getIngressSubscription()
} }
override fun getHead(): Head { override fun getHead(): Head {
@@ -88,8 +86,8 @@ open class EthereumRpcUpstream(
return connector.isRunning() return connector.isRunning()
} }
override fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse> { override fun getIngressReader(): JsonRpcReader {
return connector.getApi() return connector.getIngressReader()
} }
override fun isGrpc(): Boolean { override fun isGrpc(): Boolean {

View File

@@ -45,5 +45,5 @@ abstract class EthereumUpstream(
return node?.let { listOf(it.labels) } ?: emptyList() return node?.let { listOf(it.labels) } ?: emptyList()
} }
abstract fun getUpstreamSubscriptions(): EthereumUpstreamSubscriptions abstract fun getIngressSubscription(): EthereumIngressSubscription
} }

View File

@@ -48,7 +48,7 @@ open class EthereumUpstreamValidator(
open fun validate(): Mono<UpstreamAvailability> { open fun validate(): Mono<UpstreamAvailability> {
return upstream return upstream
.getApi() .getIngressReader()
.read(JsonRpcRequest("eth_syncing", listOf())) .read(JsonRpcRequest("eth_syncing", listOf()))
.flatMap(JsonRpcResponse::requireResult) .flatMap(JsonRpcResponse::requireResult)
.map { objectMapper.readValue(it, SyncingJson::class.java) } .map { objectMapper.readValue(it, SyncingJson::class.java) }
@@ -62,7 +62,7 @@ open class EthereumUpstreamValidator(
Mono.just(UpstreamAvailability.SYNCING) Mono.just(UpstreamAvailability.SYNCING)
} else { } else {
upstream upstream
.getApi() .getIngressReader()
.read(JsonRpcRequest("net_peerCount", listOf())) .read(JsonRpcRequest("net_peerCount", listOf()))
.flatMap(JsonRpcResponse::requireStringResult) .flatMap(JsonRpcResponse::requireStringResult)
.map(Integer::decode) .map(Integer::decode)

View File

@@ -20,7 +20,7 @@ import io.emeraldpay.dshackle.Defaults
import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.SilentException import io.emeraldpay.dshackle.SilentException
import io.emeraldpay.dshackle.data.BlockContainer import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.BlockValidator import io.emeraldpay.dshackle.upstream.BlockValidator
import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice
@@ -40,7 +40,7 @@ class EthereumWsHead(
upstreamId: String, upstreamId: String,
forkChoice: ForkChoice, forkChoice: ForkChoice,
blockValidator: BlockValidator, blockValidator: BlockValidator,
private val api: Reader<JsonRpcRequest, JsonRpcResponse>, private val api: JsonRpcReader,
private val wsSubscriptions: WsSubscriptions, private val wsSubscriptions: WsSubscriptions,
) : DefaultEthereumHead(upstreamId, forkChoice, blockValidator), Lifecycle { ) : DefaultEthereumHead(upstreamId, forkChoice, blockValidator), Lifecycle {

View File

@@ -15,14 +15,14 @@
*/ */
package io.emeraldpay.dshackle.upstream.ethereum package io.emeraldpay.dshackle.upstream.ethereum
import io.emeraldpay.dshackle.upstream.NoUpstreamSubscriptions import io.emeraldpay.dshackle.upstream.NoIngressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.PendingTxesSource import io.emeraldpay.dshackle.upstream.ethereum.subscribe.PendingTxesSource
class NoEthereumUpstreamSubscriptions : NoUpstreamSubscriptions(), EthereumUpstreamSubscriptions { class NoEthereumIngressSubscription : NoIngressSubscription(), EthereumIngressSubscription {
companion object { companion object {
@JvmStatic @JvmStatic
val DEFAULT = NoEthereumUpstreamSubscriptions() val DEFAULT = NoEthereumIngressSubscription()
} }
override fun getPendingTxes(): PendingTxesSource? { override fun getPendingTxes(): PendingTxesSource? {

View File

@@ -1,16 +1,14 @@
package io.emeraldpay.dshackle.upstream.ethereum.connectors package io.emeraldpay.dshackle.upstream.ethereum.connectors
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.ethereum.EthereumUpstreamSubscriptions import io.emeraldpay.dshackle.upstream.ethereum.EthereumIngressSubscription
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
interface EthereumConnector : Lifecycle { interface EthereumConnector : Lifecycle {
fun getHead(): Head fun getHead(): Head
fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse> fun getIngressReader(): JsonRpcReader
fun getUpstreamSubscriptions(): EthereumUpstreamSubscriptions fun getIngressSubscription(): EthereumIngressSubscription
} }

View File

@@ -2,7 +2,7 @@ package io.emeraldpay.dshackle.upstream.ethereum.connectors
import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.cache.CachesEnabled import io.emeraldpay.dshackle.cache.CachesEnabled
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.BlockValidator import io.emeraldpay.dshackle.upstream.BlockValidator
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Lifecycle
@@ -10,13 +10,11 @@ import io.emeraldpay.dshackle.upstream.MergedHead
import io.emeraldpay.dshackle.upstream.ethereum.* import io.emeraldpay.dshackle.upstream.ethereum.*
import io.emeraldpay.dshackle.upstream.forkchoice.AlwaysForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.AlwaysForkChoice
import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import java.time.Duration import java.time.Duration
class EthereumRpcConnector( class EthereumRpcConnector(
private val directReader: Reader<JsonRpcRequest, JsonRpcResponse>, private val directReader: JsonRpcReader,
wsFactory: EthereumWsFactory?, wsFactory: EthereumWsFactory?,
id: String, id: String,
forkChoice: ForkChoice, forkChoice: ForkChoice,
@@ -34,14 +32,14 @@ class EthereumRpcConnector(
// do not set upstream to the WS, since it doesn't control the RPC upstream // do not set upstream to the WS, since it doesn't control the RPC upstream
conn = wsFactory.create(null) conn = wsFactory.create(null)
val subscriptions = WsSubscriptionsImpl(conn) val subscriptions = WsSubscriptionsImpl(conn)
val wsHead = EthereumWsHead(id, AlwaysForkChoice(), blockValidator, getApi(), subscriptions) val wsHead = EthereumWsHead(id, AlwaysForkChoice(), blockValidator, getIngressReader(), subscriptions)
// receive all new blocks through WebSockets, but also periodically verify with RPC in case if WS failed // receive all new blocks through WebSockets, but also periodically verify with RPC in case if WS failed
val rpcHead = EthereumRpcHead(getApi(), AlwaysForkChoice(), id, blockValidator, Duration.ofSeconds(30)) val rpcHead = EthereumRpcHead(getIngressReader(), AlwaysForkChoice(), id, blockValidator, Duration.ofSeconds(30))
head = MergedHead(listOf(rpcHead, wsHead), forkChoice, "Merged for $id") head = MergedHead(listOf(rpcHead, wsHead), forkChoice, "Merged for $id")
} else { } else {
conn = null conn = null
log.warn("Setting up connector for $id upstream with RPC-only access, less effective than WS+RPC") log.warn("Setting up connector for $id upstream with RPC-only access, less effective than WS+RPC")
head = EthereumRpcHead(getApi(), forkChoice, id, blockValidator) head = EthereumRpcHead(getIngressReader(), forkChoice, id, blockValidator)
} }
} }
@@ -72,12 +70,12 @@ class EthereumRpcConnector(
conn?.close() conn?.close()
} }
override fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse> { override fun getIngressReader(): JsonRpcReader {
return directReader return directReader
} }
override fun getUpstreamSubscriptions(): EthereumUpstreamSubscriptions { override fun getIngressSubscription(): EthereumIngressSubscription {
return NoEthereumUpstreamSubscriptions.DEFAULT return NoEthereumIngressSubscription.DEFAULT
} }
override fun getHead(): Head { override fun getHead(): Head {

View File

@@ -1,14 +1,12 @@
package io.emeraldpay.dshackle.upstream.ethereum.connectors package io.emeraldpay.dshackle.upstream.ethereum.connectors
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.BlockValidator import io.emeraldpay.dshackle.upstream.BlockValidator
import io.emeraldpay.dshackle.upstream.DefaultUpstream import io.emeraldpay.dshackle.upstream.DefaultUpstream
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.ethereum.* import io.emeraldpay.dshackle.upstream.ethereum.*
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.EthereumWsSubscriptions import io.emeraldpay.dshackle.upstream.ethereum.subscribe.EthereumWsIngressSubscription
import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcWsClient import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcWsClient
class EthereumWsConnector( class EthereumWsConnector(
@@ -18,16 +16,16 @@ class EthereumWsConnector(
blockValidator: BlockValidator blockValidator: BlockValidator
) : EthereumConnector { ) : EthereumConnector {
private val conn: WsConnectionImpl private val conn: WsConnectionImpl
private val api: Reader<JsonRpcRequest, JsonRpcResponse> private val reader: JsonRpcReader
private val head: EthereumWsHead private val head: EthereumWsHead
private val subscriptions: EthereumUpstreamSubscriptions private val subscriptions: EthereumIngressSubscription
init { init {
conn = wsFactory.create(upstream) conn = wsFactory.create(upstream)
api = JsonRpcWsClient(conn) reader = JsonRpcWsClient(conn)
val wsSubscriptions = WsSubscriptionsImpl(conn) val wsSubscriptions = WsSubscriptionsImpl(conn)
head = EthereumWsHead(upstream.getId(), forkChoice, blockValidator, api, wsSubscriptions) head = EthereumWsHead(upstream.getId(), forkChoice, blockValidator, reader, wsSubscriptions)
subscriptions = EthereumWsSubscriptions(wsSubscriptions) subscriptions = EthereumWsIngressSubscription(wsSubscriptions)
} }
override fun start() { override fun start() {
@@ -44,11 +42,11 @@ class EthereumWsConnector(
head.stop() head.stop()
} }
override fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse> { override fun getIngressReader(): JsonRpcReader {
return api return reader
} }
override fun getUpstreamSubscriptions(): EthereumUpstreamSubscriptions { override fun getIngressSubscription(): EthereumIngressSubscription {
return subscriptions return subscriptions
} }

View File

@@ -20,7 +20,7 @@ import io.emeraldpay.api.proto.BlockchainOuterClass.NativeSubscribeReplyItem
import io.emeraldpay.api.proto.BlockchainOuterClass.NativeSubscribeRequest import io.emeraldpay.api.proto.BlockchainOuterClass.NativeSubscribeRequest
import io.emeraldpay.api.proto.ReactorBlockchainGrpc import io.emeraldpay.api.proto.ReactorBlockchainGrpc
import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.upstream.ethereum.EthereumSubscriptionApi import io.emeraldpay.dshackle.upstream.ethereum.EthereumEgressSubscription
import io.emeraldpay.etherjar.domain.TransactionId import io.emeraldpay.etherjar.domain.TransactionId
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
@@ -31,7 +31,7 @@ class DshacklePendingTxesSource(
private val request = NativeSubscribeRequest.newBuilder() private val request = NativeSubscribeRequest.newBuilder()
.setChainValue(blockchain.id) .setChainValue(blockchain.id)
.setMethod(EthereumSubscriptionApi.METHOD_PENDING_TXES) .setMethod(EthereumEgressSubscription.METHOD_PENDING_TXES)
.build() .build()
var available = false var available = false
@@ -53,7 +53,7 @@ class DshacklePendingTxesSource(
fun update(conf: BlockchainOuterClass.DescribeChain) { fun update(conf: BlockchainOuterClass.DescribeChain) {
available = conf.supportedMethodsList.any { available = conf.supportedMethodsList.any {
it == EthereumSubscriptionApi.METHOD_PENDING_TXES it == EthereumEgressSubscription.METHOD_PENDING_TXES
} }
} }
} }

View File

@@ -18,20 +18,29 @@ package io.emeraldpay.dshackle.upstream.ethereum.subscribe
import io.emeraldpay.api.proto.BlockchainOuterClass import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.api.proto.ReactorBlockchainGrpc import io.emeraldpay.api.proto.ReactorBlockchainGrpc
import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.upstream.IngressSubscription
import io.emeraldpay.dshackle.upstream.SubscriptionConnect import io.emeraldpay.dshackle.upstream.SubscriptionConnect
import io.emeraldpay.dshackle.upstream.UpstreamSubscriptions import io.emeraldpay.dshackle.upstream.ethereum.EthereumEgressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.EthereumSubscriptionApi import io.emeraldpay.dshackle.upstream.ethereum.EthereumIngressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.EthereumUpstreamSubscriptions import org.slf4j.LoggerFactory
class EthereumDshackleSubscriptions( class EthereumDshackleIngressSubscription(
blockchain: Chain, private val blockchain: Chain,
conn: ReactorBlockchainGrpc.ReactorBlockchainStub, private val conn: ReactorBlockchainGrpc.ReactorBlockchainStub,
) : UpstreamSubscriptions, EthereumUpstreamSubscriptions { ) : IngressSubscription, EthereumIngressSubscription {
companion object {
private val log = LoggerFactory.getLogger(EthereumDshackleIngressSubscription::class.java)
}
private val pendingTxes = DshacklePendingTxesSource(blockchain, conn) private val pendingTxes = DshacklePendingTxesSource(blockchain, conn)
override fun <T> get(method: String): SubscriptionConnect<T>? { override fun getAvailableTopics(): List<String> {
if (method == EthereumSubscriptionApi.METHOD_PENDING_TXES) { return listOf(EthereumEgressSubscription.METHOD_PENDING_TXES)
}
override fun <T> get(topic: String): SubscriptionConnect<T>? {
if (topic == EthereumEgressSubscription.METHOD_PENDING_TXES) {
return pendingTxes as SubscriptionConnect<T> return pendingTxes as SubscriptionConnect<T>
} }
return null return null

View File

@@ -15,25 +15,30 @@
*/ */
package io.emeraldpay.dshackle.upstream.ethereum.subscribe package io.emeraldpay.dshackle.upstream.ethereum.subscribe
import io.emeraldpay.dshackle.upstream.IngressSubscription
import io.emeraldpay.dshackle.upstream.SubscriptionConnect import io.emeraldpay.dshackle.upstream.SubscriptionConnect
import io.emeraldpay.dshackle.upstream.UpstreamSubscriptions import io.emeraldpay.dshackle.upstream.ethereum.EthereumEgressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.EthereumSubscriptionApi import io.emeraldpay.dshackle.upstream.ethereum.EthereumIngressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.EthereumUpstreamSubscriptions
import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptions import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptions
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
class EthereumWsSubscriptions( class EthereumWsIngressSubscription(
private val conn: WsSubscriptions private val conn: WsSubscriptions
) : UpstreamSubscriptions, EthereumUpstreamSubscriptions { ) : IngressSubscription, EthereumIngressSubscription {
companion object { companion object {
private val log = LoggerFactory.getLogger(EthereumWsSubscriptions::class.java) private val log = LoggerFactory.getLogger(EthereumWsIngressSubscription::class.java)
} }
private val pendingTxes = WebsocketPendingTxes(conn) private val pendingTxes = WebsocketPendingTxes(conn)
override fun <T> get(method: String): SubscriptionConnect<T>? { override fun getAvailableTopics(): List<String> {
if (method == EthereumSubscriptionApi.METHOD_PENDING_TXES) { return listOf(EthereumEgressSubscription.METHOD_PENDING_TXES)
}
@Suppress("UNCHECKED_CAST")
override fun <T> get(topic: String): SubscriptionConnect<T>? {
if (topic == EthereumEgressSubscription.METHOD_PENDING_TXES) {
return pendingTxes as SubscriptionConnect<T> return pendingTxes as SubscriptionConnect<T>
} }
return null return null

View File

@@ -16,7 +16,7 @@
package io.emeraldpay.dshackle.upstream.ethereum.subscribe package io.emeraldpay.dshackle.upstream.ethereum.subscribe
import io.emeraldpay.dshackle.upstream.SubscriptionConnect import io.emeraldpay.dshackle.upstream.SubscriptionConnect
import io.emeraldpay.dshackle.upstream.ethereum.EthereumSubscriptionApi import io.emeraldpay.dshackle.upstream.ethereum.EthereumEgressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptions import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptions
import io.emeraldpay.etherjar.domain.TransactionId import io.emeraldpay.etherjar.domain.TransactionId
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
@@ -33,7 +33,7 @@ class WebsocketPendingTxes(
} }
override fun createConnection(): Flux<TransactionId> { override fun createConnection(): Flux<TransactionId> {
return wsSubscriptions.subscribe(EthereumSubscriptionApi.METHOD_PENDING_TXES) return wsSubscriptions.subscribe(EthereumEgressSubscription.METHOD_PENDING_TXES)
.timeout(Duration.ofSeconds(60), Mono.empty()) .timeout(Duration.ofSeconds(60), Mono.empty())
.map { .map {
// comes as a JS string, i.e., within quotes // comes as a JS string, i.e., within quotes

View File

@@ -20,22 +20,13 @@ import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.ChainFees import io.emeraldpay.dshackle.upstream.*
import io.emeraldpay.dshackle.upstream.EmptyHead
import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.MergedHead
import io.emeraldpay.dshackle.upstream.Multistream
import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.AggregatedPendingTxes import io.emeraldpay.dshackle.upstream.ethereum.subscribe.AggregatedPendingTxes
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.NoPendingTxes import io.emeraldpay.dshackle.upstream.ethereum.subscribe.NoPendingTxes
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.PendingTxesSource import io.emeraldpay.dshackle.upstream.ethereum.subscribe.PendingTxesSource
import io.emeraldpay.dshackle.upstream.forkchoice.PriorityForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.PriorityForkChoice
import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstream import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstream
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.util.ConcurrentReferenceHashMap import org.springframework.util.ConcurrentReferenceHashMap
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
@@ -55,7 +46,7 @@ open class EthereumPosMultiStream(
private var head: Head? = null private var head: Head? = null
private val reader: EthereumCachingReader = EthereumCachingReader(this, this.caches, getMethodsFactory()) private val reader: EthereumCachingReader = EthereumCachingReader(this, this.caches, getMethodsFactory())
private var subscribe = EthereumSubscriptionApi(this, NoPendingTxes()) private var subscribe = EthereumEgressSubscription(this, NoPendingTxes())
private val feeEstimation = EthereumPriorityFees(this, reader, 256) private val feeEstimation = EthereumPriorityFees(this, reader, 256)
private val filteredHeads: MutableMap<String, Head> = private val filteredHeads: MutableMap<String, Head> =
ConcurrentReferenceHashMap(16, ConcurrentReferenceHashMap.ReferenceType.WEAK) ConcurrentReferenceHashMap(16, ConcurrentReferenceHashMap.ReferenceType.WEAK)
@@ -156,11 +147,11 @@ open class EthereumPosMultiStream(
return this as T return this as T
} }
override fun getRoutedApi(localEnabled: Boolean): Mono<Reader<JsonRpcRequest, JsonRpcResponse>> { override fun getLocalReader(localEnabled: Boolean): Mono<JsonRpcReader> {
return Mono.just(LocalCallRouter(reader, getMethods(), getHead(), localEnabled)) return Mono.just(EthereumLocalReader(reader, getMethods(), getHead(), localEnabled))
} }
override fun getSubscriptionApi(): EthereumSubscriptionApi { override fun getEgressSubscription(): EthereumEgressSubscription {
return subscribe return subscribe
} }
@@ -191,7 +182,7 @@ open class EthereumPosMultiStream(
val pendingTxes: PendingTxesSource = upstreams val pendingTxes: PendingTxesSource = upstreams
.mapNotNull { .mapNotNull {
it.getUpstreamSubscriptions().getPendingTxes() it.getIngressSubscription().getPendingTxes()
}.let { }.let {
if (it.isEmpty()) { if (it.isEmpty()) {
NoPendingTxes() NoPendingTxes()
@@ -201,6 +192,6 @@ open class EthereumPosMultiStream(
AggregatedPendingTxes(it) AggregatedPendingTxes(it)
} }
} }
subscribe = EthereumSubscriptionApi(this, pendingTxes) subscribe = EthereumEgressSubscription(this, pendingTxes)
} }
} }

View File

@@ -20,7 +20,7 @@ import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.cache.CachesEnabled import io.emeraldpay.dshackle.cache.CachesEnabled
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.startup.QuorumForLabels import io.emeraldpay.dshackle.startup.QuorumForLabels
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Lifecycle
@@ -29,8 +29,6 @@ import io.emeraldpay.dshackle.upstream.UpstreamAvailability
import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.dshackle.upstream.ethereum.connectors.ConnectorFactory import io.emeraldpay.dshackle.upstream.ethereum.connectors.ConnectorFactory
import io.emeraldpay.dshackle.upstream.ethereum.connectors.EthereumConnector import io.emeraldpay.dshackle.upstream.ethereum.connectors.EthereumConnector
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import reactor.core.Disposable import reactor.core.Disposable
@@ -83,8 +81,8 @@ open class EthereumPosRpcUpstream(
return connector.isRunning() return connector.isRunning()
} }
override fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse> { override fun getIngressReader(): JsonRpcReader {
return connector.getApi() return connector.getIngressReader()
} }
override fun isGrpc(): Boolean { override fun isGrpc(): Boolean {
@@ -99,7 +97,7 @@ open class EthereumPosRpcUpstream(
return this as T return this as T
} }
override fun getUpstreamSubscriptions(): EthereumUpstreamSubscriptions { override fun getIngressSubscription(): EthereumIngressSubscription {
return connector.getUpstreamSubscriptions() return connector.getIngressSubscription()
} }
} }

View File

@@ -45,5 +45,5 @@ abstract class EthereumPosUpstream(
return node?.let { listOf(it.labels) } ?: emptyList() return node?.let { listOf(it.labels) } ?: emptyList()
} }
abstract fun getUpstreamSubscriptions(): EthereumUpstreamSubscriptions abstract fun getIngressSubscription(): EthereumIngressSubscription
} }

View File

@@ -22,7 +22,7 @@ import io.emeraldpay.dshackle.Defaults
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.data.BlockContainer import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.data.BlockId import io.emeraldpay.dshackle.data.BlockId
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.Capability
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Lifecycle
@@ -66,7 +66,7 @@ class BitcoinGrpcUpstream(
} }
private val extractBlock = ExtractBlock() private val extractBlock = ExtractBlock()
private val defaultReader: Reader<JsonRpcRequest, JsonRpcResponse> = client.getReader() private val defaultReader: JsonRpcReader = client.getReader()
private val blockConverter: Function<BlockchainOuterClass.ChainHead, BlockContainer> = Function { value -> private val blockConverter: Function<BlockchainOuterClass.ChainHead, BlockContainer> = Function { value ->
val block = BlockContainer( val block = BlockContainer(
value.height, value.height,
@@ -112,7 +112,7 @@ class BitcoinGrpcUpstream(
return grpcHead return grpcHead
} }
override fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse> { override fun getIngressReader(): JsonRpcReader {
return defaultReader return defaultReader
} }

View File

@@ -23,7 +23,7 @@ import io.emeraldpay.dshackle.Defaults
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.data.BlockContainer import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.data.BlockId import io.emeraldpay.dshackle.data.BlockId
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.startup.QuorumForLabels import io.emeraldpay.dshackle.startup.QuorumForLabels
import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.Capability
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
@@ -31,9 +31,9 @@ import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.UpstreamAvailability
import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.dshackle.upstream.ethereum.EthereumIngressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.EthereumUpstream import io.emeraldpay.dshackle.upstream.ethereum.EthereumUpstream
import io.emeraldpay.dshackle.upstream.ethereum.EthereumUpstreamSubscriptions import io.emeraldpay.dshackle.upstream.ethereum.subscribe.EthereumDshackleIngressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.EthereumDshackleSubscriptions
import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcGrpcClient import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcGrpcClient
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
@@ -107,9 +107,9 @@ open class EthereumGrpcUpstream(
private val grpcHead = GrpcHead(chain, this, remote, blockConverter, reloadBlock, MostWorkForkChoice()) private val grpcHead = GrpcHead(chain, this, remote, blockConverter, reloadBlock, MostWorkForkChoice())
private var capabilities: Set<Capability> = emptySet() private var capabilities: Set<Capability> = emptySet()
private val defaultReader: Reader<JsonRpcRequest, JsonRpcResponse> = client.getReader() private val defaultReader: JsonRpcReader = client.getReader()
var timeout = Defaults.timeout var timeout = Defaults.timeout
private val ethereumSubscriptions = EthereumDshackleSubscriptions(chain, remote) private val ethereumSubscriptions = EthereumDshackleIngressSubscription(chain, remote)
override fun getBlockchainApi(): ReactorBlockchainStub { override fun getBlockchainApi(): ReactorBlockchainStub {
return remote return remote
@@ -144,7 +144,7 @@ open class EthereumGrpcUpstream(
return upstreamStatus.getLabels() return upstreamStatus.getLabels()
} }
override fun getUpstreamSubscriptions(): EthereumUpstreamSubscriptions { override fun getIngressSubscription(): EthereumIngressSubscription {
return ethereumSubscriptions return ethereumSubscriptions
} }
@@ -162,7 +162,7 @@ open class EthereumGrpcUpstream(
return grpcHead return grpcHead
} }
override fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse> { override fun getIngressReader(): JsonRpcReader {
return defaultReader return defaultReader
} }

View File

@@ -23,7 +23,7 @@ import io.emeraldpay.dshackle.Defaults
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.data.BlockContainer import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.data.BlockId import io.emeraldpay.dshackle.data.BlockId
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.startup.QuorumForLabels import io.emeraldpay.dshackle.startup.QuorumForLabels
import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.Capability
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
@@ -31,9 +31,9 @@ import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.UpstreamAvailability
import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.dshackle.upstream.ethereum.EthereumIngressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosUpstream import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosUpstream
import io.emeraldpay.dshackle.upstream.ethereum.EthereumUpstreamSubscriptions import io.emeraldpay.dshackle.upstream.ethereum.subscribe.EthereumDshackleIngressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.EthereumDshackleSubscriptions
import io.emeraldpay.dshackle.upstream.forkchoice.NoChoiceWithPriorityForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.NoChoiceWithPriorityForkChoice
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcGrpcClient import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcGrpcClient
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
@@ -107,9 +107,9 @@ open class EthereumPosGrpcUpstream(
private val grpcHead = GrpcHead(chain, this, remote, blockConverter, reloadBlock, NoChoiceWithPriorityForkChoice(nodeRating, parentId)) private val grpcHead = GrpcHead(chain, this, remote, blockConverter, reloadBlock, NoChoiceWithPriorityForkChoice(nodeRating, parentId))
private var capabilities: Set<Capability> = emptySet() private var capabilities: Set<Capability> = emptySet()
private val defaultReader: Reader<JsonRpcRequest, JsonRpcResponse> = client.getReader() private val defaultReader: JsonRpcReader = client.getReader()
var timeout = Defaults.timeout var timeout = Defaults.timeout
private val ethereumSubscriptions = EthereumDshackleSubscriptions(chain, remote) private val ethereumSubscriptions = EthereumDshackleIngressSubscription(chain, remote)
override fun start() { override fun start() {
} }
@@ -144,6 +144,10 @@ open class EthereumPosGrpcUpstream(
return upstreamStatus.getLabels() return upstreamStatus.getLabels()
} }
override fun getIngressSubscription(): EthereumIngressSubscription {
return ethereumSubscriptions
}
override fun getMethods(): CallMethods { override fun getMethods(): CallMethods {
return upstreamStatus.getCallMethods() return upstreamStatus.getCallMethods()
} }
@@ -158,7 +162,7 @@ open class EthereumPosGrpcUpstream(
return grpcHead return grpcHead
} }
override fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse> { override fun getIngressReader(): JsonRpcReader {
return defaultReader return defaultReader
} }
@@ -177,8 +181,4 @@ open class EthereumPosGrpcUpstream(
override fun isGrpc(): Boolean { override fun isGrpc(): Boolean {
return true return true
} }
override fun getUpstreamSubscriptions(): EthereumUpstreamSubscriptions {
return ethereumSubscriptions
}
} }

View File

@@ -21,7 +21,7 @@ import io.emeraldpay.api.proto.BlockchainOuterClass.NativeCallReplySignature
import io.emeraldpay.api.proto.ReactorBlockchainGrpc import io.emeraldpay.api.proto.ReactorBlockchainGrpc
import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
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 io.emeraldpay.etherjar.rpc.RpcResponseError import io.emeraldpay.etherjar.rpc.RpcResponseError
@@ -41,7 +41,7 @@ class JsonRpcGrpcClient(
private val log = LoggerFactory.getLogger(JsonRpcGrpcClient::class.java) private val log = LoggerFactory.getLogger(JsonRpcGrpcClient::class.java)
} }
fun getReader(): Reader<JsonRpcRequest, JsonRpcResponse> { fun getReader(): JsonRpcReader {
return Executor(stub, chain, metrics) return Executor(stub, chain, metrics)
} }
@@ -49,7 +49,7 @@ class JsonRpcGrpcClient(
private val stub: ReactorBlockchainGrpc.ReactorBlockchainStub, private val stub: ReactorBlockchainGrpc.ReactorBlockchainStub,
private val chain: Chain, private val chain: Chain,
private val metrics: RpcMetrics? private val metrics: RpcMetrics?
) : Reader<JsonRpcRequest, JsonRpcResponse> { ) : JsonRpcReader {
override fun read(key: JsonRpcRequest): Mono<JsonRpcResponse> { override fun read(key: JsonRpcRequest): Mono<JsonRpcResponse> {
val timer = StopWatch() val timer = StopWatch()

View File

@@ -16,7 +16,7 @@
package io.emeraldpay.dshackle.upstream.rpcclient package io.emeraldpay.dshackle.upstream.rpcclient
import io.emeraldpay.dshackle.config.AuthConfig import io.emeraldpay.dshackle.config.AuthConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.etherjar.rpc.RpcException import io.emeraldpay.etherjar.rpc.RpcException
import io.emeraldpay.etherjar.rpc.RpcResponseError import io.emeraldpay.etherjar.rpc.RpcResponseError
import io.netty.buffer.Unpooled import io.netty.buffer.Unpooled
@@ -48,7 +48,7 @@ class JsonRpcHttpClient(
private val metrics: RpcMetrics, private val metrics: RpcMetrics,
basicAuth: AuthConfig.ClientBasicAuth? = null, basicAuth: AuthConfig.ClientBasicAuth? = null,
tlsCAAuth: ByteArray? = null tlsCAAuth: ByteArray? = null
) : Reader<JsonRpcRequest, JsonRpcResponse> { ) : JsonRpcReader {
companion object { companion object {
private val log = LoggerFactory.getLogger(JsonRpcHttpClient::class.java) private val log = LoggerFactory.getLogger(JsonRpcHttpClient::class.java)

View File

@@ -1,6 +1,6 @@
package io.emeraldpay.dshackle.upstream.rpcclient package io.emeraldpay.dshackle.upstream.rpcclient
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
@@ -9,9 +9,9 @@ import reactor.core.publisher.Mono
* It always calls the Primary reader, and if it fails or produces an empty result, then it calls the Secondary reader. * It always calls the Primary reader, and if it fails or produces an empty result, then it calls the Secondary reader.
*/ */
class JsonRpcSwitchClient( class JsonRpcSwitchClient(
private val primary: Reader<JsonRpcRequest, JsonRpcResponse>, private val primary: JsonRpcReader,
private val secondary: Reader<JsonRpcRequest, JsonRpcResponse>, private val secondary: JsonRpcReader,
) : Reader<JsonRpcRequest, JsonRpcResponse> { ) : JsonRpcReader {
companion object { companion object {
private val log = LoggerFactory.getLogger(JsonRpcSwitchClient::class.java) private val log = LoggerFactory.getLogger(JsonRpcSwitchClient::class.java)

View File

@@ -15,14 +15,14 @@
*/ */
package io.emeraldpay.dshackle.upstream.rpcclient package io.emeraldpay.dshackle.upstream.rpcclient
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.ethereum.WsConnectionImpl import io.emeraldpay.dshackle.upstream.ethereum.WsConnectionImpl
import io.emeraldpay.etherjar.rpc.RpcResponseError import io.emeraldpay.etherjar.rpc.RpcResponseError
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
class JsonRpcWsClient( class JsonRpcWsClient(
private val ws: WsConnectionImpl private val ws: WsConnectionImpl
) : Reader<JsonRpcRequest, JsonRpcResponse> { ) : JsonRpcReader {
override fun read(key: JsonRpcRequest): Mono<JsonRpcResponse> { override fun read(key: JsonRpcRequest): Mono<JsonRpcResponse> {
if (!ws.isConnected) { if (!ws.isConnected) {

View File

@@ -17,7 +17,6 @@ package io.emeraldpay.dshackle.quorum
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.test.TestingCommons
import io.emeraldpay.dshackle.upstream.FilteredApis import io.emeraldpay.dshackle.upstream.FilteredApis
import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstream
@@ -40,7 +39,7 @@ class QuorumRpcReaderSpec extends Specification {
def up = Mock(Upstream) { def up = Mock(Upstream) {
_ * isAvailable() >> true _ * isAvailable() >> true
_ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY _ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY
1 * getApi() >> Mock(Reader) { 1 * getIngressReader() >> Mock(Reader) {
1 * read(new JsonRpcRequest("eth_test", [])) >> Mono.just(JsonRpcResponse.ok("1")) 1 * read(new JsonRpcRequest("eth_test", [])) >> Mono.just(JsonRpcResponse.ok("1"))
} }
} }
@@ -73,7 +72,7 @@ class QuorumRpcReaderSpec extends Specification {
def up = Mock(Upstream) { def up = Mock(Upstream) {
_ * isAvailable() >> true _ * isAvailable() >> true
_ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY _ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY
_ * getApi() >> api _ * getIngressReader() >> api
} }
def apis = new FilteredApis( def apis = new FilteredApis(
Chain.ETHEREUM, Chain.ETHEREUM,
@@ -110,7 +109,7 @@ class QuorumRpcReaderSpec extends Specification {
def up = Mock(Upstream) { def up = Mock(Upstream) {
_ * isAvailable() >> true _ * isAvailable() >> true
_ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY _ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY
_ * getApi() >> api _ * getIngressReader() >> api
} }
def apis = new FilteredApis( def apis = new FilteredApis(
Chain.ETHEREUM, Chain.ETHEREUM,
@@ -137,7 +136,7 @@ class QuorumRpcReaderSpec extends Specification {
def up = Mock(Upstream) { def up = Mock(Upstream) {
_ * isAvailable() >> true _ * isAvailable() >> true
_ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY _ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY
_ * getApi() >> Mock(Reader) { _ * getIngressReader() >> Mock(Reader) {
2 * read(new JsonRpcRequest("eth_test", [])) >>> [ 2 * read(new JsonRpcRequest("eth_test", [])) >>> [
Mono.just(JsonRpcResponse.ok("null")), Mono.just(JsonRpcResponse.ok("null")),
Mono.just(JsonRpcResponse.ok("1")) Mono.just(JsonRpcResponse.ok("1"))
@@ -169,7 +168,7 @@ class QuorumRpcReaderSpec extends Specification {
def up = Mock(Upstream) { def up = Mock(Upstream) {
_ * isAvailable() >> true _ * isAvailable() >> true
_ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY _ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY
_ * getApi() >> Mock(Reader) { _ * getIngressReader() >> Mock(Reader) {
2 * read(new JsonRpcRequest("eth_test", [])) >>> [ 2 * read(new JsonRpcRequest("eth_test", [])) >>> [
Mono.just(JsonRpcResponse.error(1, "test")), Mono.just(JsonRpcResponse.error(1, "test")),
Mono.just(JsonRpcResponse.ok("1")) Mono.just(JsonRpcResponse.ok("1"))
@@ -200,7 +199,7 @@ class QuorumRpcReaderSpec extends Specification {
def up = Mock(Upstream) { def up = Mock(Upstream) {
_ * isAvailable() >> true _ * isAvailable() >> true
_ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY _ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY
_ * getApi() >> Mock(Reader) { _ * getIngressReader() >> Mock(Reader) {
3 * read(new JsonRpcRequest("eth_test", [])) >>> [ 3 * read(new JsonRpcRequest("eth_test", [])) >>> [
Mono.just(JsonRpcResponse.ok("null")), Mono.just(JsonRpcResponse.ok("null")),
Mono.just(JsonRpcResponse.error(1, "test")), Mono.just(JsonRpcResponse.error(1, "test")),
@@ -239,7 +238,7 @@ class QuorumRpcReaderSpec extends Specification {
def up = Mock(Upstream) { def up = Mock(Upstream) {
_ * isAvailable() >> true _ * isAvailable() >> true
_ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY _ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY
_ * getApi() >> api _ * getIngressReader() >> api
} }
def apis = new FilteredApis( def apis = new FilteredApis(
Chain.ETHEREUM, Chain.ETHEREUM,
@@ -269,7 +268,7 @@ class QuorumRpcReaderSpec extends Specification {
def up = Mock(Upstream) { def up = Mock(Upstream) {
_ * isAvailable() >> true _ * isAvailable() >> true
_ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY _ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY
_ * getApi() >> api _ * getIngressReader() >> api
} }
def apis = new FilteredApis( def apis = new FilteredApis(
Chain.ETHEREUM, Chain.ETHEREUM,
@@ -298,7 +297,7 @@ class QuorumRpcReaderSpec extends Specification {
_ * getLag() >> 0 _ * getLag() >> 0
_ * isAvailable() >> true _ * isAvailable() >> true
_ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY _ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY
_ * getApi() >> Mock(Reader) { _ * getIngressReader() >> Mock(Reader) {
_ * read(new JsonRpcRequest("eth_test", [])) >>> [ _ * read(new JsonRpcRequest("eth_test", [])) >>> [
Mono.just(JsonRpcResponse.error(-3010, "test")), Mono.just(JsonRpcResponse.error(-3010, "test")),
] ]

View File

@@ -75,7 +75,7 @@ class NativeCallSpec extends Specification {
1 * read(new JsonRpcRequest("eth_test", [])) >> Mono.just(new JsonRpcResponse("1".bytes, null)) 1 * read(new JsonRpcRequest("eth_test", [])) >> Mono.just(new JsonRpcResponse("1".bytes, null))
} }
def upstream = Mock(Multistream) { def upstream = Mock(Multistream) {
1 * getRoutedApi(_) >> Mono.just(routedApi) 1 * getLocalReader(_) >> Mono.just(routedApi)
} }
def nativeCall = nativeCall() def nativeCall = nativeCall()
@@ -96,7 +96,7 @@ class NativeCallSpec extends Specification {
1 * read(new JsonRpcRequest("eth_test", [])) >> Mono.error(new RpcException(RpcResponseError.CODE_METHOD_NOT_EXIST, "Test message")) 1 * read(new JsonRpcRequest("eth_test", [])) >> Mono.error(new RpcException(RpcResponseError.CODE_METHOD_NOT_EXIST, "Test message"))
} }
def upstream = Mock(Multistream) { def upstream = Mock(Multistream) {
1 * getRoutedApi(_) >> Mono.just(routedApi) 1 * getLocalReader(_) >> Mono.just(routedApi)
} }
def nativeCall = nativeCall() def nativeCall = nativeCall()

View File

@@ -20,7 +20,7 @@ import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.dshackle.test.MultistreamHolderMock import io.emeraldpay.dshackle.test.MultistreamHolderMock
import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream
import io.emeraldpay.dshackle.upstream.ethereum.EthereumSubscriptionApi import io.emeraldpay.dshackle.upstream.ethereum.EthereumEgressSubscription
import io.emeraldpay.dshackle.upstream.signature.NoSigner import io.emeraldpay.dshackle.upstream.signature.NoSigner
import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.Chain
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
@@ -39,12 +39,12 @@ class NativeSubscribeSpec extends Specification {
.setMethod("newHeads") .setMethod("newHeads")
.build() .build()
def subscribe = Mock(EthereumSubscriptionApi) { def subscribe = Mock(EthereumEgressSubscription) {
1 * it.subscribe("newHeads", null, _ as Selector.AnyLabelMatcher) >> Flux.just("{}") 1 * it.subscribe("newHeads", null, _ as Selector.AnyLabelMatcher) >> Flux.just("{}")
} }
def up = Mock(EthereumPosMultiStream) { def up = Mock(EthereumPosMultiStream) {
1 * it.tryProxy(_ as Selector.AnyLabelMatcher, call) >> null 1 * it.tryProxy(_ as Selector.AnyLabelMatcher, call) >> null
1 * it.getSubscriptionApi() >> subscribe 1 * it.getEgressSubscription() >> subscribe
} }
def nativeSubscribe = new NativeSubscribe(new MultistreamHolderMock(Chain.ETHEREUM, up), signer) def nativeSubscribe = new NativeSubscribe(new MultistreamHolderMock(Chain.ETHEREUM, up), signer)
@@ -70,7 +70,7 @@ class NativeSubscribeSpec extends Specification {
)) ))
.build() .build()
def subscribe = Mock(EthereumSubscriptionApi) { def subscribe = Mock(EthereumEgressSubscription) {
1 * it.subscribe("logs", { params -> 1 * it.subscribe("logs", { params ->
println("params: $params") println("params: $params")
def ok = params instanceof Map && def ok = params instanceof Map &&
@@ -83,7 +83,7 @@ class NativeSubscribeSpec extends Specification {
} }
def up = Mock(EthereumPosMultiStream) { def up = Mock(EthereumPosMultiStream) {
1 * it.tryProxy(_ as Selector.AnyLabelMatcher, call) >> null 1 * it.tryProxy(_ as Selector.AnyLabelMatcher, call) >> null
1 * it.getSubscriptionApi() >> subscribe 1 * it.getEgressSubscription() >> subscribe
} }
def nativeSubscribe = new NativeSubscribe(new MultistreamHolderMock(Chain.ETHEREUM, up), signer) def nativeSubscribe = new NativeSubscribe(new MultistreamHolderMock(Chain.ETHEREUM, up), signer)
@@ -106,7 +106,7 @@ class NativeSubscribeSpec extends Specification {
.build() .build()
def up = Mock(EthereumPosMultiStream) { def up = Mock(EthereumPosMultiStream) {
1 * it.tryProxy(_ as Selector.AnyLabelMatcher, call) >> Flux.just("{}") 1 * it.tryProxy(_ as Selector.AnyLabelMatcher, call) >> Flux.just("{}")
0 * it.getSubscriptionApi() 0 * it.getEgressSubscription()
} }
def nativeSubscribe = new NativeSubscribe(new MultistreamHolderMock(Chain.ETHEREUM, up), signer) def nativeSubscribe = new NativeSubscribe(new MultistreamHolderMock(Chain.ETHEREUM, up), signer)

View File

@@ -8,7 +8,7 @@ import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.SubscriptionConnect import io.emeraldpay.dshackle.upstream.SubscriptionConnect
import io.emeraldpay.dshackle.upstream.ethereum.ERC20Balance import io.emeraldpay.dshackle.upstream.ethereum.ERC20Balance
import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream
import io.emeraldpay.dshackle.upstream.ethereum.EthereumSubscriptionApi import io.emeraldpay.dshackle.upstream.ethereum.EthereumEgressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.ConnectLogs import io.emeraldpay.dshackle.upstream.ethereum.subscribe.ConnectLogs
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.json.LogMessage import io.emeraldpay.dshackle.upstream.ethereum.subscribe.json.LogMessage
import io.emeraldpay.etherjar.domain.BlockHash import io.emeraldpay.etherjar.domain.BlockHash
@@ -184,11 +184,11 @@ class TrackERC20AddressSpec extends Specification {
connect connect
} }
} }
def sub = Mock(EthereumSubscriptionApi) { def sub = Mock(EthereumEgressSubscription) {
1 * getLogs() >> logs 1 * getLogs() >> logs
} }
def up = Mock(EthereumPosMultiStream) { def up = Mock(EthereumPosMultiStream) {
1 * getSubscriptionApi() >> sub 1 * getEgressSubscription() >> sub
_ * cast(EthereumPosMultiStream) >> { args -> _ * cast(EthereumPosMultiStream) >> { args ->
it it
} }

View File

@@ -2,8 +2,9 @@ package io.emeraldpay.dshackle.test
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.ethereum.EthereumUpstreamSubscriptions import io.emeraldpay.dshackle.upstream.NoIngressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.NoEthereumUpstreamSubscriptions import io.emeraldpay.dshackle.upstream.ethereum.EthereumIngressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.NoEthereumIngressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.connectors.EthereumConnector import io.emeraldpay.dshackle.upstream.ethereum.connectors.EthereumConnector
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
@@ -17,7 +18,7 @@ class EthereumConnectorMock implements EthereumConnector {
} }
@Override @Override
Reader<JsonRpcRequest, JsonRpcResponse> getApi() { Reader<JsonRpcRequest, JsonRpcResponse> getIngressReader() {
return this.api return this.api
} }
@@ -38,7 +39,7 @@ class EthereumConnectorMock implements EthereumConnector {
} }
@Override @Override
EthereumUpstreamSubscriptions getUpstreamSubscriptions() { EthereumIngressSubscription getIngressSubscription() {
return NoEthereumUpstreamSubscriptions.DEFAULT return NoEthereumIngressSubscription.DEFAULT
} }
} }

View File

@@ -259,7 +259,7 @@ class MultistreamSpec extends Specification {
@NotNull @NotNull
@Override @Override
Mono<Reader<JsonRpcRequest, JsonRpcResponse>> getRoutedApi(boolean localEnabled) { Mono<Reader<JsonRpcRequest, JsonRpcResponse>> getLocalReader(boolean localEnabled) {
return null return null
} }

View File

@@ -21,11 +21,11 @@ import io.emeraldpay.etherjar.domain.Address
import io.emeraldpay.etherjar.hex.Hex32 import io.emeraldpay.etherjar.hex.Hex32
import spock.lang.Specification import spock.lang.Specification
class EthereumSubscriptionApiSpec extends Specification { class EthereumEgressSubscriptionSpec extends Specification {
def "read empty logs request"() { def "read empty logs request"() {
setup: setup:
def ethereumSubscribe = new EthereumSubscriptionApi(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) def ethereumSubscribe = new EthereumEgressSubscription(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource))
when: when:
def act = ethereumSubscribe.readLogsRequest([:]) def act = ethereumSubscribe.readLogsRequest([:])
@@ -36,7 +36,7 @@ class EthereumSubscriptionApiSpec extends Specification {
def "read single address logs request"() { def "read single address logs request"() {
setup: setup:
def ethereumSubscribe = new EthereumSubscriptionApi(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) def ethereumSubscribe = new EthereumEgressSubscription(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource))
when: when:
def act = ethereumSubscribe.readLogsRequest([ def act = ethereumSubscribe.readLogsRequest([
address: "0x829bd824b016326a401d083b33d092293333a830" address: "0x829bd824b016326a401d083b33d092293333a830"
@@ -61,7 +61,7 @@ class EthereumSubscriptionApiSpec extends Specification {
def "ignores invalid address for logs request"() { def "ignores invalid address for logs request"() {
setup: setup:
def ethereumSubscribe = new EthereumSubscriptionApi(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) def ethereumSubscribe = new EthereumEgressSubscription(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource))
when: when:
def act = ethereumSubscribe.readLogsRequest([ def act = ethereumSubscribe.readLogsRequest([
address: "829bd824b016326a401d083b33d092293333a830" address: "829bd824b016326a401d083b33d092293333a830"
@@ -74,7 +74,7 @@ class EthereumSubscriptionApiSpec extends Specification {
def "read multi address logs request"() { def "read multi address logs request"() {
setup: setup:
def ethereumSubscribe = new EthereumSubscriptionApi(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) def ethereumSubscribe = new EthereumEgressSubscription(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource))
when: when:
def act = ethereumSubscribe.readLogsRequest([ def act = ethereumSubscribe.readLogsRequest([
address: ["0x829bd824b016326a401d083b33d092293333a830", "0x401d083b33d092293333a83829bd824b016326a0"] address: ["0x829bd824b016326a401d083b33d092293333a830", "0x401d083b33d092293333a83829bd824b016326a0"]
@@ -90,7 +90,7 @@ class EthereumSubscriptionApiSpec extends Specification {
def "read single topic logs request"() { def "read single topic logs request"() {
setup: setup:
def ethereumSubscribe = new EthereumSubscriptionApi(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) def ethereumSubscribe = new EthereumEgressSubscription(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource))
when: when:
def act = ethereumSubscribe.readLogsRequest([ def act = ethereumSubscribe.readLogsRequest([
topics: "0xddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef" topics: "0xddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef"
@@ -115,7 +115,7 @@ class EthereumSubscriptionApiSpec extends Specification {
def "read invalid topic for request"() { def "read invalid topic for request"() {
setup: setup:
def ethereumSubscribe = new EthereumSubscriptionApi(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) def ethereumSubscribe = new EthereumEgressSubscription(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource))
when: when:
def act = ethereumSubscribe.readLogsRequest([ def act = ethereumSubscribe.readLogsRequest([
topics: [ topics: [
@@ -133,7 +133,7 @@ class EthereumSubscriptionApiSpec extends Specification {
def "read multi topic logs request"() { def "read multi topic logs request"() {
setup: setup:
def ethereumSubscribe = new EthereumSubscriptionApi(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) def ethereumSubscribe = new EthereumEgressSubscription(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource))
when: when:
def act = ethereumSubscribe.readLogsRequest([ def act = ethereumSubscribe.readLogsRequest([
topics: [ topics: [
@@ -152,7 +152,7 @@ class EthereumSubscriptionApiSpec extends Specification {
def "read full logs request"() { def "read full logs request"() {
setup: setup:
def ethereumSubscribe = new EthereumSubscriptionApi(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) def ethereumSubscribe = new EthereumEgressSubscription(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource))
when: when:
def act = ethereumSubscribe.readLogsRequest([ def act = ethereumSubscribe.readLogsRequest([
address: "0x298d492e8c1d909d3f63bc4a36c66c64acb3d695", address: "0x298d492e8c1d909d3f63bc4a36c66c64acb3d695",

View File

@@ -2,7 +2,6 @@ package io.emeraldpay.dshackle.upstream.ethereum
import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.config.CacheConfig
import io.emeraldpay.dshackle.reader.EmptyReader import io.emeraldpay.dshackle.reader.EmptyReader
import io.emeraldpay.dshackle.test.TestingCommons import io.emeraldpay.dshackle.test.TestingCommons
import io.emeraldpay.dshackle.upstream.EmptyHead import io.emeraldpay.dshackle.upstream.EmptyHead
@@ -18,12 +17,12 @@ import spock.lang.Specification
import java.time.Duration import java.time.Duration
class LocalCallRouterSpec extends Specification { class EthereumLocalReaderSpec extends Specification {
def "Calls hardcoded"() { def "Calls hardcoded"() {
setup: setup:
def methods = new DefaultEthereumMethods(Chain.ETHEREUM) def methods = new DefaultEthereumMethods(Chain.ETHEREUM)
def router = new LocalCallRouter( def router = new EthereumLocalReader(
new EthereumCachingReader( new EthereumCachingReader(
TestingCommons.multistream(TestingCommons.api()), TestingCommons.multistream(TestingCommons.api()),
Caches.default(), Caches.default(),
@@ -42,7 +41,7 @@ class LocalCallRouterSpec extends Specification {
def "Returns empty if nonce set"() { def "Returns empty if nonce set"() {
setup: setup:
def methods = new DefaultEthereumMethods(Chain.ETHEREUM) def methods = new DefaultEthereumMethods(Chain.ETHEREUM)
def router = new LocalCallRouter( def router = new EthereumLocalReader(
new EthereumCachingReader( new EthereumCachingReader(
TestingCommons.multistream(TestingCommons.api()), TestingCommons.multistream(TestingCommons.api()),
Caches.default(), Caches.default(),
@@ -72,7 +71,7 @@ class LocalCallRouterSpec extends Specification {
} }
} }
def methods = new DefaultEthereumMethods(Chain.ETHEREUM) def methods = new DefaultEthereumMethods(Chain.ETHEREUM)
def router = new LocalCallRouter(reader, methods, head, true) def router = new EthereumLocalReader(reader, methods, head, true)
when: when:
def act = router.getBlockByNumber(["latest", false]) def act = router.getBlockByNumber(["latest", false])
@@ -98,7 +97,7 @@ class LocalCallRouterSpec extends Specification {
} }
} }
def methods = new DefaultEthereumMethods(Chain.ETHEREUM) def methods = new DefaultEthereumMethods(Chain.ETHEREUM)
def router = new LocalCallRouter(reader, methods, head, true) def router = new EthereumLocalReader(reader, methods, head, true)
when: when:
def act = router.getBlockByNumber(["earliest", false]) def act = router.getBlockByNumber(["earliest", false])
@@ -124,7 +123,7 @@ class LocalCallRouterSpec extends Specification {
} }
} }
def methods = new DefaultEthereumMethods(Chain.ETHEREUM) def methods = new DefaultEthereumMethods(Chain.ETHEREUM)
def router = new LocalCallRouter(reader, methods, head, true) def router = new EthereumLocalReader(reader, methods, head, true)
when: when:
def act = router.getBlockByNumber(["0x123ef", false]) def act = router.getBlockByNumber(["0x123ef", false])
@@ -148,7 +147,7 @@ class LocalCallRouterSpec extends Specification {
_ * blocksByHeightAsCont() >> new EmptyReader<>() _ * blocksByHeightAsCont() >> new EmptyReader<>()
} }
def methods = new DefaultEthereumMethods(Chain.ETHEREUM) def methods = new DefaultEthereumMethods(Chain.ETHEREUM)
def router = new LocalCallRouter(reader, methods, head, true) def router = new EthereumLocalReader(reader, methods, head, true)
when: when:
def act = router.getBlockByNumber(["0x0", true]) def act = router.getBlockByNumber(["0x0", true])