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 b348c16a23 fix(disruptor): shut down orderly worker executors (#7229)
b348c16a23 is described below

commit b348c16a237e771e5343a176fa4da59c21f54531
Author: Liming Deng <[email protected]>
AuthorDate: Wed Sep 30 11:52:57 2026 +0800

    fix(disruptor): shut down orderly worker executors (#7229)
---
 .../shenyu/disruptor/thread/OrderlyExecutor.java   | 40 +++++++++
 .../disruptor/thread/OrderlyExecutorTest.java      | 97 ++++++++++++++++++++++
 2 files changed, 137 insertions(+)

diff --git 
a/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/thread/OrderlyExecutor.java
 
b/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/thread/OrderlyExecutor.java
index f215e6d11d..d40287c063 100644
--- 
a/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/thread/OrderlyExecutor.java
+++ 
b/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/thread/OrderlyExecutor.java
@@ -20,6 +20,9 @@ package org.apache.shenyu.disruptor.thread;
 import com.google.common.hash.Hashing;
 
 import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.HashSet;
+import java.util.List;
 import java.util.SortedMap;
 import java.util.concurrent.BlockingQueue;
 import java.util.concurrent.ConcurrentSkipListMap;
@@ -94,6 +97,43 @@ public class OrderlyExecutor extends ThreadPoolExecutor {
         }
         return virtualExecutors.get(select);
     }
+
+    @Override
+    public void shutdown() {
+        new 
HashSet<>(virtualExecutors.values()).forEach(SingletonExecutor::shutdown);
+        super.shutdown();
+    }
+
+    @Override
+    public List<Runnable> shutdownNow() {
+        List<Runnable> pending = new ArrayList<>();
+        new HashSet<>(virtualExecutors.values()).forEach(executor -> 
pending.addAll(executor.shutdownNow()));
+        pending.addAll(super.shutdownNow());
+        return pending;
+    }
+
+    @Override
+    public boolean isTerminated() {
+        return super.isTerminated() && 
virtualExecutors.values().stream().allMatch(SingletonExecutor::isTerminated);
+    }
+
+    @Override
+    public boolean awaitTermination(final long timeout, final TimeUnit unit) 
throws InterruptedException {
+        long remaining = unit.toNanos(timeout);
+        long start = System.nanoTime();
+        if (!super.awaitTermination(remaining, TimeUnit.NANOSECONDS)) {
+            return false;
+        }
+        remaining -= System.nanoTime() - start;
+        for (SingletonExecutor executor : new 
HashSet<>(virtualExecutors.values())) {
+            start = System.nanoTime();
+            if (!executor.awaitTermination(Math.max(0, remaining), 
TimeUnit.NANOSECONDS)) {
+                return false;
+            }
+            remaining -= System.nanoTime() - start;
+        }
+        return true;
+    }
     
     /**
      * The type Thread selector.
diff --git 
a/shenyu-disruptor/src/test/java/org/apache/shenyu/disruptor/thread/OrderlyExecutorTest.java
 
b/shenyu-disruptor/src/test/java/org/apache/shenyu/disruptor/thread/OrderlyExecutorTest.java
new file mode 100644
index 0000000000..f232ade4af
--- /dev/null
+++ 
b/shenyu-disruptor/src/test/java/org/apache/shenyu/disruptor/thread/OrderlyExecutorTest.java
@@ -0,0 +1,97 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.shenyu.disruptor.thread;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.Collections;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.Executors;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.RejectedExecutionException;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class OrderlyExecutorTest {
+
+    @Test
+    void testShutdownWaitsForWorkerAndRejectsNewTasks() throws Exception {
+        OrderlyExecutor executor = newExecutor();
+        CountDownLatch running = new CountDownLatch(1);
+        CountDownLatch release = new CountDownLatch(1);
+        SingletonExecutor worker = executor.select("key");
+        try {
+            worker.execute(() -> {
+                running.countDown();
+                await(release);
+            });
+            assertTrue(running.await(2, TimeUnit.SECONDS));
+            executor.shutdown();
+            assertTrue(worker.isShutdown());
+            assertFalse(executor.isTerminated());
+            assertFalse(executor.awaitTermination(1, TimeUnit.MILLISECONDS));
+            assertThrows(RejectedExecutionException.class, () -> 
worker.execute(() -> { }));
+            release.countDown();
+            assertTrue(executor.awaitTermination(2, TimeUnit.SECONDS));
+            assertTrue(executor.isTerminated());
+        } finally {
+            release.countDown();
+            executor.shutdownNow();
+        }
+    }
+
+    @Test
+    void testShutdownNowReturnsWorkerQueueAndInterruptsWorker() throws 
Exception {
+        OrderlyExecutor executor = newExecutor();
+        CountDownLatch running = new CountDownLatch(1);
+        CountDownLatch release = new CountDownLatch(1);
+        Runnable pending = () -> { };
+        try {
+            SingletonExecutor worker = executor.select("key");
+            worker.execute(() -> {
+                running.countDown();
+                await(release);
+            });
+            assertTrue(running.await(2, TimeUnit.SECONDS));
+            worker.execute(pending);
+            assertEquals(Collections.singletonList(pending), 
executor.shutdownNow());
+            assertTrue(executor.awaitTermination(2, TimeUnit.SECONDS));
+        } finally {
+            release.countDown();
+            executor.shutdownNow();
+        }
+    }
+
+    private OrderlyExecutor newExecutor() {
+        return new OrderlyExecutor(true, 2, 2, 0, TimeUnit.MILLISECONDS,
+                new LinkedBlockingQueue<>(), Executors.defaultThreadFactory(), 
new ThreadPoolExecutor.AbortPolicy());
+    }
+
+    private void await(final CountDownLatch latch) {
+        try {
+            latch.await();
+        } catch (InterruptedException e) {
+            Thread.currentThread().interrupt();
+        }
+    }
+}

Reply via email to