solution: transfer actual status and availability through grpc
This commit is contained in:
@@ -4,40 +4,38 @@ import io.emeraldpay.api.proto.BlockchainOuterClass
|
|||||||
import io.emeraldpay.api.proto.Common
|
import io.emeraldpay.api.proto.Common
|
||||||
import io.emeraldpay.dshackle.upstream.UpstreamAvailability
|
import io.emeraldpay.dshackle.upstream.UpstreamAvailability
|
||||||
import io.emeraldpay.dshackle.upstream.ConfiguredUpstreams
|
import io.emeraldpay.dshackle.upstream.ConfiguredUpstreams
|
||||||
|
import io.emeraldpay.dshackle.upstream.Upstream
|
||||||
import io.emeraldpay.dshackle.upstream.Upstreams
|
import io.emeraldpay.dshackle.upstream.Upstreams
|
||||||
|
import io.emeraldpay.grpc.Chain
|
||||||
import io.grpc.stub.StreamObserver
|
import io.grpc.stub.StreamObserver
|
||||||
import org.springframework.beans.factory.annotation.Autowired
|
import org.springframework.beans.factory.annotation.Autowired
|
||||||
import org.springframework.stereotype.Service
|
import org.springframework.stereotype.Service
|
||||||
|
|
||||||
@Service
|
@Service
|
||||||
class Describe(
|
class Describe(
|
||||||
@Autowired private val upstreams: Upstreams
|
@Autowired private val upstreams: Upstreams,
|
||||||
|
@Autowired private val subscribeStatus: SubscribeStatus
|
||||||
) {
|
) {
|
||||||
|
|
||||||
fun describe(request: BlockchainOuterClass.DescribeRequest, responseObserver: StreamObserver<BlockchainOuterClass.DescribeResponse>) {
|
fun describe(request: BlockchainOuterClass.DescribeRequest, responseObserver: StreamObserver<BlockchainOuterClass.DescribeResponse>) {
|
||||||
val resp = BlockchainOuterClass.DescribeResponse.newBuilder()
|
val resp = BlockchainOuterClass.DescribeResponse.newBuilder()
|
||||||
upstreams.getAvailable().forEach { chain ->
|
upstreams.getAvailable().forEach { chain ->
|
||||||
upstreams.getUpstream(chain)?.let { chainUpstreams ->
|
upstreams.getUpstream(chain)?.let { chainUpstreams ->
|
||||||
val quorum = chainUpstreams.getAll().map { u ->
|
chainUpstreams.getAll().let { ups ->
|
||||||
if (u.getStatus() == UpstreamAvailability.OK) {
|
if (ups.isNotEmpty()) {
|
||||||
u.getOptions().quorum
|
val status = subscribeStatus.chainStatus(chain, ups)
|
||||||
} else {
|
resp.addChains(
|
||||||
0
|
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.onNext(resp.build())
|
||||||
responseObserver.onCompleted()
|
responseObserver.onCompleted()
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
@@ -2,8 +2,10 @@ package io.emeraldpay.dshackle.rpc
|
|||||||
|
|
||||||
import io.emeraldpay.api.proto.BlockchainOuterClass
|
import io.emeraldpay.api.proto.BlockchainOuterClass
|
||||||
import io.emeraldpay.api.proto.Common
|
import io.emeraldpay.api.proto.Common
|
||||||
|
import io.emeraldpay.dshackle.upstream.Upstream
|
||||||
import io.emeraldpay.dshackle.upstream.UpstreamAvailability
|
import io.emeraldpay.dshackle.upstream.UpstreamAvailability
|
||||||
import io.emeraldpay.dshackle.upstream.Upstreams
|
import io.emeraldpay.dshackle.upstream.Upstreams
|
||||||
|
import io.emeraldpay.grpc.Chain
|
||||||
import io.grpc.StatusRuntimeException
|
import io.grpc.StatusRuntimeException
|
||||||
import io.grpc.stub.StreamObserver
|
import io.grpc.stub.StreamObserver
|
||||||
import org.springframework.beans.factory.annotation.Autowired
|
import org.springframework.beans.factory.annotation.Autowired
|
||||||
@@ -18,18 +20,11 @@ class SubscribeStatus(
|
|||||||
fun subscribeStatus(request: BlockchainOuterClass.StatusRequest, responseObserver: StreamObserver<BlockchainOuterClass.ChainStatus>) {
|
fun subscribeStatus(request: BlockchainOuterClass.StatusRequest, responseObserver: StreamObserver<BlockchainOuterClass.ChainStatus>) {
|
||||||
upstreams.getAvailable().forEach { chain ->
|
upstreams.getAvailable().forEach { chain ->
|
||||||
var d: Disposable? = null
|
var d: Disposable? = null
|
||||||
d = upstreams.getUpstream(chain)?.observeStatus()?.subscribe { availability ->
|
val chainUpstream = upstreams.getUpstream(chain)
|
||||||
val chainStatus = BlockchainOuterClass.ChainStatus.newBuilder()
|
d = chainUpstream?.observeStatus()?.subscribe { availability ->
|
||||||
.setChain(Common.ChainRef.forNumber(chain.id))
|
val status = chainStatus(chain, chainUpstream.getAll())
|
||||||
.setAvailable(availability == UpstreamAvailability.OK)
|
|
||||||
.setQuorum(0)
|
|
||||||
if (availability == UpstreamAvailability.OK) {
|
|
||||||
upstreams.getUpstream(chain)?.getOptions()?.let { opts ->
|
|
||||||
chainStatus.setQuorum(opts.quorum)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
try {
|
try {
|
||||||
responseObserver.onNext(chainStatus.build())
|
responseObserver.onNext(status)
|
||||||
} catch (e: StatusRuntimeException) {
|
} catch (e: StatusRuntimeException) {
|
||||||
// gRPC channel was closed
|
// gRPC channel was closed
|
||||||
d?.dispose()
|
d?.dispose()
|
||||||
@@ -38,4 +33,19 @@ class SubscribeStatus(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fun chainStatus(chain: Chain, ups: List<Upstream>): 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
|
||||||
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
@@ -86,18 +86,14 @@ open class GrpcUpstream(
|
|||||||
}
|
}
|
||||||
|
|
||||||
fun init(conf: BlockchainOuterClass.DescribeChain) {
|
fun init(conf: BlockchainOuterClass.DescribeChain) {
|
||||||
val available = conf.available
|
conf.status?.let { status -> onStatus(status) }
|
||||||
val quorum = conf.quorum
|
|
||||||
setStatus(
|
|
||||||
if (available && quorum > 0) UpstreamAvailability.OK else UpstreamAvailability.UNAVAILABLE
|
|
||||||
)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fun onStatus(value: BlockchainOuterClass.ChainStatus) {
|
fun onStatus(value: BlockchainOuterClass.ChainStatus) {
|
||||||
val available = value.available
|
val available = value.availability
|
||||||
val quorum = value.quorum
|
val quorum = value.quorum
|
||||||
setStatus(
|
setStatus(
|
||||||
if (available && quorum > 0) UpstreamAvailability.OK else UpstreamAvailability.UNAVAILABLE
|
if (available != null) UpstreamAvailability.fromGrpc(available.number) else UpstreamAvailability.UNAVAILABLE
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1,11 +1,21 @@
|
|||||||
package io.emeraldpay.dshackle.upstream
|
package io.emeraldpay.dshackle.upstream
|
||||||
|
|
||||||
enum class UpstreamAvailability {
|
enum class UpstreamAvailability(val grpcId: Int) {
|
||||||
|
|
||||||
OK,
|
OK(1),
|
||||||
IMMATURE,
|
IMMATURE(2),
|
||||||
SYNCING,
|
SYNCING(3),
|
||||||
LAGGING,
|
LAGGING(4),
|
||||||
UNAVAILABLE
|
UNAVAILABLE(5);
|
||||||
|
|
||||||
|
companion object {
|
||||||
|
fun fromGrpc(id: Int?): UpstreamAvailability {
|
||||||
|
if (id == null) {
|
||||||
|
return UNAVAILABLE
|
||||||
|
}
|
||||||
|
return UpstreamAvailability.values().find {
|
||||||
|
it.grpcId == id
|
||||||
|
} ?: UNAVAILABLE
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
Reference in New Issue
Block a user