From f298b99251db1a3a27c6c61f57ec0258d4c5e206 Mon Sep 17 00:00:00 2001 From: Termina1 Date: Thu, 12 Jan 2023 13:05:47 +0200 Subject: [PATCH] subscribe status send current upstream state on subscription --- .../dshackle/rpc/SubscribeNodeStatus.kt | 24 ++++++------------- 1 file changed, 7 insertions(+), 17 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt index 85f5345d..744be837 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt @@ -29,20 +29,6 @@ class SubscribeNodeStatus( 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() @@ -102,7 +88,7 @@ class SubscribeNodeStatus( } ) - Flux.concat(descriptions, Flux.merge(upstreamUpdates, multiStreamUpdates, removals)) + Flux.merge(upstreamUpdates, multiStreamUpdates, removals) } private fun subscribeUpstreamUpdates( @@ -120,8 +106,12 @@ class SubscribeNodeStatus( .setDescription(buildDescription(chain, upstream)) .build() } - - return statuses + val currentState = Mono.just(NodeStatusResponse.newBuilder() + .setNodeId(upstream.getId()) + .setDescription(buildDescription(chain, upstream)) + .setStatus(buildStatus(upstream.getStatus(), upstream.getHead().getCurrentHeight())) + .build()) + return Flux.concat(currentState, statuses) } private fun buildDescription(chain: Chain, up: Upstream): NodeDescription.Builder =