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 a448dcdb8c Fixes #7356: Reconcile websocket data sync across
standalone admin nodes sharing one database (#7358)
a448dcdb8c is described below
commit a448dcdb8cb448d11ad39400d56634bd660a7255
Author: BobSong <[email protected]>
AuthorDate: Thu Oct 1 11:02:54 2026 +0800
Fixes #7356: Reconcile websocket data sync across standalone admin nodes
sharing one database (#7358)
---
.../admin/config/WebSocketSyncConfiguration.java | 47 +++
.../config/properties/WebsocketSyncProperties.java | 76 ++++-
.../listener/websocket/WebsocketCollector.java | 40 ++-
.../websocket/WebsocketDataReconciler.java | 274 +++++++++++++++++
shenyu-admin/src/main/resources/application.yml | 10 +
.../listener/websocket/WebsocketCollectorTest.java | 11 +
.../WebsocketDataChangedListenerTest.java | 8 +-
.../websocket/WebsocketDataReconcilerTest.java | 323 +++++++++++++++++++++
.../apache/shenyu/common/dto/WebsocketData.java | 41 ++-
.../base/cache/CommonPluginDataSubscriber.java | 26 ++
.../shenyu/sync/data/api/PluginDataSubscriber.java | 31 ++
.../shenyu-sync-data-websocket/pom.xml | 6 +
.../websocket/client/ShenyuWebsocketClient.java | 9 +-
.../websocket/handler/AbstractDataHandler.java | 22 ++
.../data/websocket/handler/PluginDataHandler.java | 9 +
.../data/websocket/handler/RuleDataHandler.java | 9 +
.../websocket/handler/SelectorDataHandler.java | 9 +
.../websocket/handler/WebsocketDataHandler.java | 16 +
.../handler/WebsocketDataHandlerTest.java | 136 +++++++++
19 files changed, 1086 insertions(+), 17 deletions(-)
diff --git
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/config/WebSocketSyncConfiguration.java
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/config/WebSocketSyncConfiguration.java
index a914a493c4..b58aea5ba9 100644
---
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/config/WebSocketSyncConfiguration.java
+++
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/config/WebSocketSyncConfiguration.java
@@ -17,10 +17,22 @@
package org.apache.shenyu.admin.config;
+import org.apache.shenyu.admin.config.properties.ClusterProperties;
import org.apache.shenyu.admin.config.properties.WebsocketSyncProperties;
import org.apache.shenyu.admin.listener.DataChangedListener;
import org.apache.shenyu.admin.listener.websocket.WebsocketCollector;
import org.apache.shenyu.admin.listener.websocket.WebsocketDataChangedListener;
+import org.apache.shenyu.admin.listener.websocket.WebsocketDataReconciler;
+import org.apache.shenyu.admin.mode.cluster.service.ClusterSelectMasterService;
+import org.apache.shenyu.admin.service.AiProxyApiKeyService;
+import org.apache.shenyu.admin.service.AppAuthService;
+import org.apache.shenyu.admin.service.DiscoveryUpstreamService;
+import org.apache.shenyu.admin.service.MetaDataService;
+import org.apache.shenyu.admin.service.NamespacePluginService;
+import org.apache.shenyu.admin.service.ProxySelectorService;
+import org.apache.shenyu.admin.service.RuleService;
+import org.apache.shenyu.admin.service.SelectorService;
+import org.springframework.beans.factory.ObjectProvider;
import
org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import
org.springframework.boot.context.properties.EnableConfigurationProperties;
@@ -58,6 +70,41 @@ public class WebSocketSyncConfiguration {
return new WebsocketCollector();
}
+ /**
+ * Websocket data reconciler, converges gateways connected to this admin
node
+ * with configuration written through other admin nodes sharing the same
database.
+ *
+ * @param websocketSyncProperties the websocket sync properties
+ * @param clusterProperties the cluster properties
+ * @param masterServiceProvider the cluster master service provider
+ * @param appAuthService the app auth service
+ * @param namespacePluginService the namespace plugin service
+ * @param selectorService the selector service
+ * @param ruleService the rule service
+ * @param metaDataService the meta data service
+ * @param proxySelectorService the proxy selector service
+ * @param discoveryUpstreamService the discovery upstream service
+ * @param aiProxyApiKeyService the ai proxy api key service
+ * @return the websocket data reconciler
+ */
+ @Bean
+ @ConditionalOnMissingBean(WebsocketDataReconciler.class)
+ public WebsocketDataReconciler websocketDataReconciler(final
WebsocketSyncProperties websocketSyncProperties,
+ final
ClusterProperties clusterProperties,
+ final
ObjectProvider<ClusterSelectMasterService> masterServiceProvider,
+ final
AppAuthService appAuthService,
+ final
NamespacePluginService namespacePluginService,
+ final
SelectorService selectorService,
+ final RuleService
ruleService,
+ final
MetaDataService metaDataService,
+ final
ProxySelectorService proxySelectorService,
+ final
DiscoveryUpstreamService discoveryUpstreamService,
+ final
AiProxyApiKeyService aiProxyApiKeyService) {
+ return new WebsocketDataReconciler(websocketSyncProperties,
clusterProperties, masterServiceProvider,
+ appAuthService, namespacePluginService, selectorService,
ruleService, metaDataService,
+ proxySelectorService, discoveryUpstreamService,
aiProxyApiKeyService);
+ }
+
/**
* Server endpoint exporter server endpoint exporter.
*
diff --git
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/config/properties/WebsocketSyncProperties.java
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/config/properties/WebsocketSyncProperties.java
index 4dc730ddbd..897c8429f6 100644
---
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/config/properties/WebsocketSyncProperties.java
+++
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/config/properties/WebsocketSyncProperties.java
@@ -19,6 +19,8 @@ package org.apache.shenyu.admin.config.properties;
import org.springframework.boot.context.properties.ConfigurationProperties;
+import java.time.Duration;
+
/**
* the websocket sync strategy properties.
*/
@@ -45,6 +47,8 @@ public class WebsocketSyncProperties {
*/
private String token;
+ private final Reconciliation reconciliation = new Reconciliation();
+
/**
* Gets the value of enabled.
*
@@ -82,7 +86,8 @@ public class WebsocketSyncProperties {
}
/**
- * set allowOrigins.
+ * get allowOrigins.
+ *
* @return allowOrigins
*/
public String getAllowOrigins() {
@@ -90,7 +95,7 @@ public class WebsocketSyncProperties {
}
/**
- * get allowOrigins.
+ * set allowOrigins.
* @param allowOrigins allowOrigins
*/
public void setAllowOrigins(final String allowOrigins) {
@@ -114,4 +119,71 @@ public class WebsocketSyncProperties {
public void setToken(final String token) {
this.token = token;
}
+
+ /**
+ * get reconciliation settings.
+ *
+ * @return reconciliation
+ */
+ public Reconciliation getReconciliation() {
+ return reconciliation;
+ }
+
+ /**
+ * Reconciliation settings for deployments where several standalone admin
nodes
+ * share one database: each node periodically compares a digest of plugin,
selector and rule
+ * configuration groups with the database state and pushes a full refresh
of the
+ * changed groups to the gateway sessions connected to it, so gateways
converge
+ * even when the change was written by another admin node.
+ */
+ public static class Reconciliation {
+
+ /**
+ * Whether reconciliation is enabled, default: false.
+ */
+ private boolean enabled;
+
+ /**
+ * Fixed delay between reconciliation cycles, default: 60s.
+ * Larger values reduce the database polling cost but increase
+ * the worst-case consistency delay for cross-admin changes.
+ */
+ private Duration interval = Duration.ofSeconds(60);
+
+ /**
+ * Whether reconciliation is enabled.
+ *
+ * @return enabled
+ */
+ public boolean isEnabled() {
+ return enabled;
+ }
+
+ /**
+ * Set enabled.
+ *
+ * @param enabled enabled
+ */
+ public void setEnabled(final boolean enabled) {
+ this.enabled = enabled;
+ }
+
+ /**
+ * Gets the fixed delay between reconciliation cycles.
+ *
+ * @return interval
+ */
+ public Duration getInterval() {
+ return interval;
+ }
+
+ /**
+ * Sets the fixed delay between reconciliation cycles.
+ *
+ * @param interval interval
+ */
+ public void setInterval(final Duration interval) {
+ this.interval = interval;
+ }
+ }
}
diff --git
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/websocket/WebsocketCollector.java
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/websocket/WebsocketCollector.java
index e100045639..57570604cf 100644
---
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/websocket/WebsocketCollector.java
+++
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/websocket/WebsocketCollector.java
@@ -50,6 +50,7 @@ import jakarta.websocket.server.ServerEndpoint;
import java.util.ArrayDeque;
import java.util.ArrayList;
+import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
@@ -254,6 +255,15 @@ public class WebsocketCollector {
}
+ /**
+ * Snapshot the namespace ids that currently have at least one registered
session.
+ *
+ * @return the namespace ids with active sessions
+ */
+ public static Set<String> getActiveNamespaceIds() {
+ return Set.copyOf(NAMESPACE_SESSION_MAP.keySet());
+ }
+
/**
* On close.
*
@@ -422,17 +432,31 @@ public class WebsocketCollector {
try {
Map<String, Object> map = JsonUtils.jsonToMap(json);
if (Objects.nonNull(map)) {
- if (map.containsKey("apiKey")) {
- map.put("apiKey", "******");
- }
- if (map.containsKey("realApiKey")) {
- map.put("realApiKey", "******");
- }
+ redactSensitive(map);
return JsonUtils.toJson(map);
}
- return json;
+ return "[unparseable websocket payload]";
} catch (Exception e) {
- return json;
+ return "[unparseable websocket payload]";
+ }
+ }
+
+ private static void redactSensitive(final Object value) {
+ if (value instanceof Map) {
+ Map<?, ?> map = (Map<?, ?>) value;
+ for (Map.Entry<?, ?> entry : map.entrySet()) {
+ if (entry.getKey() instanceof String
+ && ("apiKey".equals(entry.getKey()) ||
"realApiKey".equals(entry.getKey())
+ || "proxyApiKey".equals(entry.getKey()))) {
+ ((Map<Object, Object>) map).put(entry.getKey(), "******");
+ } else {
+ redactSensitive(entry.getValue());
+ }
+ }
+ } else if (value instanceof List) {
+ for (Object item : (List<?>) value) {
+ redactSensitive(item);
+ }
}
}
diff --git
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/websocket/WebsocketDataReconciler.java
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/websocket/WebsocketDataReconciler.java
new file mode 100644
index 0000000000..b4018bc3d1
--- /dev/null
+++
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/websocket/WebsocketDataReconciler.java
@@ -0,0 +1,274 @@
+/*
+ * 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.admin.listener.websocket;
+
+import org.apache.commons.codec.digest.DigestUtils;
+import org.apache.shenyu.admin.config.properties.ClusterProperties;
+import org.apache.shenyu.admin.config.properties.WebsocketSyncProperties;
+import org.apache.shenyu.admin.mode.cluster.service.ClusterSelectMasterService;
+import org.apache.shenyu.admin.service.AiProxyApiKeyService;
+import org.apache.shenyu.admin.service.AppAuthService;
+import org.apache.shenyu.admin.service.DiscoveryUpstreamService;
+import org.apache.shenyu.admin.service.MetaDataService;
+import org.apache.shenyu.admin.service.NamespacePluginService;
+import org.apache.shenyu.admin.service.ProxySelectorService;
+import org.apache.shenyu.admin.service.RuleService;
+import org.apache.shenyu.admin.service.SelectorService;
+import org.apache.shenyu.common.concurrent.ShenyuThreadFactory;
+import org.apache.shenyu.common.dto.WebsocketData;
+import org.apache.shenyu.common.enums.ConfigGroupEnum;
+import org.apache.shenyu.common.enums.DataEventTypeEnum;
+import org.apache.shenyu.common.utils.GsonUtils;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.DisposableBean;
+import org.springframework.beans.factory.InitializingBean;
+import org.springframework.beans.factory.ObjectProvider;
+
+import java.util.Collections;
+import com.google.gson.JsonElement;
+import com.google.gson.JsonParser;
+import java.util.stream.Collectors;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
+import java.util.concurrent.ScheduledThreadPoolExecutor;
+import java.util.concurrent.ThreadLocalRandom;
+import java.util.concurrent.TimeUnit;
+
+/**
+ * Reconciles the configuration pushed through WebSocket data sync with the
state
+ * stored in the database.
+ *
+ * <p>{@code DataChangedEvent} is a local Spring event and {@link
WebsocketCollector}
+ * sessions are local JVM state, so when several standalone admin nodes
(cluster mode
+ * disabled) share one database, a change written through one admin node never
reaches
+ * the gateway sessions connected to the other admin nodes. This task
periodically
+ * compares a digest of the plugin, selector and rule groups per namespace
with the database state
+ * and pushes a full {@link DataEventTypeEnum#REFRESH} of the changed groups
to the
+ * gateway sessions connected to this admin node, so all gateways converge
within the
+ * configured interval without manual synchronization. Reconciliation is
opt-in;
+ * app auth, metadata, proxy selectors, discovery upstreams and AI proxy keys
+ * remain on their existing sync paths and are not reconciled by this task.</p>
+ *
+ * <p>Namespaces without connected gateway sessions are skipped, because there
is
+ * nothing to converge on this node. When cluster mode is enabled, non-master
nodes
+ * skip every cycle and reconciliation is owned by the master node. The digest
cursor
+ * is only advanced after a successful load and push, so failures are retried
in the
+ * next cycle. A full refresh does not modify the database, so a pushed cycle
can
+ * never trigger itself again.</p>
+ */
+public class WebsocketDataReconciler implements InitializingBean,
DisposableBean {
+
+ private static final Logger LOG =
LoggerFactory.getLogger(WebsocketDataReconciler.class);
+
+ private final ScheduledThreadPoolExecutor executor;
+
+ private final WebsocketSyncProperties websocketSyncProperties;
+
+ private final ClusterProperties clusterProperties;
+
+ private final ObjectProvider<ClusterSelectMasterService>
masterServiceProvider;
+
+ private final AppAuthService appAuthService;
+
+ private final NamespacePluginService namespacePluginService;
+
+ private final SelectorService selectorService;
+
+ private final RuleService ruleService;
+
+ private final MetaDataService metaDataService;
+
+ private final ProxySelectorService proxySelectorService;
+
+ private final DiscoveryUpstreamService discoveryUpstreamService;
+
+ private final AiProxyApiKeyService aiProxyApiKeyService;
+
+ /**
+ * Digest of the last pushed configuration group, keyed by {@code
namespaceId:group}.
+ * The cursor is in-memory only: after an admin restart the first cycle
pushes a
+ * full refresh for every connected namespace, which converges any drift
that
+ * accumulated while this admin node was down.
+ */
+ private final ConcurrentMap<String, String> digestCursor = new
ConcurrentHashMap<>();
+
+ public WebsocketDataReconciler(final WebsocketSyncProperties
websocketSyncProperties,
+ final ClusterProperties clusterProperties,
+ final
ObjectProvider<ClusterSelectMasterService> masterServiceProvider,
+ final AppAuthService appAuthService,
+ final NamespacePluginService
namespacePluginService,
+ final SelectorService selectorService,
+ final RuleService ruleService,
+ final MetaDataService metaDataService,
+ final ProxySelectorService
proxySelectorService,
+ final DiscoveryUpstreamService
discoveryUpstreamService,
+ final AiProxyApiKeyService
aiProxyApiKeyService) {
+ this.executor = new ScheduledThreadPoolExecutor(1,
+ ShenyuThreadFactory.create("websocket-reconciliation", true));
+ this.websocketSyncProperties = websocketSyncProperties;
+ this.clusterProperties = clusterProperties;
+ this.masterServiceProvider = masterServiceProvider;
+ this.appAuthService = appAuthService;
+ this.namespacePluginService = namespacePluginService;
+ this.selectorService = selectorService;
+ this.ruleService = ruleService;
+ this.metaDataService = metaDataService;
+ this.proxySelectorService = proxySelectorService;
+ this.discoveryUpstreamService = discoveryUpstreamService;
+ this.aiProxyApiKeyService = aiProxyApiKeyService;
+ }
+
+ @Override
+ public void afterPropertiesSet() {
+ WebsocketSyncProperties.Reconciliation reconciliation =
websocketSyncProperties.getReconciliation();
+ if (!reconciliation.isEnabled()) {
+ LOG.info("websocket data reconciliation is disabled");
+ return;
+ }
+ long intervalMillis = reconciliation.getInterval().toMillis();
+ // jitter the initial delay so several admin nodes sharing one
database do not
+ // poll it simultaneously; scheduleWithFixedDelay also prevents
overlapping runs
+ long jitterMillis = ThreadLocalRandom.current().nextLong(Math.max(1L,
intervalMillis));
+ executor.scheduleWithFixedDelay(this::reconcileSafely, intervalMillis
+ jitterMillis,
+ intervalMillis, TimeUnit.MILLISECONDS);
+ LOG.info("websocket data reconciliation started, interval: {}ms,
initial jitter: {}ms",
+ intervalMillis, jitterMillis);
+ }
+
+ @Override
+ public void destroy() {
+ executor.shutdownNow();
+ }
+
+ /**
+ * Run one reconciliation cycle synchronously, visible for tests.
+ */
+ void reconcileSafely() {
+ try {
+ reconcile();
+ } catch (Exception e) {
+ LOG.error("websocket data reconciliation cycle failed", e);
+ }
+ }
+
+ private void reconcile() {
+ if (skipForClusterNonMaster()) {
+ LOG.debug("websocket data reconciliation skipped: cluster mode
enabled and this node is not the master");
+ return;
+ }
+ Set<String> namespaceIds = activeNamespaceIds();
+ for (String namespaceId : namespaceIds) {
+ for (ConfigGroupEnum group : List.of(ConfigGroupEnum.PLUGIN,
ConfigGroupEnum.SELECTOR, ConfigGroupEnum.RULE)) {
+ reconcileGroup(namespaceId, group);
+ }
+ }
+ }
+
+ private boolean skipForClusterNonMaster() {
+ if (!clusterProperties.isEnabled()) {
+ return false;
+ }
+ ClusterSelectMasterService masterService =
masterServiceProvider.getIfAvailable();
+ return Objects.nonNull(masterService) && !masterService.isMaster();
+ }
+
+ /**
+ * Snapshot the namespaces this node currently serves, visible for tests.
+ *
+ * @return the namespace ids with active websocket sessions
+ */
+ Set<String> activeNamespaceIds() {
+ return WebsocketCollector.getActiveNamespaceIds();
+ }
+
+ private void reconcileGroup(final String namespaceId, final
ConfigGroupEnum group) {
+ String cursorKey = namespaceId + ":" + group.name();
+ try {
+ List<String> rows = load(namespaceId, group).stream()
+ .map(GsonUtils.getInstance()::toJson)
+ .sorted()
+ .collect(Collectors.toList());
+ String digest = DigestUtils.md5Hex(String.join("\n", rows));
+ if (digest.equals(digestCursor.get(cursorKey))) {
+ LOG.debug("websocket reconciliation group {} in namespace {}
is unchanged, skip push",
+ group, namespaceId);
+ return;
+ }
+ List<JsonElement> dataList =
rows.stream().map(JsonParser::parseString).collect(Collectors.toList());
+ WebsocketData<?> websocketData =
+ new WebsocketData<>(group.name(),
DataEventTypeEnum.REFRESH.name(), dataList);
+ websocketData.setNamespaceId(namespaceId);
+ websocketData.setFullSnapshot(true);
+ push(namespaceId, GsonUtils.getInstance().toJson(websocketData));
+ digestCursor.put(cursorKey, digest);
+ LOG.info("websocket reconciliation pushed group {} for namespace
{}, size: {}",
+ group, namespaceId, dataList.size());
+ } catch (Exception e) {
+ // the cursor is not advanced, so the next cycle retries this group
+ LOG.error("websocket reconciliation failed for group {} in
namespace {}, cursor not advanced",
+ group, namespaceId, e);
+ }
+ }
+
+ /**
+ * Push one refresh message to the sessions of a namespace, visible for
tests.
+ *
+ * @param namespaceId the namespace id
+ * @param message the message
+ */
+ void push(final String namespaceId, final String message) {
+ WebsocketCollector.send(namespaceId, message,
DataEventTypeEnum.REFRESH);
+ }
+
+ private List<?> load(final String namespaceId, final ConfigGroupEnum
group) {
+ switch (group) {
+ case APP_AUTH:
+ return appAuthService.listAllByNamespaceId(namespaceId);
+ case PLUGIN:
+ return namespacePluginService.listAll(namespaceId);
+ case RULE:
+ return ruleService.listAllByNamespaceId(namespaceId);
+ case SELECTOR:
+ return selectorService.listAllByNamespaceId(namespaceId);
+ case META_DATA:
+ return metaDataService.listAllByNamespaceId(namespaceId);
+ case PROXY_SELECTOR:
+ return proxySelectorService.listAllByNamespaceId(namespaceId);
+ case DISCOVER_UPSTREAM:
+ return
discoveryUpstreamService.listAllByNamespaceId(namespaceId);
+ case AI_PROXY_API_KEY:
+ return aiProxyApiKeyService.listAllByNamespaceId(namespaceId);
+ default:
+ return Collections.emptyList();
+ }
+ }
+
+ /**
+ * Expose the digest cursor for tests.
+ *
+ * @return the digest cursor
+ */
+ Map<String, String> getDigestCursor() {
+ return digestCursor;
+ }
+}
diff --git a/shenyu-admin/src/main/resources/application.yml
b/shenyu-admin/src/main/resources/application.yml
index f1d503625b..eb976ea4cb 100755
--- a/shenyu-admin/src/main/resources/application.yml
+++ b/shenyu-admin/src/main/resources/application.yml
@@ -72,6 +72,16 @@ shenyu:
messageMaxSize: 10240
token: ${SHENYU_SYNC_WEBSOCKET_TOKEN:}
allowOrigins: ws://localhost:9095;ws://localhost:9195;
+ # Reconciles gateway sessions connected to this admin node with the
shared
+ # database, so gateways converge even when the change was written through
+ # another standalone admin node on the same database. Namespaces without
+ # connected gateway sessions are not polled; a larger interval reduces
the
+ # database polling cost but increases the worst-case consistency delay.
+ # Opt in for shared-database standalone nodes. Only plugin/selector/rule
+ # groups support namespace-scoped replacement; other groups use normal
sync.
+ reconciliation:
+ enabled: false
+ interval: 60s
# apollo:
# meta: http://localhost:8080
# appId: shenyu
diff --git
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/websocket/WebsocketCollectorTest.java
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/websocket/WebsocketCollectorTest.java
index 38d5d7fbb4..0c8e6f00bd 100644
---
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/websocket/WebsocketCollectorTest.java
+++
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/websocket/WebsocketCollectorTest.java
@@ -146,6 +146,17 @@ public final class WebsocketCollectorTest {
}
}
+ @Test
+ public void testNestedApiKeysAreRedacted() {
+ String message = "{\"data\":[{\"proxyApiKey\":\"proxy-secret\","
+ +
"\"nested\":{\"realApiKey\":\"real-secret\",\"apiKey\":\"api-secret\"}}]}";
+ String masked =
ReflectionTestUtils.invokeMethod(WebsocketCollector.class, "maskSensitive",
message);
+ assertFalse(masked.contains("proxy-secret"));
+ assertFalse(masked.contains("real-secret"));
+ assertFalse(masked.contains("api-secret"));
+ assertTrue(masked.contains("******"));
+ }
+
@Test
void testOnOpen() {
websocketCollector.onOpen(session);
diff --git
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/websocket/WebsocketDataChangedListenerTest.java
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/websocket/WebsocketDataChangedListenerTest.java
index 3e1fb5a074..8acd4321c7 100644
---
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/websocket/WebsocketDataChangedListenerTest.java
+++
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/websocket/WebsocketDataChangedListenerTest.java
@@ -84,7 +84,7 @@ public final class WebsocketDataChangedListenerTest {
@Test
public void testOnPluginChanged() {
String message =
"{\"groupType\":\"PLUGIN\",\"eventType\":\"UPDATE\",\"data\":[{\"config\":\"{\\\\\\\"model\\\\\\\":\\\\\\\"black\\\\\\\"}\","
- +
"\"role\":\"1\",\"id\":\"2\",\"name\":\"waf\",\"enabled\":true,\"namespaceId\":\"649330b6-c2d7-4edc-be8e-8a54df9eb385\"}]}";
+ +
"\"role\":\"1\",\"id\":\"2\",\"name\":\"waf\",\"enabled\":true,\"namespaceId\":\"649330b6-c2d7-4edc-be8e-8a54df9eb385\"}],\"fullSnapshot\":false}";
try (MockedStatic<WebsocketCollector> mockedStatic =
mockStatic(WebsocketCollector.class)) {
mockedStatic.when(() -> WebsocketCollector.send(anyString(),
anyString(), any()))
.thenAnswer(invocation -> null);
@@ -172,7 +172,7 @@ public final class WebsocketDataChangedListenerTest {
String message =
"{\"groupType\":\"APP_AUTH\",\"eventType\":\"UPDATE\",\"data\":[{\"appKey\":"
+
"\"D9FD95F496C9495DB5604778A13C3D08\",\"appSecret\":\"02D25048AA1E466F8920E68B08E668DE\","
+
"\"enabled\":true,\"paramDataList\":[{\"appName\":\"axiba\",\"appParam\":\"123\"}]"
- +
",\"pathDataList\":[{\"appName\":\"alibaba\",\"path\":\"/1\",\"enabled\":true}],\"namespaceId\":\"649330b6-c2d7-4edc-be8e-8a54df9eb385\"}]}";
+ +
",\"pathDataList\":[{\"appName\":\"alibaba\",\"path\":\"/1\",\"enabled\":true}],\"namespaceId\":\"649330b6-c2d7-4edc-be8e-8a54df9eb385\"}],\"fullSnapshot\":false}";
try (MockedStatic<WebsocketCollector> mockedStatic =
mockStatic(WebsocketCollector.class)) {
mockedStatic.when(() ->
WebsocketCollector.send(Constants.SYS_DEFAULT_NAMESPACE_ID, message,
DataEventTypeEnum.UPDATE))
.thenAnswer((Answer<Void>) invocation -> null);
@@ -226,7 +226,7 @@ public final class WebsocketDataChangedListenerTest {
public void testOnMetaDataChanged() {
String message =
"{\"groupType\":\"META_DATA\",\"eventType\":\"CREATE\",\"data\":[{\"appName\":\"axiba\","
+
"\"path\":\"/test/execute\",\"rpcType\":\"http\",\"serviceName\":\"execute\",\"methodName\":"
- +
"\"execute\",\"parameterTypes\":\"int\",\"rpcExt\":\"{}\",\"enabled\":true,\"namespaceId\":\"649330b6-c2d7-4edc-be8e-8a54df9eb385\"}]}";
+ +
"\"execute\",\"parameterTypes\":\"int\",\"rpcExt\":\"{}\",\"enabled\":true,\"namespaceId\":\"649330b6-c2d7-4edc-be8e-8a54df9eb385\"}],\"fullSnapshot\":false}";
try (MockedStatic<WebsocketCollector> mockedStatic =
mockStatic(WebsocketCollector.class)) {
mockedStatic.when(() -> WebsocketCollector.send(anyString(),
anyString(), any()))
.thenAnswer(invocation -> null);
@@ -392,7 +392,7 @@ public final class WebsocketDataChangedListenerTest {
private void verifyEmptySnapshot(final MockedStatic<WebsocketCollector>
mockedStatic, final String namespaceId,
final String groupType, final
DataEventTypeEnum eventType) {
String message = String.format(
- "{\"groupType\":\"%s\",\"eventType\":\"%s\",\"data\":[]}",
groupType, eventType.name());
+
"{\"groupType\":\"%s\",\"eventType\":\"%s\",\"data\":[],\"fullSnapshot\":false}",
groupType, eventType.name());
mockedStatic.verify(() -> WebsocketCollector.send(
eq(namespaceId), argThat(actualMsg -> jsonEquals(message,
actualMsg)), eq(eventType)));
}
diff --git
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/websocket/WebsocketDataReconcilerTest.java
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/websocket/WebsocketDataReconcilerTest.java
new file mode 100644
index 0000000000..37e9ece15f
--- /dev/null
+++
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/websocket/WebsocketDataReconcilerTest.java
@@ -0,0 +1,323 @@
+/*
+ * 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.admin.listener.websocket;
+
+import org.apache.shenyu.admin.config.properties.ClusterProperties;
+import org.apache.shenyu.admin.config.properties.WebsocketSyncProperties;
+import org.apache.shenyu.admin.mode.cluster.service.ClusterSelectMasterService;
+import org.apache.shenyu.admin.service.AiProxyApiKeyService;
+import org.apache.shenyu.admin.service.AppAuthService;
+import org.apache.shenyu.admin.service.DiscoveryUpstreamService;
+import org.apache.shenyu.admin.service.MetaDataService;
+import org.apache.shenyu.admin.service.NamespacePluginService;
+import org.apache.shenyu.admin.service.ProxySelectorService;
+import org.apache.shenyu.admin.service.RuleService;
+import org.apache.shenyu.admin.service.SelectorService;
+import org.apache.shenyu.common.dto.AppAuthData;
+import org.apache.shenyu.common.dto.DiscoverySyncData;
+import org.apache.shenyu.common.dto.MetaData;
+import org.apache.shenyu.common.dto.PluginData;
+import org.apache.shenyu.common.dto.ProxyApiKeyData;
+import org.apache.shenyu.common.dto.ProxySelectorData;
+import org.apache.shenyu.common.dto.RuleData;
+import org.apache.shenyu.common.dto.SelectorData;
+import org.apache.shenyu.common.enums.DataEventTypeEnum;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.ArgumentCaptor;
+import org.mockito.Mock;
+import org.mockito.MockedStatic;
+import org.mockito.junit.jupiter.MockitoExtension;
+import org.springframework.beans.factory.ObjectProvider;
+
+import java.util.Collections;
+import java.util.List;
+import java.util.Set;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.atLeast;
+import static org.mockito.Mockito.lenient;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verifyNoInteractions;
+import static org.mockito.Mockito.when;
+
+/**
+ * Tests for {@link WebsocketDataReconciler}.
+ */
+@ExtendWith(MockitoExtension.class)
+public final class WebsocketDataReconcilerTest {
+
+ private static final String NAMESPACE_1 = "namespace-1";
+
+ private static final String NAMESPACE_2 = "namespace-2";
+
+ @Mock
+ private AppAuthService appAuthService;
+
+ @Mock
+ private NamespacePluginService namespacePluginService;
+
+ @Mock
+ private SelectorService selectorService;
+
+ @Mock
+ private RuleService ruleService;
+
+ @Mock
+ private MetaDataService metaDataService;
+
+ @Mock
+ private ProxySelectorService proxySelectorService;
+
+ @Mock
+ private DiscoveryUpstreamService discoveryUpstreamService;
+
+ @Mock
+ private AiProxyApiKeyService aiProxyApiKeyService;
+
+ @Mock
+ private ObjectProvider<ClusterSelectMasterService> masterServiceProvider;
+
+ @Mock
+ private ClusterSelectMasterService clusterSelectMasterService;
+
+ private WebsocketSyncProperties properties;
+
+ private ClusterProperties clusterProperties;
+
+ private WebsocketDataReconciler reconciler;
+
+ @BeforeEach
+ public void setUp() {
+ properties = new WebsocketSyncProperties();
+ clusterProperties = new ClusterProperties();
+ reconciler = new WebsocketDataReconciler(properties,
clusterProperties, masterServiceProvider,
+ appAuthService, namespacePluginService, selectorService,
ruleService,
+ metaDataService, proxySelectorService,
discoveryUpstreamService, aiProxyApiKeyService);
+ }
+
+ @Test
+ public void testReconciliationDefaultsToDisabled() {
+
org.junit.jupiter.api.Assertions.assertFalse(properties.getReconciliation().isEnabled());
+ }
+
+ @Test
+ public void testUnsupportedGroupsAreNotPolled() {
+ stubAllGroups();
+ try (ReconciledCollector ignored = new
ReconciledCollector(Set.of(NAMESPACE_1))) {
+ reconciler.reconcileSafely();
+ verifyNoInteractions(appAuthService, metaDataService,
proxySelectorService,
+ discoveryUpstreamService, aiProxyApiKeyService);
+ }
+ }
+
+ @Test
+ public void testLifecycleMethodsDoNotThrow() {
+ properties.getReconciliation().setEnabled(false);
+ assertDoesNotThrow(() -> reconciler.afterPropertiesSet());
+ reconciler.destroy();
+ // a fresh reconciler starts and stops cleanly with scheduling enabled
+ WebsocketDataReconciler started = new
WebsocketDataReconciler(properties, clusterProperties,
+ masterServiceProvider, appAuthService, namespacePluginService,
selectorService, ruleService,
+ metaDataService, proxySelectorService,
discoveryUpstreamService, aiProxyApiKeyService);
+ properties.getReconciliation().setEnabled(true);
+ assertDoesNotThrow(started::afterPropertiesSet);
+ started.destroy();
+ }
+
+ @Test
+ public void testNoActiveSessionsSkipsAllLoads() {
+ try (ReconciledCollector ignored = new
ReconciledCollector(Collections.emptySet())) {
+ reconciler.reconcileSafely();
+ verifyNoInteractions(appAuthService, namespacePluginService,
selectorService, ruleService,
+ metaDataService, proxySelectorService,
discoveryUpstreamService, aiProxyApiKeyService);
+ }
+ }
+
+ @Test
+ public void testChangedGroupsArePushedOncePerCycle() {
+ stubAllGroups();
+ try (ReconciledCollector mocked = new
ReconciledCollector(Set.of(NAMESPACE_1))) {
+ reconciler.reconcileSafely();
+ mocked.verifySends(NAMESPACE_1, 3);
+ // unchanged state: the second cycle pushes nothing
+ reconciler.reconcileSafely();
+ mocked.verifySends(NAMESPACE_1, 3);
+ }
+ }
+
+ @Test
+ public void testEmptyGroupIsPushedForClearing() {
+ stubAllGroups();
+
when(ruleService.listAllByNamespaceId(NAMESPACE_1)).thenReturn(Collections.emptyList());
+ try (ReconciledCollector mocked = new
ReconciledCollector(Set.of(NAMESPACE_1))) {
+ reconciler.reconcileSafely();
+ List<String> messages = mocked.capturedMessages(NAMESPACE_1);
+ assertEquals(3, messages.size());
+ String ruleMessage = messages.stream()
+ .filter(m -> m.contains("\"groupType\":\"RULE\""))
+ .findFirst()
+ .orElse("");
+ assertTrue(ruleMessage.contains("\"eventType\":\"REFRESH\""),
ruleMessage);
+ assertTrue(ruleMessage.contains("\"data\":[]"), ruleMessage);
+ }
+ }
+
+ @Test
+ public void testFailedLoadIsRetriedNextCycle() {
+ stubAllGroups();
+ when(ruleService.listAllByNamespaceId(NAMESPACE_1))
+ .thenThrow(new RuntimeException("db down"))
+ .thenReturn(Collections.singletonList(new
RuleData().setId("rule-1")));
+ try (ReconciledCollector mocked = new
ReconciledCollector(Set.of(NAMESPACE_1))) {
+ // the failed group is not pushed and does not fail the whole cycle
+ reconciler.reconcileSafely();
+ mocked.verifySends(NAMESPACE_1, 2);
+ // the cursor was not advanced, the next cycle retries the group
+ reconciler.reconcileSafely();
+ mocked.verifySends(NAMESPACE_1, 3);
+ assertEquals(1, mocked.capturedMessages(NAMESPACE_1).stream()
+ .filter(m -> m.contains("\"groupType\":\"RULE\""))
+ .count());
+ }
+ }
+
+ @Test
+ public void testOnlyChangedNamespaceIsRepulsed() {
+ stubAllGroups();
+ when(ruleService.listAllByNamespaceId(NAMESPACE_1))
+ .thenReturn(Collections.singletonList(new
RuleData().setId("rule-v1")))
+ .thenReturn(Collections.singletonList(new
RuleData().setId("rule-v2")));
+ when(ruleService.listAllByNamespaceId(NAMESPACE_2))
+ .thenReturn(Collections.singletonList(new
RuleData().setId("rule-ns2")));
+ try (ReconciledCollector mocked = new
ReconciledCollector(Set.of(NAMESPACE_1, NAMESPACE_2))) {
+ reconciler.reconcileSafely();
+ mocked.verifySends(NAMESPACE_1, 3);
+ mocked.verifySends(NAMESPACE_2, 3);
+ // only the changed group of namespace-1 is pushed in the second
cycle
+ reconciler.reconcileSafely();
+ mocked.verifySends(NAMESPACE_1, 4);
+ mocked.verifySends(NAMESPACE_2, 3);
+ }
+ }
+
+ @Test
+ public void testClusterNonMasterSkipsCycle() {
+ clusterProperties.setEnabled(true);
+
when(masterServiceProvider.getIfAvailable()).thenReturn(clusterSelectMasterService);
+ when(clusterSelectMasterService.isMaster()).thenReturn(false);
+ stubAllGroups();
+ try (ReconciledCollector mocked = new
ReconciledCollector(Set.of(NAMESPACE_1))) {
+ reconciler.reconcileSafely();
+ mocked.verifyNoSends();
+ }
+ }
+
+ @Test
+ public void testClusterMasterReconciles() {
+ clusterProperties.setEnabled(true);
+
when(masterServiceProvider.getIfAvailable()).thenReturn(clusterSelectMasterService);
+ when(clusterSelectMasterService.isMaster()).thenReturn(true);
+ stubAllGroups();
+ try (ReconciledCollector mocked = new
ReconciledCollector(Set.of(NAMESPACE_1))) {
+ reconciler.reconcileSafely();
+ mocked.verifySends(NAMESPACE_1, 3);
+ }
+ }
+
+ @Test
+ public void testReorderingDoesNotPushAgain() {
+ stubAllGroups();
+ RuleData first = new RuleData().setId("first");
+ RuleData second = new RuleData().setId("second");
+ when(ruleService.listAllByNamespaceId(NAMESPACE_1))
+ .thenReturn(List.of(first, second)).thenReturn(List.of(second,
first));
+ try (ReconciledCollector mocked = new
ReconciledCollector(Set.of(NAMESPACE_1))) {
+ reconciler.reconcileSafely();
+ reconciler.reconcileSafely();
+ mocked.verifySends(NAMESPACE_1, 3);
+ assertTrue(mocked.capturedMessages(NAMESPACE_1).stream()
+ .allMatch(message ->
message.contains("\"fullSnapshot\":true")
+ &&
message.contains("\"namespaceId\":\"namespace-1\"")));
+ }
+ }
+
+ private void stubAllGroups() {
+ lenient().when(appAuthService.listAllByNamespaceId(anyString()))
+ .thenReturn(Collections.singletonList(new AppAuthData()));
+ lenient().when(namespacePluginService.listAll(anyString()))
+ .thenReturn(Collections.singletonList(new PluginData()));
+ lenient().when(selectorService.listAllByNamespaceId(anyString()))
+ .thenReturn(Collections.singletonList(new SelectorData()));
+ lenient().when(ruleService.listAllByNamespaceId(anyString()))
+ .thenReturn(Collections.singletonList(new RuleData()));
+ lenient().when(metaDataService.listAllByNamespaceId(anyString()))
+ .thenReturn(Collections.singletonList(new MetaData()));
+ lenient().when(proxySelectorService.listAllByNamespaceId(anyString()))
+ .thenReturn(Collections.singletonList(new
ProxySelectorData()));
+
lenient().when(discoveryUpstreamService.listAllByNamespaceId(anyString()))
+ .thenReturn(Collections.singletonList(new
DiscoverySyncData()));
+ lenient().when(aiProxyApiKeyService.listAllByNamespaceId(anyString()))
+ .thenReturn(Collections.singletonList(new ProxyApiKeyData()));
+ }
+
+ /**
+ * Wrapper around the static {@link WebsocketCollector} mock that records
and
+ * verifies the refresh messages pushed by the reconciler.
+ */
+ private static final class ReconciledCollector implements AutoCloseable {
+
+ private final MockedStatic<WebsocketCollector> mocked;
+
+ private final ArgumentCaptor<String> messageCaptor;
+
+ private ReconciledCollector(final Set<String> activeNamespaces) {
+ this.mocked = mockStatic(WebsocketCollector.class);
+ this.messageCaptor = ArgumentCaptor.forClass(String.class);
+
mocked.when(WebsocketCollector::getActiveNamespaceIds).thenReturn(activeNamespaces);
+ }
+
+ private void verifySends(final String namespaceId, final long
expected) {
+ mocked.verify(() -> WebsocketCollector.send(eq(namespaceId),
anyString(),
+ eq(DataEventTypeEnum.REFRESH)), times((int) expected));
+ }
+
+ private void verifyNoSends() {
+ mocked.verify(() -> WebsocketCollector.send(anyString(),
anyString(), any()), never());
+ }
+
+ private List<String> capturedMessages(final String namespaceId) {
+ mocked.verify(() -> WebsocketCollector.send(eq(namespaceId),
messageCaptor.capture(),
+ eq(DataEventTypeEnum.REFRESH)), atLeast(0));
+ return messageCaptor.getAllValues();
+ }
+
+ @Override
+ public void close() {
+ mocked.close();
+ }
+ }
+}
diff --git
a/shenyu-common/src/main/java/org/apache/shenyu/common/dto/WebsocketData.java
b/shenyu-common/src/main/java/org/apache/shenyu/common/dto/WebsocketData.java
index beb690c723..fbd749a7a5 100644
---
a/shenyu-common/src/main/java/org/apache/shenyu/common/dto/WebsocketData.java
+++
b/shenyu-common/src/main/java/org/apache/shenyu/common/dto/WebsocketData.java
@@ -49,6 +49,10 @@ public class WebsocketData<T> {
*/
private List<T> data;
+ private String namespaceId;
+
+ private boolean fullSnapshot;
+
/**
* no args constructor.
*/
@@ -68,6 +72,38 @@ public class WebsocketData<T> {
this.data = data;
}
+ /**
+ * Get the namespace of an authoritative snapshot.
+ * @return namespace id
+ */
+ public String getNamespaceId() {
+ return namespaceId;
+ }
+
+ /**
+ * Set the snapshot namespace.
+ * @param namespaceId namespace id
+ */
+ public void setNamespaceId(final String namespaceId) {
+ this.namespaceId = namespaceId;
+ }
+
+ /**
+ * Whether this message replaces the entire group in the namespace.
+ * @return whether the snapshot is complete
+ */
+ public boolean isFullSnapshot() {
+ return fullSnapshot;
+ }
+
+ /**
+ * Mark a complete namespace snapshot.
+ * @param fullSnapshot whether the snapshot is complete
+ */
+ public void setFullSnapshot(final boolean fullSnapshot) {
+ this.fullSnapshot = fullSnapshot;
+ }
+
/**
* get groupType.
*
@@ -137,12 +173,13 @@ public class WebsocketData<T> {
return false;
}
WebsocketData<?> that = (WebsocketData<?>) o;
- return Objects.equals(groupType, that.groupType) &&
Objects.equals(eventType, that.eventType) && Objects.equals(data, that.data);
+ return Objects.equals(groupType, that.groupType) &&
Objects.equals(eventType, that.eventType) && Objects.equals(data, that.data)
+ && Objects.equals(namespaceId, that.namespaceId) &&
fullSnapshot == that.fullSnapshot;
}
@Override
public int hashCode() {
- return Objects.hash(groupType, eventType, data);
+ return Objects.hash(groupType, eventType, data, namespaceId,
fullSnapshot);
}
@Override
diff --git
a/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriber.java
b/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriber.java
index 0ff78bc865..a8526ed151 100644
---
a/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriber.java
+++
b/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriber.java
@@ -119,6 +119,32 @@ public class CommonPluginDataSubscriber implements
PluginDataSubscriber {
subscribeDataHandler(pluginData, DataEventTypeEnum.DELETE);
}
+ @Override
+ public void refreshPluginDataNamespace(final String namespaceId) {
+ List<PluginData> stale =
BaseDataCache.getInstance().getPluginMap().values().stream()
+ .filter(data -> namespaceId.equals(data.getNamespaceId()))
+ .collect(Collectors.toList());
+ stale.forEach(this::unSubscribe);
+ }
+
+ @Override
+ public void refreshSelectorDataNamespace(final String namespaceId) {
+ List<SelectorData> stale =
BaseDataCache.getInstance().getSelectorMap().values().stream()
+ .flatMap(List::stream)
+ .filter(data -> namespaceId.equals(data.getNamespaceId()))
+ .collect(Collectors.toList());
+ stale.forEach(this::unSelectorSubscribe);
+ }
+
+ @Override
+ public void refreshRuleDataNamespace(final String namespaceId) {
+ List<RuleData> stale =
BaseDataCache.getInstance().getRuleMap().values().stream()
+ .flatMap(List::stream)
+ .filter(data -> namespaceId.equals(data.getNamespaceId()))
+ .collect(Collectors.toList());
+ stale.forEach(this::unRuleSubscribe);
+ }
+
@Override
public void refreshPluginDataAll() {
BaseDataCache.getInstance().cleanPluginData();
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/api/PluginDataSubscriber.java
b/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/api/PluginDataSubscriber.java
index ef52e2fd9a..10026fdb3d 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/api/PluginDataSubscriber.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/api/PluginDataSubscriber.java
@@ -146,4 +146,35 @@ public interface PluginDataSubscriber {
refreshRuleDataSelf(dataList);
dataList.forEach(this::onRuleSubscribe);
}
+
+ /**
+ * Remove cached plugin rows belonging to one namespace before applying a
snapshot.
+ * Custom subscribers must implement this operation before enabling
reconciliation.
+ *
+ * @param namespaceId namespace to replace
+ */
+ default void refreshPluginDataNamespace(final String namespaceId) {
+ throw new UnsupportedOperationException("Namespace-scoped plugin
snapshots are not supported");
+ }
+
+ /**
+ * Remove cached selector rows belonging to one namespace before applying
a snapshot.
+ * Custom subscribers must implement this operation before enabling
reconciliation.
+ *
+ * @param namespaceId namespace to replace
+ */
+ default void refreshSelectorDataNamespace(final String namespaceId) {
+ throw new UnsupportedOperationException("Namespace-scoped selector
snapshots are not supported");
+ }
+
+ /**
+ * Remove cached rule rows belonging to one namespace before applying a
snapshot.
+ * Custom subscribers must implement this operation before enabling
reconciliation.
+ *
+ * @param namespaceId namespace to replace
+ */
+ default void refreshRuleDataNamespace(final String namespaceId) {
+ throw new UnsupportedOperationException("Namespace-scoped rule
snapshots are not supported");
+ }
+
}
diff --git a/shenyu-sync-data-center/shenyu-sync-data-websocket/pom.xml
b/shenyu-sync-data-center/shenyu-sync-data-websocket/pom.xml
index d1d3ecc8b4..ab2ed12e24 100644
--- a/shenyu-sync-data-center/shenyu-sync-data-websocket/pom.xml
+++ b/shenyu-sync-data-center/shenyu-sync-data-websocket/pom.xml
@@ -26,6 +26,12 @@
<artifactId>shenyu-sync-data-websocket</artifactId>
<dependencies>
+ <dependency>
+ <groupId>org.apache.shenyu</groupId>
+ <artifactId>shenyu-plugin-base</artifactId>
+ <version>${project.version}</version>
+ <scope>test</scope>
+ </dependency>
<dependency>
<groupId>org.apache.shenyu</groupId>
<artifactId>shenyu-sync-data-api</artifactId>
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/client/ShenyuWebsocketClient.java
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/client/ShenyuWebsocketClient.java
index 35f9b23253..1d4fe3bb9a 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/client/ShenyuWebsocketClient.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/client/ShenyuWebsocketClient.java
@@ -425,7 +425,14 @@ public final class ShenyuWebsocketClient extends
WebSocketClient {
ConfigGroupEnum groupEnum =
ConfigGroupEnum.acquireByName(websocketData.getGroupType());
String eventType = websocketData.getEventType();
String json = GsonUtils.getInstance().toJson(websocketData.getData());
- websocketDataHandler.executor(groupEnum, json, eventType);
+ if (websocketData.isFullSnapshot()) {
+ if (!DataEventTypeEnum.REFRESH.name().equals(eventType) &&
!DataEventTypeEnum.MYSELF.name().equals(eventType)) {
+ throw new IllegalArgumentException("Snapshot requires a
refresh event");
+ }
+ websocketDataHandler.snapshot(groupEnum, json,
websocketData.getNamespaceId(), namespaceId);
+ } else {
+ websocketDataHandler.executor(groupEnum, json, eventType);
+ }
}
/**
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/AbstractDataHandler.java
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/AbstractDataHandler.java
index a75e443ee1..eb8c55bf1b 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/AbstractDataHandler.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/AbstractDataHandler.java
@@ -59,6 +59,28 @@ public abstract class AbstractDataHandler<T> implements
DataHandler {
*/
protected abstract void doDelete(List<T> dataList);
+ /**
+ * Apply an authoritative snapshot for the connection's namespace.
+ * @param json snapshot array, including an empty array
+ * @param namespaceId namespace to replace
+ */
+ public void handleSnapshot(final String json, final String namespaceId) {
+ List<T> dataList = convert(json);
+ if (java.util.Objects.isNull(dataList)) {
+ throw new IllegalArgumentException("A snapshot must contain a data
array");
+ }
+ doSnapshot(dataList, namespaceId);
+ }
+
+ /**
+ * Replace the complete group, not just rows present in the payload.
+ * @param dataList complete group
+ * @param namespaceId namespace to replace
+ */
+ protected void doSnapshot(final List<T> dataList, final String
namespaceId) {
+ throw new IllegalArgumentException("Namespace snapshots are supported
only for plugins, selectors and rules");
+ }
+
@Override
public void handle(final String json, final String eventType) {
List<T> dataList = convert(json);
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/PluginDataHandler.java
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/PluginDataHandler.java
index 719414fe54..055c95009d 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/PluginDataHandler.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/PluginDataHandler.java
@@ -45,6 +45,15 @@ public class PluginDataHandler extends
AbstractDataHandler<PluginData> {
pluginDataSubscriber.onPluginRefresh(dataList);
}
+ @Override
+ protected void doSnapshot(final List<PluginData> dataList, final String
namespaceId) {
+ if (dataList.stream().anyMatch(data ->
!namespaceId.equals(data.getNamespaceId()))) {
+ throw new IllegalArgumentException("Snapshot row namespace does
not match the connection");
+ }
+ pluginDataSubscriber.refreshPluginDataNamespace(namespaceId);
+ doUpdate(dataList);
+ }
+
@Override
protected void doUpdate(final List<PluginData> dataList) {
dataList.forEach(pluginDataSubscriber::onSubscribe);
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/RuleDataHandler.java
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/RuleDataHandler.java
index a447b19ee9..0fbdbed978 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/RuleDataHandler.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/RuleDataHandler.java
@@ -44,6 +44,15 @@ public class RuleDataHandler extends
AbstractDataHandler<RuleData> {
pluginDataSubscriber.onRuleRefresh(dataList);
}
+ @Override
+ protected void doSnapshot(final List<RuleData> dataList, final String
namespaceId) {
+ if (dataList.stream().anyMatch(data ->
!namespaceId.equals(data.getNamespaceId()))) {
+ throw new IllegalArgumentException("Snapshot row namespace does
not match the connection");
+ }
+ pluginDataSubscriber.refreshRuleDataNamespace(namespaceId);
+ doUpdate(dataList);
+ }
+
@Override
protected void doUpdate(final List<RuleData> dataList) {
dataList.forEach(pluginDataSubscriber::onRuleSubscribe);
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/SelectorDataHandler.java
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/SelectorDataHandler.java
index 7db231633e..3b142c49c6 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/SelectorDataHandler.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/SelectorDataHandler.java
@@ -45,6 +45,15 @@ public class SelectorDataHandler extends
AbstractDataHandler<SelectorData> {
pluginDataSubscriber.onSelectorRefresh(dataList);
}
+ @Override
+ protected void doSnapshot(final List<SelectorData> dataList, final String
namespaceId) {
+ if (dataList.stream().anyMatch(data ->
!namespaceId.equals(data.getNamespaceId()))) {
+ throw new IllegalArgumentException("Snapshot row namespace does
not match the connection");
+ }
+ pluginDataSubscriber.refreshSelectorDataNamespace(namespaceId);
+ doUpdate(dataList);
+ }
+
@Override
protected void doUpdate(final List<SelectorData> dataList) {
dataList.forEach(pluginDataSubscriber::onSelectorSubscribe);
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/WebsocketDataHandler.java
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/WebsocketDataHandler.java
index f3c5fe91bb..2acf0fa3cb 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/WebsocketDataHandler.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/WebsocketDataHandler.java
@@ -68,4 +68,20 @@ public class WebsocketDataHandler {
handlers.get(type).handle(json, eventType);
}
+ /**
+ * Apply a complete group snapshot after verifying the connection
namespace.
+ * @param type configuration group
+ * @param json snapshot array
+ * @param snapshotNamespace namespace supplied by Admin
+ * @param connectionNamespace namespace configured on this connection
+ */
+ public void snapshot(final ConfigGroupEnum type, final String json,
+ final String snapshotNamespace, final String
connectionNamespace) {
+ if (java.util.Objects.isNull(snapshotNamespace) ||
snapshotNamespace.isEmpty()
+ || !snapshotNamespace.equals(connectionNamespace)) {
+ throw new IllegalArgumentException("Snapshot namespace does not
match the connection");
+ }
+ ((AbstractDataHandler<?>) handlers.get(type)).handleSnapshot(json,
snapshotNamespace);
+ }
+
}
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/WebsocketDataHandlerTest.java
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/WebsocketDataHandlerTest.java
index afec30bcee..9dd8efaec6 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/WebsocketDataHandlerTest.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/WebsocketDataHandlerTest.java
@@ -20,10 +20,16 @@ package
org.apache.shenyu.plugin.sync.data.websocket.handler;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
+import java.util.Collections;
import java.util.LinkedList;
import java.util.List;
+import org.apache.shenyu.common.config.ShenyuConfig;
import org.apache.shenyu.common.dto.PluginData;
+import org.apache.shenyu.common.dto.RuleData;
+import org.apache.shenyu.common.dto.SelectorData;
+import org.apache.shenyu.plugin.base.cache.BaseDataCache;
+import org.apache.shenyu.plugin.base.cache.CommonPluginDataSubscriber;
import org.apache.shenyu.common.enums.ConfigGroupEnum;
import org.apache.shenyu.common.enums.DataEventTypeEnum;
import org.apache.shenyu.common.utils.GsonUtils;
@@ -102,11 +108,141 @@ public final class WebsocketDataHandlerTest {
pluginDataList.forEach(verify(pluginDataSubscriber)::unSubscribe);
}
+ @Test
+ public void testEmptySnapshotClearsOnlyItsGroup() {
+ websocketDataHandler.snapshot(ConfigGroupEnum.RULE, "[]",
"namespace-a", "namespace-a");
+ verify(pluginDataSubscriber).refreshRuleDataNamespace("namespace-a");
+ Mockito.verifyNoMoreInteractions(pluginDataSubscriber);
+ }
+
+ @Test
+ public void testWrongNamespaceCannotClearCache() {
+
org.junit.jupiter.api.Assertions.assertThrows(IllegalArgumentException.class,
+ () -> websocketDataHandler.snapshot(ConfigGroupEnum.PLUGIN,
"[]", "namespace-b", "namespace-a"));
+ Mockito.verifyNoInteractions(pluginDataSubscriber);
+ }
+
+ @Test
+ public void testNullSnapshotCannotClearCache() {
+
org.junit.jupiter.api.Assertions.assertThrows(IllegalArgumentException.class,
+ () -> websocketDataHandler.snapshot(ConfigGroupEnum.PLUGIN,
"null", "namespace-a", "namespace-a"));
+ Mockito.verifyNoInteractions(pluginDataSubscriber);
+ }
+
+ @Test
+ public void testSnapshotReplacesStalePluginBeforeSubscribing() {
+ websocketDataHandler.snapshot(ConfigGroupEnum.PLUGIN, getJson(),
"namespace-a", "namespace-a");
+ org.mockito.InOrder order = Mockito.inOrder(pluginDataSubscriber);
+
order.verify(pluginDataSubscriber).refreshPluginDataNamespace("namespace-a");
+
order.verify(pluginDataSubscriber).onSubscribe(Mockito.any(PluginData.class));
+ }
+
+ @Test
+ public void testEmptySnapshotsRemoveStaleGatewayCache() {
+ BaseDataCache cache = BaseDataCache.getInstance();
+ PluginDataSubscriber subscriber = new
CommonPluginDataSubscriber(Collections.emptyList(),
+ new ShenyuConfig.SelectorMatchCache(), new
ShenyuConfig.RuleMatchCache());
+ WebsocketDataHandler handler = new WebsocketDataHandler(subscriber,
Collections.emptyList(), Collections.emptyList(),
+ Collections.emptyList(), Collections.emptyList(),
Collections.emptyList());
+
cache.cachePluginData(PluginData.builder().name("snapshot-plugin").namespaceId("namespace-a").build());
+
cache.cacheSelectData(SelectorData.builder().id("snapshot-selector").pluginName("snapshot-plugin").namespaceId("namespace-a").sort(1).build());
+
cache.cacheRuleData(RuleData.builder().id("snapshot-rule").selectorId("snapshot-selector").namespaceId("namespace-a").sort(1).build());
+ try {
+ handler.snapshot(ConfigGroupEnum.RULE, "[]", "namespace-a",
"namespace-a");
+
org.junit.jupiter.api.Assertions.assertNull(cache.obtainRuleData("snapshot-selector"));
+
org.junit.jupiter.api.Assertions.assertNotNull(cache.obtainSelectorData("snapshot-plugin"));
+ handler.snapshot(ConfigGroupEnum.SELECTOR, "[]", "namespace-a",
"namespace-a");
+
org.junit.jupiter.api.Assertions.assertNull(cache.obtainSelectorData("snapshot-plugin"));
+
org.junit.jupiter.api.Assertions.assertNotNull(cache.obtainPluginData("snapshot-plugin"));
+ handler.snapshot(ConfigGroupEnum.PLUGIN, "[]", "namespace-a",
"namespace-a");
+
org.junit.jupiter.api.Assertions.assertNull(cache.obtainPluginData("snapshot-plugin"));
+ } finally {
+ cache.cleanRuleData();
+ cache.cleanSelectorData();
+ cache.cleanPluginData();
+ }
+ }
+
+ @Test
+ public void testUnsupportedNamespaceSnapshotsDoNotTouchOtherGroups() {
+ MetaDataSubscriber metadata = mock(MetaDataSubscriber.class);
+ AuthDataSubscriber auth = mock(AuthDataSubscriber.class);
+ ProxySelectorDataSubscriber proxy =
mock(ProxySelectorDataSubscriber.class);
+ DiscoveryUpstreamDataSubscriber discovery =
mock(DiscoveryUpstreamDataSubscriber.class);
+ AiProxyApiKeyDataSubscriber apiKey =
mock(AiProxyApiKeyDataSubscriber.class);
+ WebsocketDataHandler handler = new
WebsocketDataHandler(pluginDataSubscriber, List.of(metadata), List.of(auth),
+ List.of(proxy), List.of(discovery), List.of(apiKey));
+ for (ConfigGroupEnum group : List.of(ConfigGroupEnum.META_DATA,
ConfigGroupEnum.APP_AUTH,
+ ConfigGroupEnum.PROXY_SELECTOR,
ConfigGroupEnum.DISCOVER_UPSTREAM, ConfigGroupEnum.AI_PROXY_API_KEY)) {
+
org.junit.jupiter.api.Assertions.assertThrows(IllegalArgumentException.class,
+ () -> handler.snapshot(group, "[]", "namespace-a",
"namespace-a"));
+ }
+ Mockito.verifyNoInteractions(metadata, auth, proxy, discovery, apiKey);
+ }
+
+ @Test
+ public void testEmptyNamespaceSnapshotPreservesOtherNamespaces() {
+ BaseDataCache cache = BaseDataCache.getInstance();
+ PluginDataSubscriber subscriber = new
CommonPluginDataSubscriber(Collections.emptyList(),
+ new ShenyuConfig.SelectorMatchCache(), new
ShenyuConfig.RuleMatchCache());
+ WebsocketDataHandler handler = new WebsocketDataHandler(subscriber,
Collections.emptyList(), Collections.emptyList(),
+ Collections.emptyList(), Collections.emptyList(),
Collections.emptyList());
+ PluginData other =
PluginData.builder().name("other-plugin").namespaceId("namespace-b").build();
+ SelectorData selector =
SelectorData.builder().id("other-selector").pluginName("other-plugin")
+ .namespaceId("namespace-b").sort(1).build();
+ RuleData rule =
RuleData.builder().id("other-rule").selectorId("other-selector")
+ .namespaceId("namespace-b").sort(1).build();
+ cache.cachePluginData(other);
+ cache.cacheSelectData(selector);
+ cache.cacheRuleData(rule);
+ try {
+ for (ConfigGroupEnum group : List.of(ConfigGroupEnum.PLUGIN,
ConfigGroupEnum.SELECTOR, ConfigGroupEnum.RULE)) {
+ handler.snapshot(group, "[]", "namespace-a", "namespace-a");
+ }
+ org.junit.jupiter.api.Assertions.assertEquals(other,
cache.obtainPluginData("other-plugin"));
+ org.junit.jupiter.api.Assertions.assertEquals(List.of(selector),
cache.obtainSelectorData("other-plugin"));
+ org.junit.jupiter.api.Assertions.assertEquals(List.of(rule),
cache.obtainRuleData("other-selector"));
+
org.junit.jupiter.api.Assertions.assertThrows(IllegalArgumentException.class,
+ () -> handler.snapshot(ConfigGroupEnum.PLUGIN, getJson(),
"namespace-b", "namespace-b"));
+ org.junit.jupiter.api.Assertions.assertEquals(other,
cache.obtainPluginData("other-plugin"));
+ } finally {
+ cache.cleanPluginData();
+ cache.cleanSelectorData();
+ cache.cleanRuleData();
+ }
+ }
+
+ @Test
+ public void
testNonEmptySnapshotRemovesMissingRulesAndPreservesOtherNamespace() {
+ BaseDataCache cache = BaseDataCache.getInstance();
+ PluginDataSubscriber subscriber = new
CommonPluginDataSubscriber(Collections.emptyList(),
+ new ShenyuConfig.SelectorMatchCache(), new
ShenyuConfig.RuleMatchCache());
+ WebsocketDataHandler handler = new WebsocketDataHandler(subscriber,
Collections.emptyList(), Collections.emptyList(),
+ Collections.emptyList(), Collections.emptyList(),
Collections.emptyList());
+ RuleData stale =
RuleData.builder().id("stale-rule").selectorId("selector-a")
+ .namespaceId("namespace-a").sort(1).build();
+ RuleData replacement =
RuleData.builder().id("replacement-rule").selectorId("selector-a")
+ .namespaceId("namespace-a").sort(2).build();
+ RuleData other =
RuleData.builder().id("other-rule").selectorId("selector-b")
+ .namespaceId("namespace-b").sort(1).build();
+ cache.cacheRuleData(stale);
+ cache.cacheRuleData(other);
+ try {
+ handler.snapshot(ConfigGroupEnum.RULE,
GsonUtils.getInstance().toJson(List.of(replacement)),
+ "namespace-a", "namespace-a");
+
org.junit.jupiter.api.Assertions.assertEquals(List.of(replacement),
cache.obtainRuleData("selector-a"));
+ org.junit.jupiter.api.Assertions.assertEquals(List.of(other),
cache.obtainRuleData("selector-b"));
+ } finally {
+ cache.cleanRuleData();
+ }
+ }
+
private String getJson() {
PluginData pluginData = new PluginData();
pluginData.setId("1397952341475799040");
pluginData.setName("plugin_test");
pluginData.setConfig("config_test");
+ pluginData.setNamespaceId("namespace-a");
pluginData.setEnabled(true);
pluginData.setRole("1");
LinkedList<PluginData> list = new LinkedList<>();