From 34d5592c76571ddaaa564f172dd3599fbe5fbffc Mon Sep 17 00:00:00 2001 From: Vyacheslav Date: Fri, 22 Sep 2023 10:29:56 +0300 Subject: [PATCH] emit unavailable from all multistreams on shutdown (#303) --- .../io/emeraldpay/dshackle/upstream/Multistream.kt | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt index 2d41711e..73491f96 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/Multistream.kt @@ -69,6 +69,7 @@ abstract class Multistream( private var callMethodsFactory: Factory = Factory { return@Factory callMethods ?: throw FunctorException("Not initialized yet") } + private var stopSignal = Sinks.many().multicast().directBestEffort() private var seq = 0 protected var lagObserver: HeadLagObserver? = null private var subscription: Disposable? = null @@ -256,9 +257,11 @@ abstract class Multistream( up.observeStatus() ).map { UpstreamStatus(up, it) } } - return Flux.merge(upstreamsFluxes) - .map(FilterBestAvailability()) - .distinct() + val onShutdown = stopSignal.asFlux().map { UpstreamAvailability.UNAVAILABLE } + return Flux.merge( + Flux.merge(upstreamsFluxes).map(FilterBestAvailability()).takeUntilOther(stopSignal.asFlux()), + onShutdown + ).distinct() } override fun isAvailable(): Boolean { @@ -330,6 +333,7 @@ abstract class Multistream( } } lagObserver?.stop() + stopSignal.tryEmitNext(true) started = false }