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