fix node status subsciption for gRPC upstreams
This commit is contained in:
@@ -119,9 +119,11 @@ class SubscribeNodeStatus(
|
|||||||
|
|
||||||
val statuses = upstream.observeStatus()
|
val statuses = upstream.observeStatus()
|
||||||
.distinctUntilChanged()
|
.distinctUntilChanged()
|
||||||
.takeUntil { it == UpstreamAvailability.UNAVAILABLE }
|
.takeUntil {
|
||||||
|
it == UpstreamAvailability.UNAVAILABLE && !upstream.isGrpc()
|
||||||
|
}
|
||||||
.map {
|
.map {
|
||||||
if (it == UpstreamAvailability.UNAVAILABLE) {
|
if (it == UpstreamAvailability.UNAVAILABLE && !upstream.isGrpc()) {
|
||||||
onUnavailable.accept(upstream.getId())
|
onUnavailable.accept(upstream.getId())
|
||||||
// cancel head subscription & reconnections when upstream becomes unavailable
|
// cancel head subscription & reconnections when upstream becomes unavailable
|
||||||
cancel.tryEmitNext(true)
|
cancel.tryEmitNext(true)
|
||||||
|
|||||||
Reference in New Issue
Block a user