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

Reply via email to