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 a12aa3e595 fix(register): propagate partial multi-server failures
(#7260)
a12aa3e595 is described below
commit a12aa3e595e59f2b82ea2e73835770fe1fdcb88a
Author: Liming Deng <[email protected]>
AuthorDate: Thu Oct 1 15:03:52 2026 +0800
fix(register): propagate partial multi-server failures (#7260)
---
.../ShenyuClientURIExecutorSubscriber.java | 10 ++-
.../subcriber/HeartbeatFailureIsolationTest.java | 70 ++++++++++++++++++
.../subcriber/UriReadinessTimeoutTest.java | 57 ++++++++-------
.../register/client/beat/HeartbeatListener.java | 8 +--
.../client/beat/HeartbeatListenerTest.java | 35 +++++++--
.../client/http/HttpClientRegisterRepository.java | 25 ++++---
.../http/HttpClientRegisterRepositoryTest.java | 84 ++++++++++++++++++++++
7 files changed, 238 insertions(+), 51 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 5b361e782d..743f5bcbe5 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
@@ -155,8 +155,14 @@ public class ShenyuClientURIExecutorSubscriber implements
ExecutorTypeSubscriber
}
private void sendHeartbeat(final URIRegisterDTO uriRegisterDTO) {
- uriRegisterDTO.setInstanceInfo(SystemInfoUtils.getSystemInfo());
- shenyuClientRegisterRepository.sendHeartbeat(uriRegisterDTO);
+ try {
+ uriRegisterDTO.setInstanceInfo(SystemInfoUtils.getSystemInfo());
+ shenyuClientRegisterRepository.sendHeartbeat(uriRegisterDTO);
+ } catch (Exception ex) {
+ // One unavailable admin must not suppress other URIs or future
scheduled executions.
+ LOG.warn("Heartbeat failed for host:{}, port:{}, will retry on the
next tick",
+ uriRegisterDTO.getHost(), uriRegisterDTO.getPort(), ex);
+ }
}
private void addUriIfAbsent(final URIRegisterDTO uriRegisterDTO) {
diff --git
a/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/HeartbeatFailureIsolationTest.java
b/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/HeartbeatFailureIsolationTest.java
new file mode 100644
index 0000000000..5a7520a0a5
--- /dev/null
+++
b/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/HeartbeatFailureIsolationTest.java
@@ -0,0 +1,70 @@
+/*
+ * 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.client.core.disruptor.subcriber;
+
+import org.apache.shenyu.common.utils.SystemInfoUtils;
+import org.apache.shenyu.register.client.api.ShenyuClientRegisterRepository;
+import org.apache.shenyu.register.common.dto.URIRegisterDTO;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+import org.springframework.test.util.ReflectionTestUtils;
+
+import java.util.List;
+import java.util.concurrent.RunnableScheduledFuture;
+import java.util.concurrent.ScheduledThreadPoolExecutor;
+
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+
+/**
+ * Verify periodic heartbeats survive repository failures.
+ */
+public final class HeartbeatFailureIsolationTest {
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void failureDoesNotSuppressOtherUrisOrSubsequentTicks() {
+ ShenyuClientRegisterRepository repository =
mock(ShenyuClientRegisterRepository.class);
+ ShenyuClientURIExecutorSubscriber subscriber = new
ShenyuClientURIExecutorSubscriber(repository);
+ ScheduledThreadPoolExecutor executor = (ScheduledThreadPoolExecutor)
ReflectionTestUtils.getField(subscriber, "executor");
+ List<URIRegisterDTO> uris = (List<URIRegisterDTO>)
ReflectionTestUtils.getField(subscriber, "uris");
+ URIRegisterDTO failing =
URIRegisterDTO.builder().host("localhost").port(18080).build();
+ URIRegisterDTO healthy =
URIRegisterDTO.builder().host("localhost").port(18081).build();
+ doThrow(new IllegalStateException("admin
unavailable")).when(repository).sendHeartbeat(failing);
+ try (MockedStatic<SystemInfoUtils> ignored =
mockStatic(SystemInfoUtils.class)) {
+ uris.clear();
+ uris.add(failing);
+ uris.add(healthy);
+ RunnableScheduledFuture<?> task = (RunnableScheduledFuture<?>)
executor.getQueue().iterator().next();
+ for (int tick = 0; tick < 2; tick++) {
+ executor.getQueue().remove(task);
+ task.run();
+ assertFalse(task.isDone(), "Periodic task must remain
schedulable after a failed heartbeat");
+ }
+ verify(repository, times(2)).sendHeartbeat(failing);
+ verify(repository, times(2)).sendHeartbeat(healthy);
+ } finally {
+ executor.shutdownNow();
+ uris.clear();
+ }
+ }
+}
diff --git
a/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/UriReadinessTimeoutTest.java
b/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/UriReadinessTimeoutTest.java
index bb95242e89..4cddf5d35d 100644
---
a/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/UriReadinessTimeoutTest.java
+++
b/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/UriReadinessTimeoutTest.java
@@ -25,6 +25,8 @@ import org.junit.jupiter.api.Test;
import org.springframework.test.util.ReflectionTestUtils;
import java.net.ServerSocket;
+import java.net.InetSocketAddress;
+import java.net.Socket;
import java.time.Duration;
import java.util.List;
import java.util.Objects;
@@ -84,12 +86,10 @@ public final class UriReadinessTimeoutTest {
public void testUnreachableUriDoesNotBlockNextUri() throws Exception {
subscriber = new ShenyuClientURIExecutorSubscriber(repository, 100);
ShenyuClientShutdownHook.set(repository, new Properties());
- int closedPort;
- try (ServerSocket closed = new ServerSocket(0)) {
- closedPort = closed.getLocalPort();
- }
- try (ServerSocket ready = new ServerSocket(0)) {
- URIRegisterDTO unavailable = uri(closedPort);
+ try (Socket unavailableSocket = new Socket(); ServerSocket ready = new
ServerSocket(0)) {
+ // Reserve a port without listening so the ready server cannot
reuse the unavailable endpoint.
+ unavailableSocket.bind(new InetSocketAddress("127.0.0.1", 0));
+ URIRegisterDTO unavailable = uri(unavailableSocket.getLocalPort());
URIRegisterDTO available = uri(ready.getLocalPort());
assertTimeoutPreemptively(Duration.ofSeconds(3), () ->
subscriber.executor(List.of(unavailable, available)));
verify(repository, never()).persistURI(unavailable);
@@ -100,29 +100,28 @@ public final class UriReadinessTimeoutTest {
@Test
public void testInterruptionStopsWaitingAndPreservesFlag() throws
Exception {
subscriber = new ShenyuClientURIExecutorSubscriber(repository, 30000);
- int closedPort;
- try (ServerSocket closed = new ServerSocket(0)) {
- closedPort = closed.getLocalPort();
- }
- URIRegisterDTO unavailable = uri(closedPort);
- AtomicBoolean interrupted = new AtomicBoolean();
- CountDownLatch started = new CountDownLatch(1);
- Thread worker = new Thread(() -> {
- started.countDown();
- subscriber.executor(List.of(unavailable));
- interrupted.set(Thread.currentThread().isInterrupted());
- });
- worker.start();
- try {
- assertTrue(started.await(1, TimeUnit.SECONDS));
- worker.interrupt();
- worker.join(2000);
- assertFalse(worker.isAlive());
- assertTrue(interrupted.get());
- verify(repository, never()).persistURI(unavailable);
- } finally {
- worker.interrupt();
- worker.join(2000);
+ try (Socket unavailableSocket = new Socket()) {
+ unavailableSocket.bind(new InetSocketAddress("127.0.0.1", 0));
+ URIRegisterDTO unavailable = uri(unavailableSocket.getLocalPort());
+ AtomicBoolean interrupted = new AtomicBoolean();
+ CountDownLatch started = new CountDownLatch(1);
+ Thread worker = new Thread(() -> {
+ started.countDown();
+ subscriber.executor(List.of(unavailable));
+ interrupted.set(Thread.currentThread().isInterrupted());
+ });
+ worker.start();
+ try {
+ assertTrue(started.await(1, TimeUnit.SECONDS));
+ worker.interrupt();
+ worker.join(2000);
+ assertFalse(worker.isAlive());
+ assertTrue(interrupted.get());
+ verify(repository, never()).persistURI(unavailable);
+ } finally {
+ worker.interrupt();
+ worker.join(2000);
+ }
}
}
diff --git
a/shenyu-register-center/shenyu-register-client-beat/src/main/java/org/apache/shenyu/register/client/beat/HeartbeatListener.java
b/shenyu-register-center/shenyu-register-client-beat/src/main/java/org/apache/shenyu/register/client/beat/HeartbeatListener.java
index 36a54e2fcf..67e971a065 100644
---
a/shenyu-register-center/shenyu-register-client-beat/src/main/java/org/apache/shenyu/register/client/beat/HeartbeatListener.java
+++
b/shenyu-register-center/shenyu-register-client-beat/src/main/java/org/apache/shenyu/register/client/beat/HeartbeatListener.java
@@ -112,9 +112,7 @@ public class HeartbeatListener {
}
private void sendHeartbeat(final InstanceBeatInfoDTO instanceBeatInfoDTO) {
- int i = 0;
for (String server : serverList) {
- i++;
String concat = server.concat(Constants.BEAT_URI_PATH);
try {
String accessToken = this.accessToken.get(server);
@@ -124,9 +122,7 @@ public class HeartbeatListener {
RegisterUtils.doHeartBeat(GsonUtils.getInstance().toJson(instanceBeatInfoDTO),
concat, Constants.HEARTBEAT, accessToken);
} catch (Exception e) {
LOG.error("HeartBeat admin url :{} is fail, will retry.",
server, e);
- if (i == serverList.size()) {
- throw new RuntimeException(e);
- }
+ // This is a periodic reporter, not a failback registration:
retry on the next tick.
}
}
}
@@ -142,4 +138,4 @@ public class HeartbeatListener {
Thread.currentThread().interrupt();
}
}
-}
\ No newline at end of file
+}
diff --git
a/shenyu-register-center/shenyu-register-client-beat/src/test/java/org/apache/shenyu/register/client/beat/HeartbeatListenerTest.java
b/shenyu-register-center/shenyu-register-client-beat/src/test/java/org/apache/shenyu/register/client/beat/HeartbeatListenerTest.java
index 17929133f8..18b2f1b313 100644
---
a/shenyu-register-center/shenyu-register-client-beat/src/test/java/org/apache/shenyu/register/client/beat/HeartbeatListenerTest.java
+++
b/shenyu-register-center/shenyu-register-client-beat/src/test/java/org/apache/shenyu/register/client/beat/HeartbeatListenerTest.java
@@ -19,26 +19,28 @@ package org.apache.shenyu.register.client.beat;
import org.apache.shenyu.common.config.ShenyuConfig;
import org.apache.shenyu.common.constant.Constants;
+import org.apache.shenyu.common.utils.SystemInfoUtils;
import org.apache.shenyu.register.client.http.utils.RegisterUtils;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.MockedStatic;
+import org.mockito.MockedConstruction;
+import org.mockito.ArgumentCaptor;
import org.mockito.Mockito;
import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.boot.autoconfigure.web.ServerProperties;
import java.lang.reflect.Field;
-import java.lang.reflect.InvocationTargetException;
import java.util.Optional;
import java.util.Properties;
import java.util.concurrent.ScheduledThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import java.io.IOException;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
-import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertNotNull;
-import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.anyString;
@@ -89,6 +91,26 @@ class HeartbeatListenerTest {
return properties;
}
+ @Test
+ void testPeriodicHeartbeatSurvivesFailuresOnEveryServer() throws
IOException {
+ try (MockedConstruction<ScheduledThreadPoolExecutor> executors =
Mockito.mockConstruction(ScheduledThreadPoolExecutor.class);
+ MockedStatic<SystemInfoUtils> ignored =
Mockito.mockStatic(SystemInfoUtils.class);
+ MockedStatic<RegisterUtils> register =
Mockito.mockStatic(RegisterUtils.class)) {
+ register.when(() -> RegisterUtils.doLogin(anyString(),
anyString(), anyString())).thenReturn(Optional.of("mock-token"));
+ register.when(() -> RegisterUtils.doHeartBeat(anyString(),
anyString(), anyString(), anyString()))
+ .thenThrow(new IOException("admin unavailable"));
+ heartbeatListener = new HeartbeatListener(config, shenyuConfig,
serverProperties);
+ ArgumentCaptor<Runnable> periodic =
ArgumentCaptor.forClass(Runnable.class);
+
Mockito.verify(executors.constructed().get(0)).scheduleAtFixedRate(periodic.capture(),
Mockito.eq(0L), Mockito.eq(5L), Mockito.eq(TimeUnit.SECONDS));
+ assertDoesNotThrow(periodic.getValue()::run);
+ assertDoesNotThrow(periodic.getValue()::run);
+ register.verify(() -> RegisterUtils.doHeartBeat(anyString(),
Mockito.eq("http://localhost:9095" + Constants.BEAT_URI_PATH),
+ Mockito.eq(Constants.HEARTBEAT),
Mockito.eq("mock-token")), Mockito.times(2));
+ register.verify(() -> RegisterUtils.doHeartBeat(anyString(),
Mockito.eq("http://localhost:9096" + Constants.BEAT_URI_PATH),
+ Mockito.eq(Constants.HEARTBEAT),
Mockito.eq("mock-token")), Mockito.times(2));
+ }
+ }
+
@Test
void testHeartbeatListenerCreation() {
@@ -186,9 +208,10 @@ class HeartbeatListenerTest {
org.apache.shenyu.register.common.dto.InstanceBeatInfoDTO beatInfo
=
new
org.apache.shenyu.register.common.dto.InstanceBeatInfoDTO();
- InvocationTargetException exception =
assertThrows(InvocationTargetException.class,
- () -> sendHeartbeatMethod.invoke(heartbeatListener,
beatInfo));
- assertInstanceOf(RuntimeException.class, exception.getCause());
+ assertDoesNotThrow(() ->
sendHeartbeatMethod.invoke(heartbeatListener, beatInfo));
+ assertDoesNotThrow(() ->
sendHeartbeatMethod.invoke(heartbeatListener, beatInfo));
+ registerUtilsMockedStatic.verify(() ->
RegisterUtils.doHeartBeat(anyString(), anyString(), anyString(), anyString()),
+ Mockito.never());
}
}
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 00bd4580de..659ac9a2f6 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
@@ -205,9 +205,8 @@ public class HttpClientRegisterRepository extends
FailbackRegistryRepository {
}
private <T> void doRegister(final T t, final String path, final String
type) {
- int i = 0;
+ RuntimeException failure = null;
for (String server : serverList) {
- i++;
String concat = server.concat(path);
try {
String accessToken = this.accessToken.get(server);
@@ -218,17 +217,22 @@ public class HttpClientRegisterRepository extends
FailbackRegistryRepository {
// considering the situation of multiple clusters, we should
continue to execute here
} catch (Exception e) {
LOGGER.error("Register admin url :{} is fail, will retry.
cause:{}", server, e.getMessage());
- if (i == serverList.size()) {
- throw new RuntimeException(e);
+ if (Objects.isNull(failure)) {
+ failure = new RuntimeException(e);
+ } else {
+ failure.addSuppressed(e);
}
}
}
+ if (Objects.nonNull(failure)) {
+ // Failback replays the registration to every server, so admin
registration endpoints must be idempotent.
+ throw failure;
+ }
}
private <T> void doHeartbeat(final T t, final String path) {
- int i = 0;
+ RuntimeException failure = null;
for (String server : serverList) {
- i++;
String concat = server.concat(path);
try {
String accessToken = this.accessToken.get(server);
@@ -238,11 +242,16 @@ public class HttpClientRegisterRepository extends
FailbackRegistryRepository {
RegisterUtils.doHeartBeat(GsonUtils.getInstance().toJson(t),
concat, Constants.HEARTBEAT, accessToken);
} catch (Exception e) {
LOGGER.error("HeartBeat admin url :{} is fail, will retry.",
server, e);
- if (i == serverList.size()) {
- throw new RuntimeException(e);
+ if (Objects.isNull(failure)) {
+ failure = new RuntimeException(e);
+ } else {
+ failure.addSuppressed(e);
}
}
}
+ if (Objects.nonNull(failure)) {
+ throw failure;
+ }
}
private <T> void doUnregister(final T t) {
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 2c96875cf2..ce7d8e6fe4 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
@@ -24,12 +24,15 @@ import
org.apache.shenyu.register.client.http.utils.RuntimeUtils;
import org.apache.shenyu.register.common.config.ShenyuRegisterCenterConfig;
import org.apache.shenyu.register.common.dto.ApiDocRegisterDTO;
import org.apache.shenyu.register.common.dto.DiscoveryConfigRegisterDTO;
+import org.apache.shenyu.register.common.dto.InstanceBeatInfoDTO;
import org.apache.shenyu.register.common.dto.McpToolsRegisterDTO;
import org.apache.shenyu.register.common.dto.MetaDataRegisterDTO;
import org.apache.shenyu.register.common.dto.URIRegisterDTO;
import org.apache.shenyu.register.common.enums.EventType;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
import org.mockito.MockedStatic;
import java.io.IOException;
@@ -86,6 +89,54 @@ public final class HttpClientRegisterRepositoryTest {
}
}
+ @ParameterizedTest
+ @ValueSource(strings = {FIRST_SERVER, SECOND_SERVER})
+ public void partialRegistrationFailureMustPropagate(final String
failedServer) throws IOException {
+ HttpClientRegisterRepository multiServerRepository = new
HttpClientRegisterRepository(config(FIRST_SERVER + "," + SECOND_SERVER));
+ try (MockedStatic<RegisterUtils> registerUtils =
mockStatic(RegisterUtils.class);
+ MockedStatic<RuntimeUtils> runtimeUtils =
mockStatic(RuntimeUtils.class)) {
+ runtimeUtils.when(() ->
RuntimeUtils.listenByOther(anyInt())).thenReturn(false);
+ registerUtils.when(() -> RegisterUtils.doLogin(anyString(),
anyString(), anyString())).thenReturn(Optional.of(TOKEN));
+ registerUtils.when(() -> RegisterUtils.doRegister(anyString(),
eq(failedServer + Constants.URI_PATH), anyString(), anyString()))
+ .thenThrow(new IOException("unavailable"));
+ assertThrows(RuntimeException.class, () ->
multiServerRepository.doPersistURI(uriRegisterDTO()));
+ for (String server : new String[]{FIRST_SERVER, SECOND_SERVER}) {
+ registerUtils.verify(() ->
RegisterUtils.doRegister(anyString(), eq(server + Constants.URI_PATH),
eq(Constants.URI), eq(TOKEN)));
+ }
+ }
+ }
+
+ @ParameterizedTest
+ @ValueSource(strings = {FIRST_SERVER, SECOND_SERVER})
+ public void partialHeartbeatFailureMustPropagate(final String
failedServer) throws IOException {
+ HttpClientRegisterRepository multiServerRepository = new
HttpClientRegisterRepository(config(FIRST_SERVER + "," + SECOND_SERVER));
+ try (MockedStatic<RegisterUtils> registerUtils =
mockStatic(RegisterUtils.class);
+ MockedStatic<RuntimeUtils> runtimeUtils =
mockStatic(RuntimeUtils.class)) {
+ runtimeUtils.when(() ->
RuntimeUtils.listenByOther(anyInt())).thenReturn(false);
+ registerUtils.when(() -> RegisterUtils.doLogin(anyString(),
anyString(), anyString())).thenReturn(Optional.of(TOKEN));
+ registerUtils.when(() -> RegisterUtils.doHeartBeat(anyString(),
eq(failedServer + Constants.URI_PATH), anyString(), anyString()))
+ .thenThrow(new IOException("unavailable"));
+ assertThrows(RuntimeException.class, () ->
multiServerRepository.sendHeartbeat(uriRegisterDTO()));
+ for (String server : new String[]{FIRST_SERVER, SECOND_SERVER}) {
+ registerUtils.verify(() ->
RegisterUtils.doHeartBeat(anyString(), eq(server + Constants.URI_PATH),
eq(Constants.HEARTBEAT), eq(TOKEN)));
+ }
+ }
+ }
+
+ @Test
+ public void retainsAllRegistrationFailures() throws IOException {
+ HttpClientRegisterRepository multiServerRepository = new
HttpClientRegisterRepository(config(FIRST_SERVER + "," + SECOND_SERVER));
+ try (MockedStatic<RegisterUtils> registerUtils =
mockStatic(RegisterUtils.class);
+ MockedStatic<RuntimeUtils> runtimeUtils =
mockStatic(RuntimeUtils.class)) {
+ runtimeUtils.when(() ->
RuntimeUtils.listenByOther(anyInt())).thenReturn(false);
+ registerUtils.when(() -> RegisterUtils.doLogin(anyString(),
anyString(), anyString())).thenReturn(Optional.of(TOKEN));
+ registerUtils.when(() -> RegisterUtils.doRegister(anyString(),
anyString(), anyString(), anyString())).thenThrow(new
IOException("unavailable"));
+ RuntimeException failure = assertThrows(RuntimeException.class, ()
-> multiServerRepository.doPersistURI(uriRegisterDTO()));
+ assertEquals(1, failure.getSuppressed().length);
+ assertTrue(failure.getCause() instanceof IOException);
+ }
+ }
+
@Test
public void
persistUriShouldSkipRegistrationWhenPortIsUsedByAnotherProcess() {
URIRegisterDTO uriRegisterDTO = uriRegisterDTO();
@@ -291,6 +342,39 @@ public final class HttpClientRegisterRepositoryTest {
}
}
+ @ParameterizedTest
+ @ValueSource(strings = {FIRST_SERVER, SECOND_SERVER})
+ public void partialInstanceHeartbeatFailureMustPropagate(final String
failedServer) 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.doHeartBeat(anyString(),
eq(failedServer + Constants.BEAT_URI_PATH), anyString(), anyString()))
+ .thenThrow(new IOException("unavailable"));
+ assertThrows(RuntimeException.class, () ->
multiServerRepository.sendHeartbeat(new InstanceBeatInfoDTO()));
+ for (String server : new String[]{FIRST_SERVER, SECOND_SERVER}) {
+ registerUtils.verify(() ->
RegisterUtils.doHeartBeat(anyString(), eq(server + Constants.BEAT_URI_PATH),
eq(Constants.HEARTBEAT), eq(TOKEN)));
+ }
+ }
+ }
+
+ @Test
+ public void retainsAllHeartbeatFailures() throws IOException {
+ HttpClientRegisterRepository multiServerRepository = new
HttpClientRegisterRepository(config(FIRST_SERVER + "," + SECOND_SERVER));
+ try (MockedStatic<RegisterUtils> registerUtils =
mockStatic(RegisterUtils.class);
+ MockedStatic<RuntimeUtils> runtimeUtils =
mockStatic(RuntimeUtils.class)) {
+ runtimeUtils.when(() ->
RuntimeUtils.listenByOther(anyInt())).thenReturn(false);
+ registerUtils.when(() -> RegisterUtils.doLogin(anyString(),
anyString(), anyString())).thenReturn(Optional.of(TOKEN));
+ IOException first = new IOException("first unavailable");
+ IOException second = new IOException("second unavailable");
+ registerUtils.when(() -> RegisterUtils.doHeartBeat(anyString(),
eq(FIRST_SERVER + Constants.URI_PATH), anyString(),
anyString())).thenThrow(first);
+ registerUtils.when(() -> RegisterUtils.doHeartBeat(anyString(),
eq(SECOND_SERVER + Constants.URI_PATH), anyString(),
anyString())).thenThrow(second);
+ RuntimeException failure = assertThrows(RuntimeException.class, ()
-> multiServerRepository.sendHeartbeat(uriRegisterDTO()));
+ assertEquals(first, failure.getCause());
+ assertEquals(1, failure.getSuppressed().length);
+ assertEquals(second, failure.getSuppressed()[0]);
+ }
+ }
+
private ShenyuRegisterCenterConfig config(final String serverLists) {
Properties props = new Properties();
props.setProperty(Constants.USER_NAME, "admin");