From db97b656cfe42ab5b61158be3c46a716f5a7b2d1 Mon Sep 17 00:00:00 2001 From: Maksim Fomenkov Date: Mon, 19 Dec 2022 11:41:05 +0300 Subject: [PATCH] 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