This is an automated email from the ASF dual-hosted git repository.

Aias00 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/shenyu.git


The following commit(s) were added to refs/heads/master by this push:
     new 795cec97b0 fix: invoke RedisRateLimiter error callback/logging inside 
onErrorResume (#6645) (#7029)
795cec97b0 is described below

commit 795cec97b0ada7a5b1c050c9a1565615f5dc109e
Author: Daming <[email protected]>
AuthorDate: Sat Sep 5 07:23:59 2026 +0800

    fix: invoke RedisRateLimiter error callback/logging inside onErrorResume 
(#6645) (#7029)
    
    * 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.
    
    * Improve error handling in RedisRateLimiter
    
    Co-authored-by: Copilot Autofix powered by AI 
<[email protected]>
    
    ---------
    
    Co-authored-by: Luke Haochao Zhuang <[email protected]>
    Co-authored-by: Copilot Autofix powered by AI 
<[email protected]>
    Co-authored-by: aias00 <[email protected]>
---
 .../ratelimiter/executor/RedisRateLimiter.java     | 16 +++++++++----
 .../ratelimiter/executor/RedisRateLimiterTest.java | 28 ++++++++++++++++++++++
 2 files changed, 39 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..29869bca47 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,17 @@ 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 -> {
+                    try {
+                        rateLimiterAlgorithm.callback(script, keys, 
scriptArgs);
+                    } catch (final Exception callbackEx) {
+                        LOG.warn("RateLimiterAlgorithm callback failed after 
RedisRateLimiter error", callbackEx);
+                    }
+                    LOG.error("Error occurred while judging if user is allowed 
by RedisRateLimiter; fail-open and allow the request", throwable);
+                    return Flux.just(Arrays.asList(1L, -1L));
+                })
                 .reduce(new ArrayList<Long>(), (longs, l) -> {
                     longs.addAll(l);
                     return longs;
@@ -68,10 +78,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.
      *

Reply via email to