This is an automated email from the ASF dual-hosted git repository. dmsolr pushed a commit to branch fix/#6645 in repository https://gitbox.apache.org/repos/asf/shenyu.git
commit 92b0e7e1417ba7244947383260d1511b697344ef Author: Luke Haochao Zhuang <[email protected]> AuthorDate: Thu Sep 3 12:01:03 2026 +0800 fix: invoke RedisRateLimiter error callback/logging inside onErrorResume (#6645) onErrorResume absorbed every Redis error into a normal Flux.just(1L, -1L) emission, so the trailing doOnError (which logged the error and invoked the rate limiter algorithm's callback, e.g. cleanup of optimistically written keys) was unreachable dead code. On a Redis outage every request was silently fail-opened with no log signal and no cleanup callback. Move the callback invocation and error logging into the onErrorResume lambda itself so they actually execute when Redis errors occur. --- .../ratelimiter/executor/RedisRateLimiter.java | 12 ++++++---- .../ratelimiter/executor/RedisRateLimiterTest.java | 28 ++++++++++++++++++++++ 2 files changed, 35 insertions(+), 5 deletions(-) diff --git a/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/main/java/org/apache/shenyu/plugin/ratelimiter/executor/RedisRateLimiter.java b/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/main/java/org/apache/shenyu/plugin/ratelimiter/executor/RedisRateLimiter.java index c55187f6dd..4a1522360e 100644 --- a/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/main/java/org/apache/shenyu/plugin/ratelimiter/executor/RedisRateLimiter.java +++ b/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/main/java/org/apache/shenyu/plugin/ratelimiter/executor/RedisRateLimiter.java @@ -60,7 +60,13 @@ public class RedisRateLimiter { List<String> keys = rateLimiterAlgorithm.getKeys(id); List<String> scriptArgs = Stream.of(replenishRate, burstCapacity, Instant.now().getEpochSecond(), requestCount).map(String::valueOf).collect(Collectors.toList()); Flux<List<Long>> resultFlux = Singleton.INST.get(ReactiveRedisTemplate.class).execute(script, keys, scriptArgs); - return resultFlux.onErrorResume(throwable -> Flux.just(Arrays.asList(1L, -1L))) + // the error is absorbed right here, so callback/logging must happen in this lambda, not in a downstream doOnError + return resultFlux + .onErrorResume(throwable -> { + rateLimiterAlgorithm.callback(rateLimiterAlgorithm.getScript(), keys, scriptArgs); + LOG.error("Error occurred while judging if user is allowed by RedisRateLimiter, fail-open and allow the request:{}", throwable.getMessage()); + return Flux.just(Arrays.asList(1L, -1L)); + }) .reduce(new ArrayList<Long>(), (longs, l) -> { longs.addAll(l); return longs; @@ -68,10 +74,6 @@ public class RedisRateLimiter { boolean allowed = ((Number) results.get(0)).longValue() == 1L; long tokensLeft = ((Number) results.get(1)).longValue(); return new RateLimiterResponse(allowed, tokensLeft, keys); - }) - .doOnError(throwable -> { - rateLimiterAlgorithm.callback(rateLimiterAlgorithm.getScript(), keys, scriptArgs); - LOG.error("Error occurred while judging if user is allowed by RedisRateLimiter:{}", throwable.getMessage()); }); } diff --git a/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/test/java/org/apache/shenyu/plugin/ratelimiter/executor/RedisRateLimiterTest.java b/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/test/java/org/apache/shenyu/plugin/ratelimiter/executor/RedisRateLimiterTest.java index 5ad9345dcb..b69dc1f02e 100644 --- a/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/test/java/org/apache/shenyu/plugin/ratelimiter/executor/RedisRateLimiterTest.java +++ b/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/test/java/org/apache/shenyu/plugin/ratelimiter/executor/RedisRateLimiterTest.java @@ -26,6 +26,7 @@ import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; import org.mockito.junit.jupiter.MockitoExtension; import org.springframework.data.redis.core.ReactiveRedisTemplate; +import org.springframework.data.redis.core.ReactiveZSetOperations; import org.springframework.data.redis.core.script.RedisScript; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @@ -39,6 +40,7 @@ import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyList; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; /** @@ -163,6 +165,32 @@ public final class RedisRateLimiterTest { }).verifyComplete(); } + /** + * redisRateLimiter.isAllowed exception case must still invoke the algorithm's error callback + * (e.g. cleanup of optimistically-written keys), since the error is absorbed by onErrorResume + * before the response reaches any downstream operator. + */ + @Test + @SuppressWarnings({"unchecked", "rawtypes"}) + public void allowedThrowableInvokesAlgorithmCallbackTest() { + ReactiveRedisTemplate reactiveRedisTemplate = mock(ReactiveRedisTemplate.class); + Singleton.INST.single(ReactiveRedisTemplate.class, reactiveRedisTemplate); + when(reactiveRedisTemplate.execute(any(RedisScript.class), anyList(), anyList())).thenReturn( + Flux.error(new IllegalStateException("redis unavailable"))); + ReactiveZSetOperations reactiveZSetOperations = mock(ReactiveZSetOperations.class); + when(reactiveRedisTemplate.opsForZSet()).thenReturn(reactiveZSetOperations); + when(reactiveZSetOperations.remove(any(), any())).thenReturn(Mono.just(1L)); + + rateLimiterHandle.setAlgorithmName("concurrent"); + Mono<RateLimiterResponse> responseMono = redisRateLimiter.isAllowed(DEFAULT_TEST_ID, rateLimiterHandle); + StepVerifier.create(responseMono).assertNext(r -> { + assertEquals(-1, r.getTokensRemaining()); + assertTrue(r.isAllowed()); + }).verifyComplete(); + + verify(reactiveZSetOperations).remove(any(), any()); + } + /** * redisRateLimiter.isAllowed test pre init. *
