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"));
+ }
}