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 6fb0361d5c  fix: concurrent access to discovery sync data listeners 
(#7197)
6fb0361d5c is described below

commit 6fb0361d5cbe3733043e2fef9a651f7185db3e94
Author: Southern <[email protected]>
AuthorDate: Wed Sep 30 09:58:56 2026 +0800

     fix: concurrent access to discovery sync data listeners (#7197)
    
    Replace the non-thread-safe ArrayList in 
DiscoveryDataChangedEventSyncListener with CopyOnWriteArrayList and add a
      deterministic concurrency test covering listener registration during 
onChange iteration.
    
    Co-authored-by: aias00 <[email protected]>
---
 .../DiscoveryDataChangedEventSyncListener.java     |  4 +--
 .../DiscoveryDataChangedEventSyncListenerTest.java | 41 ++++++++++++++++++++++
 2 files changed, 43 insertions(+), 2 deletions(-)

diff --git 
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/discovery/DiscoveryDataChangedEventSyncListener.java
 
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/discovery/DiscoveryDataChangedEventSyncListener.java
index 78f47d40e0..d3ba6f0916 100644
--- 
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/discovery/DiscoveryDataChangedEventSyncListener.java
+++ 
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/discovery/DiscoveryDataChangedEventSyncListener.java
@@ -38,10 +38,10 @@ import org.springframework.dao.DuplicateKeyException;
 import org.springframework.transaction.annotation.Transactional;
 
 import java.sql.Timestamp;
-import java.util.ArrayList;
 import java.util.Collections;
 import java.util.List;
 import java.util.Objects;
+import java.util.concurrent.CopyOnWriteArrayList;
 import java.util.stream.Collectors;
 
 /**
@@ -66,7 +66,7 @@ public class DiscoveryDataChangedEventSyncListener implements 
DataChangedEventLi
                                                  final KeyValueParser 
keyValueParser,
                                                  final DiscoverySyncData 
contextInfo,
                                                  final String discoveryId) {
-        this.discoverySyncDataList = new ArrayList<>();
+        this.discoverySyncDataList = new CopyOnWriteArrayList<>();
         this.eventPublisher = eventPublisher;
         this.keyValueParser = keyValueParser;
         this.discoveryId = discoveryId;
diff --git 
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/discovery/DiscoveryDataChangedEventSyncListenerTest.java
 
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/discovery/DiscoveryDataChangedEventSyncListenerTest.java
index 3924428c27..1a4286f533 100644
--- 
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/discovery/DiscoveryDataChangedEventSyncListenerTest.java
+++ 
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/discovery/DiscoveryDataChangedEventSyncListenerTest.java
@@ -39,6 +39,10 @@ import org.springframework.context.ApplicationEventPublisher;
 import java.util.ArrayList;
 import java.util.Collections;
 import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
 
 import static 
org.apache.shenyu.common.constant.Constants.SYS_DEFAULT_NAMESPACE_ID;
 import static org.mockito.ArgumentMatchers.any;
@@ -115,4 +119,41 @@ public class DiscoveryDataChangedEventSyncListenerTest {
         Assertions.assertEquals(namespaceId, 
discoveryUpstreamCaptor.getValue().getNamespaceId());
     }
 
+    @Test
+    public void testOnChangeIsSafeWhenListenerIsAddedConcurrently() throws 
Exception {
+        DiscoverySyncData additionalContext = 
org.mockito.Mockito.mock(DiscoverySyncData.class);
+        
when(contextInfo.getNamespaceId()).thenReturn(SYS_DEFAULT_NAMESPACE_ID);
+        
when(contextInfo.getDiscoveryHandlerId()).thenReturn("discoveryHandlerId");
+        when(contextInfo.getSelectorId()).thenReturn("selector-1");
+        
when(additionalContext.getNamespaceId()).thenReturn(SYS_DEFAULT_NAMESPACE_ID);
+        
when(additionalContext.getDiscoveryHandlerId()).thenReturn("discoveryHandlerId");
+        when(additionalContext.getSelectorId()).thenReturn("selector-2");
+        DiscoveryUpstreamData upstreamData = new DiscoveryUpstreamData();
+        upstreamData.setProtocol("http://";);
+        upstreamData.setNamespaceId(SYS_DEFAULT_NAMESPACE_ID);
+        upstreamData.setUrl("127.0.0.1:8080");
+        final CountDownLatch processingStarted = new CountDownLatch(1);
+        final CountDownLatch continueProcessing = new CountDownLatch(1);
+        org.mockito.Mockito.doAnswer(invocation -> {
+            processingStarted.countDown();
+            continueProcessing.await();
+            return Collections.singletonList(upstreamData);
+        }).when(keyValueParser).parseValue(anyString());
+
+        ExecutorService executor = Executors.newSingleThreadExecutor();
+        try {
+            final java.util.concurrent.Future<?> change = executor.submit(() 
-> discoveryDataChangedEventSyncListener.onChange(
+                    new DiscoveryDataChangedEvent("key", "value", 
DiscoveryDataChangedEvent.Event.ADDED)));
+            Assertions.assertTrue(processingStarted.await(1, 
TimeUnit.SECONDS));
+            
discoveryDataChangedEventSyncListener.addListener(additionalContext);
+            continueProcessing.countDown();
+            Assertions.assertDoesNotThrow(() -> change.get(1, 
TimeUnit.SECONDS));
+        } finally {
+            continueProcessing.countDown();
+            executor.shutdownNow();
+        }
+
+        verify(discoveryUpstreamMapper).insert(any(DiscoveryUpstreamDO.class));
+    }
+
 }

Reply via email to