From e61997a301964c992e8cac9c8e5323bb8eb5cdeb Mon Sep 17 00:00:00 2001 From: HY-love-sleep <583699747@qq.com> Date: Tue, 22 Sep 2026 10:54:03 +0800 Subject: [PATCH] fix: release the redis client that is replaced (#7156) - RedisConnectionFactory implements DisposableBean: destroy() delegates to the lettuce connection factory, and destroyQuietly(ReactiveRedisConnectionFactory) lets the callers that keep only the reactive template release the client they replace - the four sites that rebuild a redis client on a configuration change now destroy the previous one: AiTokenLimiterPluginHandler, RateLimiterPluginDataHandler, RedisCache.close() (which closed a connection but not the factory that owns the pool) and SensitiveWordPluginDataHandler, which also releases the client in removePlugin - tests: the lettuce factory stops running after destroy, destroyQuietly ignores a factory without a lifecycle (and a null one) and swallows a failure, and every handler destroys the client it replaces while keeping the new one running and without rebuilding an unchanged configuration --- .../infra/redis/RedisConnectionFactory.java | 36 ++++++- .../redis/RedisConnectionFactoryTest.java | 41 +++++++ .../SensitiveWordPluginDataHandler.java | 11 +- .../SensitiveWordPluginDataHandlerTest.java | 54 ++++++++++ .../handler/AiTokenLimiterPluginHandler.java | 6 ++ .../AiTokenLimiterPluginHandlerTest.java | 102 ++++++++++++++++++ .../shenyu/plugin/cache/redis/RedisCache.java | 2 + .../handler/RateLimiterPluginDataHandler.java | 5 + .../RateLimiterPluginDataHandlerTest.java | 29 +++++ 9 files changed, 284 insertions(+), 2 deletions(-) create mode 100644 shenyu-plugin/shenyu-plugin-ai/shenyu-plugin-ai-token-limiter/src/test/java/org/apache/shenyu/plugin/ai/token/limiter/handler/AiTokenLimiterPluginHandlerTest.java diff --git a/shenyu-infra/shenyu-infra-redis/src/main/java/org/apache/shenyu/infra/redis/RedisConnectionFactory.java b/shenyu-infra/shenyu-infra-redis/src/main/java/org/apache/shenyu/infra/redis/RedisConnectionFactory.java index f8f68bb269b3..90f666434e2f 100644 --- a/shenyu-infra/shenyu-infra-redis/src/main/java/org/apache/shenyu/infra/redis/RedisConnectionFactory.java +++ b/shenyu-infra/shenyu-infra-redis/src/main/java/org/apache/shenyu/infra/redis/RedisConnectionFactory.java @@ -21,6 +21,10 @@ import com.google.common.collect.Lists; import org.apache.commons.pool2.impl.GenericObjectPoolConfig; import org.apache.shenyu.common.enums.RedisModeEnum; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.DisposableBean; +import org.springframework.data.redis.connection.ReactiveRedisConnectionFactory; import org.springframework.data.redis.connection.RedisClusterConfiguration; import org.springframework.data.redis.connection.RedisNode; import org.springframework.data.redis.connection.RedisPassword; @@ -38,7 +42,9 @@ /** * RedisConnectionFactory. */ -public class RedisConnectionFactory { +public class RedisConnectionFactory implements DisposableBean { + + private static final Logger LOG = LoggerFactory.getLogger(RedisConnectionFactory.class); private final LettuceConnectionFactory lettuceConnectionFactory; @@ -56,6 +62,34 @@ public LettuceConnectionFactory getLettuceConnectionFactory() { return this.lettuceConnectionFactory; } + /** + * Destroy the lettuce connection factory and the connection pool it owns. The client this factory + * was built for must not be used afterwards. + */ + @Override + public void destroy() { + lettuceConnectionFactory.destroy(); + } + + /** + * Destroy a connection factory that was created by this class or by {@link #getLettuceConnectionFactory()}. + * The handlers of the plugins that rebuild their client on a configuration change keep the reactive + * template only, so this is how they release the client they replace. A null factory, or one that has + * no lifecycle, is ignored; a failure is logged rather than thrown, because a client that cannot be + * released must not break the configuration update that replaces it. + * + * @param connectionFactory the connection factory to destroy, may be null + */ + public static void destroyQuietly(final ReactiveRedisConnectionFactory connectionFactory) { + if (connectionFactory instanceof DisposableBean) { + try { + ((DisposableBean) connectionFactory).destroy(); + } catch (Exception e) { + LOG.warn("failed to destroy the redis connection factory", e); + } + } + } + private LettuceConnectionFactory createLettuceConnectionFactory(final RedisConfigProperties redisConfigProperties) { LettuceClientConfiguration lettuceClientConfiguration = getLettuceClientConfiguration(redisConfigProperties); if (RedisModeEnum.SENTINEL.getName().equals(redisConfigProperties.getMode())) { diff --git a/shenyu-infra/shenyu-infra-redis/src/test/java/org/apache/shenyu/infra/redis/RedisConnectionFactoryTest.java b/shenyu-infra/shenyu-infra-redis/src/test/java/org/apache/shenyu/infra/redis/RedisConnectionFactoryTest.java index 0f3f3cbe8da9..fe18982ed828 100644 --- a/shenyu-infra/shenyu-infra-redis/src/test/java/org/apache/shenyu/infra/redis/RedisConnectionFactoryTest.java +++ b/shenyu-infra/shenyu-infra-redis/src/test/java/org/apache/shenyu/infra/redis/RedisConnectionFactoryTest.java @@ -21,7 +21,10 @@ import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; import org.mockito.Mockito; +import org.springframework.beans.factory.DisposableBean; +import org.springframework.data.redis.connection.ReactiveRedisConnectionFactory; import org.springframework.data.redis.connection.RedisNode; +import org.springframework.data.redis.connection.lettuce.LettuceConnectionFactory; import java.time.Duration; import java.lang.reflect.Method; @@ -53,6 +56,38 @@ public void redisConnectionFactoryTest() { Assertions.assertDoesNotThrow(() -> new RedisConnectionFactory(redisConfigProperties)); } + @Test + public void destroyDestroysTheLettuceFactory() { + RedisConfigProperties redisConfigProperties = new RedisConfigProperties(); + redisConfigProperties.setUrl("localhost:6379"); + redisConfigProperties.setMode(RedisModeEnum.STANDALONE.getName()); + RedisConnectionFactory factory = new RedisConnectionFactory(redisConfigProperties); + LettuceConnectionFactory lettuceConnectionFactory = factory.getLettuceConnectionFactory(); + Assertions.assertTrue(lettuceConnectionFactory.isRunning()); + factory.destroy(); + Assertions.assertFalse(lettuceConnectionFactory.isRunning()); + } + + @Test + public void destroyQuietlyDestroysTheFactory() throws Exception { + DisposableReactiveFactory connectionFactory = Mockito.mock(DisposableReactiveFactory.class); + RedisConnectionFactory.destroyQuietly(connectionFactory); + Mockito.verify(connectionFactory).destroy(); + } + + @Test + public void destroyQuietlyIgnoresWhatItCannotDestroy() throws Exception { + // nothing to destroy + Assertions.assertDoesNotThrow(() -> RedisConnectionFactory.destroyQuietly(null)); + // a factory of another type, without a lifecycle, is left alone + Assertions.assertDoesNotThrow(() -> RedisConnectionFactory.destroyQuietly( + Mockito.mock(ReactiveRedisConnectionFactory.class))); + // a failure is logged instead of thrown: replacing a client must not fail because of it + DisposableReactiveFactory failing = Mockito.mock(DisposableReactiveFactory.class); + Mockito.doThrow(new IllegalStateException("boom")).when(failing).destroy(); + Assertions.assertDoesNotThrow(() -> RedisConnectionFactory.destroyQuietly(failing)); + } + @Test public void parseRedisNodeValidInputs() throws Exception { RedisConnectionFactory factory = createFactoryWithDefaultUrl(); @@ -118,4 +153,10 @@ private void assertInvalidNode(final Method parseMethod, final RedisConnectionFa } }); } + + /** + * A reactive factory that has a lifecycle, which is what the lettuce factory is in production. + */ + private interface DisposableReactiveFactory extends ReactiveRedisConnectionFactory, DisposableBean { + } } diff --git a/shenyu-plugin/shenyu-plugin-ai/shenyu-plugin-ai-sensitive-word/src/main/java/org/apache/shenyu/plugin/ai/sensitive/word/handler/SensitiveWordPluginDataHandler.java b/shenyu-plugin/shenyu-plugin-ai/shenyu-plugin-ai-sensitive-word/src/main/java/org/apache/shenyu/plugin/ai/sensitive/word/handler/SensitiveWordPluginDataHandler.java index 699f83367cca..f155162434aa 100644 --- a/shenyu-plugin/shenyu-plugin-ai/shenyu-plugin-ai-sensitive-word/src/main/java/org/apache/shenyu/plugin/ai/sensitive/word/handler/SensitiveWordPluginDataHandler.java +++ b/shenyu-plugin/shenyu-plugin-ai/shenyu-plugin-ai-sensitive-word/src/main/java/org/apache/shenyu/plugin/ai/sensitive/word/handler/SensitiveWordPluginDataHandler.java @@ -89,19 +89,28 @@ public void handlerPlugin(final PluginData pluginData) { return; } RedisConfigProperties cachedProperties = REDIS_PROPERTIES.get().obtainHandle(PLUGIN_NAME); - if (Objects.isNull(REDIS_TEMPLATES.get().obtainHandle(PLUGIN_NAME)) || !redisConfig.equals(cachedProperties)) { + ReactiveRedisTemplate cachedTemplate = REDIS_TEMPLATES.get().obtainHandle(PLUGIN_NAME); + if (Objects.isNull(cachedTemplate) || !redisConfig.equals(cachedProperties)) { RedisConnectionFactory connectionFactory = new RedisConnectionFactory(redisConfig); ReactiveRedisTemplate redisTemplate = new ShenyuReactiveRedisTemplate<>( connectionFactory.getLettuceConnectionFactory(), ShenyuRedisSerializationContext.stringSerializationContext()); REDIS_TEMPLATES.get().cachedHandle(PLUGIN_NAME, redisTemplate); REDIS_PROPERTIES.get().cachedHandle(PLUGIN_NAME, redisConfig); + // the client that is replaced must not keep its connection pool and its threads alive + if (Objects.nonNull(cachedTemplate)) { + RedisConnectionFactory.destroyQuietly(cachedTemplate.getConnectionFactory()); + } LOG.info("sensitive word plugin: cached the reactive redis template"); } } @Override public void removePlugin(final PluginData pluginData) { + ReactiveRedisTemplate cachedTemplate = REDIS_TEMPLATES.get().obtainHandle(PLUGIN_NAME); + if (Objects.nonNull(cachedTemplate)) { + RedisConnectionFactory.destroyQuietly(cachedTemplate.getConnectionFactory()); + } REDIS_TEMPLATES.get().removeHandle(PLUGIN_NAME); REDIS_PROPERTIES.get().removeHandle(PLUGIN_NAME); LOG.info("sensitive word plugin: released the cached redis template"); diff --git a/shenyu-plugin/shenyu-plugin-ai/shenyu-plugin-ai-sensitive-word/src/test/java/org/apache/shenyu/plugin/ai/sensitive/word/handler/SensitiveWordPluginDataHandlerTest.java b/shenyu-plugin/shenyu-plugin-ai/shenyu-plugin-ai-sensitive-word/src/test/java/org/apache/shenyu/plugin/ai/sensitive/word/handler/SensitiveWordPluginDataHandlerTest.java index 767da72c7d1d..10b2f6698737 100644 --- a/shenyu-plugin/shenyu-plugin-ai/shenyu-plugin-ai-sensitive-word/src/test/java/org/apache/shenyu/plugin/ai/sensitive/word/handler/SensitiveWordPluginDataHandlerTest.java +++ b/shenyu-plugin/shenyu-plugin-ai/shenyu-plugin-ai-sensitive-word/src/test/java/org/apache/shenyu/plugin/ai/sensitive/word/handler/SensitiveWordPluginDataHandlerTest.java @@ -27,10 +27,15 @@ import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.springframework.data.redis.connection.lettuce.LettuceConnectionFactory; +import org.springframework.data.redis.core.ReactiveRedisTemplate; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNotSame; import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertSame; import static org.junit.jupiter.api.Assertions.assertTrue; /** @@ -146,6 +151,44 @@ public void testRemovePluginReleasesTheRedisTemplate() { .obtainHandle(SensitiveWordPluginDataHandler.PLUGIN_NAME)); } + @Test + public void testHandlerPluginDestroysTheClientItReplaces() { + handler.handlerPlugin(pluginData("127.0.0.1:6379")); + ReactiveRedisTemplate first = SensitiveWordPluginDataHandler.REDIS_TEMPLATES.get() + .obtainHandle(SensitiveWordPluginDataHandler.PLUGIN_NAME); + assertNotNull(first); + assertTrue(lettuceFactory(first).isRunning()); + + handler.handlerPlugin(pluginData("127.0.0.1:6380")); + ReactiveRedisTemplate second = SensitiveWordPluginDataHandler.REDIS_TEMPLATES.get() + .obtainHandle(SensitiveWordPluginDataHandler.PLUGIN_NAME); + assertNotSame(first, second); + // the client that was replaced must not keep its connection pool and its threads alive + assertFalse(lettuceFactory(first).isRunning()); + assertTrue(lettuceFactory(second).isRunning()); + } + + @Test + public void testHandlerPluginKeepsTheClientWhenTheConfigurationIsUnchanged() { + handler.handlerPlugin(pluginData("127.0.0.1:6379")); + ReactiveRedisTemplate first = SensitiveWordPluginDataHandler.REDIS_TEMPLATES.get() + .obtainHandle(SensitiveWordPluginDataHandler.PLUGIN_NAME); + handler.handlerPlugin(pluginData("127.0.0.1:6379")); + assertSame(first, SensitiveWordPluginDataHandler.REDIS_TEMPLATES.get() + .obtainHandle(SensitiveWordPluginDataHandler.PLUGIN_NAME)); + assertTrue(lettuceFactory(first).isRunning()); + } + + @Test + public void testRemovePluginDestroysTheClient() { + handler.handlerPlugin(pluginData("127.0.0.1:6379")); + ReactiveRedisTemplate cached = SensitiveWordPluginDataHandler.REDIS_TEMPLATES.get() + .obtainHandle(SensitiveWordPluginDataHandler.PLUGIN_NAME); + assertNotNull(cached); + handler.removePlugin(new PluginData()); + assertFalse(lettuceFactory(cached).isRunning()); + } + @Test public void testCachedDictionaryExpires() { CachedDictionary dictionary = new CachedDictionary(AhoCorasick.empty()); @@ -154,6 +197,17 @@ public void testCachedDictionaryExpires() { assertTrue(!dictionary.isExpired(300L)); } + private PluginData pluginData(final String url) { + PluginData pluginData = new PluginData(); + pluginData.setEnabled(true); + pluginData.setConfig("{\"url\":\"" + url + "\"}"); + return pluginData; + } + + private LettuceConnectionFactory lettuceFactory(final ReactiveRedisTemplate template) { + return (LettuceConnectionFactory) template.getConnectionFactory(); + } + private RuleData ruleData(final String handle) { RuleData ruleData = new RuleData(); ruleData.setId("rule-1"); diff --git a/shenyu-plugin/shenyu-plugin-ai/shenyu-plugin-ai-token-limiter/src/main/java/org/apache/shenyu/plugin/ai/token/limiter/handler/AiTokenLimiterPluginHandler.java b/shenyu-plugin/shenyu-plugin-ai/shenyu-plugin-ai-token-limiter/src/main/java/org/apache/shenyu/plugin/ai/token/limiter/handler/AiTokenLimiterPluginHandler.java index b97738471e03..aba6e5275786 100644 --- a/shenyu-plugin/shenyu-plugin-ai/shenyu-plugin-ai-token-limiter/src/main/java/org/apache/shenyu/plugin/ai/token/limiter/handler/AiTokenLimiterPluginHandler.java +++ b/shenyu-plugin/shenyu-plugin-ai/shenyu-plugin-ai-token-limiter/src/main/java/org/apache/shenyu/plugin/ai/token/limiter/handler/AiTokenLimiterPluginHandler.java @@ -62,12 +62,18 @@ public void handlerPlugin(final PluginData pluginData) { if (Objects.isNull(REDIS_CACHED_HANDLE.get().obtainHandle(PluginEnum.AI_TOKEN_LIMITER.getName())) || Objects.isNull(REDIS_PROPERTIES_CACHED_HANDLE.get().obtainHandle(PluginEnum.AI_TOKEN_LIMITER.getName())) || !redisConfigProperties.equals(REDIS_PROPERTIES_CACHED_HANDLE.get().obtainHandle(PluginEnum.AI_TOKEN_LIMITER.getName()))) { + final ReactiveRedisTemplate previousRedisTemplate = REDIS_CACHED_HANDLE.get() + .obtainHandle(PluginEnum.AI_TOKEN_LIMITER.getName()); final RedisConnectionFactory redisConnectionFactory = new RedisConnectionFactory(redisConfigProperties); ReactiveRedisTemplate reactiveRedisTemplate = new ShenyuReactiveRedisTemplate<>( redisConnectionFactory.getLettuceConnectionFactory(), ShenyuRedisSerializationContext.stringSerializationContext()); REDIS_CACHED_HANDLE.get().cachedHandle(PluginEnum.AI_TOKEN_LIMITER.getName(), reactiveRedisTemplate); REDIS_PROPERTIES_CACHED_HANDLE.get().cachedHandle(PluginEnum.AI_TOKEN_LIMITER.getName(), redisConfigProperties); + // The client that is replaced must not keep its connection pool and its threads alive. + if (Objects.nonNull(previousRedisTemplate)) { + RedisConnectionFactory.destroyQuietly(previousRedisTemplate.getConnectionFactory()); + } } } } diff --git a/shenyu-plugin/shenyu-plugin-ai/shenyu-plugin-ai-token-limiter/src/test/java/org/apache/shenyu/plugin/ai/token/limiter/handler/AiTokenLimiterPluginHandlerTest.java b/shenyu-plugin/shenyu-plugin-ai/shenyu-plugin-ai-token-limiter/src/test/java/org/apache/shenyu/plugin/ai/token/limiter/handler/AiTokenLimiterPluginHandlerTest.java new file mode 100644 index 000000000000..eb7e3000a7ca --- /dev/null +++ b/shenyu-plugin/shenyu-plugin-ai/shenyu-plugin-ai-token-limiter/src/test/java/org/apache/shenyu/plugin/ai/token/limiter/handler/AiTokenLimiterPluginHandlerTest.java @@ -0,0 +1,102 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.shenyu.plugin.ai.token.limiter.handler; + +import org.apache.shenyu.common.dto.PluginData; +import org.apache.shenyu.common.enums.PluginEnum; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; +import org.springframework.data.redis.connection.lettuce.LettuceConnectionFactory; +import org.springframework.data.redis.core.ReactiveRedisTemplate; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNotSame; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Test cases for {@link AiTokenLimiterPluginHandler}. + */ +public final class AiTokenLimiterPluginHandlerTest { + + @AfterEach + public void tearDown() { + AiTokenLimiterPluginHandler.REDIS_CACHED_HANDLE.get().removeHandle(PluginEnum.AI_TOKEN_LIMITER.getName()); + AiTokenLimiterPluginHandler.REDIS_PROPERTIES_CACHED_HANDLE.get() + .removeHandle(PluginEnum.AI_TOKEN_LIMITER.getName()); + } + + @Test + public void testHandlerPluginCachesTheRedisTemplate() { + new AiTokenLimiterPluginHandler().handlerPlugin(pluginData("127.0.0.1:6379")); + ReactiveRedisTemplate template = redisTemplate(); + assertNotNull(template); + assertTrue(lettuceFactory(template).isRunning()); + } + + @Test + public void testHandlerPluginDestroysTheClientItReplaces() { + AiTokenLimiterPluginHandler handler = new AiTokenLimiterPluginHandler(); + handler.handlerPlugin(pluginData("127.0.0.1:6379")); + ReactiveRedisTemplate first = redisTemplate(); + assertNotNull(first); + + handler.handlerPlugin(pluginData("127.0.0.1:6380")); + ReactiveRedisTemplate second = redisTemplate(); + assertNotSame(first, second); + // the client that was replaced must not keep its connection pool and its threads alive + assertFalse(lettuceFactory(first).isRunning()); + assertTrue(lettuceFactory(second).isRunning()); + } + + @Test + public void testHandlerPluginKeepsTheClientWhenTheConfigurationIsUnchanged() { + AiTokenLimiterPluginHandler handler = new AiTokenLimiterPluginHandler(); + handler.handlerPlugin(pluginData("127.0.0.1:6379")); + ReactiveRedisTemplate first = redisTemplate(); + handler.handlerPlugin(pluginData("127.0.0.1:6379")); + assertSame(first, redisTemplate()); + assertTrue(lettuceFactory(first).isRunning()); + } + + @Test + public void testHandlerPluginDisabledDoesNothing() { + PluginData pluginData = pluginData("127.0.0.1:6379"); + pluginData.setEnabled(false); + new AiTokenLimiterPluginHandler().handlerPlugin(pluginData); + assertNull(redisTemplate()); + } + + private ReactiveRedisTemplate redisTemplate() { + return AiTokenLimiterPluginHandler.REDIS_CACHED_HANDLE.get() + .obtainHandle(PluginEnum.AI_TOKEN_LIMITER.getName()); + } + + private LettuceConnectionFactory lettuceFactory(final ReactiveRedisTemplate template) { + return (LettuceConnectionFactory) template.getConnectionFactory(); + } + + private PluginData pluginData(final String url) { + PluginData pluginData = new PluginData(); + pluginData.setEnabled(true); + pluginData.setConfig("{\"url\":\"" + url + "\"}"); + return pluginData; + } +} diff --git a/shenyu-plugin/shenyu-plugin-cache/shenyu-plugin-cache-redis/src/main/java/org/apache/shenyu/plugin/cache/redis/RedisCache.java b/shenyu-plugin/shenyu-plugin-cache/shenyu-plugin-cache-redis/src/main/java/org/apache/shenyu/plugin/cache/redis/RedisCache.java index b5e54e2ee67f..7e2518a20611 100644 --- a/shenyu-plugin/shenyu-plugin-cache/shenyu-plugin-cache-redis/src/main/java/org/apache/shenyu/plugin/cache/redis/RedisCache.java +++ b/shenyu-plugin/shenyu-plugin-cache/shenyu-plugin-cache-redis/src/main/java/org/apache/shenyu/plugin/cache/redis/RedisCache.java @@ -89,5 +89,7 @@ public void close() { connection.close(); } catch (Exception ignored) { } + // the factory owns the connection pool and its threads, closing a connection does not release them + RedisConnectionFactory.destroyQuietly(connectionFactory); } } diff --git a/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/main/java/org/apache/shenyu/plugin/ratelimiter/handler/RateLimiterPluginDataHandler.java b/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/main/java/org/apache/shenyu/plugin/ratelimiter/handler/RateLimiterPluginDataHandler.java index c0203e201429..246ce6dceee9 100644 --- a/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/main/java/org/apache/shenyu/plugin/ratelimiter/handler/RateLimiterPluginDataHandler.java +++ b/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/main/java/org/apache/shenyu/plugin/ratelimiter/handler/RateLimiterPluginDataHandler.java @@ -55,12 +55,17 @@ public void handlerPlugin(final PluginData pluginData) { if (Objects.isNull(Singleton.INST.get(ReactiveRedisTemplate.class)) || Objects.isNull(Singleton.INST.get(RedisConfigProperties.class)) || !redisConfigProperties.equals(Singleton.INST.get(RedisConfigProperties.class))) { + final ReactiveRedisTemplate previousRedisTemplate = Singleton.INST.get(ReactiveRedisTemplate.class); final RedisConnectionFactory redisConnectionFactory = new RedisConnectionFactory(redisConfigProperties); ReactiveRedisTemplate reactiveRedisTemplate = new ShenyuReactiveRedisTemplate<>( redisConnectionFactory.getLettuceConnectionFactory(), ShenyuRedisSerializationContext.stringSerializationContext()); Singleton.INST.single(ReactiveRedisTemplate.class, reactiveRedisTemplate); Singleton.INST.single(RedisConfigProperties.class, redisConfigProperties); + // The client that is replaced must not keep its connection pool and its threads alive. + if (Objects.nonNull(previousRedisTemplate)) { + RedisConnectionFactory.destroyQuietly(previousRedisTemplate.getConnectionFactory()); + } } } } diff --git a/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/test/java/org/apache/shenyu/plugin/ratelimiter/handler/RateLimiterPluginDataHandlerTest.java b/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/test/java/org/apache/shenyu/plugin/ratelimiter/handler/RateLimiterPluginDataHandlerTest.java index c040496f5505..0e46e1a1c7d6 100644 --- a/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/test/java/org/apache/shenyu/plugin/ratelimiter/handler/RateLimiterPluginDataHandlerTest.java +++ b/shenyu-plugin/shenyu-plugin-fault-tolerance/shenyu-plugin-ratelimiter/src/test/java/org/apache/shenyu/plugin/ratelimiter/handler/RateLimiterPluginDataHandlerTest.java @@ -36,14 +36,18 @@ import org.springframework.data.redis.connection.RedisPassword; import org.springframework.data.redis.connection.RedisSentinelConfiguration; import org.springframework.data.redis.connection.RedisStandaloneConfiguration; +import org.springframework.data.redis.connection.lettuce.LettuceConnectionFactory; import org.springframework.data.redis.core.ReactiveRedisTemplate; import org.springframework.test.util.ReflectionTestUtils; import java.util.Collections; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNotSame; import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; /** * RateLimiterPluginDataHandler test. @@ -91,6 +95,31 @@ public void handlerPluginTest() { assertNotNull(Singleton.INST.get(ReactiveRedisTemplate.class)); } + /** + * the client that is replaced must not keep its connection pool and its threads alive. + */ + @Test + public void handlerPluginDestroysTheClientItReplaces() { + RateLimiterPluginDataHandler handler = new RateLimiterPluginDataHandler(); + handler.handlerPlugin(pluginData(generateRedisConfig("localhost:6379"))); + ReactiveRedisTemplate first = Singleton.INST.get(ReactiveRedisTemplate.class); + assertNotNull(first); + assertTrue(((LettuceConnectionFactory) first.getConnectionFactory()).isRunning()); + + handler.handlerPlugin(pluginData(generateRedisConfig("localhost:6380"))); + ReactiveRedisTemplate second = Singleton.INST.get(ReactiveRedisTemplate.class); + assertNotSame(first, second); + assertFalse(((LettuceConnectionFactory) first.getConnectionFactory()).isRunning()); + assertTrue(((LettuceConnectionFactory) second.getConnectionFactory()).isRunning()); + } + + private PluginData pluginData(final RedisConfigProperties redisConfigProperties) { + PluginData pluginData = new PluginData(); + pluginData.setEnabled(true); + pluginData.setConfig(GsonUtils.getInstance().toJson(redisConfigProperties)); + return pluginData; + } + /** * parts parse result null test case. */