From 0f7fd35b123529d45472cc60ae28971f5843c6e0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=9A=D0=B8=D1=80=D0=B8=D0=BB=D0=BB?= Date: Mon, 23 Jan 2023 17:22:43 +0400 Subject: [PATCH] Add supported sub methods to SubscribeNodeStatus --- emerald-grpc | 2 +- .../dshackle/rpc/SubscribeNodeStatus.kt | 21 ++++++++++--------- 2 files changed, 12 insertions(+), 11 deletions(-) diff --git a/emerald-grpc b/emerald-grpc index f5543c96..fddf154c 160000 --- a/emerald-grpc +++ b/emerald-grpc @@ -1 +1 @@ -Subproject commit f5543c960330757f9ea1e729407ec334b163f232 +Subproject commit fddf154c34b7a59bee790cc14eff9eff69e8339f diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt index 318548bb..306ecf05 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt @@ -6,8 +6,8 @@ import io.emeraldpay.api.proto.BlockchainOuterClass.NodeStatus import io.emeraldpay.api.proto.BlockchainOuterClass.NodeStatusResponse import io.emeraldpay.api.proto.BlockchainOuterClass.SubscribeNodeStatusRequest import io.emeraldpay.api.proto.Common -import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.upstream.CurrentMultistreamHolder +import io.emeraldpay.dshackle.upstream.Multistream import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.UpstreamAvailability import org.slf4j.LoggerFactory @@ -35,7 +35,7 @@ class SubscribeNodeStatus( .flatMap { ms -> ms.getAll().map { up -> knownUpstreams[up.getId()] = Sinks.many().multicast().directBestEffort() - subscribeUpstreamUpdates(ms.chain, up, knownUpstreams[up.getId()]!!) + subscribeUpstreamUpdates(ms, up, knownUpstreams[up.getId()]!!) } } ) @@ -53,7 +53,7 @@ class SubscribeNodeStatus( knownUpstreams.remove(up.getId()) NodeStatusResponse.newBuilder() .setNodeId(up.getId()) - .setDescription(buildDescription(ms.chain, up)) + .setDescription(buildDescription(ms, up)) .setStatus(buildStatus(UpstreamAvailability.UNAVAILABLE, up.getHead().getCurrentHeight())) .build() } @@ -78,11 +78,11 @@ class SubscribeNodeStatus( Mono.just( NodeStatusResponse.newBuilder() .setNodeId(it.getId()) - .setDescription(buildDescription(ms.chain, it)) + .setDescription(buildDescription(ms, it)) .setStatus(buildStatus(it.getStatus(), it.getHead().getCurrentHeight())) .build() ), - subscribeUpstreamUpdates(ms.chain, it, knownUpstreams[it.getId()]!!) + subscribeUpstreamUpdates(ms, it, knownUpstreams[it.getId()]!!) ) } } @@ -92,7 +92,7 @@ class SubscribeNodeStatus( } private fun subscribeUpstreamUpdates( - chain: Chain, + ms: Multistream, upstream: Upstream, cancel: Sinks.Many ): Flux { @@ -103,22 +103,22 @@ class SubscribeNodeStatus( NodeStatusResponse.newBuilder() .setNodeId(upstream.getId()) .setStatus(buildStatus(it, upstream.getHead().getCurrentHeight())) - .setDescription(buildDescription(chain, upstream)) + .setDescription(buildDescription(ms, upstream)) .build() } val currentState = Mono.just( NodeStatusResponse.newBuilder() .setNodeId(upstream.getId()) - .setDescription(buildDescription(chain, upstream)) + .setDescription(buildDescription(ms, upstream)) .setStatus(buildStatus(upstream.getStatus(), upstream.getHead().getCurrentHeight())) .build() ) 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() - .setChain(Common.ChainRef.forNumber(chain.id)) + .setChain(Common.ChainRef.forNumber(ms.chain.id)) .setNodeId(up.nodeId().toInt()) .addAllNodeLabels( up.getLabels().map { nodeLabels -> @@ -134,6 +134,7 @@ class SubscribeNodeStatus( .build() } ) + .addAllSupportedSubscriptions(ms.getEgressSubscription().getAvailableTopics()) .addAllSupportedMethods(up.getMethods().getSupportedMethods()) private fun buildStatus(status: UpstreamAvailability, height: Long?): NodeStatus.Builder =