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 7ec00f5733 fix(admin): scope long-poll cache refreshes (#7160)
7ec00f5733 is described below

commit 7ec00f57339eda7f22a8ce1ae731f337f2ef29d4
Author: Liming Deng <[email protected]>
AuthorDate: Wed Sep 30 09:58:28 2026 +0800

    fix(admin): scope long-poll cache refreshes (#7160)
    
    * fix(admin): scope long-poll cache refreshes
    
    * fix(admin): keep long polling responses serialized
    
    ---------
    
    Co-authored-by: aias00 <[email protected]>
---
 .../listener/AbstractDataChangedListener.java      | 26 ++++++++-----
 .../http/HttpLongPollingDataChangedListener.java   | 21 ++++++-----
 .../HttpLongPollingDataChangedListenerTest.java    | 43 +++++++++++++++++++++-
 3 files changed, 69 insertions(+), 21 deletions(-)

diff --git 
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/AbstractDataChangedListener.java
 
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/AbstractDataChangedListener.java
index e546fb71df..3045829b0d 100644
--- 
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/AbstractDataChangedListener.java
+++ 
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/AbstractDataChangedListener.java
@@ -335,18 +335,26 @@ public abstract class AbstractDataChangedListener 
implements DataChangedListener
     protected void refreshLocalCache() {
         List<NamespaceVO> namespaceList = namespaceService.listAll();
         for (NamespaceVO namespace : namespaceList) {
-            String namespaceId = namespace.getNamespaceId();
-            this.updatePluginCache(namespaceId);
-            this.updateAppAuthCache(namespaceId);
-            this.updateRuleCache(namespaceId);
-            this.updateSelectorCache(namespaceId);
-            this.updateMetaDataCache(namespaceId);
-            this.updateProxySelectorDataCache(namespaceId);
-            this.updateDiscoveryUpstreamDataCache(namespaceId);
-            this.updateAiProxyApiKeyCache(namespaceId);
+            this.refreshLocalCache(namespace.getNamespaceId());
         }
     }
 
+    /**
+     * Refresh local cache for one namespace.
+     *
+     * @param namespaceId namespace id
+     */
+    protected void refreshLocalCache(final String namespaceId) {
+        this.updatePluginCache(namespaceId);
+        this.updateAppAuthCache(namespaceId);
+        this.updateRuleCache(namespaceId);
+        this.updateSelectorCache(namespaceId);
+        this.updateMetaDataCache(namespaceId);
+        this.updateProxySelectorDataCache(namespaceId);
+        this.updateDiscoveryUpstreamDataCache(namespaceId);
+        this.updateAiProxyApiKeyCache(namespaceId);
+    }
+
     /**
      * Update selector cache.
      */
diff --git 
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/http/HttpLongPollingDataChangedListener.java
 
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/http/HttpLongPollingDataChangedListener.java
index c269664d18..75ee6860d7 100644
--- 
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/http/HttpLongPollingDataChangedListener.java
+++ 
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/http/HttpLongPollingDataChangedListener.java
@@ -55,7 +55,6 @@ import java.io.IOException;
 import java.util.ArrayList;
 import java.util.Collection;
 import java.util.Collections;
-import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
 import java.util.Objects;
@@ -97,6 +96,8 @@ public class HttpLongPollingDataChangedListener extends 
AbstractDataChangedListe
      */
     private final Map<String, BlockingQueue<LongPollingClient>> clientsMap;
 
+    private final Map<String, Object> refreshLocks;
+
     private final ScheduledExecutorService scheduler;
 
     private final HttpSyncProperties httpSyncProperties;
@@ -108,6 +109,7 @@ public class HttpLongPollingDataChangedListener extends 
AbstractDataChangedListe
      */
     public HttpLongPollingDataChangedListener(final HttpSyncProperties 
httpSyncProperties) {
         this.clientsMap = new ConcurrentHashMap<>();
+        this.refreshLocks = new ConcurrentHashMap<>();
         this.scheduler = new ScheduledThreadPoolExecutor(1,
                 ShenyuThreadFactory.create("long-polling", true));
         this.httpSyncProperties = httpSyncProperties;
@@ -258,12 +260,13 @@ public class HttpLongPollingDataChangedListener extends 
AbstractDataChangedListe
         if (latest != serverCache) {
             return !StringUtils.equals(clientMd5, latest.getMd5());
         }
-        synchronized (this) {
+        Object refreshLock = 
refreshLocks.computeIfAbsent(serverCache.getNamespaceId(), key -> new Object());
+        synchronized (refreshLock) {
             latest = CACHE.get(configDataCacheKey);
             if (latest != serverCache) {
                 return !StringUtils.equals(clientMd5, latest.getMd5());
             }
-            super.refreshLocalCache();
+            this.refreshLocalCache(serverCache.getNamespaceId());
             latest = CACHE.get(configDataCacheKey);
             return !StringUtils.equals(clientMd5, latest.getMd5());
         }
@@ -353,20 +356,18 @@ public class HttpLongPollingDataChangedListener extends 
AbstractDataChangedListe
             if (CollectionUtils.isEmpty(namespaceClients)) {
                 return;
             }
-            if (namespaceClients.size() > 
httpSyncProperties.getNotifyBatchSize()) {
-                List<LongPollingClient> targetClients = new 
ArrayList<>(namespaceClients.size());
-                namespaceClients.drainTo(targetClients);
+            List<LongPollingClient> targetClients = new 
ArrayList<>(namespaceClients.size());
+            namespaceClients.drainTo(targetClients);
+            if (targetClients.size() > 
httpSyncProperties.getNotifyBatchSize()) {
                 List<List<LongPollingClient>> partitionClients = 
Lists.partition(targetClients, httpSyncProperties.getNotifyBatchSize());
                 partitionClients.forEach(item -> scheduler.execute(() -> 
doRun(item)));
             } else {
-                doRun(namespaceClients);
+                doRun(targetClients);
             }
         }
 
         private void doRun(final Collection<LongPollingClient> clients) {
-            for (Iterator<LongPollingClient> iter = clients.iterator(); 
iter.hasNext();) {
-                LongPollingClient client = iter.next();
-                iter.remove();
+            for (LongPollingClient client : clients) {
                 client.sendResponse(Collections.singletonList(groupKey));
                 LOG.info("send response with the changed group,ip={}, 
group={}, changeTime={}", client.ip, groupKey, changeTime);
             }
diff --git 
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/http/HttpLongPollingDataChangedListenerTest.java
 
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/http/HttpLongPollingDataChangedListenerTest.java
index 8fbd127453..a9b33f00d3 100644
--- 
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/http/HttpLongPollingDataChangedListenerTest.java
+++ 
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/http/HttpLongPollingDataChangedListenerTest.java
@@ -347,6 +347,23 @@ public final class HttpLongPollingDataChangedListenerTest {
         getCache().remove(cacheKey);
     }
 
+    @Test
+    public void testClientNewerRefreshesOnlyItsNamespace() throws Exception {
+        RecordingHttpLongPollingDataChangedListener recordingListener =
+                new 
RecordingHttpLongPollingDataChangedListener(httpSyncProperties);
+        String namespaceId = "namespace-one";
+        String group = ConfigGroupEnum.PLUGIN.name();
+        String cacheKey = 
HttpLongPollingDataChangedListener.buildCacheKey(namespaceId, group);
+        ConfigDataCache serverCache = new ConfigDataCache(group, "{}", 
"serverMd5", 1000L, namespaceId);
+        getCache().put(cacheKey, serverCache);
+
+        boolean result = invokeCheckCacheDelayAndUpdate(recordingListener, 
serverCache, "clientMd5", 2000L);
+
+        assertEquals(true, result);
+        assertEquals(namespaceId, recordingListener.refreshedNamespace);
+        getCache().remove(cacheKey);
+    }
+
     /**
      * test doLongPolling with changed groups.
      */
@@ -753,10 +770,16 @@ public final class HttpLongPollingDataChangedListenerTest 
{
 
     private boolean invokeCheckCacheDelayAndUpdate(final ConfigDataCache 
serverCache,
                                                    final String clientMd5, 
final long clientModifyTime) throws Exception {
+        return invokeCheckCacheDelayAndUpdate(listener, serverCache, 
clientMd5, clientModifyTime);
+    }
+
+    private boolean invokeCheckCacheDelayAndUpdate(final 
HttpLongPollingDataChangedListener target,
+                                                   final ConfigDataCache 
serverCache,
+                                                   final String clientMd5, 
final long clientModifyTime) throws Exception {
         Method method = 
HttpLongPollingDataChangedListener.class.getDeclaredMethod(
                 "checkCacheDelayAndUpdate", ConfigDataCache.class, 
String.class, long.class);
         method.setAccessible(true);
-        return (boolean) method.invoke(listener, serverCache, clientMd5, 
clientModifyTime);
+        return (boolean) method.invoke(target, serverCache, clientMd5, 
clientModifyTime);
     }
 
     @SuppressWarnings("unchecked")
@@ -786,4 +809,20 @@ public final class HttpLongPollingDataChangedListenerTest {
         method.setAccessible(true);
         return (String) method.invoke(null, request);
     }
-}
\ No newline at end of file
+
+    private static final class RecordingHttpLongPollingDataChangedListener 
extends HttpLongPollingDataChangedListener {
+
+        private String refreshedNamespace;
+
+        private RecordingHttpLongPollingDataChangedListener(final 
HttpSyncProperties httpSyncProperties) {
+            super(httpSyncProperties);
+        }
+
+        @Override
+        protected void refreshLocalCache(final String namespaceId) {
+            refreshedNamespace = namespaceId;
+            String cacheKey = buildCacheKey(namespaceId, 
ConfigGroupEnum.PLUGIN.name());
+            CACHE.put(cacheKey, new 
ConfigDataCache(ConfigGroupEnum.PLUGIN.name(), "{}", "refreshedMd5", 3000L, 
namespaceId));
+        }
+    }
+}

Reply via email to