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

Reply via email to