From 25dde8cfdfbba9d10bb08a75dd238a7cde42e1b2 Mon Sep 17 00:00:00 2001 From: a10zn8 Date: Tue, 27 Dec 2022 17:09:02 +0400 Subject: [PATCH] blockchain supported subscriptions data --- emerald-grpc | 2 +- .../io/emeraldpay/dshackle/rpc/Describe.kt | 3 +- .../dshackle/rpc/NativeSubscribe.kt | 2 +- .../dshackle/rpc/TrackERC20Address.kt | 2 +- .../dshackle/upstream/EgressSubscription.kt | 33 ++++++++++++++++ ...scriptions.kt => HasEgressSubscription.kt} | 11 +----- ...ubscriptions.kt => IngressSubscription.kt} | 5 ++- .../dshackle/upstream/Multistream.kt | 2 +- .../upstream/NoIngressSubscription.kt | 31 +++++++++++++++ .../upstream/bitcoin/BitcoinMultistream.kt | 4 ++ ...onApi.kt => EthereumEgressSubscription.kt} | 39 +++++++++++++------ ...ions.kt => EthereumIngressSubscription.kt} | 4 +- .../ethereum/EthereumLikeMultistream.kt | 4 +- .../upstream/ethereum/EthereumMultistream.kt | 14 +++---- .../upstream/ethereum/EthereumRpcUpstream.kt | 4 +- .../upstream/ethereum/EthereumUpstream.kt | 2 +- ...ns.kt => NoEthereumIngressSubscription.kt} | 6 +-- .../ethereum/connectors/EthereumConnector.kt | 4 +- .../connectors/EthereumRpcConnector.kt | 4 +- .../connectors/EthereumWsConnector.kt | 8 ++-- .../subscribe/DshacklePendingTxesSource.kt | 6 +-- ...=> EthereumDshackleIngressSubscription.kt} | 27 ++++++++----- ...ns.kt => EthereumWsIngressSubscription.kt} | 21 ++++++---- .../subscribe/WebsocketPendingTxes.kt | 4 +- .../ethereum_pos/EthereumPosMultiStream.kt | 17 +++----- .../ethereum_pos/EthereumPosRpcUpstream.kt | 4 +- .../ethereum_pos/EthereumPosUpstream.kt | 2 +- .../upstream/grpc/EthereumGrpcUpstream.kt | 8 ++-- .../upstream/grpc/EthereumPosGrpcUpstream.kt | 14 +++---- .../dshackle/rpc/NativeSubscribeSpec.groovy | 12 +++--- .../dshackle/rpc/TrackERC20AddressSpec.groovy | 6 +-- .../test/EthereumConnectorMock.groovy | 9 +++-- ... => EthereumEgressSubscriptionSpec.groovy} | 18 ++++----- 33 files changed, 209 insertions(+), 123 deletions(-) create mode 100644 src/main/kotlin/io/emeraldpay/dshackle/upstream/EgressSubscription.kt rename src/main/kotlin/io/emeraldpay/dshackle/upstream/{NoUpstreamSubscriptions.kt => HasEgressSubscription.kt} (73%) rename src/main/kotlin/io/emeraldpay/dshackle/upstream/{UpstreamSubscriptions.kt => IngressSubscription.kt} (84%) create mode 100644 src/main/kotlin/io/emeraldpay/dshackle/upstream/NoIngressSubscription.kt rename src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/{EthereumSubscriptionApi.kt => EthereumEgressSubscription.kt} (79%) rename src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/{EthereumUpstreamSubscriptions.kt => EthereumIngressSubscription.kt} (86%) rename src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/{NoEthereumUpstreamSubscriptions.kt => NoEthereumIngressSubscription.kt} (79%) rename src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/{EthereumDshackleSubscriptions.kt => EthereumDshackleIngressSubscription.kt} (59%) rename src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/{EthereumWsSubscriptions.kt => EthereumWsIngressSubscription.kt} (61%) rename src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/{EthereumSubscriptionApiSpec.groovy => EthereumEgressSubscriptionSpec.groovy} (80%) diff --git a/emerald-grpc b/emerald-grpc index 812e9e6d..c4aa5d05 160000 --- a/emerald-grpc +++ b/emerald-grpc @@ -1 +1 @@ -Subproject commit 812e9e6ddbbd0fb3f21936b26b82fc76475bce1c +Subproject commit c4aa5d052a53c17377477d5740b670049f86d4c3 diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/Describe.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/Describe.kt index 775d16fc..fd7828ee 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/Describe.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/Describe.kt @@ -36,13 +36,14 @@ class Describe( return requestMono.map { _ -> val resp = BlockchainOuterClass.DescribeResponse.newBuilder() multistreamHolder.getAvailable().forEach { chain -> - multistreamHolder.getUpstream(chain)?.let { chainUpstreams -> + multistreamHolder.getUpstream(chain).let { chainUpstreams -> val status = subscribeStatus.chainStatus(chain, chainUpstreams.getStatus(), chainUpstreams) val targets = chainUpstreams.getMethods().getSupportedMethods() val capabilities: MutableSet = mutableSetOf() val chainDescription = BlockchainOuterClass.DescribeChain.newBuilder() .setChain(Common.ChainRef.forNumber(chain.id)) .addAllSupportedMethods(targets) + .addAllSupportedSubscriptions(chainUpstreams.getEgressSubscription().getAvailableTopics()) .setStatus(status) .setCurrentHeight(chainUpstreams.getHead().getCurrentHeight() ?: 0) chainUpstreams.getAll().let { ups -> diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeSubscribe.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeSubscribe.kt index 8b5408c8..1a95f671 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeSubscribe.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeSubscribe.kt @@ -97,7 +97,7 @@ open class NativeSubscribe( } open fun subscribe(chain: Chain, method: String, params: Any?, matcher: Selector.Matcher): Flux = - getUpstream(chain).getSubscriptionApi().subscribe(method, params, matcher) + getUpstream(chain).getEgressSubscription().subscribe(method, params, matcher) private fun getUpstream(chain: Chain): EthereumLikeMultistream = multistreamHolder.getUpstream(chain).let { it as EthereumLikeMultistream } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackERC20Address.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackERC20Address.kt index 94bfa5a3..7c62e12b 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackERC20Address.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackERC20Address.kt @@ -88,7 +88,7 @@ class TrackERC20Address( val asset = request.asset.code.lowercase(Locale.getDefault()) val tokenDefinition = tokens[TokenId(chain, asset)] ?: return Flux.empty() val logs = getUpstream(chain) - .getSubscriptionApi().logs + .getEgressSubscription().logs .create( listOf(tokenDefinition.token.contract), listOf(EventId.fromSignature("Transfer", "address", "address", "uint256")), diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/EgressSubscription.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/EgressSubscription.kt new file mode 100644 index 00000000..4c83c4de --- /dev/null +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/EgressSubscription.kt @@ -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 + fun subscribe(topic: String, params: Any?, matcher: Selector.Matcher): Flux +} + +class EmptyEgressSubscription : EgressSubscription { + override fun getAvailableTopics(): List { + return emptyList() + } + + override fun subscribe(topic: String, params: Any?, matcher: Selector.Matcher): Flux { + return Flux.empty() + } +} diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/NoUpstreamSubscriptions.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/HasEgressSubscription.kt similarity index 73% rename from src/main/kotlin/io/emeraldpay/dshackle/upstream/NoUpstreamSubscriptions.kt rename to src/main/kotlin/io/emeraldpay/dshackle/upstream/HasEgressSubscription.kt index 245823d5..4fd0f34d 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/NoUpstreamSubscriptions.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/HasEgressSubscription.kt @@ -15,13 +15,6 @@ */ package io.emeraldpay.dshackle.upstream -open class NoUpstreamSubscriptions : UpstreamSubscriptions { - - companion object { - val DEFAULT = NoUpstreamSubscriptions() - } - - override fun get(method: String): SubscriptionConnect? { - return null - } +interface HasEgressSubscription { + fun getEgressSubscription(): EgressSubscription } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/UpstreamSubscriptions.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/IngressSubscription.kt similarity index 84% rename from src/main/kotlin/io/emeraldpay/dshackle/upstream/UpstreamSubscriptions.kt rename to src/main/kotlin/io/emeraldpay/dshackle/upstream/IngressSubscription.kt index b1159c5f..d597587d 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/UpstreamSubscriptions.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/IngressSubscription.kt @@ -18,7 +18,8 @@ package io.emeraldpay.dshackle.upstream /** * Subscriptions available on the current upstream */ -interface UpstreamSubscriptions { +interface IngressSubscription { - fun get(method: String): SubscriptionConnect? + fun getAvailableTopics(): List + fun get(topic: String): SubscriptionConnect? } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt index aa7d5a9d..9c655c26 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt @@ -55,7 +55,7 @@ abstract class Multistream( val chain: Chain, private val upstreams: MutableList, val caches: Caches, -) : Upstream, Lifecycle { +) : Upstream, Lifecycle, HasEgressSubscription { companion object { private val log = LoggerFactory.getLogger(Multistream::class.java) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/NoIngressSubscription.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/NoIngressSubscription.kt new file mode 100644 index 00000000..9fe184a7 --- /dev/null +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/NoIngressSubscription.kt @@ -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 { + return listOf() + } + + override fun get(topic: String): SubscriptionConnect? { + return null + } +} diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinMultistream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinMultistream.kt index a36b4a1c..0400debd 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinMultistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinMultistream.kt @@ -144,6 +144,10 @@ open class BitcoinMultistream( return this as T } + override fun getEgressSubscription(): EgressSubscription { + return EmptyEgressSubscription() + } + override fun isRunning(): Boolean { return super.isRunning() || reader.isRunning() } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumSubscriptionApi.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumEgressSubscription.kt similarity index 79% rename from src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumSubscriptionApi.kt rename to src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumEgressSubscription.kt index f24ade24..41fe519c 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumSubscriptionApi.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumEgressSubscription.kt @@ -1,5 +1,6 @@ package io.emeraldpay.dshackle.upstream.ethereum +import io.emeraldpay.dshackle.upstream.EgressSubscription import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.ethereum.subscribe.ConnectLogs import io.emeraldpay.dshackle.upstream.ethereum.subscribe.ConnectNewHeads @@ -10,13 +11,13 @@ import io.emeraldpay.etherjar.hex.Hex32 import org.slf4j.LoggerFactory import reactor.core.publisher.Flux -open class EthereumSubscriptionApi( +open class EthereumEgressSubscription( val upstream: EthereumLikeMultistream, - val pendingTxesSource: PendingTxesSource -) { + val pendingTxesSource: PendingTxesSource? +) : EgressSubscription { 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_LOGS = "logs" @@ -24,16 +25,30 @@ open class EthereumSubscriptionApi( 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) open val logs = ConnectLogs(upstream) private val syncing = ConnectSyncing(upstream) + override fun getAvailableTopics() = availableTopics + @Suppress("UNCHECKED_CAST") - open fun subscribe(method: String, params: Any?, matcher: Selector.Matcher): Flux { - if (method == METHOD_NEW_HEADS) { + override fun subscribe(topic: String, params: Any?, matcher: Selector.Matcher): Flux { + if (topic == METHOD_NEW_HEADS) { return newHeads.connect(matcher) } - if (method == METHOD_LOGS) { + if (topic == METHOD_LOGS) { val paramsMap = try { if (params != null && Map::class.java.isAssignableFrom(params.javaClass)) { readLogsRequest(params as Map) @@ -41,17 +56,17 @@ open class EthereumSubscriptionApi( LogsRequest(emptyList(), emptyList()) } } 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) } - if (method == METHOD_SYNCING) { + if (topic == METHOD_SYNCING) { return syncing.connect(matcher) } - if (method == METHOD_PENDING_TXES) { - return pendingTxesSource.connect(matcher) + if (topic == METHOD_PENDING_TXES) { + 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( diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumUpstreamSubscriptions.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumIngressSubscription.kt similarity index 86% rename from src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumUpstreamSubscriptions.kt rename to src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumIngressSubscription.kt index f646ff68..5cdf4165 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumUpstreamSubscriptions.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumIngressSubscription.kt @@ -15,10 +15,10 @@ */ 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 -interface EthereumUpstreamSubscriptions : UpstreamSubscriptions { +interface EthereumIngressSubscription : IngressSubscription { fun getPendingTxes(): PendingTxesSource? } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumLikeMultistream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumLikeMultistream.kt index a83a01fb..6bd40ebb 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumLikeMultistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumLikeMultistream.kt @@ -1,14 +1,14 @@ package io.emeraldpay.dshackle.upstream.ethereum import io.emeraldpay.api.proto.BlockchainOuterClass +import io.emeraldpay.dshackle.upstream.HasEgressSubscription import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Upstream import reactor.core.publisher.Flux -interface EthereumLikeMultistream : Upstream { +interface EthereumLikeMultistream : Upstream, HasEgressSubscription { fun getReader(): EthereumCachingReader - fun getSubscriptionApi(): EthereumSubscriptionApi fun getHead(mather: Selector.Matcher): Head diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumMultistream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumMultistream.kt index 0ffa9e27..60df1920 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumMultistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumMultistream.kt @@ -51,7 +51,7 @@ open class EthereumMultistream( ConcurrentReferenceHashMap(16, ConcurrentReferenceHashMap.ReferenceType.WEAK) 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) { Chain.ETHEREUM, Chain.TESTNET_ROPSTEN, Chain.TESTNET_GOERLI, Chain.TESTNET_RINKEBY -> true @@ -76,7 +76,7 @@ open class EthereumMultistream( val pendingTxes: PendingTxesSource = upstreams .mapNotNull { - it.getUpstreamSubscriptions().getPendingTxes() + it.getIngressSubscription().getPendingTxes() }.let { if (it.isEmpty()) { NoPendingTxes() @@ -86,7 +86,7 @@ open class EthereumMultistream( AggregatedPendingTxes(it) } } - subscribe = EthereumSubscriptionApi(this, pendingTxes) + subscribe = EthereumEgressSubscription(this, pendingTxes) } override fun start() { @@ -175,12 +175,12 @@ open class EthereumMultistream( return this as T } - override fun getRoutedApi(localEnabled: Boolean): Mono> { - return Mono.just(LocalCallRouter(reader, getMethods(), getHead(), localEnabled)) + override fun getEgressSubscription(): EgressSubscription { + return subscribe } - override fun getSubscriptionApi(): EthereumSubscriptionApi { - return subscribe + override fun getRoutedApi(localEnabled: Boolean): Mono> { + return Mono.just(LocalCallRouter(reader, getMethods(), getHead(), localEnabled)) } override fun getHead(mather: Selector.Matcher): Head = diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumRpcUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumRpcUpstream.kt index 1e059b59..c1dac856 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumRpcUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumRpcUpstream.kt @@ -70,8 +70,8 @@ open class EthereumRpcUpstream( } } - override fun getUpstreamSubscriptions(): EthereumUpstreamSubscriptions { - return connector.getUpstreamSubscriptions() + override fun getIngressSubscription(): EthereumIngressSubscription { + return connector.getIngressSubscription() } override fun getHead(): Head { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumUpstream.kt index 65c38e32..ba9046a9 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumUpstream.kt @@ -45,5 +45,5 @@ abstract class EthereumUpstream( return node?.let { listOf(it.labels) } ?: emptyList() } - abstract fun getUpstreamSubscriptions(): EthereumUpstreamSubscriptions + abstract fun getIngressSubscription(): EthereumIngressSubscription } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/NoEthereumUpstreamSubscriptions.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/NoEthereumIngressSubscription.kt similarity index 79% rename from src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/NoEthereumUpstreamSubscriptions.kt rename to src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/NoEthereumIngressSubscription.kt index 56b701a3..0d420aa6 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/NoEthereumUpstreamSubscriptions.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/NoEthereumIngressSubscription.kt @@ -15,14 +15,14 @@ */ 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 -class NoEthereumUpstreamSubscriptions : NoUpstreamSubscriptions(), EthereumUpstreamSubscriptions { +class NoEthereumIngressSubscription : NoIngressSubscription(), EthereumIngressSubscription { companion object { @JvmStatic - val DEFAULT = NoEthereumUpstreamSubscriptions() + val DEFAULT = NoEthereumIngressSubscription() } override fun getPendingTxes(): PendingTxesSource? { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/connectors/EthereumConnector.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/connectors/EthereumConnector.kt index 1eba7166..74acddda 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/connectors/EthereumConnector.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/connectors/EthereumConnector.kt @@ -3,7 +3,7 @@ package io.emeraldpay.dshackle.upstream.ethereum.connectors import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.upstream.Head 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 @@ -12,5 +12,5 @@ interface EthereumConnector : Lifecycle { fun getApi(): Reader - fun getUpstreamSubscriptions(): EthereumUpstreamSubscriptions + fun getIngressSubscription(): EthereumIngressSubscription } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/connectors/EthereumRpcConnector.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/connectors/EthereumRpcConnector.kt index fa79ae50..15d95a8c 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/connectors/EthereumRpcConnector.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/connectors/EthereumRpcConnector.kt @@ -76,8 +76,8 @@ class EthereumRpcConnector( return directReader } - override fun getUpstreamSubscriptions(): EthereumUpstreamSubscriptions { - return NoEthereumUpstreamSubscriptions.DEFAULT + override fun getIngressSubscription(): EthereumIngressSubscription { + return NoEthereumIngressSubscription.DEFAULT } override fun getHead(): Head { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/connectors/EthereumWsConnector.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/connectors/EthereumWsConnector.kt index 4c67acff..ab4d18da 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/connectors/EthereumWsConnector.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/connectors/EthereumWsConnector.kt @@ -5,7 +5,7 @@ import io.emeraldpay.dshackle.upstream.BlockValidator import io.emeraldpay.dshackle.upstream.DefaultUpstream import io.emeraldpay.dshackle.upstream.Head 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.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse @@ -20,14 +20,14 @@ class EthereumWsConnector( private val conn: WsConnectionImpl private val api: Reader private val head: EthereumWsHead - private val subscriptions: EthereumUpstreamSubscriptions + private val subscriptions: EthereumIngressSubscription init { conn = wsFactory.create(upstream) api = JsonRpcWsClient(conn) val wsSubscriptions = WsSubscriptionsImpl(conn) head = EthereumWsHead(upstream.getId(), forkChoice, blockValidator, api, wsSubscriptions) - subscriptions = EthereumWsSubscriptions(wsSubscriptions) + subscriptions = EthereumWsIngressSubscription(wsSubscriptions) } override fun start() { @@ -48,7 +48,7 @@ class EthereumWsConnector( return api } - override fun getUpstreamSubscriptions(): EthereumUpstreamSubscriptions { + override fun getIngressSubscription(): EthereumIngressSubscription { return subscriptions } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/DshacklePendingTxesSource.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/DshacklePendingTxesSource.kt index 324d5743..88ed72bc 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/DshacklePendingTxesSource.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/DshacklePendingTxesSource.kt @@ -20,7 +20,7 @@ import io.emeraldpay.api.proto.BlockchainOuterClass.NativeSubscribeReplyItem import io.emeraldpay.api.proto.BlockchainOuterClass.NativeSubscribeRequest import io.emeraldpay.api.proto.ReactorBlockchainGrpc 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 reactor.core.publisher.Flux @@ -31,7 +31,7 @@ class DshacklePendingTxesSource( private val request = NativeSubscribeRequest.newBuilder() .setChainValue(blockchain.id) - .setMethod(EthereumSubscriptionApi.METHOD_PENDING_TXES) + .setMethod(EthereumEgressSubscription.METHOD_PENDING_TXES) .build() var available = false @@ -53,7 +53,7 @@ class DshacklePendingTxesSource( fun update(conf: BlockchainOuterClass.DescribeChain) { available = conf.supportedMethodsList.any { - it == EthereumSubscriptionApi.METHOD_PENDING_TXES + it == EthereumEgressSubscription.METHOD_PENDING_TXES } } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/EthereumDshackleSubscriptions.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/EthereumDshackleIngressSubscription.kt similarity index 59% rename from src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/EthereumDshackleSubscriptions.kt rename to src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/EthereumDshackleIngressSubscription.kt index da651838..f6887199 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/EthereumDshackleSubscriptions.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/EthereumDshackleIngressSubscription.kt @@ -18,20 +18,29 @@ package io.emeraldpay.dshackle.upstream.ethereum.subscribe import io.emeraldpay.api.proto.BlockchainOuterClass import io.emeraldpay.api.proto.ReactorBlockchainGrpc import io.emeraldpay.dshackle.Chain +import io.emeraldpay.dshackle.upstream.IngressSubscription import io.emeraldpay.dshackle.upstream.SubscriptionConnect -import io.emeraldpay.dshackle.upstream.UpstreamSubscriptions -import io.emeraldpay.dshackle.upstream.ethereum.EthereumSubscriptionApi -import io.emeraldpay.dshackle.upstream.ethereum.EthereumUpstreamSubscriptions +import io.emeraldpay.dshackle.upstream.ethereum.EthereumEgressSubscription +import io.emeraldpay.dshackle.upstream.ethereum.EthereumIngressSubscription +import org.slf4j.LoggerFactory -class EthereumDshackleSubscriptions( - blockchain: Chain, - conn: ReactorBlockchainGrpc.ReactorBlockchainStub, -) : UpstreamSubscriptions, EthereumUpstreamSubscriptions { +class EthereumDshackleIngressSubscription( + private val blockchain: Chain, + private val conn: ReactorBlockchainGrpc.ReactorBlockchainStub, +) : IngressSubscription, EthereumIngressSubscription { + + companion object { + private val log = LoggerFactory.getLogger(EthereumDshackleIngressSubscription::class.java) + } private val pendingTxes = DshacklePendingTxesSource(blockchain, conn) - override fun get(method: String): SubscriptionConnect? { - if (method == EthereumSubscriptionApi.METHOD_PENDING_TXES) { + override fun getAvailableTopics(): List { + return listOf(EthereumEgressSubscription.METHOD_PENDING_TXES) + } + + override fun get(topic: String): SubscriptionConnect? { + if (topic == EthereumEgressSubscription.METHOD_PENDING_TXES) { return pendingTxes as SubscriptionConnect } return null diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/EthereumWsSubscriptions.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/EthereumWsIngressSubscription.kt similarity index 61% rename from src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/EthereumWsSubscriptions.kt rename to src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/EthereumWsIngressSubscription.kt index b60ed3ab..89e2d2d7 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/EthereumWsSubscriptions.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/EthereumWsIngressSubscription.kt @@ -15,25 +15,30 @@ */ package io.emeraldpay.dshackle.upstream.ethereum.subscribe +import io.emeraldpay.dshackle.upstream.IngressSubscription import io.emeraldpay.dshackle.upstream.SubscriptionConnect -import io.emeraldpay.dshackle.upstream.UpstreamSubscriptions -import io.emeraldpay.dshackle.upstream.ethereum.EthereumSubscriptionApi -import io.emeraldpay.dshackle.upstream.ethereum.EthereumUpstreamSubscriptions +import io.emeraldpay.dshackle.upstream.ethereum.EthereumEgressSubscription +import io.emeraldpay.dshackle.upstream.ethereum.EthereumIngressSubscription import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptions import org.slf4j.LoggerFactory -class EthereumWsSubscriptions( +class EthereumWsIngressSubscription( private val conn: WsSubscriptions -) : UpstreamSubscriptions, EthereumUpstreamSubscriptions { +) : IngressSubscription, EthereumIngressSubscription { companion object { - private val log = LoggerFactory.getLogger(EthereumWsSubscriptions::class.java) + private val log = LoggerFactory.getLogger(EthereumWsIngressSubscription::class.java) } private val pendingTxes = WebsocketPendingTxes(conn) - override fun get(method: String): SubscriptionConnect? { - if (method == EthereumSubscriptionApi.METHOD_PENDING_TXES) { + override fun getAvailableTopics(): List { + return listOf(EthereumEgressSubscription.METHOD_PENDING_TXES) + } + + @Suppress("UNCHECKED_CAST") + override fun get(topic: String): SubscriptionConnect? { + if (topic == EthereumEgressSubscription.METHOD_PENDING_TXES) { return pendingTxes as SubscriptionConnect } return null diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/WebsocketPendingTxes.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/WebsocketPendingTxes.kt index 3e59fc18..de52c906 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/WebsocketPendingTxes.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/WebsocketPendingTxes.kt @@ -16,7 +16,7 @@ package io.emeraldpay.dshackle.upstream.ethereum.subscribe 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.etherjar.domain.TransactionId import org.slf4j.LoggerFactory @@ -33,7 +33,7 @@ class WebsocketPendingTxes( } override fun createConnection(): Flux { - return wsSubscriptions.subscribe(EthereumSubscriptionApi.METHOD_PENDING_TXES) + return wsSubscriptions.subscribe(EthereumEgressSubscription.METHOD_PENDING_TXES) .timeout(Duration.ofSeconds(60), Mono.empty()) .map { // comes as a JS string, i.e., within quotes diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosMultiStream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosMultiStream.kt index cac43b8a..ad473770 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosMultiStream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosMultiStream.kt @@ -21,14 +21,7 @@ import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.reader.Reader -import io.emeraldpay.dshackle.upstream.ChainFees -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.* import io.emeraldpay.dshackle.upstream.ethereum.subscribe.AggregatedPendingTxes import io.emeraldpay.dshackle.upstream.ethereum.subscribe.NoPendingTxes import io.emeraldpay.dshackle.upstream.ethereum.subscribe.PendingTxesSource @@ -55,7 +48,7 @@ open class EthereumPosMultiStream( private var head: Head? = null 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 filteredHeads: MutableMap = ConcurrentReferenceHashMap(16, ConcurrentReferenceHashMap.ReferenceType.WEAK) @@ -160,7 +153,7 @@ open class EthereumPosMultiStream( return Mono.just(LocalCallRouter(reader, getMethods(), getHead(), localEnabled)) } - override fun getSubscriptionApi(): EthereumSubscriptionApi { + override fun getEgressSubscription(): EthereumEgressSubscription { return subscribe } @@ -191,7 +184,7 @@ open class EthereumPosMultiStream( val pendingTxes: PendingTxesSource = upstreams .mapNotNull { - it.getUpstreamSubscriptions().getPendingTxes() + it.getIngressSubscription().getPendingTxes() }.let { if (it.isEmpty()) { NoPendingTxes() @@ -201,6 +194,6 @@ open class EthereumPosMultiStream( AggregatedPendingTxes(it) } } - subscribe = EthereumSubscriptionApi(this, pendingTxes) + subscribe = EthereumEgressSubscription(this, pendingTxes) } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosRpcUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosRpcUpstream.kt index 9440faa0..cf34f76b 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosRpcUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosRpcUpstream.kt @@ -99,7 +99,7 @@ open class EthereumPosRpcUpstream( return this as T } - override fun getUpstreamSubscriptions(): EthereumUpstreamSubscriptions { - return connector.getUpstreamSubscriptions() + override fun getIngressSubscription(): EthereumIngressSubscription { + return connector.getIngressSubscription() } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosUpstream.kt index 5e915b59..a3f7a84f 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosUpstream.kt @@ -45,5 +45,5 @@ abstract class EthereumPosUpstream( return node?.let { listOf(it.labels) } ?: emptyList() } - abstract fun getUpstreamSubscriptions(): EthereumUpstreamSubscriptions + abstract fun getIngressSubscription(): EthereumIngressSubscription } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcUpstream.kt index 3be65b97..8006e078 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcUpstream.kt @@ -31,9 +31,9 @@ import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.UpstreamAvailability 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.EthereumUpstreamSubscriptions -import io.emeraldpay.dshackle.upstream.ethereum.subscribe.EthereumDshackleSubscriptions +import io.emeraldpay.dshackle.upstream.ethereum.subscribe.EthereumDshackleIngressSubscription import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcGrpcClient import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest @@ -109,7 +109,7 @@ open class EthereumGrpcUpstream( private val defaultReader: Reader = client.getReader() var timeout = Defaults.timeout - private val ethereumSubscriptions = EthereumDshackleSubscriptions(chain, remote) + private val ethereumSubscriptions = EthereumDshackleIngressSubscription(chain, remote) override fun getBlockchainApi(): ReactorBlockchainStub { return remote @@ -144,7 +144,7 @@ open class EthereumGrpcUpstream( return upstreamStatus.getLabels() } - override fun getUpstreamSubscriptions(): EthereumUpstreamSubscriptions { + override fun getIngressSubscription(): EthereumIngressSubscription { return ethereumSubscriptions } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumPosGrpcUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumPosGrpcUpstream.kt index 73724d1f..75f456ba 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumPosGrpcUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumPosGrpcUpstream.kt @@ -31,9 +31,9 @@ import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.UpstreamAvailability 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.EthereumUpstreamSubscriptions -import io.emeraldpay.dshackle.upstream.ethereum.subscribe.EthereumDshackleSubscriptions +import io.emeraldpay.dshackle.upstream.ethereum.subscribe.EthereumDshackleIngressSubscription import io.emeraldpay.dshackle.upstream.forkchoice.NoChoiceWithPriorityForkChoice import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcGrpcClient import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest @@ -109,7 +109,7 @@ open class EthereumPosGrpcUpstream( private val defaultReader: Reader = client.getReader() var timeout = Defaults.timeout - private val ethereumSubscriptions = EthereumDshackleSubscriptions(chain, remote) + private val ethereumSubscriptions = EthereumDshackleIngressSubscription(chain, remote) override fun start() { } @@ -144,6 +144,10 @@ open class EthereumPosGrpcUpstream( return upstreamStatus.getLabels() } + override fun getIngressSubscription(): EthereumIngressSubscription { + return ethereumSubscriptions + } + override fun getMethods(): CallMethods { return upstreamStatus.getCallMethods() } @@ -177,8 +181,4 @@ open class EthereumPosGrpcUpstream( override fun isGrpc(): Boolean { return true } - - override fun getUpstreamSubscriptions(): EthereumUpstreamSubscriptions { - return ethereumSubscriptions - } } diff --git a/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeSubscribeSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeSubscribeSpec.groovy index 6ca1d0bf..35950350 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeSubscribeSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeSubscribeSpec.groovy @@ -20,7 +20,7 @@ import io.emeraldpay.api.proto.BlockchainOuterClass import io.emeraldpay.dshackle.test.MultistreamHolderMock import io.emeraldpay.dshackle.upstream.Selector 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.Chain import reactor.core.publisher.Flux @@ -39,12 +39,12 @@ class NativeSubscribeSpec extends Specification { .setMethod("newHeads") .build() - def subscribe = Mock(EthereumSubscriptionApi) { + def subscribe = Mock(EthereumEgressSubscription) { 1 * it.subscribe("newHeads", null, _ as Selector.AnyLabelMatcher) >> Flux.just("{}") } def up = Mock(EthereumPosMultiStream) { 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) @@ -70,7 +70,7 @@ class NativeSubscribeSpec extends Specification { )) .build() - def subscribe = Mock(EthereumSubscriptionApi) { + def subscribe = Mock(EthereumEgressSubscription) { 1 * it.subscribe("logs", { params -> println("params: $params") def ok = params instanceof Map && @@ -83,7 +83,7 @@ class NativeSubscribeSpec extends Specification { } def up = Mock(EthereumPosMultiStream) { 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) @@ -106,7 +106,7 @@ class NativeSubscribeSpec extends Specification { .build() def up = Mock(EthereumPosMultiStream) { 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) diff --git a/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackERC20AddressSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackERC20AddressSpec.groovy index 6be4170c..e46274a6 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackERC20AddressSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackERC20AddressSpec.groovy @@ -8,7 +8,7 @@ import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.SubscriptionConnect import io.emeraldpay.dshackle.upstream.ethereum.ERC20Balance 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.json.LogMessage import io.emeraldpay.etherjar.domain.BlockHash @@ -184,11 +184,11 @@ class TrackERC20AddressSpec extends Specification { connect } } - def sub = Mock(EthereumSubscriptionApi) { + def sub = Mock(EthereumEgressSubscription) { 1 * getLogs() >> logs } def up = Mock(EthereumPosMultiStream) { - 1 * getSubscriptionApi() >> sub + 1 * getEgressSubscription() >> sub _ * cast(EthereumPosMultiStream) >> { args -> it } diff --git a/src/test/groovy/io/emeraldpay/dshackle/test/EthereumConnectorMock.groovy b/src/test/groovy/io/emeraldpay/dshackle/test/EthereumConnectorMock.groovy index 18c31581..010f6f9d 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/test/EthereumConnectorMock.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/test/EthereumConnectorMock.groovy @@ -2,8 +2,9 @@ package io.emeraldpay.dshackle.test import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.upstream.Head -import io.emeraldpay.dshackle.upstream.ethereum.EthereumUpstreamSubscriptions -import io.emeraldpay.dshackle.upstream.ethereum.NoEthereumUpstreamSubscriptions +import io.emeraldpay.dshackle.upstream.NoIngressSubscription +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.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse @@ -38,7 +39,7 @@ class EthereumConnectorMock implements EthereumConnector { } @Override - EthereumUpstreamSubscriptions getUpstreamSubscriptions() { - return NoEthereumUpstreamSubscriptions.DEFAULT + EthereumIngressSubscription getIngressSubscription() { + return NoEthereumIngressSubscription.DEFAULT } } \ No newline at end of file diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/EthereumSubscriptionApiSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/EthereumEgressSubscriptionSpec.groovy similarity index 80% rename from src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/EthereumSubscriptionApiSpec.groovy rename to src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/EthereumEgressSubscriptionSpec.groovy index 05b5cef2..efc02a7d 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/EthereumSubscriptionApiSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/EthereumEgressSubscriptionSpec.groovy @@ -21,11 +21,11 @@ import io.emeraldpay.etherjar.domain.Address import io.emeraldpay.etherjar.hex.Hex32 import spock.lang.Specification -class EthereumSubscriptionApiSpec extends Specification { +class EthereumEgressSubscriptionSpec extends Specification { def "read empty logs request"() { setup: - def ethereumSubscribe = new EthereumSubscriptionApi(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) + def ethereumSubscribe = new EthereumEgressSubscription(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) when: def act = ethereumSubscribe.readLogsRequest([:]) @@ -36,7 +36,7 @@ class EthereumSubscriptionApiSpec extends Specification { def "read single address logs request"() { setup: - def ethereumSubscribe = new EthereumSubscriptionApi(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) + def ethereumSubscribe = new EthereumEgressSubscription(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) when: def act = ethereumSubscribe.readLogsRequest([ address: "0x829bd824b016326a401d083b33d092293333a830" @@ -61,7 +61,7 @@ class EthereumSubscriptionApiSpec extends Specification { def "ignores invalid address for logs request"() { setup: - def ethereumSubscribe = new EthereumSubscriptionApi(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) + def ethereumSubscribe = new EthereumEgressSubscription(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) when: def act = ethereumSubscribe.readLogsRequest([ address: "829bd824b016326a401d083b33d092293333a830" @@ -74,7 +74,7 @@ class EthereumSubscriptionApiSpec extends Specification { def "read multi address logs request"() { setup: - def ethereumSubscribe = new EthereumSubscriptionApi(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) + def ethereumSubscribe = new EthereumEgressSubscription(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) when: def act = ethereumSubscribe.readLogsRequest([ address: ["0x829bd824b016326a401d083b33d092293333a830", "0x401d083b33d092293333a83829bd824b016326a0"] @@ -90,7 +90,7 @@ class EthereumSubscriptionApiSpec extends Specification { def "read single topic logs request"() { setup: - def ethereumSubscribe = new EthereumSubscriptionApi(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) + def ethereumSubscribe = new EthereumEgressSubscription(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) when: def act = ethereumSubscribe.readLogsRequest([ topics: "0xddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef" @@ -115,7 +115,7 @@ class EthereumSubscriptionApiSpec extends Specification { def "read invalid topic for request"() { setup: - def ethereumSubscribe = new EthereumSubscriptionApi(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) + def ethereumSubscribe = new EthereumEgressSubscription(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) when: def act = ethereumSubscribe.readLogsRequest([ topics: [ @@ -133,7 +133,7 @@ class EthereumSubscriptionApiSpec extends Specification { def "read multi topic logs request"() { setup: - def ethereumSubscribe = new EthereumSubscriptionApi(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) + def ethereumSubscribe = new EthereumEgressSubscription(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) when: def act = ethereumSubscribe.readLogsRequest([ topics: [ @@ -152,7 +152,7 @@ class EthereumSubscriptionApiSpec extends Specification { def "read full logs request"() { setup: - def ethereumSubscribe = new EthereumSubscriptionApi(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) + def ethereumSubscribe = new EthereumEgressSubscription(TestingCommons.emptyMultistream() as EthereumPosMultiStream, Stub(PendingTxesSource)) when: def act = ethereumSubscribe.readLogsRequest([ address: "0x298d492e8c1d909d3f63bc4a36c66c64acb3d695",