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 4bf86dd14c fix(disruptor): make provider startup idempotent (#7106)
4bf86dd14c is described below
commit 4bf86dd14cf43dc8e1e8442dcd2b49d5b724b082
Author: hengyuss <[email protected]>
AuthorDate: Wed Sep 30 09:57:51 2026 +0800
fix(disruptor): make provider startup idempotent (#7106)
Guard startup with a synchronized started flag to prevent resource leaks
from repeated or concurrent initialization. Add tests for repeated
startup,
concurrent calls, and retry after failure.
Fixes #6776
Co-authored-by: aias00 <[email protected]>
---
.../shenyu/disruptor/DisruptorProviderManage.java | 13 ++-
.../disruptor/DisruptorProviderManagerTest.java | 111 +++++++++++++++++++++
2 files changed, 121 insertions(+), 3 deletions(-)
diff --git
a/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/DisruptorProviderManage.java
b/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/DisruptorProviderManage.java
index 5f13b3f66b..6512d01d49 100644
---
a/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/DisruptorProviderManage.java
+++
b/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/DisruptorProviderManage.java
@@ -54,7 +54,9 @@ public class DisruptorProviderManage<T> {
private final QueueConsumerFactory<T> consumerFactory;
- private DisruptorProvider<T> provider;
+ private volatile DisruptorProvider<T> provider;
+
+ private boolean started;
/**
* Instantiates a new Disruptor provider manage.
@@ -100,11 +102,15 @@ public class DisruptorProviderManage<T> {
}
/**
- * start disruptor..
+ * Start disruptor once. Subsequent calls retain the provider and
execution mode
+ * from the first successful startup.
*
* @param isOrderly the orderly Whether to execute sequentially.
*/
- public void startup(final boolean isOrderly) {
+ public synchronized void startup(final boolean isOrderly) {
+ if (started) {
+ return;
+ }
OrderlyExecutor executor = new OrderlyExecutor(isOrderly,
consumerSize, consumerSize, 0, TimeUnit.MILLISECONDS,
new LinkedBlockingQueue<>(),
DisruptorThreadFactory.create("shenyu_disruptor_consumer_",
false), new ThreadPoolExecutor.AbortPolicy());
@@ -131,6 +137,7 @@ public class DisruptorProviderManage<T> {
disruptor.start();
RingBuffer<DataEvent<T>> ringBuffer = disruptor.getRingBuffer();
provider = new DisruptorProvider<>(ringBuffer, disruptor, isOrderly,
executor);
+ started = true;
}
/**
diff --git
a/shenyu-disruptor/src/test/java/org/apache/shenyu/disruptor/DisruptorProviderManagerTest.java
b/shenyu-disruptor/src/test/java/org/apache/shenyu/disruptor/DisruptorProviderManagerTest.java
index 87c3e46c90..725db8b521 100644
---
a/shenyu-disruptor/src/test/java/org/apache/shenyu/disruptor/DisruptorProviderManagerTest.java
+++
b/shenyu-disruptor/src/test/java/org/apache/shenyu/disruptor/DisruptorProviderManagerTest.java
@@ -17,5 +17,116 @@
package org.apache.shenyu.disruptor;
+import org.apache.shenyu.disruptor.consumer.QueueConsumerFactory;
+import org.apache.shenyu.disruptor.provider.DisruptorProvider;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
+
+import java.util.ArrayList;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Objects;
+import java.util.Set;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
class DisruptorProviderManagerTest {
+
+ @ParameterizedTest
+ @ValueSource(booleans = {false, true})
+ void testRepeatedStartupRetainsProviderAndMode(final boolean orderly) {
+ QueueConsumerFactory<String> factory =
mock(QueueConsumerFactory.class);
+ DisruptorProviderManage<String> manager = new
DisruptorProviderManage<>(factory, 1, 16);
+ Set<DisruptorProvider<String>> providers = new HashSet<>();
+ try {
+ manager.startup(orderly);
+ DisruptorProvider<String> provider = manager.getProvider();
+ providers.add(provider);
+ assertNotNull(provider);
+ manager.startup();
+ providers.add(manager.getProvider());
+ assertSame(provider, manager.getProvider());
+ manager.startup(true);
+ providers.add(manager.getProvider());
+ assertSame(provider, manager.getProvider());
+ if (orderly) {
+ assertThrows(IllegalArgumentException.class, () ->
provider.onData("data"));
+ } else {
+ assertThrows(IllegalArgumentException.class, () ->
provider.onOrderlyData("data", "key"));
+ }
+ verify(factory, times(1)).fixName();
+ } finally {
+ providers.forEach(DisruptorProvider::shutdown);
+ }
+ }
+
+ @ParameterizedTest
+ @ValueSource(booleans = {false, true})
+ void testConcurrentStartup(final boolean orderly) throws Exception {
+ QueueConsumerFactory<String> factory =
mock(QueueConsumerFactory.class);
+ DisruptorProviderManage<String> manager = new
DisruptorProviderManage<>(factory, 1, 16);
+ int callers = 8;
+ ExecutorService executor = Executors.newFixedThreadPool(callers);
+ CountDownLatch ready = new CountDownLatch(callers);
+ CountDownLatch start = new CountDownLatch(1);
+ List<Future<DisruptorProvider<String>>> results = new ArrayList<>();
+ Set<DisruptorProvider<String>> providers = new HashSet<>();
+ try {
+ for (int i = 0; i < callers; i++) {
+ results.add(executor.submit(() -> {
+ ready.countDown();
+ assertTrue(start.await(5, TimeUnit.SECONDS));
+ manager.startup(orderly);
+ return manager.getProvider();
+ }));
+ }
+ assertTrue(ready.await(5, TimeUnit.SECONDS));
+ start.countDown();
+ for (Future<DisruptorProvider<String>> result : results) {
+ providers.add(result.get(5, TimeUnit.SECONDS));
+ }
+ assertNotNull(manager.getProvider());
+ for (DisruptorProvider<String> provider : providers) {
+ assertSame(manager.getProvider(), provider);
+ }
+ verify(factory, times(1)).fixName();
+ } finally {
+ start.countDown();
+ executor.shutdownNow();
+ assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS));
+ providers.forEach(DisruptorProvider::shutdown);
+ }
+ }
+
+ @Test
+ void testStartupCanRetryAfterFailure() {
+ QueueConsumerFactory<String> factory =
mock(QueueConsumerFactory.class);
+ when(factory.fixName()).thenThrow(new IllegalStateException("startup
failed")).thenReturn("retry");
+ DisruptorProviderManage<String> manager = new
DisruptorProviderManage<>(factory, 1, 16);
+ try {
+ assertThrows(IllegalStateException.class, manager::startup);
+ assertNull(manager.getProvider());
+ manager.startup();
+ assertNotNull(manager.getProvider());
+ verify(factory, times(2)).fixName();
+ } finally {
+ if (Objects.nonNull(manager.getProvider())) {
+ manager.getProvider().shutdown();
+ }
+ }
+ }
}