diff --git a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/rate/LinearRateLimiter.java b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/rate/LinearRateLimiter.java index 25af46b5d8..d2a7a94eb0 100644 --- a/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/rate/LinearRateLimiter.java +++ b/operator-framework-core/src/main/java/io/javaoperatorsdk/operator/processing/event/rate/LinearRateLimiter.java @@ -64,7 +64,10 @@ public Optional isLimited(RateLimitState rateLimitState) { actualState.increaseCount(); return Optional.empty(); } else { - return Optional.of(Duration.between(actualState.getLastRefreshTime(), LocalDateTime.now())); + var remaining = + Duration.between( + LocalDateTime.now(), actualState.getLastRefreshTime().plus(refreshPeriod)); + return Optional.of(remaining.isNegative() ? Duration.ZERO : remaining); } } diff --git a/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/event/rate/LinearRateLimiterTest.java b/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/event/rate/LinearRateLimiterTest.java index 43b30f9de7..7c5bdac479 100644 --- a/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/event/rate/LinearRateLimiterTest.java +++ b/operator-framework-core/src/test/java/io/javaoperatorsdk/operator/processing/event/rate/LinearRateLimiterTest.java @@ -53,7 +53,22 @@ void returnsMinimalDurationToAcquirePermission() { res = rl.isLimited(state); assertThat(res).isPresent(); - assertThat(res.get()).isLessThan(REFRESH_PERIOD); + assertThat(res.get()).isLessThanOrEqualTo(REFRESH_PERIOD); + // the whole period is still ahead of us, so the reported wait must be close to it + assertThat(res.get()).isGreaterThan(REFRESH_PERIOD.dividedBy(2)); + } + + @Test + void reportedDurationIsTheTimeRemainingNotTheTimeElapsed() throws InterruptedException { + var rl = new LinearRateLimiter(REFRESH_PERIOD, 1); + assertThat(rl.isLimited(state)).isEmpty(); + + var justAfterLimit = rl.isLimited(state).orElseThrow(); + Thread.sleep(REFRESH_PERIOD.toMillis() / 2); + var halfWayThroughPeriod = rl.isLimited(state).orElseThrow(); + + // as the refresh period elapses, less time remains until a permission is available again + assertThat(halfWayThroughPeriod).isLessThan(justAfterLimit); } @Test