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 c7fb00e12d fix(discovery): remove stale upstream cache after selector
deletion (#7289)
c7fb00e12d is described below
commit c7fb00e12d945c2231668c987fe5e9dc7e73707d
Author: lymerin <[email protected]>
AuthorDate: Thu Oct 1 06:27:46 2026 +0800
fix(discovery): remove stale upstream cache after selector deletion (#7289)
---
RELEASE-NOTES.md | 4 +
.../listener/AbstractPathDataChangedListener.java | 6 ++
.../AbstractPathDataChangedListenerTest.java | 115 +++++++++++++++++++++
.../admin/service/impl/SelectorServiceImpl.java | 45 +++++++-
.../shenyu/admin/service/SelectorServiceTest.java | 75 ++++++++++++++
.../CommonDiscoveryUpstreamDataSubscriber.java | 13 ++-
.../base/handler/DiscoveryUpstreamDataHandler.java | 9 ++
.../CommonDiscoveryUpstreamDataSubscriberTest.java | 21 ++--
.../divide/handler/DivideUpstreamDataHandler.java | 9 ++
.../handler/DivideUpstreamDataHandlerTest.java | 13 +++
.../handler/GrpcDiscoveryUpstreamDataHandler.java | 9 ++
.../GrpcDiscoveryUpstreamDataHandlerTest.java | 45 ++++++++
.../plugin/tcp/handler/TcpBootstrapFactory.java | 3 +
.../plugin/tcp/handler/TcpUpstreamDataHandler.java | 36 ++++++-
.../tcp/handler/TcpBootstrapFactoryTest.java | 3 +
.../handler/TcpProxySelectorDataHandlerTest.java | 22 ++++
.../tcp/handler/TcpUpstreamDataHandlerTest.java | 80 +++++++++++++-
.../handler/WebSocketUpstreamDataHandler.java | 9 ++
.../plugin/websocket/WebSocketPluginTest.java | 14 +++
.../shenyu/protocol/tcp/UpstreamProvider.java | 56 ++++++++++
.../data/api/DiscoveryUpstreamDataSubscriber.java | 4 +-
...taSubscriber.java => DiscoveryUpstreamKey.java} | 30 +++---
.../data/core/AbstractNodeDataSyncService.java | 7 +-
.../data/core/AbstractPathDataSyncService.java | 24 +++--
.../data/core/AbstractNodeDataSyncServiceTest.java | 14 ++-
.../data/core/AbstractPathDataSyncServiceTest.java | 19 +++-
.../http/refresh/DiscoveryUpstreamDataRefresh.java | 37 ++++++-
.../refresh/DiscoveryUpstreamDataRefreshTest.java | 52 +++++++++-
.../handler/DiscoveryUpstreamDataHandler.java | 5 +-
.../handler/DiscoveryUpstreamDataHandlerTest.java | 34 +++++-
.../zookeeper/ZookeeperSyncDataServiceTest.java | 32 +++++-
31 files changed, 775 insertions(+), 70 deletions(-)
diff --git a/RELEASE-NOTES.md b/RELEASE-NOTES.md
index c573a17e0f..66027048f7 100644
--- a/RELEASE-NOTES.md
+++ b/RELEASE-NOTES.md
@@ -1,5 +1,9 @@
## Unreleased
+### API Changes
+
+- `DiscoveryUpstreamDataSubscriber#unSubscribe(DiscoverySyncData)` has been
replaced by `unSubscribe(DiscoveryUpstreamKey)`. Downstream implementations
must update their method signature; this is a source- and binary-incompatible
change. (#7289)
+
### Behavior Changes
- HTTP retry strategies budget the entire sequence separately from each
attempt:
diff --git
a/shenyu-admin-listener/shenyu-admin-listener-api/src/main/java/org/apache/shenyu/admin/listener/AbstractPathDataChangedListener.java
b/shenyu-admin-listener/shenyu-admin-listener-api/src/main/java/org/apache/shenyu/admin/listener/AbstractPathDataChangedListener.java
index 587d277fcb..949ea8edf8 100644
---
a/shenyu-admin-listener/shenyu-admin-listener-api/src/main/java/org/apache/shenyu/admin/listener/AbstractPathDataChangedListener.java
+++
b/shenyu-admin-listener/shenyu-admin-listener-api/src/main/java/org/apache/shenyu/admin/listener/AbstractPathDataChangedListener.java
@@ -18,6 +18,7 @@
package org.apache.shenyu.admin.listener;
import org.apache.commons.collections4.CollectionUtils;
+import org.apache.commons.lang3.StringUtils;
import org.apache.shenyu.common.constant.DefaultPathConstants;
import org.apache.shenyu.common.dto.AppAuthData;
import org.apache.shenyu.common.dto.PluginData;
@@ -97,6 +98,11 @@ public abstract class AbstractPathDataChangedListener
implements DataChangedList
@Override
public void onDiscoveryUpstreamChanged(final List<DiscoverySyncData>
changed, final DataEventTypeEnum eventType) {
for (DiscoverySyncData data : changed) {
+ if (StringUtils.isAnyBlank(data.getPluginName(),
data.getSelectorId())) {
+ LOG.warn("[DataChangedListener] ignore discoveryUpstream
change with empty pluginName or selectorId, namespaceId={}, pluginName={},
selectorId={}",
+ data.getNamespaceId(), data.getPluginName(),
data.getSelectorId());
+ continue;
+ }
String upstreamPath =
DefaultPathConstants.buildDiscoveryUpstreamPath(data.getNamespaceId(),
data.getPluginName(), data.getSelectorId());
// delete
if (eventType == DataEventTypeEnum.DELETE) {
diff --git
a/shenyu-admin-listener/shenyu-admin-listener-api/src/test/java/org/apache/shenyu/admin/listener/AbstractPathDataChangedListenerTest.java
b/shenyu-admin-listener/shenyu-admin-listener-api/src/test/java/org/apache/shenyu/admin/listener/AbstractPathDataChangedListenerTest.java
new file mode 100644
index 0000000000..fa92563e55
--- /dev/null
+++
b/shenyu-admin-listener/shenyu-admin-listener-api/src/test/java/org/apache/shenyu/admin/listener/AbstractPathDataChangedListenerTest.java
@@ -0,0 +1,115 @@
+/*
+ * 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.common.dto.DiscoverySyncData;
+import org.apache.shenyu.common.enums.DataEventTypeEnum;
+import org.junit.jupiter.api.Test;
+
+import java.util.Collections;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+
+public final class AbstractPathDataChangedListenerTest {
+
+ @Test
+ public void
testOnDiscoveryUpstreamChangedIgnoresDeleteWithEmptyPluginName() {
+ DiscoverySyncData discoverySyncData = new DiscoverySyncData();
+ discoverySyncData.setNamespaceId("default");
+ discoverySyncData.setSelectorId("selector-id");
+ discoverySyncData.setPluginName("");
+ TestPathDataChangedListener listener = new
TestPathDataChangedListener();
+
+
listener.onDiscoveryUpstreamChanged(Collections.singletonList(discoverySyncData),
DataEventTypeEnum.DELETE);
+
+ assertNull(listener.deletedPath);
+ }
+
+ @Test
+ public void
testOnDiscoveryUpstreamChangedIgnoresUpdateWithEmptyPluginName() {
+ DiscoverySyncData discoverySyncData = new DiscoverySyncData();
+ discoverySyncData.setNamespaceId("default");
+ discoverySyncData.setSelectorId("selector-id");
+ discoverySyncData.setPluginName("");
+ TestPathDataChangedListener listener = new
TestPathDataChangedListener();
+
+
listener.onDiscoveryUpstreamChanged(Collections.singletonList(discoverySyncData),
DataEventTypeEnum.UPDATE);
+
+ assertNull(listener.updatedPath);
+ }
+
+ @Test
+ public void
testOnDiscoveryUpstreamChangedIgnoresDeleteWithMissingSelectorId() {
+ DiscoverySyncData discoverySyncData = new DiscoverySyncData();
+ discoverySyncData.setNamespaceId("default");
+ discoverySyncData.setPluginName("divide");
+ TestPathDataChangedListener listener = new
TestPathDataChangedListener();
+
+
listener.onDiscoveryUpstreamChanged(Collections.singletonList(discoverySyncData),
DataEventTypeEnum.DELETE);
+
+ assertNull(listener.deletedPath);
+ }
+
+ @Test
+ public void
testOnDiscoveryUpstreamChangedIgnoresUpdateWithBlankSelectorId() {
+ DiscoverySyncData discoverySyncData = new DiscoverySyncData();
+ discoverySyncData.setNamespaceId("default");
+ discoverySyncData.setSelectorId(" ");
+ discoverySyncData.setPluginName("divide");
+ TestPathDataChangedListener listener = new
TestPathDataChangedListener();
+
+
listener.onDiscoveryUpstreamChanged(Collections.singletonList(discoverySyncData),
DataEventTypeEnum.UPDATE);
+
+ assertNull(listener.updatedPath);
+ }
+
+ @Test
+ public void testOnDiscoveryUpstreamChangedDeletesValidPath() {
+ DiscoverySyncData discoverySyncData = new DiscoverySyncData();
+ discoverySyncData.setNamespaceId("default");
+ discoverySyncData.setSelectorId("selector-id");
+ discoverySyncData.setPluginName("divide");
+ TestPathDataChangedListener listener = new
TestPathDataChangedListener();
+
+
listener.onDiscoveryUpstreamChanged(Collections.singletonList(discoverySyncData),
DataEventTypeEnum.DELETE);
+
+ assertEquals("/default/shenyu/discoveryUpstream/divide/selector-id",
listener.deletedPath);
+ }
+
+ private static final class TestPathDataChangedListener extends
AbstractPathDataChangedListener {
+
+ private String deletedPath;
+
+ private String updatedPath;
+
+ @Override
+ public void createOrUpdate(final String pluginPath, final Object data)
{
+ updatedPath = pluginPath;
+ }
+
+ @Override
+ public void deleteNode(final String pluginPath) {
+ deletedPath = pluginPath;
+ }
+
+ @Override
+ public void deletePathRecursive(final String selectorParentPath) {
+ }
+ }
+}
diff --git
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/SelectorServiceImpl.java
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/SelectorServiceImpl.java
index c1434b3864..0b373de389 100644
---
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/SelectorServiceImpl.java
+++
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/SelectorServiceImpl.java
@@ -29,6 +29,7 @@ import org.apache.shenyu.admin.aspect.annotation.Pageable;
import org.apache.shenyu.admin.discovery.DiscoveryLevel;
import org.apache.shenyu.admin.discovery.DiscoveryProcessor;
import org.apache.shenyu.admin.discovery.DiscoveryProcessorHolder;
+import org.apache.shenyu.admin.exception.ShenyuAdminException;
import org.apache.shenyu.admin.listener.DataChangedEvent;
import org.apache.shenyu.admin.mapper.DiscoveryHandlerMapper;
import org.apache.shenyu.admin.mapper.DiscoveryMapper;
@@ -90,6 +91,7 @@ import
org.springframework.transaction.annotation.Transactional;
import java.util.ArrayList;
import java.util.Collections;
+import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
@@ -297,7 +299,7 @@ public class SelectorServiceImpl implements SelectorService
{
@Transactional(rollbackFor = Exception.class)
public int deleteByNamespaceId(final List<String> ids, final String
namespaceId) {
final List<SelectorDO> selectors = selectorMapper.selectByIdSet(new
TreeSet<>(ids));
- List<PluginDO> pluginDOS =
pluginMapper.selectByIds(ListUtil.map(selectors, SelectorDO::getPluginId));
+ List<PluginDO> pluginDOS = new
ArrayList<>(pluginMapper.selectByIds(ListUtil.map(selectors,
SelectorDO::getPluginId)));
unbindDiscovery(selectors, pluginDOS);
return deleteSelector(selectors, pluginDOS);
}
@@ -308,22 +310,54 @@ public class SelectorServiceImpl implements
SelectorService {
* @param selectors selectors
*/
private void unbindDiscovery(final List<SelectorDO> selectors, final
List<PluginDO> pluginDOS) {
- Map<String, String> pluginMap = ListUtil.toMap(pluginDOS,
PluginDO::getId, PluginDO::getName);
+ Map<String, String> pluginMap = new
HashMap<>(ListUtil.toMap(pluginDOS, PluginDO::getId, PluginDO::getName));
+ List<ResolvedDiscovery> resolvedDiscoveries = new ArrayList<>();
+ // Validate the whole batch before deleting any discovery rows or
publishing removal events.
for (SelectorDO selector : selectors) {
DiscoveryHandlerDO discoveryHandlerDO =
discoveryHandlerMapper.selectBySelectorId(selector.getId());
if (Objects.isNull(discoveryHandlerDO)) {
continue;
}
+ DiscoveryDO discoveryDO =
discoveryMapper.selectById(discoveryHandlerDO.getDiscoveryId());
+ String pluginName = null;
+ if (Objects.nonNull(discoveryDO)) {
+ pluginName = pluginMap.get(selector.getPluginId());
+ if (StringUtils.isBlank(pluginName)) {
+ PluginDO pluginDO =
pluginMapper.selectById(selector.getPluginId());
+ pluginName = Objects.isNull(pluginDO) ? null :
pluginDO.getName();
+ }
+ if (StringUtils.isBlank(pluginName)) {
+ pluginName = discoveryDO.getPluginName();
+ }
+ if (StringUtils.isBlank(pluginName)) {
+ throw new ShenyuAdminException("Cannot delete selector
batch: plugin name for selector " + selector.getId()
+ + " with discovery upstream data cannot be
resolved. No selectors in this batch were deleted; restore the plugin name and
retry");
+ }
+ if (!Objects.equals(pluginMap.get(selector.getPluginId()),
pluginName)) {
+ // The selector deletion event also resolves plugin names
from this list.
+ pluginDOS.removeIf(plugin ->
Objects.equals(plugin.getId(), selector.getPluginId()));
+ PluginDO pluginDO = new PluginDO();
+ pluginDO.setId(selector.getPluginId());
+ pluginDO.setName(pluginName);
+ pluginDOS.add(pluginDO);
+ pluginMap.put(selector.getPluginId(), pluginName);
+ }
+ }
+ resolvedDiscoveries.add(new ResolvedDiscovery(selector,
discoveryHandlerDO, discoveryDO, pluginName));
+ }
+ for (ResolvedDiscovery resolved : resolvedDiscoveries) {
+ SelectorDO selector = resolved.selector();
+ DiscoveryHandlerDO discoveryHandlerDO = resolved.handler();
discoveryHandlerMapper.delete(discoveryHandlerDO.getId());
discoveryRelMapper.deleteByDiscoveryHandlerId(discoveryHandlerDO.getId());
discoveryUpstreamMapper.deleteByDiscoveryHandlerId(discoveryHandlerDO.getId());
- DiscoveryDO discoveryDO =
discoveryMapper.selectById(discoveryHandlerDO.getDiscoveryId());
+ DiscoveryDO discoveryDO = resolved.discovery();
if (Objects.nonNull(discoveryDO)) {
final DiscoveryProcessor discoveryProcessor =
discoveryProcessorHolder.chooseProcessor(discoveryDO.getDiscoveryType());
ProxySelectorDTO proxySelectorDTO = new ProxySelectorDTO();
proxySelectorDTO.setId(selector.getId());
proxySelectorDTO.setName(selector.getSelectorName());
-
proxySelectorDTO.setPluginName(pluginMap.getOrDefault(selector.getPluginId(),
""));
+ proxySelectorDTO.setPluginName(resolved.pluginName());
proxySelectorDTO.setNamespaceId(selector.getNamespaceId());
discoveryProcessor.removeProxySelector(DiscoveryTransfer.INSTANCE.mapToDTO(discoveryHandlerDO),
proxySelectorDTO);
if
(DiscoveryLevel.SELECTOR.getCode().equals(discoveryDO.getDiscoveryLevel())) {
@@ -720,4 +754,7 @@ public class SelectorServiceImpl implements SelectorService
{
.collect(Collectors.toList());
}
+ private record ResolvedDiscovery(SelectorDO selector, DiscoveryHandlerDO
handler, DiscoveryDO discovery, String pluginName) {
+ }
+
}
diff --git
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/service/SelectorServiceTest.java
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/service/SelectorServiceTest.java
index 5396ebb0ba..c2d5557d75 100644
---
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/service/SelectorServiceTest.java
+++
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/service/SelectorServiceTest.java
@@ -20,6 +20,7 @@ package org.apache.shenyu.admin.service;
import org.apache.commons.lang3.StringUtils;
import org.apache.shenyu.admin.discovery.DiscoveryProcessor;
import org.apache.shenyu.admin.discovery.DiscoveryProcessorHolder;
+import org.apache.shenyu.admin.exception.ShenyuAdminException;
import org.apache.shenyu.admin.mapper.DataPermissionMapper;
import org.apache.shenyu.admin.mapper.DiscoveryHandlerMapper;
import org.apache.shenyu.admin.mapper.DiscoveryMapper;
@@ -70,6 +71,7 @@ import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.Random;
+import java.util.TreeSet;
import java.util.stream.Collectors;
import java.util.stream.Stream;
@@ -79,6 +81,9 @@ import static org.hamcrest.Matchers.greaterThan;
import static org.hamcrest.Matchers.notNullValue;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.argThat;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.eq;
@@ -86,6 +91,9 @@ import static org.mockito.BDDMockito.given;
import static org.mockito.Mockito.doNothing;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
/**
@@ -199,6 +207,73 @@ public final class SelectorServiceTest {
assertEquals(selectorService.deleteByNamespaceId(ids, any()),
ids.size());
}
+ @Test
+ public void testDeleteResolvesDiscoveryPluginNameWhenBatchLookupMisses() {
+ SelectorDO selector = buildSelectorDO();
+
given(selectorMapper.selectByIdSet(Collections.singleton(selector.getId()))).willReturn(Collections.singletonList(selector));
+
given(pluginMapper.selectByIds(Collections.singletonList(selector.getPluginId()))).willReturn(Collections.emptyList());
+
given(pluginMapper.selectById(selector.getPluginId())).willReturn(buildPluginDO());
+
given(selectorMapper.deleteByIds(Collections.singletonList(selector.getId()))).willReturn(1);
+
+
selectorService.deleteByNamespaceId(Collections.singletonList(selector.getId()),
SYS_DEFAULT_NAMESPACE_ID);
+
+ verify(discoveryProcessor).removeProxySelector(any(), argThat(data ->
"test".equals(data.getPluginName())));
+ verify(selectorEventPublisher).onDeleted(any(), argThat(plugins ->
"test".equals(plugins.get(0).getName())));
+ }
+
+ @Test
+ public void testDeleteUsesDiscoveryPluginNameWhenPluginRowIsMissing() {
+ DiscoveryDO discovery = new DiscoveryDO();
+ discovery.setId("1");
+ discovery.setDiscoveryType("local");
+ discovery.setPluginName("test");
+ SelectorDO selector = buildSelectorDO();
+
given(selectorMapper.selectByIdSet(Collections.singleton(selector.getId()))).willReturn(Collections.singletonList(selector));
+
given(pluginMapper.selectByIds(Collections.singletonList(selector.getPluginId()))).willReturn(Collections.emptyList());
+ given(discoveryMapper.selectById("1")).willReturn(discovery);
+
given(selectorMapper.deleteByIds(Collections.singletonList(selector.getId()))).willReturn(1);
+
+ assertEquals(1,
selectorService.deleteByNamespaceId(Collections.singletonList(selector.getId()),
SYS_DEFAULT_NAMESPACE_ID));
+
+ verify(discoveryProcessor).removeProxySelector(any(), argThat(data ->
"test".equals(data.getPluginName())));
+ verify(selectorEventPublisher).onDeleted(any(), argThat(plugins ->
"test".equals(plugins.get(0).getName())));
+ }
+
+ @Test
+ public void testDeleteRejectsDiscoveryWithoutPluginName() {
+ SelectorDO selector = buildSelectorDO();
+
given(selectorMapper.selectByIdSet(Collections.singleton(selector.getId()))).willReturn(Collections.singletonList(selector));
+
given(pluginMapper.selectByIds(Collections.singletonList(selector.getPluginId()))).willReturn(Collections.emptyList());
+
+ ShenyuAdminException exception =
assertThrows(ShenyuAdminException.class,
+ () ->
selectorService.deleteByNamespaceId(Collections.singletonList(selector.getId()),
SYS_DEFAULT_NAMESPACE_ID));
+ assertTrue(exception.getMessage().contains("Cannot delete selector
batch"));
+ assertTrue(exception.getMessage().contains(selector.getId()));
+ assertTrue(exception.getMessage().contains("No selectors in this batch
were deleted"));
+ verifyNoInteractions(discoveryProcessor, selectorEventPublisher);
+ verify(discoveryHandlerMapper, never()).delete(any());
+ verify(selectorMapper, never()).deleteByIds(any());
+ }
+
+ @Test
+ public void testDeleteValidatesWholeBatchBeforeUnbindingDiscovery() {
+ SelectorDO first = buildSelectorDO();
+ SelectorDO second = buildSelectorDO();
+ second.setId("other-selector");
+ second.setPluginId("missing-plugin");
+ List<String> ids = Arrays.asList(first.getId(), second.getId());
+ given(selectorMapper.selectByIdSet(new
TreeSet<>(ids))).willReturn(Arrays.asList(first, second));
+ given(pluginMapper.selectByIds(Arrays.asList(first.getPluginId(),
second.getPluginId()))).willReturn(Collections.singletonList(buildPluginDO()));
+
+ ShenyuAdminException exception =
assertThrows(ShenyuAdminException.class, () ->
selectorService.deleteByNamespaceId(ids, SYS_DEFAULT_NAMESPACE_ID));
+ assertTrue(exception.getMessage().contains(second.getId()));
+ assertTrue(exception.getMessage().contains("No selectors in this batch
were deleted"));
+
+ verifyNoInteractions(discoveryProcessor, selectorEventPublisher);
+ verify(discoveryHandlerMapper, never()).delete(any());
+ verify(selectorMapper, never()).deleteByIds(any());
+ }
+
@Test
public void testFindById() {
SelectorDO selectorDO = buildSelectorDO();
diff --git
a/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonDiscoveryUpstreamDataSubscriber.java
b/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonDiscoveryUpstreamDataSubscriber.java
index aa37d7cb8c..b98aecf229 100644
---
a/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonDiscoveryUpstreamDataSubscriber.java
+++
b/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonDiscoveryUpstreamDataSubscriber.java
@@ -20,9 +20,11 @@ package org.apache.shenyu.plugin.base.cache;
import org.apache.shenyu.common.dto.DiscoverySyncData;
import org.apache.shenyu.plugin.base.handler.DiscoveryUpstreamDataHandler;
import org.apache.shenyu.sync.data.api.DiscoveryUpstreamDataSubscriber;
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamKey;
import java.util.List;
import java.util.Map;
+import java.util.Objects;
import java.util.Optional;
import java.util.stream.Collectors;
@@ -36,13 +38,20 @@ public class CommonDiscoveryUpstreamDataSubscriber
implements DiscoveryUpstreamD
@Override
public void onSubscribe(final DiscoverySyncData upstreamDataList) {
+ if (Objects.isNull(upstreamDataList) ||
Objects.isNull(upstreamDataList.getPluginName())) {
+ return;
+ }
Optional.ofNullable(handlerMap.get(upstreamDataList.getPluginName()))
.ifPresent(handler ->
handler.handlerDiscoveryUpstreamData(upstreamDataList));
}
@Override
- public void unSubscribe(final DiscoverySyncData upstreamDataList) {
- //ignore
+ public void unSubscribe(final DiscoveryUpstreamKey key) {
+ if (Objects.isNull(key) || Objects.isNull(key.pluginName())) {
+ return;
+ }
+ Optional.ofNullable(handlerMap.get(key.pluginName()))
+ .ifPresent(handler ->
handler.removeDiscoveryUpstreamData(key));
}
}
diff --git
a/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/handler/DiscoveryUpstreamDataHandler.java
b/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/handler/DiscoveryUpstreamDataHandler.java
index 476040b2e8..e0e55cf54b 100644
---
a/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/handler/DiscoveryUpstreamDataHandler.java
+++
b/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/handler/DiscoveryUpstreamDataHandler.java
@@ -18,6 +18,7 @@
package org.apache.shenyu.plugin.base.handler;
import org.apache.shenyu.common.dto.DiscoverySyncData;
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamKey;
public interface DiscoveryUpstreamDataHandler {
@@ -28,6 +29,14 @@ public interface DiscoveryUpstreamDataHandler {
*/
void handlerDiscoveryUpstreamData(DiscoverySyncData discoverySyncData);
+ /**
+ * Remove discovery upstream data.
+ *
+ * @param key discovery upstream selector identity
+ */
+ default void removeDiscoveryUpstreamData(final DiscoveryUpstreamKey key) {
+ }
+
/**
* pluginName.
*
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/DiscoveryUpstreamDataRefreshTest.java
b/shenyu-plugin/shenyu-plugin-base/src/test/java/org/apache/shenyu/plugin/base/cache/CommonDiscoveryUpstreamDataSubscriberTest.java
similarity index 55%
copy from
shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/DiscoveryUpstreamDataRefreshTest.java
copy to
shenyu-plugin/shenyu-plugin-base/src/test/java/org/apache/shenyu/plugin/base/cache/CommonDiscoveryUpstreamDataSubscriberTest.java
index a58bd4605e..4b60c425bd 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/DiscoveryUpstreamDataRefreshTest.java
+++
b/shenyu-plugin/shenyu-plugin-base/src/test/java/org/apache/shenyu/plugin/base/cache/CommonDiscoveryUpstreamDataSubscriberTest.java
@@ -15,25 +15,30 @@
* limitations under the License.
*/
-package org.apache.shenyu.sync.data.http.refresh;
+package org.apache.shenyu.plugin.base.cache;
-import org.apache.shenyu.sync.data.api.DiscoveryUpstreamDataSubscriber;
+import org.apache.shenyu.plugin.base.handler.DiscoveryUpstreamDataHandler;
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamKey;
import org.junit.jupiter.api.Test;
import java.util.Collections;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
-public final class DiscoveryUpstreamDataRefreshTest {
+public final class CommonDiscoveryUpstreamDataSubscriberTest {
@Test
- public void testRefreshWithEmptyData() {
- DiscoveryUpstreamDataSubscriber subscriber =
mock(DiscoveryUpstreamDataSubscriber.class);
- DiscoveryUpstreamDataRefresh dataRefresh = new
DiscoveryUpstreamDataRefresh(Collections.singletonList(subscriber));
+ public void testUnSubscribe() {
+ DiscoveryUpstreamDataHandler handler =
mock(DiscoveryUpstreamDataHandler.class);
+ when(handler.pluginName()).thenReturn("divide");
+ CommonDiscoveryUpstreamDataSubscriber subscriber = new
CommonDiscoveryUpstreamDataSubscriber(Collections.singletonList(handler));
+ DiscoveryUpstreamKey key = new DiscoveryUpstreamKey("divide",
"selector-id", null);
- dataRefresh.refresh(Collections.emptyList());
+ subscriber.unSubscribe(key);
- verify(subscriber).refresh();
+ verify(handler).removeDiscoveryUpstreamData(key);
}
+
}
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-divide/src/main/java/org/apache/shenyu/plugin/divide/handler/DivideUpstreamDataHandler.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-divide/src/main/java/org/apache/shenyu/plugin/divide/handler/DivideUpstreamDataHandler.java
index d6b66bbb6c..7617a83250 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-divide/src/main/java/org/apache/shenyu/plugin/divide/handler/DivideUpstreamDataHandler.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-divide/src/main/java/org/apache/shenyu/plugin/divide/handler/DivideUpstreamDataHandler.java
@@ -25,6 +25,7 @@ import
org.apache.shenyu.loadbalancer.cache.UpstreamCacheManager;
import org.apache.shenyu.loadbalancer.entity.Upstream;
import org.apache.shenyu.plugin.base.cache.MetaDataCache;
import org.apache.shenyu.plugin.base.handler.DiscoveryUpstreamDataHandler;
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamKey;
import org.springframework.util.ObjectUtils;
import java.sql.Timestamp;
@@ -59,6 +60,14 @@ public class DivideUpstreamDataHandler implements
DiscoveryUpstreamDataHandler {
MetaDataCache.getInstance().clean();
}
+ @Override
+ public void removeDiscoveryUpstreamData(final DiscoveryUpstreamKey key) {
+ if (Objects.isNull(key) || Objects.isNull(key.selectorId())) {
+ return;
+ }
+ UpstreamCacheManager.getInstance().removeByKey(key.selectorId());
+ }
+
@Override
public String pluginName() {
return PluginEnum.DIVIDE.getName();
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-divide/src/test/java/org/apache/shenyu/plugin/divide/handler/DivideUpstreamDataHandlerTest.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-divide/src/test/java/org/apache/shenyu/plugin/divide/handler/DivideUpstreamDataHandlerTest.java
index 95888e944b..a4e39078a0 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-divide/src/test/java/org/apache/shenyu/plugin/divide/handler/DivideUpstreamDataHandlerTest.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-divide/src/test/java/org/apache/shenyu/plugin/divide/handler/DivideUpstreamDataHandlerTest.java
@@ -18,6 +18,7 @@
package org.apache.shenyu.plugin.divide.handler;
import org.apache.shenyu.common.dto.DiscoverySyncData;
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamKey;
import org.apache.shenyu.common.dto.DiscoveryUpstreamData;
import org.apache.shenyu.common.enums.PluginEnum;
import org.apache.shenyu.common.utils.GsonUtils;
@@ -40,6 +41,8 @@ import java.util.stream.Collectors;
import java.util.stream.Stream;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.Mockito.mock;
@@ -93,6 +96,16 @@ public class DivideUpstreamDataHandlerTest {
divideUpstreamDataHandler.handlerDiscoveryUpstreamData(discoverySyncData);
}
+ @Test
+ public void removeDiscoveryUpstreamDataTest() {
+
divideUpstreamDataHandler.handlerDiscoveryUpstreamData(discoverySyncData);
+
assertNotNull(UpstreamCacheManager.getInstance().findUpstreamListBySelectorId("handler"));
+
+
divideUpstreamDataHandler.removeDiscoveryUpstreamData(DiscoveryUpstreamKey.from(discoverySyncData));
+
+
assertNull(UpstreamCacheManager.getInstance().findUpstreamListBySelectorId("handler"));
+ }
+
/**
* Handler discovery upstream props metadata test.
*/
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-grpc/src/main/java/org/apache/shenyu/plugin/grpc/handler/GrpcDiscoveryUpstreamDataHandler.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-grpc/src/main/java/org/apache/shenyu/plugin/grpc/handler/GrpcDiscoveryUpstreamDataHandler.java
index d2e604e8a0..9af925d983 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-grpc/src/main/java/org/apache/shenyu/plugin/grpc/handler/GrpcDiscoveryUpstreamDataHandler.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-grpc/src/main/java/org/apache/shenyu/plugin/grpc/handler/GrpcDiscoveryUpstreamDataHandler.java
@@ -26,6 +26,7 @@ import org.apache.shenyu.common.utils.JsonUtils;
import org.apache.shenyu.plugin.base.handler.DiscoveryUpstreamDataHandler;
import org.apache.shenyu.plugin.grpc.cache.ApplicationConfigCache;
import org.apache.shenyu.plugin.grpc.cache.GrpcClientCache;
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamKey;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.util.ObjectUtils;
@@ -62,6 +63,14 @@ public class GrpcDiscoveryUpstreamDataHandler implements
DiscoveryUpstreamDataHa
GrpcClientCache.initGrpcClient(selectorId);
}
+ @Override
+ public void removeDiscoveryUpstreamData(final DiscoveryUpstreamKey key) {
+ if (Objects.isNull(key) || Objects.isNull(key.selectorId())) {
+ return;
+ }
+ ApplicationConfigCache.getInstance().invalidate(key.selectorId());
+ }
+
private List<GrpcUpstream> convertUpstreamList(final
List<DiscoveryUpstreamData> upstreamList) {
if (ObjectUtils.isEmpty(upstreamList)) {
return Collections.emptyList();
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-grpc/src/test/java/org/apache/shenyu/plugin/grpc/handler/GrpcDiscoveryUpstreamDataHandlerTest.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-grpc/src/test/java/org/apache/shenyu/plugin/grpc/handler/GrpcDiscoveryUpstreamDataHandlerTest.java
new file mode 100644
index 0000000000..90b8421281
--- /dev/null
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-grpc/src/test/java/org/apache/shenyu/plugin/grpc/handler/GrpcDiscoveryUpstreamDataHandlerTest.java
@@ -0,0 +1,45 @@
+/*
+ * 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.plugin.grpc.handler;
+
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamKey;
+import org.apache.shenyu.plugin.grpc.cache.GrpcClientCache;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+
+public final class GrpcDiscoveryUpstreamDataHandlerTest {
+
+ private static final String SELECTOR_ID = "selector-id";
+
+ @AfterEach
+ public void tearDown() {
+ GrpcClientCache.removeClient(SELECTOR_ID);
+ }
+
+ @Test
+ public void testRemoveDiscoveryUpstreamData() {
+ GrpcClientCache.initGrpcClient(SELECTOR_ID);
+ assertNotNull(GrpcClientCache.getGrpcClient(SELECTOR_ID));
+ new GrpcDiscoveryUpstreamDataHandler().removeDiscoveryUpstreamData(new
DiscoveryUpstreamKey("grpc", SELECTOR_ID, null));
+
+ assertNull(GrpcClientCache.getGrpcClient(SELECTOR_ID));
+ }
+}
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/main/java/org/apache/shenyu/plugin/tcp/handler/TcpBootstrapFactory.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/main/java/org/apache/shenyu/plugin/tcp/handler/TcpBootstrapFactory.java
index fe9f43a5d6..cf0fe0ee8d 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/main/java/org/apache/shenyu/plugin/tcp/handler/TcpBootstrapFactory.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/main/java/org/apache/shenyu/plugin/tcp/handler/TcpBootstrapFactory.java
@@ -104,6 +104,7 @@ public final class TcpBootstrapFactory {
creation.complete(bootstrapServer);
return true;
} catch (RuntimeException ex) {
+ UpstreamProvider.getSingleton().removeUpstreams(selectorName);
creation.completeExceptionally(ex);
throw ex;
} finally {
@@ -164,6 +165,7 @@ public final class TcpBootstrapFactory {
*/
public boolean removeAndShutdown(final String selectorName) {
BootstrapServer bootstrapServer = cache.remove(selectorName);
+ UpstreamProvider.getSingleton().removeUpstreams(selectorName);
if (Objects.isNull(bootstrapServer)) {
return false;
}
@@ -184,6 +186,7 @@ public final class TcpBootstrapFactory {
}
}
});
+ UpstreamProvider.getSingleton().clear();
}
/**
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/main/java/org/apache/shenyu/plugin/tcp/handler/TcpUpstreamDataHandler.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/main/java/org/apache/shenyu/plugin/tcp/handler/TcpUpstreamDataHandler.java
index 7bc3c44bea..bf7a1ad02a 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/main/java/org/apache/shenyu/plugin/tcp/handler/TcpUpstreamDataHandler.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/main/java/org/apache/shenyu/plugin/tcp/handler/TcpUpstreamDataHandler.java
@@ -23,6 +23,7 @@ import org.apache.shenyu.common.enums.PluginEnum;
import org.apache.shenyu.plugin.base.handler.DiscoveryUpstreamDataHandler;
import org.apache.shenyu.protocol.tcp.BootstrapServer;
import org.apache.shenyu.protocol.tcp.UpstreamProvider;
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamKey;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -43,13 +44,40 @@ public class TcpUpstreamDataHandler implements
DiscoveryUpstreamDataHandler {
@Override
public void handlerDiscoveryUpstreamData(final DiscoverySyncData
discoverySyncData) {
- List<DiscoveryUpstreamData> removed =
UpstreamProvider.getSingleton().refreshCache(discoverySyncData.getSelectorName(),
discoverySyncData.getUpstreamDataList());
- BootstrapServer bootstrapServer =
TcpBootstrapFactory.getSingleton().getCache(discoverySyncData.getSelectorName());
+ if (Objects.isNull(discoverySyncData) ||
Objects.isNull(discoverySyncData.getSelectorName())) {
+ return;
+ }
+ final String selectorName = discoverySyncData.getSelectorName();
+ UpstreamProvider upstreamProvider = UpstreamProvider.getSingleton();
+ List<DiscoveryUpstreamData> removed =
upstreamProvider.refreshCache(selectorName,
discoverySyncData.getUpstreamDataList());
+ if (upstreamProvider.inCache(selectorName)) {
+
upstreamProvider.registerSelector(discoverySyncData.getSelectorId(),
selectorName);
+ }
+ BootstrapServer bootstrapServer =
TcpBootstrapFactory.getSingleton().getCache(selectorName);
if (Objects.nonNull(bootstrapServer)) {
bootstrapServer.removeCommonUpstream(removed);
- LOG.info("shenyu update TcpBootstrapServer [{}] success upstream
is {}", discoverySyncData.getSelectorName(),
discoverySyncData.getUpstreamDataList());
+ LOG.info("shenyu update TcpBootstrapServer [{}] success upstream
is {}", selectorName, discoverySyncData.getUpstreamDataList());
} else {
- LOG.warn("shenyu update TcpBootstrapServer don't find name is {}",
discoverySyncData.getSelectorName());
+ LOG.warn("shenyu update TcpBootstrapServer don't find name is {}",
selectorName);
+ }
+ }
+
+ @Override
+ public void removeDiscoveryUpstreamData(final DiscoveryUpstreamKey key) {
+ if (Objects.isNull(key)) {
+ return;
+ }
+ final String selectorId = key.selectorId();
+ final String localSelectorName =
UpstreamProvider.getSingleton().getSelectorName(selectorId);
+ final String selectorName = Objects.nonNull(localSelectorName) ?
localSelectorName : key.selectorName();
+ if (Objects.isNull(selectorName)) {
+ LOG.warn("shenyu remove TcpBootstrapServer upstreams don't find
selectorName by selectorId {}", selectorId);
+ return;
+ }
+ List<DiscoveryUpstreamData> removed =
UpstreamProvider.getSingleton().removeUpstreams(selectorName);
+ BootstrapServer bootstrapServer =
TcpBootstrapFactory.getSingleton().getCache(selectorName);
+ if (Objects.nonNull(bootstrapServer)) {
+ bootstrapServer.removeCommonUpstream(removed);
}
}
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/test/java/org/apache/shenyu/plugin/tcp/handler/TcpBootstrapFactoryTest.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/test/java/org/apache/shenyu/plugin/tcp/handler/TcpBootstrapFactoryTest.java
index 6bbc02fd72..95b6d7e05c 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/test/java/org/apache/shenyu/plugin/tcp/handler/TcpBootstrapFactoryTest.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/test/java/org/apache/shenyu/plugin/tcp/handler/TcpBootstrapFactoryTest.java
@@ -19,6 +19,7 @@ package org.apache.shenyu.plugin.tcp.handler;
import org.apache.shenyu.protocol.tcp.BootstrapServer;
import org.apache.shenyu.protocol.tcp.TcpServerConfiguration;
+import org.apache.shenyu.protocol.tcp.UpstreamProvider;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
@@ -37,6 +38,7 @@ import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertSame;
@@ -103,6 +105,7 @@ public final class TcpBootstrapFactoryTest {
configuration = configuration(FIRST_SELECTOR,
occupiedPort.getLocalPort());
assertThrows(RuntimeException.class, () ->
factory.createBootstrapServerIfAbsent(configuration));
assertNull(factory.getCache(FIRST_SELECTOR));
+
assertFalse(UpstreamProvider.getSingleton().inCache(FIRST_SELECTOR));
}
configuration.setPort(getFreePort());
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/test/java/org/apache/shenyu/plugin/tcp/handler/TcpProxySelectorDataHandlerTest.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/test/java/org/apache/shenyu/plugin/tcp/handler/TcpProxySelectorDataHandlerTest.java
index 76f9f2241e..ad0d8359dc 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/test/java/org/apache/shenyu/plugin/tcp/handler/TcpProxySelectorDataHandlerTest.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/test/java/org/apache/shenyu/plugin/tcp/handler/TcpProxySelectorDataHandlerTest.java
@@ -19,6 +19,7 @@ package org.apache.shenyu.plugin.tcp.handler;
import org.apache.shenyu.plugin.base.cache.CommonProxySelectorDataSubscriber;
import org.apache.shenyu.protocol.tcp.BootstrapServer;
+import org.apache.shenyu.protocol.tcp.UpstreamProvider;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
@@ -27,6 +28,7 @@ import java.util.Collections;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
@@ -37,6 +39,10 @@ public final class TcpProxySelectorDataHandlerTest {
private static final String SECOND_SELECTOR = "second";
+ private static final String FIRST_SELECTOR_ID = "first-id";
+
+ private static final String SECOND_SELECTOR_ID = "second-id";
+
private final TcpBootstrapFactory factory =
TcpBootstrapFactory.getSingleton();
@BeforeEach
@@ -55,6 +61,10 @@ public final class TcpProxySelectorDataHandlerTest {
BootstrapServer secondServer = mock(BootstrapServer.class);
factory.cache(FIRST_SELECTOR, firstServer);
factory.cache(SECOND_SELECTOR, secondServer);
+ UpstreamProvider.getSingleton().createUpstreams(FIRST_SELECTOR,
Collections.emptyList());
+ UpstreamProvider.getSingleton().createUpstreams(SECOND_SELECTOR,
Collections.emptyList());
+ UpstreamProvider.getSingleton().registerSelector(FIRST_SELECTOR_ID,
FIRST_SELECTOR);
+ UpstreamProvider.getSingleton().registerSelector(SECOND_SELECTOR_ID,
SECOND_SELECTOR);
new CommonProxySelectorDataSubscriber(Collections.singletonList(new
TcpProxySelectorDataHandler())).refresh();
@@ -62,6 +72,10 @@ public final class TcpProxySelectorDataHandlerTest {
verify(secondServer).shutdown();
assertFalse(factory.inCache(FIRST_SELECTOR));
assertFalse(factory.inCache(SECOND_SELECTOR));
+ assertFalse(UpstreamProvider.getSingleton().inCache(FIRST_SELECTOR));
+ assertFalse(UpstreamProvider.getSingleton().inCache(SECOND_SELECTOR));
+
assertNull(UpstreamProvider.getSingleton().getSelectorName(FIRST_SELECTOR_ID));
+
assertNull(UpstreamProvider.getSingleton().getSelectorName(SECOND_SELECTOR_ID));
}
@Test
@@ -71,6 +85,8 @@ public final class TcpProxySelectorDataHandlerTest {
doThrow(new IllegalStateException("shutdown
failed")).when(failingServer).shutdown();
factory.cache(FIRST_SELECTOR, failingServer);
factory.cache(SECOND_SELECTOR, secondServer);
+ UpstreamProvider.getSingleton().createUpstreams(FIRST_SELECTOR,
Collections.emptyList());
+ UpstreamProvider.getSingleton().registerSelector(FIRST_SELECTOR_ID,
FIRST_SELECTOR);
assertDoesNotThrow(() -> new TcpProxySelectorDataHandler().refresh());
@@ -78,18 +94,24 @@ public final class TcpProxySelectorDataHandlerTest {
verify(secondServer).shutdown();
assertFalse(factory.inCache(FIRST_SELECTOR));
assertFalse(factory.inCache(SECOND_SELECTOR));
+ assertFalse(UpstreamProvider.getSingleton().inCache(FIRST_SELECTOR));
+
assertNull(UpstreamProvider.getSingleton().getSelectorName(FIRST_SELECTOR_ID));
}
@Test
public void testRemoveProxySelector() {
BootstrapServer bootstrapServer = mock(BootstrapServer.class);
factory.cache(FIRST_SELECTOR, bootstrapServer);
+ UpstreamProvider.getSingleton().createUpstreams(FIRST_SELECTOR,
Collections.emptyList());
+ UpstreamProvider.getSingleton().registerSelector(FIRST_SELECTOR_ID,
FIRST_SELECTOR);
TcpProxySelectorDataHandler handler = new
TcpProxySelectorDataHandler();
handler.removeProxySelector(FIRST_SELECTOR);
verify(bootstrapServer).shutdown();
assertFalse(factory.inCache(FIRST_SELECTOR));
+ assertFalse(UpstreamProvider.getSingleton().inCache(FIRST_SELECTOR));
+
assertNull(UpstreamProvider.getSingleton().getSelectorName(FIRST_SELECTOR_ID));
assertDoesNotThrow(() -> handler.removeProxySelector(FIRST_SELECTOR));
}
}
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/test/java/org/apache/shenyu/plugin/tcp/handler/TcpUpstreamDataHandlerTest.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/test/java/org/apache/shenyu/plugin/tcp/handler/TcpUpstreamDataHandlerTest.java
index 4a7b375c9e..f9c2594d59 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/test/java/org/apache/shenyu/plugin/tcp/handler/TcpUpstreamDataHandlerTest.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/test/java/org/apache/shenyu/plugin/tcp/handler/TcpUpstreamDataHandlerTest.java
@@ -18,6 +18,7 @@
package org.apache.shenyu.plugin.tcp.handler;
import org.apache.shenyu.common.dto.DiscoverySyncData;
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamKey;
import org.apache.shenyu.common.dto.DiscoveryUpstreamData;
import org.apache.shenyu.protocol.tcp.BootstrapServer;
import org.apache.shenyu.protocol.tcp.UpstreamProvider;
@@ -33,9 +34,13 @@ import java.util.stream.Collectors;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.clearInvocations;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
/**
* Test cases for {@link TcpUpstreamDataHandler}.
@@ -44,6 +49,8 @@ public final class TcpUpstreamDataHandlerTest {
private static final String SELECTOR = "tcp-upstream-selector";
+ private static final String SELECTOR_ID = "tcp-upstream-selector-id";
+
private final TcpBootstrapFactory factory =
TcpBootstrapFactory.getSingleton();
private final TcpUpstreamDataHandler dataHandler = new
TcpUpstreamDataHandler();
@@ -51,13 +58,13 @@ public final class TcpUpstreamDataHandlerTest {
@BeforeEach
public void setUp() {
factory.clearCache();
+ // TcpBootstrapFactory initializes an empty entry before the first
discovery upstream update.
UpstreamProvider.getSingleton().createUpstreams(SELECTOR,
Collections.emptyList());
}
@AfterEach
public void tearDown() {
factory.clearCache();
- UpstreamProvider.getSingleton().createUpstreams(SELECTOR,
Collections.emptyList());
}
@Test
@@ -103,8 +110,79 @@ public final class TcpUpstreamDataHandlerTest {
assertTrue(urls.contains("127.0.0.1:10003"));
}
+ @Test
+ public void removeShouldResolveSelectorNameByIdAndEvictUpstreams() {
+ DiscoveryUpstreamData upstream = upstream("127.0.0.1:10001");
+ BootstrapServer bootstrapServer = mock(BootstrapServer.class);
+ factory.cache(SELECTOR, bootstrapServer);
+ dataHandler.handlerDiscoveryUpstreamData(syncData(upstream));
+ assertTrue(UpstreamProvider.getSingleton().inCache(SELECTOR));
+ assertEquals(SELECTOR,
UpstreamProvider.getSingleton().getSelectorName(SELECTOR_ID));
+ DiscoveryUpstreamKey deleteData = new DiscoveryUpstreamKey("tcp",
SELECTOR_ID, SELECTOR);
+
+ dataHandler.removeDiscoveryUpstreamData(deleteData);
+
+
assertTrue(UpstreamProvider.getSingleton().provide(SELECTOR).isEmpty());
+ assertFalse(UpstreamProvider.getSingleton().inCache(SELECTOR));
+
assertNull(UpstreamProvider.getSingleton().getSelectorName(SELECTOR_ID));
+
verify(bootstrapServer).removeCommonUpstream(Collections.singletonList(upstream));
+ }
+
+ @Test
+ public void removeShouldPreferLocalSelectorNameOverEventName() {
+ DiscoveryUpstreamData upstream = upstream("127.0.0.1:10001");
+ BootstrapServer bootstrapServer = mock(BootstrapServer.class);
+ BootstrapServer otherServer = mock(BootstrapServer.class);
+ factory.cache(SELECTOR, bootstrapServer);
+ factory.cache("new-name", otherServer);
+ UpstreamProvider.getSingleton().createUpstreams("new-name",
Collections.singletonList(upstream));
+ dataHandler.handlerDiscoveryUpstreamData(syncData(upstream));
+
+ dataHandler.removeDiscoveryUpstreamData(new
DiscoveryUpstreamKey("tcp", SELECTOR_ID, "new-name"));
+
+ assertFalse(UpstreamProvider.getSingleton().inCache(SELECTOR));
+ assertEquals(Collections.singletonList(upstream),
UpstreamProvider.getSingleton().provide("new-name"));
+
verify(bootstrapServer).removeCommonUpstream(Collections.singletonList(upstream));
+ verifyNoInteractions(otherServer);
+ }
+
+ @Test
+ public void removeShouldFallBackToEventNameWithoutLocalMapping() {
+ DiscoveryUpstreamData upstream = upstream("127.0.0.1:10001");
+ BootstrapServer bootstrapServer = mock(BootstrapServer.class);
+ factory.cache(SELECTOR, bootstrapServer);
+ UpstreamProvider.getSingleton().createUpstreams(SELECTOR,
Collections.singletonList(upstream));
+
+ dataHandler.removeDiscoveryUpstreamData(new
DiscoveryUpstreamKey("tcp", "unknown", SELECTOR));
+
+ assertFalse(UpstreamProvider.getSingleton().inCache(SELECTOR));
+
verify(bootstrapServer).removeCommonUpstream(Collections.singletonList(upstream));
+ }
+
+ @Test
+ public void removeShouldIgnoreUnknownSelectorId() {
+ DiscoveryUpstreamData upstream = upstream("127.0.0.1:10001");
+ BootstrapServer bootstrapServer = mock(BootstrapServer.class);
+ factory.cache(SELECTOR, bootstrapServer);
+ dataHandler.handlerDiscoveryUpstreamData(syncData(upstream));
+ clearInvocations(bootstrapServer);
+ DiscoveryUpstreamKey deleteData = new DiscoveryUpstreamKey("tcp",
"unknown", null);
+
+ dataHandler.removeDiscoveryUpstreamData(deleteData);
+
+ assertEquals(Collections.singletonList(upstream),
UpstreamProvider.getSingleton().provide(SELECTOR));
+ assertEquals(SELECTOR,
UpstreamProvider.getSingleton().getSelectorName(SELECTOR_ID));
+ verifyNoInteractions(bootstrapServer);
+ }
+
+ @Test
+ public void removeShouldIgnoreMissingSelectorIdentity() {
+ assertDoesNotThrow(() -> dataHandler.removeDiscoveryUpstreamData(new
DiscoveryUpstreamKey("tcp", null, null)));
+ }
+
private DiscoverySyncData syncData(final DiscoveryUpstreamData...
upstreams) {
DiscoverySyncData discoverySyncData = new DiscoverySyncData();
+ discoverySyncData.setSelectorId(SELECTOR_ID);
discoverySyncData.setSelectorName(SELECTOR);
discoverySyncData.setUpstreamDataList(Arrays.asList(upstreams));
return discoverySyncData;
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-websocket/src/main/java/org/apache/shenyu/plugin/websocket/handler/WebSocketUpstreamDataHandler.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-websocket/src/main/java/org/apache/shenyu/plugin/websocket/handler/WebSocketUpstreamDataHandler.java
index 923cdf67c5..5619f18e52 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-websocket/src/main/java/org/apache/shenyu/plugin/websocket/handler/WebSocketUpstreamDataHandler.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-websocket/src/main/java/org/apache/shenyu/plugin/websocket/handler/WebSocketUpstreamDataHandler.java
@@ -25,6 +25,7 @@ import
org.apache.shenyu.loadbalancer.cache.UpstreamCacheManager;
import org.apache.shenyu.loadbalancer.entity.Upstream;
import org.apache.shenyu.plugin.base.cache.MetaDataCache;
import org.apache.shenyu.plugin.base.handler.DiscoveryUpstreamDataHandler;
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamKey;
import org.springframework.util.ObjectUtils;
import java.sql.Timestamp;
@@ -56,6 +57,14 @@ public class WebSocketUpstreamDataHandler implements
DiscoveryUpstreamDataHandle
MetaDataCache.getInstance().clean();
}
+ @Override
+ public void removeDiscoveryUpstreamData(final DiscoveryUpstreamKey key) {
+ if (Objects.isNull(key) || Objects.isNull(key.selectorId())) {
+ return;
+ }
+ UpstreamCacheManager.getInstance().removeByKey(key.selectorId());
+ }
+
private List<Upstream> convertUpstreamList(final
List<DiscoveryUpstreamData> upstreamList) {
if (ObjectUtils.isEmpty(upstreamList)) {
return Collections.emptyList();
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-websocket/src/test/java/org/apache/shenyu/plugin/websocket/WebSocketPluginTest.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-websocket/src/test/java/org/apache/shenyu/plugin/websocket/WebSocketPluginTest.java
index c0841d389b..b03f75ee90 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-websocket/src/test/java/org/apache/shenyu/plugin/websocket/WebSocketPluginTest.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-websocket/src/test/java/org/apache/shenyu/plugin/websocket/WebSocketPluginTest.java
@@ -19,6 +19,7 @@ package org.apache.shenyu.plugin.websocket;
import org.apache.shenyu.common.constant.Constants;
import org.apache.shenyu.common.dto.DiscoverySyncData;
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamKey;
import org.apache.shenyu.common.dto.DiscoveryUpstreamData;
import org.apache.shenyu.common.dto.RuleData;
import org.apache.shenyu.common.dto.SelectorData;
@@ -27,6 +28,7 @@ import org.apache.shenyu.common.enums.PluginEnum;
import org.apache.shenyu.common.enums.RpcTypeEnum;
import org.apache.shenyu.common.utils.GsonUtils;
import org.apache.shenyu.common.utils.UpstreamCheckUtils;
+import org.apache.shenyu.loadbalancer.cache.UpstreamCacheManager;
import org.apache.shenyu.plugin.api.ShenyuPluginChain;
import org.apache.shenyu.plugin.api.context.ShenyuContext;
import org.apache.shenyu.plugin.websocket.handler.WebSocketPluginDataHandler;
@@ -54,6 +56,8 @@ import java.util.stream.Collectors;
import java.util.stream.Stream;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyInt;
@@ -148,6 +152,16 @@ public class WebSocketPluginTest {
assertEquals(PluginEnum.WEB_SOCKET.getName(), webSocketPlugin.named());
}
+ @Test
+ public void removeDiscoveryUpstreamDataTest() {
+ initMockInfo();
+
assertNotNull(UpstreamCacheManager.getInstance().findUpstreamListBySelectorId("mock"));
+
+ new
WebSocketUpstreamDataHandler().removeDiscoveryUpstreamData(DiscoveryUpstreamKey.from(discoverySyncData));
+
+
assertNull(UpstreamCacheManager.getInstance().findUpstreamListBySelectorId("mock"));
+ }
+
/**
* GetOrder default value test case.
*/
diff --git
a/shenyu-protocol/shenyu-protocol-tcp/src/main/java/org/apache/shenyu/protocol/tcp/UpstreamProvider.java
b/shenyu-protocol/shenyu-protocol-tcp/src/main/java/org/apache/shenyu/protocol/tcp/UpstreamProvider.java
index bdf2eeabe5..a99edd9002 100644
---
a/shenyu-protocol/shenyu-protocol-tcp/src/main/java/org/apache/shenyu/protocol/tcp/UpstreamProvider.java
+++
b/shenyu-protocol/shenyu-protocol-tcp/src/main/java/org/apache/shenyu/protocol/tcp/UpstreamProvider.java
@@ -38,6 +38,8 @@ public final class UpstreamProvider {
private final Map<String, List<DiscoveryUpstreamData>> cache = new
ConcurrentHashMap<>();
+ private final Map<String, String> selectorNames = new
ConcurrentHashMap<>();
+
private UpstreamProvider() {
}
@@ -60,6 +62,38 @@ public final class UpstreamProvider {
return cache.getOrDefault(pluginSelectorName, new ArrayList<>());
}
+ /**
+ * Whether the selector has an upstream cache entry.
+ *
+ * @param pluginSelectorName pluginSelectorName
+ * @return true if the selector has an upstream cache entry
+ */
+ public boolean inCache(final String pluginSelectorName) {
+ return cache.containsKey(pluginSelectorName);
+ }
+
+ /**
+ * Register selector name by selector id.
+ *
+ * @param selectorId selectorId
+ * @param selectorName selectorName
+ */
+ public void registerSelector(final String selectorId, final String
selectorName) {
+ if (Objects.nonNull(selectorId) && Objects.nonNull(selectorName)) {
+ selectorNames.put(selectorId, selectorName);
+ }
+ }
+
+ /**
+ * Get selector name by selector id.
+ *
+ * @param selectorId selectorId
+ * @return selectorName
+ */
+ public String getSelectorName(final String selectorId) {
+ return Objects.isNull(selectorId) ? null :
selectorNames.get(selectorId);
+ }
+
/**
* createUpstreams.
*
@@ -88,4 +122,26 @@ public final class UpstreamProvider {
Set<String> urlSet =
discoveryUpstreamDataList.stream().map(DiscoveryUpstreamData::getUrl).collect(Collectors.toSet());
return remove.stream().filter(r ->
!urlSet.contains(r.getUrl())).collect(Collectors.toList());
}
+
+ /**
+ * Remove upstreams.
+ *
+ * @param pluginSelectorName pluginSelectorName
+ * @return removed upstreams
+ */
+ public List<DiscoveryUpstreamData> removeUpstreams(final String
pluginSelectorName) {
+ if (Objects.isNull(pluginSelectorName)) {
+ return Collections.emptyList();
+ }
+ selectorNames.entrySet().removeIf(entry ->
pluginSelectorName.equals(entry.getValue()));
+ return
Optional.ofNullable(cache.remove(pluginSelectorName)).orElseGet(Collections::emptyList);
+ }
+
+ /**
+ * Clear upstreams and selector names.
+ */
+ public void clear() {
+ cache.clear();
+ selectorNames.clear();
+ }
}
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/api/DiscoveryUpstreamDataSubscriber.java
b/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/api/DiscoveryUpstreamDataSubscriber.java
index 347aa41ce1..40e3fd6e5e 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/api/DiscoveryUpstreamDataSubscriber.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/api/DiscoveryUpstreamDataSubscriber.java
@@ -31,9 +31,9 @@ public interface DiscoveryUpstreamDataSubscriber {
/**
* Un subscribe.
*
- * @param upstreamDataList the upstreamData data
+ * @param key the discovery upstream selector to remove
*/
- void unSubscribe(DiscoverySyncData upstreamDataList);
+ void unSubscribe(DiscoveryUpstreamKey key);
/**
* Refresh.
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/api/DiscoveryUpstreamDataSubscriber.java
b/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/api/DiscoveryUpstreamKey.java
similarity index 60%
copy from
shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/api/DiscoveryUpstreamDataSubscriber.java
copy to
shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/api/DiscoveryUpstreamKey.java
index 347aa41ce1..260194cb7c 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/api/DiscoveryUpstreamDataSubscriber.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/api/DiscoveryUpstreamKey.java
@@ -19,26 +19,22 @@ package org.apache.shenyu.sync.data.api;
import org.apache.shenyu.common.dto.DiscoverySyncData;
-public interface DiscoveryUpstreamDataSubscriber {
-
- /**
- * On subscribe.
- *
- * @param upstreamDataList the discoveryUpstream data
- */
- void onSubscribe(DiscoverySyncData upstreamDataList);
+/**
+ * Identity of a discovery upstream selector to remove.
+ *
+ * @param pluginName plugin name
+ * @param selectorId selector id
+ * @param selectorName optional selector name
+ */
+public record DiscoveryUpstreamKey(String pluginName, String selectorId,
String selectorName) {
/**
- * Un subscribe.
+ * Extract the deletion identity from a discovery upstream event.
*
- * @param upstreamDataList the upstreamData data
+ * @param data discovery upstream event
+ * @return deletion identity
*/
- void unSubscribe(DiscoverySyncData upstreamDataList);
-
- /**
- * Refresh.
- */
- default void refresh() {
+ public static DiscoveryUpstreamKey from(final DiscoverySyncData data) {
+ return new DiscoveryUpstreamKey(data.getPluginName(),
data.getSelectorId(), data.getSelectorName());
}
-
}
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/core/AbstractNodeDataSyncService.java
b/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/core/AbstractNodeDataSyncService.java
index 98e6926143..4b81b7429b 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/core/AbstractNodeDataSyncService.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/core/AbstractNodeDataSyncService.java
@@ -33,6 +33,7 @@ import org.apache.shenyu.common.exception.ShenyuException;
import org.apache.shenyu.common.utils.GsonUtils;
import org.apache.shenyu.sync.data.api.AuthDataSubscriber;
import org.apache.shenyu.sync.data.api.DiscoveryUpstreamDataSubscriber;
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamKey;
import org.apache.shenyu.sync.data.api.MetaDataSubscriber;
import org.apache.shenyu.sync.data.api.PluginDataSubscriber;
import org.apache.shenyu.sync.data.api.ProxySelectorDataSubscriber;
@@ -320,11 +321,9 @@ public abstract class AbstractNodeDataSyncService {
}
protected void unCacheDiscoveryUpstreamData(final String removeKey) {
- DiscoverySyncData proxySelectorData = new DiscoverySyncData();
final String[] proxySelectorKeys = StringUtils.split(removeKey,
DefaultNodeConstants.JOIN_POINT);
- proxySelectorData.setPluginName(proxySelectorKeys[2]);
- proxySelectorData.setSelectorId(proxySelectorKeys[3]);
- discoveryUpstreamDataSubscribers.forEach(e ->
e.unSubscribe(proxySelectorData));
+ DiscoveryUpstreamKey key = new
DiscoveryUpstreamKey(proxySelectorKeys[2], proxySelectorKeys[3], null);
+ discoveryUpstreamDataSubscribers.forEach(e -> e.unSubscribe(key));
removeListener(removeKey);
}
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/core/AbstractPathDataSyncService.java
b/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/core/AbstractPathDataSyncService.java
index 10c86a1aa1..7f31ab0fee 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/core/AbstractPathDataSyncService.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/core/AbstractPathDataSyncService.java
@@ -31,10 +31,13 @@ import org.apache.shenyu.common.dto.SelectorData;
import org.apache.shenyu.common.utils.GsonUtils;
import org.apache.shenyu.sync.data.api.AuthDataSubscriber;
import org.apache.shenyu.sync.data.api.DiscoveryUpstreamDataSubscriber;
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamKey;
import org.apache.shenyu.sync.data.api.MetaDataSubscriber;
import org.apache.shenyu.sync.data.api.PluginDataSubscriber;
import org.apache.shenyu.sync.data.api.ProxySelectorDataSubscriber;
import org.apache.shenyu.sync.data.api.SyncDataService;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
import java.net.URLDecoder;
import java.nio.charset.StandardCharsets;
@@ -47,6 +50,8 @@ import java.util.Optional;
*/
public abstract class AbstractPathDataSyncService implements SyncDataService {
+ private static final Logger LOG =
LoggerFactory.getLogger(AbstractPathDataSyncService.class);
+
private final PluginDataSubscriber pluginDataSubscriber;
private final List<MetaDataSubscriber> metaDataSubscribers;
@@ -130,17 +135,14 @@ public abstract class AbstractPathDataSyncService
implements SyncDataService {
}
private void discoveryUpstreamHandlerEvent(final String updatePath, final
String updateData, final EventType eventType) {
- String[] pathInfoArray = updatePath.split("/");
- if (pathInfoArray.length != 5) {
+ String[] pathInfoArray2 = updatePath.split("/");
+ if (pathInfoArray2.length != 5) {
+ LOG.warn("Ignore invalid discovery upstream path: {}", updatePath);
return;
}
- String pluginName = pathInfoArray[pathInfoArray.length - 2];
- String selectorId = pathInfoArray[pathInfoArray.length - 1];
if (EventType.DELETE.equals(eventType)) {
- DiscoverySyncData discoverySyncData = new DiscoverySyncData();
- discoverySyncData.setPluginName(pluginName);
- discoverySyncData.setSelectorId(selectorId);
- unCacheDiscoveryUpstreamData(discoverySyncData);
+ unCacheDiscoveryUpstreamData(new
DiscoveryUpstreamKey(pathInfoArray2[pathInfoArray2.length - 2],
+ pathInfoArray2[pathInfoArray2.length - 1], null));
return;
}
Optional.ofNullable(updateData)
@@ -287,9 +289,9 @@ public abstract class AbstractPathDataSyncService
implements SyncDataService {
.ifPresent(data -> discoveryUpstreamDataSubscribers.forEach(e
-> e.onSubscribe(upstreamDataList)));
}
- protected void unCacheDiscoveryUpstreamData(final DiscoverySyncData
discoverySyncData) {
- Optional.ofNullable(discoverySyncData)
- .ifPresent(data -> discoveryUpstreamDataSubscribers.forEach(e
-> e.unSubscribe(data)));
+ protected void unCacheDiscoveryUpstreamData(final DiscoveryUpstreamKey
key) {
+ Optional.ofNullable(discoveryUpstreamDataSubscribers)
+ .ifPresent(data -> discoveryUpstreamDataSubscribers.forEach(e
-> e.unSubscribe(key)));
}
protected void unCacheMetaData(final MetaData metaData) {
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-api/src/test/java/org/apache/shenyu/sync/data/core/AbstractNodeDataSyncServiceTest.java
b/shenyu-sync-data-center/shenyu-sync-data-api/src/test/java/org/apache/shenyu/sync/data/core/AbstractNodeDataSyncServiceTest.java
index 0673fd5079..075a8f4a6b 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-api/src/test/java/org/apache/shenyu/sync/data/core/AbstractNodeDataSyncServiceTest.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-api/src/test/java/org/apache/shenyu/sync/data/core/AbstractNodeDataSyncServiceTest.java
@@ -24,6 +24,7 @@ import org.apache.shenyu.common.dto.PluginData;
import org.apache.shenyu.common.dto.ProxySelectorData;
import org.apache.shenyu.sync.data.api.AuthDataSubscriber;
import org.apache.shenyu.sync.data.api.DiscoveryUpstreamDataSubscriber;
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamKey;
import org.apache.shenyu.sync.data.api.MetaDataSubscriber;
import org.apache.shenyu.sync.data.api.PluginDataSubscriber;
import org.apache.shenyu.sync.data.api.ProxySelectorDataSubscriber;
@@ -60,9 +61,11 @@ public class AbstractNodeDataSyncServiceTest {
private List<ProxySelectorDataSubscriber> proxySelectorDataSubscribers;
- @Mock
private List<DiscoveryUpstreamDataSubscriber>
discoveryUpstreamDataSubscribers;
+ @Mock
+ private DiscoveryUpstreamDataSubscriber discoveryUpstreamDataSubscriber;
+
@Mock
private ShenyuConfig shenyuConfig;
@@ -88,6 +91,8 @@ public class AbstractNodeDataSyncServiceTest {
metaDataSubscribers.add(metaDataSubscriber);
proxySelectorDataSubscribers = new ArrayList<>();
proxySelectorDataSubscribers.add(proxySelectorDataSubscriber);
+ discoveryUpstreamDataSubscribers = new ArrayList<>();
+ discoveryUpstreamDataSubscribers.add(discoveryUpstreamDataSubscriber);
nodeDataSyncService = new AbstractNodeDataSyncServiceImpl(
changeData,
@@ -161,6 +166,13 @@ public class AbstractNodeDataSyncServiceTest {
assertEquals("selectorName", captor.getValue().getName());
}
+ @Test
+ public void testUnCacheDiscoveryUpstreamData() {
+
nodeDataSyncService.unCacheDiscoveryUpstreamData("namespace.discoveryUpstream.divide.selector-id");
+
+ verify(discoveryUpstreamDataSubscriber).unSubscribe(new
DiscoveryUpstreamKey("divide", "selector-id", null));
+ }
+
@Test
public void testUnCacheDataWithInvalidKey() {
assertDoesNotThrow(() ->
nodeDataSyncService.unCachePluginData("namespace"));
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-api/src/test/java/org/apache/shenyu/sync/data/core/AbstractPathDataSyncServiceTest.java
b/shenyu-sync-data-center/shenyu-sync-data-api/src/test/java/org/apache/shenyu/sync/data/core/AbstractPathDataSyncServiceTest.java
index c1c25e9deb..8156da5a58 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-api/src/test/java/org/apache/shenyu/sync/data/core/AbstractPathDataSyncServiceTest.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-api/src/test/java/org/apache/shenyu/sync/data/core/AbstractPathDataSyncServiceTest.java
@@ -19,10 +19,10 @@ package org.apache.shenyu.sync.data.core;
import org.apache.shenyu.common.constant.DefaultPathConstants;
import org.apache.shenyu.common.dto.AppAuthData;
-import org.apache.shenyu.common.dto.DiscoverySyncData;
import org.apache.shenyu.common.utils.GsonUtils;
import org.apache.shenyu.sync.data.api.AuthDataSubscriber;
import org.apache.shenyu.sync.data.api.DiscoveryUpstreamDataSubscriber;
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamKey;
import org.apache.shenyu.sync.data.api.MetaDataSubscriber;
import org.apache.shenyu.sync.data.api.PluginDataSubscriber;
import org.apache.shenyu.sync.data.api.ProxySelectorDataSubscriber;
@@ -89,6 +89,17 @@ public class AbstractPathDataSyncServiceTest {
verify(authDataSubscriber).unSubscribe(any());
}
+ @Test
+ public void testUnCacheDiscoveryUpstreamData() {
+ pathDataSyncService.event("/default",
"/default/shenyu/discoveryUpstream/divide/selector-id", null,
+ "/default/shenyu/discoveryUpstream",
AbstractPathDataSyncService.EventType.DELETE);
+
+ ArgumentCaptor<DiscoveryUpstreamKey> captor =
ArgumentCaptor.forClass(DiscoveryUpstreamKey.class);
+ verify(discoveryUpstreamDataSubscriber).unSubscribe(captor.capture());
+ assertEquals("divide", captor.getValue().pluginName());
+ assertEquals("selector-id", captor.getValue().selectorId());
+ }
+
@Test
public void testDiscoveryUpstreamHandlerEvent() {
@@ -101,10 +112,10 @@ public class AbstractPathDataSyncServiceTest {
verify(discoveryUpstreamDataSubscriber).onSubscribe(any());
pathDataSyncService.event(namespaceId, updatePath, null, registerPath,
AbstractPathDataSyncService.EventType.DELETE);
- ArgumentCaptor<DiscoverySyncData> captor =
ArgumentCaptor.forClass(DiscoverySyncData.class);
+ ArgumentCaptor<DiscoveryUpstreamKey> captor =
ArgumentCaptor.forClass(DiscoveryUpstreamKey.class);
verify(discoveryUpstreamDataSubscriber).unSubscribe(captor.capture());
- assertEquals("divide", captor.getValue().getPluginName());
- assertEquals("testSelectorId", captor.getValue().getSelectorId());
+ assertEquals("divide", captor.getValue().pluginName());
+ assertEquals("testSelectorId", captor.getValue().selectorId());
}
// Mock implementation
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-http/src/main/java/org/apache/shenyu/sync/data/http/refresh/DiscoveryUpstreamDataRefresh.java
b/shenyu-sync-data-center/shenyu-sync-data-http/src/main/java/org/apache/shenyu/sync/data/http/refresh/DiscoveryUpstreamDataRefresh.java
index 6887412461..807d9a0f19 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-http/src/main/java/org/apache/shenyu/sync/data/http/refresh/DiscoveryUpstreamDataRefresh.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-http/src/main/java/org/apache/shenyu/sync/data/http/refresh/DiscoveryUpstreamDataRefresh.java
@@ -20,15 +20,20 @@ package org.apache.shenyu.sync.data.http.refresh;
import com.google.gson.JsonObject;
import com.google.gson.reflect.TypeToken;
import org.apache.commons.collections4.CollectionUtils;
+import org.apache.commons.lang3.StringUtils;
import org.apache.shenyu.common.dto.ConfigData;
import org.apache.shenyu.common.dto.DiscoverySyncData;
import org.apache.shenyu.common.enums.ConfigGroupEnum;
import org.apache.shenyu.common.utils.GsonUtils;
import org.apache.shenyu.sync.data.api.DiscoveryUpstreamDataSubscriber;
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamKey;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.util.HashMap;
import java.util.List;
+import java.util.Map;
+import java.util.Objects;
public class DiscoveryUpstreamDataRefresh extends
AbstractDataRefresh<DiscoverySyncData> {
@@ -36,6 +41,9 @@ public class DiscoveryUpstreamDataRefresh extends
AbstractDataRefresh<DiscoveryS
private final List<DiscoveryUpstreamDataSubscriber>
discoveryUpstreamDataSubscribers;
+ // Only the synchronized refresh method accesses this shared mutable
snapshot.
+ private final Map<String, DiscoveryUpstreamKey> previousSnapshot = new
HashMap<>();
+
public DiscoveryUpstreamDataRefresh(final
List<DiscoveryUpstreamDataSubscriber> discoveryUpstreamDataSubscribers) {
this.discoveryUpstreamDataSubscribers =
discoveryUpstreamDataSubscribers;
}
@@ -52,13 +60,34 @@ public class DiscoveryUpstreamDataRefresh extends
AbstractDataRefresh<DiscoveryS
}
@Override
- protected void refresh(final List<DiscoverySyncData> data) {
+ protected synchronized void refresh(final List<DiscoverySyncData> data) {
if (CollectionUtils.isEmpty(data)) {
- LOG.info("clear all discovery upstream data cache");
+ LOG.info("clear discovery upstream data from the HTTP snapshot");
discoveryUpstreamDataSubscribers.forEach(DiscoveryUpstreamDataSubscriber::refresh);
- return;
}
- data.forEach(d -> discoveryUpstreamDataSubscribers.forEach(dus ->
dus.onSubscribe(d)));
+ Map<String, DiscoverySyncData> currentSnapshot = new HashMap<>();
+ if (CollectionUtils.isNotEmpty(data)) {
+ data.forEach(item -> {
+ if (Objects.isNull(item) ||
StringUtils.isAnyBlank(item.getPluginName(), item.getSelectorId())) {
+ LOG.warn("ignore discovery upstream snapshot item without
pluginName or selectorId");
+ return;
+ }
+ currentSnapshot.put(selectorKey(item), item);
+ });
+ }
+ previousSnapshot.forEach((key, previous) -> {
+ DiscoverySyncData current = currentSnapshot.get(key);
+ if (Objects.isNull(current) ||
!Objects.equals(previous.selectorName(), current.getSelectorName())) {
+ discoveryUpstreamDataSubscribers.forEach(subscriber ->
subscriber.unSubscribe(previous));
+ }
+ });
+ currentSnapshot.values().forEach(item ->
discoveryUpstreamDataSubscribers.forEach(subscriber ->
subscriber.onSubscribe(item)));
+ previousSnapshot.clear();
+ currentSnapshot.forEach((key, item) -> previousSnapshot.put(key,
DiscoveryUpstreamKey.from(item)));
+ }
+
+ private static String selectorKey(final DiscoverySyncData data) {
+ return Objects.toString(data.getNamespaceId(), "") + '/' +
data.getPluginName() + '/' + data.getSelectorId();
}
@Override
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/DiscoveryUpstreamDataRefreshTest.java
b/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/DiscoveryUpstreamDataRefreshTest.java
index a58bd4605e..b6ffaabd05 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/DiscoveryUpstreamDataRefreshTest.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/DiscoveryUpstreamDataRefreshTest.java
@@ -17,12 +17,17 @@
package org.apache.shenyu.sync.data.http.refresh;
+import org.apache.shenyu.common.dto.DiscoverySyncData;
import org.apache.shenyu.sync.data.api.DiscoveryUpstreamDataSubscriber;
import org.junit.jupiter.api.Test;
import java.util.Collections;
+import java.util.List;
+import static org.mockito.ArgumentMatchers.argThat;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
public final class DiscoveryUpstreamDataRefreshTest {
@@ -31,9 +36,54 @@ public final class DiscoveryUpstreamDataRefreshTest {
public void testRefreshWithEmptyData() {
DiscoveryUpstreamDataSubscriber subscriber =
mock(DiscoveryUpstreamDataSubscriber.class);
DiscoveryUpstreamDataRefresh dataRefresh = new
DiscoveryUpstreamDataRefresh(Collections.singletonList(subscriber));
+ DiscoverySyncData discoverySyncData = discoverySyncData("selector-id",
"selector-name");
+ dataRefresh.refresh(Collections.singletonList(discoverySyncData));
dataRefresh.refresh(Collections.emptyList());
+ dataRefresh.refresh(Collections.emptyList());
+
+ verify(subscriber).onSubscribe(discoverySyncData);
+ verify(subscriber, times(2)).refresh();
+ verify(subscriber, times(1)).unSubscribe(argThat(item ->
"selector-id".equals(item.selectorId())
+ && "selector-name".equals(item.selectorName())));
+ }
+
+ @Test
+ public void testRefreshRemovesOnlyMissingSelector() {
+ DiscoveryUpstreamDataSubscriber subscriber =
mock(DiscoveryUpstreamDataSubscriber.class);
+ DiscoveryUpstreamDataRefresh dataRefresh = new
DiscoveryUpstreamDataRefresh(Collections.singletonList(subscriber));
+ DiscoverySyncData removed = discoverySyncData("removed",
"removed-name");
+ DiscoverySyncData retained = discoverySyncData("retained",
"retained-name");
+ dataRefresh.refresh(List.of(removed, retained));
+
+ dataRefresh.refresh(Collections.singletonList(retained));
+
+ verify(subscriber).onSubscribe(removed);
+ verify(subscriber, times(2)).onSubscribe(retained);
+ verify(subscriber).unSubscribe(argThat(item ->
"removed".equals(item.selectorId())));
+ verify(subscriber, never()).unSubscribe(argThat(item ->
"retained".equals(item.selectorId())));
+ }
+
+ @Test
+ public void testRefreshRemovesOldSelectorNameBeforeReplacement() {
+ DiscoveryUpstreamDataSubscriber subscriber =
mock(DiscoveryUpstreamDataSubscriber.class);
+ DiscoveryUpstreamDataRefresh dataRefresh = new
DiscoveryUpstreamDataRefresh(Collections.singletonList(subscriber));
+ DiscoverySyncData previous = discoverySyncData("selector-id",
"old-name");
+ DiscoverySyncData replacement = discoverySyncData("selector-id",
"new-name");
+ dataRefresh.refresh(Collections.singletonList(previous));
+
+ dataRefresh.refresh(Collections.singletonList(replacement));
+
+ verify(subscriber).unSubscribe(argThat(item ->
"old-name".equals(item.selectorName())));
+ verify(subscriber).onSubscribe(replacement);
+ }
- verify(subscriber).refresh();
+ private static DiscoverySyncData discoverySyncData(final String
selectorId, final String selectorName) {
+ DiscoverySyncData data = new DiscoverySyncData();
+ data.setNamespaceId("default");
+ data.setPluginName("tcp");
+ data.setSelectorId(selectorId);
+ data.setSelectorName(selectorName);
+ return data;
}
}
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/DiscoveryUpstreamDataHandler.java
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/DiscoveryUpstreamDataHandler.java
index 295f93114d..5f9a9be627 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/DiscoveryUpstreamDataHandler.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/DiscoveryUpstreamDataHandler.java
@@ -20,6 +20,7 @@ package org.apache.shenyu.plugin.sync.data.websocket.handler;
import org.apache.shenyu.common.dto.DiscoverySyncData;
import org.apache.shenyu.common.utils.GsonUtils;
import org.apache.shenyu.sync.data.api.DiscoveryUpstreamDataSubscriber;
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamKey;
import java.util.List;
@@ -39,7 +40,7 @@ public class DiscoveryUpstreamDataHandler extends
AbstractDataHandler<DiscoveryS
@Override
protected void doRefresh(final List<DiscoverySyncData> dataList) {
discoveryUpstreamDataSubscribers.forEach(DiscoveryUpstreamDataSubscriber::refresh);
- dataList.forEach(data -> discoveryUpstreamDataSubscribers.forEach(p ->
p.onSubscribe(data)));
+ doUpdate(dataList);
}
@Override
@@ -49,7 +50,7 @@ public class DiscoveryUpstreamDataHandler extends
AbstractDataHandler<DiscoveryS
@Override
protected void doDelete(final List<DiscoverySyncData> dataList) {
- dataList.forEach(data -> discoveryUpstreamDataSubscribers.forEach(p ->
p.unSubscribe(data)));
+ dataList.forEach(data -> discoveryUpstreamDataSubscribers.forEach(p ->
p.unSubscribe(DiscoveryUpstreamKey.from(data))));
}
}
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/DiscoveryUpstreamDataHandlerTest.java
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/DiscoveryUpstreamDataHandlerTest.java
index 7f0c204077..0ac98c6a8f 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/DiscoveryUpstreamDataHandlerTest.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/DiscoveryUpstreamDataHandlerTest.java
@@ -19,6 +19,7 @@ package org.apache.shenyu.plugin.sync.data.websocket.handler;
import org.apache.shenyu.common.dto.DiscoverySyncData;
import org.apache.shenyu.sync.data.api.DiscoveryUpstreamDataSubscriber;
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamKey;
import org.junit.jupiter.api.Test;
import org.mockito.InOrder;
@@ -27,6 +28,8 @@ import java.util.Collections;
import static org.mockito.Mockito.inOrder;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoMoreInteractions;
public final class DiscoveryUpstreamDataHandlerTest {
@@ -35,14 +38,37 @@ public final class DiscoveryUpstreamDataHandlerTest {
DiscoveryUpstreamDataSubscriber firstSubscriber =
mock(DiscoveryUpstreamDataSubscriber.class);
DiscoveryUpstreamDataSubscriber secondSubscriber =
mock(DiscoveryUpstreamDataSubscriber.class);
DiscoveryUpstreamDataHandler handler = new
DiscoveryUpstreamDataHandler(Arrays.asList(firstSubscriber, secondSubscriber));
- DiscoverySyncData data = new DiscoverySyncData();
+ DiscoverySyncData firstSelector = new DiscoverySyncData();
+ firstSelector.setSelectorId("first");
+ DiscoverySyncData secondSelector = new DiscoverySyncData();
+ secondSelector.setSelectorId("second");
- handler.doRefresh(Collections.singletonList(data));
+ handler.doRefresh(Collections.singletonList(firstSelector));
+ handler.doRefresh(Collections.singletonList(secondSelector));
InOrder inOrder = inOrder(firstSubscriber, secondSubscriber);
inOrder.verify(firstSubscriber).refresh();
inOrder.verify(secondSubscriber).refresh();
- inOrder.verify(firstSubscriber).onSubscribe(data);
- inOrder.verify(secondSubscriber).onSubscribe(data);
+ inOrder.verify(firstSubscriber).onSubscribe(firstSelector);
+ inOrder.verify(secondSubscriber).onSubscribe(firstSelector);
+ inOrder.verify(firstSubscriber).refresh();
+ inOrder.verify(secondSubscriber).refresh();
+ inOrder.verify(firstSubscriber).onSubscribe(secondSelector);
+ inOrder.verify(secondSubscriber).onSubscribe(secondSelector);
+ verifyNoMoreInteractions(firstSubscriber, secondSubscriber);
+ }
+
+ @Test
+ public void testDoDeleteUsesSelectorIdentity() {
+ DiscoveryUpstreamDataSubscriber subscriber =
mock(DiscoveryUpstreamDataSubscriber.class);
+ DiscoverySyncData data = new DiscoverySyncData();
+ data.setPluginName("tcp");
+ data.setSelectorId("selector-id");
+ data.setSelectorName("selector-name");
+ DiscoveryUpstreamDataHandler handler = new
DiscoveryUpstreamDataHandler(Collections.singletonList(subscriber));
+
+ handler.doDelete(Collections.singletonList(data));
+
+ verify(subscriber).unSubscribe(new DiscoveryUpstreamKey("tcp",
"selector-id", "selector-name"));
}
}
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/test/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataServiceTest.java
b/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/test/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataServiceTest.java
index 78d906d00f..076c604b81 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/test/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataServiceTest.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/test/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataServiceTest.java
@@ -24,18 +24,21 @@ import
org.apache.curator.framework.recipes.cache.TreeCacheEvent;
import org.apache.curator.framework.recipes.cache.TreeCacheListener;
import org.apache.shenyu.common.config.ShenyuConfig;
import org.apache.shenyu.common.constant.Constants;
+import org.apache.shenyu.common.constant.DefaultPathConstants;
import org.apache.shenyu.infra.zookeeper.client.ZookeeperClient;
import org.apache.shenyu.sync.data.api.AuthDataSubscriber;
+import org.apache.shenyu.sync.data.api.DiscoveryUpstreamDataSubscriber;
import org.apache.shenyu.sync.data.api.MetaDataSubscriber;
import org.apache.shenyu.sync.data.api.PluginDataSubscriber;
import org.apache.shenyu.sync.data.api.ProxySelectorDataSubscriber;
-import org.apache.shenyu.sync.data.api.DiscoveryUpstreamDataSubscriber;
import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor;
import java.util.ArrayList;
import java.util.Collections;
+import java.util.HashMap;
import java.util.List;
+import java.util.Map;
import java.util.Objects;
import static org.mockito.ArgumentMatchers.any;
@@ -43,12 +46,39 @@ import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.argThat;
import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
public final class ZookeeperSyncDataServiceTest {
+ @Test
+ public void testDiscoveryUpstreamDeleteUsesOldNodeWithoutPayload() {
+ ZookeeperClient zkClient = mock(ZookeeperClient.class);
+ Map<String, CuratorCacheListener> listeners = new HashMap<>();
+ doAnswer(invocation -> {
+ CuratorCacheListener registered = invocation.getArgument(1);
+ listeners.put(invocation.getArgument(0), registered);
+ return null;
+ }).when(zkClient).addCuratorCache(any(),
any(CuratorCacheListener[].class));
+ ShenyuConfig config = mock(ShenyuConfig.class);
+
when(config.getNamespace()).thenReturn(Constants.SYS_DEFAULT_NAMESPACE_ID);
+ DiscoveryUpstreamDataSubscriber subscriber =
mock(DiscoveryUpstreamDataSubscriber.class);
+ new ZookeeperSyncDataService(config, zkClient,
mock(PluginDataSubscriber.class), Collections.emptyList(),
+ Collections.emptyList(), Collections.emptyList(),
Collections.singletonList(subscriber));
+
+ String namespacePath = Constants.PATH_SEPARATOR +
Constants.SYS_DEFAULT_NAMESPACE_ID;
+ String registerPath = namespacePath +
DefaultPathConstants.DISCOVERY_UPSTREAM;
+ ChildData oldData = mock(ChildData.class);
+ when(oldData.getPath()).thenReturn(registerPath +
"/divide/selector-id");
+
listeners.get(registerPath).event(CuratorCacheListener.Type.NODE_DELETED,
oldData, null);
+
+ verify(subscriber).unSubscribe(argThat(key ->
"divide".equals(key.pluginName())
+ && "selector-id".equals(key.selectorId())));
+ verify(subscriber, never()).onSubscribe(any());
+ }
+
@Test
public void testDeletedNodeUsesOldData() {
ZookeeperClient zkClient = mock(ZookeeperClient.class);