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");

Reply via email to