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 70def0197e fix(response): release WebClient bodies on chain failure
(#7274)
70def0197e is described below
commit 70def0197e9008f167e2c353f9f0b010786b109f
Author: Liming Deng <[email protected]>
AuthorDate: Thu Oct 1 06:27:31 2026 +0800
fix(response): release WebClient bodies on chain failure (#7274)
---
.../response/strategy/WebClientMessageWriter.java | 18 +++-
.../strategy/WebClientChainCleanupTest.java | 103 +++++++++++++++++++++
2 files changed, 120 insertions(+), 1 deletion(-)
diff --git
a/shenyu-plugin/shenyu-plugin-response/src/main/java/org/apache/shenyu/plugin/response/strategy/WebClientMessageWriter.java
b/shenyu-plugin/shenyu-plugin-response/src/main/java/org/apache/shenyu/plugin/response/strategy/WebClientMessageWriter.java
index 9f4c84b06c..5751eb8a6e 100644
---
a/shenyu-plugin/shenyu-plugin-response/src/main/java/org/apache/shenyu/plugin/response/strategy/WebClientMessageWriter.java
+++
b/shenyu-plugin/shenyu-plugin-response/src/main/java/org/apache/shenyu/plugin/response/strategy/WebClientMessageWriter.java
@@ -24,6 +24,8 @@ import org.apache.shenyu.plugin.api.ShenyuPluginChain;
import org.apache.shenyu.plugin.api.result.ShenyuResultEnum;
import org.apache.shenyu.plugin.api.result.ShenyuResultWrap;
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.DataBufferUtils;
import org.springframework.http.HttpHeaders;
@@ -48,6 +50,8 @@ import java.util.regex.Pattern;
*/
public class WebClientMessageWriter implements MessageWriter {
+ private static final Logger LOG =
LoggerFactory.getLogger(WebClientMessageWriter.class);
+
/**
* the common binary media type regex.
*/
@@ -67,7 +71,11 @@ public class WebClientMessageWriter implements MessageWriter
{
@Override
public Mono<Void> writeWith(final ServerWebExchange exchange, final
ShenyuPluginChain chain) {
- return chain.execute(exchange).then(Mono.defer(() -> {
+ // Invoke the chain on subscription, converting synchronous failures
into cleanable error signals.
+ Mono<Void> chainResult = Mono.defer(() -> chain.execute(exchange))
+ .doOnError(error -> clean(exchange))
+ .doOnCancel(() -> clean(exchange));
+ return chainResult.then(Mono.defer(() -> {
ServerHttpResponse response = exchange.getResponse();
ResponseEntity<Flux<DataBuffer>> fluxResponseEntity =
exchange.getAttribute(Constants.CLIENT_RESPONSE_ATTR);
@@ -121,6 +129,14 @@ public class WebClientMessageWriter implements
MessageWriter {
response.getHeaders().putAll(httpHeaders);
}
+ private void clean(final ServerWebExchange exchange) {
+ ResponseEntity<Flux<DataBuffer>> fluxResponseEntity =
exchange.getAttribute(Constants.CLIENT_RESPONSE_ATTR);
+ if (Objects.nonNull(fluxResponseEntity) &&
Objects.nonNull(fluxResponseEntity.getBody())) {
+ fluxResponseEntity.getBody().map(DataBufferUtils::release).then()
+ .subscribe(ignored -> { }, error -> LOG.debug("Unable to
drain upstream response body during cleanup", error));
+ }
+ }
+
static {
// https://www.iana.org/assignments/media-types/media-types.xhtml
//
https://developer.mozilla.org/en-US/docs/Web/HTTP/Basics_of_HTTP/MIME_types
diff --git
a/shenyu-plugin/shenyu-plugin-response/src/test/java/org/apache/shenyu/plugin/response/strategy/WebClientChainCleanupTest.java
b/shenyu-plugin/shenyu-plugin-response/src/test/java/org/apache/shenyu/plugin/response/strategy/WebClientChainCleanupTest.java
new file mode 100644
index 0000000000..bd9bc0e249
--- /dev/null
+++
b/shenyu-plugin/shenyu-plugin-response/src/test/java/org/apache/shenyu/plugin/response/strategy/WebClientChainCleanupTest.java
@@ -0,0 +1,103 @@
+/*
+ * 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.response.strategy;
+
+import io.netty.buffer.ByteBufAllocator;
+import org.apache.shenyu.common.constant.Constants;
+import org.apache.shenyu.plugin.api.ShenyuPluginChain;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
+import org.springframework.core.io.buffer.DataBuffer;
+import org.springframework.core.io.buffer.NettyDataBuffer;
+import org.springframework.core.io.buffer.NettyDataBufferFactory;
+import org.springframework.http.ResponseEntity;
+import org.springframework.mock.http.server.reactive.MockServerHttpRequest;
+import org.springframework.mock.web.server.MockServerWebExchange;
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Hooks;
+import reactor.core.publisher.Mono;
+import reactor.test.StepVerifier;
+
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+
+class WebClientChainCleanupTest {
+
+ @ParameterizedTest
+ @ValueSource(booleans = {false, true})
+ void releasesUnsubscribedBodyAfterReactiveOrSynchronousChainFailure(final
boolean synchronous) {
+ MockServerWebExchange exchange =
MockServerWebExchange.from(MockServerHttpRequest.get("/"));
+ NettyDataBuffer buffer = new
NettyDataBufferFactory(ByteBufAllocator.DEFAULT).allocateBuffer();
+ AtomicInteger subscriptions = new AtomicInteger();
+ Flux<DataBuffer> body =
Flux.<DataBuffer>just(buffer).doOnSubscribe(subscription ->
subscriptions.incrementAndGet());
+ exchange.getAttributes().put(Constants.CLIENT_RESPONSE_ATTR,
ResponseEntity.ok(body));
+ ShenyuPluginChain chain = ignored -> {
+ if (synchronous) {
+ throw new IllegalStateException("chain failed");
+ }
+ return Mono.error(new IllegalStateException("chain failed"));
+ };
+
+ StepVerifier.create(new WebClientMessageWriter().writeWith(exchange,
chain)).expectErrorMessage("chain failed").verify();
+
+ assertEquals(1, subscriptions.get());
+ assertEquals(0, buffer.getNativeBuffer().refCnt());
+ }
+
+ @Test
+ void releasesUnsubscribedBodyWhenChainIsCancelled() {
+ MockServerWebExchange exchange =
MockServerWebExchange.from(MockServerHttpRequest.get("/"));
+ NettyDataBuffer buffer = new
NettyDataBufferFactory(ByteBufAllocator.DEFAULT).allocateBuffer();
+ AtomicInteger subscriptions = new AtomicInteger();
+ Flux<DataBuffer> body =
Flux.<DataBuffer>just(buffer).doOnSubscribe(subscription ->
subscriptions.incrementAndGet());
+ exchange.getAttributes().put(Constants.CLIENT_RESPONSE_ATTR,
ResponseEntity.ok(body));
+
+ StepVerifier.create(new WebClientMessageWriter().writeWith(exchange,
ignored -> Mono.never())).thenCancel().verify();
+
+ assertEquals(1, subscriptions.get());
+ assertEquals(0, buffer.getNativeBuffer().refCnt());
+ }
+
+ @Test
+ void cleanupFailureDoesNotDropErrorsOrReplaceChainFailure() {
+ MockServerWebExchange exchange =
MockServerWebExchange.from(MockServerHttpRequest.get("/"));
+ exchange.getAttributes().put(Constants.CLIENT_RESPONSE_ATTR,
+ ResponseEntity.ok(Flux.<DataBuffer>error(new
IllegalStateException("body unavailable"))));
+ AtomicReference<Throwable> dropped = new AtomicReference<>();
+ Hooks.onErrorDropped(dropped::set);
+ try {
+ StepVerifier.create(new
WebClientMessageWriter().writeWith(exchange, ignored -> Mono.error(new
IllegalStateException("chain failed"))))
+ .expectErrorMessage("chain failed").verify();
+ assertNull(dropped.get());
+ } finally {
+ Hooks.resetOnErrorDropped();
+ }
+ }
+
+ @Test
+ void preservesFailureWhenNoUpstreamResponseExists() {
+ MockServerWebExchange exchange =
MockServerWebExchange.from(MockServerHttpRequest.get("/"));
+ StepVerifier.create(new WebClientMessageWriter().writeWith(exchange,
ignored -> Mono.error(new IllegalStateException("chain failed"))))
+ .expectErrorMessage("chain failed").verify();
+ }
+}