From 96109d40131b7c8c42e1ba81e2007fef84d87121 Mon Sep 17 00:00:00 2001 From: Igor Artamonov Date: Sat, 31 Jul 2021 22:06:02 -0400 Subject: [PATCH] problem: gives no status updates if chain is not configured --- .../io/emeraldpay/dshackle/rpc/Describe.kt | 2 +- .../dshackle/rpc/SubscribeStatus.kt | 47 +++++++++------ .../dshackle/upstream/Multistream.kt | 2 +- .../dshackle/rpc/SubscribeStatusSpec.groovy | 60 +++++++++++++++++++ 4 files changed, 92 insertions(+), 19 deletions(-) create mode 100644 src/test/groovy/io/emeraldpay/dshackle/rpc/SubscribeStatusSpec.groovy diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/Describe.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/Describe.kt index 3e995499..7163853d 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/Describe.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/Describe.kt @@ -35,7 +35,7 @@ class Describe( val resp = BlockchainOuterClass.DescribeResponse.newBuilder() multistreamHolder.getAvailable().forEach { chain -> multistreamHolder.getUpstream(chain)?.let { chainUpstreams -> - val status = subscribeStatus.chainStatus(chain, chainUpstreams.getAll()) + val status = subscribeStatus.chainStatus(chain, chainUpstreams.getStatus(), chainUpstreams) val targets = chainUpstreams.getMethods().getSupportedMethods() val capabilities: MutableSet = mutableSetOf() val chainDescription = BlockchainOuterClass.DescribeChain.newBuilder() diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeStatus.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeStatus.kt index f22b2b17..dbe9aca1 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeStatus.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeStatus.kt @@ -31,28 +31,41 @@ class SubscribeStatus( ) { fun subscribeStatus(requestMono: Mono): Flux { - return requestMono.flatMapMany { - val ups = multistreamHolder.getAvailable().mapNotNull { chain -> - val chainUpstream = multistreamHolder.getUpstream(chain) - chainUpstream?.observeStatus()?.map { avail -> - ChainSubscription(chain, chainUpstream, avail) + return requestMono.flatMapMany { req -> + // check status for all requested chains + val all = req.chainsList.map { + val chain = Chain.byId(it.number) + val up = multistreamHolder.getUpstream(chain) + if (up == null) { + // when the chain is not configured return just UNAVAILABLE + Mono.just(chainUnavailable(chain)).flux() + } else { + // when configured subscribe to its updates + up.observeStatus().map { availability -> + chainStatus(chain, availability, up) + } } } - - Flux.merge(ups) - .map { - chainStatus(it.chain, it.up.getAll()) - } + Flux.merge(all) } } - fun chainStatus(chain: Chain, ups: List): BlockchainOuterClass.ChainStatus { - val available = ups.map { u -> - u.getStatus() - }.min() ?: UpstreamAvailability.UNAVAILABLE - val quorum = ups.filter { - it.getStatus() > UpstreamAvailability.UNAVAILABLE - }.count() + fun chainUnavailable(chain: Chain): BlockchainOuterClass.ChainStatus { + return BlockchainOuterClass.ChainStatus.newBuilder() + .setAvailability(BlockchainOuterClass.AvailabilityEnum.AVAIL_UNAVAILABLE) + .setChain(Common.ChainRef.forNumber(chain.id)) + .setQuorum(0) + .build() + } + + fun chainStatus(chain: Chain, available: UpstreamAvailability, ups: Multistream): BlockchainOuterClass.ChainStatus { + val quorum = if (available != UpstreamAvailability.UNAVAILABLE) { + ups.getAll().count { + it.getStatus() > UpstreamAvailability.UNAVAILABLE + } + } else { + 0 + } return BlockchainOuterClass.ChainStatus.newBuilder() .setAvailability(BlockchainOuterClass.AvailabilityEnum.forNumber(available.grpcId)) .setChain(Common.ChainRef.forNumber(chain.id)) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt index 3122ada9..2904e15e 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt @@ -98,7 +98,7 @@ abstract class Multistream( /** * Get list of all underlying upstreams */ - fun getAll(): List { + open fun getAll(): List { return upstreams } diff --git a/src/test/groovy/io/emeraldpay/dshackle/rpc/SubscribeStatusSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/rpc/SubscribeStatusSpec.groovy new file mode 100644 index 00000000..cb1b952d --- /dev/null +++ b/src/test/groovy/io/emeraldpay/dshackle/rpc/SubscribeStatusSpec.groovy @@ -0,0 +1,60 @@ +package io.emeraldpay.dshackle.rpc + +import io.emeraldpay.api.proto.BlockchainOuterClass +import io.emeraldpay.api.proto.Common +import io.emeraldpay.dshackle.upstream.Multistream +import io.emeraldpay.dshackle.upstream.MultistreamHolder +import io.emeraldpay.dshackle.upstream.Upstream +import io.emeraldpay.dshackle.upstream.UpstreamAvailability +import io.emeraldpay.grpc.Chain +import reactor.core.publisher.Mono +import reactor.test.StepVerifier +import spock.lang.Specification + +import java.time.Duration + +class SubscribeStatusSpec extends Specification { + + def "gives UNAVAIL if non-configured chain is requested"() { + setup: + def ethereumUp = Mock(Upstream) { + _ * getStatus() >> UpstreamAvailability.OK + } + def ethereumUpAll = Mock(Multistream) { + _ * it.observeStatus() >> Mono.just(UpstreamAvailability.OK).repeat() + .delayElements(Duration.ofMillis(100)) + _ * it.getAll() >> [ethereumUp] + } + def ups = Mock(MultistreamHolder) { + _ * it.getAvailable() >> [Chain.ETHEREUM] + 1 * it.getUpstream(Chain.ETHEREUM) >> ethereumUpAll + 1 * it.getUpstream(Chain.BITCOIN) >> null + } + def ctrl = new SubscribeStatus(ups) + + when: + def req = BlockchainOuterClass.StatusRequest.newBuilder() + .addChains(Common.ChainRef.CHAIN_ETHEREUM) + .addChains(Common.ChainRef.CHAIN_BITCOIN) + .build() + def act = ctrl.subscribeStatus(Mono.just(req)) + .take(2) + // sort just for testing + .sort(new Comparator() { + @Override + int compare(BlockchainOuterClass.ChainStatus o1, BlockchainOuterClass.ChainStatus o2) { + return o1.chainValue <=> o2.chainValue + } + }) + then: + StepVerifier.create(act) + .expectNextMatches { + it.chainValue == Chain.BITCOIN.id && it.availability == BlockchainOuterClass.AvailabilityEnum.AVAIL_UNAVAILABLE + } + .expectNextMatches { + it.chainValue == Chain.ETHEREUM.id && it.availability == BlockchainOuterClass.AvailabilityEnum.AVAIL_OK + } + .expectComplete() + .verify(Duration.ofSeconds(3)) + } +}