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