From 4aa54cb9070f61c0fd56dd22ef7fac4c352cd2a5 Mon Sep 17 00:00:00 2001 From: Termina1 Date: Tue, 10 Jan 2023 13:33:35 +0200 Subject: [PATCH] subscribe node status - remove head subscription and remove sampling, afraid of missing updates --- .../dshackle/rpc/SubscribeNodeStatus.kt | 34 +++---------------- 1 file changed, 5 insertions(+), 29 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt index 2ad44345..d5059ead 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt @@ -30,9 +30,8 @@ class SubscribeNodeStatus( } fun subscribe(req: Mono): Flux = - req.flatMapMany { request -> + req.flatMapMany { val knownUpstreams = ConcurrentHashMap>() - val duration = Duration.ofMillis(request.timespan) // send known upstreams details immediately val descriptions = Flux.fromIterable( multistreams.all() @@ -53,7 +52,7 @@ class SubscribeNodeStatus( .flatMap { ms -> ms.getAll().map { up -> knownUpstreams[up.getId()] = Sinks.many().multicast().directBestEffort() - subscribeUpstreamUpdates(ms.chain, up, duration, knownUpstreams[up.getId()]!!) + subscribeUpstreamUpdates(ms.chain, up, knownUpstreams[up.getId()]!!) } } ) @@ -100,7 +99,7 @@ class SubscribeNodeStatus( .setStatus(buildStatus(it.getStatus(), it.getHead().getCurrentHeight())) .build() ), - subscribeUpstreamUpdates(ms.chain, it, duration, knownUpstreams[it.getId()]!!) + subscribeUpstreamUpdates(ms.chain, it, knownUpstreams[it.getId()]!!) ) } } @@ -112,31 +111,8 @@ class SubscribeNodeStatus( private fun subscribeUpstreamUpdates( chain: Chain, upstream: Upstream, - timespan: Duration, cancel: Sinks.Many ): Flux { - val retry = 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() .takeUntilOther(cancel.asFlux()) @@ -146,9 +122,9 @@ class SubscribeNodeStatus( .setStatus(buildStatus(it, upstream.getHead().getCurrentHeight())) .setDescription(buildDescription(chain, upstream)) .build() - }.sample(timespan) + } - return Flux.merge(heads, statuses) + return statuses } private fun buildDescription(chain: Chain, up: Upstream): NodeDescription.Builder =