problem: doesn't start head subscription for a single grpc based upstream; no validation for second bitcoin upstream

This commit is contained in:
Igor Artamonov
2021-11-18 17:46:43 -05:00
parent f918e1b138
commit e51f4861d3
3 changed files with 14 additions and 5 deletions

View File

@@ -72,13 +72,18 @@ open class BitcoinMultistream(
val head = if (upstreams.size == 1) { val head = if (upstreams.size == 1) {
val upstream = upstreams.first() val upstream = upstreams.first()
upstream.setLag(0) upstream.setLag(0)
upstream.getHead() upstream.getHead().apply {
if (this is Lifecycle) {
this.start()
}
}
} else { } else {
val newHead = MergedHead(upstreams.map { it.getHead() }).apply { val newHead = MergedHead(upstreams.map { it.getHead() }).apply {
this.start() this.start()
} }
val lagObserver = BitcoinHeadLagObserver(newHead, upstreams) val lagObserver = BitcoinHeadLagObserver(newHead, upstreams)
this.lagObserver = lagObserver this.lagObserver = lagObserver
lagObserver.start()
newHead newHead
} }
onHeadUpdated(head) onHeadUpdated(head)

View File

@@ -95,16 +95,19 @@ open class EthereumMultistream(
val head = if (upstreams.size == 1) { val head = if (upstreams.size == 1) {
val upstream = upstreams.first() val upstream = upstreams.first()
upstream.setLag(0) upstream.setLag(0)
upstream.getHead() upstream.getHead().apply {
if (this is Lifecycle) {
this.start()
}
}
} else { } else {
val heads = upstreams.map { it.getHead() } val heads = upstreams.map { it.getHead() }
val newHead = MergedHead(heads).apply { val newHead = MergedHead(heads).apply {
this.start() this.start()
} }
val lagObserver = EthereumHeadLagObserver(newHead, upstreams as Collection<Upstream>).apply { val lagObserver = EthereumHeadLagObserver(newHead, upstreams as Collection<Upstream>)
this.start()
}
this.lagObserver = lagObserver this.lagObserver = lagObserver
lagObserver.start()
newHead newHead
} }
onHeadUpdated(head) onHeadUpdated(head)

View File

@@ -43,6 +43,7 @@ class EthereumRpcHead(
private var refreshSubscription: Disposable? = null private var refreshSubscription: Disposable? = null
override fun start() { override fun start() {
refreshSubscription?.dispose()
val base = Flux.interval(interval) val base = Flux.interval(interval)
.publishOn(scheduler) .publishOn(scheduler)
.flatMap { .flatMap {