From 1f611c826827ce4a84d1223fe72414af72e0c0af Mon Sep 17 00:00:00 2001 From: Maksim Fomenkov Date: Mon, 19 Dec 2022 11:12:12 +0300 Subject: [PATCH 1/7] subscribe node status --- emerald-grpc | 2 +- .../emeraldpay/dshackle/rpc/BlockchainRpc.kt | 5 + .../dshackle/rpc/SubscribeNodeStatus.kt | 113 ++++++++++++++++++ .../upstream/CurrentMultistreamHolder.kt | 3 + .../dshackle/upstream/Multistream.kt | 8 ++ .../dshackle/upstream/MultistreamHolder.kt | 2 + 6 files changed, 132 insertions(+), 1 deletion(-) create mode 100644 src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt diff --git a/emerald-grpc b/emerald-grpc index c4c9ea44..6b012f3c 160000 --- a/emerald-grpc +++ b/emerald-grpc @@ -1 +1 @@ -Subproject commit c4c9ea44d67f867852aef399aa81c13908ecbfcf +Subproject commit 6b012f3c89f3b311113b3f7a71cdba3f4fee0a3d diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/BlockchainRpc.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/BlockchainRpc.kt index d3321b46..ac3143e9 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/BlockchainRpc.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/BlockchainRpc.kt @@ -47,6 +47,7 @@ class BlockchainRpc( private val describe: Describe, private val subscribeStatus: SubscribeStatus, private val estimateFee: EstimateFee, + private val subscribeNodeStatus: SubscribeNodeStatus, @Qualifier("rpcScheduler") private val scheduler: Scheduler ) : ReactorBlockchainGrpc.BlockchainImplBase() { @@ -220,6 +221,10 @@ class BlockchainRpc( .doOnError { failMetric.increment() } } + override fun subscribeNodeStatus(request: Mono): Flux { + return subscribeNodeStatus.subscribe(request).subscribeOn(scheduler).doOnError { failMetric.increment() } + } + class RequestMetrics(val chain: Chain) { val nativeCallMetric = Counter.builder("request.grpc.request") .tag("type", "nativeCall") diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt new file mode 100644 index 00000000..1ab303ad --- /dev/null +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt @@ -0,0 +1,113 @@ +package io.emeraldpay.dshackle.rpc + +import io.emeraldpay.api.proto.BlockchainOuterClass +import io.emeraldpay.api.proto.BlockchainOuterClass.NodeDescription +import io.emeraldpay.api.proto.BlockchainOuterClass.NodeStatus +import io.emeraldpay.api.proto.BlockchainOuterClass.NodeStatusResponse +import io.emeraldpay.api.proto.BlockchainOuterClass.SubscribeNodeStatusRequest +import io.emeraldpay.api.proto.Common +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.springframework.stereotype.Service +import reactor.core.publisher.Flux +import reactor.core.publisher.Mono +import reactor.core.publisher.SignalType +import reactor.core.publisher.Sinks +import java.time.Duration +import java.util.concurrent.ConcurrentHashMap + +@Service +class SubscribeNodeStatus( + private val multistreams: CurrentMultistreamHolder +) { + + fun subscribe(req: Mono): Flux = + req.flatMapMany { + val knownUpstreams = ConcurrentHashMap() + // send known upstreams details immediately + val descriptions = Flux.fromIterable(multistreams.all().flatMap { multiStream -> + multiStream.getAll().map { up -> + NodeStatusResponse.newBuilder() + .setNodeId(up.getId()) + .setDescription(buildDescription(multiStream.chain, up)) + .setStatus(buildStatus(up.getStatus(), up.getHead().getCurrentHeight())) + .build() + } + }) + + //subscribe on head/status updates for known upstreams + val upstreamUpdates = Flux.merge(multistreams.all().flatMap { ms -> + ms.getAll().map { up -> + knownUpstreams[up.getId()] = true + subscribeUpstreamUpdates(ms.chain, up) + } + }) + + //subscribe on head/status updates for just added upstreams + val muliStreamUpdates = Flux.merge(multistreams.all().map { ms -> + ms.subscribeAddedUpstreams() + .distinctUntilChanged { + it.getId() + } + .filter { knownUpstreams[it.getId()] != true }.flatMap { + knownUpstreams[it.getId()] = true + Flux.concat( + Mono.just( + NodeStatusResponse.newBuilder() + .setNodeId(it.getId()) + .setDescription(buildDescription(ms.chain, it)) + .setStatus(buildStatus(it.getStatus(), it.getHead().getCurrentHeight())) + .build() + ), subscribeUpstreamUpdates(ms.chain, it) + ) + } + }) + + Flux.concat(descriptions, Flux.merge(upstreamUpdates, muliStreamUpdates)) + } + + private fun subscribeUpstreamUpdates(chain: Chain, upstream: Upstream): Flux { + val r = Sinks.many().multicast().directBestEffort() + val heads = Mono.just(upstream).repeatWhen { r.asFlux() }.sample(Duration.ofSeconds(10)).flatMap { up -> + up.getHead().getFlux().map { block -> + NodeStatusResponse.newBuilder() + .setNodeId(up.getId()) + .setStatus(buildStatus(up.getStatus(), block.height)) + .build() + }.doFinally { + if (it == SignalType.ON_COMPLETE) { + r.tryEmitNext(true) + } + } + } + + val statuses = upstream.observeStatus().map { + NodeStatusResponse.newBuilder() + .setNodeId(upstream.getId()) + .setStatus(buildStatus(it, upstream.getHead().getCurrentHeight())) + .setDescription(buildDescription(chain, upstream)) + .build() + } + return Flux.merge(heads.distinctUntilChanged(), statuses.distinctUntilChanged()) + } + + private fun buildDescription(chain: Chain, up: Upstream): NodeDescription.Builder = + NodeDescription.newBuilder() + .setChain(Common.ChainRef.forNumber(chain.id)) + .addAllLabels(up.getLabels().flatMap { labels -> + labels.map { + BlockchainOuterClass.Label.newBuilder() + .setName(it.key) + .setValue(it.value) + .build() + } + }) + .addAllSupportedMethods(up.getMethods().getSupportedMethods()) + + private fun buildStatus(status: UpstreamAvailability, height: Long?): NodeStatus.Builder = + NodeStatus.newBuilder() + .setAvailability(BlockchainOuterClass.AvailabilityEnum.forNumber(status.grpcId)) + .setCurrentHeight(height ?: 0) +} diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt index 165c20fe..842639a3 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentMultistreamHolder.kt @@ -45,6 +45,9 @@ open class CurrentMultistreamHolder( return chainMapping.getValue(chain).isAvailable() } + override fun all(): List = + chainMapping.values.toList() + @PreDestroy fun shutdown() { log.info("Closing upstream connections...") diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt index 81741a2f..aa7d5a9d 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt @@ -39,6 +39,7 @@ import org.springframework.core.annotation.Order import reactor.core.Disposable import reactor.core.publisher.Flux import reactor.core.publisher.Mono +import reactor.core.publisher.Sinks import java.time.Duration import java.time.Instant import java.util.Locale @@ -76,6 +77,9 @@ abstract class Multistream( private var capabilities: Set = emptySet() private val removed: MutableMap = HashMap() private val meters: MutableMap = HashMap() + private val addedUpstreams = Sinks.many() + .multicast() + .directBestEffort() init { UpstreamAvailability.values().forEach { status -> @@ -350,6 +354,7 @@ abstract class Multistream( if (!started) { start() } + addedUpstreams.tryEmitNext(event.upstream) log.info("Upstream ${event.upstream.getId()} with chain $chain has been added") } } @@ -364,6 +369,9 @@ abstract class Multistream( return upstreams.any { matcher.matches(it) } } + fun subscribeAddedUpstreams(): Flux = + addedUpstreams.asFlux() + // -------------------------------------------------------------------------------------------------------- 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 da409d1c..84d47144 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/MultistreamHolder.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/MultistreamHolder.kt @@ -25,4 +25,6 @@ interface MultistreamHolder { fun getUpstream(chain: Chain): Multistream fun getAvailable(): List fun isAvailable(chain: Chain): Boolean + + fun all(): List } From db97b656cfe42ab5b61158be3c46a716f5a7b2d1 Mon Sep 17 00:00:00 2001 From: Maksim Fomenkov Date: Mon, 19 Dec 2022 11:41:05 +0300 Subject: [PATCH 2/7] fix style and tests --- .../dshackle/rpc/SubscribeNodeStatus.kt | 96 +++++++++++-------- .../test/MultistreamHolderMock.groovy | 5 + 2 files changed, 59 insertions(+), 42 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt index 1ab303ad..e47037ad 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt @@ -27,43 +27,53 @@ class SubscribeNodeStatus( req.flatMapMany { val knownUpstreams = ConcurrentHashMap() // send known upstreams details immediately - val descriptions = Flux.fromIterable(multistreams.all().flatMap { multiStream -> - multiStream.getAll().map { up -> - NodeStatusResponse.newBuilder() - .setNodeId(up.getId()) - .setDescription(buildDescription(multiStream.chain, up)) - .setStatus(buildStatus(up.getStatus(), up.getHead().getCurrentHeight())) - .build() - } - }) - - //subscribe on head/status updates for known upstreams - val upstreamUpdates = Flux.merge(multistreams.all().flatMap { ms -> - ms.getAll().map { up -> - knownUpstreams[up.getId()] = true - subscribeUpstreamUpdates(ms.chain, up) - } - }) - - //subscribe on head/status updates for just added upstreams - val muliStreamUpdates = Flux.merge(multistreams.all().map { ms -> - ms.subscribeAddedUpstreams() - .distinctUntilChanged { - it.getId() + val descriptions = Flux.fromIterable( + multistreams.all() + .flatMap { multiStream -> + multiStream.getAll().map { up -> + NodeStatusResponse.newBuilder() + .setNodeId(up.getId()) + .setDescription(buildDescription(multiStream.chain, up)) + .setStatus(buildStatus(up.getStatus(), up.getHead().getCurrentHeight())) + .build() + } } - .filter { knownUpstreams[it.getId()] != true }.flatMap { - knownUpstreams[it.getId()] = true - Flux.concat( - Mono.just( - NodeStatusResponse.newBuilder() - .setNodeId(it.getId()) - .setDescription(buildDescription(ms.chain, it)) - .setStatus(buildStatus(it.getStatus(), it.getHead().getCurrentHeight())) - .build() - ), subscribeUpstreamUpdates(ms.chain, it) - ) + ) + + // subscribe on head/status updates for known upstreams + val upstreamUpdates = Flux.merge( + multistreams.all() + .flatMap { ms -> + ms.getAll().map { up -> + knownUpstreams[up.getId()] = true + subscribeUpstreamUpdates(ms.chain, up) + } } - }) + ) + + // subscribe on head/status updates for just added upstreams + val muliStreamUpdates = Flux.merge( + multistreams.all() + .map { ms -> + ms.subscribeAddedUpstreams() + .distinctUntilChanged { + it.getId() + } + .filter { knownUpstreams[it.getId()] != true }.flatMap { + knownUpstreams[it.getId()] = true + Flux.concat( + Mono.just( + NodeStatusResponse.newBuilder() + .setNodeId(it.getId()) + .setDescription(buildDescription(ms.chain, it)) + .setStatus(buildStatus(it.getStatus(), it.getHead().getCurrentHeight())) + .build() + ), + subscribeUpstreamUpdates(ms.chain, it) + ) + } + } + ) Flux.concat(descriptions, Flux.merge(upstreamUpdates, muliStreamUpdates)) } @@ -96,14 +106,16 @@ class SubscribeNodeStatus( private fun buildDescription(chain: Chain, up: Upstream): NodeDescription.Builder = NodeDescription.newBuilder() .setChain(Common.ChainRef.forNumber(chain.id)) - .addAllLabels(up.getLabels().flatMap { labels -> - labels.map { - BlockchainOuterClass.Label.newBuilder() - .setName(it.key) - .setValue(it.value) - .build() + .addAllLabels( + up.getLabels().flatMap { labels -> + labels.map { + BlockchainOuterClass.Label.newBuilder() + .setName(it.key) + .setValue(it.value) + .build() + } } - }) + ) .addAllSupportedMethods(up.getMethods().getSupportedMethods()) private fun buildStatus(status: UpstreamAvailability, height: Long?): NodeStatus.Builder = diff --git a/src/test/groovy/io/emeraldpay/dshackle/test/MultistreamHolderMock.groovy b/src/test/groovy/io/emeraldpay/dshackle/test/MultistreamHolderMock.groovy index 777adc05..56942007 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/test/MultistreamHolderMock.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/test/MultistreamHolderMock.groovy @@ -86,6 +86,11 @@ class MultistreamHolderMock implements MultistreamHolder { return upstreams.containsKey(chain) } + @Override + List all() { + return upstreams.values().toList() + } + static class EthereumMultistreamMock extends EthereumPosMultiStream { EthereumCachingReader customReader = null From 8f0c9d11873829968c42eaac0bb18ed966650448 Mon Sep 17 00:00:00 2001 From: Maksim Fomenkov Date: Tue, 20 Dec 2022 12:01:34 +0300 Subject: [PATCH 3/7] added sampling --- .../dshackle/rpc/SubscribeNodeStatus.kt | 32 ++++++++++++------- 1 file changed, 21 insertions(+), 11 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt index e47037ad..633c90f3 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt @@ -23,9 +23,14 @@ class SubscribeNodeStatus( private val multistreams: CurrentMultistreamHolder ) { + companion object { + private val RETRY_TIMEOUT = Duration.ofSeconds(10) + } + fun subscribe(req: Mono): Flux = - req.flatMapMany { + req.flatMapMany { request -> val knownUpstreams = ConcurrentHashMap() + val duration = Duration.ofMillis(request.timespan) // send known upstreams details immediately val descriptions = Flux.fromIterable( multistreams.all() @@ -46,7 +51,7 @@ class SubscribeNodeStatus( .flatMap { ms -> ms.getAll().map { up -> knownUpstreams[up.getId()] = true - subscribeUpstreamUpdates(ms.chain, up) + subscribeUpstreamUpdates(ms.chain, up, duration) } } ) @@ -69,7 +74,7 @@ class SubscribeNodeStatus( .setStatus(buildStatus(it.getStatus(), it.getHead().getCurrentHeight())) .build() ), - subscribeUpstreamUpdates(ms.chain, it) + subscribeUpstreamUpdates(ms.chain, it, duration) ) } } @@ -78,29 +83,34 @@ class SubscribeNodeStatus( Flux.concat(descriptions, Flux.merge(upstreamUpdates, muliStreamUpdates)) } - private fun subscribeUpstreamUpdates(chain: Chain, upstream: Upstream): Flux { - val r = Sinks.many().multicast().directBestEffort() - val heads = Mono.just(upstream).repeatWhen { r.asFlux() }.sample(Duration.ofSeconds(10)).flatMap { up -> + private fun subscribeUpstreamUpdates( + chain: Chain, + upstream: Upstream, + timespan: Duration + ): Flux { + val retry = Sinks.many().multicast().directBestEffort() + val heads = Mono.just(upstream).repeatWhen { retry.asFlux() }.sample(RETRY_TIMEOUT).flatMap { up -> up.getHead().getFlux().map { block -> NodeStatusResponse.newBuilder() .setNodeId(up.getId()) .setStatus(buildStatus(up.getStatus(), block.height)) .build() }.doFinally { + // retry when subscribed head stopped if (it == SignalType.ON_COMPLETE) { - r.tryEmitNext(true) + retry.tryEmitNext(true) } } - } + }.sample(timespan) - val statuses = upstream.observeStatus().map { + val statuses = upstream.observeStatus().distinctUntilChanged().map { NodeStatusResponse.newBuilder() .setNodeId(upstream.getId()) .setStatus(buildStatus(it, upstream.getHead().getCurrentHeight())) .setDescription(buildDescription(chain, upstream)) .build() - } - return Flux.merge(heads.distinctUntilChanged(), statuses.distinctUntilChanged()) + }.sample(timespan) + return Flux.merge(heads, statuses) } private fun buildDescription(chain: Chain, up: Upstream): NodeDescription.Builder = From 5894501f0c4ee7854147339fbca32fb43010e6d8 Mon Sep 17 00:00:00 2001 From: Maksim Fomenkov Date: Tue, 20 Dec 2022 15:00:23 +0300 Subject: [PATCH 4/7] fixed labels grouping --- emerald-grpc | 2 +- .../dshackle/rpc/SubscribeNodeStatus.kt | 20 +++++++++++-------- 2 files changed, 13 insertions(+), 9 deletions(-) diff --git a/emerald-grpc b/emerald-grpc index 6b012f3c..812e9e6d 160000 --- a/emerald-grpc +++ b/emerald-grpc @@ -1 +1 @@ -Subproject commit 6b012f3c89f3b311113b3f7a71cdba3f4fee0a3d +Subproject commit 812e9e6ddbbd0fb3f21936b26b82fc76475bce1c diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt index 633c90f3..1f814661 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt @@ -116,14 +116,18 @@ class SubscribeNodeStatus( private fun buildDescription(chain: Chain, up: Upstream): NodeDescription.Builder = NodeDescription.newBuilder() .setChain(Common.ChainRef.forNumber(chain.id)) - .addAllLabels( - up.getLabels().flatMap { labels -> - labels.map { - BlockchainOuterClass.Label.newBuilder() - .setName(it.key) - .setValue(it.value) - .build() - } + .addAllNodeLabels( + up.getLabels().map { nodeLabels -> + BlockchainOuterClass.NodeLabels.newBuilder() + .addAllLabels( + nodeLabels.map { + BlockchainOuterClass.Label.newBuilder() + .setName(it.key) + .setValue(it.value) + .build() + } + ) + .build() } ) .addAllSupportedMethods(up.getMethods().getSupportedMethods()) From f6c6c8d4ef41b78dec94aa97b6b372e73b21bfb1 Mon Sep 17 00:00:00 2001 From: Maksim Fomenkov Date: Tue, 20 Dec 2022 15:36:20 +0300 Subject: [PATCH 5/7] fixed review comments --- .../io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt index 1f814661..77a8c296 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt @@ -57,14 +57,14 @@ class SubscribeNodeStatus( ) // subscribe on head/status updates for just added upstreams - val muliStreamUpdates = Flux.merge( + val multiStreamUpdates = Flux.merge( multistreams.all() .map { ms -> ms.subscribeAddedUpstreams() .distinctUntilChanged { it.getId() } - .filter { knownUpstreams[it.getId()] != true }.flatMap { + .filter { knownUpstreams.getOrDefault(it.getId(), false) }.flatMap { knownUpstreams[it.getId()] = true Flux.concat( Mono.just( @@ -80,7 +80,7 @@ class SubscribeNodeStatus( } ) - Flux.concat(descriptions, Flux.merge(upstreamUpdates, muliStreamUpdates)) + Flux.concat(descriptions, Flux.merge(upstreamUpdates, multiStreamUpdates)) } private fun subscribeUpstreamUpdates( From a1c234d3bfd2e7a11c5e08bac73269c53ab625f5 Mon Sep 17 00:00:00 2001 From: Maksim Fomenkov Date: Tue, 20 Dec 2022 16:33:25 +0300 Subject: [PATCH 6/7] fixed re-subscription --- .../dshackle/rpc/SubscribeNodeStatus.kt | 70 ++++++++++++------- 1 file changed, 46 insertions(+), 24 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt index 77a8c296..9967b156 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt @@ -17,6 +17,7 @@ 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( @@ -51,7 +52,7 @@ class SubscribeNodeStatus( .flatMap { ms -> ms.getAll().map { up -> knownUpstreams[up.getId()] = true - subscribeUpstreamUpdates(ms.chain, up, duration) + subscribeUpstreamUpdates(ms.chain, up, duration) { r -> knownUpstreams.remove(r) } } } ) @@ -64,7 +65,10 @@ class SubscribeNodeStatus( .distinctUntilChanged { it.getId() } - .filter { knownUpstreams.getOrDefault(it.getId(), false) }.flatMap { + .filter { + !knownUpstreams.getOrDefault(it.getId(), false) + } + .flatMap { knownUpstreams[it.getId()] = true Flux.concat( Mono.just( @@ -74,7 +78,7 @@ class SubscribeNodeStatus( .setStatus(buildStatus(it.getStatus(), it.getHead().getCurrentHeight())) .build() ), - subscribeUpstreamUpdates(ms.chain, it, duration) + subscribeUpstreamUpdates(ms.chain, it, duration) { r -> knownUpstreams.remove(r) } ) } } @@ -86,30 +90,48 @@ class SubscribeNodeStatus( private fun subscribeUpstreamUpdates( chain: Chain, upstream: Upstream, - timespan: Duration + timespan: Duration, + onUnavailable: Consumer ): Flux { val retry = Sinks.many().multicast().directBestEffort() - val heads = Mono.just(upstream).repeatWhen { retry.asFlux() }.sample(RETRY_TIMEOUT).flatMap { up -> - up.getHead().getFlux().map { block -> - NodeStatusResponse.newBuilder() - .setNodeId(up.getId()) - .setStatus(buildStatus(up.getStatus(), block.height)) - .build() - }.doFinally { - // retry when subscribed head stopped - if (it == SignalType.ON_COMPLETE) { - retry.tryEmitNext(true) - } - } - }.sample(timespan) + val cancel = Sinks.many().multicast().directBestEffort() + + val heads = Mono.just(upstream) + .repeatWhen { + retry.asFlux() + } + .sample(RETRY_TIMEOUT) + .takeUntilOther(cancel.asFlux()).flatMap { up -> + up.getHead().getFlux() + .takeUntilOther(cancel.asFlux()) + .map { block -> + NodeStatusResponse.newBuilder() + .setNodeId(up.getId()) + .setStatus(buildStatus(up.getStatus(), block.height)) + .build() + }.doFinally { + // retry when subscribed head stopped + if (it == SignalType.ON_COMPLETE) { + retry.tryEmitNext(true) + } + } + }.sample(timespan) + + val statuses = upstream.observeStatus() + .distinctUntilChanged() + .takeUntil{it == UpstreamAvailability.UNAVAILABLE} + .map { + if (it == UpstreamAvailability.UNAVAILABLE) { + onUnavailable.accept(upstream.getId()) + cancel.tryEmitNext(true) + } + NodeStatusResponse.newBuilder() + .setNodeId(upstream.getId()) + .setStatus(buildStatus(it, upstream.getHead().getCurrentHeight())) + .setDescription(buildDescription(chain, upstream)) + .build() + }.sample(timespan) - val statuses = upstream.observeStatus().distinctUntilChanged().map { - NodeStatusResponse.newBuilder() - .setNodeId(upstream.getId()) - .setStatus(buildStatus(it, upstream.getHead().getCurrentHeight())) - .setDescription(buildDescription(chain, upstream)) - .build() - }.sample(timespan) return Flux.merge(heads, statuses) } From 214cb4fd8e7086b93a5aa8c1c5a462f08a26755a Mon Sep 17 00:00:00 2001 From: Maksim Fomenkov Date: Tue, 20 Dec 2022 16:39:26 +0300 Subject: [PATCH 7/7] fixed code style --- .../kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt index 9967b156..93013e66 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt @@ -119,10 +119,11 @@ class SubscribeNodeStatus( val statuses = upstream.observeStatus() .distinctUntilChanged() - .takeUntil{it == UpstreamAvailability.UNAVAILABLE} + .takeUntil { it == UpstreamAvailability.UNAVAILABLE } .map { if (it == UpstreamAvailability.UNAVAILABLE) { onUnavailable.accept(upstream.getId()) + // cancel head subscription & reconnections when upstream becomes unavailable cancel.tryEmitNext(true) } NodeStatusResponse.newBuilder()