Revert "Revert "Native subscribe filter""

This commit is contained in:
Vyacheslav Shebanov
2022-09-22 13:26:55 +03:00
committed by GitHub
parent 283f44ac49
commit ac715861cb
19 changed files with 119 additions and 72 deletions

4
.gitmodules vendored
View File

@@ -1,6 +1,6 @@
[submodule "emerald-java-client"] [submodule "emerald-java-client"]
path = 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"] [submodule "dshackle-cli/grpc"]
path = dshackle-cli/grpc path = dshackle-cli/grpc
url = https://github.com/emeraldpay/emerald-grpc.git url = https://github.com/p2p-org/emerald-grpc.git

View File

@@ -31,7 +31,7 @@ cglib-nodep = "cglib:cglib-nodep:3.3.0"
detekt-formatting = { module = "io.gitlab.arturbosch.detekt:detekt-formatting", version.ref = "detekt" } 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" 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-buffer = { module = "io.netty:netty-buffer", version.ref = "netty" }
netty-tcnative-core = { module = "io.netty:netty-tcnative", version.ref = "netty-tcnative" } 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-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" zeromq = "org.zeromq:jeromq:0.5.2"

View File

@@ -3,6 +3,6 @@ enableFeaturePreview("TYPESAFE_PROJECT_ACCESSORS")
includeBuild('./emerald-java-client') { includeBuild('./emerald-java-client') {
dependencySubstitution { 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(':')
} }
} }

View File

@@ -17,6 +17,7 @@ package io.emeraldpay.dshackle.proxy
import com.google.protobuf.ByteString import com.google.protobuf.ByteString
import io.emeraldpay.api.proto.BlockchainOuterClass import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.api.proto.BlockchainOuterClass.Selector
import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.config.ProxyConfig import io.emeraldpay.dshackle.config.ProxyConfig
import io.emeraldpay.dshackle.monitoring.accesslog.AccessHandlerHttp import io.emeraldpay.dshackle.monitoring.accesslog.AccessHandlerHttp
@@ -139,7 +140,7 @@ class WebsocketHandler(
} }
// produce actual responses // produce actual responses
val responses = nativeSubscribe val responses = nativeSubscribe
.subscribe(blockchain, methodParams.first, methodParams.second) .subscribe(blockchain, methodParams.first, methodParams.second, io.emeraldpay.dshackle.upstream.Selector.empty)
.map { event -> .map { event ->
WsSubscriptionResponse(params = WsSubscriptionData(event, subscriptionId)) WsSubscriptionResponse(params = WsSubscriptionData(event, subscriptionId))
} }

View File

