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 366de5f99d fix(register): propagate unregister failures (#7144)
366de5f99d is described below
commit 366de5f99d496694a15b325e7006b9909ef00399
Author: Liming Deng <[email protected]>
AuthorDate: Thu Oct 1 06:29:05 2026 +0800
fix(register): propagate unregister failures (#7144)
---
.../ShenyuClientURIExecutorSubscriber.java | 26 +++++++++-------
.../ShenyuClientURIExecutorSubscriberTest.java | 16 ++++++++++
.../client/http/HttpClientRegisterRepository.java | 7 ++++-
.../http/HttpClientRegisterRepositoryTest.java | 36 ++++++++++++++++++++++
4 files changed, 73 insertions(+), 12 deletions(-)
diff --git
a/shenyu-client/shenyu-client-core/src/main/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientURIExecutorSubscriber.java
b/shenyu-client/shenyu-client-core/src/main/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientURIExecutorSubscriber.java
index 228d929b9a..5b361e782d 100644
---
a/shenyu-client/shenyu-client-core/src/main/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientURIExecutorSubscriber.java
+++
b/shenyu-client/shenyu-client-core/src/main/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientURIExecutorSubscriber.java
@@ -107,17 +107,21 @@ public class ShenyuClientURIExecutorSubscriber implements
ExecutorTypeSubscriber
addUriIfAbsent(uriRegisterDTO);
- ShutdownHookManager.get().addShutdownHook(new Thread(() -> {
- final URIRegisterDTO offlineDTO = new URIRegisterDTO();
- BeanUtils.copyProperties(uriRegisterDTO, offlineDTO);
- offlineDTO.setEventType(EventType.OFFLINE);
- shenyuClientRegisterRepository.offline(offlineDTO);
-
- // shutdown heartbeat executor
- if (!executor.isTerminated()) {
- executor.shutdown();
- }
- }), 2);
+ ShutdownHookManager.get().addShutdownHook(new Thread(() ->
offlineAndShutdown(uriRegisterDTO)), 2);
+ }
+ }
+
+ void offlineAndShutdown(final URIRegisterDTO uriRegisterDTO) {
+ final URIRegisterDTO offlineDTO = new URIRegisterDTO();
+ BeanUtils.copyProperties(uriRegisterDTO, offlineDTO);
+ offlineDTO.setEventType(EventType.OFFLINE);
+ try {
+ shenyuClientRegisterRepository.offline(offlineDTO);
+ } finally {
+ // shutdown heartbeat executor
+ if (!executor.isTerminated()) {
+ executor.shutdown();
+ }
}
}
diff --git
a/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientURIExecutorSubscriberTest.java
b/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientURIExecutorSubscriberTest.java
index 6cd1949a45..adc490a86f 100644
---
a/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientURIExecutorSubscriberTest.java
+++
b/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientURIExecutorSubscriberTest.java
@@ -33,10 +33,13 @@ import java.util.Collection;
import java.util.Collections;
import java.util.List;
import java.util.Properties;
+import java.util.concurrent.ScheduledThreadPoolExecutor;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
@@ -92,6 +95,19 @@ public class ShenyuClientURIExecutorSubscriberTest {
verify(shenyuClientRegisterRepository,
times(1)).persistURI(uriRegisterDTO);
}
+ @Test
+ public void testOfflineFailureStillShutsDownHeartbeatExecutor() throws
Exception {
+ URIRegisterDTO uriRegisterDTO =
URIRegisterDTO.builder().host("localhost").port(9527).build();
+ doThrow(new RuntimeException("offline
failed")).when(shenyuClientRegisterRepository).offline(any());
+
+ assertThrows(RuntimeException.class, () ->
executorSubscriber.offlineAndShutdown(uriRegisterDTO));
+
+ Field executorField =
ShenyuClientURIExecutorSubscriber.class.getDeclaredField("executor");
+ executorField.setAccessible(true);
+ ScheduledThreadPoolExecutor executor = (ScheduledThreadPoolExecutor)
executorField.get(executorSubscriber);
+ assertTrue(executor.isShutdown());
+ }
+
@Test
public void testExecutorDeduplicatesSameUri() throws Exception {
try (ServerSocket socket = new ServerSocket(0)) {
diff --git
a/shenyu-register-center/shenyu-register-client/shenyu-register-client-http/src/main/java/org/apache/shenyu/register/client/http/HttpClientRegisterRepository.java
b/shenyu-register-center/shenyu-register-client/shenyu-register-client-http/src/main/java/org/apache/shenyu/register/client/http/HttpClientRegisterRepository.java
index 288f9e1690..00bd4580de 100644
---
a/shenyu-register-center/shenyu-register-client/shenyu-register-client-http/src/main/java/org/apache/shenyu/register/client/http/HttpClientRegisterRepository.java
+++
b/shenyu-register-center/shenyu-register-client/shenyu-register-client-http/src/main/java/org/apache/shenyu/register/client/http/HttpClientRegisterRepository.java
@@ -246,6 +246,7 @@ public class HttpClientRegisterRepository extends
FailbackRegistryRepository {
}
private <T> void doUnregister(final T t) {
+ int failureCount = 0;
for (String server : serverList) {
String concat = server.concat(Constants.OFFLINE_PATH);
try {
@@ -256,7 +257,11 @@ public class HttpClientRegisterRepository extends
FailbackRegistryRepository {
RegisterUtils.doUnregister(GsonUtils.getInstance().toJson(t),
concat, accessToken);
// considering the situation of multiple clusters, we should
continue to execute here
} catch (Exception e) {
- LOGGER.error("Unregister admin url :{} is fail. cause:{}",
server, e.getMessage());
+ failureCount++;
+ LOGGER.error("Unregister admin url :{} is fail.", server, e);
+ if (failureCount == serverList.size()) {
+ throw new RuntimeException(e);
+ }
}
}
}
diff --git
a/shenyu-register-center/shenyu-register-client/shenyu-register-client-http/src/test/java/org/apache/shenyu/register/client/http/HttpClientRegisterRepositoryTest.java
b/shenyu-register-center/shenyu-register-client/shenyu-register-client-http/src/test/java/org/apache/shenyu/register/client/http/HttpClientRegisterRepositoryTest.java
index 38a2ce3474..2c96875cf2 100644
---
a/shenyu-register-center/shenyu-register-client/shenyu-register-client-http/src/test/java/org/apache/shenyu/register/client/http/HttpClientRegisterRepositoryTest.java
+++
b/shenyu-register-center/shenyu-register-client/shenyu-register-client-http/src/test/java/org/apache/shenyu/register/client/http/HttpClientRegisterRepositoryTest.java
@@ -36,6 +36,7 @@ import java.io.IOException;
import java.util.Optional;
import java.util.Properties;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -239,6 +240,41 @@ public final class HttpClientRegisterRepositoryTest {
}
}
+ @Test
+ public void offlineShouldThrowWhenEveryServerFails() throws IOException {
+ try (MockedStatic<RegisterUtils> registerUtils =
mockStatic(RegisterUtils.class)) {
+ registerUtils.when(() -> RegisterUtils.doLogin(anyString(),
anyString(), anyString()))
+ .thenReturn(Optional.of(TOKEN));
+ registerUtils.when(() -> RegisterUtils.doUnregister(anyString(),
anyString(), anyString()))
+ .thenThrow(new IOException("unregister failed"));
+
+ RuntimeException exception = assertThrows(RuntimeException.class,
() -> repository.offline(uriRegisterDTO()));
+
+ assertTrue(exception.getCause() instanceof IOException);
+ assertEquals("unregister failed",
exception.getCause().getMessage());
+ }
+ }
+
+ @Test
+ public void offlineShouldNotThrowWhenOnlyLastServerFails() throws
IOException {
+ HttpClientRegisterRepository multiServerRepository = new
HttpClientRegisterRepository(config(
+ FIRST_SERVER + "," + SECOND_SERVER));
+ try (MockedStatic<RegisterUtils> registerUtils =
mockStatic(RegisterUtils.class)) {
+ registerUtils.when(() -> RegisterUtils.doLogin(anyString(),
anyString(), anyString()))
+ .thenReturn(Optional.of(TOKEN));
+ registerUtils.when(() -> RegisterUtils.doUnregister(anyString(),
+ eq(SECOND_SERVER + Constants.OFFLINE_PATH),
anyString()))
+ .thenThrow(new IOException("unregister failed"));
+
+ assertDoesNotThrow(() ->
multiServerRepository.offline(uriRegisterDTO()));
+
+ registerUtils.verify(() -> RegisterUtils.doUnregister(anyString(),
+ eq(FIRST_SERVER + Constants.OFFLINE_PATH), eq(TOKEN)));
+ registerUtils.verify(() -> RegisterUtils.doUnregister(anyString(),
+ eq(SECOND_SERVER + Constants.OFFLINE_PATH), eq(TOKEN)));
+ }
+ }
+
@Test
public void loginFailureShouldSkipRegistration() {
URIRegisterDTO uriRegisterDTO = uriRegisterDTO();