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 e27fb02f60 Fixes #7312: Forward data-change events from non-master
admin nodes to the master (#7342)
e27fb02f60 is described below
commit e27fb02f60560a9587331f68e4f04a9931baaf48
Author: BobSong <[email protected]>
AuthorDate: Thu Oct 1 12:00:27 2026 +0800
Fixes #7312: Forward data-change events from non-master admin nodes to the
master (#7342)
---
.../shenyu/admin/config/ClusterConfiguration.java | 19 +-
.../admin/config/properties/ClusterProperties.java | 21 ++
.../ClusterDataChangedEventController.java | 152 ++++++++++++++
.../listener/ClusterDataChangedEventForwarder.java | 132 +++++++++++++
.../admin/listener/DataChangedEventDispatcher.java | 68 +++++--
.../model/dto/ClusterDataChangedEventPayload.java | 114 +++++++++++
.../admin/shiro/bean/ClusterEventAuthFilter.java | 57 ++++++
.../admin/shiro/config/ShiroConfiguration.java | 10 +-
shenyu-admin/src/main/resources/application.yml | 3 +
.../ClusterDataChangedEventControllerTest.java | 150 ++++++++++++++
.../ClusterDataChangedEventForwarderTest.java | 220 +++++++++++++++++++++
.../listener/DataChangedEventDispatcherTest.java | 101 +++++++++-
.../shiro/bean/ClusterEventAuthFilterTest.java | 63 ++++++
.../admin/shiro/config/ShiroConfigurationTest.java | 32 ++-
14 files changed, 1125 insertions(+), 17 deletions(-)
diff --git
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/config/ClusterConfiguration.java
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/config/ClusterConfiguration.java
index 561b46cfd6..26846bf6f9 100644
---
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/config/ClusterConfiguration.java
+++
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/config/ClusterConfiguration.java
@@ -19,6 +19,7 @@ package org.apache.shenyu.admin.config;
import org.apache.shenyu.admin.config.properties.ClusterProperties;
import org.apache.shenyu.admin.config.properties.ClusterZookeeperProperties;
+import org.apache.shenyu.admin.listener.ClusterDataChangedEventForwarder;
import org.apache.shenyu.admin.mode.ShenyuRunningModeService;
import org.apache.shenyu.admin.mode.cluster.filter.ClusterForwardFilter;
import org.apache.shenyu.admin.mode.cluster.service.ClusterSelectMasterService;
@@ -81,5 +82,21 @@ public class ClusterConfiguration {
factory.setReadTimeout(clusterProperties.getReadTimeout());
return new ClusterForwardFilter(new RestTemplate(factory));
}
-
+
+ /**
+ * Shenyu cluster data changed event forwarder.
+ *
+ * @param clusterProperties cluster properties
+ * @param clusterSelectMasterService the cluster select master service
+ * @return the Shenyu cluster data changed event forwarder
+ */
+ @Bean
+ public ClusterDataChangedEventForwarder
clusterDataChangedEventForwarder(final ClusterProperties clusterProperties,
+
final ClusterSelectMasterService clusterSelectMasterService) {
+ SimpleClientHttpRequestFactory factory = new
SimpleClientHttpRequestFactory();
+ factory.setConnectTimeout(clusterProperties.getConnectionTimeout());
+ factory.setReadTimeout(clusterProperties.getReadTimeout());
+ return new ClusterDataChangedEventForwarder(new RestTemplate(factory),
clusterProperties, clusterSelectMasterService);
+ }
+
}
diff --git
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/config/properties/ClusterProperties.java
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/config/properties/ClusterProperties.java
index 6c5da4a408..7615b9c6fd 100644
---
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/config/properties/ClusterProperties.java
+++
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/config/properties/ClusterProperties.java
@@ -73,6 +73,27 @@ public class ClusterProperties {
*/
private Long lockTtl = 30L;
+ /**
+ * Dedicated node credential; configure the same high-entropy value on
every node.
+ */
+ private String eventSecret;
+
+ /**
+ * Get the event forwarding credential.
+ * @return node credential
+ */
+ public String getEventSecret() {
+ return eventSecret;
+ }
+
+ /**
+ * Set the event forwarding credential.
+ * @param eventSecret node credential
+ */
+ public void setEventSecret(final String eventSecret) {
+ this.eventSecret = eventSecret;
+ }
+
/**
* Gets the value of enabled.
*
diff --git
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/controller/ClusterDataChangedEventController.java
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/controller/ClusterDataChangedEventController.java
new file mode 100644
index 0000000000..6ef047f2f7
--- /dev/null
+++
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/controller/ClusterDataChangedEventController.java
@@ -0,0 +1,152 @@
+/*
+ * 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.controller;
+
+import org.apache.shenyu.admin.listener.DataChangedEvent;
+import org.apache.shenyu.admin.mode.cluster.service.ClusterSelectMasterService;
+import org.apache.shenyu.admin.model.dto.ClusterDataChangedEventPayload;
+import org.apache.shenyu.admin.model.result.ShenyuAdminResult;
+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.ConfigGroupEnum;
+import org.apache.shenyu.common.enums.DataEventTypeEnum;
+import org.apache.shenyu.common.utils.GsonUtils;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.context.ApplicationEventPublisher;
+import org.springframework.http.HttpStatus;
+import org.springframework.http.ResponseEntity;
+import org.springframework.lang.Nullable;
+import org.springframework.web.bind.annotation.PostMapping;
+import org.springframework.web.bind.annotation.RequestBody;
+import org.springframework.web.bind.annotation.RequestMapping;
+import org.springframework.web.bind.annotation.RestController;
+
+import java.util.List;
+import java.util.Objects;
+
+/**
+ * Receives data change events handed over by non-master admin nodes and
+ * re-publishes them locally, so configuration committed on any node reaches
+ * the listeners (websocket push, long polling cache, registry writers) that
+ * run on the current master.
+ */
+@RestController
+@RequestMapping("/cluster")
+@ConditionalOnProperty(value = "shenyu.cluster.enabled", havingValue = "true")
+public class ClusterDataChangedEventController {
+
+ private final ApplicationEventPublisher eventPublisher;
+
+ private final ClusterSelectMasterService clusterSelectMasterService;
+
+ /**
+ * Instantiates a new cluster data changed event controller.
+ *
+ * @param eventPublisher the application event publisher
+ * @param clusterSelectMasterService the cluster select master service
+ */
+ @Autowired
+ public ClusterDataChangedEventController(final ApplicationEventPublisher
eventPublisher,
+ @Nullable final
ClusterSelectMasterService clusterSelectMasterService) {
+ this.eventPublisher = eventPublisher;
+ this.clusterSelectMasterService = clusterSelectMasterService;
+ }
+
+ /**
+ * Accept a data change event forwarded by a non-master admin node and
+ * re-publish it locally. Only the current master accepts events; other
+ * nodes reject the request so the event can be retried by the sender or
+ * picked up by the new master.
+ *
+ * @param payload the forwarded event payload
+ * @return the acceptance result
+ */
+ @PostMapping("/data-change-event")
+ public ResponseEntity<ShenyuAdminResult> receive(@RequestBody final
ClusterDataChangedEventPayload payload) {
+ if (Objects.isNull(clusterSelectMasterService) ||
!clusterSelectMasterService.isMaster()) {
+ return ResponseEntity.status(HttpStatus.CONFLICT)
+ .body(ShenyuAdminResult.error("this node is not the
cluster master, data change event not accepted"));
+ }
+ if (Objects.isNull(payload) || Objects.isNull(payload.getGroupKey())
+ || Objects.isNull(payload.getEventType()) ||
Objects.isNull(payload.getSource())) {
+ return
ResponseEntity.badRequest().body(ShenyuAdminResult.error("event fields must not
be null"));
+ }
+ final ConfigGroupEnum groupKey;
+ final DataEventTypeEnum eventType;
+ try {
+ groupKey = ConfigGroupEnum.valueOf(payload.getGroupKey());
+ eventType = DataEventTypeEnum.valueOf(payload.getEventType());
+ } catch (RuntimeException ex) {
+ return ResponseEntity.status(HttpStatus.BAD_REQUEST)
+ .body(ShenyuAdminResult.error("unknown data change event
group or type: "
+ + payload.getGroupKey() + "/" +
payload.getEventType()));
+ }
+ final List<?> source;
+ try {
+ source = deserializeSource(groupKey, payload.getSource());
+ if (Objects.isNull(source) ||
source.stream().anyMatch(Objects::isNull)) {
+ return
ResponseEntity.badRequest().body(ShenyuAdminResult.error("source must be an
array of non-null records"));
+ }
+ } catch (RuntimeException ex) {
+ return ResponseEntity.status(HttpStatus.BAD_REQUEST)
+ .body(ShenyuAdminResult.error("unknown data change event
group: " + groupKey.name()));
+ }
+ eventPublisher.publishEvent(new DataChangedEvent(groupKey, eventType,
source));
+ return ResponseEntity.ok(ShenyuAdminResult.success());
+ }
+
+ private List<?> deserializeSource(final ConfigGroupEnum groupKey, final
String sourceJson) {
+ final Class<?> targetClass;
+ switch (groupKey) {
+ case APP_AUTH:
+ targetClass = AppAuthData.class;
+ break;
+ case PLUGIN:
+ targetClass = PluginData.class;
+ break;
+ case RULE:
+ targetClass = RuleData.class;
+ break;
+ case SELECTOR:
+ targetClass = SelectorData.class;
+ break;
+ case META_DATA:
+ targetClass = MetaData.class;
+ break;
+ case PROXY_SELECTOR:
+ targetClass = ProxySelectorData.class;
+ break;
+ case AI_PROXY_API_KEY:
+ targetClass = ProxyApiKeyData.class;
+ break;
+ case DISCOVER_UPSTREAM:
+ targetClass = DiscoverySyncData.class;
+ break;
+ default:
+ throw new IllegalArgumentException("Unexpected value: " +
groupKey);
+ }
+ return GsonUtils.getInstance().fromList(sourceJson, targetClass);
+ }
+}
diff --git
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/ClusterDataChangedEventForwarder.java
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/ClusterDataChangedEventForwarder.java
new file mode 100644
index 0000000000..03f1cbcbb7
--- /dev/null
+++
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/ClusterDataChangedEventForwarder.java
@@ -0,0 +1,132 @@
+/*
+ * 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;
+
+import org.apache.commons.lang3.StringUtils;
+import org.apache.shenyu.admin.config.properties.ClusterProperties;
+import org.apache.shenyu.admin.mode.cluster.service.ClusterSelectMasterService;
+import org.apache.shenyu.admin.model.dto.ClusterDataChangedEventPayload;
+import org.apache.shenyu.admin.model.dto.ClusterMasterDTO;
+import org.apache.shenyu.common.utils.GsonUtils;
+import org.apache.shenyu.admin.shiro.bean.ClusterEventAuthFilter;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.http.HttpEntity;
+import org.springframework.http.HttpHeaders;
+import org.springframework.http.ResponseEntity;
+import org.springframework.web.client.RestTemplate;
+
+import java.util.Objects;
+
+/**
+ * Forwards locally committed {@link DataChangedEvent}s to the current master
admin node.
+ *
+ * <p>{@code DataChangedEvent} is a local Spring application event. In cluster
mode a
+ * configuration write accepted by a non-master node would otherwise never
reach the
+ * listeners running on the master node (websocket push, long polling cache,
registry
+ * writers), so the change is handed over to the master over HTTP and
re-published there.
+ *
+ * <p>The forward authenticates with a dedicated node credential, independent
of request
+ * context. Configure HTTPS to protect the credential in transit.
+ *
+ * <p>Forwarding runs synchronously on the publishing thread, so an
unreachable master adds up
+ * to the configured connect/read timeout to the calling write API; delivery
is best-effort
+ * and failures are logged with the master identity and outcome.
+ */
+public class ClusterDataChangedEventForwarder {
+
+ private static final Logger LOG =
LoggerFactory.getLogger(ClusterDataChangedEventForwarder.class);
+
+ private static final String FORWARD_PATH = "/cluster/data-change-event";
+
+ private final RestTemplate restTemplate;
+
+ private final ClusterProperties clusterProperties;
+
+ private final ClusterSelectMasterService clusterSelectMasterService;
+
+ /**
+ * Instantiates a new cluster data changed event forwarder.
+ *
+ * @param restTemplate the rest template used to reach the
master node
+ * @param clusterProperties the cluster properties
+ * @param clusterSelectMasterService the cluster select master service
+ */
+ public ClusterDataChangedEventForwarder(final RestTemplate restTemplate,
+ final ClusterProperties
clusterProperties,
+ final ClusterSelectMasterService
clusterSelectMasterService) {
+ this.restTemplate = restTemplate;
+ this.clusterProperties = clusterProperties;
+ this.clusterSelectMasterService = clusterSelectMasterService;
+ }
+
+ /**
+ * Forward a committed data change event to the current master node,
authenticated with
+ * the dedicated cluster credential.
+ *
+ * @param event the locally committed data change event
+ * @return true if the master accepted the event
+ */
+ public boolean forward(final DataChangedEvent event) {
+ final ClusterMasterDTO master = clusterSelectMasterService.getMaster();
+ if (Objects.isNull(master) ||
StringUtils.isBlank(master.getMasterHost()) ||
StringUtils.isBlank(master.getMasterPort())) {
+ LOG.warn("no master available, cannot forward DataChangedEvent,
group={}, type={}, size={}",
+ event.getGroupKey(), event.getEventType(),
sourceSize(event));
+ return false;
+ }
+ final String accessToken = clusterProperties.getEventSecret();
+ if (StringUtils.isBlank(accessToken)) {
+ LOG.warn("Cluster event credential is not configured; refusing to
forward configuration");
+ return false;
+ }
+ final String url = buildMasterUrl(master);
+ final ClusterDataChangedEventPayload payload = new
ClusterDataChangedEventPayload(
+ event.getGroupKey().name(), event.getEventType().name(),
+ GsonUtils.getInstance().toJson(event.getSource()));
+ final HttpHeaders headers = new HttpHeaders();
+ headers.set(ClusterEventAuthFilter.HEADER, accessToken);
+ try {
+ final ResponseEntity<String> response =
+ restTemplate.postForEntity(url, new HttpEntity<>(payload,
headers), String.class);
+ final boolean accepted =
response.getStatusCode().is2xxSuccessful();
+ LOG.info("forwarded DataChangedEvent to master {}:{}, group={},
type={}, size={}, outcome={}",
+ master.getMasterHost(), master.getMasterPort(),
+ event.getGroupKey(), event.getEventType(),
sourceSize(event),
+ accepted ? "delivered" : "rejected");
+ return accepted;
+ } catch (final RuntimeException ex) {
+ LOG.warn("failed to forward DataChangedEvent to master {}:{},
group={}, type={}, size={}, outcome=failed",
+ master.getMasterHost(), master.getMasterPort(),
+ event.getGroupKey(), event.getEventType(),
sourceSize(event), ex);
+ return false;
+ }
+ }
+
+ private int sourceSize(final DataChangedEvent event) {
+ return Objects.isNull(event.getSource()) ? 0 :
event.getSource().size();
+ }
+
+ private String buildMasterUrl(final ClusterMasterDTO master) {
+ String contextPath =
StringUtils.defaultString(master.getContextPath());
+ if (StringUtils.isNotEmpty(contextPath) &&
!contextPath.startsWith("/")) {
+ contextPath = "/" + contextPath;
+ }
+ return clusterProperties.getSchema() + "://" + master.getMasterHost()
+ ":" + master.getMasterPort()
+ + contextPath + FORWARD_PATH;
+ }
+}
diff --git
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/DataChangedEventDispatcher.java
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/DataChangedEventDispatcher.java
index 0962677b25..6d634f9e88 100644
---
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/DataChangedEventDispatcher.java
+++
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/DataChangedEventDispatcher.java
@@ -17,6 +17,8 @@
package org.apache.shenyu.admin.listener;
+import org.springframework.transaction.support.TransactionSynchronization;
+import
org.springframework.transaction.support.TransactionSynchronizationManager;
import org.apache.shenyu.admin.config.properties.ClusterProperties;
import org.apache.shenyu.admin.mode.cluster.service.ClusterSelectMasterService;
import org.apache.shenyu.admin.service.manager.LoadServiceDocEntry;
@@ -57,28 +59,45 @@ public class DataChangedEventDispatcher implements
ApplicationListener<DataChang
private final ApplicationContext applicationContext;
private List<DataChangedListener> listeners;
-
+
+ private ClusterDataChangedEventForwarder eventForwarder;
+
@Resource
private ClusterProperties clusterProperties;
-
+
@Resource
@Nullable
private ClusterSelectMasterService shenyuClusterSelectMasterService;
-
+
public DataChangedEventDispatcher(final ApplicationContext
applicationContext) {
this.applicationContext = applicationContext;
}
-
+
@Override
- @SuppressWarnings("unchecked")
public void onApplicationEvent(@NotNull final DataChangedEvent event) {
+ if (clusterProperties.isEnabled() &&
TransactionSynchronizationManager.isActualTransactionActive()
+ &&
TransactionSynchronizationManager.isSynchronizationActive()) {
+ TransactionSynchronizationManager.registerSynchronization(new
TransactionSynchronization() {
+ @Override
+ public void afterCommit() {
+ dispatch(event);
+ }
+ });
+ return;
+ }
+ dispatch(event);
+ }
+
+ @SuppressWarnings("unchecked")
+ private void dispatch(final DataChangedEvent event) {
+ final boolean master = isMasterOrStandalone();
+ if (!master) {
+ forwardEventToMaster(event);
+ }
for (DataChangedListener listener : listeners) {
- if (!(listener instanceof AbstractDataChangedListener)
- && clusterProperties.isEnabled()
- && Objects.nonNull(shenyuClusterSelectMasterService)
- && !shenyuClusterSelectMasterService.isMaster()) {
- LOG.info("received DataChangedEvent, not master, pass");
- return;
+ if (!master && !(listener instanceof AbstractDataChangedListener))
{
+ // push listeners run on the master only; the event was handed
over above
+ continue;
}
final int size = event.getSource() instanceof java.util.Collection
? ((java.util.Collection<?>) event.getSource()).size() : 1;
@@ -130,5 +149,32 @@ public class DataChangedEventDispatcher implements
ApplicationListener<DataChang
Collection<DataChangedListener> listenerBeans =
applicationContext.getBeansOfType(DataChangedListener.class)
.values();
this.listeners = Collections.unmodifiableList(new
ArrayList<>(listenerBeans));
+ this.eventForwarder =
applicationContext.getBeanProvider(ClusterDataChangedEventForwarder.class)
+ .getIfAvailable();
+ }
+
+ private boolean isMasterOrStandalone() {
+ if (!clusterProperties.isEnabled() ||
Objects.isNull(shenyuClusterSelectMasterService)) {
+ return true;
+ }
+ return shenyuClusterSelectMasterService.isMaster();
+ }
+
+ private void forwardEventToMaster(final DataChangedEvent event) {
+ if (Objects.isNull(eventForwarder)) {
+ LOG.warn("received DataChangedEvent, not master, no forwarder
available, group={}, type={},"
+ + " push listeners will be skipped on this node",
+ event.getGroupKey(), event.getEventType());
+ return;
+ }
+ final boolean forwarded = eventForwarder.forward(event);
+ if (forwarded) {
+ LOG.info("received DataChangedEvent, not master, forwarded to
master, group={}, type={}",
+ event.getGroupKey(), event.getEventType());
+ } else {
+ LOG.warn("received DataChangedEvent, not master, forward to master
failed, group={}, type={},"
+ + " push listeners will be skipped on this node",
+ event.getGroupKey(), event.getEventType());
+ }
}
}
diff --git
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/model/dto/ClusterDataChangedEventPayload.java
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/model/dto/ClusterDataChangedEventPayload.java
new file mode 100644
index 0000000000..3903ca82e8
--- /dev/null
+++
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/model/dto/ClusterDataChangedEventPayload.java
@@ -0,0 +1,114 @@
+/*
+ * 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.model.dto;
+
+import java.io.Serializable;
+
+/**
+ * Payload used to hand a committed {@link
org.apache.shenyu.admin.listener.DataChangedEvent}
+ * over from a non-master admin node to the current master node.
+ */
+public class ClusterDataChangedEventPayload implements Serializable {
+
+ private static final long serialVersionUID = -6354049813404241267L;
+
+ /**
+ * the config group key name, see {@link
org.apache.shenyu.common.enums.ConfigGroupEnum}.
+ */
+ private String groupKey;
+
+ /**
+ * the event type name, see {@link
org.apache.shenyu.common.enums.DataEventTypeEnum}.
+ */
+ private String eventType;
+
+ /**
+ * the JSON serialized source list of the event.
+ */
+ private String source;
+
+ public ClusterDataChangedEventPayload() {
+ }
+
+ /**
+ * Instantiates a new cluster data changed event payload.
+ *
+ * @param groupKey the group key name
+ * @param eventType the event type name
+ * @param source the JSON serialized source list
+ */
+ public ClusterDataChangedEventPayload(final String groupKey, final String
eventType, final String source) {
+ this.groupKey = groupKey;
+ this.eventType = eventType;
+ this.source = source;
+ }
+
+ /**
+ * Gets the group key name.
+ *
+ * @return the group key name
+ */
+ public String getGroupKey() {
+ return groupKey;
+ }
+
+ /**
+ * Sets the group key name.
+ *
+ * @param groupKey the group key name
+ */
+ public void setGroupKey(final String groupKey) {
+ this.groupKey = groupKey;
+ }
+
+ /**
+ * Gets the event type name.
+ *
+ * @return the event type name
+ */
+ public String getEventType() {
+ return eventType;
+ }
+
+ /**
+ * Sets the event type name.
+ *
+ * @param eventType the event type name
+ */
+ public void setEventType(final String eventType) {
+ this.eventType = eventType;
+ }
+
+ /**
+ * Gets the JSON serialized source list.
+ *
+ * @return the JSON serialized source list
+ */
+ public String getSource() {
+ return source;
+ }
+
+ /**
+ * Sets the JSON serialized source list.
+ *
+ * @param source the JSON serialized source list
+ */
+ public void setSource(final String source) {
+ this.source = source;
+ }
+}
diff --git
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/shiro/bean/ClusterEventAuthFilter.java
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/shiro/bean/ClusterEventAuthFilter.java
new file mode 100644
index 0000000000..17c7163927
--- /dev/null
+++
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/shiro/bean/ClusterEventAuthFilter.java
@@ -0,0 +1,57 @@
+/*
+ * 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.shiro.bean;
+
+import jakarta.servlet.ServletRequest;
+import jakarta.servlet.ServletResponse;
+import jakarta.servlet.http.HttpServletRequest;
+import jakarta.servlet.http.HttpServletResponse;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.shiro.web.filter.AccessControlFilter;
+import org.apache.shenyu.admin.config.properties.ClusterProperties;
+
+import java.nio.charset.StandardCharsets;
+import java.security.MessageDigest;
+
+/**
+ * Dedicated authentication for the internal configuration event endpoint.
+ */
+public final class ClusterEventAuthFilter extends AccessControlFilter {
+
+ public static final String HEADER = "X-Shenyu-Cluster-Event-Secret";
+
+ private final ClusterProperties properties;
+
+ public ClusterEventAuthFilter(final ClusterProperties properties) {
+ this.properties = properties;
+ }
+
+ @Override
+ protected boolean isAccessAllowed(final ServletRequest request, final
ServletResponse response, final Object mappedValue) {
+ String expected = properties.getEventSecret();
+ String supplied = ((HttpServletRequest) request).getHeader(HEADER);
+ return properties.isEnabled() && StringUtils.isNotBlank(expected) &&
StringUtils.isNotBlank(supplied)
+ &&
MessageDigest.isEqual(expected.getBytes(StandardCharsets.UTF_8),
supplied.getBytes(StandardCharsets.UTF_8));
+ }
+
+ @Override
+ protected boolean onAccessDenied(final ServletRequest request, final
ServletResponse response) {
+ ((HttpServletResponse)
response).setStatus(HttpServletResponse.SC_FORBIDDEN);
+ return false;
+ }
+}
diff --git
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/shiro/config/ShiroConfiguration.java
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/shiro/config/ShiroConfiguration.java
index 21093d72b3..59b50bf39e 100644
---
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/shiro/config/ShiroConfiguration.java
+++
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/shiro/config/ShiroConfiguration.java
@@ -17,6 +17,8 @@
package org.apache.shenyu.admin.shiro.config;
+import org.apache.shenyu.admin.config.properties.ClusterProperties;
+import org.apache.shenyu.admin.shiro.bean.ClusterEventAuthFilter;
import org.apache.shenyu.admin.config.properties.ShiroProperties;
import org.apache.shenyu.admin.shiro.bean.StatelessAuthFilter;
import org.apache.shiro.realm.AuthorizingRealm;
@@ -58,22 +60,26 @@ public class ShiroConfiguration {
*
* @param securityManager {@linkplain DefaultWebSecurityManager}
* @param shiroProperties {@linkplain ShiroProperties}
+ * @param clusterProperties cluster authentication settings
* @return {@linkplain ShiroFilterFactoryBean}
*/
@Bean
public ShiroFilterFactoryBean shiroFilterFactoryBean(
@Qualifier("shiroSecurityManager") final DefaultWebSecurityManager
securityManager,
- @Qualifier("shiroProperties") final ShiroProperties
shiroProperties) {
+ @Qualifier("shiroProperties") final ShiroProperties
shiroProperties,
+ final ClusterProperties clusterProperties) {
ShiroFilterFactoryBean factoryBean = new ShiroFilterFactoryBean();
factoryBean.setSecurityManager(securityManager);
Map<String, Filter> filterMap = new LinkedHashMap<>();
filterMap.put("statelessAuth", new StatelessAuthFilter());
+ filterMap.put("clusterEventAuth", new
ClusterEventAuthFilter(clusterProperties));
factoryBean.setFilters(filterMap);
Map<String, String> filterChainDefinitionMap = new LinkedHashMap<>();
+ filterChainDefinitionMap.put("/cluster/data-change-event",
"clusterEventAuth");
for (String s : shiroProperties.getWhiteList()) {
- filterChainDefinitionMap.put(s, "anon");
+ filterChainDefinitionMap.putIfAbsent(s, "anon");
}
filterChainDefinitionMap.put("/**", "statelessAuth");
diff --git a/shenyu-admin/src/main/resources/application.yml
b/shenyu-admin/src/main/resources/application.yml
index eb976ea4cb..d85297cd9c 100755
--- a/shenyu-admin/src/main/resources/application.yml
+++ b/shenyu-admin/src/main/resources/application.yml
@@ -131,6 +131,9 @@ shenyu:
# secret-key: "your-secret-key-here"
expired-seconds: 86400000
cluster:
+ # Same high-entropy node-only credential on every Admin. Missing
credentials fail closed.
+ # Use HTTPS (schema: https) or a trusted TLS service mesh to protect the
credential in transit.
+ event-secret: ${SHENYU_CLUSTER_EVENT_SECRET:}
enabled: false
type: jdbc
connectionTimeout: 15000
diff --git
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/controller/ClusterDataChangedEventControllerTest.java
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/controller/ClusterDataChangedEventControllerTest.java
new file mode 100644
index 0000000000..e86ea9e989
--- /dev/null
+++
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/controller/ClusterDataChangedEventControllerTest.java
@@ -0,0 +1,150 @@
+/*
+ * 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.controller;
+
+import org.apache.shenyu.admin.listener.DataChangedEvent;
+import org.apache.shenyu.admin.mode.cluster.service.ClusterSelectMasterService;
+import org.apache.shenyu.admin.model.dto.ClusterDataChangedEventPayload;
+import org.apache.shenyu.common.dto.PluginData;
+import org.apache.shenyu.common.enums.ConfigGroupEnum;
+import org.apache.shenyu.common.enums.DataEventTypeEnum;
+import org.apache.shenyu.common.utils.GsonUtils;
+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.junit.jupiter.MockitoExtension;
+import org.mockito.junit.jupiter.MockitoSettings;
+import org.mockito.quality.Strictness;
+import org.springframework.context.ApplicationEventPublisher;
+import org.springframework.http.HttpStatus;
+import org.springframework.http.ResponseEntity;
+
+import java.util.Collections;
+import java.util.List;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * Test cases for {@link ClusterDataChangedEventController}.
+ */
+@ExtendWith(MockitoExtension.class)
+@MockitoSettings(strictness = Strictness.LENIENT)
+public final class ClusterDataChangedEventControllerTest {
+
+ @Mock
+ private ApplicationEventPublisher eventPublisher;
+
+ @Mock
+ private ClusterSelectMasterService clusterSelectMasterService;
+
+ private ClusterDataChangedEventController controller;
+
+ @BeforeEach
+ public void setUp() {
+ controller = new ClusterDataChangedEventController(eventPublisher,
clusterSelectMasterService);
+ }
+
+ /**
+ * The master accepts a forwarded event and re-publishes it locally with
typed source data.
+ */
+ @Test
+ public void receiveOnMasterPublishesEventTest() {
+ when(clusterSelectMasterService.isMaster()).thenReturn(true);
+ PluginData pluginData = new PluginData();
+ pluginData.setName("mockPlugin");
+ String sourceJson =
GsonUtils.getInstance().toJson(Collections.singletonList(pluginData));
+ ClusterDataChangedEventPayload payload =
+ new
ClusterDataChangedEventPayload(ConfigGroupEnum.PLUGIN.name(),
DataEventTypeEnum.UPDATE.name(), sourceJson);
+
+ ResponseEntity<?> response = controller.receive(payload);
+
+ assertEquals(HttpStatus.OK, response.getStatusCode());
+ ArgumentCaptor<DataChangedEvent> captor =
ArgumentCaptor.forClass(DataChangedEvent.class);
+ verify(eventPublisher, times(1)).publishEvent(captor.capture());
+ assertEquals(ConfigGroupEnum.PLUGIN, captor.getValue().getGroupKey());
+ assertEquals(1, captor.getValue().getSource().size());
+ assertInstanceOf(PluginData.class,
captor.getValue().getSource().get(0));
+ }
+
+ /**
+ * A non-master node rejects the forwarded event with 409 and does not
publish.
+ */
+ @Test
+ public void receiveOnNonMasterRejectsWithConflictTest() {
+ when(clusterSelectMasterService.isMaster()).thenReturn(false);
+ ClusterDataChangedEventPayload payload =
+ new
ClusterDataChangedEventPayload(ConfigGroupEnum.PLUGIN.name(),
DataEventTypeEnum.UPDATE.name(), "[]");
+
+ ResponseEntity<?> response = controller.receive(payload);
+
+ assertEquals(HttpStatus.CONFLICT, response.getStatusCode());
+ verify(eventPublisher,
never()).publishEvent(any(DataChangedEvent.class));
+ }
+
+ /**
+ * An unknown group key is rejected with a debuggable 400 before
publishing.
+ */
+ @Test
+ public void receiveWithUnknownGroupKeyReturnsBadRequestTest() {
+ when(clusterSelectMasterService.isMaster()).thenReturn(true);
+ ClusterDataChangedEventPayload payload =
+ new ClusterDataChangedEventPayload("UNKNOWN_GROUP",
DataEventTypeEnum.UPDATE.name(), "[]");
+
+ ResponseEntity<?> response = controller.receive(payload);
+
+ assertEquals(HttpStatus.BAD_REQUEST, response.getStatusCode());
+ verify(eventPublisher,
never()).publishEvent(any(DataChangedEvent.class));
+ }
+
+ /**
+ * All config groups map to a typed source deserialization.
+ */
+ @Test
+ public void receiveMapsEveryConfigGroupTest() {
+ when(clusterSelectMasterService.isMaster()).thenReturn(true);
+ List<ConfigGroupEnum> groups = List.of(ConfigGroupEnum.values());
+ for (ConfigGroupEnum group : groups) {
+ ClusterDataChangedEventPayload payload =
+ new ClusterDataChangedEventPayload(group.name(),
DataEventTypeEnum.UPDATE.name(), "[]");
+ ResponseEntity<?> response = controller.receive(payload);
+ assertEquals(HttpStatus.OK, response.getStatusCode());
+ }
+ verify(eventPublisher,
times(groups.size())).publishEvent(any(DataChangedEvent.class));
+ }
+
+ @Test
+ public void testNullAndMalformedFieldsReturnBadRequest() {
+ when(clusterSelectMasterService.isMaster()).thenReturn(true);
+ assertEquals(HttpStatus.BAD_REQUEST,
controller.receive(null).getStatusCode());
+ assertEquals(HttpStatus.BAD_REQUEST, controller.receive(new
ClusterDataChangedEventPayload()).getStatusCode());
+ for (String source : java.util.List.of("null", "[null]", "not-json",
"{}")) {
+ assertEquals(HttpStatus.BAD_REQUEST,
+ controller.receive(new
ClusterDataChangedEventPayload("PLUGIN", "UPDATE", source)).getStatusCode());
+ }
+ org.mockito.Mockito.verifyNoInteractions(eventPublisher);
+ }
+
+}
diff --git
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/ClusterDataChangedEventForwarderTest.java
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/ClusterDataChangedEventForwarderTest.java
new file mode 100644
index 0000000000..139ab3ca51
--- /dev/null
+++
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/ClusterDataChangedEventForwarderTest.java
@@ -0,0 +1,220 @@
+/*
+ * 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;
+
+import org.apache.shenyu.admin.config.properties.ClusterProperties;
+import org.apache.shenyu.admin.mode.cluster.service.ClusterSelectMasterService;
+import org.apache.shenyu.admin.model.dto.ClusterDataChangedEventPayload;
+import org.apache.shenyu.admin.model.dto.ClusterMasterDTO;
+import org.apache.shenyu.common.constant.Constants;
+import org.apache.shenyu.common.dto.PluginData;
+import org.apache.shenyu.common.enums.ConfigGroupEnum;
+import org.apache.shenyu.common.enums.DataEventTypeEnum;
+import org.junit.jupiter.api.AfterEach;
+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.junit.jupiter.MockitoExtension;
+import org.mockito.junit.jupiter.MockitoSettings;
+import org.mockito.quality.Strictness;
+import org.springframework.http.HttpEntity;
+import org.springframework.http.HttpStatus;
+import org.springframework.http.ResponseEntity;
+import org.springframework.mock.web.MockHttpServletRequest;
+import org.springframework.web.client.RestClientException;
+import org.springframework.web.client.RestTemplate;
+import org.springframework.web.context.request.RequestContextHolder;
+import org.springframework.web.context.request.ServletRequestAttributes;
+
+import java.util.Collections;
+import java.util.List;
+import java.util.Objects;
+
+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.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * Test cases for {@link ClusterDataChangedEventForwarder}.
+ */
+@ExtendWith(MockitoExtension.class)
+@MockitoSettings(strictness = Strictness.LENIENT)
+public final class ClusterDataChangedEventForwarderTest {
+
+ private static final String OPERATOR_TOKEN = "operator-token";
+
+ @Mock
+ private RestTemplate restTemplate;
+
+ @Mock
+ private ClusterProperties clusterProperties;
+
+ @Mock
+ private ClusterSelectMasterService clusterSelectMasterService;
+
+ private ClusterDataChangedEventForwarder forwarder;
+
+ @BeforeEach
+ public void setUp() {
+ when(clusterProperties.getEventSecret()).thenReturn("node-secret");
+ forwarder = new ClusterDataChangedEventForwarder(restTemplate,
clusterProperties, clusterSelectMasterService);
+ when(clusterProperties.getSchema()).thenReturn("http");
+ MockHttpServletRequest request = new MockHttpServletRequest();
+ request.addHeader(Constants.X_ACCESS_TOKEN, OPERATOR_TOKEN);
+ RequestContextHolder.setRequestAttributes(new
ServletRequestAttributes(request));
+ }
+
+ @AfterEach
+ public void tearDown() {
+ RequestContextHolder.resetRequestAttributes();
+ }
+
+ /**
+ * Forward posts the serialized event, authenticated with the dedicated
node credential,
+ * to the master url built from the master dto.
+ */
+ @Test
+ public void forwardPostsPayloadToMasterUrlTest() {
+ ClusterMasterDTO master = new ClusterMasterDTO();
+ master.setMasterHost("10.0.0.2");
+ master.setMasterPort("9095");
+ master.setContextPath("/admin");
+ when(clusterSelectMasterService.getMaster()).thenReturn(master);
+ when(restTemplate.postForEntity(any(String.class), any(Object.class),
eq(String.class)))
+ .thenReturn(ResponseEntity.ok("ok"));
+
+ DataChangedEvent event = new DataChangedEvent(ConfigGroupEnum.PLUGIN,
DataEventTypeEnum.UPDATE,
+ Collections.singletonList(new PluginData()));
+ boolean forwarded = forwarder.forward(event);
+
+ assertTrue(forwarded);
+ ArgumentCaptor<HttpEntity<ClusterDataChangedEventPayload>> captor =
ArgumentCaptor.forClass(HttpEntity.class);
+
verify(restTemplate).postForEntity(eq("http://10.0.0.2:9095/admin/cluster/data-change-event"),
+ captor.capture(), eq(String.class));
+ assertEquals("node-secret", captor.getValue().getHeaders().getFirst(
+
org.apache.shenyu.admin.shiro.bean.ClusterEventAuthFilter.HEADER));
+
org.junit.jupiter.api.Assertions.assertNull(captor.getValue().getHeaders().getFirst(Constants.X_ACCESS_TOKEN));
+ assertEquals(ConfigGroupEnum.PLUGIN.name(),
captor.getValue().getBody().getGroupKey());
+ assertEquals(DataEventTypeEnum.UPDATE.name(),
captor.getValue().getBody().getEventType());
+ assertTrue(captor.getValue().getBody().getSource().startsWith("["));
+ }
+
+ /**
+ * Forward without a master available returns false without posting.
+ */
+ @Test
+ public void forwardWithoutMasterReturnsFalseTest() {
+ when(clusterSelectMasterService.getMaster()).thenReturn(null);
+
+ DataChangedEvent event = new DataChangedEvent(ConfigGroupEnum.PLUGIN,
DataEventTypeEnum.UPDATE,
+ Collections.singletonList(new PluginData()));
+ boolean forwarded = forwarder.forward(event);
+
+ assertFalse(forwarded);
+ verify(restTemplate, never()).postForEntity(any(String.class),
any(Object.class), eq(String.class));
+ }
+
+ /**
+ * An event published off a request thread uses the configured node
credential.
+ */
+ @Test
+ public void forwardWithoutRequestContextUsesNodeCredentialTest() {
+ ClusterMasterDTO master = new ClusterMasterDTO();
+ master.setMasterHost("10.0.0.2");
+ master.setMasterPort("9095");
+ when(clusterSelectMasterService.getMaster()).thenReturn(master);
+ RequestContextHolder.resetRequestAttributes();
+ when(restTemplate.postForEntity(any(String.class), any(Object.class),
eq(String.class)))
+ .thenReturn(ResponseEntity.ok("ok"));
+
+ DataChangedEvent event = new DataChangedEvent(ConfigGroupEnum.PLUGIN,
DataEventTypeEnum.UPDATE,
+ Collections.singletonList(new PluginData()));
+ boolean forwarded = forwarder.forward(event);
+
+ assertTrue(forwarded);
+ verify(restTemplate).postForEntity(any(String.class),
any(Object.class), eq(String.class));
+ }
+
+ /**
+ * Forward rejected by the master (non-2xx) returns false.
+ */
+ @Test
+ public void forwardRejectedByMasterReturnsFalseTest() {
+ ClusterMasterDTO master = new ClusterMasterDTO();
+ master.setMasterHost("10.0.0.2");
+ master.setMasterPort("9095");
+ when(clusterSelectMasterService.getMaster()).thenReturn(master);
+ when(restTemplate.postForEntity(any(String.class), any(Object.class),
eq(String.class)))
+
.thenReturn(ResponseEntity.status(HttpStatus.CONFLICT).body("not master"));
+
+ DataChangedEvent event = new DataChangedEvent(ConfigGroupEnum.RULE,
DataEventTypeEnum.UPDATE,
+ Collections.singletonList(new
org.apache.shenyu.common.dto.RuleData()));
+ boolean forwarded = forwarder.forward(event);
+
+ assertFalse(forwarded);
+ }
+
+ /**
+ * Forward failure (connection error) returns false instead of throwing.
+ */
+ @Test
+ public void forwardConnectionFailureReturnsFalseTest() {
+ ClusterMasterDTO master = new ClusterMasterDTO();
+ master.setMasterHost("10.0.0.2");
+ master.setMasterPort("9095");
+ when(clusterSelectMasterService.getMaster()).thenReturn(master);
+ when(restTemplate.postForEntity(any(String.class), any(Object.class),
eq(String.class)))
+ .thenThrow(new RestClientException("connection refused"));
+
+ DataChangedEvent event = new
DataChangedEvent(ConfigGroupEnum.META_DATA, DataEventTypeEnum.UPDATE,
+ Collections.singletonList(new
org.apache.shenyu.common.dto.MetaData()));
+ boolean forwarded = forwarder.forward(event);
+
+ assertFalse(forwarded);
+ }
+
+ /**
+ * The source list is serialized into the payload body as JSON.
+ */
+ @Test
+ public void forwardSerializesSourceAsJsonTest() {
+ ClusterMasterDTO master = new ClusterMasterDTO();
+ master.setMasterHost("10.0.0.2");
+ master.setMasterPort("9095");
+ when(clusterSelectMasterService.getMaster()).thenReturn(master);
+ when(restTemplate.postForEntity(any(String.class), any(Object.class),
eq(String.class)))
+ .thenReturn(ResponseEntity.ok("ok"));
+
+ PluginData pluginData = new PluginData();
+ pluginData.setName("mockPlugin");
+ List<PluginData> source = Collections.singletonList(pluginData);
+ forwarder.forward(new DataChangedEvent(ConfigGroupEnum.PLUGIN,
DataEventTypeEnum.UPDATE, source));
+
+ ArgumentCaptor<HttpEntity<ClusterDataChangedEventPayload>> captor =
ArgumentCaptor.forClass(HttpEntity.class);
+ verify(restTemplate).postForEntity(any(String.class),
captor.capture(), eq(String.class));
+ assertTrue(Objects.nonNull(captor.getValue().getBody()));
+
assertTrue(captor.getValue().getBody().getSource().contains("mockPlugin"));
+ }
+}
diff --git
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/DataChangedEventDispatcherTest.java
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/DataChangedEventDispatcherTest.java
index 52236191bd..70a5e9cd32 100644
---
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/DataChangedEventDispatcherTest.java
+++
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/DataChangedEventDispatcherTest.java
@@ -33,6 +33,7 @@ import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.mockito.junit.jupiter.MockitoSettings;
import org.mockito.quality.Strictness;
+import org.springframework.beans.factory.ObjectProvider;
import org.springframework.context.ApplicationContext;
import org.springframework.test.util.ReflectionTestUtils;
@@ -84,10 +85,16 @@ public final class DataChangedEventDispatcherTest {
@Mock
private ClusterProperties clusterProperties;
-
+
@Mock
private ClusterSelectMasterService shenyuClusterSelectMasterService;
+ @Mock
+ private ClusterDataChangedEventForwarder clusterDataChangedEventForwarder;
+
+ @Mock
+ private ObjectProvider<ClusterDataChangedEventForwarder> forwarderProvider;
+
@BeforeEach
public void setUp() throws NoSuchFieldException, IllegalAccessException {
Map<String, DataChangedListener> listenerMap = new HashMap<>();
@@ -99,7 +106,10 @@ public final class DataChangedEventDispatcherTest {
when(applicationContext.getBean(LoadServiceDocEntry.class)).thenReturn(loadServiceDocEntry);
applicationContext.getBean(LoadServiceDocEntry.class);
-
+
+
when(applicationContext.getBeanProvider(ClusterDataChangedEventForwarder.class)).thenReturn(forwarderProvider);
+
when(forwarderProvider.getIfAvailable()).thenReturn(clusterDataChangedEventForwarder);
+
Field shenyuClusterSelectMasterServiceField =
DataChangedEventDispatcher.class.getDeclaredField("shenyuClusterSelectMasterService");
shenyuClusterSelectMasterServiceField.setAccessible(true);
shenyuClusterSelectMasterServiceField.set(dataChangedEventDispatcher,
shenyuClusterSelectMasterService);
@@ -311,4 +321,91 @@ public final class DataChangedEventDispatcherTest {
verify(websocketDataChangedListener,
times(1)).onPluginChanged(anyList(), any(), any());
verify(zookeeperDataChangedListener,
times(1)).onPluginChanged(anyList(), any(), any());
}
+
+ /**
+ * When not master, the event is forwarded once and push listeners are
skipped,
+ * while AbstractDataChangedListener listeners run regardless of their
position
+ * in the listener iteration order.
+ */
+ @Test
+ void
onApplicationEventNotMasterForwardsEventAndKeepsAbstractListenersTest() {
+ when(clusterProperties.isEnabled()).thenReturn(true);
+ when(shenyuClusterSelectMasterService.isMaster()).thenReturn(false);
+
when(clusterDataChangedEventForwarder.forward(any(DataChangedEvent.class))).thenReturn(true);
+ List<DataChangedListener> orderedListeners = new ArrayList<>();
+ orderedListeners.add(nacosDataChangedListener);
+ orderedListeners.add(httpLongPollingDataChangedListener);
+ orderedListeners.add(websocketDataChangedListener);
+ ReflectionTestUtils.setField(dataChangedEventDispatcher, "listeners",
Collections.unmodifiableList(orderedListeners));
+ DataChangedEvent dataChangedEvent = new
DataChangedEvent(ConfigGroupEnum.PLUGIN, null, new ArrayList<>());
+ dataChangedEventDispatcher.onApplicationEvent(dataChangedEvent);
+ verify(clusterDataChangedEventForwarder,
times(1)).forward(dataChangedEvent);
+ verify(httpLongPollingDataChangedListener,
times(1)).onPluginChanged(anyList(), any());
+ verify(nacosDataChangedListener, never()).onPluginChanged(anyList(),
any());
+ verify(websocketDataChangedListener,
never()).onPluginChanged(anyList(), any());
+ }
+
+ /**
+ * When not master and forwarding to the master fails, local cache
listeners still run.
+ */
+ @Test
+ void onApplicationEventNotMasterForwardFailedStillUpdatesLocalCachesTest()
{
+ when(clusterProperties.isEnabled()).thenReturn(true);
+ when(shenyuClusterSelectMasterService.isMaster()).thenReturn(false);
+
when(clusterDataChangedEventForwarder.forward(any(DataChangedEvent.class))).thenReturn(false);
+ DataChangedEvent dataChangedEvent = new
DataChangedEvent(ConfigGroupEnum.PLUGIN, null, new ArrayList<>());
+ dataChangedEventDispatcher.onApplicationEvent(dataChangedEvent);
+ verify(clusterDataChangedEventForwarder,
times(1)).forward(dataChangedEvent);
+ verify(httpLongPollingDataChangedListener,
times(1)).onPluginChanged(anyList(), any());
+ verify(websocketDataChangedListener,
never()).onPluginChanged(anyList(), any());
+ }
+
+ /**
+ * When cluster is disabled, events are never forwarded even if a
forwarder is present.
+ */
+ @Test
+ void onApplicationEventStandaloneNeverForwardsTest() {
+ when(clusterProperties.isEnabled()).thenReturn(false);
+ DataChangedEvent dataChangedEvent = new
DataChangedEvent(ConfigGroupEnum.PLUGIN, null, new ArrayList<>());
+ dataChangedEventDispatcher.onApplicationEvent(dataChangedEvent);
+ verify(clusterDataChangedEventForwarder,
never()).forward(any(DataChangedEvent.class));
+ verify(websocketDataChangedListener,
times(1)).onPluginChanged(anyList(), any());
+ verify(httpLongPollingDataChangedListener,
times(1)).onPluginChanged(anyList(), any());
+ }
+
+ @Test
+ void testTransactionDefersForwardingUntilCommit() {
+ when(clusterProperties.isEnabled()).thenReturn(true);
+ when(shenyuClusterSelectMasterService.isMaster()).thenReturn(false);
+
org.springframework.transaction.support.TransactionSynchronizationManager.initSynchronization();
+
org.springframework.transaction.support.TransactionSynchronizationManager.setActualTransactionActive(true);
+ try {
+ DataChangedEvent event = new
DataChangedEvent(ConfigGroupEnum.PLUGIN, null, new ArrayList<>());
+ dataChangedEventDispatcher.onApplicationEvent(event);
+ verify(clusterDataChangedEventForwarder, never()).forward(any());
+ verify(httpLongPollingDataChangedListener,
never()).onPluginChanged(anyList(), any());
+
org.springframework.transaction.support.TransactionSynchronizationManager.getSynchronizations()
+
.forEach(org.springframework.transaction.support.TransactionSynchronization::afterCommit);
+ verify(clusterDataChangedEventForwarder).forward(event);
+ } finally {
+
org.springframework.transaction.support.TransactionSynchronizationManager.clear();
+ }
+ }
+
+ @Test
+ void testRollbackDoesNotForwardOrUpdateCache() {
+ when(clusterProperties.isEnabled()).thenReturn(true);
+
org.springframework.transaction.support.TransactionSynchronizationManager.initSynchronization();
+
org.springframework.transaction.support.TransactionSynchronizationManager.setActualTransactionActive(true);
+ try {
+ dataChangedEventDispatcher.onApplicationEvent(new
DataChangedEvent(ConfigGroupEnum.PLUGIN, null, new ArrayList<>()));
+
org.springframework.transaction.support.TransactionSynchronizationManager.getSynchronizations()
+ .forEach(sync ->
sync.afterCompletion(org.springframework.transaction.support.TransactionSynchronization.STATUS_ROLLED_BACK));
+ verify(clusterDataChangedEventForwarder, never()).forward(any());
+ verify(httpLongPollingDataChangedListener,
never()).onPluginChanged(anyList(), any());
+ } finally {
+
org.springframework.transaction.support.TransactionSynchronizationManager.clear();
+ }
+ }
+
}
diff --git
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/shiro/bean/ClusterEventAuthFilterTest.java
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/shiro/bean/ClusterEventAuthFilterTest.java
new file mode 100644
index 0000000000..0c960ad755
--- /dev/null
+++
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/shiro/bean/ClusterEventAuthFilterTest.java
@@ -0,0 +1,63 @@
+/*
+ * 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.shiro.bean;
+
+import org.apache.shenyu.admin.config.properties.ClusterProperties;
+import org.junit.jupiter.api.Test;
+import org.springframework.mock.web.MockHttpServletRequest;
+import org.springframework.mock.web.MockHttpServletResponse;
+
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Exercise the node authentication filter, not a direct controller call.
+ */
+public final class ClusterEventAuthFilterTest {
+
+ @Test
+ public void testOrdinaryUserCannotPublishButNodeCan() throws Exception {
+ ClusterProperties properties = new ClusterProperties();
+ properties.setEnabled(true);
+ properties.setEventSecret("dedicated-node-secret");
+ ClusterEventAuthFilter filter = new ClusterEventAuthFilter(properties);
+ filter.processPathConfig("/cluster/data-change-event", null);
+ MockHttpServletRequest request = new MockHttpServletRequest("POST",
"/cluster/data-change-event");
+ request.setServletPath("/cluster/data-change-event");
+ request.addHeader("X-Access-Token", "ordinary-user-token");
+ MockHttpServletResponse response = new MockHttpServletResponse();
+ AtomicBoolean reachedController = new AtomicBoolean();
+ filter.doFilter(request, response, (req, res) ->
reachedController.set(true));
+ assertEquals(403, response.getStatus());
+ assertFalse(reachedController.get());
+ request.addHeader(ClusterEventAuthFilter.HEADER, "wrong-secret");
+ filter.doFilter(request, new MockHttpServletResponse(), (req, res) ->
reachedController.set(true));
+ assertFalse(reachedController.get());
+ request.removeHeader(ClusterEventAuthFilter.HEADER);
+ request.addHeader(ClusterEventAuthFilter.HEADER,
"dedicated-node-secret");
+ filter.doFilter(request, new MockHttpServletResponse(), (req, res) ->
reachedController.set(true));
+ assertTrue(reachedController.get());
+ reachedController.set(false);
+ properties.setEventSecret("");
+ filter.doFilter(request, new MockHttpServletResponse(), (req, res) ->
reachedController.set(true));
+ assertFalse(reachedController.get());
+ }
+}
diff --git
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/shiro/config/ShiroConfigurationTest.java
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/shiro/config/ShiroConfigurationTest.java
index 91b1df8e40..498ae7075f 100644
---
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/shiro/config/ShiroConfigurationTest.java
+++
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/shiro/config/ShiroConfigurationTest.java
@@ -70,7 +70,8 @@ public final class ShiroConfigurationTest {
ShiroProperties shiroProperties = mock(ShiroProperties.class);
List<String> whiteList = Arrays.asList("test1", "test2");
when(shiroProperties.getWhiteList()).thenReturn(whiteList);
- ShiroFilterFactoryBean shiroFilterFactoryBean =
shiroConfiguration.shiroFilterFactoryBean(securityManager, shiroProperties);
+ ShiroFilterFactoryBean shiroFilterFactoryBean =
shiroConfiguration.shiroFilterFactoryBean(securityManager, shiroProperties,
+ new
org.apache.shenyu.admin.config.properties.ClusterProperties());
assertEquals(securityManager,
shiroFilterFactoryBean.getSecurityManager());
assertNotNull(shiroFilterFactoryBean.getFilters());
assertNotNull(shiroFilterFactoryBean.getFilters().get("statelessAuth"));
@@ -98,4 +99,33 @@ public final class ShiroConfigurationTest {
LifecycleBeanPostProcessor postProcessor =
shiroConfiguration.lifecycleBeanPostProcessor();
assertNotNull(postProcessor);
}
+
+ @Test
+ public void testClusterEndpointUsesDedicatedFilterBeforeWhitelist() throws
Exception {
+ org.apache.shenyu.admin.config.properties.ClusterProperties cluster =
new org.apache.shenyu.admin.config.properties.ClusterProperties();
+ cluster.setEnabled(true);
+ cluster.setEventSecret("dedicated-node-secret");
+ ShiroProperties properties = new ShiroProperties();
+ properties.setWhiteList(java.util.List.of("/cluster/**"));
+ DefaultWebSecurityManager manager = new
DefaultWebSecurityManager(mock(AuthorizingRealm.class));
+ ShiroFilterFactoryBean factory =
shiroConfiguration.shiroFilterFactoryBean(manager, properties, cluster);
+ jakarta.servlet.Filter filter = (jakarta.servlet.Filter)
factory.getObject();
+ org.springframework.mock.web.MockHttpServletRequest request =
+ new
org.springframework.mock.web.MockHttpServletRequest("POST",
"/cluster/data-change-event");
+ request.setServletPath("/cluster/data-change-event");
+ request.addHeader("X-Access-Token", "ordinary-user-token");
+ org.springframework.mock.web.MockHttpServletResponse response = new
org.springframework.mock.web.MockHttpServletResponse();
+ java.util.concurrent.atomic.AtomicBoolean reachedController = new
java.util.concurrent.atomic.AtomicBoolean();
+ try {
+ filter.doFilter(request, response, (req, res) ->
reachedController.set(true));
+ assertEquals(403, response.getStatus());
+
org.junit.jupiter.api.Assertions.assertFalse(reachedController.get());
+
request.addHeader(org.apache.shenyu.admin.shiro.bean.ClusterEventAuthFilter.HEADER,
"dedicated-node-secret");
+ filter.doFilter(request, new
org.springframework.mock.web.MockHttpServletResponse(), (req, res) ->
reachedController.set(true));
+ assertTrue(reachedController.get());
+ } finally {
+ manager.destroy();
+ }
+ }
+
}