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