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 8f141a1f27 fix: clear stale proxy selector data after empty HTTP
refresh. (#6909)
8f141a1f27 is described below
commit 8f141a1f277f080076bd652482c8c7fb887e08ea
Author: lymerin <[email protected]>
AuthorDate: Fri Aug 14 21:46:33 2026 +0800
fix: clear stale proxy selector data after empty HTTP refresh. (#6909)
Co-authored-by: aias00 <[email protected]>
---
.../cache/CommonProxySelectorDataSubscriber.java | 2 +-
.../base/handler/ProxySelectorDataHandler.java | 6 ++
.../CommonProxySelectorDataSubscriberTest.java | 53 ++++++++++++
.../plugin/tcp/handler/TcpBootstrapFactory.java | 19 +++++
.../tcp/handler/TcpProxySelectorDataHandler.java | 11 ++-
.../handler/TcpProxySelectorDataHandlerTest.java | 95 ++++++++++++++++++++++
.../data/http/refresh/ProxySelectorRefresh.java | 1 +
.../sync/data/http/HttpSyncDataServiceTest.java | 2 +
.../http/refresh/ProxySelectorRefreshTest.java | 74 +++++++++++++++++
9 files changed, 260 insertions(+), 3 deletions(-)
diff --git
a/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonProxySelectorDataSubscriber.java
b/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonProxySelectorDataSubscriber.java
index f9d3bef437..cf98fa3f6d 100644
---
a/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonProxySelectorDataSubscriber.java
+++
b/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonProxySelectorDataSubscriber.java
@@ -51,6 +51,6 @@ public class CommonProxySelectorDataSubscriber implements
ProxySelectorDataSubsc
@Override
public void refresh() {
- ProxySelectorDataSubscriber.super.refresh();
+ handlerMap.values().forEach(ProxySelectorDataHandler::refresh);
}
}
diff --git
a/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/handler/ProxySelectorDataHandler.java
b/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/handler/ProxySelectorDataHandler.java
index fb03a41895..8ce1e0bca5 100644
---
a/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/handler/ProxySelectorDataHandler.java
+++
b/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/handler/ProxySelectorDataHandler.java
@@ -39,6 +39,12 @@ public interface ProxySelectorDataHandler {
*/
void removeProxySelector(String proxySelectorName);
+ /**
+ * Refresh.
+ */
+ default void refresh() {
+ }
+
/**
* pluginName.
*
diff --git
a/shenyu-plugin/shenyu-plugin-base/src/test/java/org/apache/shenyu/plugin/base/cache/CommonProxySelectorDataSubscriberTest.java
b/shenyu-plugin/shenyu-plugin-base/src/test/java/org/apache/shenyu/plugin/base/cache/CommonProxySelectorDataSubscriberTest.java
new file mode 100644
index 0000000000..dec32d6a1d
--- /dev/null
+++
b/shenyu-plugin/shenyu-plugin-base/src/test/java/org/apache/shenyu/plugin/base/cache/CommonProxySelectorDataSubscriberTest.java
@@ -0,0 +1,53 @@
+/*
+ * 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.base.cache;
+
+import org.apache.shenyu.plugin.base.handler.ProxySelectorDataHandler;
+import org.junit.jupiter.api.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+public final class CommonProxySelectorDataSubscriberTest {
+
+ @Test
+ public void testRefresh() {
+ ProxySelectorDataHandler firstHandler =
mock(ProxySelectorDataHandler.class);
+ ProxySelectorDataHandler secondHandler =
mock(ProxySelectorDataHandler.class);
+ when(firstHandler.pluginName()).thenReturn("first");
+ when(secondHandler.pluginName()).thenReturn("second");
+ CommonProxySelectorDataSubscriber subscriber = new
CommonProxySelectorDataSubscriber(Arrays.asList(firstHandler, secondHandler));
+
+ subscriber.refresh();
+
+ verify(firstHandler).refresh();
+ verify(secondHandler).refresh();
+ }
+
+ @Test
+ public void testRefreshWithoutHandlers() {
+ CommonProxySelectorDataSubscriber subscriber = new
CommonProxySelectorDataSubscriber(Collections.emptyList());
+
+ assertDoesNotThrow(subscriber::refresh);
+ }
+}
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 7c4739f19d..070ada7064 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
@@ -21,6 +21,8 @@ import com.google.common.eventbus.EventBus;
import org.apache.shenyu.protocol.tcp.BootstrapServer;
import org.apache.shenyu.protocol.tcp.TcpBootstrapServer;
import org.apache.shenyu.protocol.tcp.TcpServerConfiguration;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@@ -30,6 +32,8 @@ import java.util.concurrent.ConcurrentHashMap;
*/
public final class TcpBootstrapFactory {
+ private static final Logger LOG =
LoggerFactory.getLogger(TcpBootstrapFactory.class);
+
private static final TcpBootstrapFactory SINGLETON = new
TcpBootstrapFactory();
private final Map<String, BootstrapServer> cache = new
ConcurrentHashMap<>();
@@ -89,6 +93,21 @@ public final class TcpBootstrapFactory {
return cache.remove(selectorName);
}
+ /**
+ * Clear cache.
+ */
+ public void clearCache() {
+ cache.forEach((selectorName, bootstrapServer) -> {
+ if (cache.remove(selectorName, bootstrapServer)) {
+ try {
+ bootstrapServer.shutdown();
+ } catch (RuntimeException ex) {
+ LOG.error("Failed to shutdown TcpBootstrapServer for
selector {}", selectorName, ex);
+ }
+ }
+ });
+ }
+
/**
* getCache.
*
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/main/java/org/apache/shenyu/plugin/tcp/handler/TcpProxySelectorDataHandler.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/main/java/org/apache/shenyu/plugin/tcp/handler/TcpProxySelectorDataHandler.java
index 96a912c91b..78d4f2aefd 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/main/java/org/apache/shenyu/plugin/tcp/handler/TcpProxySelectorDataHandler.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/main/java/org/apache/shenyu/plugin/tcp/handler/TcpProxySelectorDataHandler.java
@@ -27,6 +27,7 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.Collections;
+import java.util.Objects;
public class TcpProxySelectorDataHandler implements ProxySelectorDataHandler {
@@ -52,12 +53,18 @@ public class TcpProxySelectorDataHandler implements
ProxySelectorDataHandler {
@Override
public void removeProxySelector(final String proxySelectorName) {
- if (TcpBootstrapFactory.getSingleton().inCache(proxySelectorName)) {
-
TcpBootstrapFactory.getSingleton().removeCache(proxySelectorName).shutdown();
+ BootstrapServer bootstrapServer =
TcpBootstrapFactory.getSingleton().removeCache(proxySelectorName);
+ if (Objects.nonNull(bootstrapServer)) {
+ bootstrapServer.shutdown();
LOG.info("shenyu shutdown {}", proxySelectorName);
}
}
+ @Override
+ public void refresh() {
+ TcpBootstrapFactory.getSingleton().clearCache();
+ }
+
@Override
public String pluginName() {
return PluginEnum.TCP.getName();
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
new file mode 100644
index 0000000000..76f9f2241e
--- /dev/null
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/test/java/org/apache/shenyu/plugin/tcp/handler/TcpProxySelectorDataHandlerTest.java
@@ -0,0 +1,95 @@
+/*
+ * 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.tcp.handler;
+
+import org.apache.shenyu.plugin.base.cache.CommonProxySelectorDataSubscriber;
+import org.apache.shenyu.protocol.tcp.BootstrapServer;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.util.Collections;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+
+public final class TcpProxySelectorDataHandlerTest {
+
+ private static final String FIRST_SELECTOR = "first";
+
+ private static final String SECOND_SELECTOR = "second";
+
+ private final TcpBootstrapFactory factory =
TcpBootstrapFactory.getSingleton();
+
+ @BeforeEach
+ public void setUp() {
+ factory.clearCache();
+ }
+
+ @AfterEach
+ public void tearDown() {
+ factory.clearCache();
+ }
+
+ @Test
+ public void testRefreshThroughSubscriber() {
+ BootstrapServer firstServer = mock(BootstrapServer.class);
+ BootstrapServer secondServer = mock(BootstrapServer.class);
+ factory.cache(FIRST_SELECTOR, firstServer);
+ factory.cache(SECOND_SELECTOR, secondServer);
+
+ new CommonProxySelectorDataSubscriber(Collections.singletonList(new
TcpProxySelectorDataHandler())).refresh();
+
+ verify(firstServer).shutdown();
+ verify(secondServer).shutdown();
+ assertFalse(factory.inCache(FIRST_SELECTOR));
+ assertFalse(factory.inCache(SECOND_SELECTOR));
+ }
+
+ @Test
+ public void testRefreshContinuesWhenShutdownFails() {
+ BootstrapServer failingServer = mock(BootstrapServer.class);
+ BootstrapServer secondServer = mock(BootstrapServer.class);
+ doThrow(new IllegalStateException("shutdown
failed")).when(failingServer).shutdown();
+ factory.cache(FIRST_SELECTOR, failingServer);
+ factory.cache(SECOND_SELECTOR, secondServer);
+
+ assertDoesNotThrow(() -> new TcpProxySelectorDataHandler().refresh());
+
+ verify(failingServer).shutdown();
+ verify(secondServer).shutdown();
+ assertFalse(factory.inCache(FIRST_SELECTOR));
+ assertFalse(factory.inCache(SECOND_SELECTOR));
+ }
+
+ @Test
+ public void testRemoveProxySelector() {
+ BootstrapServer bootstrapServer = mock(BootstrapServer.class);
+ factory.cache(FIRST_SELECTOR, bootstrapServer);
+ TcpProxySelectorDataHandler handler = new
TcpProxySelectorDataHandler();
+
+ handler.removeProxySelector(FIRST_SELECTOR);
+
+ verify(bootstrapServer).shutdown();
+ assertFalse(factory.inCache(FIRST_SELECTOR));
+ assertDoesNotThrow(() -> handler.removeProxySelector(FIRST_SELECTOR));
+ }
+}
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-http/src/main/java/org/apache/shenyu/sync/data/http/refresh/ProxySelectorRefresh.java
b/shenyu-sync-data-center/shenyu-sync-data-http/src/main/java/org/apache/shenyu/sync/data/http/refresh/ProxySelectorRefresh.java
index f2387dde20..e4b8ae8407 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-http/src/main/java/org/apache/shenyu/sync/data/http/refresh/ProxySelectorRefresh.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-http/src/main/java/org/apache/shenyu/sync/data/http/refresh/ProxySelectorRefresh.java
@@ -55,6 +55,7 @@ public class ProxySelectorRefresh extends
AbstractDataRefresh<ProxySelectorData>
protected void refresh(final List<ProxySelectorData> data) {
if (CollectionUtils.isEmpty(data)) {
LOG.info("clear all ProxySelector data cache");
+
proxySelectorDataSubscribers.forEach(ProxySelectorDataSubscriber::refresh);
return;
}
data.forEach(d -> proxySelectorDataSubscribers.forEach(pss ->
pss.onSubscribe(d)));
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/HttpSyncDataServiceTest.java
b/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/HttpSyncDataServiceTest.java
index 4754feabff..db14f2d93b 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/HttpSyncDataServiceTest.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/HttpSyncDataServiceTest.java
@@ -162,6 +162,7 @@ public final class HttpSyncDataServiceTest {
verify(pluginDataSubscriber, atLeastOnce()).refreshPluginDataAll();
verify(metaDataSubscriber, atLeastOnce()).refresh();
verify(authDataSubscriber, atLeastOnce()).refresh();
+ verify(proxySelectorDataSubscriber, atLeastOnce()).refresh();
}
private String getMockServerUrl() {
@@ -193,6 +194,7 @@ public final class HttpSyncDataServiceTest {
data.put(ConfigGroupEnum.APP_AUTH.name(), emptyData);
data.put(ConfigGroupEnum.SELECTOR.name(), emptyData);
data.put(ConfigGroupEnum.RULE.name(), emptyData);
+ data.put(ConfigGroupEnum.PROXY_SELECTOR.name(), emptyData);
Map<String, Object> response = new HashMap<>();
response.put("data", data);
response.put("code", 200);
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/ProxySelectorRefreshTest.java
b/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/ProxySelectorRefreshTest.java
new file mode 100644
index 0000000000..c1ec7e0ca2
--- /dev/null
+++
b/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/ProxySelectorRefreshTest.java
@@ -0,0 +1,74 @@
+/*
+ * 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.sync.data.http.refresh;
+
+import org.apache.shenyu.common.dto.ProxySelectorData;
+import org.apache.shenyu.sync.data.api.ProxySelectorDataSubscriber;
+import org.junit.jupiter.api.Test;
+
+import java.util.Collections;
+import java.util.List;
+
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Tests for {@link ProxySelectorRefresh}.
+ */
+public final class ProxySelectorRefreshTest {
+
+ private final StubProxySelectorDataSubscriber subscriber = new
StubProxySelectorDataSubscriber();
+
+ private final ProxySelectorRefresh proxySelectorRefresh =
+ new ProxySelectorRefresh(Collections.singletonList(subscriber));
+
+ @Test
+ public void testRefreshWithEmptyDataShouldClearSubscriberCache() {
+ List<ProxySelectorData> empty = Collections.emptyList();
+ proxySelectorRefresh.refresh(empty);
+ assertTrue(subscriber.refreshed, "empty snapshot must invoke
subscriber.refresh() to clear stale data");
+ assertFalse(subscriber.subscribed, "empty snapshot must not subscribe
data");
+ }
+
+ @Test
+ public void testRefreshWithDataShouldSubscribe() {
+ proxySelectorRefresh.refresh(Collections.singletonList(new
ProxySelectorData()));
+ assertTrue(subscriber.subscribed);
+ }
+
+ private static final class StubProxySelectorDataSubscriber implements
ProxySelectorDataSubscriber {
+
+ private boolean refreshed;
+
+ private boolean subscribed;
+
+ @Override
+ public void onSubscribe(final ProxySelectorData proxySelectorData) {
+ subscribed = true;
+ }
+
+ @Override
+ public void unSubscribe(final ProxySelectorData proxySelectorData) {
+ }
+
+ @Override
+ public void refresh() {
+ refreshed = true;
+ }
+ }
+}