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 09c6a52873 fix(client): make register publisher startup idempotent 
(#7082)
09c6a52873 is described below

commit 09c6a5287330c8d6ada64cd29f2a08a570b0b2ab
Author: Liming Deng <[email protected]>
AuthorDate: Thu Oct 1 20:47:03 2026 +0800

    fix(client): make register publisher startup idempotent (#7082)
    
    * fix(client): make register publisher startup idempotent
    
    * fix(client): clean up failed publisher startup
---
 .../ShenyuClientRegisterEventPublisher.java        | 47 ++++++++++++++++--
 .../ShenyuClientURIExecutorSubscriber.java         | 24 ++++++++--
 .../ShenyuClientRegisterEventPublisherTest.java    | 56 +++++++++++++++++++++-
 .../client/tars/TarsServiceBeanEventListener.java  |  1 -
 4 files changed, 117 insertions(+), 11 deletions(-)

diff --git 
a/shenyu-client/shenyu-client-core/src/main/java/org/apache/shenyu/client/core/disruptor/ShenyuClientRegisterEventPublisher.java
 
b/shenyu-client/shenyu-client-core/src/main/java/org/apache/shenyu/client/core/disruptor/ShenyuClientRegisterEventPublisher.java
index dc896149fa..2c0474ebb6 100644
--- 
a/shenyu-client/shenyu-client-core/src/main/java/org/apache/shenyu/client/core/disruptor/ShenyuClientRegisterEventPublisher.java
+++ 
b/shenyu-client/shenyu-client-core/src/main/java/org/apache/shenyu/client/core/disruptor/ShenyuClientRegisterEventPublisher.java
@@ -27,6 +27,8 @@ import org.apache.shenyu.disruptor.provider.DisruptorProvider;
 import org.apache.shenyu.register.client.api.ShenyuClientRegisterRepository;
 import org.apache.shenyu.register.common.type.DataTypeParent;
 
+import java.util.Objects;
+
 /**
  * The type shenyu client register event publisher.
  */
@@ -34,7 +36,7 @@ public class ShenyuClientRegisterEventPublisher {
 
     private static final ShenyuClientRegisterEventPublisher INSTANCE = new 
ShenyuClientRegisterEventPublisher();
 
-    private DisruptorProviderManage<DataTypeParent> providerManage;
+    private volatile DisruptorProviderManage<DataTypeParent> providerManage;
 
     /**
      * Get instance.
@@ -50,14 +52,49 @@ public class ShenyuClientRegisterEventPublisher {
      *
      * @param shenyuClientRegisterRepository shenyuClientRegisterRepository
      */
-    public void start(final ShenyuClientRegisterRepository 
shenyuClientRegisterRepository) {
+    public synchronized void start(final ShenyuClientRegisterRepository 
shenyuClientRegisterRepository) {
+        if (Objects.nonNull(providerManage)) {
+            return;
+        }
         RegisterClientExecutorFactory factory = new 
RegisterClientExecutorFactory();
         factory.addSubscribers(new 
ShenyuClientMetadataExecutorSubscriber(shenyuClientRegisterRepository));
-        factory.addSubscribers(new 
ShenyuClientURIExecutorSubscriber(shenyuClientRegisterRepository));
+        ShenyuClientURIExecutorSubscriber uriSubscriber = 
createUriSubscriber(shenyuClientRegisterRepository);
+        factory.addSubscribers(uriSubscriber);
         factory.addSubscribers(new 
ShenyuClientApiDocExecutorSubscriber(shenyuClientRegisterRepository));
         factory.addSubscribers(new 
ShenyuClientMcpExecutorSubscriber(shenyuClientRegisterRepository));
-        providerManage = new DisruptorProviderManage<>(factory);
-        providerManage.startup();
+        DisruptorProviderManage<DataTypeParent> manage = 
createProviderManage(factory);
+        try {
+            manage.startup();
+            uriSubscriber.start();
+            providerManage = manage;
+        } catch (RuntimeException ex) {
+            uriSubscriber.shutdown();
+            DisruptorProvider<DataTypeParent> provider = manage.getProvider();
+            if (Objects.nonNull(provider)) {
+                provider.shutdown();
+            }
+            throw ex;
+        }
+    }
+
+    /**
+     * Create URI subscriber.
+     *
+     * @param repository register repository
+     * @return URI subscriber
+     */
+    protected ShenyuClientURIExecutorSubscriber createUriSubscriber(final 
ShenyuClientRegisterRepository repository) {
+        return new ShenyuClientURIExecutorSubscriber(repository);
+    }
+
+    /**
+     * Create provider manager.
+     *
+     * @param factory consumer executor factory
+     * @return provider manager
+     */
+    protected DisruptorProviderManage<DataTypeParent> 
createProviderManage(final RegisterClientExecutorFactory<DataTypeParent> 
factory) {
+        return new DisruptorProviderManage<>(factory);
     }
 
     /**
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 743f5bcbe5..67ac30be01 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
@@ -41,6 +41,7 @@ import java.util.concurrent.CopyOnWriteArrayList;
 import java.util.concurrent.ScheduledThreadPoolExecutor;
 import java.util.concurrent.ThreadFactory;
 import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
 
 /**
  * The type Shenyu client uri executor subscriber.
@@ -61,6 +62,8 @@ public class ShenyuClientURIExecutorSubscriber implements 
ExecutorTypeSubscriber
     
     private final ScheduledThreadPoolExecutor executor;
 
+    private final AtomicBoolean started = new AtomicBoolean();
+
     private final long readinessTimeoutMillis;
     
     /**
@@ -85,7 +88,22 @@ public class ShenyuClientURIExecutorSubscriber implements 
ExecutorTypeSubscriber
         ThreadFactory requestFactory = 
ShenyuThreadFactory.create("heartbeat-reporter", true);
         executor = new ScheduledThreadPoolExecutor(1, requestFactory);
         
-        executor.scheduleAtFixedRate(() -> uris.forEach(this::sendHeartbeat), 
30, 10, TimeUnit.SECONDS);
+    }
+
+    /**
+     * Start reporting URI heartbeats.
+     */
+    public void start() {
+        if (started.compareAndSet(false, true)) {
+            executor.scheduleAtFixedRate(() -> 
uris.forEach(this::sendHeartbeat), 30, 10, TimeUnit.SECONDS);
+        }
+    }
+
+    /**
+     * Stop reporting URI heartbeats.
+     */
+    public void shutdown() {
+        executor.shutdown();
     }
     
     @Override
@@ -119,9 +137,7 @@ public class ShenyuClientURIExecutorSubscriber implements 
ExecutorTypeSubscriber
             shenyuClientRegisterRepository.offline(offlineDTO);
         } finally {
             // shutdown heartbeat executor
-            if (!executor.isTerminated()) {
-                executor.shutdown();
-            }
+            shutdown();
         }
     }
 
diff --git 
a/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientRegisterEventPublisherTest.java
 
b/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientRegisterEventPublisherTest.java
index a3ccf5ed64..200bf91f48 100644
--- 
a/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientRegisterEventPublisherTest.java
+++ 
b/shenyu-client/shenyu-client-core/src/test/java/org/apache/shenyu/client/core/disruptor/subcriber/ShenyuClientRegisterEventPublisherTest.java
@@ -18,6 +18,8 @@
 package org.apache.shenyu.client.core.disruptor.subcriber;
 
 import 
org.apache.shenyu.client.core.disruptor.ShenyuClientRegisterEventPublisher;
+import 
org.apache.shenyu.client.core.disruptor.executor.RegisterClientConsumerExecutor.RegisterClientExecutorFactory;
+import org.apache.shenyu.disruptor.DisruptorProviderManage;
 import org.apache.shenyu.register.client.api.ShenyuClientRegisterRepository;
 import org.apache.shenyu.register.common.type.DataTypeParent;
 import org.junit.jupiter.api.Assertions;
@@ -25,7 +27,11 @@ import org.junit.jupiter.api.Test;
 import org.mockito.Mock;
 
 import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
 
 public class ShenyuClientRegisterEventPublisherTest {
     @Mock
@@ -46,7 +52,9 @@ public class ShenyuClientRegisterEventPublisherTest {
         ShenyuClientRegisterEventPublisher publisher = 
ShenyuClientRegisterEventPublisher.getInstance();
         publisher.start(shenyuClientRegisterRepository);
         Assertions.assertNotNull(publisher.getProviderManage());
-        assertDoesNotThrow(() -> publisher.getProviderManage().startup());
+        Object providerManage = publisher.getProviderManage();
+        publisher.start(shenyuClientRegisterRepository);
+        assertSame(providerManage, publisher.getProviderManage());
     }
 
     @Test
@@ -62,4 +70,50 @@ public class ShenyuClientRegisterEventPublisherTest {
         publisher.start(shenyuClientRegisterRepository);
         assertDoesNotThrow(() -> publisher.publishEvent(null));
     }
+
+    @Test
+    public void testStartupFailureCleansResourcesAndAllowsRetry() {
+        DisruptorProviderManage<DataTypeParent> failedManage = 
mock(DisruptorProviderManage.class);
+        DisruptorProviderManage<DataTypeParent> successfulManage = 
mock(DisruptorProviderManage.class);
+        ShenyuClientURIExecutorSubscriber failedSubscriber = 
mock(ShenyuClientURIExecutorSubscriber.class);
+        ShenyuClientURIExecutorSubscriber successfulSubscriber = 
mock(ShenyuClientURIExecutorSubscriber.class);
+        doThrow(new IllegalStateException("startup 
failed")).when(failedManage).startup();
+        TestPublisher publisher = new TestPublisher(failedManage, 
successfulManage, failedSubscriber, successfulSubscriber);
+
+        assertThrows(IllegalStateException.class, () -> 
publisher.start(shenyuClientRegisterRepository));
+        verify(failedSubscriber).shutdown();
+
+        publisher.start(shenyuClientRegisterRepository);
+        assertSame(successfulManage, publisher.getProviderManage());
+        verify(successfulSubscriber).start();
+    }
+
+    private static final class TestPublisher extends 
ShenyuClientRegisterEventPublisher {
+
+        private final DisruptorProviderManage<DataTypeParent>[] manages;
+
+        private final ShenyuClientURIExecutorSubscriber[] subscribers;
+
+        private int manageIndex;
+
+        private int subscriberIndex;
+
+        private TestPublisher(final DisruptorProviderManage<DataTypeParent> 
failedManage,
+                              final DisruptorProviderManage<DataTypeParent> 
successfulManage,
+                              final ShenyuClientURIExecutorSubscriber 
failedSubscriber,
+                              final ShenyuClientURIExecutorSubscriber 
successfulSubscriber) {
+            manages = new DisruptorProviderManage[]{failedManage, 
successfulManage};
+            subscribers = new 
ShenyuClientURIExecutorSubscriber[]{failedSubscriber, successfulSubscriber};
+        }
+
+        @Override
+        protected ShenyuClientURIExecutorSubscriber createUriSubscriber(final 
ShenyuClientRegisterRepository repository) {
+            return subscribers[subscriberIndex++];
+        }
+
+        @Override
+        protected DisruptorProviderManage<DataTypeParent> 
createProviderManage(final RegisterClientExecutorFactory<DataTypeParent> 
factory) {
+            return manages[manageIndex++];
+        }
+    }
 }
diff --git 
a/shenyu-client/shenyu-client-tars/src/main/java/org/apache/shenyu/client/tars/TarsServiceBeanEventListener.java
 
b/shenyu-client/shenyu-client-tars/src/main/java/org/apache/shenyu/client/tars/TarsServiceBeanEventListener.java
index 917a7ade96..e730d5baf4 100644
--- 
a/shenyu-client/shenyu-client-tars/src/main/java/org/apache/shenyu/client/tars/TarsServiceBeanEventListener.java
+++ 
b/shenyu-client/shenyu-client-tars/src/main/java/org/apache/shenyu/client/tars/TarsServiceBeanEventListener.java
@@ -77,7 +77,6 @@ public class TarsServiceBeanEventListener extends 
AbstractContextRefreshedEventL
         }
         this.contextPath = contextPath;
         this.ipAndPort = this.getHost() + ":" + port;
-        publisher.start(shenyuClientRegisterRepository);
     }
 
     @Override

Reply via email to