From 90e96ceb70199a28bc6f4dc1b164a6bea49614ef Mon Sep 17 00:00:00 2001 From: Igor Artamonov Date: Sun, 28 Jul 2019 18:04:54 -0400 Subject: [PATCH] solution: transfer actual status and availability through grpc --- .../io/emeraldpay/dshackle/rpc/Describe.kt | 30 ++++++++--------- .../dshackle/rpc/SubscribeStatus.kt | 32 ++++++++++++------- .../dshackle/upstream/GrpcUpstream.kt | 10 ++---- .../dshackle/upstream/UpstreamAvailability.kt | 22 +++++++++---- 4 files changed, 54 insertions(+), 40 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/Describe.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/Describe.kt index a1b93430..0fb3743f 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/Describe.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/Describe.kt @@ -4,40 +4,38 @@ import io.emeraldpay.api.proto.BlockchainOuterClass import io.emeraldpay.api.proto.Common import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.ConfiguredUpstreams +import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstreams +import io.emeraldpay.grpc.Chain import io.grpc.stub.StreamObserver import org.springframework.beans.factory.annotation.Autowired import org.springframework.stereotype.Service @Service class Describe( - @Autowired private val upstreams: Upstreams + @Autowired private val upstreams: Upstreams, + @Autowired private val subscribeStatus: SubscribeStatus ) { fun describe(request: BlockchainOuterClass.DescribeRequest, responseObserver: StreamObserver) { val resp = BlockchainOuterClass.DescribeResponse.newBuilder() upstreams.getAvailable().forEach { chain -> upstreams.getUpstream(chain)?.let { chainUpstreams -> - val quorum = chainUpstreams.getAll().map { u -> - if (u.getStatus() == UpstreamAvailability.OK) { - u.getOptions().quorum - } else { - 0 + chainUpstreams.getAll().let { ups -> + if (ups.isNotEmpty()) { + val status = subscribeStatus.chainStatus(chain, ups) + resp.addChains( + BlockchainOuterClass.DescribeChain.newBuilder() + .setChain(Common.ChainRef.forNumber(chain.id)) + .setStatus(status) + .build() + ) } - }.sum() - val available = chainUpstreams.getAll().any { u -> - u.getStatus() == UpstreamAvailability.OK } - resp.addChains( - BlockchainOuterClass.DescribeChain.newBuilder() - .setChain(Common.ChainRef.forNumber(chain.id)) - .setQuorum(quorum) - .setAvailable(available) - .build() - ) } } responseObserver.onNext(resp.build()) responseObserver.onCompleted() } + } \ No newline at end of file diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeStatus.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeStatus.kt index db7ab69a..0a517fa4 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeStatus.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeStatus.kt @@ -2,8 +2,10 @@ package io.emeraldpay.dshackle.rpc import io.emeraldpay.api.proto.BlockchainOuterClass import io.emeraldpay.api.proto.Common +import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.Upstreams +import io.emeraldpay.grpc.Chain import io.grpc.StatusRuntimeException import io.grpc.stub.StreamObserver import org.springframework.beans.factory.annotation.Autowired @@ -18,18 +20,11 @@ class SubscribeStatus( fun subscribeStatus(request: BlockchainOuterClass.StatusRequest, responseObserver: StreamObserver) { upstreams.getAvailable().forEach { chain -> var d: Disposable? = null - d = upstreams.getUpstream(chain)?.observeStatus()?.subscribe { availability -> - val chainStatus = BlockchainOuterClass.ChainStatus.newBuilder() - .setChain(Common.ChainRef.forNumber(chain.id)) - .setAvailable(availability == UpstreamAvailability.OK) - .setQuorum(0) - if (availability == UpstreamAvailability.OK) { - upstreams.getUpstream(chain)?.getOptions()?.let { opts -> - chainStatus.setQuorum(opts.quorum) - } - } + val chainUpstream = upstreams.getUpstream(chain) + d = chainUpstream?.observeStatus()?.subscribe { availability -> + val status = chainStatus(chain, chainUpstream.getAll()) try { - responseObserver.onNext(chainStatus.build()) + responseObserver.onNext(status) } catch (e: StatusRuntimeException) { // gRPC channel was closed d?.dispose() @@ -38,4 +33,19 @@ class SubscribeStatus( } } + fun chainStatus(chain: Chain, ups: List): BlockchainOuterClass.ChainStatus { + val available = ups.map { u -> + u.getStatus() + }.min()!! + val quorum = ups.filter { + it.getStatus() > UpstreamAvailability.UNAVAILABLE + }.count() + val status = BlockchainOuterClass.ChainStatus.newBuilder() + .setAvailability(BlockchainOuterClass.AvailabilityEnum.forNumber(available.grpcId)) + .setChain(Common.ChainRef.forNumber(chain.id)) + .setQuorum(quorum) + .build() + return status + } + } \ No newline at end of file diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/GrpcUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/GrpcUpstream.kt index 844e9381..c9a42435 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/GrpcUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/GrpcUpstream.kt @@ -86,18 +86,14 @@ open class GrpcUpstream( } fun init(conf: BlockchainOuterClass.DescribeChain) { - val available = conf.available - val quorum = conf.quorum - setStatus( - if (available && quorum > 0) UpstreamAvailability.OK else UpstreamAvailability.UNAVAILABLE - ) + conf.status?.let { status -> onStatus(status) } } fun onStatus(value: BlockchainOuterClass.ChainStatus) { - val available = value.available + val available = value.availability val quorum = value.quorum setStatus( - if (available && quorum > 0) UpstreamAvailability.OK else UpstreamAvailability.UNAVAILABLE + if (available != null) UpstreamAvailability.fromGrpc(available.number) else UpstreamAvailability.UNAVAILABLE ) } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/UpstreamAvailability.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/UpstreamAvailability.kt index 1d7214e7..f9f9eed3 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/UpstreamAvailability.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/UpstreamAvailability.kt @@ -1,11 +1,21 @@ package io.emeraldpay.dshackle.upstream -enum class UpstreamAvailability { +enum class UpstreamAvailability(val grpcId: Int) { - OK, - IMMATURE, - SYNCING, - LAGGING, - UNAVAILABLE + OK(1), + IMMATURE(2), + SYNCING(3), + LAGGING(4), + UNAVAILABLE(5); + companion object { + fun fromGrpc(id: Int?): UpstreamAvailability { + if (id == null) { + return UNAVAILABLE + } + return UpstreamAvailability.values().find { + it.grpcId == id + } ?: UNAVAILABLE + } + } } \ No newline at end of file