fixed review comments
This commit is contained in:
@@ -57,14 +57,14 @@ class SubscribeNodeStatus(
|
|||||||
)
|
)
|
||||||
|
|
||||||
// subscribe on head/status updates for just added upstreams
|
// subscribe on head/status updates for just added upstreams
|
||||||
val muliStreamUpdates = Flux.merge(
|
val multiStreamUpdates = Flux.merge(
|
||||||
multistreams.all()
|
multistreams.all()
|
||||||
.map { ms ->
|
.map { ms ->
|
||||||
ms.subscribeAddedUpstreams()
|
ms.subscribeAddedUpstreams()
|
||||||
.distinctUntilChanged {
|
.distinctUntilChanged {
|
||||||
it.getId()
|
it.getId()
|
||||||
}
|
}
|
||||||
.filter { knownUpstreams[it.getId()] != true }.flatMap {
|
.filter { knownUpstreams.getOrDefault(it.getId(), false) }.flatMap {
|
||||||
knownUpstreams[it.getId()] = true
|
knownUpstreams[it.getId()] = true
|
||||||
Flux.concat(
|
Flux.concat(
|
||||||
Mono.just(
|
Mono.just(
|
||||||
@@ -80,7 +80,7 @@ class SubscribeNodeStatus(
|
|||||||
}
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
Flux.concat(descriptions, Flux.merge(upstreamUpdates, muliStreamUpdates))
|
Flux.concat(descriptions, Flux.merge(upstreamUpdates, multiStreamUpdates))
|
||||||
}
|
}
|
||||||
|
|
||||||
private fun subscribeUpstreamUpdates(
|
private fun subscribeUpstreamUpdates(
|
||||||
|
|||||||
Reference in New Issue
Block a user