From 153e5115ff675d3ad7fdd460adb93e1cd6a0b1fb Mon Sep 17 00:00:00 2001 From: a10zn8 Date: Fri, 17 Mar 2023 18:42:09 +0300 Subject: [PATCH] head lag observer throttling - fix test (#170) --- .../io/emeraldpay/dshackle/upstream/HeadLagObserver.kt | 5 +++-- .../emeraldpay/dshackle/upstream/HeadLagObserverSpec.groovy | 4 ++-- 2 files changed, 5 insertions(+), 4 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/HeadLagObserver.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/HeadLagObserver.kt index 7138c6b2..054c4247 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/HeadLagObserver.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/HeadLagObserver.kt @@ -34,7 +34,8 @@ typealias Extractor = (top: BlockContainer, curr: BlockContainer) -> DistanceExt abstract class HeadLagObserver( private val master: Head, private val followers: Collection, - private val distanceExtractor: Extractor + private val distanceExtractor: Extractor, + private val throttling: Duration = Duration.ofSeconds(5) ) : Lifecycle { private val log = LoggerFactory.getLogger(HeadLagObserver::class.java) @@ -57,7 +58,7 @@ abstract class HeadLagObserver( fun subscription(): Flux { return master.getFlux() - .sample(Duration.ofSeconds(5)) + .sample(throttling) .flatMap(this::probeFollowers) .map { item -> item.t2.setLag(item.t1) diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/HeadLagObserverSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/HeadLagObserverSpec.groovy index 3d1f8a66..b667be41 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/HeadLagObserverSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/HeadLagObserverSpec.groovy @@ -72,7 +72,7 @@ class HeadLagObserverSpec extends Specification { HeadLagObserver observer = new TestHeadLagObserver(master, [up1, up2]) when: - def act = observer.subscription().take(Duration.ofMillis(1200)) + def act = observer.subscription().take(Duration.ofMillis(5000)) then: StepVerifier.create(act) @@ -111,7 +111,7 @@ class HeadLagObserverSpec extends Specification { class TestHeadLagObserver extends HeadLagObserver { TestHeadLagObserver(@NotNull Head master, @NotNull Collection followers) { - super(master, followers, DistanceExtractor.@Companion::extractPowDistance) + super(master, followers, DistanceExtractor.@Companion::extractPowDistance, Duration.ofNanos(1)) } @Override