Add supported sub methods to SubscribeNodeStatus

This commit is contained in:
Кирилл
2023-01-23 17:22:43 +04:00
parent b00f14d928
commit 0f7fd35b12
2 changed files with 12 additions and 11 deletions

View File

@@ -6,8 +6,8 @@ import io.emeraldpay.api.proto.BlockchainOuterClass.NodeStatus
import io.emeraldpay.api.proto.BlockchainOuterClass.NodeStatusResponse import io.emeraldpay.api.proto.BlockchainOuterClass.NodeStatusResponse
import io.emeraldpay.api.proto.BlockchainOuterClass.SubscribeNodeStatusRequest import io.emeraldpay.api.proto.BlockchainOuterClass.SubscribeNodeStatusRequest
import io.emeraldpay.api.proto.Common import io.emeraldpay.api.proto.Common
import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.upstream.CurrentMultistreamHolder import io.emeraldpay.dshackle.upstream.CurrentMultistreamHolder
import io.emeraldpay.dshackle.upstream.Multistream
import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.UpstreamAvailability
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
@@ -35,7 +35,7 @@ class SubscribeNodeStatus(
.flatMap { ms -> .flatMap { ms ->
ms.getAll().map { up -> ms.getAll().map { up ->
knownUpstreams[up.getId()] = Sinks.many().multicast().directBestEffort<Boolean>() knownUpstreams[up.getId()] = Sinks.many().multicast().directBestEffort<Boolean>()
subscribeUpstreamUpdates(ms.chain, up, knownUpstreams[up.getId()]!!) subscribeUpstreamUpdates(ms, up, knownUpstreams[up.getId()]!!)
} }
} }
) )
@@ -53,7 +53,7 @@ class SubscribeNodeStatus(
knownUpstreams.remove(up.getId()) knownUpstreams.remove(up.getId())
NodeStatusResponse.newBuilder() NodeStatusResponse.newBuilder()
.setNodeId(up.getId()) .setNodeId(up.getId())
.setDescription(buildDescription(ms.chain, up)) .setDescription(buildDescription(ms, up))
.setStatus(buildStatus(UpstreamAvailability.UNAVAILABLE, up.getHead().getCurrentHeight())) .setStatus(buildStatus(UpstreamAvailability.UNAVAILABLE, up.getHead().getCurrentHeight()))
.build() .build()
} }
@@ -78,11 +78,11 @@ class SubscribeNodeStatus(
Mono.just( Mono.just(
NodeStatusResponse.newBuilder() NodeStatusResponse.newBuilder()
.setNodeId(it.getId()) .setNodeId(it.getId())
.setDescription(buildDescription(ms.chain, it)) .setDescription(buildDescription(ms, it))
.setStatus(buildStatus(it.getStatus(), it.getHead().getCurrentHeight())) .setStatus(buildStatus(it.getStatus(), it.getHead().getCurrentHeight()))
.build() .build()
), ),
subscribeUpstreamUpdates(ms.chain, it, knownUpstreams[it.getId()]!!) subscribeUpstreamUpdates(ms, it, knownUpstreams[it.getId()]!!)
) )
} }
} }
@@ -92,7 +92,7 @@ class SubscribeNodeStatus(
} }
private fun subscribeUpstreamUpdates( private fun subscribeUpstreamUpdates(
chain: Chain, ms: Multistream,
upstream: Upstream, upstream: Upstream,
cancel: Sinks.Many<Boolean> cancel: Sinks.Many<Boolean>
): Flux<NodeStatusResponse> { ): Flux<NodeStatusResponse> {
@@ -103,22 +103,22 @@ class SubscribeNodeStatus(
NodeStatusResponse.newBuilder() NodeStatusResponse.newBuilder()
.setNodeId(upstream.getId()) .setNodeId(upstream.getId())
.setStatus(buildStatus(it, upstream.getHead().getCurrentHeight())) .setStatus(buildStatus(it, upstream.getHead().getCurrentHeight()))
.setDescription(buildDescription(chain, upstream)) .setDescription(buildDescription(ms, upstream))
.build() .build()
} }
val currentState = Mono.just( val currentState = Mono.just(
NodeStatusResponse.newBuilder() NodeStatusResponse.newBuilder()
.setNodeId(upstream.getId()) .setNodeId(upstream.getId())
.setDescription(buildDescription(chain, upstream)) .setDescription(buildDescription(ms, upstream))
.setStatus(buildStatus(upstream.getStatus(), upstream.getHead().getCurrentHeight())) .setStatus(buildStatus(upstream.getStatus(), upstream.getHead().getCurrentHeight()))
.build() .build()
) )
return Flux.concat(currentState, statuses) return Flux.concat(currentState, statuses)
} }
private fun buildDescription(chain: Chain, up: Upstream): NodeDescription.Builder = private fun buildDescription(ms: Multistream, up: Upstream): NodeDescription.Builder =
NodeDescription.newBuilder() NodeDescription.newBuilder()
.setChain(Common.ChainRef.forNumber(chain.id)) .setChain(Common.ChainRef.forNumber(ms.chain.id))
.setNodeId(up.nodeId().toInt()) .setNodeId(up.nodeId().toInt())
.addAllNodeLabels( .addAllNodeLabels(
up.getLabels().map { nodeLabels -> up.getLabels().map { nodeLabels ->
@@ -134,6 +134,7 @@ class SubscribeNodeStatus(
.build() .build()
} }
) )
.addAllSupportedSubscriptions(ms.getEgressSubscription().getAvailableTopics())
.addAllSupportedMethods(up.getMethods().getSupportedMethods()) .addAllSupportedMethods(up.getMethods().getSupportedMethods())
private fun buildStatus(status: UpstreamAvailability, height: Long?): NodeStatus.Builder = private fun buildStatus(status: UpstreamAvailability, height: Long?): NodeStatus.Builder =