From f6c6c8d4ef41b78dec94aa97b6b372e73b21bfb1 Mon Sep 17 00:00:00 2001 From: Maksim Fomenkov Date: Tue, 20 Dec 2022 15:36:20 +0300 Subject: [PATCH] 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(