use only chain map
This commit is contained in:
@@ -19,31 +19,23 @@ package io.emeraldpay.dshackle.upstream
|
|||||||
import io.emeraldpay.grpc.Chain
|
import io.emeraldpay.grpc.Chain
|
||||||
import org.slf4j.LoggerFactory
|
import org.slf4j.LoggerFactory
|
||||||
import org.springframework.stereotype.Component
|
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 javax.annotation.PreDestroy
|
||||||
import kotlin.concurrent.withLock
|
|
||||||
|
|
||||||
@Component
|
@Component
|
||||||
open class CurrentMultistreamHolder(
|
open class CurrentMultistreamHolder(
|
||||||
private val multistreams: List<Multistream>
|
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>().apply {
|
private val chainMapping = multistreams.associateBy { it.chain }
|
||||||
multistreams.forEach { this[it.chain] = it }
|
|
||||||
}
|
|
||||||
private val updateLock = ReentrantLock()
|
|
||||||
|
|
||||||
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 multistreams.asSequence()
|
return chainMapping.values.asSequence()
|
||||||
.filter { it.isAvailable() }
|
.filter { it.isAvailable() }
|
||||||
.map { it.chain }
|
.map { it.chain }
|
||||||
.toList()
|
.toList()
|
||||||
@@ -56,11 +48,8 @@ open class CurrentMultistreamHolder(
|
|||||||
@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()
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -324,6 +324,10 @@ 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
|
@EventListener
|
||||||
@Order(Ordered.HIGHEST_PRECEDENCE)
|
@Order(Ordered.HIGHEST_PRECEDENCE)
|
||||||
fun onUpstreamChange(event: UpstreamChangeEvent) {
|
fun onUpstreamChange(event: UpstreamChangeEvent) {
|
||||||
|
|||||||
Reference in New Issue
Block a user