subscribe node status - remove head subscription and remove sampling, afraid of missing updates

This commit is contained in:
Termina1
2023-01-10 13:33:35 +02:00
parent ca77d77d16
commit 4aa54cb907

View File

@@ -30,9 +30,8 @@ class SubscribeNodeStatus(
} }
fun subscribe(req: Mono<SubscribeNodeStatusRequest>): Flux<NodeStatusResponse> = fun subscribe(req: Mono<SubscribeNodeStatusRequest>): Flux<NodeStatusResponse> =
req.flatMapMany { request -> req.flatMapMany {
val knownUpstreams = ConcurrentHashMap<String, Sinks.Many<Boolean>>() val knownUpstreams = ConcurrentHashMap<String, Sinks.Many<Boolean>>()
val duration = Duration.ofMillis(request.timespan)
// send known upstreams details immediately // send known upstreams details immediately
val descriptions = Flux.fromIterable( val descriptions = Flux.fromIterable(
multistreams.all() multistreams.all()
@@ -53,7 +52,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, duration, knownUpstreams[up.getId()]!!) subscribeUpstreamUpdates(ms.chain, up, knownUpstreams[up.getId()]!!)
} }
} }
) )
@@ -100,7 +99,7 @@ class SubscribeNodeStatus(
.setStatus(buildStatus(it.getStatus(), it.getHead().getCurrentHeight())) .setStatus(buildStatus(it.getStatus(), it.getHead().getCurrentHeight()))
.build() .build()
), ),
subscribeUpstreamUpdates(ms.chain, it, duration, knownUpstreams[it.getId()]!!) subscribeUpstreamUpdates(ms.chain, it, knownUpstreams[it.getId()]!!)
) )
} }
} }
@@ -112,31 +111,8 @@ class SubscribeNodeStatus(
private fun subscribeUpstreamUpdates( private fun subscribeUpstreamUpdates(
chain: Chain, chain: Chain,
upstream: Upstream, upstream: Upstream,
timespan: Duration,
cancel: Sinks.Many<Boolean> cancel: Sinks.Many<Boolean>
): Flux<NodeStatusResponse> { ): Flux<NodeStatusResponse> {
val retry = Sinks.many().multicast().directBestEffort<Boolean>()
val heads = Mono.just(upstream)
.repeatWhen {
retry.asFlux()
}
.sample(RETRY_TIMEOUT)
.takeUntilOther(cancel.asFlux()).flatMap { up ->
up.getHead().getFlux()
.takeUntilOther(cancel.asFlux())
.map { block ->
NodeStatusResponse.newBuilder()
.setNodeId(up.getId())
.setStatus(buildStatus(up.getStatus(), block.height))
.build()
}.doFinally {
// retry when subscribed head stopped
if (it == SignalType.ON_COMPLETE) {
retry.tryEmitNext(true)
}
}
}.sample(timespan)
val statuses = upstream.observeStatus() val statuses = upstream.observeStatus()
.distinctUntilChanged() .distinctUntilChanged()
.takeUntilOther(cancel.asFlux()) .takeUntilOther(cancel.asFlux())
@@ -146,9 +122,9 @@ class SubscribeNodeStatus(
.setStatus(buildStatus(it, upstream.getHead().getCurrentHeight())) .setStatus(buildStatus(it, upstream.getHead().getCurrentHeight()))
.setDescription(buildDescription(chain, upstream)) .setDescription(buildDescription(chain, upstream))
.build() .build()
}.sample(timespan) }
return Flux.merge(heads, statuses) return statuses
} }
private fun buildDescription(chain: Chain, up: Upstream): NodeDescription.Builder = private fun buildDescription(chain: Chain, up: Upstream): NodeDescription.Builder =