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

Reply via email to