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