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, () -> {