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 154470da2d fix: apply modify response rules to streaming responses 
(#6990)
154470da2d is described below

commit 154470da2dc60133494ca9c73ffbfd1724287ad7
Author: liudawang001 <[email protected]>
AuthorDate: Thu Sep 3 23:23:12 2026 +0800

    fix: apply modify response rules to streaming responses (#6990)
    
    Co-authored-by: aias00 <[email protected]>
---
 .../modify/response/ModifyResponsePlugin.java      |  7 ++++++
 .../modify/response/ModifyResponsePluginTest.java  | 29 ++++++++++++++++++++++
 2 files changed, 36 insertions(+)

diff --git 
a/shenyu-plugin/shenyu-plugin-modify-response/src/main/java/org/apache/shenyu/plugin/modify/response/ModifyResponsePlugin.java
 
b/shenyu-plugin/shenyu-plugin-modify-response/src/main/java/org/apache/shenyu/plugin/modify/response/ModifyResponsePlugin.java
index b4698683ee..d4ff5a4003 100644
--- 
a/shenyu-plugin/shenyu-plugin-modify-response/src/main/java/org/apache/shenyu/plugin/modify/response/ModifyResponsePlugin.java
+++ 
b/shenyu-plugin/shenyu-plugin-modify-response/src/main/java/org/apache/shenyu/plugin/modify/response/ModifyResponsePlugin.java
@@ -106,6 +106,13 @@ public class ModifyResponsePlugin extends 
AbstractShenyuPlugin {
             });
         }
 
+        @Override
+        @NonNull
+        public Mono<Void> writeAndFlushWith(@NonNull final Publisher<? extends 
Publisher<? extends DataBuffer>> body) {
+            modifyResponseHeadersAndStatus();
+            return super.writeAndFlushWith(body);
+        }
+
         private void modifyResponseHeadersAndStatus() {
             HttpHeaders httpHeaders = new HttpHeaders();
             // add origin headers
diff --git 
a/shenyu-plugin/shenyu-plugin-modify-response/src/test/java/org/apache/shenyu/plugin/modify/response/ModifyResponsePluginTest.java
 
b/shenyu-plugin/shenyu-plugin-modify-response/src/test/java/org/apache/shenyu/plugin/modify/response/ModifyResponsePluginTest.java
index da50212f37..0f23982365 100644
--- 
a/shenyu-plugin/shenyu-plugin-modify-response/src/test/java/org/apache/shenyu/plugin/modify/response/ModifyResponsePluginTest.java
+++ 
b/shenyu-plugin/shenyu-plugin-modify-response/src/test/java/org/apache/shenyu/plugin/modify/response/ModifyResponsePluginTest.java
@@ -29,13 +29,19 @@ import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.extension.ExtendWith;
 import org.mockito.junit.jupiter.MockitoExtension;
+import org.reactivestreams.Publisher;
+import org.springframework.core.io.buffer.DataBuffer;
+import org.springframework.http.HttpStatus;
 import org.springframework.mock.http.server.reactive.MockServerHttpRequest;
+import org.springframework.mock.http.server.reactive.MockServerHttpResponse;
 import org.springframework.mock.web.server.MockServerWebExchange;
 import org.springframework.web.server.ServerWebExchange;
+import reactor.core.publisher.Flux;
 import reactor.core.publisher.Mono;
 import reactor.test.StepVerifier;
 
 import java.net.InetSocketAddress;
+import java.nio.charset.StandardCharsets;
 import java.util.HashMap;
 import java.util.Map;
 
@@ -100,6 +106,29 @@ public final class ModifyResponsePluginTest {
         StepVerifier.create(result).expectSubscription().verifyComplete();
     }
 
+    @Test
+    public void testWriteAndFlushWith() {
+        final ModifyResponseRuleHandle responseRuleHandle = new 
ModifyResponseRuleHandle();
+        final Map<String, String> addHeaders = new HashMap<>();
+        addHeaders.put("X-Test", "streaming");
+        responseRuleHandle.setAddHeaders(addHeaders);
+        responseRuleHandle.setStatusCode(HttpStatus.CREATED.value());
+        final ModifyResponsePlugin.ModifyResponseDecorator decorator =
+                new ModifyResponsePlugin.ModifyResponseDecorator(exchange, 
responseRuleHandle);
+        final MockServerHttpResponse response = (MockServerHttpResponse) 
exchange.getResponse();
+        final DataBuffer dataBuffer = response.bufferFactory()
+                .wrap("data: hello\n\n".getBytes(StandardCharsets.UTF_8));
+        final Publisher<? extends Publisher<? extends DataBuffer>> body = 
Flux.just(Flux.just(dataBuffer));
+
+        
StepVerifier.create(decorator.writeAndFlushWith(body)).verifyComplete();
+
+        assertEquals("streaming", response.getHeaders().getFirst("X-Test"));
+        assertEquals(HttpStatus.CREATED, response.getStatusCode());
+        StepVerifier.create(response.getBodyAsString())
+                .expectNext("data: hello\n\n")
+                .verifyComplete();
+    }
+
     @Test
     public void testGetOrder() {
         assertEquals(modifyResponsePlugin.getOrder(), 
PluginEnum.MODIFY_RESPONSE.getCode());

Reply via email to