From 3e2cae47930ab0c5f2fa7c795e1f40725ea9c2d1 Mon Sep 17 00:00:00 2001 From: KirillPamPam Date: Wed, 28 Jun 2023 13:34:13 +0400 Subject: [PATCH] Send head updates (#236) --- .../dshackle/rpc/SubscribeNodeStatus.kt | 18 +++++++++++++++++- 1 file changed, 17 insertions(+), 1 deletion(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt index 8d3ef6c6..4c400b35 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/SubscribeNodeStatus.kt @@ -16,6 +16,7 @@ import org.springframework.stereotype.Service import reactor.core.publisher.Flux import reactor.core.publisher.Mono import reactor.core.publisher.Sinks +import java.time.Duration import java.util.concurrent.ConcurrentHashMap @Service @@ -100,6 +101,21 @@ class SubscribeNodeStatus( upstream: Upstream, cancel: Sinks.Many ): Flux { + val heads = upstream.getHead().getFlux() + .takeUntilOther(cancel.asFlux()) + .sample(Duration.ofSeconds(2)) + .map { block -> + NodeStatusResponse.newBuilder() + .setDescription( + NodeDescription.newBuilder() + .setNodeId(upstream.nodeId().toInt()) + .setChain(Common.ChainRef.forNumber(ms.chain.id)) + .build() + ) + .setNodeId(upstream.getId()) + .setStatus(buildStatus(upstream.getStatus(), block.height)) + .build() + } val statuses = upstream.observeStatus() .takeUntilOther(cancel.asFlux()) .map { @@ -116,7 +132,7 @@ class SubscribeNodeStatus( .setStatus(buildStatus(upstream.getStatus(), upstream.getHead().getCurrentHeight())) .build() ) - return Flux.concat(currentState, statuses) + return Flux.concat(currentState, Flux.merge(statuses, heads)) } private fun buildDescription(ms: Multistream, up: Upstream): NodeDescription.Builder {