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

Reply via email to