@@ -20,6 +20,7 @@ import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.SilentException import io.emeraldpay.dshackle.SilentException
import io.emeraldpay.dshackle.upstream.MultistreamHolder import io.emeraldpay.dshackle.upstream.MultistreamHolder
import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.ethereum.EthereumLikeMultistream import io.emeraldpay.dshackle.upstream.ethereum.EthereumLikeMultistream
import io.emeraldpay.grpc.BlockchainType import io.emeraldpay.grpc.BlockchainType
import io.emeraldpay.grpc.Chain import io.emeraldpay.grpc.Chain
@@ -50,20 +51,17 @@ open class NativeSubscribe(
.onErrorMap(this@NativeSubscribe::convertToStatus) .onErrorMap(this@NativeSubscribe::convertToStatus)
} }
fun start(it: BlockchainOuterClass.NativeSubscribeRequest): Publisher<out Any> { fun start(request: BlockchainOuterClass.NativeSubscribeRequest): Publisher<out Any> {
val chain = Chain.byId(it.chainValue) val chain = Chain.byId(request.chainValue)
if (BlockchainType.from(chain) != BlockchainType.ETHEREUM_POS && BlockchainType.from(chain) != BlockchainType.ETHEREUM) { if (BlockchainType.from(chain) != BlockchainType.ETHEREUM_POS && BlockchainType.from(chain) != BlockchainType.ETHEREUM) {
return Mono.error(UnsupportedOperationException("Native subscribe is not supported for ${chain.chainCode}")) return Mono.error(UnsupportedOperationException("Native subscribe is not supported for ${chain.chainCode}"))
} }
val method = it.method val method = request.method
val params: Any? = it.payload?.let { payload -> val params: Any? = request.payload?.takeIf { !it.isEmpty }?.let {
if (payload.size() > 0) { objectMapper.readValue(it.newInput(), Map::class.java)
objectMapper.readValue(payload.newInput(), Map::class.java)
} else {
null
}
} }
return subscribe(chain, method, params) val matcher = Selector.convertToMatcher(request.selector)
return subscribe(chain, method, params, matcher)
} }
fun convertToStatus(t: Throwable) = when (t) { fun convertToStatus(t: Throwable) = when (t) {
@@ -81,11 +79,11 @@ open class NativeSubscribe(
} }
} }
open fun subscribe(chain: Chain, method: String, params: Any?): Flux<out Any> { open fun subscribe(chain: Chain, method: String, params: Any?, matcher: Selector.Matcher): Flux<out Any> {
val up = multistreamHolder.getUpstream(chain) ?: return Flux.error(SilentException.UnsupportedBlockchain(chain)) val up = multistreamHolder.getUpstream(chain) ?: return Flux.error(SilentException.UnsupportedBlockchain(chain))
return (up as EthereumLikeMultistream) return (up as EthereumLikeMultistream)
.getSubscribe() .getSubscribe()
.subscribe(method, params) .subscribe(method, params, matcher)
} }
fun convertToProto(value: Any): BlockchainOuterClass.NativeSubscribeReplyItem { fun convertToProto(value: Any): BlockchainOuterClass.NativeSubscribeReplyItem {

View File

@@ -20,6 +20,7 @@ import io.emeraldpay.api.proto.Common
import io.emeraldpay.dshackle.SilentException import io.emeraldpay.dshackle.SilentException
import io.emeraldpay.dshackle.config.TokensConfig import io.emeraldpay.dshackle.config.TokensConfig
import io.emeraldpay.dshackle.upstream.MultistreamHolder 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.ERC20Balance
import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream
import io.emeraldpay.etherjar.domain.Address import io.emeraldpay.etherjar.domain.Address
@@ -90,7 +91,8 @@ class TrackERC20Address(
.getSubscribe().logs .getSubscribe().logs
.start( .start(
listOf(tokenDefinition.token.contract), 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) return ethereumAddresses.extract(request.address)

View File

@@ -1,8 +1,12 @@
package io.emeraldpay.dshackle.upstream.ethereum package io.emeraldpay.dshackle.upstream.ethereum
import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstream
interface EthereumLikeMultistream : Upstream { interface EthereumLikeMultistream : Upstream {
fun getReader(): EthereumReader fun getReader(): EthereumReader
fun getSubscribe(): EthereumSubscribe fun getSubscribe(): EthereumSubscribe
fun getHead(mather: Selector.Matcher): Head
} }

View File

@@ -20,6 +20,7 @@ import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.ChainFees import io.emeraldpay.dshackle.upstream.ChainFees
import io.emeraldpay.dshackle.upstream.EmptyHead
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.MergedHead import io.emeraldpay.dshackle.upstream.MergedHead
import io.emeraldpay.dshackle.upstream.Multistream import io.emeraldpay.dshackle.upstream.Multistream
@@ -31,6 +32,7 @@ import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import io.emeraldpay.grpc.Chain import io.emeraldpay.grpc.Chain
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle import org.springframework.context.Lifecycle
import org.springframework.util.ConcurrentReferenceHashMap
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
@Suppress("UNCHECKED_CAST") @Suppress("UNCHECKED_CAST")
@@ -45,6 +47,8 @@ open class EthereumMultistream(
} }
private var head: Head? = null private var head: Head? = null
private val filteredHeads: MutableMap<String, Head> =
ConcurrentReferenceHashMap(16, ConcurrentReferenceHashMap.ReferenceType.WEAK)
private val reader: EthereumReader = EthereumReader(this, this.caches, getMethodsFactory()) private val reader: EthereumReader = EthereumReader(this, this.caches, getMethodsFactory())
private val subscribe = EthereumSubscribe(this) private val subscribe = EthereumSubscribe(this)
@@ -74,6 +78,7 @@ open class EthereumMultistream(
override fun stop() { override fun stop() {
super.stop() super.stop()
reader.stop() reader.stop()
filteredHeads.clear()
} }
override fun isRunning(): Boolean { override fun isRunning(): Boolean {
@@ -118,6 +123,7 @@ open class EthereumMultistream(
lagObserver.start() lagObserver.start()
newHead newHead
} }
filteredHeads[Selector.AnyLabelMatcher().describeInternal()] = head
onHeadUpdated(head) onHeadUpdated(head)
return head return head
} }
@@ -142,6 +148,24 @@ open class EthereumMultistream(
return subscribe 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 { override fun getFeeEstimation(): ChainFees {
return feeEstimation return feeEstimation
} }

View File

@@ -1,5 +1,6 @@
package io.emeraldpay.dshackle.upstream.ethereum 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.ConnectLogs
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.ConnectNewHeads import io.emeraldpay.dshackle.upstream.ethereum.subscribe.ConnectNewHeads
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.ConnectSyncing import io.emeraldpay.dshackle.upstream.ethereum.subscribe.ConnectSyncing
@@ -21,9 +22,9 @@ open class EthereumSubscribe(
private val syncing = ConnectSyncing(upstream) private val syncing = ConnectSyncing(upstream)
@Suppress("UNCHECKED_CAST") @Suppress("UNCHECKED_CAST")
open fun subscribe(method: String, params: Any?): Flux<out Any> { open fun subscribe(method: String, params: Any?, matcher: Selector.Matcher): Flux<out Any> {
if (method == "newHeads") { if (method == "newHeads") {
return newHeads.connect() return newHeads.connect(matcher)
} }
if (method == "logs") { if (method == "logs") {
val paramsMap = try { val paramsMap = try {
@@ -35,7 +36,7 @@ open class EthereumSubscribe(
} catch (t: Throwable) { } catch (t: Throwable) {
return Flux.error(UnsupportedOperationException("Invalid parameter for $method. Error: ${t.message}")) 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") { if (method == "syncing") {
return syncing.connect() return syncing.connect()

View File

@@ -19,16 +19,16 @@ import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.data.BlockId import io.emeraldpay.dshackle.data.BlockId
import io.emeraldpay.dshackle.data.TxId import io.emeraldpay.dshackle.data.TxId
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.ethereum.EthereumLikeMultistream import io.emeraldpay.dshackle.upstream.ethereum.EthereumLikeMultistream
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
import reactor.core.scheduler.Schedulers import reactor.core.scheduler.Schedulers
import java.time.Duration import java.time.Duration
import java.util.LinkedList import java.util.LinkedList
import java.util.concurrent.locks.ReentrantLock import java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.locks.ReentrantReadWriteLock import java.util.concurrent.locks.ReentrantReadWriteLock
import kotlin.concurrent.read import kotlin.concurrent.read
import kotlin.concurrent.withLock
import kotlin.concurrent.write import kotlin.concurrent.write
class ConnectBlockUpdates( class ConnectBlockUpdates(
@@ -46,30 +46,19 @@ class ConnectBlockUpdates(
*/ */
private val history = LinkedList<BlockContainer>() private val history = LinkedList<BlockContainer>()
private val historyUpdateLock = ReentrantReadWriteLock() private val historyUpdateLock = ReentrantReadWriteLock()
private val connected: MutableMap<String, Flux<Update>> = ConcurrentHashMap()
private var connected: Flux<Update>? = null fun connect() = connect(Selector.empty)
private val connectLock = ReentrantLock() fun connect(matcher: Selector.Matcher): Flux<Update> {
return connected.computeIfAbsent(matcher.describeInternal()) { key ->
fun connect(): Flux<Update> { extract(upstream.getHead(matcher))
val current = connected
if (current != null) {
return current
}
connectLock.withLock {
val currentRecheck = connected
if (currentRecheck != null) {
return currentRecheck
}
val created = extract(upstream.getHead())
.publishOn(Schedulers.boundedElastic()) .publishOn(Schedulers.boundedElastic())
.publish() .publish()
.refCount(1, Duration.ofSeconds(60)) .refCount(1, Duration.ofSeconds(60))
.doFinally { .doFinally {
// forget it on disconnect, so next time it's recreated // forget it on disconnect, so next time it's recreated
connected = null connected.remove(key)
} }
connected = created
return created
} }
} }

View File

@@ -15,6 +15,7 @@
*/ */
package io.emeraldpay.dshackle.upstream.ethereum.subscribe 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.EthereumLikeMultistream
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.json.LogMessage import io.emeraldpay.dshackle.upstream.ethereum.subscribe.json.LogMessage
import io.emeraldpay.etherjar.domain.Address import io.emeraldpay.etherjar.domain.Address
@@ -40,17 +41,17 @@ open class ConnectLogs(
private val produceLogs = ProduceLogs(upstream) private val produceLogs = ProduceLogs(upstream)
fun start(): Flux<LogMessage> { fun start(matcher: Selector.Matcher): Flux<LogMessage> {
return produceLogs.produce(connectBlockUpdates.connect()) return produceLogs.produce(connectBlockUpdates.connect(matcher))
} }
open fun start(addresses: List<Address>, topics: List<Hex32>): Flux<LogMessage> { open fun start(addresses: List<Address>, topics: List<Hex32>, matcher: Selector.Matcher): Flux<LogMessage> {
// shortcut to the whole output if we don't have any filters // shortcut to the whole output if we don't have any filters
if (addresses.isEmpty() && topics.isEmpty()) { if (addresses.isEmpty() && topics.isEmpty()) {
return start() return start(matcher)
} }
// filtered output // filtered output
return start() return start(matcher)
.transform(filtered(addresses, topics)) .transform(filtered(addresses, topics))
} }

View File

@@ -15,14 +15,14 @@
*/ */
package io.emeraldpay.dshackle.upstream.ethereum.subscribe 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.EthereumLikeMultistream
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.json.NewHeadMessage import io.emeraldpay.dshackle.upstream.ethereum.subscribe.json.NewHeadMessage
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
import reactor.core.scheduler.Schedulers import reactor.core.scheduler.Schedulers
import java.time.Duration import java.time.Duration
import java.util.concurrent.locks.ReentrantLock import java.util.concurrent.ConcurrentHashMap
import kotlin.concurrent.withLock
/** /**
* Connects/reconnects to the upstream to produce NewHeads messages * Connects/reconnects to the upstream to produce NewHeads messages
@@ -35,30 +35,18 @@ class ConnectNewHeads(
private val log = LoggerFactory.getLogger(ConnectNewHeads::class.java) private val log = LoggerFactory.getLogger(ConnectNewHeads::class.java)
} }
private var connected: Flux<NewHeadMessage>? = null private val connected: MutableMap<String, Flux<NewHeadMessage>> = ConcurrentHashMap()
private val connectLock = ReentrantLock()
fun connect(): Flux<NewHeadMessage> { fun connect(matcher: Selector.Matcher): Flux<NewHeadMessage> =
val current = connected connected.computeIfAbsent(matcher.describeInternal()) { key ->
if (current != null) { ProduceNewHeads(upstream.getHead(matcher))
return current
}
connectLock.withLock {
val currentRecheck = connected
if (currentRecheck != null) {
return currentRecheck
}
val created = ProduceNewHeads(upstream.getHead())
.start() .start()
.publishOn(Schedulers.boundedElastic()) .publishOn(Schedulers.boundedElastic())
.publish() .publish()
.refCount(1, Duration.ofSeconds(60)) .refCount(1, Duration.ofSeconds(60))
.doFinally { .doFinally {
// forget it on disconnect, so next time it's recreated // forget it on disconnect, so next time it's recreated
connected = null connected.remove(key)
} }
connected = created
return created
} }
}
} }

View File

@@ -20,6 +20,7 @@ import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.ChainFees import io.emeraldpay.dshackle.upstream.ChainFees
import io.emeraldpay.dshackle.upstream.EmptyHead
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.MergedHead import io.emeraldpay.dshackle.upstream.MergedHead
import io.emeraldpay.dshackle.upstream.Multistream import io.emeraldpay.dshackle.upstream.Multistream
@@ -31,6 +32,7 @@ import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import io.emeraldpay.grpc.Chain import io.emeraldpay.grpc.Chain
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle import org.springframework.context.Lifecycle
import org.springframework.util.ConcurrentReferenceHashMap
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
@Suppress("UNCHECKED_CAST") @Suppress("UNCHECKED_CAST")
@@ -49,6 +51,8 @@ open class EthereumPosMultiStream(
private val reader: EthereumReader = EthereumReader(this, this.caches, getMethodsFactory()) private val reader: EthereumReader = EthereumReader(this, this.caches, getMethodsFactory())
private val feeEstimation = EthereumPriorityFees(this, reader, 256) private val feeEstimation = EthereumPriorityFees(this, reader, 256)
private val subscribe = EthereumSubscribe(this) private val subscribe = EthereumSubscribe(this)
private val filteredHeads: MutableMap<String, Head> =
ConcurrentReferenceHashMap(16, ConcurrentReferenceHashMap.ReferenceType.WEAK)
init { init {
this.init() this.init()
@@ -69,6 +73,7 @@ open class EthereumPosMultiStream(
override fun stop() { override fun stop() {
super.stop() super.stop()
reader.stop() reader.stop()
filteredHeads.clear()
} }
override fun isRunning(): Boolean { override fun isRunning(): Boolean {
@@ -137,6 +142,24 @@ open class EthereumPosMultiStream(
return subscribe 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 { override fun getFeeEstimation(): ChainFees {
return feeEstimation return feeEstimation
} }

View File

@@ -18,6 +18,7 @@ package io.emeraldpay.dshackle.proxy
import io.emeraldpay.dshackle.monitoring.accesslog.AccessHandlerHttp import io.emeraldpay.dshackle.monitoring.accesslog.AccessHandlerHttp
import io.emeraldpay.dshackle.rpc.NativeCall import io.emeraldpay.dshackle.rpc.NativeCall
import io.emeraldpay.dshackle.rpc.NativeSubscribe import io.emeraldpay.dshackle.rpc.NativeSubscribe
import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.etherjar.rpc.json.RequestJson import io.emeraldpay.etherjar.rpc.json.RequestJson
import io.emeraldpay.grpc.Chain import io.emeraldpay.grpc.Chain
import io.micrometer.core.instrument.Counter import io.micrometer.core.instrument.Counter
@@ -108,7 +109,7 @@ class WebsocketHandlerSpec extends Specification {
def response2 = [foo: 2] def response2 = [foo: 2]
def nativeSubscribe = Mock(NativeSubscribe) { 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( def handler = new WebsocketHandler(
new ReadRpcJson(), new WriteRpcJson(), Stub(NativeCall), nativeSubscribe, requestHandlerFactory, Stub(ProxyServer.RequestMetricsFactory) new ReadRpcJson(), new WriteRpcJson(), Stub(NativeCall), nativeSubscribe, requestHandlerFactory, Stub(ProxyServer.RequestMetricsFactory)

View File

@@ -18,6 +18,7 @@ package io.emeraldpay.dshackle.rpc
import com.google.protobuf.ByteString import com.google.protobuf.ByteString
import io.emeraldpay.api.proto.BlockchainOuterClass import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.dshackle.test.MultistreamHolderMock import io.emeraldpay.dshackle.test.MultistreamHolderMock
import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream
import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream
import io.emeraldpay.dshackle.upstream.ethereum.EthereumSubscribe import io.emeraldpay.dshackle.upstream.ethereum.EthereumSubscribe
@@ -33,7 +34,7 @@ class NativeSubscribeSpec extends Specification {
def "Call with empty params when not provided"() { def "Call with empty params when not provided"() {
setup: setup:
def subscribe = Mock(EthereumSubscribe) { 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) { def up = Mock(EthereumPosMultiStream) {
1 * it.getSubscribe() >> subscribe 1 * it.getSubscribe() >> subscribe
@@ -65,7 +66,7 @@ class NativeSubscribeSpec extends Specification {
params["topics"][0] == "0x7fcf532c15f0a6db0bd6d0e038bea71d30d808c7d98cb3bf7268a95bf5081b65" params["topics"][0] == "0x7fcf532c15f0a6db0bd6d0e038bea71d30d808c7d98cb3bf7268a95bf5081b65"
println("ok: $ok") println("ok: $ok")
ok ok
}) >> Flux.just("{}") }, _ as Selector.AnyLabelMatcher) >> Flux.just("{}")
} }
def up = Mock(EthereumPosMultiStream) { def up = Mock(EthereumPosMultiStream) {
1 * it.getSubscribe() >> subscribe 1 * it.getSubscribe() >> subscribe

View File

@@ -4,6 +4,7 @@ import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.api.proto.Common import io.emeraldpay.api.proto.Common
import io.emeraldpay.dshackle.config.TokensConfig import io.emeraldpay.dshackle.config.TokensConfig
import io.emeraldpay.dshackle.upstream.MultistreamHolder 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.ERC20Balance
import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream
import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream
@@ -170,7 +171,8 @@ class TrackERC20AddressSpec extends Specification {
def logs = Mock(ConnectLogs) { def logs = Mock(ConnectLogs) {
1 * start( 1 * start(
[Address.from("0x54EedeAC495271d0F6B175474E89094C44Da98b9")], [Address.from("0x54EedeAC495271d0F6B175474E89094C44Da98b9")],
[Hex32.from("0xddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef")] [Hex32.from("0xddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef")],
Selector.empty
) >> { args -> ) >> { args ->
println("ConnectLogs.start $args") println("ConnectLogs.start $args")
Flux.fromIterable(events) Flux.fromIterable(events)

View File

@@ -20,6 +20,7 @@ package io.emeraldpay.dshackle.test
import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Multistream 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.BitcoinMultistream
import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinRpcUpstream import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinRpcUpstream
import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.dshackle.upstream.calls.CallMethods
@@ -142,6 +143,14 @@ class MultistreamHolderMock implements MultistreamHolder {
} }
return super.getHead() return super.getHead()
} }
@Override
Head getHead(@NotNull Selector.Matcher mather) {
if (customHead != null) {
return customHead
}
return super.getHead(mather)
}
} }
} }

View File

@@ -19,6 +19,7 @@ import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.data.BlockId import io.emeraldpay.dshackle.data.BlockId
import io.emeraldpay.dshackle.data.TxId import io.emeraldpay.dshackle.data.TxId
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream
import io.emeraldpay.etherjar.domain.BlockHash import io.emeraldpay.etherjar.domain.BlockHash
import io.emeraldpay.etherjar.domain.TransactionId import io.emeraldpay.etherjar.domain.TransactionId
@@ -223,7 +224,7 @@ class ConnectBlockUpdatesSpec extends Specification {
1 * getFlux() >> Flux.never() 1 * getFlux() >> Flux.never()
} }
def up = Mock(EthereumMultistream) { def up = Mock(EthereumMultistream) {
1 * getHead() >> head 1 * getHead(Selector.empty) >> head
} }
def connectBlockUpdates = new ConnectBlockUpdates(up) def connectBlockUpdates = new ConnectBlockUpdates(up)

View File

@@ -2,6 +2,7 @@ package io.emeraldpay.dshackle.upstream.ethereum.subscribe
import io.emeraldpay.dshackle.test.TestingCommons import io.emeraldpay.dshackle.test.TestingCommons
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
import reactor.test.StepVerifier import reactor.test.StepVerifier
@@ -17,12 +18,12 @@ class ConnectNewHeadsSpec extends Specification {
]) ])
} }
def up = Mock(EthereumMultistream) { def up = Mock(EthereumMultistream) {
1 * getHead() >> head 1 * getHead(Selector.empty) >> head
} }
ConnectNewHeads connectNewHeads = new ConnectNewHeads(up) ConnectNewHeads connectNewHeads = new ConnectNewHeads(up)
when: when:
def act1 = connectNewHeads.connect() def act1 = connectNewHeads.connect(Selector.empty)
def act2 = connectNewHeads.connect() def act2 = connectNewHeads.connect(Selector.empty)
then: then:
StepVerifier.create(act1) StepVerifier.create(act1)
.expectNextCount(1) .expectNextCount(1)