Merge pull request #37 from p2p-org/upstream_refactoring_2

Multistreams refactoring
This commit is contained in:
a10zn8
2022-11-22 17:03:01 +04:00
committed by GitHub
43 changed files with 333 additions and 304 deletions

View File

@@ -289,3 +289,6 @@ detekt {
tasks.withType(Detekt).configureEach { tasks.withType(Detekt).configureEach {
jvmTarget = "13" jvmTarget = "13"
} }
// formats code for each build
tasks.findByName("ktlintCheck").dependsOn("ktlintFormat")

View File

@@ -20,12 +20,10 @@ import org.slf4j.LoggerFactory
import org.springframework.boot.ResourceBanner import org.springframework.boot.ResourceBanner
import org.springframework.boot.SpringApplication import org.springframework.boot.SpringApplication
import org.springframework.boot.autoconfigure.SpringBootApplication import org.springframework.boot.autoconfigure.SpringBootApplication
import org.springframework.context.annotation.Import
import org.springframework.core.io.ClassPathResource import org.springframework.core.io.ClassPathResource
import org.springframework.core.io.support.ResourcePropertySource import org.springframework.core.io.support.ResourcePropertySource
@SpringBootApplication(scanBasePackages = ["io.emeraldpay.dshackle"]) @SpringBootApplication(scanBasePackages = ["io.emeraldpay.dshackle"])
@Import(Config::class)
open class Starter open class Starter
private val log = LoggerFactory.getLogger(Starter::class.java) private val log = LoggerFactory.getLogger(Starter::class.java)

View File

@@ -0,0 +1,78 @@
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
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(val beanFactory: ConfigurableListableBeanFactory) {
@Bean
open fun allMultistreams(
cachesFactory: CachesFactory,
callTargetsHolder: CallTargetsHolder
): List<Multistream> {
return Chain.values()
.filterNot { it == Chain.UNSPECIFIED }
.mapNotNull { chain ->
when (BlockchainType.from(chain)) {
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)
}
}

View File

@@ -26,11 +26,10 @@ import io.emeraldpay.dshackle.quorum.NotLaggingQuorum
import io.emeraldpay.dshackle.quorum.QuorumReaderFactory import io.emeraldpay.dshackle.quorum.QuorumReaderFactory
import io.emeraldpay.dshackle.quorum.QuorumRpcReader import io.emeraldpay.dshackle.quorum.QuorumRpcReader
import io.emeraldpay.dshackle.startup.ConfiguredUpstreams import io.emeraldpay.dshackle.startup.ConfiguredUpstreams
import io.emeraldpay.dshackle.upstream.ApiSource import io.emeraldpay.dshackle.startup.UpstreamChangeEvent
import io.emeraldpay.dshackle.upstream.Multistream import io.emeraldpay.dshackle.upstream.*
import io.emeraldpay.dshackle.upstream.MultistreamHolder
import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.calls.EthereumCallSelector 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.EthereumMultistream
import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream import io.emeraldpay.dshackle.upstream.ethereum.EthereumPosMultiStream
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcError import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcError
@@ -45,7 +44,7 @@ import io.emeraldpay.grpc.Chain
import io.micrometer.core.instrument.Metrics import io.micrometer.core.instrument.Metrics
import org.apache.commons.lang3.StringUtils import org.apache.commons.lang3.StringUtils
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.beans.factory.annotation.Autowired import org.springframework.context.event.EventListener
import org.springframework.stereotype.Service import org.springframework.stereotype.Service
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
@@ -55,9 +54,9 @@ import java.util.concurrent.atomic.AtomicInteger
@Service @Service
open class NativeCall( open class NativeCall(
@Autowired private val multistreamHolder: MultistreamHolder, private val multistreamHolder: MultistreamHolder,
@Autowired private val configuredUpstreams: ConfiguredUpstreams, private val configuredUpstreams: ConfiguredUpstreams,
@Autowired private val signer: ResponseSigner private val signer: ResponseSigner
) { ) {
private val log = LoggerFactory.getLogger(NativeCall::class.java) private val log = LoggerFactory.getLogger(NativeCall::class.java)
@@ -66,18 +65,19 @@ open class NativeCall(
var quorumReaderFactory: QuorumReaderFactory = QuorumReaderFactory.default() var quorumReaderFactory: QuorumReaderFactory = QuorumReaderFactory.default()
private val ethereumCallSelectors = EnumMap<Chain, EthereumCallSelector>(Chain::class.java) private val ethereumCallSelectors = EnumMap<Chain, EthereumCallSelector>(Chain::class.java)
init { companion object {
val casting = mapOf( val casting: Map<BlockchainType, Class<out EthereumLikeMultistream>> = mapOf(
BlockchainType.EVM_POS to EthereumPosMultiStream::class.java, 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()))
} }
} }
} }

View File

@@ -20,6 +20,7 @@ import io.emeraldpay.api.proto.Common
import io.emeraldpay.api.proto.ReactorBlockchainGrpc import io.emeraldpay.api.proto.ReactorBlockchainGrpc
import io.emeraldpay.dshackle.Defaults import io.emeraldpay.dshackle.Defaults
import io.emeraldpay.dshackle.SilentException import io.emeraldpay.dshackle.SilentException
import io.emeraldpay.dshackle.startup.UpstreamChangeEvent
import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.Capability
import io.emeraldpay.dshackle.upstream.MultistreamHolder import io.emeraldpay.dshackle.upstream.MultistreamHolder
import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Selector
@@ -33,12 +34,12 @@ import org.bitcoinj.params.MainNetParams
import org.bitcoinj.params.TestNet3Params import org.bitcoinj.params.TestNet3Params
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.beans.factory.annotation.Autowired import org.springframework.beans.factory.annotation.Autowired
import org.springframework.context.event.EventListener
import org.springframework.stereotype.Service import org.springframework.stereotype.Service
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
import java.math.BigInteger import java.math.BigInteger
import java.util.concurrent.ConcurrentHashMap import java.util.concurrent.ConcurrentHashMap
import javax.annotation.PostConstruct
@Service @Service
class TrackBitcoinAddress( class TrackBitcoinAddress(
@@ -67,15 +68,13 @@ class TrackBitcoinAddress(
Selector.CapabilityMatcher(Capability.BALANCE) Selector.CapabilityMatcher(Capability.BALANCE)
) )
@PostConstruct @EventListener
fun listenChains() { fun onUpstreamChangeEvent(event: UpstreamChangeEvent) {
multistreamHolder.observeChains().subscribe { chain -> multistreamHolder.getUpstream(event.chain)?.let { mup ->
multistreamHolder.getUpstream(chain)?.let { mup ->
val available = mup.getAll().any { up -> val available = mup.getAll().any { up ->
!up.isGrpc() && up.getCapabilities().contains(Capability.BALANCE) !up.isGrpc() && up.getCapabilities().contains(Capability.BALANCE)
} }
setBalanceAvailability(chain, available) setBalanceAvailability(event.chain, available)
}
} }
} }

View File

@@ -21,13 +21,7 @@ import io.emeraldpay.dshackle.FileResolver
import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.BlockValidator import io.emeraldpay.dshackle.upstream.*
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.bitcoin.BitcoinRpcHead import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinRpcHead
import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinRpcUpstream import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinRpcUpstream
import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinZMQHead import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinZMQHead
@@ -50,28 +44,29 @@ import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import io.emeraldpay.grpc.BlockchainType import io.emeraldpay.grpc.BlockchainType
import io.emeraldpay.grpc.Chain import io.emeraldpay.grpc.Chain
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.beans.factory.annotation.Autowired import org.springframework.boot.ApplicationArguments
import org.springframework.stereotype.Repository import org.springframework.boot.ApplicationRunner
import org.springframework.context.ApplicationEventPublisher
import org.springframework.stereotype.Component
import java.net.URI import java.net.URI
import java.util.concurrent.atomic.AtomicInteger import java.util.concurrent.atomic.AtomicInteger
import java.util.function.Function import java.util.function.Function
import javax.annotation.PostConstruct
import kotlin.math.abs import kotlin.math.abs
@Repository @Component
open class ConfiguredUpstreams( open class ConfiguredUpstreams(
@Autowired private val currentUpstreams: CurrentMultistreamHolder, private val fileResolver: FileResolver,
@Autowired private val fileResolver: FileResolver, private val config: UpstreamsConfig,
@Autowired private val config: UpstreamsConfig private val callTargets: CallTargetsHolder,
) { private val eventPublisher: ApplicationEventPublisher
) : ApplicationRunner {
private val log = LoggerFactory.getLogger(ConfiguredUpstreams::class.java) private val log = LoggerFactory.getLogger(ConfiguredUpstreams::class.java)
private var seq = AtomicInteger(0) private var seq = AtomicInteger(0)
private val hashes: MutableMap<Byte, Boolean> = HashMap() private val hashes: MutableMap<Byte, Boolean> = HashMap()
@PostConstruct override fun run(args: ApplicationArguments) {
fun start() {
log.debug("Starting upstreams") log.debug("Starting upstreams")
val defaultOptions = buildDefaultOptions(config) val defaultOptions = buildDefaultOptions(config)
config.upstreams.forEach { up -> config.upstreams.forEach { up ->
@@ -115,7 +110,8 @@ open class ConfiguredUpstreams(
} }
} }
upstream?.let { 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 { fun buildMethods(config: UpstreamsConfig.Upstream<*>, chain: Chain): CallMethods {
return if (config.methods != null) { return if (config.methods != null) {
ManagedCallMethods( ManagedCallMethods(
currentUpstreams.getDefaultMethods(chain), callTargets.getDefaultMethods(chain),
config.methods!!.enabled.map { it.name }.toSet(), config.methods!!.enabled.map { it.name }.toSet(),
config.methods!!.disabled.map { it.name }.toSet() config.methods!!.disabled.map { it.name }.toSet()
).also { ).also {
@@ -164,7 +160,7 @@ open class ConfiguredUpstreams(
} }
} }
} else { } else {
currentUpstreams.getDefaultMethods(chain) callTargets.getDefaultMethods(chain)
} }
} }
@@ -315,7 +311,7 @@ open class ConfiguredUpstreams(
.doOnNext { .doOnNext {
log.info("Chain ${it.chain} ${it.type} through gRPC at ${endpoint.host}:${endpoint.port}. With caps: ${it.upstream.getCapabilities()}") 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<URI>? = null): HttpRpcFactory? { private fun buildHttpFactory(conn: UpstreamsConfig.RpcConnection, urls: ArrayList<URI>? = null): HttpRpcFactory? {

View File

@@ -24,7 +24,7 @@ import io.emeraldpay.grpc.Chain
/** /**
* An update event to the list of currently available upstreams. * An update event to the list of currently available upstreams.
*/ */
class UpstreamChange( class UpstreamChangeEvent(
/** /**
* Target blockchain * Target blockchain
*/ */

View File

@@ -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<Chain, CallMethods>()
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
}
}

View File

@@ -16,156 +16,40 @@
*/ */
package io.emeraldpay.dshackle.upstream 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 io.emeraldpay.grpc.Chain
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.beans.factory.annotation.Autowired import org.springframework.stereotype.Component
import org.springframework.stereotype.Repository
import reactor.core.publisher.Flux
import reactor.core.publisher.Sinks
import java.util.Collections
import java.util.concurrent.Callable
import java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.locks.ReentrantLock
import javax.annotation.PreDestroy import javax.annotation.PreDestroy
import kotlin.concurrent.withLock
@Repository @Component
open class CurrentMultistreamHolder( open class CurrentMultistreamHolder(
@Autowired private val cachesFactory: CachesFactory multistreams: List<Multistream>
) : MultistreamHolder { ) : MultistreamHolder {
private val log = LoggerFactory.getLogger(CurrentMultistreamHolder::class.java) private val log = LoggerFactory.getLogger(CurrentMultistreamHolder::class.java)
private val chainMapping = ConcurrentHashMap<Chain, Multistream>() private val chainMapping = multistreams.associateBy { it.chain }
private val chainsBus = Sinks.many()
.multicast()
.directBestEffort<Chain>()
private val callTargets = HashMap<Chain, CallMethods>()
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[chain]
val factory = Callable<Multistream> {
EthereumMultistream(chain, ArrayList(), cachesFactory.getCaches(chain))
}
processUpdate(change, up, current, factory)
}
BlockchainType.EVM_POS -> {
val up = change.upstream.cast(EthereumPosUpstream::class.java)
val current = chainMapping[chain]
val factory = Callable<Multistream> {
EthereumPosMultiStream(chain, ArrayList(), cachesFactory.getCaches(chain))
}
processUpdate(change, up, current, factory)
}
BlockchainType.BITCOIN -> {
val up = change.upstream.cast(BitcoinUpstream::class.java)
val current = chainMapping[chain]
val factory = Callable<Multistream> {
BitcoinMultistream(chain, ArrayList(), cachesFactory.getCaches(chain))
}
processUpdate(change, up, current, factory)
}
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?, factory: Callable<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 (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 (!callTargets.containsKey(chain)) {
setupDefaultMethods(chain)
}
log.info("Upstream ${change.upstream.getId()} with chain $chain has been added")
}
}
override fun getUpstream(chain: Chain): Multistream? { override fun getUpstream(chain: Chain): Multistream? {
return chainMapping[chain] return chainMapping[chain]
} }
override fun getAvailable(): List<Chain> { override fun getAvailable(): List<Chain> {
return Collections.unmodifiableList(chainMapping.keys.toList()) return chainMapping.values.asSequence()
} .filter { it.haveUpstreams() }
.map { it.chain }
override fun observeChains(): Flux<Chain> { .toList()
return Flux.concat(
Flux.fromIterable(getAvailable()),
chainsBus.asFlux()
)
}
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 { override fun isAvailable(chain: Chain): Boolean {
return chainMapping.containsKey(chain) && callTargets.containsKey(chain) return chainMapping.getValue(chain).isAvailable()
} }
@PreDestroy @PreDestroy
fun shutdown() { fun shutdown() {
log.info("Closing upstream connections...") log.info("Closing upstream connections...")
updateLock.withLock {
chainMapping.values.forEach { chainMapping.values.forEach {
it.stop() it.stop()
} }
chainMapping.clear()
}
} }
} }

View File

@@ -0,0 +1,7 @@
package io.emeraldpay.dshackle.upstream
interface Lifecycle {
fun start()
fun stop()
fun isRunning(): Boolean
}

View File

@@ -21,7 +21,6 @@ import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.cache.CachesEnabled import io.emeraldpay.dshackle.cache.CachesEnabled
import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import reactor.core.Disposable import reactor.core.Disposable
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
@@ -43,7 +42,7 @@ class MergedHead(
override fun start() { override fun start() {
super.start() super.start()
sources.forEach { head -> sources.forEach { head ->
if (head is Lifecycle && !head.isRunning) { if (head is Lifecycle && !head.isRunning()) {
head.start() head.start()
} }
} }
@@ -56,7 +55,7 @@ class MergedHead(
override fun stop() { override fun stop() {
super.stop() super.stop()
sources.forEach { head -> sources.forEach { head ->
if (head is Lifecycle && head.isRunning) { if (head is Lifecycle && head.isRunning()) {
head.stop() head.stop()
} }
} }

View File

@@ -17,8 +17,10 @@
package io.emeraldpay.dshackle.upstream package io.emeraldpay.dshackle.upstream
import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.cache.CachesEnabled
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.startup.UpstreamChangeEvent
import io.emeraldpay.dshackle.upstream.calls.AggregatedCallMethods import io.emeraldpay.dshackle.upstream.calls.AggregatedCallMethods
import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
@@ -29,13 +31,15 @@ import io.micrometer.core.instrument.Tag
import org.apache.commons.collections4.Factory import org.apache.commons.collections4.Factory
import org.apache.commons.collections4.FunctorException import org.apache.commons.collections4.FunctorException
import org.slf4j.LoggerFactory 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.Disposable
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
import java.time.Duration import java.time.Duration
import java.time.Instant import java.time.Instant
import java.util.Locale import java.util.*
import java.util.concurrent.atomic.AtomicReference import java.util.concurrent.atomic.AtomicReference
import java.util.concurrent.locks.ReentrantLock import java.util.concurrent.locks.ReentrantLock
import java.util.function.Predicate import java.util.function.Predicate
@@ -56,6 +60,8 @@ abstract class Multistream(
private const val metrics = "upstreams" private const val metrics = "upstreams"
} }
private var started = false
private var cacheSubscription: Disposable? = null private var cacheSubscription: Disposable? = null
private val reconfigLock = ReentrantLock() private val reconfigLock = ReentrantLock()
private var callMethods: CallMethods? = null private var callMethods: CallMethods? = null
@@ -231,6 +237,7 @@ abstract class Multistream(
// print status _change_ every 15 seconds, at most; otherwise prints it on interval of 30 seconds // print status _change_ every 15 seconds, at most; otherwise prints it on interval of 30 seconds
.sample(Duration.ofSeconds(15)) .sample(Duration.ofSeconds(15))
.subscribe { printStatus() } .subscribe { printStatus() }
started = true
} }
override fun stop() { override fun stop() {
@@ -244,6 +251,7 @@ abstract class Multistream(
} }
} }
lagObserver?.stop() lagObserver?.stop()
started = false
} }
fun onHeadUpdated(head: Head) { fun onHeadUpdated(head: Head) {
@@ -316,6 +324,34 @@ abstract class Multistream(
log.info("State of ${chain.chainCode}: height=${height ?: '?'}, status=[$statuses], lag=[$lag], weak=[$weak]") 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) {
val chain = event.chain
if (this.chain == chain) {
if (event.type == UpstreamChangeEvent.ChangeType.REMOVED) {
removeUpstream(event.upstream.getId())
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)
if (!started) {
start()
}
log.error("Upstream ${event.upstream.getId()} with chain $chain has been added")
}
}
}
fun haveUpstreams(): Boolean =
upstreams.isNotEmpty()
// -------------------------------------------------------------------------------------------------------- // --------------------------------------------------------------------------------------------------------
class UpstreamStatus(val upstream: Upstream, val status: UpstreamAvailability, val ts: Instant = Instant.now()) class UpstreamStatus(val upstream: Upstream, val status: UpstreamAvailability, val ts: Instant = Instant.now())

View File

@@ -16,9 +16,7 @@
*/ */
package io.emeraldpay.dshackle.upstream package io.emeraldpay.dshackle.upstream
import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.grpc.Chain import io.emeraldpay.grpc.Chain
import reactor.core.publisher.Flux
/** /**
* Holds Multistreams configured for a chain. * Holds Multistreams configured for a chain.
@@ -26,7 +24,5 @@ import reactor.core.publisher.Flux
interface MultistreamHolder { interface MultistreamHolder {
fun getUpstream(chain: Chain): Multistream? fun getUpstream(chain: Chain): Multistream?
fun getAvailable(): List<Chain> fun getAvailable(): List<Chain>
fun observeChains(): Flux<Chain>
fun getDefaultMethods(chain: Chain): CallMethods
fun isAvailable(chain: Chain): Boolean fun isAvailable(chain: Chain): Boolean
} }

View File

@@ -18,21 +18,14 @@ package io.emeraldpay.dshackle.upstream.bitcoin
import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.ChainFees import io.emeraldpay.dshackle.upstream.*
import io.emeraldpay.dshackle.upstream.EmptyHead import io.emeraldpay.dshackle.upstream.Lifecycle
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.calls.DefaultBitcoinMethods import io.emeraldpay.dshackle.upstream.calls.DefaultBitcoinMethods
import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import io.emeraldpay.grpc.Chain import io.emeraldpay.grpc.Chain
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
@Suppress("UNCHECKED_CAST") @Suppress("UNCHECKED_CAST")
@@ -142,7 +135,7 @@ open class BitcoinMultistream(
} }
override fun isRunning(): Boolean { override fun isRunning(): Boolean {
return super.isRunning() || reader.isRunning return super.isRunning() || reader.isRunning()
} }
override fun start() { override fun start() {

View File

@@ -19,13 +19,13 @@ import com.fasterxml.jackson.databind.ObjectMapper
import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.Capability
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.bitcoin.data.SimpleUnspent import io.emeraldpay.dshackle.upstream.bitcoin.data.SimpleUnspent
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import org.bitcoinj.core.Address import org.bitcoinj.core.Address
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
import reactor.kotlin.core.publisher.cast import reactor.kotlin.core.publisher.cast
@@ -72,7 +72,7 @@ open class BitcoinReader(
} }
override fun isRunning(): Boolean { override fun isRunning(): Boolean {
return mempool.isRunning return mempool.isRunning()
} }
override fun start() { override fun start() {

View File

@@ -19,11 +19,11 @@ import io.emeraldpay.dshackle.Defaults
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.AbstractHead import io.emeraldpay.dshackle.upstream.AbstractHead
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import org.springframework.scheduling.concurrent.CustomizableThreadFactory import org.springframework.scheduling.concurrent.CustomizableThreadFactory
import reactor.core.Disposable import reactor.core.Disposable
import reactor.core.publisher.Flux import reactor.core.publisher.Flux

View File

@@ -20,6 +20,7 @@ import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.startup.QuorumForLabels import io.emeraldpay.dshackle.startup.QuorumForLabels
import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.Capability
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.UpstreamAvailability
import io.emeraldpay.dshackle.upstream.calls.CallMethods 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.dshackle.upstream.rpcclient.JsonRpcResponse
import io.emeraldpay.grpc.Chain import io.emeraldpay.grpc.Chain
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import reactor.core.Disposable import reactor.core.Disposable
open class BitcoinRpcUpstream( open class BitcoinRpcUpstream(
@@ -92,7 +92,7 @@ open class BitcoinRpcUpstream(
override fun isRunning(): Boolean { override fun isRunning(): Boolean {
var runningAny = validatorSubscription != null var runningAny = validatorSubscription != null
if (head is Lifecycle) { if (head is Lifecycle) {
runningAny = runningAny || head.isRunning runningAny = runningAny || head.isRunning()
} }
return runningAny return runningAny
} }
@@ -100,7 +100,7 @@ open class BitcoinRpcUpstream(
override fun start() { override fun start() {
log.info("Configured for ${chain.chainName}") log.info("Configured for ${chain.chainName}")
if (head is Lifecycle) { if (head is Lifecycle) {
if (!head.isRunning) { if (!head.isRunning()) {
head.start() head.start()
} }
} }

View File

@@ -5,12 +5,12 @@ import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.AbstractHead import io.emeraldpay.dshackle.upstream.AbstractHead
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import org.apache.commons.codec.binary.Hex import org.apache.commons.codec.binary.Hex
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import reactor.core.Disposable import reactor.core.Disposable
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
@@ -65,6 +65,6 @@ class BitcoinZMQHead(
} }
override fun isRunning(): Boolean { override fun isRunning(): Boolean {
return server.isRunning || refreshSubscription != null return server.isRunning() || refreshSubscription != null
} }
} }

View File

@@ -18,11 +18,11 @@ package io.emeraldpay.dshackle.upstream.bitcoin
import com.fasterxml.jackson.databind.ObjectMapper import com.fasterxml.jackson.databind.ObjectMapper
import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import reactor.core.Disposable import reactor.core.Disposable
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
import java.time.Duration import java.time.Duration

View File

@@ -1,7 +1,7 @@
package io.emeraldpay.dshackle.upstream.bitcoin package io.emeraldpay.dshackle.upstream.bitcoin
import io.emeraldpay.dshackle.upstream.Lifecycle
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import org.zeromq.SocketType import org.zeromq.SocketType
import org.zeromq.ZContext import org.zeromq.ZContext
import org.zeromq.ZMQ import org.zeromq.ZMQ

View File

@@ -20,20 +20,14 @@ import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.ChainFees import io.emeraldpay.dshackle.upstream.*
import io.emeraldpay.dshackle.upstream.EmptyHead import io.emeraldpay.dshackle.upstream.Lifecycle
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.forkchoice.MostWorkForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice
import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstream import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstream
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import io.emeraldpay.grpc.Chain import io.emeraldpay.grpc.Chain
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import org.springframework.util.ConcurrentReferenceHashMap import org.springframework.util.ConcurrentReferenceHashMap
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
@@ -85,7 +79,7 @@ open class EthereumMultistream(
} }
override fun isRunning(): Boolean { override fun isRunning(): Boolean {
return super.isRunning() || reader.isRunning return super.isRunning() || reader.isRunning()
} }
override fun getReader(): EthereumReader { override fun getReader(): EthereumReader {

View File

@@ -30,6 +30,7 @@ import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.reader.RekeyingReader import io.emeraldpay.dshackle.reader.RekeyingReader
import io.emeraldpay.dshackle.reader.RpcReader import io.emeraldpay.dshackle.reader.RpcReader
import io.emeraldpay.dshackle.reader.TransformingReader import io.emeraldpay.dshackle.reader.TransformingReader
import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.Multistream import io.emeraldpay.dshackle.upstream.Multistream
import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.etherjar.domain.Address 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 io.emeraldpay.etherjar.rpc.json.TransactionRefJson
import org.apache.commons.collections4.Factory import org.apache.commons.collections4.Factory
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import java.util.function.Function import java.util.function.Function
/** /**
@@ -174,7 +174,7 @@ open class EthereumReader(
override fun isRunning(): Boolean { override fun isRunning(): Boolean {
// TODO should be always running? // TODO should be always running?
return up.isRunning return true // up.isRunning
} }
override fun start() { override fun start() {

View File

@@ -18,11 +18,11 @@ package io.emeraldpay.dshackle.upstream.ethereum
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.BlockValidator import io.emeraldpay.dshackle.upstream.BlockValidator
import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import org.springframework.scheduling.concurrent.CustomizableThreadFactory import org.springframework.scheduling.concurrent.CustomizableThreadFactory
import reactor.core.Disposable import reactor.core.Disposable
import reactor.core.publisher.Flux import reactor.core.publisher.Flux

View File

@@ -80,7 +80,7 @@ open class EthereumRpcUpstream(
} }
override fun isRunning(): Boolean { override fun isRunning(): Boolean {
return connector.isRunning return connector.isRunning()
} }
override fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse> { override fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse> {

View File

@@ -17,10 +17,10 @@
package io.emeraldpay.dshackle.upstream.ethereum package io.emeraldpay.dshackle.upstream.ethereum
import io.emeraldpay.dshackle.upstream.BlockValidator import io.emeraldpay.dshackle.upstream.BlockValidator
import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcWsClient import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcWsClient
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import reactor.core.Disposable import reactor.core.Disposable
import reactor.core.publisher.Flux import reactor.core.publisher.Flux

View File

@@ -2,9 +2,9 @@ package io.emeraldpay.dshackle.upstream.ethereum.connectors
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.Head 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.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import org.springframework.context.Lifecycle
interface EthereumConnector : Lifecycle { interface EthereumConnector : Lifecycle {
fun getHead(): Head fun getHead(): Head

View File

@@ -5,6 +5,7 @@ import io.emeraldpay.dshackle.cache.CachesEnabled
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.BlockValidator import io.emeraldpay.dshackle.upstream.BlockValidator
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.MergedHead import io.emeraldpay.dshackle.upstream.MergedHead
import io.emeraldpay.dshackle.upstream.ethereum.EthereumRpcHead import io.emeraldpay.dshackle.upstream.ethereum.EthereumRpcHead
import io.emeraldpay.dshackle.upstream.ethereum.EthereumWsFactory 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.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import java.time.Duration import java.time.Duration
class EthereumRpcConnector( class EthereumRpcConnector(
@@ -61,7 +61,7 @@ class EthereumRpcConnector(
override fun isRunning(): Boolean { override fun isRunning(): Boolean {
if (head is Lifecycle) { if (head is Lifecycle) {
return head.isRunning return head.isRunning()
} }
return true return true
} }

View File

@@ -36,7 +36,7 @@ class EthereumWsConnector(
} }
override fun isRunning(): Boolean { override fun isRunning(): Boolean {
return head.isRunning return head.isRunning()
} }
override fun stop() { override fun stop() {

View File

@@ -20,20 +20,14 @@ import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.ChainFees import io.emeraldpay.dshackle.upstream.*
import io.emeraldpay.dshackle.upstream.EmptyHead import io.emeraldpay.dshackle.upstream.Lifecycle
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.forkchoice.PriorityForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.PriorityForkChoice
import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstream import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstream
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import io.emeraldpay.grpc.Chain import io.emeraldpay.grpc.Chain
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import org.springframework.util.ConcurrentReferenceHashMap import org.springframework.util.ConcurrentReferenceHashMap
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
@@ -80,7 +74,7 @@ open class EthereumPosMultiStream(
} }
override fun isRunning(): Boolean { override fun isRunning(): Boolean {
return super.isRunning() || reader.isRunning return super.isRunning() || reader.isRunning()
} }
override fun getReader(): EthereumReader { override fun getReader(): EthereumReader {

View File

@@ -22,6 +22,7 @@ import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.startup.QuorumForLabels import io.emeraldpay.dshackle.startup.QuorumForLabels
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.UpstreamAvailability
import io.emeraldpay.dshackle.upstream.calls.CallMethods 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.dshackle.upstream.rpcclient.JsonRpcResponse
import io.emeraldpay.grpc.Chain import io.emeraldpay.grpc.Chain
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import reactor.core.Disposable import reactor.core.Disposable
open class EthereumPosRpcUpstream( open class EthereumPosRpcUpstream(
@@ -80,7 +80,7 @@ open class EthereumPosRpcUpstream(
} }
override fun isRunning(): Boolean { override fun isRunning(): Boolean {
return connector.isRunning return connector.isRunning()
} }
override fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse> { override fun getApi(): Reader<JsonRpcRequest, JsonRpcResponse> {

View File

@@ -24,6 +24,7 @@ import io.emeraldpay.dshackle.data.BlockId
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.Capability
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.UpstreamAvailability
@@ -37,7 +38,6 @@ import io.emeraldpay.etherjar.rpc.RpcException
import io.emeraldpay.grpc.Chain import io.emeraldpay.grpc.Chain
import org.reactivestreams.Publisher import org.reactivestreams.Publisher
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
import java.math.BigInteger import java.math.BigInteger

View File

@@ -26,6 +26,7 @@ import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.startup.QuorumForLabels import io.emeraldpay.dshackle.startup.QuorumForLabels
import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.Capability
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.UpstreamAvailability
@@ -40,7 +41,6 @@ import io.emeraldpay.etherjar.rpc.RpcException
import io.emeraldpay.grpc.Chain import io.emeraldpay.grpc.Chain
import org.reactivestreams.Publisher import org.reactivestreams.Publisher
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
import java.math.BigInteger import java.math.BigInteger

View File

@@ -26,6 +26,7 @@ import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.startup.QuorumForLabels import io.emeraldpay.dshackle.startup.QuorumForLabels
import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.Capability
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.UpstreamAvailability
@@ -40,7 +41,6 @@ import io.emeraldpay.etherjar.rpc.RpcException
import io.emeraldpay.grpc.Chain import io.emeraldpay.grpc.Chain
import org.reactivestreams.Publisher import org.reactivestreams.Publisher
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
import java.math.BigInteger import java.math.BigInteger

View File

@@ -22,12 +22,12 @@ import io.emeraldpay.dshackle.Defaults
import io.emeraldpay.dshackle.data.BlockContainer import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.upstream.AbstractHead import io.emeraldpay.dshackle.upstream.AbstractHead
import io.emeraldpay.dshackle.upstream.DefaultUpstream import io.emeraldpay.dshackle.upstream.DefaultUpstream
import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.UpstreamAvailability
import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice
import io.emeraldpay.grpc.Chain import io.emeraldpay.grpc.Chain
import org.reactivestreams.Publisher import org.reactivestreams.Publisher
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import reactor.core.Disposable import reactor.core.Disposable
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
@@ -60,7 +60,7 @@ class GrpcHead(
* Initiate a new head subscription with connection to the remote * Initiate a new head subscription with connection to the remote
*/ */
private fun internalStart(remote: ReactorBlockchainGrpc.ReactorBlockchainStub) { private fun internalStart(remote: ReactorBlockchainGrpc.ReactorBlockchainStub) {
if (this.isRunning) { if (this.isRunning()) {
stop() stop()
} }
log.debug("Start Head subscription to ${parent.getId()}") log.debug("Start Head subscription to ${parent.getId()}")

View File

@@ -22,7 +22,7 @@ import io.emeraldpay.dshackle.Defaults
import io.emeraldpay.dshackle.FileResolver import io.emeraldpay.dshackle.FileResolver
import io.emeraldpay.dshackle.config.AuthConfig import io.emeraldpay.dshackle.config.AuthConfig
import io.emeraldpay.dshackle.config.UpstreamsConfig 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.DefaultUpstream
import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.UpstreamAvailability
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcGrpcClient import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcGrpcClient
@@ -69,7 +69,7 @@ class GrpcUpstreams(
private val known = HashMap<Chain, DefaultUpstream>() private val known = HashMap<Chain, DefaultUpstream>()
private val lock = ReentrantLock() private val lock = ReentrantLock()
fun start(): Flux<UpstreamChange> { fun start(): Flux<UpstreamChangeEvent> {
val channel: ManagedChannelBuilder<*> = if (auth != null && StringUtils.isNotEmpty(auth.ca)) { val channel: ManagedChannelBuilder<*> = if (auth != null && StringUtils.isNotEmpty(auth.ca)) {
NettyChannelBuilder.forAddress(host, port) NettyChannelBuilder.forAddress(host, port)
// some messages are very large. many of them in megabytes, some even in gigabytes (ex. ETH Traces) // 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 return updates
} }
fun processDescription(value: BlockchainOuterClass.DescribeResponse): Flux<UpstreamChange> { fun processDescription(value: BlockchainOuterClass.DescribeResponse): Flux<UpstreamChangeEvent> {
val current = value.chainsList.filter { val current = value.chainsList.filter {
Chain.byId(it.chain.number) != Chain.UNSPECIFIED Chain.byId(it.chain.number) != Chain.UNSPECIFIED
}.mapNotNull { chainDetails -> }.mapNotNull { chainDetails ->
@@ -138,14 +138,14 @@ class GrpcUpstreams(
} }
val added = current.filter { val added = current.filter {
it.type == UpstreamChange.ChangeType.ADDED it.type == UpstreamChangeEvent.ChangeType.ADDED
} }
val removed = known.filterNot { kv -> val removed = known.filterNot { kv ->
val stillCurrent = current.any { c -> c.chain == kv.key } val stillCurrent = current.any { c -> c.chain == kv.key }
stillCurrent stillCurrent
}.map { }.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) return Flux.fromIterable(removed + added)
} }
@@ -172,7 +172,7 @@ class GrpcUpstreams(
return sslContext.build() return sslContext.build()
} }
fun getOrCreate(chain: Chain): UpstreamChange { fun getOrCreate(chain: Chain): UpstreamChangeEvent {
val metricsTags = listOf( val metricsTags = listOf(
Tag.of("upstream", id), Tag.of("upstream", id),
Tag.of("chain", chain.chainCode) 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 { lock.withLock {
val current = known[chain] val current = known[chain]
return if (current == null) { return if (current == null) {
@@ -211,14 +211,14 @@ class GrpcUpstreams(
created.timeout = this.timeout created.timeout = this.timeout
known[chain] = created known[chain] = created
created.start() created.start()
UpstreamChange(chain, created, UpstreamChange.ChangeType.ADDED) UpstreamChangeEvent(chain, created, UpstreamChangeEvent.ChangeType.ADDED)
} else { } 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 { lock.withLock {
val current = known[chain] val current = known[chain]
return if (current == null) { return if (current == null) {
@@ -227,14 +227,14 @@ class GrpcUpstreams(
created.timeout = this.timeout created.timeout = this.timeout
known[chain] = created known[chain] = created
created.start() created.start()
UpstreamChange(chain, created, UpstreamChange.ChangeType.ADDED) UpstreamChangeEvent(chain, created, UpstreamChangeEvent.ChangeType.ADDED)
} else { } 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 { lock.withLock {
val current = known[chain] val current = known[chain]
return if (current == null) { return if (current == null) {
@@ -243,9 +243,9 @@ class GrpcUpstreams(
created.timeout = this.timeout created.timeout = this.timeout
known[chain] = created known[chain] = created
created.start() created.start()
UpstreamChange(chain, created, UpstreamChange.ChangeType.ADDED) UpstreamChangeEvent(chain, created, UpstreamChangeEvent.ChangeType.ADDED)
} else { } else {
UpstreamChange(chain, current, UpstreamChange.ChangeType.REVALIDATED) UpstreamChangeEvent(chain, current, UpstreamChangeEvent.ChangeType.REVALIDATED)
} }
} }
} }

View File

@@ -246,9 +246,7 @@ class NativeCallSpec extends Specification {
def "Returns error for unsupported chain"() { def "Returns error for unsupported chain"() {
setup: setup:
def upstreams = Mock(MultistreamHolder) { def upstreams = Mock(MultistreamHolder)
_ * it.observeChains() >> Flux.empty()
}
def nativeCall = nativeCall(upstreams) def nativeCall = nativeCall(upstreams)
def req = BlockchainOuterClass.NativeCallRequest.newBuilder() def req = BlockchainOuterClass.NativeCallRequest.newBuilder()

View File

@@ -4,21 +4,24 @@ import io.emeraldpay.dshackle.FileResolver
import io.emeraldpay.dshackle.cache.CachesFactory import io.emeraldpay.dshackle.cache.CachesFactory
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.quorum.NonEmptyQuorum import io.emeraldpay.dshackle.quorum.NonEmptyQuorum
import io.emeraldpay.dshackle.upstream.CallTargetsHolder
import io.emeraldpay.dshackle.upstream.CurrentMultistreamHolder import io.emeraldpay.dshackle.upstream.CurrentMultistreamHolder
import io.emeraldpay.dshackle.upstream.calls.DefaultEthereumMethods import io.emeraldpay.dshackle.upstream.calls.DefaultEthereumMethods
import io.emeraldpay.dshackle.upstream.calls.ManagedCallMethods import io.emeraldpay.dshackle.upstream.calls.ManagedCallMethods
import io.emeraldpay.grpc.Chain import io.emeraldpay.grpc.Chain
import org.springframework.context.ApplicationEventPublisher
import spock.lang.Specification import spock.lang.Specification
class ConfiguredUpstreamsSpec extends Specification { class ConfiguredUpstreamsSpec extends Specification {
def "Applied quorum to extra methods"() { def "Applied quorum to extra methods"() {
setup: setup:
def currentUpstreams = Mock(CurrentMultistreamHolder) { def callTargetsHolder = new CallTargetsHolder()
_ * getDefaultMethods(Chain.ETHEREUM) >> new DefaultEthereumMethods(Chain.ETHEREUM)
}
def configurer = new ConfiguredUpstreams( def configurer = new ConfiguredUpstreams(
currentUpstreams, Stub(FileResolver), Stub(UpstreamsConfig) Stub(FileResolver),
Stub(UpstreamsConfig),
callTargetsHolder,
Mock(ApplicationEventPublisher)
) )
def methods = new UpstreamsConfig.Methods( def methods = new UpstreamsConfig.Methods(
[ [
@@ -38,11 +41,12 @@ class ConfiguredUpstreamsSpec extends Specification {
def "Got static response from extra methods"() { def "Got static response from extra methods"() {
setup: setup:
def currentUpstreams = Mock(CurrentMultistreamHolder) { def callTargetsHolder = new CallTargetsHolder()
_ * getDefaultMethods(Chain.ETHEREUM) >> new DefaultEthereumMethods(Chain.ETHEREUM)
}
def configurer = new ConfiguredUpstreams( def configurer = new ConfiguredUpstreams(
currentUpstreams, Stub(FileResolver), Stub(UpstreamsConfig) Stub(FileResolver),
Stub(UpstreamsConfig),
callTargetsHolder,
Mock(ApplicationEventPublisher)
) )
def methods = new UpstreamsConfig.Methods( def methods = new UpstreamsConfig.Methods(
[ [
@@ -61,7 +65,12 @@ class ConfiguredUpstreamsSpec extends Specification {
def "Calculate node-id"() { def "Calculate node-id"() {
setup: 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: expect:
configurer.getHash(node, src) == expected configurer.getHash(node, src) == expected
@@ -75,7 +84,12 @@ class ConfiguredUpstreamsSpec extends Specification {
def "Calculate node-id conflicting results"() { def "Calculate node-id conflicting results"() {
setup: 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: when:
def h1 = configurer.getHash(null, "hohoho") def h1 = configurer.getHash(null, "hohoho")

View File

@@ -81,20 +81,6 @@ class MultistreamHolderMock implements MultistreamHolder {
return upstreams.keySet().toList() return upstreams.keySet().toList()
} }
@Override
Flux<Chain> observeChains() {
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 @Override
boolean isAvailable(@NotNull Chain chain) { boolean isAvailable(@NotNull Chain chain) {
return upstreams.containsKey(chain) return upstreams.containsKey(chain)

View File

@@ -24,8 +24,10 @@ import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.data.BlockId import io.emeraldpay.dshackle.data.BlockId
import io.emeraldpay.dshackle.reader.EmptyReader import io.emeraldpay.dshackle.reader.EmptyReader
import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.CallTargetsHolder
import io.emeraldpay.dshackle.upstream.Multistream import io.emeraldpay.dshackle.upstream.Multistream
import io.emeraldpay.dshackle.upstream.calls.DirectCallMethods 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.EthereumPosMultiStream
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
@@ -90,6 +92,21 @@ class TestingCommons {
return new CachesFactory(new CacheConfig()) return new CachesFactory(new CacheConfig())
} }
static List<Multistream> defaultMultistreams() {
return [
multistreamWithoutUpstreams(Chain.ETHEREUM),
multistreamClassicWithoutUpstreams(Chain.ETHEREUM_CLASSIC)
]
}
static Multistream multistreamWithoutUpstreams(Chain chain) {
return new EthereumPosMultiStream(chain, [], emptyCaches().getCaches(chain))
}
static Multistream multistreamClassicWithoutUpstreams(Chain chain) {
return new EthereumMultistream(chain, [], emptyCaches().getCaches(chain))
}
static FileResolver fileResolver() { static FileResolver fileResolver() {
return new FileResolver(new File("src/test/resources")) return new FileResolver(new File("src/test/resources"))
} }
@@ -126,4 +143,6 @@ class TestingCommons {
} }
static MeterRegistry meterRegistry = new LoggingMeterRegistry() static MeterRegistry meterRegistry = new LoggingMeterRegistry()
static CallTargetsHolder callTargetsHolder = new CallTargetsHolder()
} }

View File

@@ -15,7 +15,7 @@
*/ */
package io.emeraldpay.dshackle.upstream 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.EthereumPosRpcUpstreamMock
import io.emeraldpay.dshackle.test.EthereumRpcUpstreamMock import io.emeraldpay.dshackle.test.EthereumRpcUpstreamMock
import io.emeraldpay.dshackle.test.TestingCommons import io.emeraldpay.dshackle.test.TestingCommons
@@ -26,10 +26,10 @@ class CurrentMultistreamHolderSpec extends Specification {
def "add upstream"() { def "add upstream"() {
setup: setup:
def current = new CurrentMultistreamHolder(TestingCommons.emptyCaches()) def current = new CurrentMultistreamHolder(TestingCommons.defaultMultistreams())
def up = new EthereumPosRpcUpstreamMock("test", Chain.ETHEREUM, TestingCommons.api()) def up = new EthereumPosRpcUpstreamMock("test", Chain.ETHEREUM, TestingCommons.api())
when: 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: then:
current.getAvailable() == [Chain.ETHEREUM] current.getAvailable() == [Chain.ETHEREUM]
current.getUpstream(Chain.ETHEREUM).getAll()[0] == up current.getUpstream(Chain.ETHEREUM).getAll()[0] == up
@@ -37,14 +37,15 @@ class CurrentMultistreamHolderSpec extends Specification {
def "add multiple upstreams"() { def "add multiple upstreams"() {
setup: setup:
def current = new CurrentMultistreamHolder(TestingCommons.emptyCaches()) def current = new CurrentMultistreamHolder(TestingCommons.defaultMultistreams())
def up1 = new EthereumPosRpcUpstreamMock("test1", Chain.ETHEREUM, TestingCommons.api()) def up1 = new EthereumPosRpcUpstreamMock("test1", Chain.ETHEREUM, TestingCommons.api())
def up2 = new EthereumRpcUpstreamMock("test2", Chain.ETHEREUM_CLASSIC, TestingCommons.api()) def up2 = new EthereumRpcUpstreamMock("test2", Chain.ETHEREUM_CLASSIC, TestingCommons.api())
def up3 = new EthereumPosRpcUpstreamMock("test3", Chain.ETHEREUM, TestingCommons.api()) def up3 = new EthereumPosRpcUpstreamMock("test3", Chain.ETHEREUM, TestingCommons.api())
when: when:
current.update(new UpstreamChange(Chain.ETHEREUM, up1, UpstreamChange.ChangeType.ADDED)) current.getUpstream(Chain.ETHEREUM).onUpstreamChange(new UpstreamChangeEvent(Chain.ETHEREUM, up1, UpstreamChangeEvent.ChangeType.ADDED))
current.update(new UpstreamChange(Chain.ETHEREUM_CLASSIC, up2, UpstreamChange.ChangeType.ADDED)) current.getUpstream(Chain.ETHEREUM_CLASSIC).onUpstreamChange(new UpstreamChangeEvent(Chain.ETHEREUM_CLASSIC, up2, UpstreamChangeEvent.ChangeType.ADDED))
current.update(new UpstreamChange(Chain.ETHEREUM, up3, UpstreamChange.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: then:
current.getAvailable().toSet() == [Chain.ETHEREUM, Chain.ETHEREUM_CLASSIC].toSet() current.getAvailable().toSet() == [Chain.ETHEREUM, Chain.ETHEREUM_CLASSIC].toSet()
current.getUpstream(Chain.ETHEREUM).getAll().toSet() == [up1, up3].toSet() current.getUpstream(Chain.ETHEREUM).getAll().toSet() == [up1, up3].toSet()
@@ -53,16 +54,16 @@ class CurrentMultistreamHolderSpec extends Specification {
def "remove upstream"() { def "remove upstream"() {
setup: setup:
def current = new CurrentMultistreamHolder(TestingCommons.emptyCaches()) def current = new CurrentMultistreamHolder(TestingCommons.defaultMultistreams())
def up1 = new EthereumPosRpcUpstreamMock("test1", Chain.ETHEREUM, TestingCommons.api()) def up1 = new EthereumPosRpcUpstreamMock("test1", Chain.ETHEREUM, TestingCommons.api())
def up2 = new EthereumRpcUpstreamMock("test2", Chain.ETHEREUM_CLASSIC, TestingCommons.api()) def up2 = new EthereumRpcUpstreamMock("test2", Chain.ETHEREUM_CLASSIC, TestingCommons.api())
def up3 = new EthereumPosRpcUpstreamMock("test3", Chain.ETHEREUM, TestingCommons.api()) def up3 = new EthereumPosRpcUpstreamMock("test3", Chain.ETHEREUM, TestingCommons.api())
def up1_del = new EthereumPosRpcUpstreamMock("test1", Chain.ETHEREUM, TestingCommons.api()) def up1_del = new EthereumPosRpcUpstreamMock("test1", Chain.ETHEREUM, TestingCommons.api())
when: when:
current.update(new UpstreamChange(Chain.ETHEREUM, up1, UpstreamChange.ChangeType.ADDED)) current.getUpstream(Chain.ETHEREUM).onUpstreamChange(new UpstreamChangeEvent(Chain.ETHEREUM, up1, UpstreamChangeEvent.ChangeType.ADDED))
current.update(new UpstreamChange(Chain.ETHEREUM_CLASSIC, up2, UpstreamChange.ChangeType.ADDED)) current.getUpstream(Chain.ETHEREUM_CLASSIC).onUpstreamChange(new UpstreamChangeEvent(Chain.ETHEREUM_CLASSIC, up2, UpstreamChangeEvent.ChangeType.ADDED))
current.update(new UpstreamChange(Chain.ETHEREUM, up3, UpstreamChange.ChangeType.ADDED)) current.getUpstream(Chain.ETHEREUM).onUpstreamChange(new UpstreamChangeEvent(Chain.ETHEREUM, up3, UpstreamChangeEvent.ChangeType.ADDED))
current.update(new UpstreamChange(Chain.ETHEREUM, up1_del, UpstreamChange.ChangeType.REMOVED)) current.getUpstream(Chain.ETHEREUM).onUpstreamChange(new UpstreamChangeEvent(Chain.ETHEREUM, up1_del, UpstreamChangeEvent.ChangeType.REMOVED))
then: then:
current.getAvailable().toSet() == [Chain.ETHEREUM, Chain.ETHEREUM_CLASSIC].toSet() current.getAvailable().toSet() == [Chain.ETHEREUM, Chain.ETHEREUM_CLASSIC].toSet()
current.getUpstream(Chain.ETHEREUM).getAll().toSet() == [up3].toSet() current.getUpstream(Chain.ETHEREUM).getAll().toSet() == [up3].toSet()
@@ -71,7 +72,7 @@ class CurrentMultistreamHolderSpec extends Specification {
def "available after adding"() { def "available after adding"() {
setup: setup:
def current = new CurrentMultistreamHolder(TestingCommons.emptyCaches()) def current = new CurrentMultistreamHolder(TestingCommons.defaultMultistreams())
def up1 = new EthereumPosRpcUpstreamMock("test1", Chain.ETHEREUM, TestingCommons.api()) def up1 = new EthereumPosRpcUpstreamMock("test1", Chain.ETHEREUM, TestingCommons.api())
when: when:
@@ -80,7 +81,7 @@ class CurrentMultistreamHolderSpec extends Specification {
!act !act
when: 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) act = current.isAvailable(Chain.ETHEREUM)
then: then:

View File

@@ -16,7 +16,7 @@
package io.emeraldpay.dshackle.upstream package io.emeraldpay.dshackle.upstream
import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.MostWorkForkChoice
import org.springframework.context.Lifecycle import io.emeraldpay.dshackle.upstream.Lifecycle
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
import spock.lang.Specification import spock.lang.Specification

View File

@@ -7,6 +7,7 @@ tls:
cluster: cluster:
upstreams: upstreams:
- id: test-1 - id: test-1
node-id: 1
chain: ethereum chain: ethereum
methods: methods:
enabled: enabled:
@@ -15,15 +16,19 @@ cluster:
options: options:
disable-validation: true disable-validation: true
connection: connection:
ethereum: ethereum-pos:
execution:
rpc: rpc:
url: "http://localhost:18545" url: "http://localhost:18545"
- id: test-2 - id: test-2
node-id: 2
chain: ethereum chain: ethereum
options: options:
disable-validation: true disable-validation: true
connection: connection:
ethereum: execution:
ethereum-pos:
rpc: rpc:
url: "http://localhost:18546" url: "http://localhost:18546"