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);
+    }
 }

Reply via email to