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 1755647886 fix(register): preserve replacement failures after retry
exhaustion (#7400)
1755647886 is described below
commit 1755647886cc17f6d482fe87bae8dac6c4389ac6
Author: JerryčAI <[email protected]>
AuthorDate: Fri Oct 2 21:27:46 2026 +0800
fix(register): preserve replacement failures after retry exhaustion (#7400)
---
.../client/api/FailbackRegistryRepository.java | 66 +++++++++++++++++++---
.../client/api/retry/FailureRegistryTask.java | 32 ++++++++++-
.../client/api/FailbackRegistryRepositoryTest.java | 39 +++++++++++++
.../client/api/retry/FailureRegistryTaskTest.java | 29 ++++++++--
4 files changed, 150 insertions(+), 16 deletions(-)
diff --git
a/shenyu-register-center/shenyu-register-client/shenyu-register-client-api/src/main/java/org/apache/shenyu/register/client/api/FailbackRegistryRepository.java
b/shenyu-register-center/shenyu-register-client/shenyu-register-client-api/src/main/java/org/apache/shenyu/register/client/api/FailbackRegistryRepository.java
index f718815687..27773e0574 100644
---
a/shenyu-register-center/shenyu-register-client/shenyu-register-client-api/src/main/java/org/apache/shenyu/register/client/api/FailbackRegistryRepository.java
+++
b/shenyu-register-center/shenyu-register-client/shenyu-register-client-api/src/main/java/org/apache/shenyu/register/client/api/FailbackRegistryRepository.java
@@ -184,15 +184,22 @@ public abstract class FailbackRegistryRepository
implements ShenyuClientRegister
}
private void addToFail(final Holder t) {
- // Update a pending payload atomically; a failure during an in-flight
retry owns a new entry and timer.
- Holder oldObj = concurrentHashMap.put(t.getKey(), t);
+ FailureRegistryTask registryTask =
FailureRegistryTask.createOwned(t.getKey(), this);
+ Holder newHolder = new Holder(t.getObj(), t.getPath(), t.getType(),
registryTask);
+ Holder oldObj = concurrentHashMap.put(t.getKey(), newHolder);
if (Objects.nonNull(oldObj)) {
+ if (Objects.nonNull(oldObj.getRetryTask())) {
+ oldObj.getRetryTask().cancel();
+ }
logger.debug("Updated failback registration payload, {}",
t.getPath());
- return;
}
- FailureRegistryTask registryTask = new FailureRegistryTask(t.getKey(),
this);
timer.add(registryTask);
- logger.warn("Add to failback and wait for execution, {}", t.getPath());
+ if (concurrentHashMap.get(t.getKey()) != newHolder) {
+ registryTask.cancel();
+ }
+ if (Objects.isNull(oldObj)) {
+ logger.warn("Add to failback and wait for execution, {}",
t.getPath());
+ }
}
/**
@@ -205,6 +212,19 @@ public abstract class FailbackRegistryRepository
implements ShenyuClientRegister
concurrentHashMap.remove(key);
}
+ /**
+ * Remove a pending registration only when the task still owns its holder.
+ *
+ * @param key the registration key
+ * @param retryTask the retry task
+ */
+ public void remove(final String key, final FailureRegistryTask retryTask) {
+ Holder holder = concurrentHashMap.get(key);
+ if (Objects.nonNull(holder) && holder.getRetryTask() == retryTask) {
+ concurrentHashMap.remove(key, holder);
+ }
+ }
+
/**
* Attempt a pending registration without claiming it, retained for
compatibility with custom retry tasks.
* New retry tasks should use {@link #retry(String)} instead of an
accept/remove pair.
@@ -225,8 +245,21 @@ public abstract class FailbackRegistryRepository
implements ShenyuClientRegister
* @param key the registration key
*/
public void retry(final String key) {
- Holder holder = concurrentHashMap.remove(key);
- if (Objects.isNull(holder)) {
+ Holder holder = concurrentHashMap.get(key);
+ if (Objects.nonNull(holder)) {
+ retry(key, holder.getRetryTask());
+ }
+ }
+
+ /**
+ * Retry a pending registration only when the task still owns its holder.
+ *
+ * @param key the registration key
+ * @param retryTask the retry task
+ */
+ public void retry(final String key, final FailureRegistryTask retryTask) {
+ Holder holder = concurrentHashMap.get(key);
+ if (Objects.isNull(holder) || holder.getRetryTask() != retryTask ||
!concurrentHashMap.remove(key, holder)) {
return;
}
try {
@@ -288,6 +321,8 @@ public abstract class FailbackRegistryRepository implements
ShenyuClientRegister
private final String type;
+ private final FailureRegistryTask retryTask;
+
/**
* Instantiates a new Holder.
*
@@ -296,9 +331,22 @@ public abstract class FailbackRegistryRepository
implements ShenyuClientRegister
* @param type the type
*/
Holder(final Object obj, final String path, final String type) {
+ this(obj, path, type, null);
+ }
+
+ /**
+ * Instantiates a new Holder.
+ *
+ * @param obj the registration object
+ * @param path the registration path
+ * @param type the registration type
+ * @param retryTask the retry task
+ */
+ Holder(final Object obj, final String path, final String type, final
FailureRegistryTask retryTask) {
this.obj = obj;
this.path = path;
this.type = type;
+ this.retryTask = retryTask;
}
/**
@@ -328,6 +376,10 @@ public abstract class FailbackRegistryRepository
implements ShenyuClientRegister
return type;
}
+ public FailureRegistryTask getRetryTask() {
+ return retryTask;
+ }
+
private String getKey() {
return String.join(":", path, type);
}
diff --git
a/shenyu-register-center/shenyu-register-client/shenyu-register-client-api/src/main/java/org/apache/shenyu/register/client/api/retry/FailureRegistryTask.java
b/shenyu-register-center/shenyu-register-client/shenyu-register-client-api/src/main/java/org/apache/shenyu/register/client/api/retry/FailureRegistryTask.java
index ba75fb7767..4556a18229 100644
---
a/shenyu-register-center/shenyu-register-client/shenyu-register-client-api/src/main/java/org/apache/shenyu/register/client/api/retry/FailureRegistryTask.java
+++
b/shenyu-register-center/shenyu-register-client/shenyu-register-client-api/src/main/java/org/apache/shenyu/register/client/api/retry/FailureRegistryTask.java
@@ -31,7 +31,9 @@ import java.util.concurrent.TimeUnit;
public class FailureRegistryTask extends AbstractRetryTask {
private final FailbackRegistryRepository registerRepository;
-
+
+ private final boolean ownerAware;
+
/**
* Instantiates a new Timer task.
*
@@ -39,9 +41,25 @@ public class FailureRegistryTask extends AbstractRetryTask {
* @param registerRepository the register repository
*/
public FailureRegistryTask(final String key, final
FailbackRegistryRepository registerRepository) {
+ this(key, registerRepository, false);
+ }
+
+ private FailureRegistryTask(final String key, final
FailbackRegistryRepository registerRepository, final boolean ownerAware) {
//Indicates 10s to retry.
super(key, TimeUnit.SECONDS.toMillis(10), 18);
this.registerRepository = registerRepository;
+ this.ownerAware = ownerAware;
+ }
+
+ /**
+ * Create a task that can only retry and remove the pending payload stored
for that task.
+ *
+ * @param key the registration key
+ * @param registerRepository the registry repository
+ * @return an owner-aware retry task
+ */
+ public static FailureRegistryTask createOwned(final String key, final
FailbackRegistryRepository registerRepository) {
+ return new FailureRegistryTask(key, registerRepository, true);
}
/**
@@ -52,11 +70,19 @@ public class FailureRegistryTask extends AbstractRetryTask {
*/
@Override
protected void doRetry(final String key, final TimerTask timerTask) {
- this.registerRepository.retry(key);
+ if (ownerAware) {
+ this.registerRepository.retry(key, this);
+ } else {
+ this.registerRepository.retry(key);
+ }
}
@Override
protected void onRetryExhausted(final String key) {
- this.registerRepository.remove(key);
+ if (ownerAware) {
+ this.registerRepository.remove(key, this);
+ } else {
+ this.registerRepository.remove(key);
+ }
}
}
diff --git
a/shenyu-register-center/shenyu-register-client/shenyu-register-client-api/src/test/java/org/apache/shenyu/register/client/api/FailbackRegistryRepositoryTest.java
b/shenyu-register-center/shenyu-register-client/shenyu-register-client-api/src/test/java/org/apache/shenyu/register/client/api/FailbackRegistryRepositoryTest.java
index 7805ebfd32..575426dd12 100644
---
a/shenyu-register-center/shenyu-register-client/shenyu-register-client-api/src/test/java/org/apache/shenyu/register/client/api/FailbackRegistryRepositoryTest.java
+++
b/shenyu-register-center/shenyu-register-client/shenyu-register-client-api/src/test/java/org/apache/shenyu/register/client/api/FailbackRegistryRepositoryTest.java
@@ -17,8 +17,11 @@
package org.apache.shenyu.register.client.api;
+import org.apache.shenyu.common.timer.TaskEntity;
import org.apache.shenyu.common.timer.Timer;
+import org.apache.shenyu.common.timer.TimerTask;
import org.apache.shenyu.common.timer.WheelTimerFactory;
+import org.apache.shenyu.register.client.api.retry.FailureRegistryTask;
import org.apache.shenyu.register.common.dto.ApiDocRegisterDTO;
import org.apache.shenyu.register.common.dto.McpToolsRegisterDTO;
import org.apache.shenyu.register.common.dto.MetaDataRegisterDTO;
@@ -28,6 +31,7 @@ import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource;
+import org.mockito.ArgumentCaptor;
import org.mockito.MockedStatic;
import java.lang.reflect.Field;
@@ -44,6 +48,7 @@ import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.same;
+import static org.mockito.Mockito.atLeastOnce;
import static org.mockito.Mockito.clearInvocations;
import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.doNothing;
@@ -53,6 +58,7 @@ import static org.mockito.Mockito.mockStatic;
import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
/**
* Test case for {@link FailbackRegistryRepository}.
@@ -131,6 +137,39 @@ public final class FailbackRegistryRepositoryTest {
verify(repository, times(3)).doPersistURI(original);
}
+ @Test
+ void retriesNewFailureQueuedAfterRetryBudgetIsExhausted() {
+ URIRegisterDTO original = createURIRegisterDTO();
+ URIRegisterDTO newer = createURIRegisterDTO();
+ doThrow(new
IllegalStateException("offline")).when(repository).doPersistURI(same(original));
+ doThrow(new IllegalStateException("new
failure")).when(repository).doPersistURI(same(newer));
+ repository.persistURI(original);
+
+ ArgumentCaptor<TimerTask> taskCaptor =
ArgumentCaptor.forClass(TimerTask.class);
+ verify(timer, atLeastOnce()).add(taskCaptor.capture());
+ FailureRegistryTask originalTask = (FailureRegistryTask)
taskCaptor.getAllValues().get(0);
+ TaskEntity originalTaskEntity = mock(TaskEntity.class);
+ when(originalTaskEntity.getTimer()).thenReturn(timer);
+ when(originalTaskEntity.getTimerTask()).thenReturn(originalTask);
+ for (int attempt = 0; attempt < 18; attempt++) {
+ originalTask.run(originalTaskEntity);
+ }
+
+ repository.persistURI(newer);
+ originalTask.run(originalTaskEntity);
+ assertEquals(1, getFailureMapSize());
+
+ verify(timer, atLeastOnce()).add(taskCaptor.capture());
+ FailureRegistryTask newerTask = (FailureRegistryTask)
taskCaptor.getAllValues()
+ .get(taskCaptor.getAllValues().size() - 1);
+ TaskEntity newerTaskEntity = mock(TaskEntity.class);
+ when(newerTaskEntity.getTimer()).thenReturn(timer);
+ when(newerTaskEntity.getTimerTask()).thenReturn(newerTask);
+ doNothing().when(repository).doPersistURI(same(newer));
+ newerTask.run(newerTaskEntity);
+ assertEquals(0, getFailureMapSize());
+ }
+
@Test
public void testPersistInterfaceSuccess() {
MetaDataRegisterDTO metadata = createMetaDataRegisterDTO();
diff --git
a/shenyu-register-center/shenyu-register-client/shenyu-register-client-api/src/test/java/org/apache/shenyu/register/client/api/retry/FailureRegistryTaskTest.java
b/shenyu-register-center/shenyu-register-client/shenyu-register-client-api/src/test/java/org/apache/shenyu/register/client/api/retry/FailureRegistryTaskTest.java
index a2cd52f9de..cfe8fa4c0d 100644
---
a/shenyu-register-center/shenyu-register-client/shenyu-register-client-api/src/test/java/org/apache/shenyu/register/client/api/retry/FailureRegistryTaskTest.java
+++
b/shenyu-register-center/shenyu-register-client/shenyu-register-client-api/src/test/java/org/apache/shenyu/register/client/api/retry/FailureRegistryTaskTest.java
@@ -37,7 +37,8 @@ public final class FailureRegistryTaskTest {
@Test
public void delegatesToAtomicRetry() {
FailbackRegistryRepository repository =
mock(FailbackRegistryRepository.class);
- new FailureRegistryTask("key", repository).doRetry("key",
mock(TimerTask.class));
+ FailureRegistryTask task = new FailureRegistryTask("key", repository);
+ task.doRetry("key", mock(TimerTask.class));
verify(repository).retry("key");
verifyNoMoreInteractions(repository);
}
@@ -45,15 +46,16 @@ public final class FailureRegistryTaskTest {
@Test
public void propagatesFailureForRescheduling() {
FailbackRegistryRepository repository =
mock(FailbackRegistryRepository.class);
- doThrow(new
IllegalStateException("offline")).when(repository).retry("key");
FailureRegistryTask task = new FailureRegistryTask("key", repository);
+ doThrow(new
IllegalStateException("offline")).when(repository).retry("key");
assertThrows(IllegalStateException.class, () -> task.doRetry("key",
mock(TimerTask.class)));
}
@Test
public void testRetryExhaustedRemovesFailure() {
FailbackRegistryRepository repository =
mock(FailbackRegistryRepository.class);
- new FailureRegistryTask("key", repository).onRetryExhausted("key");
+ FailureRegistryTask task = new FailureRegistryTask("key", repository);
+ task.onRetryExhausted("key");
verify(repository).remove("key");
}
@@ -72,8 +74,10 @@ public final class FailureRegistryTaskTest {
@Test
public void independentTasksUseTheirOwnRegistrationKeys() {
FailbackRegistryRepository repository =
mock(FailbackRegistryRepository.class);
- new FailureRegistryTask("first", repository).doRetry("first",
mock(TimerTask.class));
- new FailureRegistryTask("second", repository).doRetry("second",
mock(TimerTask.class));
+ FailureRegistryTask firstTask = new FailureRegistryTask("first",
repository);
+ FailureRegistryTask secondTask = new FailureRegistryTask("second",
repository);
+ firstTask.doRetry("first", mock(TimerTask.class));
+ secondTask.doRetry("second", mock(TimerTask.class));
verify(repository).retry("first");
verify(repository).retry("second");
verifyNoMoreInteractions(repository);
@@ -87,8 +91,8 @@ public final class FailureRegistryTaskTest {
TaskEntity entity = mock(TaskEntity.class);
when(entity.getTimer()).thenReturn(timer);
when(entity.getTimerTask()).thenReturn(timerTask);
- doThrow(new IllegalStateException("registration
failed")).when(repository).retry("key");
FailureRegistryTask task = new FailureRegistryTask("key", repository);
+ doThrow(new IllegalStateException("registration
failed")).when(repository).retry("key");
for (int attempt = 0; attempt < 19; attempt++) {
task.run(entity);
}
@@ -96,4 +100,17 @@ public final class FailureRegistryTaskTest {
verify(repository).remove("key");
verify(timer, times(18)).add(timerTask);
}
+
+ @Test
+ public void ownedTaskUsesConditionalRetryAndCleanup() {
+ FailbackRegistryRepository repository =
mock(FailbackRegistryRepository.class);
+ FailureRegistryTask task = FailureRegistryTask.createOwned("key",
repository);
+
+ task.doRetry("key", mock(TimerTask.class));
+ task.onRetryExhausted("key");
+
+ verify(repository).retry("key", task);
+ verify(repository).remove("key", task);
+ verifyNoMoreInteractions(repository);
+ }
}