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 1641987a10 fix: complete ShenyuMcpResponseDecorator future only after
all chunks are received (#6640) (#7040)
1641987a10 is described below
commit 1641987a106b9539220106a7ad667964b1672363
Author: wy471x <[email protected]>
AuthorDate: Wed Sep 30 11:17:59 2026 +0800
fix: complete ShenyuMcpResponseDecorator future only after all chunks are
received (#6640) (#7040)
Co-authored-by: aias00 <[email protected]>
---
.../response/ShenyuMcpResponseDecorator.java | 16 ++---
.../response/ShenyuMcpResponseDecoratorTest.java | 78 ++++++++++++++++++++++
2 files changed, 84 insertions(+), 10 deletions(-)
diff --git
a/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/response/ShenyuMcpResponseDecorator.java
b/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/response/ShenyuMcpResponseDecorator.java
index 410d552a6f..6b59d8b42d 100644
---
a/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/response/ShenyuMcpResponseDecorator.java
+++
b/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/response/ShenyuMcpResponseDecorator.java
@@ -68,15 +68,7 @@ public class ShenyuMcpResponseDecorator extends
ServerHttpResponseDecorator {
synchronized (this.body) {
this.body.append(chunk);
}
- // Complete future early for efficiency, but safely check if
already done
- if (!future.isDone()) {
- synchronized (future) {
- if (!future.isDone()) {
-
future.complete(applyResponseTemplate(this.body.toString()));
- }
- }
- }
- }));
+ }).doOnComplete(() -> completeFuture()));
}
@Override
@@ -88,6 +80,11 @@ public class ShenyuMcpResponseDecorator extends
ServerHttpResponseDecorator {
@Override
public Mono<Void> setComplete() {
LOG.debug("Response completed for session: {}", sessionId);
+ completeFuture();
+ return super.setComplete();
+ }
+
+ private void completeFuture() {
String responseBody;
synchronized (this.body) {
responseBody = this.body.toString();
@@ -100,7 +97,6 @@ public class ShenyuMcpResponseDecorator extends
ServerHttpResponseDecorator {
}
}
}
- return super.setComplete();
}
private String applyResponseTemplate(final String responseBody) {
diff --git
a/shenyu-plugin/shenyu-plugin-mcp-server/src/test/java/org/apache/shenyu/plugin/mcp/server/response/ShenyuMcpResponseDecoratorTest.java
b/shenyu-plugin/shenyu-plugin-mcp-server/src/test/java/org/apache/shenyu/plugin/mcp/server/response/ShenyuMcpResponseDecoratorTest.java
new file mode 100644
index 0000000000..94f4269f71
--- /dev/null
+++
b/shenyu-plugin/shenyu-plugin-mcp-server/src/test/java/org/apache/shenyu/plugin/mcp/server/response/ShenyuMcpResponseDecoratorTest.java
@@ -0,0 +1,78 @@
+/*
+ * 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.mcp.server.response;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.Mock;
+import org.mockito.junit.jupiter.MockitoExtension;
+import org.springframework.core.io.buffer.DataBuffer;
+import org.springframework.core.io.buffer.DefaultDataBufferFactory;
+import org.springframework.http.server.reactive.ServerHttpResponse;
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Mono;
+
+import java.nio.charset.StandardCharsets;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.TimeUnit;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.when;
+
+/**
+ * Test case for {@link ShenyuMcpResponseDecorator}.
+ */
+@ExtendWith(MockitoExtension.class)
+class ShenyuMcpResponseDecoratorTest {
+
+ @Mock
+ private ServerHttpResponse delegate;
+
+ private final DefaultDataBufferFactory bufferFactory = new
DefaultDataBufferFactory();
+
+ @Test
+ void testWriteWithCompletesFutureWithAllChunks() throws Exception {
+ when(delegate.writeWith(any())).thenAnswer(invocation ->
Flux.from(invocation.getArgument(0)).then());
+
+ final CompletableFuture<String> future = new CompletableFuture<>();
+ final ShenyuMcpResponseDecorator decorator =
+ new ShenyuMcpResponseDecorator(delegate, "session-1", future,
null);
+
+ decorator.writeWith(Flux.just(buffer("part-1,"),
buffer("part-2"))).block();
+
+ assertEquals("part-1,part-2", future.get(5, TimeUnit.SECONDS));
+ }
+
+ @Test
+ void testSetCompleteCompletesFutureWithAccumulatedBody() {
+ when(delegate.setComplete()).thenReturn(Mono.empty());
+
+ final CompletableFuture<String> future = new CompletableFuture<>();
+ final ShenyuMcpResponseDecorator decorator =
+ new ShenyuMcpResponseDecorator(delegate, "session-1", future,
null);
+
+ decorator.setComplete().block();
+
+ assertEquals("", future.getNow(""));
+ }
+
+ private DataBuffer buffer(final String content) {
+ return bufferFactory.wrap(content.getBytes(StandardCharsets.UTF_8));
+ }
+}