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 0b496efd12 fix(client): scope registered URIs to the subscriber
instance (#7278)
0b496efd12 is described below
commit 0b496efd12217cedd1f5a578ca6b3652f5225e16
Author: Sean-Walker0 <[email protected]>
AuthorDate: Sat Sep 26 12:04:42 2026 +0800
fix(client): scope registered URIs to the subscriber instance (#7278)
ShenyuClientURIExecutorSubscriber kept its registered URI list in a
static CopyOnWriteArrayList that was only ever added to: entries were
never removed or deduplicated, and every subscriber instance in the
JVM shared the same list. After a client context restart (Spring
DevTools, integration tests, multiple apps in one JVM) the heartbeat
scheduler kept beating stale URIs from the previous context, each
re-registration appended a duplicate entry, and different client
types heartbeated each other's URIs through their own repositories.
Make the list an instance field so each subscriber tracks only its
own URIs, and skip adding URIs whose namespace/contextPath/host/port
are already registered.
Fixes #6787
Co-authored-by: Sean-Walker0
<[email protected]>
Co-authored-by: aias00 <[email protected]>
---
.../ShenyuClientURIExecutorSubscriber.java | 28 +++++++++--
.../ShenyuClientURIExecutorSubscriberTest.java | 55 ++++++++++++++++++++++
2 files changed, 78 insertions(+), 5 deletions(-)
diff --git
a/shenyu-client/shenyu-client-core/src/main/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientURIExecutorSubscriber.java
b/shenyu-client/shenyu-client-core/src/main/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientURIExecutorSubscriber.java
index 82a74fd377..08e75b4569 100644
---
a/shenyu-client/shenyu-client-core/src/main/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientURIExecutorSubscriber.java
+++
b/shenyu-client/shenyu-client-core/src/main/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientURIExecutorSubscriber.java
@@ -35,6 +35,7 @@ import java.io.IOException;
import java.net.Socket;
import java.util.Collection;
import java.util.List;
+import java.util.Objects;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.ScheduledThreadPoolExecutor;
import java.util.concurrent.ThreadFactory;
@@ -46,8 +47,14 @@ import java.util.concurrent.TimeUnit;
public class ShenyuClientURIExecutorSubscriber implements
ExecutorTypeSubscriber<URIRegisterDTO> {
private static final Logger LOG =
LoggerFactory.getLogger(ShenyuClientURIExecutorSubscriber.class);
-
- private static final List<URIRegisterDTO> URIS = new
CopyOnWriteArrayList<>();
+
+ /**
+ * URIs registered through this subscriber instance only. Instance-scoped
so that
+ * subscriber instances from different client contexts in the same JVM
never
+ * heartbeat or offline each other's URIs, and re-registered URIs do not
+ * accumulate duplicates in the heartbeat list.
+ */
+ private final List<URIRegisterDTO> uris = new CopyOnWriteArrayList<>();
private final ShenyuClientRegisterRepository
shenyuClientRegisterRepository;
@@ -64,7 +71,7 @@ public class ShenyuClientURIExecutorSubscriber implements
ExecutorTypeSubscriber
ThreadFactory requestFactory =
ShenyuThreadFactory.create("heartbeat-reporter", true);
executor = new ScheduledThreadPoolExecutor(1, requestFactory);
- executor.scheduleAtFixedRate(() -> URIS.forEach(this::sendHeartbeat),
30, 10, TimeUnit.SECONDS);
+ executor.scheduleAtFixedRate(() -> uris.forEach(this::sendHeartbeat),
30, 10, TimeUnit.SECONDS);
}
@Override
@@ -99,8 +106,8 @@ public class ShenyuClientURIExecutorSubscriber implements
ExecutorTypeSubscriber
}
ShenyuClientShutdownHook.delayOtherHooks();
shenyuClientRegisterRepository.persistURI(uriRegisterDTO);
-
- URIS.add(uriRegisterDTO);
+
+ addUriIfAbsent(uriRegisterDTO);
ShutdownHookManager.get().addShutdownHook(new Thread(() -> {
final URIRegisterDTO offlineDTO = new URIRegisterDTO();
@@ -120,4 +127,15 @@ public class ShenyuClientURIExecutorSubscriber implements
ExecutorTypeSubscriber
uriRegisterDTO.setInstanceInfo(SystemInfoUtils.getSystemInfo());
shenyuClientRegisterRepository.sendHeartbeat(uriRegisterDTO);
}
+
+ private void addUriIfAbsent(final URIRegisterDTO uriRegisterDTO) {
+ boolean alreadyRegistered = uris.stream().anyMatch(registered ->
+ Objects.equals(registered.getNamespaceId(),
uriRegisterDTO.getNamespaceId())
+ && Objects.equals(registered.getContextPath(),
uriRegisterDTO.getContextPath())
+ && Objects.equals(registered.getHost(),
uriRegisterDTO.getHost())
+ && Objects.equals(registered.getPort(),
uriRegisterDTO.getPort()));
+ if (!alreadyRegistered) {
+ uris.add(uriRegisterDTO);
+ }
+ }
}
diff --git
a/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientURIExecutorSubscriberTest.java
b/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientURIExecutorSubscriberTest.java
index b7e6f83604..6cd1949a45 100644
---
a/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientURIExecutorSubscriberTest.java
+++
b/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientURIExecutorSubscriberTest.java
@@ -26,12 +26,16 @@ import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import java.io.IOException;
+import java.lang.reflect.Field;
import java.net.ServerSocket;
import java.util.ArrayList;
import java.util.Collection;
+import java.util.Collections;
+import java.util.List;
import java.util.Properties;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
@@ -87,4 +91,55 @@ public class ShenyuClientURIExecutorSubscriberTest {
executorSubscriber.executor(uriRegisterDTOList);
verify(shenyuClientRegisterRepository,
times(1)).persistURI(uriRegisterDTO);
}
+
+ @Test
+ public void testExecutorDeduplicatesSameUri() throws Exception {
+ try (ServerSocket socket = new ServerSocket(0)) {
+ int port = socket.getLocalPort();
+ URIRegisterDTO first = buildUriRegisterDTO(port);
+ URIRegisterDTO second = buildUriRegisterDTO(port);
+
+ executorSubscriber.executor(Collections.singletonList(first));
+ executorSubscriber.executor(Collections.singletonList(second));
+
+ // registration is always persisted, but the heartbeat list must
not duplicate
+ verify(shenyuClientRegisterRepository,
times(2)).persistURI(any(URIRegisterDTO.class));
+ assertEquals(1, registeredUris(executorSubscriber).size());
+ }
+ }
+
+ @Test
+ public void testExecutorKeepsDifferentUris() throws Exception {
+ try (ServerSocket firstSocket = new ServerSocket(0);
+ ServerSocket secondSocket = new ServerSocket(0)) {
+
executorSubscriber.executor(Collections.singletonList(buildUriRegisterDTO(firstSocket.getLocalPort())));
+
executorSubscriber.executor(Collections.singletonList(buildUriRegisterDTO(secondSocket.getLocalPort())));
+
+ assertEquals(2, registeredUris(executorSubscriber).size());
+ }
+ }
+
+ @Test
+ public void testUriListIsInstanceScoped() throws Exception {
+ ShenyuClientURIExecutorSubscriber otherSubscriber = new
ShenyuClientURIExecutorSubscriber(mock(ShenyuClientRegisterRepository.class));
+ try (ServerSocket socket = new ServerSocket(0)) {
+
executorSubscriber.executor(Collections.singletonList(buildUriRegisterDTO(socket.getLocalPort())));
+
+ // URIs registered through one subscriber instance must not leak
into another instance
+ assertEquals(1, registeredUris(executorSubscriber).size());
+ assertTrue(registeredUris(otherSubscriber).isEmpty());
+ }
+ }
+
+ private URIRegisterDTO buildUriRegisterDTO(final int port) {
+ return
URIRegisterDTO.builder().protocol("http").contextPath("/test").rpcType("http")
+
.host("localhost").eventType(EventType.REGISTER).namespaceId("test-namespace").port(port).build();
+ }
+
+ @SuppressWarnings("unchecked")
+ private List<URIRegisterDTO> registeredUris(final
ShenyuClientURIExecutorSubscriber subscriber) throws Exception {
+ Field field =
ShenyuClientURIExecutorSubscriber.class.getDeclaredField("uris");
+ field.setAccessible(true);
+ return (List<URIRegisterDTO>) field.get(subscriber);
+ }
}