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 bdd51f1c0c fix(httpclient): cache request body for replay during retry
(#6414)
bdd51f1c0c is described below
commit bdd51f1c0c00c0c723b7f944534401d94daed9af
Author: eye-gu <[email protected]>
AuthorDate: Sun Sep 20 19:55:07 2026 +0800
fix(httpclient): cache request body for replay during retry (#6414)
---
.../apache/shenyu/common/constant/Constants.java | 6 +
.../httpclient/AbstractHttpClientPlugin.java | 60 ++++-
.../plugin/httpclient/DefaultRetryStrategy.java | 5 +-
.../plugin/httpclient/FixedRetryStrategy.java | 3 +-
.../plugin/httpclient/NettyHttpClientPlugin.java | 13 +
.../shenyu/plugin/httpclient/WebClientPlugin.java | 18 +-
.../httpclient/NettyHttpClientPluginTest.java | 2 +-
.../httpclient/RequestBodyReplayRetryTest.java | 289 +++++++++++++++++++++
.../httpclient/HttpClientPluginConfiguration.java | 8 +-
9 files changed, 390 insertions(+), 14 deletions(-)
diff --git
a/shenyu-common/src/main/java/org/apache/shenyu/common/constant/Constants.java
b/shenyu-common/src/main/java/org/apache/shenyu/common/constant/Constants.java
index ccc6c0e36c..cb96014e49 100644
---
a/shenyu-common/src/main/java/org/apache/shenyu/common/constant/Constants.java
+++
b/shenyu-common/src/main/java/org/apache/shenyu/common/constant/Constants.java
@@ -147,6 +147,12 @@ public interface Constants {
*/
String ORIGINAL_RESPONSE_CONTENT_TYPE_ATTR =
"original_response_content_type";
+ /**
+ * The constant CACHED_REQUEST_BODY.
+ * Used to cache the request body for replay during retry.
+ */
+ String CACHED_REQUEST_BODY = "cachedRequestBody";
+
/**
* The constant HTTP_URI.
*/
diff --git
a/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/AbstractHttpClientPlugin.java
b/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/AbstractHttpClientPlugin.java
index ed67c79ceb..085d38766f 100644
---
a/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/AbstractHttpClientPlugin.java
+++
b/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/AbstractHttpClientPlugin.java
@@ -45,6 +45,8 @@ import org.apache.shenyu.plugin.api.utils.WebFluxResultUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.core.io.buffer.DataBuffer;
+import org.springframework.core.io.buffer.DataBufferLimitException;
+import org.springframework.core.io.buffer.DataBufferUtils;
import org.springframework.http.HttpHeaders;
import org.springframework.http.HttpStatus;
import org.springframework.web.server.ResponseStatusException;
@@ -60,6 +62,12 @@ public abstract class AbstractHttpClientPlugin<R> implements
ShenyuPlugin {
protected static final Logger LOG =
LoggerFactory.getLogger(AbstractHttpClientPlugin.class);
+ private final long maxInMemorySize;
+
+ protected AbstractHttpClientPlugin(final long maxInMemorySize) {
+ this.maxInMemorySize = maxInMemorySize;
+ }
+
@Override
public final Mono<Void> execute(final ServerWebExchange exchange, final
ShenyuPluginChain chain) {
final ShenyuContext shenyuContext =
exchange.getAttribute(Constants.CONTEXT);
@@ -73,11 +81,13 @@ public abstract class AbstractHttpClientPlugin<R>
implements ShenyuPlugin {
final Duration duration = Duration.ofMillis(timeout);
final int retryTimes = (int)
Optional.ofNullable(exchange.getAttribute(Constants.HTTP_RETRY)).orElse(0);
final String retryStrategy = (String)
Optional.ofNullable(exchange.getAttribute(Constants.RETRY_STRATEGY)).orElseGet(RetryEnum.CURRENT::getName);
+ final String httpMethod =
Objects.nonNull(exchange.getRequest().getMethod())
+ ? exchange.getRequest().getMethod().name() : "UNKNOWN";
LogUtils.debug(LOG, () -> String.format("The request urlPath is: %s,
retryTimes is : %s, retryStrategy is : %s", uri, retryTimes, retryStrategy));
- final Mono<R> response = doRequest(exchange,
- Objects.nonNull(exchange.getRequest().getMethod()) ?
exchange.getRequest().getMethod().name() : "UNKNOWN",
- uri,
- exchange.getRequest().getBody())
+ final Flux<DataBuffer> requestBody = (retryTimes > 0 &&
isRequestBodyRequired(httpMethod))
+ ? getCachedRequestBody(exchange)
+ : exchange.getRequest().getBody();
+ final Mono<R> response = doRequest(exchange, httpMethod, uri,
requestBody)
.timeout(duration, Mono.error(() -> new
TimeoutException("Response took longer than timeout: " + duration)))
.doOnError(e -> LOG.error(e.getMessage(), e));
RetryStrategy<R> strategy;
@@ -98,9 +108,51 @@ public abstract class AbstractHttpClientPlugin<R>
implements ShenyuPlugin {
.onErrorMap(ShenyuException.class, th -> new
ResponseStatusException(HttpStatus.SERVICE_UNAVAILABLE,
ShenyuResultEnum.CANNOT_FIND_HEALTHY_UPSTREAM_URL_AFTER_FAILOVER.getMsg(), th))
.onErrorMap(java.util.concurrent.TimeoutException.class, th ->
new ResponseStatusException(HttpStatus.GATEWAY_TIMEOUT, th.getMessage(), th))
+ .onErrorMap(DataBufferLimitException.class, th -> new
ResponseStatusException(HttpStatus.PAYLOAD_TOO_LARGE,
+ "Request body exceeds the maxInMemorySize limit and
cannot be cached for retry", th))
.flatMap((Function<Object, Mono<? extends Void>>) o ->
chain.execute(exchange));
}
+ /**
+ * Cache the request body as byte[] so it can be replayed during retry.
+ *
+ * <p>The original body Flux from the Netty channel is single-use; without
caching,
+ * retry attempts would send an empty body.
+ *
+ * <p>Idempotent: the cached Flux is stored in exchange attributes so that
both
+ * CURRENT (retryWhen) and FAILOVER (resend) paths share the same
replayable instance.
+ *
+ * @param exchange the server web exchange
+ * @return a replayable Flux of DataBuffer
+ */
+ protected Flux<DataBuffer> getCachedRequestBody(final ServerWebExchange
exchange) {
+ Flux<DataBuffer> cached =
exchange.getAttribute(Constants.CACHED_REQUEST_BODY);
+ if (Objects.nonNull(cached)) {
+ return cached;
+ }
+ final int joinLimit = (int) Math.min(maxInMemorySize,
Integer.MAX_VALUE);
+ cached = DataBufferUtils.join(exchange.getRequest().getBody(),
joinLimit)
+ .map(dataBuffer -> {
+ byte[] bytes = new byte[dataBuffer.readableByteCount()];
+ dataBuffer.read(bytes);
+ DataBufferUtils.release(dataBuffer);
+ return bytes;
+ })
+ .defaultIfEmpty(new byte[0])
+ .cache()
+ .flatMapMany(bytes -> Flux.defer(() -> {
+ if (bytes.length == 0) {
+ return Flux.empty();
+ }
+ // ServerHttpRequest has no bufferFactory() API in Spring;
the response's
+ // factory shares the same allocator, and
NettyHttpClientPlugin.doRequest
+ // requires NettyDataBuffer from it.
+ return
Flux.just(exchange.getResponse().bufferFactory().wrap(bytes));
+ }));
+ exchange.getAttributes().put(Constants.CACHED_REQUEST_BODY, cached);
+ return cached;
+ }
+
/**
* Process the Web request.
diff --git
a/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/DefaultRetryStrategy.java
b/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/DefaultRetryStrategy.java
index 6339ddbee6..0242df1057 100644
---
a/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/DefaultRetryStrategy.java
+++
b/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/DefaultRetryStrategy.java
@@ -69,8 +69,9 @@ public class DefaultRetryStrategy<R> implements
RetryStrategy<R> {
.maxBackoff(Duration.ofSeconds(20L))
.transientErrors(true)
.jitter(0.5d)
- .filter(t -> t instanceof
java.util.concurrent.TimeoutException || t instanceof
io.netty.channel.ConnectTimeoutException
+ .filter(t -> (t instanceof
java.util.concurrent.TimeoutException || t instanceof
io.netty.channel.ConnectTimeoutException
|| t instanceof
io.netty.handler.timeout.ReadTimeoutException || t instanceof
IllegalStateException)
+ && !(t instanceof
org.springframework.core.io.buffer.DataBufferLimitException))
.onRetryExhaustedThrow((retryBackoffSpecErr, retrySignal)
-> {
throw new ShenyuTimeoutException("Request timeout, the
maximum number of retry times has been exceeded");
});
@@ -128,7 +129,7 @@ public class DefaultRetryStrategy<R> implements
RetryStrategy<R> {
final URI newUri = RequestUrlUtils.buildRequestUri(exchange,
upstream.buildDomain());
// in order not to affect the next retry call, newUri needs to be
excluded
exclude.add(newUri);
- return httpClientPlugin.doRequest(exchange,
exchange.getRequest().getMethod().name(), newUri,
exchange.getRequest().getBody())
+ return httpClientPlugin.doRequest(exchange,
exchange.getRequest().getMethod().name(), newUri,
httpClientPlugin.getCachedRequestBody(exchange))
.timeout(duration, Mono.error(() -> new
TimeoutException("Response took longer than timeout: " + duration)))
.doOnError(e -> LOG.error(e.getMessage(), e));
});
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 116e65e8c2..2a71030648 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
@@ -55,6 +55,7 @@ public class FixedRetryStrategy<R> implements
RetryStrategy<R> {
}
private Retry initFixedBackoff(final int retryTimes) {
- return Retry.fixedDelay(retryTimes, Duration.ofSeconds(2));
+ return Retry.fixedDelay(retryTimes, Duration.ofSeconds(2))
+ .filter(t -> !(t instanceof
org.springframework.core.io.buffer.DataBufferLimitException));
}
}
diff --git
a/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/NettyHttpClientPlugin.java
b/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/NettyHttpClientPlugin.java
index ce48076f78..1038ee3d4d 100644
---
a/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/NettyHttpClientPlugin.java
+++
b/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/NettyHttpClientPlugin.java
@@ -54,8 +54,21 @@ public class NettyHttpClientPlugin extends
AbstractHttpClientPlugin<HttpClientRe
* Instantiates a new Netty http client plugin.
*
* @param httpClient the http client
+ * @deprecated use {@link #NettyHttpClientPlugin(HttpClient, long)} to
specify the replay cache cap
*/
+ @Deprecated
public NettyHttpClientPlugin(final HttpClient httpClient) {
+ this(httpClient, Constants.BYTES_PER_MB);
+ }
+
+ /**
+ * Instantiates a new Netty http client plugin.
+ *
+ * @param httpClient the http client
+ * @param maxInMemorySize max request body size in bytes that may be
cached for retry replay
+ */
+ public NettyHttpClientPlugin(final HttpClient httpClient, final long
maxInMemorySize) {
+ super(maxInMemorySize);
this.httpClient = httpClient;
}
diff --git
a/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/WebClientPlugin.java
b/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/WebClientPlugin.java
index 3e34067a43..1bfc1450b8 100644
---
a/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/WebClientPlugin.java
+++
b/shenyu-plugin/shenyu-plugin-httpclient/src/main/java/org/apache/shenyu/plugin/httpclient/WebClientPlugin.java
@@ -49,11 +49,23 @@ public class WebClientPlugin extends
AbstractHttpClientPlugin<ResponseEntity<Flu
* Instantiates a new Web client plugin.
*
* @param webClient the web client
- * @param maxInMemorySize the maximum number of bytes to buffer in memory
+ * @deprecated use {@link #WebClientPlugin(WebClient, long)} to specify
the replay cache cap
*/
- public WebClientPlugin(final WebClient webClient, final int
maxInMemorySize) {
+ @Deprecated
+ public WebClientPlugin(final WebClient webClient) {
+ this(webClient, Constants.BYTES_PER_MB);
+ }
+
+ /**
+ * Instantiates a new Web client plugin.
+ *
+ * @param webClient the web client
+ * @param maxInMemorySize max request body size in bytes that may be
cached for retry replay
+ */
+ public WebClientPlugin(final WebClient webClient, final long
maxInMemorySize) {
+ super(maxInMemorySize);
this.webClient = webClient;
- this.maxInMemorySize = maxInMemorySize;
+ this.maxInMemorySize = (int) Math.min(maxInMemorySize,
Integer.MAX_VALUE);
}
@Override
diff --git
a/shenyu-plugin/shenyu-plugin-httpclient/src/test/java/org/apache/shenyu/plugin/httpclient/NettyHttpClientPluginTest.java
b/shenyu-plugin/shenyu-plugin-httpclient/src/test/java/org/apache/shenyu/plugin/httpclient/NettyHttpClientPluginTest.java
index ddf5410913..a37d8cff8a 100644
---
a/shenyu-plugin/shenyu-plugin-httpclient/src/test/java/org/apache/shenyu/plugin/httpclient/NettyHttpClientPluginTest.java
+++
b/shenyu-plugin/shenyu-plugin-httpclient/src/test/java/org/apache/shenyu/plugin/httpclient/NettyHttpClientPluginTest.java
@@ -68,7 +68,7 @@ public final class NettyHttpClientPluginTest {
chain = mock(ShenyuPluginChain.class);
when(chain.execute(any())).thenReturn(Mono.empty());
HttpClient httpClient = HttpClient.create();
- nettyHttpClientPlugin = new NettyHttpClientPlugin(httpClient);
+ nettyHttpClientPlugin = new NettyHttpClientPlugin(httpClient,
Constants.BYTES_PER_MB);
}
/**
diff --git
a/shenyu-plugin/shenyu-plugin-httpclient/src/test/java/org/apache/shenyu/plugin/httpclient/RequestBodyReplayRetryTest.java
b/shenyu-plugin/shenyu-plugin-httpclient/src/test/java/org/apache/shenyu/plugin/httpclient/RequestBodyReplayRetryTest.java
new file mode 100644
index 0000000000..ebf27e61dc
--- /dev/null
+++
b/shenyu-plugin/shenyu-plugin-httpclient/src/test/java/org/apache/shenyu/plugin/httpclient/RequestBodyReplayRetryTest.java
@@ -0,0 +1,289 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.shenyu.plugin.httpclient;
+
+import java.net.URI;
+import java.nio.charset.StandardCharsets;
+import java.time.Duration;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.shenyu.common.constant.Constants;
+import org.apache.shenyu.common.enums.RetryEnum;
+import org.apache.shenyu.loadbalancer.cache.UpstreamCacheManager;
+import org.apache.shenyu.loadbalancer.entity.Upstream;
+import org.apache.shenyu.plugin.api.ShenyuPluginChain;
+import org.apache.shenyu.plugin.api.context.ShenyuContext;
+import org.apache.shenyu.plugin.api.result.ShenyuResult;
+import org.apache.shenyu.plugin.api.utils.SpringBeanUtils;
+import org.apache.shenyu.plugin.api.utils.RequestUrlUtils;
+import org.apache.shenyu.plugin.base.utils.LoadbalancerUtils;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.MockedStatic;
+import org.mockito.junit.jupiter.MockitoExtension;
+import org.mockito.junit.jupiter.MockitoSettings;
+import org.mockito.quality.Strictness;
+import org.springframework.context.ConfigurableApplicationContext;
+import org.springframework.core.io.buffer.DataBuffer;
+import org.springframework.core.io.buffer.DataBufferUtils;
+import org.springframework.http.HttpHeaders;
+import org.springframework.http.HttpStatus;
+import org.springframework.http.MediaType;
+import org.springframework.http.server.reactive.ServerHttpRequest;
+import org.springframework.http.server.reactive.ServerHttpRequestDecorator;
+import org.springframework.mock.http.server.reactive.MockServerHttpRequest;
+import org.springframework.mock.web.server.MockServerWebExchange;
+import org.springframework.web.server.ResponseStatusException;
+import org.springframework.web.server.ServerWebExchange;
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Mono;
+import reactor.test.StepVerifier;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+/**
+ * Tests that the request body is correctly replayed during retry.
+ *
+ * <p>The original body Flux from the Netty channel is single-use; without
caching,
+ * retry attempts would send an empty body.
+ */
+@ExtendWith(MockitoExtension.class)
+@MockitoSettings(strictness = Strictness.LENIENT)
+public class RequestBodyReplayRetryTest {
+
+ private ShenyuPluginChain chain;
+
+ @BeforeEach
+ public void setUp() {
+ ConfigurableApplicationContext context =
mock(ConfigurableApplicationContext.class);
+ SpringBeanUtils.getInstance().setApplicationContext(context);
+
when(context.getBean(ShenyuResult.class)).thenReturn(mock(ShenyuResult.class));
+ chain = mock(ShenyuPluginChain.class);
+ when(chain.execute(any())).thenReturn(Mono.empty());
+ }
+
+ @Test
+ void testBodyReplayedOnRetry() {
+ RecordingPlugin plugin = new RecordingPlugin(1);
+ ServerWebExchange exchange =
createExchangeWithSingleUseBody("{\"name\":\"hello\"}");
+
+ StepVerifier.create(plugin.execute(exchange, chain))
+ .expectComplete()
+ .verify(Duration.ofSeconds(10));
+
+ assertEquals(2, plugin.getCapturedBodies().size(), "Should have 2
attempts (1 fail + 1 success)");
+ assertEquals("{\"name\":\"hello\"}",
plugin.getCapturedBodies().get(0), "First attempt should receive full body");
+ assertEquals("{\"name\":\"hello\"}",
plugin.getCapturedBodies().get(1), "Retry attempt should receive replayed
body");
+ assertNotNull(exchange.getAttribute(Constants.CACHED_REQUEST_BODY),
+ "Body should be cached when retry is enabled");
+ }
+
+ @Test
+ void testBodyReplayedOnFailoverRetry() {
+ RecordingPlugin plugin = new RecordingPlugin(1);
+ ServerWebExchange exchange =
createExchangeWithSingleUseBody("{\"name\":\"hello\"}");
+ // switch to failover strategy: DefaultRetryStrategy.resend picks a
new upstream
+ exchange.getAttributes().put(Constants.RETRY_STRATEGY,
RetryEnum.FAILOVER.getName());
+ exchange.getAttributes().put(Constants.DIVIDE_SELECTOR_ID,
"selector-1");
+ exchange.getAttributes().put(Constants.LOAD_BALANCE, "roundRobin");
+
+ Upstream standby = Upstream.builder().url("localhost:8081").build();
+ UpstreamCacheManager cacheManager = mock(UpstreamCacheManager.class);
+
when(cacheManager.findUpstreamListBySelectorId(anyString())).thenReturn(Collections.singletonList(standby));
+
+ try (MockedStatic<UpstreamCacheManager> cacheMock =
org.mockito.Mockito.mockStatic(UpstreamCacheManager.class);
+ MockedStatic<LoadbalancerUtils> lbMock =
org.mockito.Mockito.mockStatic(LoadbalancerUtils.class);
+ MockedStatic<RequestUrlUtils> urlMock =
org.mockito.Mockito.mockStatic(RequestUrlUtils.class)) {
+
cacheMock.when(UpstreamCacheManager::getInstance).thenReturn(cacheManager);
+ lbMock.when(() -> LoadbalancerUtils.getForExchange(any(),
anyString(), any())).thenReturn(standby);
+ urlMock.when(() -> RequestUrlUtils.buildRequestUri(any(),
anyString())).thenReturn(URI.create("http://localhost:8081/test"));
+
+ StepVerifier.create(plugin.execute(exchange, chain))
+ .expectComplete()
+ .verify(Duration.ofSeconds(10));
+ }
+
+ assertEquals(2, plugin.getCapturedBodies().size(), "Should have 2
attempts (1 fail + 1 failover success)");
+ assertEquals("{\"name\":\"hello\"}",
plugin.getCapturedBodies().get(0), "First attempt should receive full body");
+ assertEquals("{\"name\":\"hello\"}",
plugin.getCapturedBodies().get(1), "Failover attempt should receive replayed
body");
+ assertNotNull(exchange.getAttribute(Constants.CACHED_REQUEST_BODY),
+ "Body should be cached so failover resend can replay it");
+ }
+
+ @Test
+ void testGetRequestRetriesWithoutBody() {
+ RecordingPlugin plugin = new RecordingPlugin(1);
+ ServerWebExchange exchange = createGetExchangeWithRetry();
+
+ StepVerifier.create(plugin.execute(exchange, chain))
+ .expectComplete()
+ .verify(Duration.ofSeconds(10));
+
+ assertEquals(2, plugin.getCapturedBodies().size(), "GET should retry
(2 attempts: 1 fail + 1 success)");
+ assertNull(exchange.getAttribute(Constants.CACHED_REQUEST_BODY),
+ "GET must not cache a body it never reads");
+ }
+
+ @Test
+ void testOversizeBodyThrowsDataBufferLimitException() {
+ // body (9 bytes) exceeds maxInMemorySize (4 bytes) during aggregation
+ RecordingPlugin plugin = new RecordingPlugin(0, 4);
+ ServerWebExchange exchange =
createExchangeWithSingleUseBody("test-body");
+
+ StepVerifier.create(plugin.execute(exchange, chain))
+ .expectErrorSatisfies(err -> {
+ // Body aggregation overflow must be mapped to 413, not a
raw/generic 500
+ assertTrue(err instanceof ResponseStatusException,
+ "Oversize body should surface as
ResponseStatusException");
+ ResponseStatusException rse = (ResponseStatusException)
err;
+ assertEquals(HttpStatus.PAYLOAD_TOO_LARGE,
rse.getStatusCode(),
+ "Oversize body should map to 413 Payload Too
Large");
+ assertTrue(rse.getCause() instanceof
org.springframework.core.io.buffer.DataBufferLimitException,
+ "Underlying cause should be
DataBufferLimitException");
+ })
+ .verify(Duration.ofSeconds(10));
+ assertEquals(0, plugin.getCapturedBodies().size(), "Oversize body must
never reach doRequest");
+ }
+
+ @Test
+ void testOversizeBodyNotRetriedUnderFixedStrategy() {
+ RecordingPlugin plugin = new RecordingPlugin(0, 4);
+ ServerWebExchange exchange =
createExchangeWithSingleUseBody("test-body");
+ exchange.getAttributes().put(Constants.HTTP_RETRY_BACK_OFF_SPEC,
"fixed");
+
+ StepVerifier.create(plugin.execute(exchange, chain))
+ .expectErrorSatisfies(err -> {
+ assertTrue(err instanceof ResponseStatusException);
+ assertEquals(HttpStatus.PAYLOAD_TOO_LARGE,
((ResponseStatusException) err).getStatusCode());
+ })
+ .verify(Duration.ofSeconds(10));
+
+ assertEquals(0, plugin.getCapturedBodies().size(),
+ "Fixed strategy must not retry oversize body (it can never be
cached for replay)");
+ }
+
+ private ServerWebExchange createExchangeWithSingleUseBody(final String
body) {
+ final AtomicBoolean consumed = new AtomicBoolean(false);
+ final MockServerHttpRequest mockRequest = MockServerHttpRequest
+ .post("/test")
+ .header(HttpHeaders.CONTENT_TYPE,
MediaType.APPLICATION_JSON_VALUE)
+ .body(body);
+ ServerWebExchange exchange = MockServerWebExchange.from(mockRequest);
+ ServerHttpRequest singleUseRequest = new
ServerHttpRequestDecorator(exchange.getRequest()) {
+ @Override
+ public Flux<DataBuffer> getBody() {
+ return Flux.defer(() -> {
+ if (consumed.compareAndSet(false, true)) {
+ return mockRequest.getBody();
+ }
+ return Flux.empty();
+ });
+ }
+ };
+ exchange = exchange.mutate().request(singleUseRequest).build();
+ exchange.getAttributes().put(Constants.CONTEXT,
mock(ShenyuContext.class));
+ exchange.getAttributes().put(Constants.HTTP_URI,
URI.create("http://localhost/test"));
+ exchange.getAttributes().put(Constants.HTTP_TIME_OUT, 30000L);
+ exchange.getAttributes().put(Constants.HTTP_RETRY, 3);
+ return exchange;
+ }
+
+ private ServerWebExchange createGetExchangeWithRetry() {
+ final MockServerHttpRequest mockRequest =
MockServerHttpRequest.get("/test").build();
+ ServerWebExchange exchange = MockServerWebExchange.from(mockRequest);
+ exchange.getAttributes().put(Constants.CONTEXT,
mock(ShenyuContext.class));
+ exchange.getAttributes().put(Constants.HTTP_URI,
URI.create("http://localhost/test"));
+ exchange.getAttributes().put(Constants.HTTP_TIME_OUT, 30000L);
+ exchange.getAttributes().put(Constants.HTTP_RETRY, 3);
+ return exchange;
+ }
+
+ /**
+ * A test plugin that records the body received on each doRequest call
+ * and fails the first N attempts with TimeoutException (retryable by
DefaultRetryStrategy).
+ */
+ static class RecordingPlugin extends AbstractHttpClientPlugin<String> {
+
+ private final List<String> capturedBodies =
Collections.synchronizedList(new ArrayList<>());
+
+ private final AtomicInteger attempts = new AtomicInteger();
+
+ private final int failFirstN;
+
+ RecordingPlugin(final int failFirstN) {
+ this(failFirstN, Constants.BYTES_PER_MB);
+ }
+
+ RecordingPlugin(final int failFirstN, final long maxInMemorySize) {
+ super(maxInMemorySize);
+ this.failFirstN = failFirstN;
+ }
+
+ List<String> getCapturedBodies() {
+ return capturedBodies;
+ }
+
+ @Override
+ protected Mono<String> doRequest(final ServerWebExchange exchange,
final String httpMethod,
+ final URI uri, final Flux<DataBuffer>
body) {
+ return DataBufferUtils.join(body)
+ .map(buffer -> {
+ byte[] bytes = new byte[buffer.readableByteCount()];
+ buffer.read(bytes);
+ DataBufferUtils.release(buffer);
+ return new String(bytes, StandardCharsets.UTF_8);
+ })
+ .defaultIfEmpty("")
+ .flatMap(bodyStr -> {
+ capturedBodies.add(bodyStr);
+ if (attempts.incrementAndGet() <= failFirstN) {
+ return Mono.error(new
java.util.concurrent.TimeoutException(
+ "Simulated timeout, attempt " +
attempts.get()));
+ }
+ return Mono.just("success");
+ });
+ }
+
+ @Override
+ public int getOrder() {
+ return 0;
+ }
+
+ @Override
+ public boolean skip(final ServerWebExchange exchange) {
+ return false;
+ }
+
+ @Override
+ public String named() {
+ return "RecordingPlugin";
+ }
+ }
+}
diff --git
a/shenyu-spring-boot-starter/shenyu-spring-boot-starter-plugin/shenyu-spring-boot-starter-plugin-httpclient/src/main/java/org/apache/shenyu/springboot/starter/plugin/httpclient/HttpClientPluginConfiguration.java
b/shenyu-spring-boot-starter/shenyu-spring-boot-starter-plugin/shenyu-spring-boot-starter-plugin-httpclient/src/main/java/org/apache/shenyu/springboot/starter/plugin/httpclient/HttpClientPluginConfiguration.java
index 6ebc42ec5c..03938e714f 100644
---
a/shenyu-spring-boot-starter/shenyu-spring-boot-starter-plugin/shenyu-spring-boot-starter-plugin-httpclient/src/main/java/org/apache/shenyu/springboot/starter/plugin/httpclient/HttpClientPluginConfiguration.java
+++
b/shenyu-spring-boot-starter/shenyu-spring-boot-starter-plugin/shenyu-spring-boot-starter-plugin-httpclient/src/main/java/org/apache/shenyu/springboot/starter/plugin/httpclient/HttpClientPluginConfiguration.java
@@ -112,7 +112,7 @@ public class HttpClientPluginConfiguration {
.build())
.clientConnector(new
ReactorClientHttpConnector(Objects.requireNonNull(httpClient.getIfAvailable())))
.build();
- return new WebClientPlugin(webClient, maxInMemorySize);
+ return new WebClientPlugin(webClient, (long)
properties.getMaxInMemorySize() * Constants.BYTES_PER_MB);
}
}
@@ -127,11 +127,13 @@ public class HttpClientPluginConfiguration {
* Netty http client plugin.
*
* @param httpClient the http client
+ * @param properties the http client properties
* @return the shenyu plugin
*/
@Bean
- public ShenyuPlugin nettyHttpClientPlugin(final
ObjectProvider<HttpClient> httpClient) {
- return new NettyHttpClientPlugin(httpClient.getIfAvailable());
+ public ShenyuPlugin nettyHttpClientPlugin(final
ObjectProvider<HttpClient> httpClient,
+ final HttpClientProperties
properties) {
+ return new NettyHttpClientPlugin(httpClient.getIfAvailable(),
(long) properties.getMaxInMemorySize() * Constants.BYTES_PER_MB);
}
}
}