fix style and tests
This commit is contained in:
@@ -27,43 +27,53 @@ class SubscribeNodeStatus(
|
|||||||
req.flatMapMany {
|
req.flatMapMany {
|
||||||
val knownUpstreams = ConcurrentHashMap<String, Boolean>()
|
val knownUpstreams = ConcurrentHashMap<String, Boolean>()
|
||||||
// send known upstreams details immediately
|
// send known upstreams details immediately
|
||||||
val descriptions = Flux.fromIterable(multistreams.all().flatMap { multiStream ->
|
val descriptions = Flux.fromIterable(
|
||||||
multiStream.getAll().map { up ->
|
multistreams.all()
|
||||||
NodeStatusResponse.newBuilder()
|
.flatMap { multiStream ->
|
||||||
.setNodeId(up.getId())
|
multiStream.getAll().map { up ->
|
||||||
.setDescription(buildDescription(multiStream.chain, up))
|
NodeStatusResponse.newBuilder()
|
||||||
.setStatus(buildStatus(up.getStatus(), up.getHead().getCurrentHeight()))
|
.setNodeId(up.getId())
|
||||||
.build()
|
.setDescription(buildDescription(multiStream.chain, up))
|
||||||
}
|
.setStatus(buildStatus(up.getStatus(), up.getHead().getCurrentHeight()))
|
||||||
})
|
.build()
|
||||||
|
}
|
||||||
//subscribe on head/status updates for known upstreams
|
|
||||||
val upstreamUpdates = Flux.merge(multistreams.all().flatMap { ms ->
|
|
||||||
ms.getAll().map { up ->
|
|
||||||
knownUpstreams[up.getId()] = true
|
|
||||||
subscribeUpstreamUpdates(ms.chain, up)
|
|
||||||
}
|
|
||||||
})
|
|
||||||
|
|
||||||
//subscribe on head/status updates for just added upstreams
|
|
||||||
val muliStreamUpdates = Flux.merge(multistreams.all().map { ms ->
|
|
||||||
ms.subscribeAddedUpstreams()
|
|
||||||
.distinctUntilChanged {
|
|
||||||
it.getId()
|
|
||||||
}
|
}
|
||||||
.filter { knownUpstreams[it.getId()] != true }.flatMap {
|
)
|
||||||
knownUpstreams[it.getId()] = true
|
|
||||||
Flux.concat(
|
// subscribe on head/status updates for known upstreams
|
||||||
Mono.just(
|
val upstreamUpdates = Flux.merge(
|
||||||
NodeStatusResponse.newBuilder()
|
multistreams.all()
|
||||||
.setNodeId(it.getId())
|
.flatMap { ms ->
|
||||||
.setDescription(buildDescription(ms.chain, it))
|
ms.getAll().map { up ->
|
||||||
.setStatus(buildStatus(it.getStatus(), it.getHead().getCurrentHeight()))
|
knownUpstreams[up.getId()] = true
|
||||||
.build()
|
subscribeUpstreamUpdates(ms.chain, up)
|
||||||
), subscribeUpstreamUpdates(ms.chain, it)
|
}
|
||||||
)
|
|
||||||
}
|
}
|
||||||
})
|
)
|
||||||
|
|
||||||
|
// subscribe on head/status updates for just added upstreams
|
||||||
|
val muliStreamUpdates = Flux.merge(
|
||||||
|
multistreams.all()
|
||||||
|
.map { ms ->
|
||||||
|
ms.subscribeAddedUpstreams()
|
||||||
|
.distinctUntilChanged {
|
||||||
|
it.getId()
|
||||||
|
}
|
||||||
|
.filter { knownUpstreams[it.getId()] != true }.flatMap {
|
||||||
|
knownUpstreams[it.getId()] = true
|
||||||
|
Flux.concat(
|
||||||
|
Mono.just(
|
||||||
|
NodeStatusResponse.newBuilder()
|
||||||
|
.setNodeId(it.getId())
|
||||||
|
.setDescription(buildDescription(ms.chain, it))
|
||||||
|
.setStatus(buildStatus(it.getStatus(), it.getHead().getCurrentHeight()))
|
||||||
|
.build()
|
||||||
|
),
|
||||||
|
subscribeUpstreamUpdates(ms.chain, it)
|
||||||
|
)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
Flux.concat(descriptions, Flux.merge(upstreamUpdates, muliStreamUpdates))
|
Flux.concat(descriptions, Flux.merge(upstreamUpdates, muliStreamUpdates))
|
||||||
}
|
}
|
||||||
@@ -96,14 +106,16 @@ class SubscribeNodeStatus(
|
|||||||
private fun buildDescription(chain: Chain, up: Upstream): NodeDescription.Builder =
|
private fun buildDescription(chain: Chain, up: Upstream): NodeDescription.Builder =
|
||||||
NodeDescription.newBuilder()
|
NodeDescription.newBuilder()
|
||||||
.setChain(Common.ChainRef.forNumber(chain.id))
|
.setChain(Common.ChainRef.forNumber(chain.id))
|
||||||
.addAllLabels(up.getLabels().flatMap { labels ->
|
.addAllLabels(
|
||||||
labels.map {
|
up.getLabels().flatMap { labels ->
|
||||||
BlockchainOuterClass.Label.newBuilder()
|
labels.map {
|
||||||
.setName(it.key)
|
BlockchainOuterClass.Label.newBuilder()
|
||||||
.setValue(it.value)
|
.setName(it.key)
|
||||||
.build()
|
.setValue(it.value)
|
||||||
|
.build()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
})
|
)
|
||||||
.addAllSupportedMethods(up.getMethods().getSupportedMethods())
|
.addAllSupportedMethods(up.getMethods().getSupportedMethods())
|
||||||
|
|
||||||
private fun buildStatus(status: UpstreamAvailability, height: Long?): NodeStatus.Builder =
|
private fun buildStatus(status: UpstreamAvailability, height: Long?): NodeStatus.Builder =
|
||||||
|
|||||||
@@ -86,6 +86,11 @@ class MultistreamHolderMock implements MultistreamHolder {
|
|||||||
return upstreams.containsKey(chain)
|
return upstreams.containsKey(chain)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
List<Multistream> all() {
|
||||||
|
return upstreams.values().toList()
|
||||||
|
}
|
||||||
|
|
||||||
static class EthereumMultistreamMock extends EthereumPosMultiStream {
|
static class EthereumMultistreamMock extends EthereumPosMultiStream {
|
||||||
|
|
||||||
EthereumCachingReader customReader = null
|
EthereumCachingReader customReader = null
|
||||||
|
|||||||
Reference in New Issue
Block a user