Fix update methods (#221)
This commit is contained in:
@@ -278,12 +278,7 @@ abstract class Multistream(
|
|||||||
|
|
||||||
private fun observeUpstreamsStatuses() {
|
private fun observeUpstreamsStatuses() {
|
||||||
stateStream.asFlux()
|
stateStream.asFlux()
|
||||||
.distinctUntilChanged(
|
.subscribe {
|
||||||
{ it },
|
|
||||||
{ prev, current ->
|
|
||||||
prev.status == current.status || prev.equals(current)
|
|
||||||
}
|
|
||||||
).subscribe {
|
|
||||||
upstreams.filter { it.isAvailable() }.map { it.getMethods() }.let {
|
upstreams.filter { it.isAvailable() }.map { it.getMethods() }.let {
|
||||||
callMethods = AggregatedCallMethods(it)
|
callMethods = AggregatedCallMethods(it)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -264,11 +264,13 @@ class MultistreamSpec extends Specification {
|
|||||||
then:
|
then:
|
||||||
StepVerifier.create(states)
|
StepVerifier.create(states)
|
||||||
.then {
|
.then {
|
||||||
up1.onStatus(status(BlockchainOuterClass.AvailabilityEnum.AVAIL_OK))
|
up1.onStatus(status(BlockchainOuterClass.AvailabilityEnum.AVAIL_UNAVAILABLE))
|
||||||
up2.onStatus(status(BlockchainOuterClass.AvailabilityEnum.AVAIL_OK))
|
up2.onStatus(status(BlockchainOuterClass.AvailabilityEnum.AVAIL_OK))
|
||||||
|
up1.onStatus(status(BlockchainOuterClass.AvailabilityEnum.AVAIL_OK))
|
||||||
}
|
}
|
||||||
.expectNext(new Multistream.UpstreamChangeState(up1.getId(), UpstreamAvailability.OK))
|
.expectNext(new Multistream.UpstreamChangeState(up1.getId(), UpstreamAvailability.UNAVAILABLE))
|
||||||
.expectNext(new Multistream.UpstreamChangeState(up2.getId(), UpstreamAvailability.OK))
|
.expectNext(new Multistream.UpstreamChangeState(up2.getId(), UpstreamAvailability.OK))
|
||||||
|
.expectNext(new Multistream.UpstreamChangeState(up1.getId(), UpstreamAvailability.OK))
|
||||||
.then {
|
.then {
|
||||||
assert ms.getMethods().supportedMethods == Set.of("eth_test1", "eth_test2", "eth_test3")
|
assert ms.getMethods().supportedMethods == Set.of("eth_test1", "eth_test2", "eth_test3")
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user