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 903f7d71b1 fix(ratelimiter): expire stale concurrent entries (#7204)
903f7d71b1 is described below

commit 903f7d71b1d803e15778bcf6701bd0cb454f5679
Author: Liming Deng <[email protected]>
AuthorDate: Thu Sep 24 16:52:30 2026 +0800

    fix(ratelimiter): expire stale concurrent entries (#7204)
    
    Co-authored-by: shown <[email protected]>
    Co-authored-by: aias00 <[email protected]>
---
 .../ratelimiter/algorithm/ConcurrentRateLimiterAlgorithm.java     | 7 ++++++-
 .../META-INF/scripts/concurrent_request_rate_limiter.lua          | 4 +++-
 .../ratelimiter/algorithm/ConcurrentRateLimiterAlgorithmTest.java | 8 ++++++++
 3 files changed, 17 insertions(+), 2 deletions(-)

diff --git 
a/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/main/java/org/apache/shenyu/plugin/ratelimiter/algorithm/ConcurrentRateLimiterAlgorithm.java
 
b/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/main/java/org/apache/shenyu/plugin/ratelimiter/algorithm/ConcurrentRateLimiterAlgorithm.java
index 0fbe8756db..b18519093e 100644
--- 
a/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/main/java/org/apache/shenyu/plugin/ratelimiter/algorithm/ConcurrentRateLimiterAlgorithm.java
+++ 
b/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/main/java/org/apache/shenyu/plugin/ratelimiter/algorithm/ConcurrentRateLimiterAlgorithm.java
@@ -21,6 +21,8 @@ import org.apache.shenyu.common.enums.RateLimitEnum;
 import org.apache.shenyu.common.utils.UUIDUtils;
 import org.apache.shenyu.common.utils.Singleton;
 import org.apache.shenyu.spi.Join;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 import org.springframework.data.redis.core.ReactiveRedisTemplate;
 import org.springframework.data.redis.core.script.RedisScript;
 
@@ -36,6 +38,8 @@ import java.util.List;
 @Join
 public class ConcurrentRateLimiterAlgorithm extends 
AbstractRateLimiterAlgorithm {
 
+    private static final Logger LOG = 
LoggerFactory.getLogger(ConcurrentRateLimiterAlgorithm.class);
+
     public ConcurrentRateLimiterAlgorithm() {
         super(RateLimitEnum.CONCURRENT.getScriptName());
     }
@@ -56,6 +60,7 @@ public class ConcurrentRateLimiterAlgorithm extends 
AbstractRateLimiterAlgorithm
     @Override
     @SuppressWarnings("unchecked")
     public void callback(final RedisScript<?> script, final List<String> keys, 
final List<?> scriptArgs) {
-        
Singleton.INST.get(ReactiveRedisTemplate.class).opsForZSet().remove(keys.get(0),
 keys.get(1)).subscribe();
+        
Singleton.INST.get(ReactiveRedisTemplate.class).opsForZSet().remove(keys.get(0),
 keys.get(1))
+                .subscribe(ignored -> { }, error -> LOG.warn("Failed to remove 
concurrent rate limiter entry", error));
     }
 }
diff --git 
a/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/main/resources/META-INF/scripts/concurrent_request_rate_limiter.lua
 
b/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/main/resources/META-INF/scripts/concurrent_request_rate_limiter.lua
index 64a73c2f2f..b5a2321cf4 100644
--- 
a/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/main/resources/META-INF/scripts/concurrent_request_rate_limiter.lua
+++ 
b/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/main/resources/META-INF/scripts/concurrent_request_rate_limiter.lua
@@ -20,14 +20,16 @@ local key = KEYS[1]
 local capacity = tonumber(ARGV[2])
 local timestamp = tonumber(ARGV[3])
 local id = KEYS[2]
+local stale_after_seconds = 86400
 
+redis.call("zremrangebyscore", key, 0, timestamp - stale_after_seconds)
 local count = redis.call("zcard", key)
 local allowed = 0
 
 if count < capacity then
   redis.call("zadd", key, timestamp, id)
+  redis.call("expire", key, stale_after_seconds)
   allowed = 1
   count = count + 1
 end
--- redis.call("setex", key, timestamp)
 return { allowed, count }
diff --git 
a/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/test/java/org/apache/shenyu/plugin/ratelimiter/algorithm/ConcurrentRateLimiterAlgorithmTest.java
 
b/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/test/java/org/apache/shenyu/plugin/ratelimiter/algorithm/ConcurrentRateLimiterAlgorithmTest.java
index ed96abf97e..b4ced486c8 100644
--- 
a/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/test/java/org/apache/shenyu/plugin/ratelimiter/algorithm/ConcurrentRateLimiterAlgorithmTest.java
+++ 
b/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/test/java/org/apache/shenyu/plugin/ratelimiter/algorithm/ConcurrentRateLimiterAlgorithmTest.java
@@ -23,6 +23,7 @@ import org.junit.jupiter.api.extension.ExtendWith;
 import org.mockito.junit.jupiter.MockitoExtension;
 
 import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.containsString;
 import static org.hamcrest.Matchers.is;
 
 /**
@@ -47,4 +48,11 @@ public final class ConcurrentRateLimiterAlgorithmTest {
     public void getKeyNameTest() {
         assertThat("concurrent_request_rate_limiter", 
is(concurrentRateLimiterAlgorithm.getKeyName()));
     }
+
+    @Test
+    public void scriptExpiresAndRemovesStaleEntriesTest() {
+        String script = 
concurrentRateLimiterAlgorithm.getScript().getScriptAsString();
+        assertThat(script, containsString("zremrangebyscore"));
+        assertThat(script, containsString("expire"));
+    }
 }

Reply via email to