head lag observer throttling - fix test (#170)
This commit is contained in:
@@ -34,7 +34,8 @@ typealias Extractor = (top: BlockContainer, curr: BlockContainer) -> DistanceExt
|
||||
abstract class HeadLagObserver(
|
||||
private val master: Head,
|
||||
private val followers: Collection<Upstream>,
|
||||
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<Unit> {
|
||||
return master.getFlux()
|
||||
.sample(Duration.ofSeconds(5))
|
||||
.sample(throttling)
|
||||
.flatMap(this::probeFollowers)
|
||||
.map { item ->
|
||||
item.t2.setLag(item.t1)
|
||||
|
||||
@@ -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<? extends Upstream> followers) {
|
||||
super(master, followers, DistanceExtractor.@Companion::extractPowDistance)
|
||||
super(master, followers, DistanceExtractor.@Companion::extractPowDistance, Duration.ofNanos(1))
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
Reference in New Issue
Block a user