This is an automated email from the ASF dual-hosted git repository.

dengliming 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 88570bc1d9 fix: shut down Consul registry executors on close (#7044)
88570bc1d9 is described below

commit 88570bc1d9baf08b2d3db50c2d1b247bcf998d49
Author: Limbo <[email protected]>
AuthorDate: Wed Sep 16 15:22:48 2026 +0800

    fix: shut down Consul registry executors on close (#7044)
---
 .../consul/ConsulInstanceRegisterRepository.java   | 20 ++++++----
 .../shenyu/registry/consul/TtlScheduler.java       |  7 ++++
 .../ConsulInstanceRegisterRepositoryTest.java      | 45 ++++++++++++++++++++++
 3 files changed, 65 insertions(+), 7 deletions(-)

diff --git 
a/shenyu-registry/shenyu-registry-consul/src/main/java/org/apache/shenyu/registry/consul/ConsulInstanceRegisterRepository.java
 
b/shenyu-registry/shenyu-registry-consul/src/main/java/org/apache/shenyu/registry/consul/ConsulInstanceRegisterRepository.java
index 73ef786628..af5c2785cc 100644
--- 
a/shenyu-registry/shenyu-registry-consul/src/main/java/org/apache/shenyu/registry/consul/ConsulInstanceRegisterRepository.java
+++ 
b/shenyu-registry/shenyu-registry-consul/src/main/java/org/apache/shenyu/registry/consul/ConsulInstanceRegisterRepository.java
@@ -139,14 +139,20 @@ public class ConsulInstanceRegisterRepository implements 
ShenyuInstanceRegisterR
 
     @Override
     public void close() {
-        if (this.running.compareAndSet(true, false) && 
!ObjectUtils.isEmpty(this.watchFutures)) {
-            this.watchFutures.forEach(watchFuture -> watchFuture.cancel(true));
-        }
-        if (!ObjectUtils.isEmpty(newService)) {
-            consulClient.agentServiceDeregister(newService.getId(), token);
-            ttlScheduler.remove(newService.getId());
+        try {
+            if (this.running.compareAndSet(true, false) && 
!ObjectUtils.isEmpty(this.watchFutures)) {
+                this.watchFutures.forEach(watchFuture -> 
watchFuture.cancel(true));
+            }
+            if (!ObjectUtils.isEmpty(newService)) {
+                consulClient.agentServiceDeregister(newService.getId(), token);
+                ttlScheduler.remove(newService.getId());
+            }
+        } finally {
+            executor.shutdownNow();
+            if (Objects.nonNull(ttlScheduler)) {
+                ttlScheduler.shutdown();
+            }
         }
-
     }
 
     private String buildInstanceNodeName(final InstanceEntity instance) {
diff --git 
a/shenyu-registry/shenyu-registry-consul/src/main/java/org/apache/shenyu/registry/consul/TtlScheduler.java
 
b/shenyu-registry/shenyu-registry-consul/src/main/java/org/apache/shenyu/registry/consul/TtlScheduler.java
index 1b425c1d42..583390fa50 100644
--- 
a/shenyu-registry/shenyu-registry-consul/src/main/java/org/apache/shenyu/registry/consul/TtlScheduler.java
+++ 
b/shenyu-registry/shenyu-registry-consul/src/main/java/org/apache/shenyu/registry/consul/TtlScheduler.java
@@ -75,6 +75,13 @@ public class TtlScheduler {
         this.serviceHeartbeats.remove(instanceId);
     }
 
+    /**
+     * Shutdown the scheduler.
+     */
+    public void shutdown() {
+        this.scheduler.shutdownNow();
+    }
+
     private class ConsulHeartbeatTask implements Runnable {
 
         private String checkId;
diff --git 
a/shenyu-registry/shenyu-registry-consul/src/test/java/org/apache/shenyu/registry/consul/ConsulInstanceRegisterRepositoryTest.java
 
b/shenyu-registry/shenyu-registry-consul/src/test/java/org/apache/shenyu/registry/consul/ConsulInstanceRegisterRepositoryTest.java
index 7badd62d2f..8d5bb3df2e 100644
--- 
a/shenyu-registry/shenyu-registry-consul/src/test/java/org/apache/shenyu/registry/consul/ConsulInstanceRegisterRepositoryTest.java
+++ 
b/shenyu-registry/shenyu-registry-consul/src/test/java/org/apache/shenyu/registry/consul/ConsulInstanceRegisterRepositoryTest.java
@@ -30,9 +30,14 @@ import java.lang.reflect.Field;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.Properties;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.ScheduledThreadPoolExecutor;
 
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyString;
+import static org.junit.jupiter.api.Assertions.assertAll;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.Mockito.doAnswer;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.mockConstruction;
@@ -115,4 +120,44 @@ public final class ConsulInstanceRegisterRepositoryTest {
             repository.close();
         }
     }
+
+    @Test
+    public void testCloseShutsDownExecutors() throws NoSuchFieldException, 
IllegalAccessException {
+        final TtlScheduler ttlScheduler = new TtlScheduler(60, 
mock(ConsulClient.class));
+        final Field ttlExecutorField = 
TtlScheduler.class.getDeclaredField("scheduler");
+        ttlExecutorField.setAccessible(true);
+        final ScheduledExecutorService ttlExecutor = 
(ScheduledExecutorService) ttlExecutorField.get(ttlScheduler);
+
+        final Field executorField = 
ConsulInstanceRegisterRepository.class.getDeclaredField("executor");
+        executorField.setAccessible(true);
+        final ScheduledThreadPoolExecutor executor = 
(ScheduledThreadPoolExecutor) executorField.get(repository);
+
+        final NewService service = new NewService();
+        service.setId("test-service");
+        final Field serviceField = 
ConsulInstanceRegisterRepository.class.getDeclaredField("newService");
+        serviceField.setAccessible(true);
+        serviceField.set(repository, service);
+        final Field ttlSchedulerField = 
ConsulInstanceRegisterRepository.class.getDeclaredField("ttlScheduler");
+        ttlSchedulerField.setAccessible(true);
+        ttlSchedulerField.set(repository, ttlScheduler);
+        final Field watchDelayField = 
ConsulInstanceRegisterRepository.class.getDeclaredField("watchDelay");
+        watchDelayField.setAccessible(true);
+        watchDelayField.set(repository, "60");
+
+        ttlScheduler.add(service.getId());
+        repository.watcherStart("test-service");
+        try {
+            assertFalse(executor.isShutdown());
+            assertFalse(ttlExecutor.isShutdown());
+
+            repository.close();
+
+            assertAll(
+                    () -> assertTrue(executor.isShutdown()),
+                    () -> assertTrue(ttlExecutor.isShutdown()));
+        } finally {
+            executor.shutdownNow();
+            ttlExecutor.shutdownNow();
+        }
+    }
 }

Reply via email to