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 1e5580216b issue-6888 (#6937)
1e5580216b is described below

commit 1e5580216b4d520fe8a486574627b72024586f8c
Author: SouthwestAsiaFloat <[email protected]>
AuthorDate: Wed Sep 30 09:56:06 2026 +0800

    issue-6888 (#6937)
    
    Co-authored-by: aias00 <[email protected]>
---
 .../mcp/server/callback/ShenyuToolCallback.java    | 19 +++++++++++++++--
 .../server/callback/ShenyuToolCallbackTest.java    | 24 ++++++++++++++++++++++
 2 files changed, 41 insertions(+), 2 deletions(-)

diff --git 
a/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/callback/ShenyuToolCallback.java
 
b/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/callback/ShenyuToolCallback.java
index dad414af28..721c46556d 100644
--- 
a/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/callback/ShenyuToolCallback.java
+++ 
b/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/callback/ShenyuToolCallback.java
@@ -49,6 +49,7 @@ import org.springframework.lang.NonNull;
 import org.springframework.util.Assert;
 import org.springframework.util.StringUtils;
 import org.springframework.web.server.ServerWebExchange;
+import reactor.core.Disposable;
 
 import java.net.URI;
 import java.net.URISyntaxException;
@@ -56,6 +57,7 @@ import java.util.Map;
 import java.util.Objects;
 import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
 import java.util.regex.Matcher;
 import java.util.regex.Pattern;
 
@@ -219,6 +221,15 @@ public class ShenyuToolCallback implements ToolCallback {
                                    final String sessionId,
                                    final String configStr,
                                    final String input) {
+        return executeToolCall(originExchange, chain, sessionId, configStr, 
input, DEFAULT_TIMEOUT_SECONDS);
+    }
+
+    String executeToolCall(final ServerWebExchange originExchange,
+                           final ShenyuPluginChain chain,
+                           final String sessionId,
+                           final String configStr,
+                           final String input,
+                           final long timeoutSeconds) {
 
         final RequestConfigHelper configHelper = new 
RequestConfigHelper(configStr);
         final String toolMethod = configHelper.getMethod();
@@ -236,7 +247,7 @@ public class ShenyuToolCallback implements ToolCallback {
         final boolean isTemporarySession = sessionId.startsWith("temp_");
 
         // Execute the plugin chain asynchronously
-        chain.execute(decoratedExchange)
+        final Disposable disposable = chain.execute(decoratedExchange)
                 .doOnSubscribe(s -> LOG.debug("Plugin chain subscribed for 
session: {}", sessionId))
                 .doOnError(e -> {
                     LOG.error("Plugin chain execution failed for session {}: 
{}", sessionId, e.getMessage(), e);
@@ -267,12 +278,16 @@ public class ShenyuToolCallback implements ToolCallback {
 
         // Wait for the response with timeout
         try {
-            final String result = responseFuture.get(DEFAULT_TIMEOUT_SECONDS, 
TimeUnit.SECONDS);
+            final String result = responseFuture.get(timeoutSeconds, 
TimeUnit.SECONDS);
             LOG.debug("Tool call completed successfully for session: {}", 
sessionId);
             return result;
         } catch (Exception e) {
             LOG.error("Timeout or error waiting for response for session {}: 
{}", sessionId, e.getMessage(), e);
 
+            if (e instanceof TimeoutException) {
+                disposable.dispose();
+            }
+
             // Ensure cleanup on error for temporary sessions
             if (isTemporarySession) {
                 LOG.debug("Emergency cleanup of temporary session on error: 
{}", sessionId);
diff --git 
a/shenyu-plugin/shenyu-plugin-mcp-server/src/test/java/org/apache/shenyu/plugin/mcp/server/callback/ShenyuToolCallbackTest.java
 
b/shenyu-plugin/shenyu-plugin-mcp-server/src/test/java/org/apache/shenyu/plugin/mcp/server/callback/ShenyuToolCallbackTest.java
index 69acb2848f..eb2302558a 100644
--- 
a/shenyu-plugin/shenyu-plugin-mcp-server/src/test/java/org/apache/shenyu/plugin/mcp/server/callback/ShenyuToolCallbackTest.java
+++ 
b/shenyu-plugin/shenyu-plugin-mcp-server/src/test/java/org/apache/shenyu/plugin/mcp/server/callback/ShenyuToolCallbackTest.java
@@ -34,9 +34,13 @@ import org.mockito.Mockito;
 import org.mockito.junit.jupiter.MockitoExtension;
 import org.springframework.ai.chat.model.ToolContext;
 import org.springframework.http.server.reactive.ServerHttpRequest;
+import org.springframework.mock.http.server.reactive.MockServerHttpRequest;
+import org.springframework.mock.web.server.MockServerWebExchange;
 import org.springframework.web.server.ServerWebExchange;
+import reactor.core.publisher.Mono;
 
 import java.util.HashMap;
+import java.util.concurrent.atomic.AtomicBoolean;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertFalse;
@@ -252,6 +256,26 @@ class ShenyuToolCallbackTest {
         });
     }
 
+    @Test
+    void testExecuteToolCallDisposesPluginChainOnTimeout() {
+        shenyuToolCallback = new ShenyuToolCallback(toolDefinition);
+
+        final String sessionId = "session123";
+        final MockServerWebExchange webExchange = MockServerWebExchange.from(
+                
MockServerHttpRequest.get("http://localhost/original";).build());
+        final AtomicBoolean cancelled = new AtomicBoolean();
+
+        when(chain.execute(any())).thenReturn(Mono.<Void>never().doOnCancel(() 
-> cancelled.set(true)));
+
+        final RuntimeException exception = assertThrows(RuntimeException.class,
+                () -> shenyuToolCallback.executeToolCall(webExchange, chain, 
sessionId,
+                        
"{\"requestTemplate\":{\"url\":\"/test\",\"method\":\"GET\"},\"argsPosition\":{}}",
+                        "{}", 0));
+
+        assertTrue(exception.getMessage().contains("Tool execution timeout or 
error"));
+        assertTrue(cancelled.get());
+    }
+
     @Test
     void testConstructorWithNullToolDefinition() {
         assertThrows(NullPointerException.class, () -> {

Reply via email to