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 3af9611570 fix(register): preserve failures queued during retries 
(#7259)
3af9611570 is described below

commit 3af9611570c79c0118b8325bc19bbb82a8ea32ee
Author: Liming Deng <[email protected]>
AuthorDate: Thu Oct 1 20:45:38 2026 +0800

    fix(register): preserve failures queued during retries (#7259)
    
    * fix(register): preserve failures queued during retries
    
    * test(register): restore retry delegation coverage and document 
compatibility
---
 RELEASE-NOTES.md                                   |   3 +
 .../client/api/FailbackRegistryRepository.java     |  31 +++-
 .../client/api/retry/FailureRegistryTask.java      |   4 +-
 .../client/api/FailbackRegistryRepositoryTest.java |  82 +++++++++
 .../client/api/retry/FailureRegistryTaskTest.java  | 194 +++++----------------
 5 files changed, 163 insertions(+), 151 deletions(-)

diff --git a/RELEASE-NOTES.md b/RELEASE-NOTES.md
index 620145c143..aede17dfee 100644
--- a/RELEASE-NOTES.md
+++ b/RELEASE-NOTES.md
@@ -6,6 +6,9 @@
 
 ### Behavior Changes
 
+- Custom registration retry tasks should call 
`FailbackRegistryRepository.retry(key)`.
+  The legacy `accept(key)` followed by `remove(key)` remains available for 
compatibility,
+  but can discard a newer registration failure arriving between those calls.
 - HTTP retry strategies budget the entire sequence separately from each 
attempt:
   `(retryTimes + 1) * attemptTimeout + retryTimes * maximumBackoff`.
   With a 3-second attempt timeout and 3 retries, the `current` strategy has a
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 7b4b76f41c..f718815687 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,6 +184,7 @@ 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);
         if (Objects.nonNull(oldObj)) {
             logger.debug("Updated failback registration payload, {}", 
t.getPath());
@@ -195,7 +196,8 @@ public abstract class FailbackRegistryRepository implements 
ShenyuClientRegister
     }
 
     /**
-     * Remove.
+     * Unconditionally remove a pending registration, retained for 
compatibility with custom retry tasks.
+     * Do not pair this with {@link #accept(String)}: use {@link 
#retry(String)} to preserve concurrent failures.
      *
      * @param key the key
      */
@@ -204,7 +206,8 @@ public abstract class FailbackRegistryRepository implements 
ShenyuClientRegister
     }
 
     /**
-     * Accpet.
+     * 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.
      *
      * @param key the key
      */
@@ -213,6 +216,30 @@ public abstract class FailbackRegistryRepository 
implements ShenyuClientRegister
         if (Objects.isNull(holder)) {
             return;
         }
