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 a2e26408a2 fix(httpclient): limit fixed retries to transient GET
failures (#7180)
a2e26408a2 is described below
commit a2e26408a219a226fc827c5539901845acac0681
Author: Liming Deng <[email protected]>
AuthorDate: Thu Oct 1 20:46:29 2026 +0800
fix(httpclient): limit fixed retries to transient GET failures (#7180)
---
.../plugin/httpclient/FixedRetryStrategy.java | 33 +++++++++--
.../plugin/httpclient/RetryStrategyTest.java | 64 +++++++++++++++++-----
.../plugin/httpclient/RetryTimeoutBudgetTest.java | 2 +-
3 files changed, 81 insertions(+), 18 deletions(-)
diff --git
a/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/FixedRetryStrategy.java
b/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/FixedRetryStrategy.java
index dad179211b..078ab2a862 100644
---
a/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/FixedRetryStrategy.java
+++
b/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/FixedRetryStrategy.java
@@ -17,13 +17,21 @@
package org.apache.shenyu.plugin.httpclient;
+import io.netty.channel.ConnectTimeoutException;
+import io.netty.handler.timeout.ReadTimeoutException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import org.springframework.http.HttpMethod;
import org.springframework.web.server.ServerWebExchange;
import reactor.core.publisher.Mono;
+import reactor.netty.http.client.PrematureCloseException;
import reactor.util.retry.Retry;
+import java.net.ConnectException;
+import java.net.SocketTimeoutException;
import java.time.Duration;
+import java.util.Objects;
+import java.util.concurrent.TimeoutException;
/**
* Fixed Retry Policy Class.
@@ -48,15 +56,32 @@ public class FixedRetryStrategy<R> implements
RetryStrategy<R> {
* @return Response Mono object after retry processing
*/
public Mono<R> execute(final Mono<R> response, final ServerWebExchange
exchange, final Duration duration, final int retryTimes) {
- Retry retrySpec = initFixedBackoff(retryTimes);
+ Retry retrySpec = initFixedBackoff(exchange, retryTimes);
Duration totalTimeout = RetryTimeoutUtils.totalTimeout(duration,
retryTimes, Duration.ofSeconds(2));
return response.retryWhen(retrySpec)
- .timeout(totalTimeout, Mono.error(() -> new
java.util.concurrent.TimeoutException("Retry sequence took longer than timeout:
" + totalTimeout)))
+ .timeout(totalTimeout, Mono.error(() -> new
TimeoutException("Retry sequence took longer than timeout: " + totalTimeout)))
.doOnError(e -> LOG.error(e.getMessage(), e));
}
- private Retry initFixedBackoff(final int retryTimes) {
+ private Retry initFixedBackoff(final ServerWebExchange exchange, final int
retryTimes) {
return Retry.fixedDelay(retryTimes, Duration.ofSeconds(2))
- .filter(t -> !(t instanceof
org.springframework.core.io.buffer.DataBufferLimitException));
+ .filter(throwable ->
HttpMethod.GET.equals(exchange.getRequest().getMethod()) &&
isTransientFailure(throwable))
+ .onRetryExhaustedThrow((retrySpec, retrySignal) ->
retrySignal.failure());
+ }
+
+ private boolean isTransientFailure(final Throwable throwable) {
+ Throwable cause = throwable;
+ while (Objects.nonNull(cause)) {
+ if (cause instanceof TimeoutException
+ || cause instanceof ConnectTimeoutException
+ || cause instanceof ReadTimeoutException
+ || cause instanceof ConnectException
+ || cause instanceof SocketTimeoutException
+ || cause instanceof PrematureCloseException) {
+ return true;
+ }
+ cause = cause.getCause();
+ }
+ return false;
}
}
diff --git
a/shenyu-plugin/shenyu-plugin-httpclient/src/test/java/org/apache/shenyu/plugin/httpclient/RetryStrategyTest.java
b/shenyu-plugin/shenyu-plugin-httpclient/src/test/java/org/apache/shenyu/plugin/httpclient/RetryStrategyTest.java
index 14daf1f322..e59214edda 100644
---
a/shenyu-plugin/shenyu-plugin-httpclient/src/test/java/org/apache/shenyu/plugin/httpclient/RetryStrategyTest.java
+++
b/shenyu-plugin/shenyu-plugin-httpclient/src/test/java/org/apache/shenyu/plugin/httpclient/RetryStrategyTest.java
@@ -17,13 +17,19 @@
package org.apache.shenyu.plugin.httpclient;
+import io.netty.channel.ConnectTimeoutException;
import org.junit.jupiter.api.Test;
+import org.springframework.mock.http.server.reactive.MockServerHttpRequest;
+import org.springframework.mock.web.server.MockServerWebExchange;
import org.springframework.web.server.ServerWebExchange;
import reactor.core.publisher.Mono;
import reactor.test.StepVerifier;
import java.time.Duration;
+import java.util.concurrent.TimeoutException;
+import java.util.concurrent.atomic.AtomicInteger;
+import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.mockito.Mockito.mock;
/**
@@ -80,25 +86,57 @@ public class RetryStrategyTest {
}
@Test
- void testFixedRetryStrategyExecute() {
- // Create a simulated AbstractHttpClientPlugin
+ void testFixedRetryStrategyRetriesTransientGetFailures() {
AbstractHttpClientPlugin<String> httpClientPlugin =
mock(AbstractHttpClientPlugin.class);
FixedRetryStrategy<String> strategy = new
FixedRetryStrategy<>(httpClientPlugin);
+ ServerWebExchange exchange =
MockServerWebExchange.from(MockServerHttpRequest.get("/").build());
+ AtomicInteger attempts = new AtomicInteger();
+ Mono<String> response = Mono.defer(() -> {
+ attempts.incrementAndGet();
+ return Mono.error(new ConnectTimeoutException("connection timed
out"));
+ });
+
+ StepVerifier.withVirtualTime(() -> strategy.execute(response,
exchange, Duration.ofSeconds(10), 2))
+ .thenAwait(Duration.ofSeconds(4))
+ .expectError(ConnectTimeoutException.class)
+ .verify();
- // Create a simulated ServerWebExchange
- ServerWebExchange exchange = mock(ServerWebExchange.class);
- Duration duration = Duration.ofSeconds(5);
- int retryTimes = 3;
+ assertEquals(3, attempts.get());
+ }
- // Create a mock response Mono that throws an exception
- Mono<String> response = Mono.error(new RuntimeException("Test error"));
+ @Test
+ void testFixedRetryStrategyDoesNotRetryPermanentFailures() {
+ AbstractHttpClientPlugin<String> httpClientPlugin =
mock(AbstractHttpClientPlugin.class);
+ FixedRetryStrategy<String> strategy = new
FixedRetryStrategy<>(httpClientPlugin);
+ ServerWebExchange exchange =
MockServerWebExchange.from(MockServerHttpRequest.get("/").build());
+ AtomicInteger attempts = new AtomicInteger();
+ Mono<String> response = Mono.defer(() -> {
+ attempts.incrementAndGet();
+ return Mono.error(new IllegalArgumentException("permanent
failure"));
+ });
+
+ StepVerifier.create(strategy.execute(response, exchange,
Duration.ofSeconds(5), 3))
+ .expectError(IllegalArgumentException.class)
+ .verify();
- // Execute retry policy
- Mono<String> result = strategy.execute(response, exchange, duration,
retryTimes);
+ assertEquals(1, attempts.get());
+ }
- // Use StepVerifier to verify results
- StepVerifier.create(result)
- .expectErrorMatches(reactor.core.Exceptions::isRetryExhausted)
+ @Test
+ void testFixedRetryStrategyDoesNotRetryPostRequests() {
+ AbstractHttpClientPlugin<String> httpClientPlugin =
mock(AbstractHttpClientPlugin.class);
+ FixedRetryStrategy<String> strategy = new
FixedRetryStrategy<>(httpClientPlugin);
+ ServerWebExchange exchange =
MockServerWebExchange.from(MockServerHttpRequest.post("/").build());
+ AtomicInteger attempts = new AtomicInteger();
+ Mono<String> response = Mono.defer(() -> {
+ attempts.incrementAndGet();
+ return Mono.error(new TimeoutException("request timed out"));
+ });
+
+ StepVerifier.create(strategy.execute(response, exchange,
Duration.ofSeconds(5), 3))
+ .expectError(TimeoutException.class)
.verify();
+
+ assertEquals(1, attempts.get());
}
}
diff --git
a/shenyu-plugin/shenyu-plugin-httpclient/src/test/java/org/apache/shenyu/plugin/httpclient/RetryTimeoutBudgetTest.java
b/shenyu-plugin/shenyu-plugin-httpclient/src/test/java/org/apache/shenyu/plugin/httpclient/RetryTimeoutBudgetTest.java
index 2ec12e1c1d..bb2f2db7ff 100644
---
a/shenyu-plugin/shenyu-plugin-httpclient/src/test/java/org/apache/shenyu/plugin/httpclient/RetryTimeoutBudgetTest.java
+++
b/shenyu-plugin/shenyu-plugin-httpclient/src/test/java/org/apache/shenyu/plugin/httpclient/RetryTimeoutBudgetTest.java
@@ -40,7 +40,7 @@ class RetryTimeoutBudgetTest {
void allowsEveryConfiguredRetryDespiteBackoffExceedingAttemptTimeout(final
String type) {
AtomicInteger attempts = new AtomicInteger();
StepVerifier.withVirtualTime(() ->
strategy(type).execute(Mono.defer(() -> attempts.incrementAndGet() < 4
- ? Mono.error(new IllegalStateException("retry")) :
Mono.just("success")),
+ ? Mono.error(new IllegalStateException("retry", new
TimeoutException("transient failure"))) : Mono.just("success")),
MockServerWebExchange.from(MockServerHttpRequest.get("/")),
Duration.ofMillis(100), 3))
.thenAwait(Duration.ofMinutes(2))
.expectNext("success")