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