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 c7486ac5dd [type:fix] Bound URI readiness waits and respect 
interruption (#7257)
c7486ac5dd is described below

commit c7486ac5dd5eebbaf6c2a1515e77b7485552bc4e
Author: Liming Deng <[email protected]>
AuthorDate: Wed Sep 30 11:36:46 2026 +0800

    [type:fix] Bound URI readiness waits and respect interruption (#7257)
    
    * fix(client): bound URI readiness waits
    
    * fix(client): reduce URI readiness default and clarify startup tuning
---
 .../ShenyuClientURIExecutorSubscriber.java         |  67 +++++++----
 .../subcriber/UriReadinessTimeoutTest.java         | 132 +++++++++++++++++++++
 2 files changed, 179 insertions(+), 20 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 08e75b4569..228d929b9a 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
@@ -32,6 +32,7 @@ import org.slf4j.LoggerFactory;
 import org.springframework.beans.BeanUtils;
 
 import java.io.IOException;
+import java.net.InetSocketAddress;
 import java.net.Socket;
 import java.util.Collection;
 import java.util.List;
@@ -59,13 +60,26 @@ public class ShenyuClientURIExecutorSubscriber implements 
ExecutorTypeSubscriber
     private final ShenyuClientRegisterRepository 
shenyuClientRegisterRepository;
     
     private final ScheduledThreadPoolExecutor executor;
+
+    private final long readinessTimeoutMillis;
     
     /**
      * Instantiates a new Shenyu client uri executor subscriber.
+     * URI readiness is bounded by {@code 
shenyu.client.uri.readyTimeoutMillis}, defaulting to thirty seconds.
+     * The system property is read once when this subscriber is constructed.
+     * Unready URIs are logged and skipped so subsequent registration events 
can be processed.
      *
      * @param shenyuClientRegisterRepository the shenyu client register 
repository
      */
     public ShenyuClientURIExecutorSubscriber(final 
ShenyuClientRegisterRepository shenyuClientRegisterRepository) {
+        this(shenyuClientRegisterRepository, 
Long.getLong("shenyu.client.uri.readyTimeoutMillis", 
TimeUnit.SECONDS.toMillis(30)));
+    }
+
+    ShenyuClientURIExecutorSubscriber(final ShenyuClientRegisterRepository 
shenyuClientRegisterRepository, final long readinessTimeoutMillis) {
+        if (readinessTimeoutMillis <= 0) {
+            throw new IllegalArgumentException("URI readiness timeout must be 
positive");
+        }
+        this.readinessTimeoutMillis = readinessTimeoutMillis;
         this.shenyuClientRegisterRepository = shenyuClientRegisterRepository;
         // executor for send heartbeat
         ThreadFactory requestFactory = 
ShenyuThreadFactory.create("heartbeat-reporter", true);
@@ -82,27 +96,11 @@ public class ShenyuClientURIExecutorSubscriber implements 
ExecutorTypeSubscriber
     @Override
     public void executor(final Collection<URIRegisterDTO> dataList) {
         for (URIRegisterDTO uriRegisterDTO : dataList) {
-            Stopwatch stopwatch = Stopwatch.createStarted();
-            while (true) {
-                try (Socket ignored = new Socket(uriRegisterDTO.getHost(), 
uriRegisterDTO.getPort())) {
-                    break;
-                } catch (IOException e) {
-                    long sleepTime = 1000;
-                    // maybe the port is delay exposed
-                    if (stopwatch.elapsed(TimeUnit.SECONDS) > 5) {
-                        LOG.error("host:{}, port:{} connection failed, will 
retry",
-                                uriRegisterDTO.getHost(), 
uriRegisterDTO.getPort());
-                        // If the connection fails for a long time, Increase 
sleep time
-                        if (stopwatch.elapsed(TimeUnit.SECONDS) > 180) {
-                            sleepTime = 10000;
-                        }
-                    }
-                    try {
-                        TimeUnit.MILLISECONDS.sleep(sleepTime);
-                    } catch (InterruptedException ex) {
-                        LOG.error("interrupted when sleep", ex);
-                    }
+            if (!awaitReadiness(uriRegisterDTO)) {
+                if (Thread.currentThread().isInterrupted()) {
+                    return;
                 }
+                continue;
             }
             ShenyuClientShutdownHook.delayOtherHooks();
             shenyuClientRegisterRepository.persistURI(uriRegisterDTO);
@@ -122,6 +120,35 @@ public class ShenyuClientURIExecutorSubscriber implements 
ExecutorTypeSubscriber
             }), 2);
         }
     }
+
+    private boolean awaitReadiness(final URIRegisterDTO uri) {
+        Stopwatch stopwatch = Stopwatch.createStarted();
+        while (!Thread.currentThread().isInterrupted()) {
+            long remaining = readinessTimeoutMillis - 
stopwatch.elapsed(TimeUnit.MILLISECONDS);
+            if (remaining <= 0) {
+                LOG.warn("Skipping URI registration for {}:{} after waiting 
{}ms for readiness; configure shenyu.client.uri.readyTimeoutMillis before 
startup",
+                        uri.getHost(), uri.getPort(), readinessTimeoutMillis);
+                return false;
+            }
+            try (Socket socket = new Socket()) {
+                socket.connect(new InetSocketAddress(uri.getHost(), 
uri.getPort()), (int) Math.min(1000, remaining));
+                return true;
+            } catch (IOException e) {
+                LOG.debug("URI {}:{} is not ready", uri.getHost(), 
uri.getPort(), e);
+            }
+            remaining = readinessTimeoutMillis - 
stopwatch.elapsed(TimeUnit.MILLISECONDS);
+            if (remaining > 0) {
+                try {
+                    TimeUnit.MILLISECONDS.sleep(Math.min(1000, remaining));
+                } catch (InterruptedException e) {
+                    Thread.currentThread().interrupt();
+                    LOG.warn("Interrupted while waiting for URI {}:{} 
readiness", uri.getHost(), uri.getPort());
+                    return false;
+                }
+            }
+        }
+        return false;
+    }
     
     private void sendHeartbeat(final URIRegisterDTO uriRegisterDTO) {
         uriRegisterDTO.setInstanceInfo(SystemInfoUtils.getSystemInfo());
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
new file mode 100644
index 0000000000..bb95242e89
--- /dev/null
+++ 
b/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/UriReadinessTimeoutTest.java
@@ -0,0 +1,132 @@
+/*
+ * 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.client.core.shutdown.ShenyuClientShutdownHook;
+import org.apache.shenyu.register.client.api.ShenyuClientRegisterRepository;
+import org.apache.shenyu.register.common.dto.URIRegisterDTO;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Test;
+import org.springframework.test.util.ReflectionTestUtils;
+
+import java.net.ServerSocket;
+import java.time.Duration;
+import java.util.List;
+import java.util.Objects;
+import java.util.Properties;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ScheduledThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTimeoutPreemptively;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+
+/**
+ * Readiness waits must not block subsequent URI registration indefinitely.
+ */
+public final class UriReadinessTimeoutTest {
+
+    private final ShenyuClientRegisterRepository repository = 
mock(ShenyuClientRegisterRepository.class);
+
+    private ShenyuClientURIExecutorSubscriber subscriber;
+
+    @AfterEach
+    public void cleanup() {
+        if (Objects.nonNull(subscriber)) {
+            ((ScheduledThreadPoolExecutor) 
ReflectionTestUtils.getField(subscriber, "executor")).shutdownNow();
+        }
+    }
+
+    @Test
+    public void testReadinessDefaultAndStartupOverride() {
+        String property = "shenyu.client.uri.readyTimeoutMillis";
+        String previous = System.getProperty(property);
+        try {
+            System.clearProperty(property);
+            subscriber = new ShenyuClientURIExecutorSubscriber(repository);
+            assertEquals(30000L, ReflectionTestUtils.getField(subscriber, 
"readinessTimeoutMillis"));
+            System.setProperty(property, "5000");
+            assertEquals(30000L, ReflectionTestUtils.getField(subscriber, 
"readinessTimeoutMillis"));
+            cleanup();
+            subscriber = new ShenyuClientURIExecutorSubscriber(repository);
+            assertEquals(5000L, ReflectionTestUtils.getField(subscriber, 
"readinessTimeoutMillis"));
+        } finally {
+            if (Objects.isNull(previous)) {
+                System.clearProperty(property);
+            } else {
+                System.setProperty(property, previous);
+            }
+        }
+    }
+
+    @Test
+    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);
+            URIRegisterDTO available = uri(ready.getLocalPort());
+            assertTimeoutPreemptively(Duration.ofSeconds(3), () -> 
subscriber.executor(List.of(unavailable, available)));
+            verify(repository, never()).persistURI(unavailable);
+            verify(repository).persistURI(available);
+        }
+    }
+
+    @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);
+        }
+    }
+
+    private URIRegisterDTO uri(final int port) {
+        return 
URIRegisterDTO.builder().host("127.0.0.1").port(port).rpcType("http").contextPath("/readiness").build();
+    }
+}

Reply via email to