diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt index 93013e66..2ad44345 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt @@ -10,6 +10,7 @@ import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.upstream.CurrentMultistreamHolder import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.UpstreamAvailability +import org.slf4j.LoggerFactory import org.springframework.stereotype.Service import reactor.core.publisher.Flux import reactor.core.publisher.Mono @@ -17,7 +18,6 @@ import reactor.core.publisher.SignalType import reactor.core.publisher.Sinks import java.time.Duration import java.util.concurrent.ConcurrentHashMap -import java.util.function.Consumer @Service class SubscribeNodeStatus( @@ -26,11 +26,12 @@ class SubscribeNodeStatus( companion object { private val RETRY_TIMEOUT = Duration.ofSeconds(10) + private val log = LoggerFactory.getLogger(SubscribeNodeStatus::class.java) } fun subscribe(req: Mono): Flux = req.flatMapMany { request -> - val knownUpstreams = ConcurrentHashMap() + val knownUpstreams = ConcurrentHashMap>() val duration = Duration.ofMillis(request.timespan) // send known upstreams details immediately val descriptions = Flux.fromIterable( @@ -51,8 +52,29 @@ class SubscribeNodeStatus( multistreams.all() .flatMap { ms -> ms.getAll().map { up -> - knownUpstreams[up.getId()] = true - subscribeUpstreamUpdates(ms.chain, up, duration) { r -> knownUpstreams.remove(r) } + knownUpstreams[up.getId()] = Sinks.many().multicast().directBestEffort() + subscribeUpstreamUpdates(ms.chain, up, duration, knownUpstreams[up.getId()]!!) + } + } + ) + + // stop removed upstreams update fluxes + val removals = Flux.merge( + multistreams.all() + .map { ms -> + ms.subscribeRemovedUpstreams().mapNotNull { up -> + knownUpstreams[up.getId()]?.let { + val result = it.tryEmitNext(true) + if (result.isFailure) { + log.warn("Unable to emit event about removal of an upstream - $result") + } + knownUpstreams.remove(up.getId()) + NodeStatusResponse.newBuilder() + .setNodeId(up.getId()) + .setDescription(buildDescription(ms.chain, up)) + .setStatus(buildStatus(UpstreamAvailability.UNAVAILABLE, up.getHead().getCurrentHeight())) + .build() + } } } ) @@ -66,10 +88,10 @@ class SubscribeNodeStatus( it.getId() } .filter { - !knownUpstreams.getOrDefault(it.getId(), false) + !knownUpstreams.contains(it.getId()) } .flatMap { - knownUpstreams[it.getId()] = true + knownUpstreams[it.getId()] = Sinks.many().multicast().directBestEffort() Flux.concat( Mono.just( NodeStatusResponse.newBuilder() @@ -78,24 +100,22 @@ class SubscribeNodeStatus( .setStatus(buildStatus(it.getStatus(), it.getHead().getCurrentHeight())) .build() ), - subscribeUpstreamUpdates(ms.chain, it, duration) { r -> knownUpstreams.remove(r) } + subscribeUpstreamUpdates(ms.chain, it, duration, knownUpstreams[it.getId()]!!) ) } } ) - Flux.concat(descriptions, Flux.merge(upstreamUpdates, multiStreamUpdates)) + Flux.concat(descriptions, Flux.merge(upstreamUpdates, multiStreamUpdates, removals)) } private fun subscribeUpstreamUpdates( chain: Chain, upstream: Upstream, timespan: Duration, - onUnavailable: Consumer + cancel: Sinks.Many ): Flux { val retry = Sinks.many().multicast().directBestEffort() - val cancel = Sinks.many().multicast().directBestEffort() - val heads = Mono.just(upstream) .repeatWhen { retry.asFlux() @@ -119,13 +139,8 @@ class SubscribeNodeStatus( val statuses = upstream.observeStatus() .distinctUntilChanged() - .takeUntil { it == UpstreamAvailability.UNAVAILABLE } + .takeUntilOther(cancel.asFlux()) .map { - if (it == UpstreamAvailability.UNAVAILABLE) { - onUnavailable.accept(upstream.getId()) - // cancel head subscription & reconnections when upstream becomes unavailable - cancel.tryEmitNext(true) - } NodeStatusResponse.newBuilder() .setNodeId(upstream.getId()) .setStatus(buildStatus(it, upstream.getHead().getCurrentHeight())) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt index 6305c7c9..c8e460af 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt @@ -78,6 +78,9 @@ abstract class Multistream( private val addedUpstreams = Sinks.many() .multicast() .directBestEffort() + private val removedUpstreams = Sinks.many() + .multicast() + .directBestEffort() init { UpstreamAvailability.values().forEach { status -> @@ -342,6 +345,7 @@ abstract class Multistream( eventLock.withLock { if (event.type == UpstreamChangeEvent.ChangeType.REMOVED) { removeUpstream(event.upstream.getId()).takeIf { it }?.let { + removedUpstreams.tryEmitNext(event.upstream) log.warn("Upstream ${event.upstream.getId()} with chain $chain has been removed") } } else { @@ -370,6 +374,9 @@ abstract class Multistream( fun subscribeAddedUpstreams(): Flux = addedUpstreams.asFlux() + fun subscribeRemovedUpstreams(): Flux = + removedUpstreams.asFlux() + // -------------------------------------------------------------------------------------------------------- class UpstreamStatus(val upstream: Upstream, val status: UpstreamAvailability, val ts: Instant = Instant.now())