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