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);
+    }
 }

Reply via email to