From 3ef5a4455ce3eede5b4fe22bb01abf238408b151 Mon Sep 17 00:00:00 2001 From: Igor Artamonov Date: Sun, 4 Jul 2021 18:37:17 -0400 Subject: [PATCH] problem: doesn't publish upstream status updates --- .../io/emeraldpay/dshackle/upstream/DefaultUpstream.kt | 10 ++++++++-- .../io/emeraldpay/dshackle/upstream/Multistream.kt | 7 ++++++- 2 files changed, 14 insertions(+), 3 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/DefaultUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/DefaultUpstream.kt index f443c2f9..aba5ab05 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/DefaultUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/DefaultUpstream.kt @@ -21,6 +21,7 @@ import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.startup.QuorumForLabels import io.emeraldpay.dshackle.upstream.calls.CallMethods import reactor.core.publisher.Flux +import reactor.core.publisher.Sinks import reactor.extra.processor.TopicProcessor import java.util.concurrent.atomic.AtomicReference @@ -41,7 +42,9 @@ abstract class DefaultUpstream( this(id, Long.MAX_VALUE, UpstreamAvailability.UNAVAILABLE, options, role, targets, node) private val status = AtomicReference(Status(defaultLag, defaultAvail, statusByLag(defaultLag, defaultAvail))) - private val statusStream: TopicProcessor = TopicProcessor.create() + private val statusStream = Sinks.many() + .multicast() + .directBestEffort() override fun isAvailable(): Boolean { return getStatus() == UpstreamAvailability.OK @@ -63,6 +66,7 @@ abstract class DefaultUpstream( status.updateAndGet { curr -> Status(curr.lag, avail, statusByLag(curr.lag, avail)) } + statusStream.tryEmitNext(status.get().status) } fun statusByLag(lag: Long, proposed: UpstreamAvailability): UpstreamAvailability { @@ -76,7 +80,8 @@ abstract class DefaultUpstream( } override fun observeStatus(): Flux { - return Flux.from(statusStream) + return statusStream.asFlux() + .distinctUntilChanged() } override fun setLag(lag: Long) { @@ -86,6 +91,7 @@ abstract class DefaultUpstream( status.updateAndGet { curr -> Status(lag, curr.avail, statusByLag(lag, curr.avail)) } + statusStream.tryEmitNext(status.get().status) } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt index 5fe53e10..3122ada9 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt @@ -168,7 +168,12 @@ abstract class Multistream( } override fun observeStatus(): Flux { - val upstreamsFluxes = getAll().map { up -> up.observeStatus().map { UpstreamStatus(up, it) } } + val upstreamsFluxes = getAll().map { up -> + Flux.concat( + Mono.just(up.getStatus()), + up.observeStatus() + ).map { UpstreamStatus(up, it) } + } return Flux.merge(upstreamsFluxes) .filter(FilterBestAvailability()) .map { it.status }