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