diff --git a/.gitmodules b/.gitmodules index 6f37b1a9..152400db 100644 --- a/.gitmodules +++ b/.gitmodules @@ -1,6 +1,6 @@ [submodule "emerald-java-client"] path = emerald-java-client - url = git@github.com:p2p-org/emerald-java-client.git + url = https://github.com/p2p-org/emerald-java-client.git [submodule "dshackle-cli/grpc"] path = dshackle-cli/grpc - url = https://github.com/emeraldpay/emerald-grpc.git + url = https://github.com/p2p-org/emerald-grpc.git diff --git a/gradle/libs.versions.toml b/gradle/libs.versions.toml index 5924c7f7..02076f63 100644 --- a/gradle/libs.versions.toml +++ b/gradle/libs.versions.toml @@ -31,7 +31,7 @@ cglib-nodep = "cglib:cglib-nodep:3.3.0" detekt-formatting = { module = "io.gitlab.arturbosch.detekt:detekt-formatting", version.ref = "detekt" } -emerald-api = "io.emeraldpay:emerald-api:0.12-alpha.1" +emerald-api = "io.emeraldpay:emerald-api:0.12-alpha.3" equals-verifier = "nl.jqno.equalsverifier:equalsverifier:3.3" @@ -85,6 +85,7 @@ netty-codec-http2 = { module = "io.netty:netty-codec-http2", version.ref = "nett netty-buffer = { module = "io.netty:netty-buffer", version.ref = "netty" } netty-tcnative-core = { module = "io.netty:netty-tcnative", version.ref = "netty-tcnative" } netty-tcnative-boringssl = { module = "io.netty:netty-tcnative-boringssl-static", version.ref = "netty-tcnative" } +netty-macos = "io.netty:netty-resolver-dns-native-macos:4.1.72.Final" zeromq = "org.zeromq:jeromq:0.5.2" diff --git a/settings.gradle b/settings.gradle index c0b940b8..470d0974 100644 --- a/settings.gradle +++ b/settings.gradle @@ -3,6 +3,6 @@ enableFeaturePreview("TYPESAFE_PROJECT_ACCESSORS") includeBuild('./emerald-java-client') { dependencySubstitution { - substitute module('io.emeraldpay:emerald-api:0.12-alpha.1') using project(':') + substitute module('io.emeraldpay:emerald-api:0.12-alpha.3') using project(':') } } \ No newline at end of file diff --git a/src/main/kotlin/io/emeraldpay/dshackle/proxy/WebsocketHandler.kt b/src/main/kotlin/io/emeraldpay/dshackle/proxy/WebsocketHandler.kt index 0de59c53..e15be1ba 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/proxy/WebsocketHandler.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/proxy/WebsocketHandler.kt @@ -17,6 +17,7 @@ package io.emeraldpay.dshackle.proxy import com.google.protobuf.ByteString import io.emeraldpay.api.proto.BlockchainOuterClass +import io.emeraldpay.api.proto.BlockchainOuterClass.Selector import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.config.ProxyConfig import io.emeraldpay.dshackle.monitoring.accesslog.AccessHandlerHttp @@ -139,7 +140,7 @@ class WebsocketHandler( } // produce actual responses val responses = nativeSubscribe - .subscribe(blockchain, methodParams.first, methodParams.second) + .subscribe(blockchain, methodParams.first, methodParams.second, io.emeraldpay.dshackle.upstream.Selector.empty) .map { event -> WsSubscriptionResponse(params = WsSubscriptionData(event, subscriptionId)) } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeSubscribe.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeSubscribe.kt index 20edb41a..1169b03b 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeSubscribe.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeSubscribe.kt @@ -20,6 +20,7 @@ import io.emeraldpay.api.proto.BlockchainOuterClass import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.SilentException import io.emeraldpay.dshackle.upstream.MultistreamHolder +import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.ethereum.EthereumLikeMultistream import io.emeraldpay.grpc.BlockchainType import io.emeraldpay.grpc.Chain @@ -50,20 +51,17 @@ open class NativeSubscribe( .onErrorMap(this@NativeSubscribe::convertToStatus) } - fun start(it: BlockchainOuterClass.NativeSubscribeRequest): Publisher { - val chain = Chain.byId(it.chainValue) + fun start(request: BlockchainOuterClass.NativeSubscribeRequest): Publisher { + val chain = Chain.byId(request.chainValue) if (BlockchainType.from(chain) != BlockchainType.ETHEREUM_POS && BlockchainType.from(chain) != BlockchainType.ETHEREUM) { return Mono.error(UnsupportedOperationException("Native subscribe is not supported for ${chain.chainCode}")) } - val method = it.method - val params: Any? = it.payload?.let { payload -> - if (payload.size() > 0) { - objectMapper.readValue(payload.newInput(), Map::class.java) - } else { - null - } + val method = request.method + val params: Any? = request.payload?.takeIf { !it.isEmpty }?.let { + objectMapper.readValue(it.newInput(), Map::class.java) } - return subscribe(chain, method, params) + val matcher = Selector.convertToMatcher(request.selector) + return subscribe(chain, method, params, matcher) } fun convertToStatus(t: Throwable) = when (t) { @@ -81,11 +79,11 @@ open class NativeSubscribe( } } - open fun subscribe(chain: Chain, method: String, params: Any?): Flux { + open fun subscribe(chain: Chain, method: String, params: Any?, matcher: Selector.Matcher): Flux { val up = multistreamHolder.getUpstream(chain) ?: return Flux.error(SilentException.UnsupportedBlockchain(chain)) return (up as EthereumLikeMultistream) .getSubscribe() - .subscribe(method, params) + .subscribe(method, params, matcher) } fun convertToProto(value: Any): BlockchainOuterClass.NativeSubscribeReplyItem { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackERC20Address.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackERC20Address.kt index 0b40d706..0f47b6e2 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackERC20Address.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackERC20Address.kt @@ -20,6 +20,7 @@ import io.emeraldpay.api.proto.Common import io.emeraldpay.dshackle.SilentException import io.emeraldpay.dshackle.config.TokensConfig import io.emeraldpay.dshackle.upstream.MultistreamHolder +import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.ethereum.ERC20Balance import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream import io.emeraldpay.etherjar.domain.Address @@ -90,7 +91,8 @@ class TrackERC20Address( .getSubscribe().logs .start( listOf(tokenDefinition.token.contract), - listOf(EventId.fromSignature("Transfer", "address", "address", "uint256")) + listOf(EventId.fromSignature("Transfer", "address", "address", "uint256")), + Selector.empty ) return ethereumAddresses.extract(request.address) 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 40d3a437..33f596fa 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumLikeMultistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumLikeMultistream.kt @@ -1,8 +1,12 @@ package io.emeraldpay.dshackle.upstream.ethereum +import io.emeraldpay.dshackle.upstream.Head +import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Upstream interface EthereumLikeMultistream : Upstream { fun getReader(): EthereumReader fun getSubscribe(): EthereumSubscribe + + 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 62df3041..3605976b 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumMultistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumMultistream.kt @@ -20,6 +20,7 @@ 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.MergedHead import io.emeraldpay.dshackle.upstream.Multistream @@ -31,6 +32,7 @@ import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.grpc.Chain import org.slf4j.LoggerFactory import org.springframework.context.Lifecycle +import org.springframework.util.ConcurrentReferenceHashMap import reactor.core.publisher.Mono @Suppress("UNCHECKED_CAST") @@ -45,6 +47,8 @@ open class EthereumMultistream( } private var head: Head? = null + private val filteredHeads: MutableMap = + ConcurrentReferenceHashMap(16, ConcurrentReferenceHashMap.ReferenceType.WEAK) private val reader: EthereumReader = EthereumReader(this, this.caches, getMethodsFactory()) private val subscribe = EthereumSubscribe(this) @@ -74,6 +78,7 @@ open class EthereumMultistream( override fun stop() { super.stop() reader.stop() + filteredHeads.clear() } override fun isRunning(): Boolean { @@ -118,6 +123,7 @@ open class EthereumMultistream( lagObserver.start() newHead } + filteredHeads[Selector.AnyLabelMatcher().describeInternal()] = head onHeadUpdated(head) return head } @@ -142,6 +148,24 @@ open class EthereumMultistream( return subscribe } + override fun getHead(mather: Selector.Matcher): Head = + filteredHeads.computeIfAbsent(mather.describeInternal().intern()) { _ -> + upstreams.filter { mather.matches(it) } + .apply { + log.debug("Found $size upstreams matching [${mather.describeInternal()}]") + } + .map { it.getHead() } + .let { + when (it.size) { + 0 -> EmptyHead() + 1 -> it.first() + else -> MergedHead(it, MostWorkForkChoice()).apply { + start() + } + } + } + } + override fun getFeeEstimation(): ChainFees { return feeEstimation } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumSubscribe.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumSubscribe.kt index cb160fa3..f598a548 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumSubscribe.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumSubscribe.kt @@ -1,5 +1,6 @@ package io.emeraldpay.dshackle.upstream.ethereum +import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.ethereum.subscribe.ConnectLogs import io.emeraldpay.dshackle.upstream.ethereum.subscribe.ConnectNewHeads import io.emeraldpay.dshackle.upstream.ethereum.subscribe.ConnectSyncing @@ -21,9 +22,9 @@ open class EthereumSubscribe( private val syncing = ConnectSyncing(upstream) @Suppress("UNCHECKED_CAST") - open fun subscribe(method: String, params: Any?): Flux { + open fun subscribe(method: String, params: Any?, matcher: Selector.Matcher): Flux { if (method == "newHeads") { - return newHeads.connect() + return newHeads.connect(matcher) } if (method == "logs") { val paramsMap = try { @@ -35,7 +36,7 @@ open class EthereumSubscribe( } catch (t: Throwable) { return Flux.error(UnsupportedOperationException("Invalid parameter for $method. Error: ${t.message}")) } - return logs.start(paramsMap.address, paramsMap.topics) + return logs.start(paramsMap.address, paramsMap.topics, matcher) } if (method == "syncing") { return syncing.connect() diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectBlockUpdates.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectBlockUpdates.kt index 0a424faf..7d635b62 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectBlockUpdates.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectBlockUpdates.kt @@ -19,16 +19,16 @@ import io.emeraldpay.dshackle.data.BlockContainer import io.emeraldpay.dshackle.data.BlockId import io.emeraldpay.dshackle.data.TxId import io.emeraldpay.dshackle.upstream.Head +import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.ethereum.EthereumLikeMultistream import org.slf4j.LoggerFactory import reactor.core.publisher.Flux import reactor.core.scheduler.Schedulers import java.time.Duration import java.util.LinkedList -import java.util.concurrent.locks.ReentrantLock +import java.util.concurrent.ConcurrentHashMap import java.util.concurrent.locks.ReentrantReadWriteLock import kotlin.concurrent.read -import kotlin.concurrent.withLock import kotlin.concurrent.write class ConnectBlockUpdates( @@ -46,30 +46,19 @@ class ConnectBlockUpdates( */ private val history = LinkedList() private val historyUpdateLock = ReentrantReadWriteLock() + private val connected: MutableMap> = ConcurrentHashMap() - private var connected: Flux? = null - private val connectLock = ReentrantLock() - - fun connect(): Flux { - val current = connected - if (current != null) { - return current - } - connectLock.withLock { - val currentRecheck = connected - if (currentRecheck != null) { - return currentRecheck - } - val created = extract(upstream.getHead()) + fun connect() = connect(Selector.empty) + fun connect(matcher: Selector.Matcher): Flux { + return connected.computeIfAbsent(matcher.describeInternal()) { key -> + extract(upstream.getHead(matcher)) .publishOn(Schedulers.boundedElastic()) .publish() .refCount(1, Duration.ofSeconds(60)) .doFinally { // forget it on disconnect, so next time it's recreated - connected = null + connected.remove(key) } - connected = created - return created } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectLogs.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectLogs.kt index 9482a10f..d03be557 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectLogs.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectLogs.kt @@ -15,6 +15,7 @@ */ package io.emeraldpay.dshackle.upstream.ethereum.subscribe +import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.ethereum.EthereumLikeMultistream import io.emeraldpay.dshackle.upstream.ethereum.subscribe.json.LogMessage import io.emeraldpay.etherjar.domain.Address @@ -40,17 +41,17 @@ open class ConnectLogs( private val produceLogs = ProduceLogs(upstream) - fun start(): Flux { - return produceLogs.produce(connectBlockUpdates.connect()) + fun start(matcher: Selector.Matcher): Flux { + return produceLogs.produce(connectBlockUpdates.connect(matcher)) } - open fun start(addresses: List
, topics: List): Flux { + open fun start(addresses: List
, topics: List, matcher: Selector.Matcher): Flux { // shortcut to the whole output if we don't have any filters if (addresses.isEmpty() && topics.isEmpty()) { - return start() + return start(matcher) } // filtered output - return start() + return start(matcher) .transform(filtered(addresses, topics)) } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectNewHeads.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectNewHeads.kt index c3b7bb4a..b91c6e8b 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectNewHeads.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectNewHeads.kt @@ -15,14 +15,14 @@ */ package io.emeraldpay.dshackle.upstream.ethereum.subscribe +import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.ethereum.EthereumLikeMultistream import io.emeraldpay.dshackle.upstream.ethereum.subscribe.json.NewHeadMessage import org.slf4j.LoggerFactory import reactor.core.publisher.Flux import reactor.core.scheduler.Schedulers import java.time.Duration -import java.util.concurrent.locks.ReentrantLock -import kotlin.concurrent.withLock +import java.util.concurrent.ConcurrentHashMap /** * Connects/reconnects to the upstream to produce NewHeads messages @@ -35,30 +35,18 @@ class ConnectNewHeads( private val log = LoggerFactory.getLogger(ConnectNewHeads::class.java) } - private var connected: Flux? = null - private val connectLock = ReentrantLock() + private val connected: MutableMap> = ConcurrentHashMap() - fun connect(): Flux { - val current = connected - if (current != null) { - return current - } - connectLock.withLock { - val currentRecheck = connected - if (currentRecheck != null) { - return currentRecheck - } - val created = ProduceNewHeads(upstream.getHead()) + fun connect(matcher: Selector.Matcher): Flux = + connected.computeIfAbsent(matcher.describeInternal()) { key -> + ProduceNewHeads(upstream.getHead(matcher)) .start() .publishOn(Schedulers.boundedElastic()) .publish() .refCount(1, Duration.ofSeconds(60)) .doFinally { // forget it on disconnect, so next time it's recreated - connected = null + connected.remove(key) } - connected = created - return created } - } } 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 d263ff1c..60326bc3 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 @@ -20,6 +20,7 @@ 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.MergedHead import io.emeraldpay.dshackle.upstream.Multistream @@ -31,6 +32,7 @@ import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.grpc.Chain import org.slf4j.LoggerFactory import org.springframework.context.Lifecycle +import org.springframework.util.ConcurrentReferenceHashMap import reactor.core.publisher.Mono @Suppress("UNCHECKED_CAST") @@ -49,6 +51,8 @@ open class EthereumPosMultiStream( private val reader: EthereumReader = EthereumReader(this, this.caches, getMethodsFactory()) private val feeEstimation = EthereumPriorityFees(this, reader, 256) private val subscribe = EthereumSubscribe(this) + private val filteredHeads: MutableMap = + ConcurrentReferenceHashMap(16, ConcurrentReferenceHashMap.ReferenceType.WEAK) init { this.init() @@ -69,6 +73,7 @@ open class EthereumPosMultiStream( override fun stop() { super.stop() reader.stop() + filteredHeads.clear() } override fun isRunning(): Boolean { @@ -137,6 +142,24 @@ open class EthereumPosMultiStream( return subscribe } + override fun getHead(mather: Selector.Matcher): Head = + filteredHeads.computeIfAbsent(mather.describeInternal().intern()) { _ -> + upstreams.filter { mather.matches(it) } + .apply { + log.debug("Found $size upstreams matching [${mather.describeInternal()}]") + } + .map { it.getHead() } + .let { + when (it.size) { + 0 -> EmptyHead() + 1 -> it.first() + else -> MergedHead(it, PriorityForkChoice()).apply { + start() + } + } + } + } + override fun getFeeEstimation(): ChainFees { return feeEstimation } diff --git a/src/test/groovy/io/emeraldpay/dshackle/proxy/WebsocketHandlerSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/proxy/WebsocketHandlerSpec.groovy index 815bc4a6..235f4778 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/proxy/WebsocketHandlerSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/proxy/WebsocketHandlerSpec.groovy @@ -18,6 +18,7 @@ package io.emeraldpay.dshackle.proxy import io.emeraldpay.dshackle.monitoring.accesslog.AccessHandlerHttp import io.emeraldpay.dshackle.rpc.NativeCall import io.emeraldpay.dshackle.rpc.NativeSubscribe +import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.etherjar.rpc.json.RequestJson import io.emeraldpay.grpc.Chain import io.micrometer.core.instrument.Counter @@ -108,7 +109,7 @@ class WebsocketHandlerSpec extends Specification { def response2 = [foo: 2] def nativeSubscribe = Mock(NativeSubscribe) { - 1 * it.subscribe(Chain.ETHEREUM, "foo_test", null) >> Flux.fromIterable([response1, response2]) + 1 * it.subscribe(Chain.ETHEREUM, "foo_test", null, Selector.empty) >> Flux.fromIterable([response1, response2]) } def handler = new WebsocketHandler( new ReadRpcJson(), new WriteRpcJson(), Stub(NativeCall), nativeSubscribe, requestHandlerFactory, Stub(ProxyServer.RequestMetricsFactory) diff --git a/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeSubscribeSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeSubscribeSpec.groovy index 888df4b8..b4e94ff9 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeSubscribeSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeSubscribeSpec.groovy @@ -18,6 +18,7 @@ package io.emeraldpay.dshackle.rpc import com.google.protobuf.ByteString import io.emeraldpay.api.proto.BlockchainOuterClass import io.emeraldpay.dshackle.test.MultistreamHolderMock +import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream import io.emeraldpay.dshackle.upstream.ethereum.EthereumSubscribe @@ -33,7 +34,7 @@ class NativeSubscribeSpec extends Specification { def "Call with empty params when not provided"() { setup: def subscribe = Mock(EthereumSubscribe) { - 1 * it.subscribe("newHeads", null) >> Flux.just("{}") + 1 * it.subscribe("newHeads", null, _ as Selector.AnyLabelMatcher) >> Flux.just("{}") } def up = Mock(EthereumPosMultiStream) { 1 * it.getSubscribe() >> subscribe @@ -65,7 +66,7 @@ class NativeSubscribeSpec extends Specification { params["topics"][0] == "0x7fcf532c15f0a6db0bd6d0e038bea71d30d808c7d98cb3bf7268a95bf5081b65" println("ok: $ok") ok - }) >> Flux.just("{}") + }, _ as Selector.AnyLabelMatcher) >> Flux.just("{}") } def up = Mock(EthereumPosMultiStream) { 1 * it.getSubscribe() >> subscribe diff --git a/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackERC20AddressSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackERC20AddressSpec.groovy index ecd9d7f1..2a347a9a 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackERC20AddressSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackERC20AddressSpec.groovy @@ -4,6 +4,7 @@ import io.emeraldpay.api.proto.BlockchainOuterClass import io.emeraldpay.api.proto.Common import io.emeraldpay.dshackle.config.TokensConfig import io.emeraldpay.dshackle.upstream.MultistreamHolder +import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.ethereum.ERC20Balance import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream @@ -170,7 +171,8 @@ class TrackERC20AddressSpec extends Specification { def logs = Mock(ConnectLogs) { 1 * start( [Address.from("0x54EedeAC495271d0F6B175474E89094C44Da98b9")], - [Hex32.from("0xddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef")] + [Hex32.from("0xddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef")], + Selector.empty ) >> { args -> println("ConnectLogs.start $args") Flux.fromIterable(events) diff --git a/src/test/groovy/io/emeraldpay/dshackle/test/MultistreamHolderMock.groovy b/src/test/groovy/io/emeraldpay/dshackle/test/MultistreamHolderMock.groovy index cd75af82..59414f04 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/test/MultistreamHolderMock.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/test/MultistreamHolderMock.groovy @@ -20,6 +20,7 @@ package io.emeraldpay.dshackle.test import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Multistream +import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinMultistream import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinRpcUpstream import io.emeraldpay.dshackle.upstream.calls.CallMethods @@ -142,6 +143,14 @@ class MultistreamHolderMock implements MultistreamHolder { } return super.getHead() } + + @Override + Head getHead(@NotNull Selector.Matcher mather) { + if (customHead != null) { + return customHead + } + return super.getHead(mather) + } } } diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectBlockUpdatesSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectBlockUpdatesSpec.groovy index 697fd263..e9557335 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectBlockUpdatesSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectBlockUpdatesSpec.groovy @@ -19,6 +19,7 @@ import io.emeraldpay.dshackle.data.BlockContainer import io.emeraldpay.dshackle.data.BlockId import io.emeraldpay.dshackle.data.TxId import io.emeraldpay.dshackle.upstream.Head +import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream import io.emeraldpay.etherjar.domain.BlockHash import io.emeraldpay.etherjar.domain.TransactionId @@ -223,7 +224,7 @@ class ConnectBlockUpdatesSpec extends Specification { 1 * getFlux() >> Flux.never() } def up = Mock(EthereumMultistream) { - 1 * getHead() >> head + 1 * getHead(Selector.empty) >> head } def connectBlockUpdates = new ConnectBlockUpdates(up) diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectNewHeadsSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectNewHeadsSpec.groovy index 6a62bf6b..25781ccc 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectNewHeadsSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/subscribe/ConnectNewHeadsSpec.groovy @@ -2,6 +2,7 @@ package io.emeraldpay.dshackle.upstream.ethereum.subscribe import io.emeraldpay.dshackle.test.TestingCommons import io.emeraldpay.dshackle.upstream.Head +import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream import reactor.core.publisher.Flux import reactor.test.StepVerifier @@ -17,12 +18,12 @@ class ConnectNewHeadsSpec extends Specification { ]) } def up = Mock(EthereumMultistream) { - 1 * getHead() >> head + 1 * getHead(Selector.empty) >> head } ConnectNewHeads connectNewHeads = new ConnectNewHeads(up) when: - def act1 = connectNewHeads.connect() - def act2 = connectNewHeads.connect() + def act1 = connectNewHeads.connect(Selector.empty) + def act2 = connectNewHeads.connect(Selector.empty) then: StepVerifier.create(act1) .expectNextCount(1)