problem: produces repeated status update on same chain when it has many upstreams
This commit is contained in:
@@ -117,11 +117,6 @@ open class CurrentMultistreamHolder(
|
|||||||
return chainMapping[chain]
|
return chainMapping[chain]
|
||||||
}
|
}
|
||||||
|
|
||||||
@Scheduled(fixedRate = 15000)
|
|
||||||
fun printStatuses() {
|
|
||||||
chainMapping.forEach { it.value.printStatus() }
|
|
||||||
}
|
|
||||||
|
|
||||||
override fun getAvailable(): List<Chain> {
|
override fun getAvailable(): List<Chain> {
|
||||||
return Collections.unmodifiableList(chainMapping.keys.toList())
|
return Collections.unmodifiableList(chainMapping.keys.toList())
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -33,6 +33,7 @@ import org.springframework.context.Lifecycle
|
|||||||
import reactor.core.Disposable
|
import reactor.core.Disposable
|
||||||
import reactor.core.publisher.Flux
|
import reactor.core.publisher.Flux
|
||||||
import reactor.core.publisher.Mono
|
import reactor.core.publisher.Mono
|
||||||
|
import reactor.core.scheduler.Schedulers
|
||||||
import java.time.Duration
|
import java.time.Duration
|
||||||
import java.time.Instant
|
import java.time.Instant
|
||||||
import java.util.concurrent.atomic.AtomicReference
|
import java.util.concurrent.atomic.AtomicReference
|
||||||
@@ -208,8 +209,12 @@ abstract class Multistream(
|
|||||||
}
|
}
|
||||||
|
|
||||||
override fun start() {
|
override fun start() {
|
||||||
subscription = observeStatus()
|
val repeated = Flux.interval(Duration.ofSeconds(30))
|
||||||
|
val whenChanged = observeStatus()
|
||||||
.distinctUntilChanged()
|
.distinctUntilChanged()
|
||||||
|
subscription = Flux.merge(repeated, whenChanged)
|
||||||
|
// print status _change_ every 15 seconds, at most; otherwise prints it on interval of 30 seconds
|
||||||
|
.sample(Duration.ofSeconds(15))
|
||||||
.subscribe { printStatus() }
|
.subscribe { printStatus() }
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -289,15 +294,18 @@ abstract class Multistream(
|
|||||||
private val lastRef = AtomicReference<UpstreamStatus>()
|
private val lastRef = AtomicReference<UpstreamStatus>()
|
||||||
|
|
||||||
override fun test(t: UpstreamStatus): Boolean {
|
override fun test(t: UpstreamStatus): Boolean {
|
||||||
val last = lastRef.get()
|
val curr = lastRef.updateAndGet { last ->
|
||||||
val changed = last == null
|
val changed = last == null
|
||||||
|| t.status > last.status
|
|| t.status > last.status
|
||||||
|| (last.upstream == t.upstream && t.status != last.status)
|
|| (last.upstream == t.upstream && t.status != last.status)
|
||||||
|| last.ts.isBefore(Instant.now() - Duration.ofSeconds(60))
|
|| last.ts.isBefore(Instant.now() - Duration.ofSeconds(60))
|
||||||
if (changed) {
|
if (changed) {
|
||||||
lastRef.set(t)
|
t
|
||||||
|
} else {
|
||||||
|
last
|
||||||
|
}
|
||||||
}
|
}
|
||||||
return changed
|
return curr == t
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user