From 6ec0c7ac27f48256cd82284f13595e6d54196294 Mon Sep 17 00:00:00 2001 From: Maxksim Fomenkov Date: Mon, 12 Sep 2022 10:46:21 +0300 Subject: [PATCH] added filter for native subscribe --- .gitmodules | 4 +-- gradle/libs.versions.toml | 3 ++- settings.gradle | 2 +- .../dshackle/rpc/NativeSubscribe.kt | 22 +++++++-------- .../ethereum/EthereumLikeMultistream.kt | 4 +++ .../upstream/ethereum/EthereumMultistream.kt | 23 ++++++++++++++++ .../upstream/ethereum/EthereumSubscribe.kt | 7 ++--- .../ethereum/subscribe/ConnectBlockUpdates.kt | 24 +++++------------ .../ethereum/subscribe/ConnectLogs.kt | 11 ++++---- .../ethereum/subscribe/ConnectNewHeads.kt | 27 ++++++------------- .../ethereum_pos/EthereumPosMultiStream.kt | 3 +++ 11 files changed, 70 insertions(+), 60 deletions(-) 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/rpc/NativeSubscribe.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeSubscribe.kt index a2b7b3b2..73d21734 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) { 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 = Selector.empty): 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/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..37fb2583 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 @@ -32,6 +33,7 @@ import io.emeraldpay.grpc.Chain import org.slf4j.LoggerFactory import org.springframework.context.Lifecycle import reactor.core.publisher.Mono +import java.util.concurrent.ConcurrentHashMap @Suppress("UNCHECKED_CAST") open class EthereumMultistream( @@ -45,6 +47,7 @@ open class EthereumMultistream( } private var head: Head? = null + private val filteredHeads: MutableMap = ConcurrentHashMap() private val reader: EthereumReader = EthereumReader(this, this.caches, getMethodsFactory()) private val subscribe = EthereumSubscribe(this) @@ -74,6 +77,7 @@ open class EthereumMultistream( override fun stop() { super.stop() reader.stop() + filteredHeads.clear() } override fun isRunning(): Boolean { @@ -118,6 +122,7 @@ open class EthereumMultistream( lagObserver.start() newHead } + filteredHeads[Selector.AnyLabelMatcher().describeInternal()] = head onHeadUpdated(head) return head } @@ -142,6 +147,24 @@ open class EthereumMultistream( return subscribe } + override fun getHead(mather: Selector.Matcher): Head = + filteredHeads.computeIfAbsent(mather.describeInternal()) { _ -> + 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() + } + } + } + }//TODO track unused heads and remove + 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..757650e4 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,12 +19,14 @@ 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.ConcurrentHashMap import java.util.concurrent.locks.ReentrantLock import java.util.concurrent.locks.ReentrantReadWriteLock import kotlin.concurrent.read @@ -46,30 +48,18 @@ 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(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..a2886fc8 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,19 @@ 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..f73d9c9d 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 @@ -137,6 +137,9 @@ open class EthereumPosMultiStream( return subscribe } + override fun getHead(mather: Selector.Matcher): Head = + getHead() //TODO + override fun getFeeEstimation(): ChainFees { return feeEstimation }