Merge pull request #108 from p2p-org/fix-missin-events-subscribe-status
Fix missing events subscribe status
This commit is contained in:
@@ -29,20 +29,6 @@ class SubscribeNodeStatus(
|
|||||||
fun subscribe(req: Mono<SubscribeNodeStatusRequest>): Flux<NodeStatusResponse> =
|
fun subscribe(req: Mono<SubscribeNodeStatusRequest>): Flux<NodeStatusResponse> =
|
||||||
req.flatMapMany {
|
req.flatMapMany {
|
||||||
val knownUpstreams = ConcurrentHashMap<String, Sinks.Many<Boolean>>()
|
val knownUpstreams = ConcurrentHashMap<String, Sinks.Many<Boolean>>()
|
||||||
// send known upstreams details immediately
|
|
||||||
val descriptions = Flux.fromIterable(
|
|
||||||
multistreams.all()
|
|
||||||
.flatMap { multiStream ->
|
|
||||||
multiStream.getAll().map { up ->
|
|
||||||
NodeStatusResponse.newBuilder()
|
|
||||||
.setNodeId(up.getId())
|
|
||||||
.setDescription(buildDescription(multiStream.chain, up))
|
|
||||||
.setStatus(buildStatus(up.getStatus(), up.getHead().getCurrentHeight()))
|
|
||||||
.build()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
)
|
|
||||||
|
|
||||||
// subscribe on head/status updates for known upstreams
|
// subscribe on head/status updates for known upstreams
|
||||||
val upstreamUpdates = Flux.merge(
|
val upstreamUpdates = Flux.merge(
|
||||||
multistreams.all()
|
multistreams.all()
|
||||||
@@ -102,7 +88,7 @@ class SubscribeNodeStatus(
|
|||||||
}
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
Flux.concat(descriptions, Flux.merge(upstreamUpdates, multiStreamUpdates, removals))
|
Flux.merge(upstreamUpdates, multiStreamUpdates, removals)
|
||||||
}
|
}
|
||||||
|
|
||||||
private fun subscribeUpstreamUpdates(
|
private fun subscribeUpstreamUpdates(
|
||||||
@@ -120,8 +106,14 @@ class SubscribeNodeStatus(
|
|||||||
.setDescription(buildDescription(chain, upstream))
|
.setDescription(buildDescription(chain, upstream))
|
||||||
.build()
|
.build()
|
||||||
}
|
}
|
||||||
|
val currentState = Mono.just(
|
||||||
return statuses
|
NodeStatusResponse.newBuilder()
|
||||||
|
.setNodeId(upstream.getId())
|
||||||
|
.setDescription(buildDescription(chain, 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(chain: Chain, up: Upstream): NodeDescription.Builder =
|
||||||
|
|||||||
Reference in New Issue
Block a user