+        persist(holder);
+    }
+
+    /**
+     * Retry a pending registration without removing failures queued during 
the attempt.
+     *
+     * @param key the registration key
+     */
+    public void retry(final String key) {
+        Holder holder = concurrentHashMap.remove(key);
+        if (Objects.isNull(holder)) {
+            return;
+        }
+        try {
+            persist(holder);
+        } catch (RuntimeException ex) {
+            // A newer failure has its own timer task; otherwise retain this 
task's retry.
+            if (Objects.isNull(concurrentHashMap.putIfAbsent(key, holder))) {
+                throw ex;
+            }
+        }
+    }
+
+    private void persist(final Holder holder) {
         String type = holder.getType();
         switch (type) {
             case Constants.URI:
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 7d173c377d..ba75fb7767 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
@@ -52,9 +52,7 @@ public class FailureRegistryTask extends AbstractRetryTask {
      */
     @Override
     protected void doRetry(final String key, final TimerTask timerTask) {
-        this.registerRepository.accept(key);
-        //Because accept requires an exception to be thrown. Only normal can 
remove.
-        this.registerRepository.remove(key);
+        this.registerRepository.retry(key);
     }
 
     @Override
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 f4128ec487..7805ebfd32 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,22 +17,39 @@
 
 package org.apache.shenyu.register.client.api;
 
+import org.apache.shenyu.common.timer.Timer;
+import org.apache.shenyu.common.timer.WheelTimerFactory;
 import org.apache.shenyu.register.common.dto.ApiDocRegisterDTO;
 import org.apache.shenyu.register.common.dto.McpToolsRegisterDTO;
 import org.apache.shenyu.register.common.dto.MetaDataRegisterDTO;
 import org.apache.shenyu.register.common.dto.URIRegisterDTO;
 import org.junit.jupiter.api.BeforeEach;
+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.MockedStatic;
 
 import java.lang.reflect.Field;
 import java.util.Map;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
+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.clearInvocations;
+import static org.mockito.Mockito.doAnswer;
 import static org.mockito.Mockito.doNothing;
 import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
 import static org.mockito.Mockito.spy;
 import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
@@ -44,11 +61,76 @@ public final class FailbackRegistryRepositoryTest {
 
     private TestFailbackRegistryRepository repository;
 
+    private MockedStatic<WheelTimerFactory> timerFactory;
+
+    private Timer timer;
+
     @BeforeEach
     public void setUp() {
+        timer = mock(Timer.class);
+        timerFactory = mockStatic(WheelTimerFactory.class);
+        timerFactory.when(WheelTimerFactory::getSharedTimer).thenReturn(timer);
         repository = spy(new TestFailbackRegistryRepository());
     }
 
+    @AfterEach
+    void closeTimerMock() {
+        timerFactory.close();
+    }
+
+    @ParameterizedTest
+    @ValueSource(booleans = {false, true})
+    void preservesNewFailureDuringRetry(final boolean retryFails) throws 
Exception {
+        URIRegisterDTO original = createURIRegisterDTO();
+        doThrow(new 
IllegalStateException("offline")).when(repository).doPersistURI(same(original));
+        repository.persistURI(original);
+        String key = getFirstKeyFromFailureMap();
+        CountDownLatch entered = new CountDownLatch(1);
+        CountDownLatch release = new CountDownLatch(1);
+        doAnswer(invocation -> {
+            entered.countDown();
+            assertTrue(release.await(5, TimeUnit.SECONDS));
+            if (retryFails) {
+                throw new IllegalStateException("still offline");
+            }
+            return null;
+        }).when(repository).doPersistURI(same(original));
+        ExecutorService executor = Executors.newSingleThreadExecutor();
+        try {
+            final Future<?> retry = executor.submit(() -> 
repository.retry(key));
+            assertTrue(entered.await(5, TimeUnit.SECONDS));
+            URIRegisterDTO newer = createURIRegisterDTO();
+            doThrow(new IllegalStateException("new 
failure")).when(repository).doPersistURI(same(newer));
+            repository.persistURI(newer);
+            release.countDown();
+            retry.get(5, TimeUnit.SECONDS);
+            assertEquals(1, getFailureMapSize());
+            verify(timer, times(2)).add(any());
+            doNothing().when(repository).doPersistURI(same(newer));
+            repository.retry(key);
+            verify(repository, times(2)).doPersistURI(same(newer));
+            assertEquals(0, getFailureMapSize());
+        } finally {
+            release.countDown();
+            executor.shutdownNow();
+        }
+    }
+
+    @Test
+    void restoresFailedRetryWhenNoNewFailureExists() {
+        URIRegisterDTO original = createURIRegisterDTO();
+        doThrow(new 
IllegalStateException("offline")).when(repository).doPersistURI(original);
+        repository.persistURI(original);
+        String key = getFirstKeyFromFailureMap();
+        assertThrows(IllegalStateException.class, () -> repository.retry(key));
+        assertEquals(1, getFailureMapSize());
+        doNothing().when(repository).doPersistURI(original);
+        repository.retry(key);
+        assertEquals(0, getFailureMapSize());
+        repository.retry(key);
+        verify(repository, times(3)).doPersistURI(original);
+    }
+
     @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 a75033a26d..a2cd52f9de 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
@@ -15,183 +15,85 @@
  * limitations under the License.
  */
 
+
 package org.apache.shenyu.register.client.api.retry;
 
+import org.apache.shenyu.common.timer.TimerTask;
 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.register.client.api.FailbackRegistryRepository;
-import org.apache.shenyu.register.common.dto.ApiDocRegisterDTO;
-import org.apache.shenyu.register.common.dto.McpToolsRegisterDTO;
-import org.apache.shenyu.register.common.dto.MetaDataRegisterDTO;
-import org.apache.shenyu.register.common.dto.URIRegisterDTO;
-import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
-import org.mockito.Mock;
-import org.mockito.MockitoAnnotations;
 
-import static org.junit.jupiter.api.Assertions.assertTrue;
-import static org.mockito.ArgumentMatchers.anyString;
-import static org.mockito.Mockito.doNothing;
+import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoMoreInteractions;
 import static org.mockito.Mockito.when;
 
-/**
- * Test case for {@link FailureRegistryTask}.
- */
 public final class FailureRegistryTaskTest {
 
-    private static final String TEST_KEY = "test-key";
-
-    @Mock
-    private FailbackRegistryRepository mockRepository;
-
-    @Mock
-    private TimerTask mockTimerTask;
-
-    @Mock
-    private TaskEntity mockTaskEntity;
-
-    @Mock
-    private Timer mockTimer;
-
-    private FailureRegistryTask failureRegistryTask;
-
-    @BeforeEach
-    public void setUp() {
-        MockitoAnnotations.openMocks(this);
-        failureRegistryTask = new FailureRegistryTask(TEST_KEY, 
mockRepository);
-    }
-
-
     @Test
-    public void testDoRetry() {
-        doNothing().when(mockRepository).accept(anyString());
-        doNothing().when(mockRepository).remove(anyString());
-        
-        failureRegistryTask.doRetry(TEST_KEY, mockTimerTask);
-        
-        verify(mockRepository, times(1)).accept(TEST_KEY);
-        verify(mockRepository, times(1)).remove(TEST_KEY);
+    public void delegatesToAtomicRetry() {
+        FailbackRegistryRepository repository = 
mock(FailbackRegistryRepository.class);
+        new FailureRegistryTask("key", repository).doRetry("key", 
mock(TimerTask.class));
+        verify(repository).retry("key");
+        verifyNoMoreInteractions(repository);
     }
 
     @Test
-    public void testRetryExhaustedRemovesFailure() {
-        failureRegistryTask.onRetryExhausted(TEST_KEY);
-
-        verify(mockRepository, times(1)).remove(TEST_KEY);
-    }
-
-    @Test
-    public void testDoRetryWithException() {
-
-        doNothing().when(mockRepository).accept(anyString());
-        doNothing().when(mockRepository).remove(anyString());
-        
-        // This should not throw an exception
-        failureRegistryTask.doRetry(TEST_KEY, mockTimerTask);
-        
-        verify(mockRepository, times(1)).accept(TEST_KEY);
-        verify(mockRepository, times(1)).remove(TEST_KEY);
+    public void propagatesFailureForRescheduling() {
+        FailbackRegistryRepository repository = 
mock(FailbackRegistryRepository.class);
+        doThrow(new 
IllegalStateException("offline")).when(repository).retry("key");
+        FailureRegistryTask task = new FailureRegistryTask("key", repository);
+        assertThrows(IllegalStateException.class, () -> task.doRetry("key", 
mock(TimerTask.class)));
     }
 
     @Test
-    public void testMultipleRetries() {
-
-        doNothing().when(mockRepository).accept(anyString());
-        doNothing().when(mockRepository).remove(anyString());
-        
-        // Test multiple retry calls
-        for (int i = 0; i < 3; i++) {
-            failureRegistryTask.doRetry(TEST_KEY, mockTimerTask);
-        }
-        
-        verify(mockRepository, times(3)).accept(TEST_KEY);
-        verify(mockRepository, times(3)).remove(TEST_KEY);
+    public void testRetryExhaustedRemovesFailure() {
+        FailbackRegistryRepository repository = 
mock(FailbackRegistryRepository.class);
+        new FailureRegistryTask("key", repository).onRetryExhausted("key");
+        verify(repository).remove("key");
     }
 
     @Test
-    public void testRemoveAfterRetriesExhausted() {
-        when(mockTaskEntity.getTimer()).thenReturn(mockTimer);
-        when(mockTaskEntity.getTimerTask()).thenReturn(mockTimerTask);
-        doThrow(new IllegalStateException("registration 
failed")).when(mockRepository).accept(TEST_KEY);
-
-        for (int i = 0; i < 19; i++) {
-            failureRegistryTask.run(mockTaskEntity);
+    public void repeatedAttemptsKeepDelegatingToTheSameKey() {
+        FailbackRegistryRepository repository = 
mock(FailbackRegistryRepository.class);
+        FailureRegistryTask task = new FailureRegistryTask("key", repository);
+        TimerTask timerTask = mock(TimerTask.class);
+        for (int attempt = 0; attempt < 3; attempt++) {
+            task.doRetry("key", timerTask);
         }
-
-        verify(mockRepository, times(18)).accept(TEST_KEY);
-        verify(mockRepository).remove(TEST_KEY);
-        verify(mockTimer, times(18)).add(mockTimerTask);
+        verify(repository, times(3)).retry("key");
+        verifyNoMoreInteractions(repository);
     }
 
     @Test
-    public void testDifferentKeys() {
-        final String key1 = "key1";
-        final String key2 = "key2";
-        
-        doNothing().when(mockRepository).accept(anyString());
-        doNothing().when(mockRepository).remove(anyString());
-        
-        failureRegistryTask.doRetry(key1, mockTimerTask);
-        failureRegistryTask.doRetry(key2, mockTimerTask);
-        
-        verify(mockRepository, times(1)).accept(key1);
-        verify(mockRepository, times(1)).remove(key1);
-        verify(mockRepository, times(1)).accept(key2);
-        verify(mockRepository, times(1)).remove(key2);
+    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));
+        verify(repository).retry("first");
+        verify(repository).retry("second");
+        verifyNoMoreInteractions(repository);
     }
 
     @Test
-    public void testTaskWithDifferentRepository() {
-        TestFailbackRegistryRepository testRepository = new 
TestFailbackRegistryRepository();
-        FailureRegistryTask task = new FailureRegistryTask("test", 
testRepository);
-        
-        task.doRetry("test", mockTimerTask);
-        
-        assertTrue(testRepository.acceptCalled);
-        assertTrue(testRepository.removeCalled);
-    }
-
-    /**
-     * Test implementation of FailbackRegistryRepository for testing.
-     */
-    private static class TestFailbackRegistryRepository extends 
FailbackRegistryRepository {
-        
-        private boolean acceptCalled;
-        
-        private boolean removeCalled;
-
-        @Override
-        public void accept(final String key) {
-            acceptCalled = true;
-        }
-
-        @Override
-        public void remove(final String key) {
-            removeCalled = true;
-        }
-
-        @Override
-        protected void doPersistApiDoc(final ApiDocRegisterDTO 
apiDocRegisterDTO) {
-            /* Test implementation */
-        }
-
-        @Override
-        protected void doPersistURI(final URIRegisterDTO registerDTO) {
-            /* Test implementation */
-        }
-
-        @Override
-        protected void doPersistInterface(final MetaDataRegisterDTO 
registerDTO) {
-            /* Test implementation */
-        }
-
-        @Override
-        protected void doPersistMcpTools(final McpToolsRegisterDTO 
registerDTO) {
-            /* Test implementation */
+    public void removesFailureAfterRetriesAreExhausted() {
+        final FailbackRegistryRepository repository = 
mock(FailbackRegistryRepository.class);
+        final TimerTask timerTask = mock(TimerTask.class);
+        final Timer timer = mock(Timer.class);
+        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);
+        for (int attempt = 0; attempt < 19; attempt++) {
+            task.run(entity);
         }
+        verify(repository, times(18)).retry("key");
+        verify(repository).remove("key");
+        verify(timer, times(18)).add(timerTask);
     }
 }

Reply via email to