Fix subscription event update (#517)
This commit is contained in:
@@ -220,7 +220,7 @@ abstract class Multistream(
|
|||||||
|
|
||||||
protected open fun onUpstreamsUpdated() {
|
protected open fun onUpstreamsUpdated() {
|
||||||
val upstreams = getAll()
|
val upstreams = getAll()
|
||||||
state.updateState(upstreams, getSubscriptionTopics())
|
state.updateState(upstreams, getEgressSubscription())
|
||||||
|
|
||||||
when {
|
when {
|
||||||
upstreams.size == 1 -> {
|
upstreams.size == 1 -> {
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import io.emeraldpay.dshackle.Defaults
|
|||||||
import io.emeraldpay.dshackle.startup.QuorumForLabels
|
import io.emeraldpay.dshackle.startup.QuorumForLabels
|
||||||
import io.emeraldpay.dshackle.upstream.Capability
|
import io.emeraldpay.dshackle.upstream.Capability
|
||||||
import io.emeraldpay.dshackle.upstream.DefaultUpstream
|
import io.emeraldpay.dshackle.upstream.DefaultUpstream
|
||||||
|
import io.emeraldpay.dshackle.upstream.EgressSubscription
|
||||||
import io.emeraldpay.dshackle.upstream.Upstream
|
import io.emeraldpay.dshackle.upstream.Upstream
|
||||||
import io.emeraldpay.dshackle.upstream.UpstreamAvailability
|
import io.emeraldpay.dshackle.upstream.UpstreamAvailability
|
||||||
import io.emeraldpay.dshackle.upstream.calls.AggregatedCallMethods
|
import io.emeraldpay.dshackle.upstream.calls.AggregatedCallMethods
|
||||||
@@ -59,7 +60,7 @@ class MultistreamState(
|
|||||||
return status
|
return status
|
||||||
}
|
}
|
||||||
|
|
||||||
fun updateState(upstreams: List<Upstream>, subs: List<String>) {
|
fun updateState(upstreams: List<Upstream>, egressSubscription: EgressSubscription) {
|
||||||
val oldState = CurrentMultistreamState(this)
|
val oldState = CurrentMultistreamState(this)
|
||||||
|
|
||||||
val availableUpstreams = upstreams.filter { it.isAvailable() }
|
val availableUpstreams = upstreams.filter { it.isAvailable() }
|
||||||
@@ -68,7 +69,7 @@ class MultistreamState(
|
|||||||
updateQuorumLabels(availableUpstreams)
|
updateQuorumLabels(availableUpstreams)
|
||||||
updateUpstreamBounds(availableUpstreams)
|
updateUpstreamBounds(availableUpstreams)
|
||||||
status = if (upstreams.isEmpty()) UpstreamAvailability.UNAVAILABLE else upstreams.minOf { it.getStatus() }
|
status = if (upstreams.isEmpty()) UpstreamAvailability.UNAVAILABLE else upstreams.minOf { it.getStatus() }
|
||||||
this.subs = subs
|
this.subs = egressSubscription.getAvailableTopics()
|
||||||
|
|
||||||
stateEvents.emitNext(
|
stateEvents.emitNext(
|
||||||
stateHandler.compareStates(oldState, CurrentMultistreamState(this)),
|
stateHandler.compareStates(oldState, CurrentMultistreamState(this)),
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import io.emeraldpay.dshackle.config.UpstreamsConfig
|
|||||||
import io.emeraldpay.dshackle.startup.QuorumForLabels
|
import io.emeraldpay.dshackle.startup.QuorumForLabels
|
||||||
import io.emeraldpay.dshackle.upstream.Capability
|
import io.emeraldpay.dshackle.upstream.Capability
|
||||||
import io.emeraldpay.dshackle.upstream.DefaultUpstream
|
import io.emeraldpay.dshackle.upstream.DefaultUpstream
|
||||||
|
import io.emeraldpay.dshackle.upstream.EgressSubscription
|
||||||
import io.emeraldpay.dshackle.upstream.UpstreamAvailability
|
import io.emeraldpay.dshackle.upstream.UpstreamAvailability
|
||||||
import io.emeraldpay.dshackle.upstream.calls.CallMethods
|
import io.emeraldpay.dshackle.upstream.calls.CallMethods
|
||||||
import io.emeraldpay.dshackle.upstream.finalization.FinalizationData
|
import io.emeraldpay.dshackle.upstream.finalization.FinalizationData
|
||||||
@@ -57,11 +58,14 @@ class MultistreamStateTest {
|
|||||||
listOf(LowerBoundData(1, LowerBoundType.STATE), LowerBoundData(1, LowerBoundType.BLOCK)),
|
listOf(LowerBoundData(1, LowerBoundType.STATE), LowerBoundData(1, LowerBoundType.BLOCK)),
|
||||||
listOf(FinalizationData(990, FinalizationType.SAFE_BLOCK), FinalizationData(880, FinalizationType.FINALIZED_BLOCK)),
|
listOf(FinalizationData(990, FinalizationType.SAFE_BLOCK), FinalizationData(880, FinalizationType.FINALIZED_BLOCK)),
|
||||||
)
|
)
|
||||||
|
val egressSub = mock<EgressSubscription> {
|
||||||
|
on { getAvailableTopics() } doReturn listOf("heads", "notHeads")
|
||||||
|
}
|
||||||
|
|
||||||
val state = MultistreamState {}
|
val state = MultistreamState {}
|
||||||
|
|
||||||
StepVerifier.create(state.stateEvents())
|
StepVerifier.create(state.stateEvents())
|
||||||
.then { state.updateState(listOf(up1, up2, up3), listOf("heads", "notHeads")) }
|
.then { state.updateState(listOf(up1, up2, up3), egressSub) }
|
||||||
.assertNext {
|
.assertNext {
|
||||||
assertThat(it).hasSize(7)
|
assertThat(it).hasSize(7)
|
||||||
assertThat(it.toList())
|
assertThat(it.toList())
|
||||||
@@ -93,11 +97,11 @@ class MultistreamStateTest {
|
|||||||
),
|
),
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
.then { state.updateState(listOf(up1, up2, up3), listOf("heads", "notHeads")) }
|
.then { state.updateState(listOf(up1, up2, up3), egressSub) }
|
||||||
.assertNext {
|
.assertNext {
|
||||||
assertThat(it).hasSize(0)
|
assertThat(it).hasSize(0)
|
||||||
}
|
}
|
||||||
.then { state.updateState(listOf(up1, up2, up3, up4), listOf("heads", "notHeads")) }
|
.then { state.updateState(listOf(up1, up2, up3, up4), egressSub) }
|
||||||
.assertNext {
|
.assertNext {
|
||||||
assertThat(it).hasSize(4)
|
assertThat(it).hasSize(4)
|
||||||
assertThat(it.toList())
|
assertThat(it.toList())
|
||||||
|
|||||||
Reference in New Issue
Block a user