From 443839062f8306ee9d8b210fe3f4572a91facc9a Mon Sep 17 00:00:00 2001 From: a10zn8 Date: Tue, 15 Nov 2022 20:52:15 +0400 Subject: [PATCH 1/7] multistreams as beans --- .../config/context/MultistreamsConfig.kt | 27 ++++++++ .../upstream/CurrentMultistreamHolder.kt | 69 +++++++------------ .../dshackle/test/TestingCommons.groovy | 12 ++++ .../CurrentMultistreamHolderSpec.groovy | 8 +-- 4 files changed, 69 insertions(+), 47 deletions(-) create mode 100644 src/main/kotlin/io/emeraldpay/dshackle/config/context/MultistreamsConfig.kt diff --git a/src/main/kotlin/io/emeraldpay/dshackle/config/context/MultistreamsConfig.kt b/src/main/kotlin/io/emeraldpay/dshackle/config/context/MultistreamsConfig.kt new file mode 100644 index 00000000..7ffdd93d --- /dev/null +++ b/src/main/kotlin/io/emeraldpay/dshackle/config/context/MultistreamsConfig.kt @@ -0,0 +1,27 @@ +package io.emeraldpay.dshackle.config.context + +import io.emeraldpay.dshackle.cache.CachesFactory +import io.emeraldpay.dshackle.upstream.Multistream +import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinMultistream +import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream +import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream +import io.emeraldpay.grpc.BlockchainType +import io.emeraldpay.grpc.Chain +import org.springframework.context.annotation.Bean +import org.springframework.context.annotation.Configuration + +@Configuration +class MultistreamsConfig { + @Bean + fun allMultistreams(cachesFactory: CachesFactory): List { + return Chain.values() + .mapNotNull { chain -> + when (BlockchainType.from(chain)) { + BlockchainType.EVM_POS -> EthereumPosMultiStream(chain, ArrayList(), cachesFactory.getCaches(chain)) + BlockchainType.EVM_POW -> EthereumMultistream(chain, ArrayList(), cachesFactory.getCaches(chain)) + BlockchainType.BITCOIN -> BitcoinMultistream(chain, ArrayList(), cachesFactory.getCaches(chain)) + else -> null + } + } + } +} diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt index 7587d6c2..2de935f6 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt @@ -17,39 +17,35 @@ package io.emeraldpay.dshackle.upstream import io.emeraldpay.dshackle.cache.CachesEnabled -import io.emeraldpay.dshackle.cache.CachesFactory import io.emeraldpay.dshackle.startup.UpstreamChange -import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinMultistream import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinUpstream import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.dshackle.upstream.calls.DefaultBitcoinMethods import io.emeraldpay.dshackle.upstream.calls.DefaultEthereumMethods -import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream -import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosUpstream import io.emeraldpay.dshackle.upstream.ethereum.EthereumUpstream import io.emeraldpay.grpc.BlockchainType import io.emeraldpay.grpc.Chain import org.slf4j.LoggerFactory -import org.springframework.beans.factory.annotation.Autowired -import org.springframework.stereotype.Repository +import org.springframework.stereotype.Component import reactor.core.publisher.Flux import reactor.core.publisher.Sinks -import java.util.Collections -import java.util.concurrent.Callable +import java.util.* import java.util.concurrent.ConcurrentHashMap import java.util.concurrent.locks.ReentrantLock import javax.annotation.PreDestroy import kotlin.concurrent.withLock -@Repository +@Component open class CurrentMultistreamHolder( - @Autowired private val cachesFactory: CachesFactory + private val multistreams: List ) : MultistreamHolder { private val log = LoggerFactory.getLogger(CurrentMultistreamHolder::class.java) - private val chainMapping = ConcurrentHashMap() + private val chainMapping = ConcurrentHashMap().apply { + multistreams.forEach { this[it.chain] = it } + } private val chainsBus = Sinks.many() .multicast() .directBestEffort() @@ -64,28 +60,22 @@ open class CurrentMultistreamHolder( when (BlockchainType.from(chain)) { BlockchainType.EVM_POW -> { val up = change.upstream.cast(EthereumUpstream::class.java) - val current = chainMapping[chain] - val factory = Callable { - EthereumMultistream(chain, ArrayList(), cachesFactory.getCaches(chain)) - } - processUpdate(change, up, current, factory) + val current = chainMapping.getValue(chain) + processUpdate(change, up, current) } + BlockchainType.EVM_POS -> { val up = change.upstream.cast(EthereumPosUpstream::class.java) - val current = chainMapping[chain] - val factory = Callable { - EthereumPosMultiStream(chain, ArrayList(), cachesFactory.getCaches(chain)) - } - processUpdate(change, up, current, factory) + val current = chainMapping.getValue(chain) + processUpdate(change, up, current) } + BlockchainType.BITCOIN -> { val up = change.upstream.cast(BitcoinUpstream::class.java) - val current = chainMapping[chain] - val factory = Callable { - BitcoinMultistream(chain, ArrayList(), cachesFactory.getCaches(chain)) - } - processUpdate(change, up, current, factory) + val current = chainMapping.getValue(chain) + processUpdate(change, up, current) } + else -> { log.error("Update for unsupported chain: $chain") } @@ -96,27 +86,17 @@ open class CurrentMultistreamHolder( } } - fun processUpdate(change: UpstreamChange, up: Upstream, current: Multistream?, factory: Callable) { + fun processUpdate(change: UpstreamChange, up: Upstream, current: Multistream) { val chain = change.chain if (change.type == UpstreamChange.ChangeType.REMOVED) { - current?.removeUpstream(up.getId()) + current.removeUpstream(up.getId()) log.info("Upstream ${change.upstream.getId()} with chain $chain has been removed") } else { - if (current == null) { - val created = factory.call() - if (up is CachesEnabled) { - up.setCaches(created.caches) - } - created.addUpstream(up) - created.start() - chainMapping[chain] = created - chainsBus.tryEmitNext(chain) - } else { - if (up is CachesEnabled) { - up.setCaches(current.caches) - } - current.addUpstream(up) + if (up is CachesEnabled) { + up.setCaches(current.caches) } + current.addUpstream(up) + if (!callTargets.containsKey(chain)) { setupDefaultMethods(chain) } @@ -129,7 +109,10 @@ open class CurrentMultistreamHolder( } override fun getAvailable(): List { - return Collections.unmodifiableList(chainMapping.keys.toList()) + return multistreams.asSequence() + .filter { it.isAvailable() } + .map { it.chain } + .toList() } override fun observeChains(): Flux { diff --git a/src/test/groovy/io/emeraldpay/dshackle/test/TestingCommons.groovy b/src/test/groovy/io/emeraldpay/dshackle/test/TestingCommons.groovy index fbec4d77..6ad08167 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/test/TestingCommons.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/test/TestingCommons.groovy @@ -27,6 +27,7 @@ import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.upstream.Multistream import io.emeraldpay.dshackle.upstream.calls.DirectCallMethods import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream +import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosUpstream import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.etherjar.domain.BlockHash @@ -90,6 +91,17 @@ class TestingCommons { return new CachesFactory(new CacheConfig()) } + static List defaultMultistreams() { + return [ + multistreamWithoutUpstreams(Chain.ETHEREUM), + multistreamWithoutUpstreams(Chain.ETHEREUM_CLASSIC) + ] + } + + static Multistream multistreamWithoutUpstreams(Chain chain) { + return new EthereumPosMultiStream(chain, [], emptyCaches().getCaches(chain)) + } + static FileResolver fileResolver() { return new FileResolver(new File("src/test/resources")) } diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolderSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolderSpec.groovy index 2f38eecd..0ec23584 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolderSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolderSpec.groovy @@ -26,7 +26,7 @@ class CurrentMultistreamHolderSpec extends Specification { def "add upstream"() { setup: - def current = new CurrentMultistreamHolder(TestingCommons.emptyCaches()) + def current = new CurrentMultistreamHolder(TestingCommons.defaultMultistreams()) def up = new EthereumPosRpcUpstreamMock("test", Chain.ETHEREUM, TestingCommons.api()) when: current.update(new UpstreamChange(Chain.ETHEREUM, up, UpstreamChange.ChangeType.ADDED)) @@ -37,7 +37,7 @@ class CurrentMultistreamHolderSpec extends Specification { def "add multiple upstreams"() { setup: - def current = new CurrentMultistreamHolder(TestingCommons.emptyCaches()) + def current = new CurrentMultistreamHolder(TestingCommons.defaultMultistreams()) def up1 = new EthereumPosRpcUpstreamMock("test1", Chain.ETHEREUM, TestingCommons.api()) def up2 = new EthereumRpcUpstreamMock("test2", Chain.ETHEREUM_CLASSIC, TestingCommons.api()) def up3 = new EthereumPosRpcUpstreamMock("test3", Chain.ETHEREUM, TestingCommons.api()) @@ -53,7 +53,7 @@ class CurrentMultistreamHolderSpec extends Specification { def "remove upstream"() { setup: - def current = new CurrentMultistreamHolder(TestingCommons.emptyCaches()) + def current = new CurrentMultistreamHolder(TestingCommons.defaultMultistreams()) def up1 = new EthereumPosRpcUpstreamMock("test1", Chain.ETHEREUM, TestingCommons.api()) def up2 = new EthereumRpcUpstreamMock("test2", Chain.ETHEREUM_CLASSIC, TestingCommons.api()) def up3 = new EthereumPosRpcUpstreamMock("test3", Chain.ETHEREUM, TestingCommons.api()) @@ -71,7 +71,7 @@ class CurrentMultistreamHolderSpec extends Specification { def "available after adding"() { setup: - def current = new CurrentMultistreamHolder(TestingCommons.emptyCaches()) + def current = new CurrentMultistreamHolder(TestingCommons.defaultMultistreams()) def up1 = new EthereumPosRpcUpstreamMock("test1", Chain.ETHEREUM, TestingCommons.api()) when: From 53e170b0af50d2366970821e826b08b8574b699b Mon Sep 17 00:00:00 2001 From: a10zn8 Date: Wed, 16 Nov 2022 14:11:16 +0400 Subject: [PATCH 2/7] introduction of upstream update event --- build.gradle | 3 + .../config/context/MultistreamsConfig.kt | 27 ++++++- .../dshackle/startup/ConfiguredUpstreams.kt | 30 +++---- ...streamChange.kt => UpstreamChangeEvent.kt} | 2 +- .../dshackle/upstream/CallTargetsHolder.kt | 29 +++++++ .../upstream/CurrentMultistreamHolder.kt | 79 +------------------ .../dshackle/upstream/Multistream.kt | 23 +++++- .../dshackle/upstream/MultistreamHolder.kt | 2 - .../upstream/bitcoin/BitcoinMultistream.kt | 14 +--- .../upstream/ethereum/EthereumMultistream.kt | 13 +-- .../ethereum_pos/EthereumPosMultiStream.kt | 13 +-- .../dshackle/upstream/grpc/GrpcUpstreams.kt | 30 +++---- .../startup/ConfiguredUpstreamsSpec.groovy | 34 +++++--- .../test/MultistreamHolderMock.groovy | 15 +--- .../dshackle/test/TestingCommons.groovy | 14 +++- .../CurrentMultistreamHolderSpec.groovy | 21 ++--- .../dshackle/upstream/MultistreamSpec.groovy | 8 +- 17 files changed, 172 insertions(+), 185 deletions(-) rename src/main/kotlin/io/emeraldpay/dshackle/startup/{UpstreamChange.kt => UpstreamChangeEvent.kt} (98%) create mode 100644 src/main/kotlin/io/emeraldpay/dshackle/upstream/CallTargetsHolder.kt diff --git a/build.gradle b/build.gradle index 03cc3cb6..d9363198 100644 --- a/build.gradle +++ b/build.gradle @@ -289,3 +289,6 @@ detekt { tasks.withType(Detekt).configureEach { jvmTarget = "13" } + +// formats code for each build +tasks.findByName("ktlintCheck").dependsOn("ktlintFormat") diff --git a/src/main/kotlin/io/emeraldpay/dshackle/config/context/MultistreamsConfig.kt b/src/main/kotlin/io/emeraldpay/dshackle/config/context/MultistreamsConfig.kt index 7ffdd93d..cfa80f3a 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/config/context/MultistreamsConfig.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/config/context/MultistreamsConfig.kt @@ -1,6 +1,7 @@ package io.emeraldpay.dshackle.config.context import io.emeraldpay.dshackle.cache.CachesFactory +import io.emeraldpay.dshackle.upstream.CallTargetsHolder import io.emeraldpay.dshackle.upstream.Multistream import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinMultistream import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream @@ -13,13 +14,31 @@ import org.springframework.context.annotation.Configuration @Configuration class MultistreamsConfig { @Bean - fun allMultistreams(cachesFactory: CachesFactory): List { + fun allMultistreams( + cachesFactory: CachesFactory, + callTargetsHolder: CallTargetsHolder + ): List { return Chain.values() .mapNotNull { chain -> when (BlockchainType.from(chain)) { - BlockchainType.EVM_POS -> EthereumPosMultiStream(chain, ArrayList(), cachesFactory.getCaches(chain)) - BlockchainType.EVM_POW -> EthereumMultistream(chain, ArrayList(), cachesFactory.getCaches(chain)) - BlockchainType.BITCOIN -> BitcoinMultistream(chain, ArrayList(), cachesFactory.getCaches(chain)) + BlockchainType.EVM_POS -> EthereumPosMultiStream( + chain, + ArrayList(), + cachesFactory.getCaches(chain), + callTargetsHolder + ) + BlockchainType.EVM_POW -> EthereumMultistream( + chain, + ArrayList(), + cachesFactory.getCaches(chain), + callTargetsHolder + ) + BlockchainType.BITCOIN -> BitcoinMultistream( + chain, + ArrayList(), + cachesFactory.getCaches(chain), + callTargetsHolder + ) else -> null } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt b/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt index 3a86c6ac..8fc7e036 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt @@ -21,13 +21,7 @@ import io.emeraldpay.dshackle.FileResolver import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.reader.Reader -import io.emeraldpay.dshackle.upstream.BlockValidator -import io.emeraldpay.dshackle.upstream.CurrentMultistreamHolder -import io.emeraldpay.dshackle.upstream.Head -import io.emeraldpay.dshackle.upstream.HttpRpcFactory -import io.emeraldpay.dshackle.upstream.MergedHead -import io.emeraldpay.dshackle.upstream.Selector -import io.emeraldpay.dshackle.upstream.Upstream +import io.emeraldpay.dshackle.upstream.* import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinRpcHead import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinRpcUpstream import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinZMQHead @@ -50,19 +44,20 @@ import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.grpc.BlockchainType import io.emeraldpay.grpc.Chain import org.slf4j.LoggerFactory -import org.springframework.beans.factory.annotation.Autowired -import org.springframework.stereotype.Repository +import org.springframework.context.ApplicationEventPublisher +import org.springframework.stereotype.Component import java.net.URI import java.util.concurrent.atomic.AtomicInteger import java.util.function.Function import javax.annotation.PostConstruct import kotlin.math.abs -@Repository +@Component open class ConfiguredUpstreams( - @Autowired private val currentUpstreams: CurrentMultistreamHolder, - @Autowired private val fileResolver: FileResolver, - @Autowired private val config: UpstreamsConfig + private val fileResolver: FileResolver, + private val config: UpstreamsConfig, + private val callTargets: CallTargetsHolder, + private val eventPublisher: ApplicationEventPublisher ) { private val log = LoggerFactory.getLogger(ConfiguredUpstreams::class.java) @@ -115,7 +110,8 @@ open class ConfiguredUpstreams( } } upstream?.let { - currentUpstreams.update(UpstreamChange(chain, upstream, UpstreamChange.ChangeType.ADDED)) + val event = UpstreamChangeEvent(chain, upstream, UpstreamChangeEvent.ChangeType.ADDED) + eventPublisher.publishEvent(event) } } } @@ -150,7 +146,7 @@ open class ConfiguredUpstreams( fun buildMethods(config: UpstreamsConfig.Upstream<*>, chain: Chain): CallMethods { return if (config.methods != null) { ManagedCallMethods( - currentUpstreams.getDefaultMethods(chain), + callTargets.getDefaultMethods(chain), config.methods!!.enabled.map { it.name }.toSet(), config.methods!!.disabled.map { it.name }.toSet() ).also { @@ -164,7 +160,7 @@ open class ConfiguredUpstreams( } } } else { - currentUpstreams.getDefaultMethods(chain) + callTargets.getDefaultMethods(chain) } } @@ -315,7 +311,7 @@ open class ConfiguredUpstreams( .doOnNext { log.info("Chain ${it.chain} ${it.type} through gRPC at ${endpoint.host}:${endpoint.port}. With caps: ${it.upstream.getCapabilities()}") } - .subscribe(currentUpstreams::update) + .subscribe(eventPublisher::publishEvent) } private fun buildHttpFactory(conn: UpstreamsConfig.RpcConnection, urls: ArrayList? = null): HttpRpcFactory? { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/startup/UpstreamChange.kt b/src/main/kotlin/io/emeraldpay/dshackle/startup/UpstreamChangeEvent.kt similarity index 98% rename from src/main/kotlin/io/emeraldpay/dshackle/startup/UpstreamChange.kt rename to src/main/kotlin/io/emeraldpay/dshackle/startup/UpstreamChangeEvent.kt index 213e55a1..34bcc8bf 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/startup/UpstreamChange.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/startup/UpstreamChangeEvent.kt @@ -24,7 +24,7 @@ import io.emeraldpay.grpc.Chain /** * An update event to the list of currently available upstreams. */ -class UpstreamChange( +class UpstreamChangeEvent( /** * Target blockchain */ diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/CallTargetsHolder.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/CallTargetsHolder.kt new file mode 100644 index 00000000..2ad03f39 --- /dev/null +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/CallTargetsHolder.kt @@ -0,0 +1,29 @@ +package io.emeraldpay.dshackle.upstream + +import io.emeraldpay.dshackle.upstream.calls.CallMethods +import io.emeraldpay.dshackle.upstream.calls.DefaultBitcoinMethods +import io.emeraldpay.dshackle.upstream.calls.DefaultEthereumMethods +import io.emeraldpay.grpc.BlockchainType +import io.emeraldpay.grpc.Chain +import org.springframework.stereotype.Component +import java.util.HashMap + +@Component +class CallTargetsHolder { + private val callTargets = HashMap() + + fun getDefaultMethods(chain: Chain): CallMethods { + return callTargets[chain] ?: return setupDefaultMethods(chain) + } + + private fun setupDefaultMethods(chain: Chain): CallMethods { + val created = when (BlockchainType.from(chain)) { + BlockchainType.EVM_POW -> DefaultEthereumMethods(chain) + BlockchainType.BITCOIN -> DefaultBitcoinMethods() + BlockchainType.EVM_POS -> DefaultEthereumMethods(chain) + else -> throw IllegalStateException("Unsupported chain: $chain") + } + callTargets[chain] = created + return created + } +} diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt index 2de935f6..0402c727 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt @@ -16,15 +16,6 @@ */ package io.emeraldpay.dshackle.upstream -import io.emeraldpay.dshackle.cache.CachesEnabled -import io.emeraldpay.dshackle.startup.UpstreamChange -import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinUpstream -import io.emeraldpay.dshackle.upstream.calls.CallMethods -import io.emeraldpay.dshackle.upstream.calls.DefaultBitcoinMethods -import io.emeraldpay.dshackle.upstream.calls.DefaultEthereumMethods -import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosUpstream -import io.emeraldpay.dshackle.upstream.ethereum.EthereumUpstream -import io.emeraldpay.grpc.BlockchainType import io.emeraldpay.grpc.Chain import org.slf4j.LoggerFactory import org.springframework.stereotype.Component @@ -49,61 +40,8 @@ open class CurrentMultistreamHolder( private val chainsBus = Sinks.many() .multicast() .directBestEffort() - private val callTargets = HashMap() private val updateLock = ReentrantLock() - fun update(change: UpstreamChange) { - updateLock.withLock { - log.debug("Upstream update: ${change.type} ${change.chain} via ${change.upstream.getId()}") - val chain = change.chain - try { - when (BlockchainType.from(chain)) { - BlockchainType.EVM_POW -> { - val up = change.upstream.cast(EthereumUpstream::class.java) - val current = chainMapping.getValue(chain) - processUpdate(change, up, current) - } - - BlockchainType.EVM_POS -> { - val up = change.upstream.cast(EthereumPosUpstream::class.java) - val current = chainMapping.getValue(chain) - processUpdate(change, up, current) - } - - BlockchainType.BITCOIN -> { - val up = change.upstream.cast(BitcoinUpstream::class.java) - val current = chainMapping.getValue(chain) - processUpdate(change, up, current) - } - - else -> { - log.error("Update for unsupported chain: $chain") - } - } - } catch (e: Throwable) { - log.error("Failed to update upstream", e) - } - } - } - - fun processUpdate(change: UpstreamChange, up: Upstream, current: Multistream) { - val chain = change.chain - if (change.type == UpstreamChange.ChangeType.REMOVED) { - current.removeUpstream(up.getId()) - log.info("Upstream ${change.upstream.getId()} with chain $chain has been removed") - } else { - if (up is CachesEnabled) { - up.setCaches(current.caches) - } - current.addUpstream(up) - - if (!callTargets.containsKey(chain)) { - setupDefaultMethods(chain) - } - log.info("Upstream ${change.upstream.getId()} with chain $chain has been added") - } - } - override fun getUpstream(chain: Chain): Multistream? { return chainMapping[chain] } @@ -122,23 +60,8 @@ open class CurrentMultistreamHolder( ) } - override fun getDefaultMethods(chain: Chain): CallMethods { - return callTargets[chain] ?: return setupDefaultMethods(chain) - } - - fun setupDefaultMethods(chain: Chain): CallMethods { - val created = when (BlockchainType.from(chain)) { - BlockchainType.EVM_POW -> DefaultEthereumMethods(chain) - BlockchainType.BITCOIN -> DefaultBitcoinMethods() - BlockchainType.EVM_POS -> DefaultEthereumMethods(chain) - else -> throw IllegalStateException("Unsupported chain: $chain") - } - callTargets[chain] = created - return created - } - override fun isAvailable(chain: Chain): Boolean { - return chainMapping.containsKey(chain) && callTargets.containsKey(chain) + return chainMapping.getValue(chain).isAvailable() } @PreDestroy diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt index c04dd9e0..80e19481 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt @@ -17,8 +17,10 @@ package io.emeraldpay.dshackle.upstream import io.emeraldpay.dshackle.cache.Caches +import io.emeraldpay.dshackle.cache.CachesEnabled import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.reader.Reader +import io.emeraldpay.dshackle.startup.UpstreamChangeEvent import io.emeraldpay.dshackle.upstream.calls.AggregatedCallMethods import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest @@ -30,6 +32,7 @@ import org.apache.commons.collections4.Factory import org.apache.commons.collections4.FunctorException import org.slf4j.LoggerFactory import org.springframework.context.Lifecycle +import org.springframework.context.event.EventListener import reactor.core.Disposable import reactor.core.publisher.Flux import reactor.core.publisher.Mono @@ -48,7 +51,8 @@ abstract class Multistream( val chain: Chain, private val upstreams: MutableList, val caches: Caches, - val postprocessor: RequestPostprocessor + val postprocessor: RequestPostprocessor, + val callTargetsHolder: CallTargetsHolder ) : Upstream, Lifecycle { companion object { @@ -316,6 +320,23 @@ abstract class Multistream( log.info("State of ${chain.chainCode}: height=${height ?: '?'}, status=[$statuses], lag=[$lag], weak=[$weak]") } + @EventListener + fun onUpstreamChange(event: UpstreamChangeEvent) { + val chain = event.chain + if (this.chain == chain) { + if (event.type == UpstreamChangeEvent.ChangeType.REMOVED) { + removeUpstream(event.upstream.getId()) + log.info("Upstream ${event.upstream.getId()} with chain $chain has been removed") + } else { + if (event.upstream is CachesEnabled) { + event.upstream.setCaches(caches) + } + addUpstream(event.upstream) + log.info("Upstream ${event.upstream.getId()} with chain $chain has been added") + } + } + } + // -------------------------------------------------------------------------------------------------------- class UpstreamStatus(val upstream: Upstream, val status: UpstreamAvailability, val ts: Instant = Instant.now()) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/MultistreamHolder.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/MultistreamHolder.kt index 852ffdee..6e1c4ab0 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/MultistreamHolder.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/MultistreamHolder.kt @@ -16,7 +16,6 @@ */ package io.emeraldpay.dshackle.upstream -import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.grpc.Chain import reactor.core.publisher.Flux @@ -27,6 +26,5 @@ interface MultistreamHolder { fun getUpstream(chain: Chain): Multistream? fun getAvailable(): List fun observeChains(): Flux - fun getDefaultMethods(chain: Chain): CallMethods fun isAvailable(chain: Chain): Boolean } 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 d72a3fb5..14ebfda0 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinMultistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinMultistream.kt @@ -18,14 +18,7 @@ package io.emeraldpay.dshackle.upstream.bitcoin 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 -import io.emeraldpay.dshackle.upstream.RequestPostprocessor -import io.emeraldpay.dshackle.upstream.Selector -import io.emeraldpay.dshackle.upstream.Upstream +import io.emeraldpay.dshackle.upstream.* import io.emeraldpay.dshackle.upstream.calls.DefaultBitcoinMethods import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest @@ -39,8 +32,9 @@ import reactor.core.publisher.Mono open class BitcoinMultistream( chain: Chain, private val sourceUpstreams: MutableList, - caches: Caches -) : Multistream(chain, sourceUpstreams as MutableList, caches, RequestPostprocessor.Empty()), Lifecycle { + caches: Caches, + callTargetsHolder: CallTargetsHolder +) : Multistream(chain, sourceUpstreams as MutableList, caches, RequestPostprocessor.Empty(), callTargetsHolder), Lifecycle { companion object { private val log = LoggerFactory.getLogger(BitcoinMultistream::class.java) 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 259e88a6..e2fb57f5 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumMultistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumMultistream.kt @@ -20,13 +20,7 @@ import io.emeraldpay.api.proto.BlockchainOuterClass 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 -import io.emeraldpay.dshackle.upstream.Selector -import io.emeraldpay.dshackle.upstream.Upstream +import io.emeraldpay.dshackle.upstream.* import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstream import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest @@ -42,8 +36,9 @@ import reactor.core.publisher.Mono open class EthereumMultistream( chain: Chain, val upstreams: MutableList, - caches: Caches -) : Multistream(chain, upstreams as MutableList, caches, CacheRequested(caches)), EthereumLikeMultistream { + caches: Caches, + callTargetsHolder: CallTargetsHolder +) : Multistream(chain, upstreams as MutableList, caches, CacheRequested(caches), callTargetsHolder), EthereumLikeMultistream { companion object { private val log = LoggerFactory.getLogger(EthereumMultistream::class.java) 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 cd872db7..3e93a81a 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,13 +20,7 @@ import io.emeraldpay.api.proto.BlockchainOuterClass 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 -import io.emeraldpay.dshackle.upstream.Selector -import io.emeraldpay.dshackle.upstream.Upstream +import io.emeraldpay.dshackle.upstream.* import io.emeraldpay.dshackle.upstream.forkchoice.PriorityForkChoice import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstream import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest @@ -42,8 +36,9 @@ import reactor.core.publisher.Mono open class EthereumPosMultiStream( chain: Chain, val upstreams: MutableList, - caches: Caches -) : Multistream(chain, upstreams as MutableList, caches, CacheRequested(caches)), EthereumLikeMultistream { + caches: Caches, + callTargetsHolder: CallTargetsHolder +) : Multistream(chain, upstreams as MutableList, caches, CacheRequested(caches), callTargetsHolder), EthereumLikeMultistream { companion object { private val log = LoggerFactory.getLogger(EthereumPosMultiStream::class.java) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreams.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreams.kt index c92fd624..f74cb62c 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreams.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreams.kt @@ -22,7 +22,7 @@ import io.emeraldpay.dshackle.Defaults import io.emeraldpay.dshackle.FileResolver import io.emeraldpay.dshackle.config.AuthConfig import io.emeraldpay.dshackle.config.UpstreamsConfig -import io.emeraldpay.dshackle.startup.UpstreamChange +import io.emeraldpay.dshackle.startup.UpstreamChangeEvent import io.emeraldpay.dshackle.upstream.DefaultUpstream import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcGrpcClient @@ -69,7 +69,7 @@ class GrpcUpstreams( private val known = HashMap() private val lock = ReentrantLock() - fun start(): Flux { + fun start(): Flux { val channel: ManagedChannelBuilder<*> = if (auth != null && StringUtils.isNotEmpty(auth.ca)) { NettyChannelBuilder.forAddress(host, port) // some messages are very large. many of them in megabytes, some even in gigabytes (ex. ETH Traces) @@ -122,7 +122,7 @@ class GrpcUpstreams( return updates } - fun processDescription(value: BlockchainOuterClass.DescribeResponse): Flux { + fun processDescription(value: BlockchainOuterClass.DescribeResponse): Flux { val current = value.chainsList.filter { Chain.byId(it.chain.number) != Chain.UNSPECIFIED }.mapNotNull { chainDetails -> @@ -138,14 +138,14 @@ class GrpcUpstreams( } val added = current.filter { - it.type == UpstreamChange.ChangeType.ADDED + it.type == UpstreamChangeEvent.ChangeType.ADDED } val removed = known.filterNot { kv -> val stillCurrent = current.any { c -> c.chain == kv.key } stillCurrent }.map { - UpstreamChange(it.key, known.remove(it.key)!!, UpstreamChange.ChangeType.REMOVED) + UpstreamChangeEvent(it.key, known.remove(it.key)!!, UpstreamChangeEvent.ChangeType.REMOVED) } return Flux.fromIterable(removed + added) } @@ -172,7 +172,7 @@ class GrpcUpstreams( return sslContext.build() } - fun getOrCreate(chain: Chain): UpstreamChange { + fun getOrCreate(chain: Chain): UpstreamChangeEvent { val metricsTags = listOf( Tag.of("upstream", id), Tag.of("chain", chain.chainCode) @@ -202,7 +202,7 @@ class GrpcUpstreams( } } - fun getOrCreateEthereum(chain: Chain, metrics: RpcMetrics): UpstreamChange { + fun getOrCreateEthereum(chain: Chain, metrics: RpcMetrics): UpstreamChangeEvent { lock.withLock { val current = known[chain] return if (current == null) { @@ -211,14 +211,14 @@ class GrpcUpstreams( created.timeout = this.timeout known[chain] = created created.start() - UpstreamChange(chain, created, UpstreamChange.ChangeType.ADDED) + UpstreamChangeEvent(chain, created, UpstreamChangeEvent.ChangeType.ADDED) } else { - UpstreamChange(chain, current, UpstreamChange.ChangeType.REVALIDATED) + UpstreamChangeEvent(chain, current, UpstreamChangeEvent.ChangeType.REVALIDATED) } } } - fun getOrCreateEthereumPos(chain: Chain, metrics: RpcMetrics): UpstreamChange { + fun getOrCreateEthereumPos(chain: Chain, metrics: RpcMetrics): UpstreamChangeEvent { lock.withLock { val current = known[chain] return if (current == null) { @@ -227,14 +227,14 @@ class GrpcUpstreams( created.timeout = this.timeout known[chain] = created created.start() - UpstreamChange(chain, created, UpstreamChange.ChangeType.ADDED) + UpstreamChangeEvent(chain, created, UpstreamChangeEvent.ChangeType.ADDED) } else { - UpstreamChange(chain, current, UpstreamChange.ChangeType.REVALIDATED) + UpstreamChangeEvent(chain, current, UpstreamChangeEvent.ChangeType.REVALIDATED) } } } - fun getOrCreateBitcoin(chain: Chain, metrics: RpcMetrics): UpstreamChange { + fun getOrCreateBitcoin(chain: Chain, metrics: RpcMetrics): UpstreamChangeEvent { lock.withLock { val current = known[chain] return if (current == null) { @@ -243,9 +243,9 @@ class GrpcUpstreams( created.timeout = this.timeout known[chain] = created created.start() - UpstreamChange(chain, created, UpstreamChange.ChangeType.ADDED) + UpstreamChangeEvent(chain, created, UpstreamChangeEvent.ChangeType.ADDED) } else { - UpstreamChange(chain, current, UpstreamChange.ChangeType.REVALIDATED) + UpstreamChangeEvent(chain, current, UpstreamChangeEvent.ChangeType.REVALIDATED) } } } diff --git a/src/test/groovy/io/emeraldpay/dshackle/startup/ConfiguredUpstreamsSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/startup/ConfiguredUpstreamsSpec.groovy index 4e6f6eea..f799729d 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/startup/ConfiguredUpstreamsSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/startup/ConfiguredUpstreamsSpec.groovy @@ -4,21 +4,24 @@ import io.emeraldpay.dshackle.FileResolver import io.emeraldpay.dshackle.cache.CachesFactory import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.quorum.NonEmptyQuorum +import io.emeraldpay.dshackle.upstream.CallTargetsHolder import io.emeraldpay.dshackle.upstream.CurrentMultistreamHolder import io.emeraldpay.dshackle.upstream.calls.DefaultEthereumMethods import io.emeraldpay.dshackle.upstream.calls.ManagedCallMethods import io.emeraldpay.grpc.Chain +import org.springframework.context.ApplicationEventPublisher import spock.lang.Specification class ConfiguredUpstreamsSpec extends Specification { def "Applied quorum to extra methods"() { setup: - def currentUpstreams = Mock(CurrentMultistreamHolder) { - _ * getDefaultMethods(Chain.ETHEREUM) >> new DefaultEthereumMethods(Chain.ETHEREUM) - } + def callTargetsHolder = new CallTargetsHolder() def configurer = new ConfiguredUpstreams( - currentUpstreams, Stub(FileResolver), Stub(UpstreamsConfig) + Stub(FileResolver), + Stub(UpstreamsConfig), + callTargetsHolder, + Mock(ApplicationEventPublisher) ) def methods = new UpstreamsConfig.Methods( [ @@ -38,11 +41,12 @@ class ConfiguredUpstreamsSpec extends Specification { def "Got static response from extra methods"() { setup: - def currentUpstreams = Mock(CurrentMultistreamHolder) { - _ * getDefaultMethods(Chain.ETHEREUM) >> new DefaultEthereumMethods(Chain.ETHEREUM) - } + def callTargetsHolder = new CallTargetsHolder() def configurer = new ConfiguredUpstreams( - currentUpstreams, Stub(FileResolver), Stub(UpstreamsConfig) + Stub(FileResolver), + Stub(UpstreamsConfig), + callTargetsHolder, + Mock(ApplicationEventPublisher) ) def methods = new UpstreamsConfig.Methods( [ @@ -61,7 +65,12 @@ class ConfiguredUpstreamsSpec extends Specification { def "Calculate node-id"() { setup: - def configurer = new ConfiguredUpstreams(Stub(CurrentMultistreamHolder), Stub(FileResolver), Stub(UpstreamsConfig) + def callTargetsHolder = new CallTargetsHolder() + def configurer = new ConfiguredUpstreams( + Stub(FileResolver), + Stub(UpstreamsConfig), + callTargetsHolder, + Mock(ApplicationEventPublisher) ) expect: configurer.getHash(node, src) == expected @@ -75,7 +84,12 @@ class ConfiguredUpstreamsSpec extends Specification { def "Calculate node-id conflicting results"() { setup: - def configurer = new ConfiguredUpstreams(Stub(CurrentMultistreamHolder), Stub(FileResolver), Stub(UpstreamsConfig) + def callTargetsHolder = new CallTargetsHolder() + def configurer = new ConfiguredUpstreams( + Stub(FileResolver), + Stub(UpstreamsConfig), + callTargetsHolder, + Mock(ApplicationEventPublisher) ) when: def h1 = configurer.getHash(null, "hohoho") diff --git a/src/test/groovy/io/emeraldpay/dshackle/test/MultistreamHolderMock.groovy b/src/test/groovy/io/emeraldpay/dshackle/test/MultistreamHolderMock.groovy index 660b7611..9e34c546 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/test/MultistreamHolderMock.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/test/MultistreamHolderMock.groovy @@ -50,7 +50,7 @@ class MultistreamHolderMock implements MultistreamHolder { if (up instanceof EthereumPosMultiStream) { upstreams[chain] = up } else if (up instanceof EthereumPosRpcUpstream) { - upstreams[chain] = new EthereumPosMultiStream(chain, [up as EthereumPosRpcUpstream], Caches.default()) + upstreams[chain] = new EthereumPosMultiStream(chain, [up as EthereumPosRpcUpstream], Caches.default(), TestingCommons.callTargetsHolder) } else { throw new IllegalArgumentException("Unsupported upstream type ${up.class}") } @@ -59,7 +59,7 @@ class MultistreamHolderMock implements MultistreamHolder { if (up instanceof BitcoinMultistream) { upstreams[chain] = up } else if (up instanceof BitcoinRpcUpstream) { - upstreams[chain] = new BitcoinMultistream(chain, [up as BitcoinRpcUpstream], Caches.default()) + upstreams[chain] = new BitcoinMultistream(chain, [up as BitcoinRpcUpstream], Caches.default(), TestingCommons.callTargetsHolder) } else { throw new IllegalArgumentException("Unsupported upstream type ${up.class}") } @@ -86,15 +86,6 @@ class MultistreamHolderMock implements MultistreamHolder { return Flux.fromIterable(getAvailable()) } - @Override - DefaultEthereumMethods getDefaultMethods(@NotNull Chain chain) { - if (target[chain] == null) { - DefaultEthereumMethods targets = new DefaultEthereumMethods(chain) - target[chain] = targets - } - return target[chain] - } - @Override boolean isAvailable(@NotNull Chain chain) { return upstreams.containsKey(chain) @@ -107,7 +98,7 @@ class MultistreamHolderMock implements MultistreamHolder { Head customHead = null EthereumMultistreamMock(@NotNull Chain chain, @NotNull List upstreams, @NotNull Caches caches) { - super(chain, upstreams, caches) + super(chain, upstreams, caches, TestingCommons.callTargetsHolder) } EthereumMultistreamMock(@NotNull Chain chain, @NotNull List upstreams) { diff --git a/src/test/groovy/io/emeraldpay/dshackle/test/TestingCommons.groovy b/src/test/groovy/io/emeraldpay/dshackle/test/TestingCommons.groovy index 6ad08167..f2c640d7 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/test/TestingCommons.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/test/TestingCommons.groovy @@ -24,8 +24,10 @@ import io.emeraldpay.dshackle.data.BlockContainer import io.emeraldpay.dshackle.data.BlockId import io.emeraldpay.dshackle.reader.EmptyReader import io.emeraldpay.dshackle.reader.Reader +import io.emeraldpay.dshackle.upstream.CallTargetsHolder import io.emeraldpay.dshackle.upstream.Multistream import io.emeraldpay.dshackle.upstream.calls.DirectCallMethods +import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosUpstream import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest @@ -78,7 +80,7 @@ class TestingCommons { } static Multistream multistream(EthereumPosRpcUpstreamMock up) { - return new EthereumPosMultiStream(Chain.ETHEREUM, [up], Caches.default()).tap { + return new EthereumPosMultiStream(Chain.ETHEREUM, [up], Caches.default(), callTargetsHolder).tap { start() } } @@ -94,12 +96,16 @@ class TestingCommons { static List defaultMultistreams() { return [ multistreamWithoutUpstreams(Chain.ETHEREUM), - multistreamWithoutUpstreams(Chain.ETHEREUM_CLASSIC) + multistreamClassicWithoutUpstreams(Chain.ETHEREUM_CLASSIC) ] } static Multistream multistreamWithoutUpstreams(Chain chain) { - return new EthereumPosMultiStream(chain, [], emptyCaches().getCaches(chain)) + return new EthereumPosMultiStream(chain, [], emptyCaches().getCaches(chain), callTargetsHolder) + } + + static Multistream multistreamClassicWithoutUpstreams(Chain chain) { + return new EthereumMultistream(chain, [], emptyCaches().getCaches(chain), callTargetsHolder) } static FileResolver fileResolver() { @@ -138,4 +144,6 @@ class TestingCommons { } static MeterRegistry meterRegistry = new LoggingMeterRegistry() + + static CallTargetsHolder callTargetsHolder = new CallTargetsHolder() } diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolderSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolderSpec.groovy index 0ec23584..4e488e9d 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolderSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolderSpec.groovy @@ -15,7 +15,7 @@ */ package io.emeraldpay.dshackle.upstream -import io.emeraldpay.dshackle.startup.UpstreamChange +import io.emeraldpay.dshackle.startup.UpstreamChangeEvent import io.emeraldpay.dshackle.test.EthereumPosRpcUpstreamMock import io.emeraldpay.dshackle.test.EthereumRpcUpstreamMock import io.emeraldpay.dshackle.test.TestingCommons @@ -29,7 +29,7 @@ class CurrentMultistreamHolderSpec extends Specification { def current = new CurrentMultistreamHolder(TestingCommons.defaultMultistreams()) def up = new EthereumPosRpcUpstreamMock("test", Chain.ETHEREUM, TestingCommons.api()) when: - current.update(new UpstreamChange(Chain.ETHEREUM, up, UpstreamChange.ChangeType.ADDED)) + current.getUpstream(Chain.ETHEREUM).onUpstreamChange(new UpstreamChangeEvent(Chain.ETHEREUM, up, UpstreamChangeEvent.ChangeType.ADDED)) then: current.getAvailable() == [Chain.ETHEREUM] current.getUpstream(Chain.ETHEREUM).getAll()[0] == up @@ -42,9 +42,10 @@ class CurrentMultistreamHolderSpec extends Specification { def up2 = new EthereumRpcUpstreamMock("test2", Chain.ETHEREUM_CLASSIC, TestingCommons.api()) def up3 = new EthereumPosRpcUpstreamMock("test3", Chain.ETHEREUM, TestingCommons.api()) when: - current.update(new UpstreamChange(Chain.ETHEREUM, up1, UpstreamChange.ChangeType.ADDED)) - current.update(new UpstreamChange(Chain.ETHEREUM_CLASSIC, up2, UpstreamChange.ChangeType.ADDED)) - current.update(new UpstreamChange(Chain.ETHEREUM, up3, UpstreamChange.ChangeType.ADDED)) + current.getUpstream(Chain.ETHEREUM).onUpstreamChange(new UpstreamChangeEvent(Chain.ETHEREUM, up1, UpstreamChangeEvent.ChangeType.ADDED)) + current.getUpstream(Chain.ETHEREUM_CLASSIC).onUpstreamChange(new UpstreamChangeEvent(Chain.ETHEREUM_CLASSIC, up2, UpstreamChangeEvent.ChangeType.ADDED)) + current.getUpstream(Chain.ETHEREUM).onUpstreamChange(new UpstreamChangeEvent(Chain.ETHEREUM, up3, UpstreamChangeEvent.ChangeType.ADDED)) + current.getUpstream(Chain.ETHEREUM_CLASSIC).onUpstreamChange(new UpstreamChangeEvent(Chain.ETHEREUM, up3, UpstreamChangeEvent.ChangeType.ADDED)) then: current.getAvailable().toSet() == [Chain.ETHEREUM, Chain.ETHEREUM_CLASSIC].toSet() current.getUpstream(Chain.ETHEREUM).getAll().toSet() == [up1, up3].toSet() @@ -59,10 +60,10 @@ class CurrentMultistreamHolderSpec extends Specification { def up3 = new EthereumPosRpcUpstreamMock("test3", Chain.ETHEREUM, TestingCommons.api()) def up1_del = new EthereumPosRpcUpstreamMock("test1", Chain.ETHEREUM, TestingCommons.api()) when: - current.update(new UpstreamChange(Chain.ETHEREUM, up1, UpstreamChange.ChangeType.ADDED)) - current.update(new UpstreamChange(Chain.ETHEREUM_CLASSIC, up2, UpstreamChange.ChangeType.ADDED)) - current.update(new UpstreamChange(Chain.ETHEREUM, up3, UpstreamChange.ChangeType.ADDED)) - current.update(new UpstreamChange(Chain.ETHEREUM, up1_del, UpstreamChange.ChangeType.REMOVED)) + current.getUpstream(Chain.ETHEREUM).onUpstreamChange(new UpstreamChangeEvent(Chain.ETHEREUM, up1, UpstreamChangeEvent.ChangeType.ADDED)) + current.getUpstream(Chain.ETHEREUM_CLASSIC).onUpstreamChange(new UpstreamChangeEvent(Chain.ETHEREUM_CLASSIC, up2, UpstreamChangeEvent.ChangeType.ADDED)) + current.getUpstream(Chain.ETHEREUM).onUpstreamChange(new UpstreamChangeEvent(Chain.ETHEREUM, up3, UpstreamChangeEvent.ChangeType.ADDED)) + current.getUpstream(Chain.ETHEREUM).onUpstreamChange(new UpstreamChangeEvent(Chain.ETHEREUM, up1_del, UpstreamChangeEvent.ChangeType.REMOVED)) then: current.getAvailable().toSet() == [Chain.ETHEREUM, Chain.ETHEREUM_CLASSIC].toSet() current.getUpstream(Chain.ETHEREUM).getAll().toSet() == [up3].toSet() @@ -80,7 +81,7 @@ class CurrentMultistreamHolderSpec extends Specification { !act when: - current.update(new UpstreamChange(Chain.ETHEREUM, up1, UpstreamChange.ChangeType.ADDED)) + current.getUpstream(Chain.ETHEREUM).onUpstreamChange(new UpstreamChangeEvent(Chain.ETHEREUM, up1, UpstreamChangeEvent.ChangeType.ADDED)) act = current.isAvailable(Chain.ETHEREUM) then: diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/MultistreamSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/MultistreamSpec.groovy index 8a28f4b3..b432bcb8 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/MultistreamSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/MultistreamSpec.groovy @@ -50,7 +50,7 @@ class MultistreamSpec extends Specification { setup: def up1 = new EthereumPosRpcUpstreamMock("test1", Chain.ETHEREUM, TestingCommons.api(), new DirectCallMethods(["eth_test1", "eth_test2"])) def up2 = new EthereumPosRpcUpstreamMock("test1", Chain.ETHEREUM, TestingCommons.api(), new DirectCallMethods(["eth_test2", "eth_test3"])) - def aggr = new EthereumPosMultiStream(Chain.ETHEREUM, [up1, up2], Caches.default()) + def aggr = new EthereumPosMultiStream(Chain.ETHEREUM, [up1, up2], Caches.default(), TestingCommons.callTargetsHolder) when: aggr.onUpstreamsUpdated() def act = aggr.getMethods() @@ -206,7 +206,7 @@ class MultistreamSpec extends Specification { def up1 = TestingCommons.upstream("test-1", "internal") def up2 = TestingCommons.upstream("test-2", "external") def up3 = TestingCommons.upstream("test-3", "external") - def multistream = new EthereumPosMultiStream(Chain.ETHEREUM, [up1, up2, up3], Caches.default()) + def multistream = new EthereumPosMultiStream(Chain.ETHEREUM, [up1, up2, up3], Caches.default(), TestingCommons.callTargetsHolder) expect: multistream.getHead(new Selector.LabelMatcher("provider", ["internal"])).is(up1.ethereumHeadMock) @@ -345,7 +345,7 @@ class MultistreamSpec extends Specification { class TestMultistream extends Multistream { TestMultistream(List upstreams, @NotNull RequestPostprocessor postprocessor) { - super(Chain.ETHEREUM, upstreams, Caches.default(), postprocessor) + super(Chain.ETHEREUM, upstreams, Caches.default(), postprocessor, TestingCommons.callTargetsHolder) } @Override @@ -386,7 +386,7 @@ class MultistreamSpec extends Specification { class TestEthereumPosMultistream extends EthereumPosMultiStream { TestEthereumPosMultistream(@NotNull Chain chain, @NotNull List upstreams, @NotNull Caches caches) { - super(chain, upstreams, caches) + super(chain, upstreams, caches, TestingCommons.callTargetsHolder) } @Override From eaa68d606bca0f0f3ea9734e647041afa65d0214 Mon Sep 17 00:00:00 2001 From: a10zn8 Date: Wed, 16 Nov 2022 15:26:02 +0400 Subject: [PATCH 3/7] fix after tests --- .../emeraldpay/dshackle/config/context/MultistreamsConfig.kt | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/config/context/MultistreamsConfig.kt b/src/main/kotlin/io/emeraldpay/dshackle/config/context/MultistreamsConfig.kt index cfa80f3a..d8204b87 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/config/context/MultistreamsConfig.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/config/context/MultistreamsConfig.kt @@ -12,13 +12,14 @@ import org.springframework.context.annotation.Bean import org.springframework.context.annotation.Configuration @Configuration -class MultistreamsConfig { +open class MultistreamsConfig { @Bean - fun allMultistreams( + open fun allMultistreams( cachesFactory: CachesFactory, callTargetsHolder: CallTargetsHolder ): List { return Chain.values() + .filterNot { it == Chain.UNSPECIFIED } .mapNotNull { chain -> when (BlockchainType.from(chain)) { BlockchainType.EVM_POS -> EthereumPosMultiStream( From d997d148c4dbb8bba04fcac706c02ff65cc66afe Mon Sep 17 00:00:00 2001 From: a10zn8 Date: Wed, 16 Nov 2022 17:00:07 +0400 Subject: [PATCH 4/7] fixing startup --- .../kotlin/io/emeraldpay/dshackle/Starter.kt | 2 - .../config/context/MultistreamsConfig.kt | 69 ++++++++++++++----- .../io/emeraldpay/dshackle/rpc/NativeCall.kt | 34 ++++----- .../dshackle/rpc/TrackBitcoinAddress.kt | 17 +++-- .../dshackle/startup/ConfiguredUpstreams.kt | 9 +-- .../upstream/CurrentMultistreamHolder.kt | 11 --- .../emeraldpay/dshackle/upstream/Lifecycle.kt | 7 ++ .../dshackle/upstream/MergedHead.kt | 5 +- .../dshackle/upstream/Multistream.kt | 20 ++++-- .../dshackle/upstream/MultistreamHolder.kt | 2 - .../upstream/bitcoin/BitcoinMultistream.kt | 9 ++- .../upstream/bitcoin/BitcoinReader.kt | 4 +- .../upstream/bitcoin/BitcoinRpcHead.kt | 2 +- .../upstream/bitcoin/BitcoinRpcUpstream.kt | 6 +- .../upstream/bitcoin/BitcoinZMQHead.kt | 4 +- .../upstream/bitcoin/CachingMempoolData.kt | 2 +- .../dshackle/upstream/bitcoin/ZMQServer.kt | 2 +- .../upstream/ethereum/EthereumMultistream.kt | 9 ++- .../upstream/ethereum/EthereumReader.kt | 4 +- .../upstream/ethereum/EthereumRpcHead.kt | 2 +- .../upstream/ethereum/EthereumRpcUpstream.kt | 2 +- .../upstream/ethereum/EthereumWsHead.kt | 2 +- .../ethereum/connectors/EthereumConnector.kt | 2 +- .../connectors/EthereumRpcConnector.kt | 4 +- .../connectors/EthereumWsConnector.kt | 2 +- .../ethereum_pos/EthereumPosMultiStream.kt | 9 ++- .../ethereum_pos/EthereumPosRpcUpstream.kt | 4 +- .../upstream/grpc/BitcoinGrpcUpstream.kt | 2 +- .../upstream/grpc/EthereumGrpcUpstream.kt | 2 +- .../upstream/grpc/EthereumPosGrpcUpstream.kt | 2 +- .../dshackle/upstream/grpc/GrpcHead.kt | 4 +- src/main/resources/log4j2.xml | 2 +- .../dshackle/rpc/NativeCallSpec.groovy | 4 +- .../test/MultistreamHolderMock.groovy | 11 +-- .../dshackle/test/TestingCommons.groovy | 7 +- .../dshackle/upstream/MergedHeadSpec.groovy | 2 +- .../dshackle/upstream/MultistreamSpec.groovy | 8 +-- 37 files changed, 154 insertions(+), 135 deletions(-) create mode 100644 src/main/kotlin/io/emeraldpay/dshackle/upstream/Lifecycle.kt diff --git a/src/main/kotlin/io/emeraldpay/dshackle/Starter.kt b/src/main/kotlin/io/emeraldpay/dshackle/Starter.kt index 807b399b..525eda35 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/Starter.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/Starter.kt @@ -20,12 +20,10 @@ import org.slf4j.LoggerFactory import org.springframework.boot.ResourceBanner import org.springframework.boot.SpringApplication import org.springframework.boot.autoconfigure.SpringBootApplication -import org.springframework.context.annotation.Import import org.springframework.core.io.ClassPathResource import org.springframework.core.io.support.ResourcePropertySource @SpringBootApplication(scanBasePackages = ["io.emeraldpay.dshackle"]) -@Import(Config::class) open class Starter private val log = LoggerFactory.getLogger(Starter::class.java) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/config/context/MultistreamsConfig.kt b/src/main/kotlin/io/emeraldpay/dshackle/config/context/MultistreamsConfig.kt index d8204b87..3329f93e 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/config/context/MultistreamsConfig.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/config/context/MultistreamsConfig.kt @@ -8,11 +8,12 @@ import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream import io.emeraldpay.grpc.BlockchainType import io.emeraldpay.grpc.Chain +import org.springframework.beans.factory.config.ConfigurableListableBeanFactory import org.springframework.context.annotation.Bean import org.springframework.context.annotation.Configuration @Configuration -open class MultistreamsConfig { +open class MultistreamsConfig(val beanFactory: ConfigurableListableBeanFactory) { @Bean open fun allMultistreams( cachesFactory: CachesFactory, @@ -22,26 +23,56 @@ open class MultistreamsConfig { .filterNot { it == Chain.UNSPECIFIED } .mapNotNull { chain -> when (BlockchainType.from(chain)) { - BlockchainType.EVM_POS -> EthereumPosMultiStream( - chain, - ArrayList(), - cachesFactory.getCaches(chain), - callTargetsHolder - ) - BlockchainType.EVM_POW -> EthereumMultistream( - chain, - ArrayList(), - cachesFactory.getCaches(chain), - callTargetsHolder - ) - BlockchainType.BITCOIN -> BitcoinMultistream( - chain, - ArrayList(), - cachesFactory.getCaches(chain), - callTargetsHolder - ) + BlockchainType.EVM_POS -> ethereumPosMultistream(chain, cachesFactory) + BlockchainType.EVM_POW -> ethereumMultistream(chain, cachesFactory) + BlockchainType.BITCOIN -> bitcoinMultistream(chain, cachesFactory) else -> null } } } + + private fun ethereumMultistream( + chain: Chain, + cachesFactory: CachesFactory + ): EthereumMultistream { + val name = "multi-ethereum-$chain" + + return EthereumMultistream( + chain, + ArrayList(), + cachesFactory.getCaches(chain) + ).also { register(it, name) } + } + + open fun ethereumPosMultistream( + chain: Chain, + cachesFactory: CachesFactory + ): EthereumPosMultiStream { + val name = "multi-ethereum-pos-$chain" + + return EthereumPosMultiStream( + chain, + ArrayList(), + cachesFactory.getCaches(chain) + ).also { register(it, name) } + } + + open fun bitcoinMultistream( + chain: Chain, + cachesFactory: CachesFactory + ): BitcoinMultistream { + val name = "multi-bitcoin-$chain" + + return BitcoinMultistream( + chain, + ArrayList(), + cachesFactory.getCaches(chain) + ).also { register(it, name) } + } + + private fun register(bean: Any, name: String) { + beanFactory.initializeBean(bean, name) + beanFactory.autowireBean(bean) + this.beanFactory.registerSingleton(name, bean) + } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt index 956b187d..93926399 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt @@ -26,11 +26,10 @@ import io.emeraldpay.dshackle.quorum.NotLaggingQuorum import io.emeraldpay.dshackle.quorum.QuorumReaderFactory import io.emeraldpay.dshackle.quorum.QuorumRpcReader import io.emeraldpay.dshackle.startup.ConfiguredUpstreams -import io.emeraldpay.dshackle.upstream.ApiSource -import io.emeraldpay.dshackle.upstream.Multistream -import io.emeraldpay.dshackle.upstream.MultistreamHolder -import io.emeraldpay.dshackle.upstream.Selector +import io.emeraldpay.dshackle.startup.UpstreamChangeEvent +import io.emeraldpay.dshackle.upstream.* import io.emeraldpay.dshackle.upstream.calls.EthereumCallSelector +import io.emeraldpay.dshackle.upstream.ethereum.EthereumLikeMultistream import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcError @@ -45,7 +44,7 @@ import io.emeraldpay.grpc.Chain import io.micrometer.core.instrument.Metrics import org.apache.commons.lang3.StringUtils import org.slf4j.LoggerFactory -import org.springframework.beans.factory.annotation.Autowired +import org.springframework.context.event.EventListener import org.springframework.stereotype.Service import reactor.core.publisher.Flux import reactor.core.publisher.Mono @@ -55,9 +54,9 @@ import java.util.concurrent.atomic.AtomicInteger @Service open class NativeCall( - @Autowired private val multistreamHolder: MultistreamHolder, - @Autowired private val configuredUpstreams: ConfiguredUpstreams, - @Autowired private val signer: ResponseSigner + private val multistreamHolder: MultistreamHolder, + private val configuredUpstreams: ConfiguredUpstreams, + private val signer: ResponseSigner ) { private val log = LoggerFactory.getLogger(NativeCall::class.java) @@ -66,18 +65,19 @@ open class NativeCall( var quorumReaderFactory: QuorumReaderFactory = QuorumReaderFactory.default() private val ethereumCallSelectors = EnumMap(Chain::class.java) - init { - val casting = mapOf( + companion object { + val casting: Map> = mapOf( BlockchainType.EVM_POS to EthereumPosMultiStream::class.java, - BlockchainType.EVM_POW to EthereumMultistream::class.java, + BlockchainType.EVM_POW to EthereumMultistream::class.java ) + } - multistreamHolder.observeChains().subscribe { chain -> - casting[BlockchainType.from(chain)]?.let { cast -> - multistreamHolder.getUpstream(chain)?.let { up -> - val reader = up.cast(cast).getReader() - ethereumCallSelectors.putIfAbsent(chain, EthereumCallSelector(reader.heightByHash())) - } + @EventListener + fun onUpstreamChangeEvent(event: UpstreamChangeEvent) { + casting[BlockchainType.from(event.chain)]?.let { cast -> + multistreamHolder.getUpstream(event.chain)?.let { up -> + val reader = up.cast(cast).getReader() + ethereumCallSelectors.putIfAbsent(event.chain, EthereumCallSelector(reader.heightByHash())) } } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackBitcoinAddress.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackBitcoinAddress.kt index 0536d655..d64b493c 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackBitcoinAddress.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackBitcoinAddress.kt @@ -20,6 +20,7 @@ import io.emeraldpay.api.proto.Common import io.emeraldpay.api.proto.ReactorBlockchainGrpc import io.emeraldpay.dshackle.Defaults import io.emeraldpay.dshackle.SilentException +import io.emeraldpay.dshackle.startup.UpstreamChangeEvent import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.MultistreamHolder import io.emeraldpay.dshackle.upstream.Selector @@ -33,12 +34,12 @@ import org.bitcoinj.params.MainNetParams import org.bitcoinj.params.TestNet3Params import org.slf4j.LoggerFactory import org.springframework.beans.factory.annotation.Autowired +import org.springframework.context.event.EventListener import org.springframework.stereotype.Service import reactor.core.publisher.Flux import reactor.core.publisher.Mono import java.math.BigInteger import java.util.concurrent.ConcurrentHashMap -import javax.annotation.PostConstruct @Service class TrackBitcoinAddress( @@ -67,15 +68,13 @@ class TrackBitcoinAddress( Selector.CapabilityMatcher(Capability.BALANCE) ) - @PostConstruct - fun listenChains() { - multistreamHolder.observeChains().subscribe { chain -> - multistreamHolder.getUpstream(chain)?.let { mup -> - val available = mup.getAll().any { up -> - !up.isGrpc() && up.getCapabilities().contains(Capability.BALANCE) - } - setBalanceAvailability(chain, available) + @EventListener + fun onUpstreamChangeEvent(event: UpstreamChangeEvent) { + multistreamHolder.getUpstream(event.chain)?.let { mup -> + val available = mup.getAll().any { up -> + !up.isGrpc() && up.getCapabilities().contains(Capability.BALANCE) } + setBalanceAvailability(event.chain, available) } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt b/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt index 8fc7e036..b66a125a 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt @@ -44,12 +44,13 @@ import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.grpc.BlockchainType import io.emeraldpay.grpc.Chain import org.slf4j.LoggerFactory +import org.springframework.boot.ApplicationArguments +import org.springframework.boot.ApplicationRunner import org.springframework.context.ApplicationEventPublisher import org.springframework.stereotype.Component import java.net.URI import java.util.concurrent.atomic.AtomicInteger import java.util.function.Function -import javax.annotation.PostConstruct import kotlin.math.abs @Component @@ -58,15 +59,14 @@ open class ConfiguredUpstreams( private val config: UpstreamsConfig, private val callTargets: CallTargetsHolder, private val eventPublisher: ApplicationEventPublisher -) { +) : ApplicationRunner { private val log = LoggerFactory.getLogger(ConfiguredUpstreams::class.java) private var seq = AtomicInteger(0) private val hashes: MutableMap = HashMap() - @PostConstruct - fun start() { + override fun run(args: ApplicationArguments) { log.debug("Starting upstreams") val defaultOptions = buildDefaultOptions(config) config.upstreams.forEach { up -> @@ -111,6 +111,7 @@ open class ConfiguredUpstreams( } upstream?.let { val event = UpstreamChangeEvent(chain, upstream, UpstreamChangeEvent.ChangeType.ADDED) + log.error("first !!!! Upstream ${event.upstream.getId()} with chain $chain has been added") eventPublisher.publishEvent(event) } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt index 0402c727..b4865951 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt @@ -19,7 +19,6 @@ package io.emeraldpay.dshackle.upstream import io.emeraldpay.grpc.Chain import org.slf4j.LoggerFactory import org.springframework.stereotype.Component -import reactor.core.publisher.Flux import reactor.core.publisher.Sinks import java.util.* import java.util.concurrent.ConcurrentHashMap @@ -37,9 +36,6 @@ open class CurrentMultistreamHolder( private val chainMapping = ConcurrentHashMap().apply { multistreams.forEach { this[it.chain] = it } } - private val chainsBus = Sinks.many() - .multicast() - .directBestEffort() private val updateLock = ReentrantLock() override fun getUpstream(chain: Chain): Multistream? { @@ -53,13 +49,6 @@ open class CurrentMultistreamHolder( .toList() } - override fun observeChains(): Flux { - return Flux.concat( - Flux.fromIterable(getAvailable()), - chainsBus.asFlux() - ) - } - override fun isAvailable(chain: Chain): Boolean { return chainMapping.getValue(chain).isAvailable() } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Lifecycle.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Lifecycle.kt new file mode 100644 index 00000000..ff7eb41d --- /dev/null +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Lifecycle.kt @@ -0,0 +1,7 @@ +package io.emeraldpay.dshackle.upstream + +interface Lifecycle { + fun start() + fun stop() + fun isRunning(): Boolean +} diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/MergedHead.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/MergedHead.kt index 42e654e0..7ea7d7f8 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/MergedHead.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/MergedHead.kt @@ -21,7 +21,6 @@ import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.cache.CachesEnabled import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice import org.slf4j.LoggerFactory -import org.springframework.context.Lifecycle import reactor.core.Disposable import reactor.core.publisher.Flux @@ -43,7 +42,7 @@ class MergedHead( override fun start() { super.start() sources.forEach { head -> - if (head is Lifecycle && !head.isRunning) { + if (head is Lifecycle && !head.isRunning()) { head.start() } } @@ -56,7 +55,7 @@ class MergedHead( override fun stop() { super.stop() sources.forEach { head -> - if (head is Lifecycle && head.isRunning) { + if (head is Lifecycle && head.isRunning()) { head.stop() } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt index 80e19481..54fb2dfd 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt @@ -31,14 +31,15 @@ import io.micrometer.core.instrument.Tag import org.apache.commons.collections4.Factory import org.apache.commons.collections4.FunctorException import org.slf4j.LoggerFactory -import org.springframework.context.Lifecycle import org.springframework.context.event.EventListener +import org.springframework.core.Ordered +import org.springframework.core.annotation.Order import reactor.core.Disposable import reactor.core.publisher.Flux import reactor.core.publisher.Mono import java.time.Duration import java.time.Instant -import java.util.Locale +import java.util.* import java.util.concurrent.atomic.AtomicReference import java.util.concurrent.locks.ReentrantLock import java.util.function.Predicate @@ -51,8 +52,7 @@ abstract class Multistream( val chain: Chain, private val upstreams: MutableList, val caches: Caches, - val postprocessor: RequestPostprocessor, - val callTargetsHolder: CallTargetsHolder + val postprocessor: RequestPostprocessor ) : Upstream, Lifecycle { companion object { @@ -60,6 +60,8 @@ abstract class Multistream( private const val metrics = "upstreams" } + private var started = false + private var cacheSubscription: Disposable? = null private val reconfigLock = ReentrantLock() private var callMethods: CallMethods? = null @@ -235,6 +237,7 @@ abstract class Multistream( // print status _change_ every 15 seconds, at most; otherwise prints it on interval of 30 seconds .sample(Duration.ofSeconds(15)) .subscribe { printStatus() } + started = true } override fun stop() { @@ -248,6 +251,7 @@ abstract class Multistream( } } lagObserver?.stop() + started = false } fun onHeadUpdated(head: Head) { @@ -321,18 +325,22 @@ abstract class Multistream( } @EventListener + @Order(Ordered.HIGHEST_PRECEDENCE) fun onUpstreamChange(event: UpstreamChangeEvent) { val chain = event.chain if (this.chain == chain) { if (event.type == UpstreamChangeEvent.ChangeType.REMOVED) { removeUpstream(event.upstream.getId()) - log.info("Upstream ${event.upstream.getId()} with chain $chain has been removed") + log.error("Upstream ${event.upstream.getId()} with chain $chain has been removed") } else { if (event.upstream is CachesEnabled) { event.upstream.setCaches(caches) } addUpstream(event.upstream) - log.info("Upstream ${event.upstream.getId()} with chain $chain has been added") + if (!started) { + start() + } + log.error("Upstream ${event.upstream.getId()} with chain $chain has been added") } } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/MultistreamHolder.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/MultistreamHolder.kt index 6e1c4ab0..80b4e564 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/MultistreamHolder.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/MultistreamHolder.kt @@ -17,7 +17,6 @@ package io.emeraldpay.dshackle.upstream import io.emeraldpay.grpc.Chain -import reactor.core.publisher.Flux /** * Holds Multistreams configured for a chain. @@ -25,6 +24,5 @@ import reactor.core.publisher.Flux interface MultistreamHolder { fun getUpstream(chain: Chain): Multistream? fun getAvailable(): List - fun observeChains(): Flux fun isAvailable(chain: Chain): Boolean } 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 14ebfda0..aca3d567 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinMultistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinMultistream.kt @@ -19,22 +19,21 @@ import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.upstream.* +import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.calls.DefaultBitcoinMethods import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.grpc.Chain import org.slf4j.LoggerFactory -import org.springframework.context.Lifecycle import reactor.core.publisher.Mono @Suppress("UNCHECKED_CAST") open class BitcoinMultistream( chain: Chain, private val sourceUpstreams: MutableList, - caches: Caches, - callTargetsHolder: CallTargetsHolder -) : Multistream(chain, sourceUpstreams as MutableList, caches, RequestPostprocessor.Empty(), callTargetsHolder), Lifecycle { + caches: Caches +) : Multistream(chain, sourceUpstreams as MutableList, caches, RequestPostprocessor.Empty()), Lifecycle { companion object { private val log = LoggerFactory.getLogger(BitcoinMultistream::class.java) @@ -136,7 +135,7 @@ open class BitcoinMultistream( } override fun isRunning(): Boolean { - return super.isRunning() || reader.isRunning + return super.isRunning() || reader.isRunning() } override fun start() { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinReader.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinReader.kt index f9f0496c..8ae68e9f 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinReader.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinReader.kt @@ -19,13 +19,13 @@ import com.fasterxml.jackson.databind.ObjectMapper import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.Head +import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.bitcoin.data.SimpleUnspent import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import org.bitcoinj.core.Address import org.slf4j.LoggerFactory -import org.springframework.context.Lifecycle import reactor.core.publisher.Mono import reactor.kotlin.core.publisher.cast @@ -72,7 +72,7 @@ open class BitcoinReader( } override fun isRunning(): Boolean { - return mempool.isRunning + return mempool.isRunning() } override fun start() { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinRpcHead.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinRpcHead.kt index 72f500bc..5c018104 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinRpcHead.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinRpcHead.kt @@ -19,11 +19,11 @@ import io.emeraldpay.dshackle.Defaults import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.upstream.AbstractHead import io.emeraldpay.dshackle.upstream.Head +import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import org.slf4j.LoggerFactory -import org.springframework.context.Lifecycle import org.springframework.scheduling.concurrent.CustomizableThreadFactory import reactor.core.Disposable import reactor.core.publisher.Flux diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinRpcUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinRpcUpstream.kt index e5ab2d29..23b82f35 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinRpcUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinRpcUpstream.kt @@ -20,6 +20,7 @@ import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.startup.QuorumForLabels import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.Head +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 @@ -27,7 +28,6 @@ import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.grpc.Chain import org.slf4j.LoggerFactory -import org.springframework.context.Lifecycle import reactor.core.Disposable open class BitcoinRpcUpstream( @@ -92,7 +92,7 @@ open class BitcoinRpcUpstream( override fun isRunning(): Boolean { var runningAny = validatorSubscription != null if (head is Lifecycle) { - runningAny = runningAny || head.isRunning + runningAny = runningAny || head.isRunning() } return runningAny } @@ -100,7 +100,7 @@ open class BitcoinRpcUpstream( override fun start() { log.info("Configured for ${chain.chainName}") if (head is Lifecycle) { - if (!head.isRunning) { + if (!head.isRunning()) { head.start() } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinZMQHead.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinZMQHead.kt index 4ce4bcc1..07ed5837 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinZMQHead.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinZMQHead.kt @@ -5,12 +5,12 @@ import io.emeraldpay.dshackle.data.BlockContainer import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.upstream.AbstractHead import io.emeraldpay.dshackle.upstream.Head +import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import org.apache.commons.codec.binary.Hex import org.slf4j.LoggerFactory -import org.springframework.context.Lifecycle import reactor.core.Disposable import reactor.core.publisher.Flux import reactor.core.publisher.Mono @@ -65,6 +65,6 @@ class BitcoinZMQHead( } override fun isRunning(): Boolean { - return server.isRunning || refreshSubscription != null + return server.isRunning() || refreshSubscription != null } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/CachingMempoolData.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/CachingMempoolData.kt index 3e26e5a5..f4df9af9 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/CachingMempoolData.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/CachingMempoolData.kt @@ -18,11 +18,11 @@ package io.emeraldpay.dshackle.upstream.bitcoin import com.fasterxml.jackson.databind.ObjectMapper import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.upstream.Head +import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import org.slf4j.LoggerFactory -import org.springframework.context.Lifecycle import reactor.core.Disposable import reactor.core.publisher.Mono import java.time.Duration diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/ZMQServer.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/ZMQServer.kt index 1358f568..fb208deb 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/ZMQServer.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/ZMQServer.kt @@ -1,7 +1,7 @@ package io.emeraldpay.dshackle.upstream.bitcoin +import io.emeraldpay.dshackle.upstream.Lifecycle import org.slf4j.LoggerFactory -import org.springframework.context.Lifecycle import org.zeromq.SocketType import org.zeromq.ZContext import org.zeromq.ZMQ 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 e2fb57f5..fd7f4cf3 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumMultistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumMultistream.kt @@ -21,13 +21,13 @@ import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.upstream.* +import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstream import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest 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.Flux import reactor.core.publisher.Mono @@ -36,9 +36,8 @@ import reactor.core.publisher.Mono open class EthereumMultistream( chain: Chain, val upstreams: MutableList, - caches: Caches, - callTargetsHolder: CallTargetsHolder -) : Multistream(chain, upstreams as MutableList, caches, CacheRequested(caches), callTargetsHolder), EthereumLikeMultistream { + caches: Caches +) : Multistream(chain, upstreams as MutableList, caches, CacheRequested(caches)), EthereumLikeMultistream { companion object { private val log = LoggerFactory.getLogger(EthereumMultistream::class.java) @@ -80,7 +79,7 @@ open class EthereumMultistream( } override fun isRunning(): Boolean { - return super.isRunning() || reader.isRunning + return super.isRunning() || reader.isRunning() } override fun getReader(): EthereumReader { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumReader.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumReader.kt index 8501a99b..6f3ea506 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumReader.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumReader.kt @@ -30,6 +30,7 @@ import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.RekeyingReader import io.emeraldpay.dshackle.reader.RpcReader import io.emeraldpay.dshackle.reader.TransformingReader +import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Multistream import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.etherjar.domain.Address @@ -41,7 +42,6 @@ import io.emeraldpay.etherjar.rpc.json.TransactionJson import io.emeraldpay.etherjar.rpc.json.TransactionRefJson import org.apache.commons.collections4.Factory import org.slf4j.LoggerFactory -import org.springframework.context.Lifecycle import java.util.function.Function /** @@ -174,7 +174,7 @@ open class EthereumReader( override fun isRunning(): Boolean { // TODO should be always running? - return up.isRunning + return true // up.isRunning } override fun start() { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumRpcHead.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumRpcHead.kt index 00637d54..16331153 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumRpcHead.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumRpcHead.kt @@ -18,11 +18,11 @@ package io.emeraldpay.dshackle.upstream.ethereum import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.upstream.BlockValidator +import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import org.slf4j.LoggerFactory -import org.springframework.context.Lifecycle import org.springframework.scheduling.concurrent.CustomizableThreadFactory import reactor.core.Disposable import reactor.core.publisher.Flux 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 3bbe468f..bb2479df 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumRpcUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumRpcUpstream.kt @@ -80,7 +80,7 @@ open class EthereumRpcUpstream( } override fun isRunning(): Boolean { - return connector.isRunning + return connector.isRunning() } override fun getApi(): Reader { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumWsHead.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumWsHead.kt index 1532d5d8..2ed626f6 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumWsHead.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumWsHead.kt @@ -17,10 +17,10 @@ package io.emeraldpay.dshackle.upstream.ethereum import io.emeraldpay.dshackle.upstream.BlockValidator +import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcWsClient import org.slf4j.LoggerFactory -import org.springframework.context.Lifecycle import reactor.core.Disposable import reactor.core.publisher.Flux 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 85ecf596..c9b2b69d 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 @@ -2,9 +2,9 @@ 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.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse -import org.springframework.context.Lifecycle interface EthereumConnector : Lifecycle { fun getHead(): Head 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 2b46a499..21fefba1 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 @@ -5,6 +5,7 @@ import io.emeraldpay.dshackle.cache.CachesEnabled import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.upstream.BlockValidator import io.emeraldpay.dshackle.upstream.Head +import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.MergedHead import io.emeraldpay.dshackle.upstream.ethereum.EthereumRpcHead import io.emeraldpay.dshackle.upstream.ethereum.EthereumWsFactory @@ -14,7 +15,6 @@ import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import org.slf4j.LoggerFactory -import org.springframework.context.Lifecycle import java.time.Duration class EthereumRpcConnector( @@ -61,7 +61,7 @@ class EthereumRpcConnector( override fun isRunning(): Boolean { if (head is Lifecycle) { - return head.isRunning + return head.isRunning() } return true } 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 0a248e2c..609e88cf 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 @@ -36,7 +36,7 @@ class EthereumWsConnector( } override fun isRunning(): Boolean { - return head.isRunning + return head.isRunning() } override fun stop() { 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 3e93a81a..71477c37 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,13 +21,13 @@ import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.upstream.* +import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.forkchoice.PriorityForkChoice import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstream import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest 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.Flux import reactor.core.publisher.Mono @@ -36,9 +36,8 @@ import reactor.core.publisher.Mono open class EthereumPosMultiStream( chain: Chain, val upstreams: MutableList, - caches: Caches, - callTargetsHolder: CallTargetsHolder -) : Multistream(chain, upstreams as MutableList, caches, CacheRequested(caches), callTargetsHolder), EthereumLikeMultistream { + caches: Caches +) : Multistream(chain, upstreams as MutableList, caches, CacheRequested(caches)), EthereumLikeMultistream { companion object { private val log = LoggerFactory.getLogger(EthereumPosMultiStream::class.java) @@ -75,7 +74,7 @@ open class EthereumPosMultiStream( } override fun isRunning(): Boolean { - return super.isRunning() || reader.isRunning + return super.isRunning() || reader.isRunning() } override fun getReader(): EthereumReader { 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 50f201c2..78ce2c24 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 @@ -22,6 +22,7 @@ import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.startup.QuorumForLabels import io.emeraldpay.dshackle.upstream.Head +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 @@ -31,7 +32,6 @@ import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.grpc.Chain import org.slf4j.LoggerFactory -import org.springframework.context.Lifecycle import reactor.core.Disposable open class EthereumPosRpcUpstream( @@ -80,7 +80,7 @@ open class EthereumPosRpcUpstream( } override fun isRunning(): Boolean { - return connector.isRunning + return connector.isRunning() } override fun getApi(): Reader { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/BitcoinGrpcUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/BitcoinGrpcUpstream.kt index a539a21a..35397e43 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/BitcoinGrpcUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/BitcoinGrpcUpstream.kt @@ -24,6 +24,7 @@ import io.emeraldpay.dshackle.data.BlockId import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.Head +import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.UpstreamAvailability @@ -37,7 +38,6 @@ import io.emeraldpay.etherjar.rpc.RpcException import io.emeraldpay.grpc.Chain import org.reactivestreams.Publisher import org.slf4j.LoggerFactory -import org.springframework.context.Lifecycle import reactor.core.publisher.Flux import reactor.core.publisher.Mono import java.math.BigInteger 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 d407d3b5..c0427a06 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcUpstream.kt @@ -26,6 +26,7 @@ import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.startup.QuorumForLabels import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.Head +import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.UpstreamAvailability @@ -40,7 +41,6 @@ import io.emeraldpay.etherjar.rpc.RpcException import io.emeraldpay.grpc.Chain import org.reactivestreams.Publisher import org.slf4j.LoggerFactory -import org.springframework.context.Lifecycle import reactor.core.publisher.Flux import reactor.core.publisher.Mono import java.math.BigInteger 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 10a005cc..72c03c28 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumPosGrpcUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumPosGrpcUpstream.kt @@ -26,6 +26,7 @@ import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.startup.QuorumForLabels import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.Head +import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.UpstreamAvailability @@ -40,7 +41,6 @@ import io.emeraldpay.etherjar.rpc.RpcException import io.emeraldpay.grpc.Chain import org.reactivestreams.Publisher import org.slf4j.LoggerFactory -import org.springframework.context.Lifecycle import reactor.core.publisher.Flux import reactor.core.publisher.Mono import java.math.BigInteger diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcHead.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcHead.kt index ac2db057..78883e1c 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcHead.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcHead.kt @@ -22,12 +22,12 @@ import io.emeraldpay.dshackle.Defaults import io.emeraldpay.dshackle.data.BlockContainer import io.emeraldpay.dshackle.upstream.AbstractHead import io.emeraldpay.dshackle.upstream.DefaultUpstream +import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice import io.emeraldpay.grpc.Chain import org.reactivestreams.Publisher import org.slf4j.LoggerFactory -import org.springframework.context.Lifecycle import reactor.core.Disposable import reactor.core.publisher.Flux import reactor.core.publisher.Mono @@ -60,7 +60,7 @@ class GrpcHead( * Initiate a new head subscription with connection to the remote */ private fun internalStart(remote: ReactorBlockchainGrpc.ReactorBlockchainStub) { - if (this.isRunning) { + if (this.isRunning()) { stop() } log.debug("Start Head subscription to ${parent.getId()}") diff --git a/src/main/resources/log4j2.xml b/src/main/resources/log4j2.xml index 08e5d92d..807a731a 100644 --- a/src/main/resources/log4j2.xml +++ b/src/main/resources/log4j2.xml @@ -41,4 +41,4 @@ - \ No newline at end of file + diff --git a/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy index 43352b51..93cc1fbc 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy @@ -246,9 +246,7 @@ class NativeCallSpec extends Specification { def "Returns error for unsupported chain"() { setup: - def upstreams = Mock(MultistreamHolder) { - _ * it.observeChains() >> Flux.empty() - } + def upstreams = Mock(MultistreamHolder) def nativeCall = nativeCall(upstreams) def req = BlockchainOuterClass.NativeCallRequest.newBuilder() diff --git a/src/test/groovy/io/emeraldpay/dshackle/test/MultistreamHolderMock.groovy b/src/test/groovy/io/emeraldpay/dshackle/test/MultistreamHolderMock.groovy index 9e34c546..e26cc9f7 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/test/MultistreamHolderMock.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/test/MultistreamHolderMock.groovy @@ -50,7 +50,7 @@ class MultistreamHolderMock implements MultistreamHolder { if (up instanceof EthereumPosMultiStream) { upstreams[chain] = up } else if (up instanceof EthereumPosRpcUpstream) { - upstreams[chain] = new EthereumPosMultiStream(chain, [up as EthereumPosRpcUpstream], Caches.default(), TestingCommons.callTargetsHolder) + upstreams[chain] = new EthereumPosMultiStream(chain, [up as EthereumPosRpcUpstream], Caches.default()) } else { throw new IllegalArgumentException("Unsupported upstream type ${up.class}") } @@ -59,7 +59,7 @@ class MultistreamHolderMock implements MultistreamHolder { if (up instanceof BitcoinMultistream) { upstreams[chain] = up } else if (up instanceof BitcoinRpcUpstream) { - upstreams[chain] = new BitcoinMultistream(chain, [up as BitcoinRpcUpstream], Caches.default(), TestingCommons.callTargetsHolder) + upstreams[chain] = new BitcoinMultistream(chain, [up as BitcoinRpcUpstream], Caches.default()) } else { throw new IllegalArgumentException("Unsupported upstream type ${up.class}") } @@ -81,11 +81,6 @@ class MultistreamHolderMock implements MultistreamHolder { return upstreams.keySet().toList() } - @Override - Flux observeChains() { - return Flux.fromIterable(getAvailable()) - } - @Override boolean isAvailable(@NotNull Chain chain) { return upstreams.containsKey(chain) @@ -98,7 +93,7 @@ class MultistreamHolderMock implements MultistreamHolder { Head customHead = null EthereumMultistreamMock(@NotNull Chain chain, @NotNull List upstreams, @NotNull Caches caches) { - super(chain, upstreams, caches, TestingCommons.callTargetsHolder) + super(chain, upstreams, caches) } EthereumMultistreamMock(@NotNull Chain chain, @NotNull List upstreams) { diff --git a/src/test/groovy/io/emeraldpay/dshackle/test/TestingCommons.groovy b/src/test/groovy/io/emeraldpay/dshackle/test/TestingCommons.groovy index f2c640d7..7bfab88c 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/test/TestingCommons.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/test/TestingCommons.groovy @@ -29,7 +29,6 @@ import io.emeraldpay.dshackle.upstream.Multistream import io.emeraldpay.dshackle.upstream.calls.DirectCallMethods import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream -import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosUpstream import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.etherjar.domain.BlockHash @@ -80,7 +79,7 @@ class TestingCommons { } static Multistream multistream(EthereumPosRpcUpstreamMock up) { - return new EthereumPosMultiStream(Chain.ETHEREUM, [up], Caches.default(), callTargetsHolder).tap { + return new EthereumPosMultiStream(Chain.ETHEREUM, [up], Caches.default()).tap { start() } } @@ -101,11 +100,11 @@ class TestingCommons { } static Multistream multistreamWithoutUpstreams(Chain chain) { - return new EthereumPosMultiStream(chain, [], emptyCaches().getCaches(chain), callTargetsHolder) + return new EthereumPosMultiStream(chain, [], emptyCaches().getCaches(chain)) } static Multistream multistreamClassicWithoutUpstreams(Chain chain) { - return new EthereumMultistream(chain, [], emptyCaches().getCaches(chain), callTargetsHolder) + return new EthereumMultistream(chain, [], emptyCaches().getCaches(chain)) } static FileResolver fileResolver() { diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/MergedHeadSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/MergedHeadSpec.groovy index e1b72360..a6e08be3 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/MergedHeadSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/MergedHeadSpec.groovy @@ -16,7 +16,7 @@ package io.emeraldpay.dshackle.upstream import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice -import org.springframework.context.Lifecycle +import io.emeraldpay.dshackle.upstream.Lifecycle import reactor.core.publisher.Flux import spock.lang.Specification diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/MultistreamSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/MultistreamSpec.groovy index b432bcb8..8a28f4b3 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/MultistreamSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/MultistreamSpec.groovy @@ -50,7 +50,7 @@ class MultistreamSpec extends Specification { setup: def up1 = new EthereumPosRpcUpstreamMock("test1", Chain.ETHEREUM, TestingCommons.api(), new DirectCallMethods(["eth_test1", "eth_test2"])) def up2 = new EthereumPosRpcUpstreamMock("test1", Chain.ETHEREUM, TestingCommons.api(), new DirectCallMethods(["eth_test2", "eth_test3"])) - def aggr = new EthereumPosMultiStream(Chain.ETHEREUM, [up1, up2], Caches.default(), TestingCommons.callTargetsHolder) + def aggr = new EthereumPosMultiStream(Chain.ETHEREUM, [up1, up2], Caches.default()) when: aggr.onUpstreamsUpdated() def act = aggr.getMethods() @@ -206,7 +206,7 @@ class MultistreamSpec extends Specification { def up1 = TestingCommons.upstream("test-1", "internal") def up2 = TestingCommons.upstream("test-2", "external") def up3 = TestingCommons.upstream("test-3", "external") - def multistream = new EthereumPosMultiStream(Chain.ETHEREUM, [up1, up2, up3], Caches.default(), TestingCommons.callTargetsHolder) + def multistream = new EthereumPosMultiStream(Chain.ETHEREUM, [up1, up2, up3], Caches.default()) expect: multistream.getHead(new Selector.LabelMatcher("provider", ["internal"])).is(up1.ethereumHeadMock) @@ -345,7 +345,7 @@ class MultistreamSpec extends Specification { class TestMultistream extends Multistream { TestMultistream(List upstreams, @NotNull RequestPostprocessor postprocessor) { - super(Chain.ETHEREUM, upstreams, Caches.default(), postprocessor, TestingCommons.callTargetsHolder) + super(Chain.ETHEREUM, upstreams, Caches.default(), postprocessor) } @Override @@ -386,7 +386,7 @@ class MultistreamSpec extends Specification { class TestEthereumPosMultistream extends EthereumPosMultiStream { TestEthereumPosMultistream(@NotNull Chain chain, @NotNull List upstreams, @NotNull Caches caches) { - super(chain, upstreams, caches, TestingCommons.callTargetsHolder) + super(chain, upstreams, caches) } @Override From b450cc9a9407742e0cae62f2b5c197d4978205b5 Mon Sep 17 00:00:00 2001 From: a10zn8 Date: Thu, 17 Nov 2022 13:52:36 +0400 Subject: [PATCH 5/7] use only chain map --- .../upstream/CurrentMultistreamHolder.kt | 21 +++++-------------- .../dshackle/upstream/Multistream.kt | 4 ++++ 2 files changed, 9 insertions(+), 16 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt index b4865951..e4f6029d 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt @@ -19,31 +19,23 @@ package io.emeraldpay.dshackle.upstream import io.emeraldpay.grpc.Chain import org.slf4j.LoggerFactory import org.springframework.stereotype.Component -import reactor.core.publisher.Sinks -import java.util.* -import java.util.concurrent.ConcurrentHashMap -import java.util.concurrent.locks.ReentrantLock import javax.annotation.PreDestroy -import kotlin.concurrent.withLock @Component open class CurrentMultistreamHolder( - private val multistreams: List + multistreams: List ) : MultistreamHolder { private val log = LoggerFactory.getLogger(CurrentMultistreamHolder::class.java) - private val chainMapping = ConcurrentHashMap().apply { - multistreams.forEach { this[it.chain] = it } - } - private val updateLock = ReentrantLock() + private val chainMapping = multistreams.associateBy { it.chain } override fun getUpstream(chain: Chain): Multistream? { return chainMapping[chain] } override fun getAvailable(): List { - return multistreams.asSequence() + return chainMapping.values.asSequence() .filter { it.isAvailable() } .map { it.chain } .toList() @@ -56,11 +48,8 @@ open class CurrentMultistreamHolder( @PreDestroy fun shutdown() { log.info("Closing upstream connections...") - updateLock.withLock { - chainMapping.values.forEach { - it.stop() - } - chainMapping.clear() + chainMapping.values.forEach { + it.stop() } } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt index 54fb2dfd..30fa4f24 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt @@ -324,6 +324,10 @@ abstract class Multistream( log.info("State of ${chain.chainCode}: height=${height ?: '?'}, status=[$statuses], lag=[$lag], weak=[$weak]") } + fun test(event: UpstreamChangeEvent): Boolean { + return event.chain == this.chain + } + @EventListener @Order(Ordered.HIGHEST_PRECEDENCE) fun onUpstreamChange(event: UpstreamChangeEvent) { From e66bb772ad06d02f5996770ddbc94f48906088c4 Mon Sep 17 00:00:00 2001 From: a10zn8 Date: Fri, 18 Nov 2022 19:20:41 +0400 Subject: [PATCH 6/7] small fixes for integration tests --- .../dshackle/startup/ConfiguredUpstreams.kt | 1 - testing/dshackle/dshackle-basic.yaml | 17 +++++++++++------ 2 files changed, 11 insertions(+), 7 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt b/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt index b66a125a..51aed124 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt @@ -111,7 +111,6 @@ open class ConfiguredUpstreams( } upstream?.let { val event = UpstreamChangeEvent(chain, upstream, UpstreamChangeEvent.ChangeType.ADDED) - log.error("first !!!! Upstream ${event.upstream.getId()} with chain $chain has been added") eventPublisher.publishEvent(event) } } diff --git a/testing/dshackle/dshackle-basic.yaml b/testing/dshackle/dshackle-basic.yaml index 51318221..56515c78 100644 --- a/testing/dshackle/dshackle-basic.yaml +++ b/testing/dshackle/dshackle-basic.yaml @@ -7,6 +7,7 @@ tls: cluster: upstreams: - id: test-1 + node-id: 1 chain: ethereum methods: enabled: @@ -15,17 +16,21 @@ cluster: options: disable-validation: true connection: - ethereum: - rpc: - url: "http://localhost:18545" + ethereum-pos: + execution: + rpc: + url: "http://localhost:18545" + - id: test-2 + node-id: 2 chain: ethereum options: disable-validation: true connection: - ethereum: - rpc: - url: "http://localhost:18546" + execution: + ethereum-pos: + rpc: + url: "http://localhost:18546" cache: redis: From 5aefa7ea22b9ebc285b1441dc319516e5983dbce Mon Sep 17 00:00:00 2001 From: a10zn8 Date: Fri, 18 Nov 2022 19:47:35 +0400 Subject: [PATCH 7/7] describe should return all chains with upstream, not only available --- .../emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt | 2 +- src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt | 3 +++ 2 files changed, 4 insertions(+), 1 deletion(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt index e4f6029d..54a152b3 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt @@ -36,7 +36,7 @@ open class CurrentMultistreamHolder( override fun getAvailable(): List { return chainMapping.values.asSequence() - .filter { it.isAvailable() } + .filter { it.haveUpstreams() } .map { it.chain } .toList() } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt index 30fa4f24..b6de244c 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt @@ -349,6 +349,9 @@ abstract class Multistream( } } + fun haveUpstreams(): Boolean = + upstreams.isNotEmpty() + // -------------------------------------------------------------------------------------------------------- class UpstreamStatus(val upstream: Upstream, val status: UpstreamAvailability, val ts: Instant = Instant.now())