added sampling

This commit is contained in:
Maksim Fomenkov
2022-12-20 12:01:34 +03:00
parent db97b656cf
commit 8f0c9d1187

View File

@@ -23,9 +23,14 @@ class SubscribeNodeStatus(
private val multistreams: CurrentMultistreamHolder private val multistreams: CurrentMultistreamHolder
) { ) {
companion object {
private val RETRY_TIMEOUT = Duration.ofSeconds(10)
}
fun subscribe(req: Mono<SubscribeNodeStatusRequest>): Flux<NodeStatusResponse> = fun subscribe(req: Mono<SubscribeNodeStatusRequest>): Flux<NodeStatusResponse> =
req.flatMapMany { req.flatMapMany { request ->
val knownUpstreams = ConcurrentHashMap<String, Boolean>() val knownUpstreams = ConcurrentHashMap<String, 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()
@@ -46,7 +51,7 @@ class SubscribeNodeStatus(
.flatMap { ms -> .flatMap { ms ->
ms.getAll().map { up -> ms.getAll().map { up ->
knownUpstreams[up.getId()] = true knownUpstreams[up.getId()] = true
subscribeUpstreamUpdates(ms.chain, up) subscribeUpstreamUpdates(ms.chain, up, duration)
} }
} }
) )
@@ -69,7 +74,7 @@ class SubscribeNodeStatus(
.setStatus(buildStatus(it.getStatus(), it.getHead().getCurrentHeight())) .setStatus(buildStatus(it.getStatus(), it.getHead().getCurrentHeight()))
.build() .build()
), ),
subscribeUpstreamUpdates(ms.chain, it) subscribeUpstreamUpdates(ms.chain, it, duration)
) )
} }
} }
@@ -78,29 +83,34 @@ class SubscribeNodeStatus(
Flux.concat(descriptions, Flux.merge(upstreamUpdates, muliStreamUpdates)) Flux.concat(descriptions, Flux.merge(upstreamUpdates, muliStreamUpdates))
} }
private fun subscribeUpstreamUpdates(chain: Chain, upstream: Upstream): Flux<NodeStatusResponse> { private fun subscribeUpstreamUpdates(
val r = Sinks.many().multicast().directBestEffort<Boolean>() chain: Chain,
val heads = Mono.just(upstream).repeatWhen { r.asFlux() }.sample(Duration.ofSeconds(10)).flatMap { up -> upstream: Upstream,
timespan: Duration
): Flux<NodeStatusResponse> {
val retry = Sinks.many().multicast().directBestEffort<Boolean>()
val heads = Mono.just(upstream).repeatWhen { retry.asFlux() }.sample(RETRY_TIMEOUT).flatMap { up ->
up.getHead().getFlux().map { block -> up.getHead().getFlux().map { block ->
NodeStatusResponse.newBuilder() NodeStatusResponse.newBuilder()
.setNodeId(up.getId()) .setNodeId(up.getId())
.setStatus(buildStatus(up.getStatus(), block.height)) .setStatus(buildStatus(up.getStatus(), block.height))
.build() .build()
}.doFinally { }.doFinally {
// retry when subscribed head stopped
if (it == SignalType.ON_COMPLETE) { if (it == SignalType.ON_COMPLETE) {
r.tryEmitNext(true) retry.tryEmitNext(true)
} }
} }
} }.sample(timespan)
val statuses = upstream.observeStatus().map { val statuses = upstream.observeStatus().distinctUntilChanged().map {
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(chain, upstream))
.build() .build()
} }.sample(timespan)
return Flux.merge(heads.distinctUntilChanged(), statuses.distinctUntilChanged()) return Flux.merge(heads, statuses)
} }
private fun buildDescription(chain: Chain, up: Upstream): NodeDescription.Builder = private fun buildDescription(chain: Chain, up: Upstream): NodeDescription.Builder =