This is an automated email from the ASF dual-hosted git repository.
yuluo-yx 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 80f7ad00e8 fix(dubbo): match whole key segments when invalidating
reference caches (#7217)
80f7ad00e8 is described below
commit 80f7ad00e83cac0a2f7a29334b3264cd1bfaa745
Author: Sean-Walker0 <[email protected]>
AuthorDate: Thu Sep 24 15:09:38 2026 +0800
fix(dubbo): match whole key segments when invalidating reference caches
(#7217)
The apache-dubbo and sofa reference caches invalidate entries with a
plain key.contains(id) check. Because ids, registry md5 hashes and
paths are joined into one string, an id that appears as a substring of
an unrelated segment (numeric snowflake ids inside the registry md5
hash, a selector id inside another rule id, one metadata path being a
prefix of another) wrongly invalidates live references of other
routes, destroying in-flight generic calls.
Build the cache keys with a '|' separator and wrap them with the same
separator, then match ids wrapped in '|' so only whole segments hit.
The caches are in-memory only, so the key format change is safe.
Fixes #6798
Co-authored-by: Sean-Walker0
<[email protected]>
---
.../apache/dubbo/cache/ApacheDubboConfigCache.java | 37 +++++++++-------
.../dubbo/cache/ApacheDubboConfigCacheTest.java | 49 +++++++++++++++++++++
.../plugin/sofa/cache/ApplicationConfigCache.java | 32 +++++++-------
.../sofa/cache/ApplicationConfigCacheTest.java | 50 ++++++++++++++++++++++
4 files changed, 137 insertions(+), 31 deletions(-)
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/main/java/org/apache/shenyu/plugin/apache/dubbo/cache/ApacheDubboConfigCache.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/main/java/org/apache/shenyu/plugin/apache/dubbo/cache/ApacheDubboConfigCache.java
index 0f8e5804ab..6ffbd5dfd6 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/main/java/org/apache/shenyu/plugin/apache/dubbo/cache/ApacheDubboConfigCache.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/main/java/org/apache/shenyu/plugin/apache/dubbo/cache/ApacheDubboConfigCache.java
@@ -28,9 +28,7 @@ import java.util.HashMap;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
-import java.util.Set;
import java.util.StringJoiner;
-import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutionException;
import java.util.stream.Collectors;
@@ -67,6 +65,13 @@ public final class ApacheDubboConfigCache extends
DubboConfigCache {
private static final Logger LOG =
LoggerFactory.getLogger(ApacheDubboConfigCache.class);
+ /**
+ * Separator of the reference cache key segments. Never occurs inside ids,
paths,
+ * protocol, registry hash, version or group, so an id can be matched as a
whole
+ * segment by wrapping it with this separator during cache invalidation.
+ */
+ private static final String KEY_SEPARATOR = "|";
+
private ApplicationConfig applicationConfig;
private RegistryConfig registryConfig;
@@ -212,7 +217,7 @@ public final class ApacheDubboConfigCache extends
DubboConfigCache {
* @return the reference config cache key
*/
public String generateUpstreamCacheKey(final String selectorId, final
String ruleId, final String metaDataId, final String namespace, final
DubboUpstream dubboUpstream) {
- StringJoiner stringJoiner = new
StringJoiner(Constants.SEPARATOR_UNDERLINE);
+ StringJoiner stringJoiner = new StringJoiner(KEY_SEPARATOR);
if (StringUtils.isNotBlank(namespace)) {
stringJoiner.add(namespace);
}
@@ -233,7 +238,8 @@ public final class ApacheDubboConfigCache extends
DubboConfigCache {
if (StringUtils.isNotBlank(dubboUpstream.getGroup())) {
stringJoiner.add(dubboUpstream.getGroup());
}
- return stringJoiner.toString();
+ // wrap with separators so that the first and the last segments are
also matched as whole segments
+ return KEY_SEPARATOR + stringJoiner + KEY_SEPARATOR;
}
/**
@@ -529,10 +535,7 @@ public final class ApacheDubboConfigCache extends
DubboConfigCache {
* @param selectorId the selectorId
*/
public void invalidateWithSelectorId(final String selectorId) {
- ConcurrentMap<String, ReferenceConfig<GenericService>> map =
cache.asMap();
- Set<String> allKeys = map.keySet();
- Set<String> needInvalidateKeys = allKeys.stream().filter(key ->
key.contains(selectorId)).collect(Collectors.toSet());
- needInvalidateKeys.forEach(cache::invalidate);
+ invalidateByWholeSegment(selectorId);
}
/**
@@ -541,10 +544,7 @@ public final class ApacheDubboConfigCache extends
DubboConfigCache {
* @param ruleId the ruleId
*/
public void invalidateWithRuleId(final String ruleId) {
- ConcurrentMap<String, ReferenceConfig<GenericService>> map =
cache.asMap();
- Set<String> allKeys = map.keySet();
- Set<String> needInvalidateKeys = allKeys.stream().filter(key ->
key.contains(ruleId)).collect(Collectors.toSet());
- needInvalidateKeys.forEach(cache::invalidate);
+ invalidateByWholeSegment(ruleId);
}
/**
@@ -553,10 +553,15 @@ public final class ApacheDubboConfigCache extends
DubboConfigCache {
* @param metadataId the metadataId
*/
public void invalidateWithMetadataId(final String metadataId) {
- ConcurrentMap<String, ReferenceConfig<GenericService>> map =
cache.asMap();
- Set<String> allKeys = map.keySet();
- Set<String> needInvalidateKeys = allKeys.stream().filter(key ->
key.contains(metadataId)).collect(Collectors.toSet());
- needInvalidateKeys.forEach(cache::invalidate);
+ invalidateByWholeSegment(metadataId);
+ }
+
+ private void invalidateByWholeSegment(final String id) {
+ final String token = KEY_SEPARATOR + id + KEY_SEPARATOR;
+ cache.asMap().keySet().stream()
+ .filter(key -> key.contains(token))
+ .collect(Collectors.toSet())
+ .forEach(cache::invalidate);
}
/**
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/test/java/org/apache/shenyu/plugin/apache/dubbo/cache/ApacheDubboConfigCacheTest.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/test/java/org/apache/shenyu/plugin/apache/dubbo/cache/ApacheDubboConfigCacheTest.java
index 5e6e287059..b11af07d5b 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/test/java/org/apache/shenyu/plugin/apache/dubbo/cache/ApacheDubboConfigCacheTest.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/test/java/org/apache/shenyu/plugin/apache/dubbo/cache/ApacheDubboConfigCacheTest.java
@@ -17,9 +17,13 @@
package org.apache.shenyu.plugin.apache.dubbo.cache;
+import com.google.common.cache.LoadingCache;
+import org.apache.dubbo.config.ReferenceConfig;
import org.apache.dubbo.config.RegistryConfig;
+import org.apache.dubbo.rpc.service.GenericService;
import org.apache.shenyu.common.dto.MetaData;
import org.apache.shenyu.common.dto.convert.plugin.DubboRegisterConfig;
+import org.apache.shenyu.common.dto.convert.selector.DubboUpstream;
import org.apache.shenyu.common.utils.GsonUtils;
import org.apache.shenyu.plugin.dubbo.common.cache.DubboParam;
import org.junit.jupiter.api.BeforeEach;
@@ -31,8 +35,10 @@ import org.mockito.quality.Strictness;
import java.lang.reflect.Field;
+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.assertTrue;
import static org.junit.jupiter.api.Assertions.fail;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
@@ -130,4 +136,47 @@ public final class ApacheDubboConfigCacheTest {
this.apacheDubboConfigCache.invalidate("/test");
this.apacheDubboConfigCache.invalidateAll();
}
+
+ @Test
+ public void testInvalidateMatchesWholeKeySegmentOnly() throws Exception {
+ DubboUpstream dubboUpstream =
DubboUpstream.builder().protocol("dubbo").build();
+ dubboUpstream.setRegistry("zookeeper://127.0.0.1:2181");
+ dubboUpstream.setVersion("1.0.0");
+ dubboUpstream.setGroup("g1");
+ // selector id of keyA is a plain substring of the rule id and of the
other selector id
+ String keyA = apacheDubboConfigCache.generateUpstreamCacheKey("15123",
"9001", "8001", "ns1", dubboUpstream);
+ String keyB =
apacheDubboConfigCache.generateUpstreamCacheKey("91512390", "9151239", "8002",
"ns1", dubboUpstream);
+ String keyC = apacheDubboConfigCache.generateUpstreamCacheKey("70000",
"7001", "8003", "ns1", dubboUpstream);
+ // blank namespace: the selector id becomes the first key segment
+ String keyD = apacheDubboConfigCache.generateUpstreamCacheKey("15123",
"9002", "8004", "", dubboUpstream);
+ String pathBasedKey = "ns1:/some/path";
+ LoadingCache<String, ReferenceConfig<GenericService>> cache =
loadReferenceCache();
+ cache.invalidateAll();
+ cache.put(keyA, new ReferenceConfig<>());
+ cache.put(keyB, new ReferenceConfig<>());
+ cache.put(keyC, new ReferenceConfig<>());
+ cache.put(keyD, new ReferenceConfig<>());
+ cache.put(pathBasedKey, new ReferenceConfig<>());
+
+ apacheDubboConfigCache.invalidateWithSelectorId("15123");
+ assertFalse(cache.asMap().containsKey(keyA));
+ assertFalse(cache.asMap().containsKey(keyD));
+ assertTrue(cache.asMap().containsKey(keyB));
+ assertTrue(cache.asMap().containsKey(keyC));
+ assertTrue(cache.asMap().containsKey(pathBasedKey));
+
+ apacheDubboConfigCache.invalidateWithRuleId("9151239");
+ assertFalse(cache.asMap().containsKey(keyB));
+
+ apacheDubboConfigCache.invalidateWithMetadataId("8003");
+ assertFalse(cache.asMap().containsKey(keyC));
+ assertTrue(cache.asMap().containsKey(pathBasedKey));
+ }
+
+ @SuppressWarnings("unchecked")
+ private LoadingCache<String, ReferenceConfig<GenericService>>
loadReferenceCache() throws Exception {
+ Field field = ApacheDubboConfigCache.class.getDeclaredField("cache");
+ field.setAccessible(true);
+ return (LoadingCache<String, ReferenceConfig<GenericService>>)
field.get(apacheDubboConfigCache);
+ }
}
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-sofa/src/main/java/org/apache/shenyu/plugin/sofa/cache/ApplicationConfigCache.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-sofa/src/main/java/org/apache/shenyu/plugin/sofa/cache/ApplicationConfigCache.java
index 1f9195668e..d01cab2fe5 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-sofa/src/main/java/org/apache/shenyu/plugin/sofa/cache/ApplicationConfigCache.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-sofa/src/main/java/org/apache/shenyu/plugin/sofa/cache/ApplicationConfigCache.java
@@ -69,6 +69,13 @@ public final class ApplicationConfigCache {
private static final Logger LOG =
LoggerFactory.getLogger(ApplicationConfigCache.class);
+ /**
+ * Separator of the reference cache key segments. Never occurs inside ids,
paths,
+ * protocol, registry hash, version or group, so an id can be matched as a
whole
+ * segment by wrapping it with this separator during cache invalidation.
+ */
+ private static final String KEY_SEPARATOR = "|";
+
private static final Map<String, SofaUpstream> UPSTREAM_CACHE_MAP =
Maps.newConcurrentMap();
private final ThreadFactory factory =
ShenyuThreadFactory.create("shenyu-sofa", true);
@@ -230,7 +237,7 @@ public final class ApplicationConfigCache {
* @return the reference config cache key
*/
public String generateUpstreamCacheKey(final String selectorId, final
String metaDataPath, final SofaUpstream sofaUpstream) {
- StringJoiner stringJoiner = new
StringJoiner(Constants.SEPARATOR_UNDERLINE);
+ StringJoiner stringJoiner = new StringJoiner(KEY_SEPARATOR);
stringJoiner.add(selectorId);
stringJoiner.add(metaDataPath);
if (StringUtils.isNotBlank(sofaUpstream.getProtocol())) {
@@ -239,7 +246,8 @@ public final class ApplicationConfigCache {
// use registry hash to short reference cache key
String registryHash = DigestUtils.md5Hex(sofaUpstream.getRegister());
stringJoiner.add(registryHash);
- return stringJoiner.toString();
+ // wrap with separators so that the first and the last segments are
also matched as whole segments
+ return KEY_SEPARATOR + stringJoiner + KEY_SEPARATOR;
}
@@ -434,17 +442,7 @@ public final class ApplicationConfigCache {
* @param metadataPath the metadataPath
*/
public void invalidateWithMetadataPath(final String metadataPath) {
- ConcurrentMap<String, ConsumerConfig<GenericService>> map =
cache.asMap();
- if (map.isEmpty()) {
- return;
- }
- Set<String> allKeys = map.keySet();
- Set<String> needInvalidateKeys = allKeys.stream().filter(key ->
key.contains(metadataPath)).collect(Collectors.toSet());
- if (needInvalidateKeys.isEmpty()) {
- return;
- }
- needInvalidateKeys.forEach(cache::invalidate);
- needInvalidateKeys.forEach(UPSTREAM_CACHE_MAP::remove);
+ invalidateByWholeSegment(metadataPath);
}
/**
@@ -453,12 +451,16 @@ public final class ApplicationConfigCache {
* @param selectorId the selectorId
*/
public void invalidateWithSelectorId(final String selectorId) {
+ invalidateByWholeSegment(selectorId);
+ }
+
+ private void invalidateByWholeSegment(final String segment) {
ConcurrentMap<String, ConsumerConfig<GenericService>> map =
cache.asMap();
if (map.isEmpty()) {
return;
}
- Set<String> allKeys = map.keySet();
- Set<String> needInvalidateKeys = allKeys.stream().filter(key ->
key.contains(selectorId)).collect(Collectors.toSet());
+ final String token = KEY_SEPARATOR + segment + KEY_SEPARATOR;
+ Set<String> needInvalidateKeys = map.keySet().stream().filter(key ->
key.contains(token)).collect(Collectors.toSet());
if (needInvalidateKeys.isEmpty()) {
return;
}
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-sofa/src/test/java/org/apache/shenyu/plugin/sofa/cache/ApplicationConfigCacheTest.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-sofa/src/test/java/org/apache/shenyu/plugin/sofa/cache/ApplicationConfigCacheTest.java
index 3a009deaff..f492bff6c6 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-sofa/src/test/java/org/apache/shenyu/plugin/sofa/cache/ApplicationConfigCacheTest.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-sofa/src/test/java/org/apache/shenyu/plugin/sofa/cache/ApplicationConfigCacheTest.java
@@ -18,6 +18,7 @@
package org.apache.shenyu.plugin.sofa.cache;
import com.alipay.sofa.rpc.config.ConsumerConfig;
+import com.google.common.cache.LoadingCache;
import org.apache.shenyu.common.dto.MetaData;
import org.apache.shenyu.common.dto.SelectorData;
import org.apache.shenyu.common.dto.convert.plugin.SofaRegisterConfig;
@@ -31,8 +32,13 @@ import org.mockito.junit.jupiter.MockitoExtension;
import org.mockito.junit.jupiter.MockitoSettings;
import org.mockito.quality.Strictness;
+import java.lang.reflect.Field;
+import java.util.Map;
+
import static org.junit.Assert.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
@@ -94,4 +100,48 @@ public class ApplicationConfigCacheTest {
assertEquals("127.0.0.1:2182",
refConfig.getRegistry().get(0).getAddress());
}
+ @Test
+ void testInvalidateMatchesWholeKeySegmentOnly() throws Exception {
+ SofaUpstream sofaUpstream = mock(SofaUpstream.class);
+ when(sofaUpstream.getProtocol()).thenReturn("bolt");
+
when(sofaUpstream.getRegister()).thenReturn("zookeeper://127.0.0.1:2181");
+ // selector id of keyA is a plain substring of the other selector id;
+ // path of keyB is a plain prefix of the path of keyA
+ String keyA = cache.generateUpstreamCacheKey("15123", "/sofa/findAll",
sofaUpstream);
+ String keyB = cache.generateUpstreamCacheKey("91512390", "/sofa/find",
sofaUpstream);
+ LoadingCache<String,
ConsumerConfig<com.alipay.sofa.rpc.api.GenericService>> referenceCache =
loadReferenceCache();
+ referenceCache.invalidateAll();
+ referenceCache.put(keyA, new ConsumerConfig<>());
+ referenceCache.put(keyB, new ConsumerConfig<>());
+ Map<String, SofaUpstream> upstreamMap = loadUpstreamMap();
+ upstreamMap.clear();
+ upstreamMap.put(keyA, sofaUpstream);
+ upstreamMap.put(keyB, sofaUpstream);
+
+ cache.invalidateWithSelectorId("15123");
+ assertFalse(referenceCache.asMap().containsKey(keyA));
+ assertTrue(referenceCache.asMap().containsKey(keyB));
+ assertFalse(upstreamMap.containsKey(keyA));
+ assertTrue(upstreamMap.containsKey(keyB));
+
+ // "/sofa/find" must not invalidate the "/sofa/findAll" entry it is a
prefix of
+ cache.invalidateWithMetadataPath("/sofa/find");
+ assertFalse(referenceCache.asMap().containsKey(keyB));
+ assertFalse(upstreamMap.containsKey(keyB));
+ }
+
+ @SuppressWarnings("unchecked")
+ private LoadingCache<String,
ConsumerConfig<com.alipay.sofa.rpc.api.GenericService>> loadReferenceCache()
throws Exception {
+ Field field = ApplicationConfigCache.class.getDeclaredField("cache");
+ field.setAccessible(true);
+ return (LoadingCache<String,
ConsumerConfig<com.alipay.sofa.rpc.api.GenericService>>) field.get(cache);
+ }
+
+ @SuppressWarnings("unchecked")
+ private Map<String, SofaUpstream> loadUpstreamMap() throws Exception {
+ Field field =
ApplicationConfigCache.class.getDeclaredField("UPSTREAM_CACHE_MAP");
+ field.setAccessible(true);
+ return (Map<String, SofaUpstream>) field.get(null);
+ }
+
}