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