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 7c1efd8758 [type: bug] Fix concurrent tools/call response correlation
by message id (#7041) (#7046)
7c1efd8758 is described below
commit 7c1efd8758c027b70cf2758923e8772cb4e5de8e
Author: Arvin <[email protected]>
AuthorDate: Wed Sep 30 09:56:20 2026 +0800
[type: bug] Fix concurrent tools/call response correlation by message id
(#7041) (#7046)
---
...henyuStreamableHttpServerTransportProvider.java | 49 +++++++++++++++++++---
1 file changed, 44 insertions(+), 5 deletions(-)
diff --git
a/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/transport/ShenyuStreamableHttpServerTransportProvider.java
b/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/transport/ShenyuStreamableHttpServerTransportProvider.java
index 370e445fd4..cdb0dc37d5 100644
---
a/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/transport/ShenyuStreamableHttpServerTransportProvider.java
+++
b/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/transport/ShenyuStreamableHttpServerTransportProvider.java
@@ -552,11 +552,11 @@ public class ShenyuStreamableHttpServerTransportProvider
implements McpServerTra
.doOnSuccess(result -> LOGGER.debug("Successfully processed
message for session: {}", sessionId))
.then(waitForTransportResponse(transport, sessionId,
messageId))
.doOnNext(result -> {
- // Clear the captured response after each completed
message so that a
- // subsequent message on this session cannot observe a
stale response
- // from a previous request.
+ // Clear the response captured for this specific message
id after it has
+ // been delivered, so that a subsequent message on this
session cannot
+ // observe a stale response from a previous request.
if (Objects.nonNull(transport)) {
- transport.resetCapturedMessage();
+ transport.resetCapturedMessage(messageId);
}
})
.onErrorResume(error -> {
@@ -834,7 +834,12 @@ public class ShenyuStreamableHttpServerTransportProvider
implements McpServerTra
final String
sessionId,
final Object
messageId) {
return Mono.fromCallable(() -> {
- if (Objects.nonNull(transport) && transport.isResponseReady() &&
Objects.nonNull(transport.getLastSentMessage())) {
+ final McpSchema.JSONRPCMessage correlatedResponse =
Objects.nonNull(transport)
+ ? transport.getLastSentMessage(messageId) : null;
+ if (Objects.nonNull(messageId) &&
Objects.nonNull(correlatedResponse)) {
+ LOGGER.debug("Retrieved correlated response for message id {}
on session: {}", messageId, sessionId);
+ return new MessageHandlingResult(200, correlatedResponse,
sessionId);
+ } else if (Objects.nonNull(transport) &&
transport.isResponseReady() && Objects.nonNull(transport.getLastSentMessage()))
{
final McpSchema.JSONRPCMessage sentMessage =
transport.getLastSentMessage();
LOGGER.debug("Retrieved captured response from transport for
session: {}", sessionId);
return createMessageHandlingResult(200, sentMessage,
sessionId);
@@ -1055,6 +1060,8 @@ public class ShenyuStreamableHttpServerTransportProvider
implements McpServerTra
private volatile boolean responseReady;
+ private final Map<String, McpSchema.JSONRPCMessage> messageResponses =
new ConcurrentHashMap<>();
+
/**
* Creates a new session transport with auto-generated session ID.
*/
@@ -1084,6 +1091,20 @@ public class ShenyuStreamableHttpServerTransportProvider
implements McpServerTra
return lastSentMessage;
}
+ /**
+ * Gets the response message captured for the given message id,
falling back
+ * to the last sent message when no id-based correlation is available.
+ *
+ * @param messageId the JSON-RPC message id to look up
+ * @return the correlated response, or null if none has been captured
+ */
+ public McpSchema.JSONRPCMessage getLastSentMessage(final Object
messageId) {
+ if (Objects.nonNull(messageId)) {
+ return messageResponses.get(String.valueOf(messageId));
+ }
+ return lastSentMessage;
+ }
+
/**
* Checks if a response is ready for retrieval.
*
@@ -1098,6 +1119,12 @@ public class ShenyuStreamableHttpServerTransportProvider
implements McpServerTra
if (!closed) {
this.lastSentMessage = message;
this.responseReady = true;
+ if (message instanceof McpSchema.JSONRPCResponse) {
+ final Object responseId = ((McpSchema.JSONRPCResponse)
message).id();
+ if (Objects.nonNull(responseId)) {
+ this.messageResponses.put(String.valueOf(responseId),
message);
+ }
+ }
LOGGER.debug("Captured response message for session: {}",
sessionId);
}
return Mono.empty();
@@ -1137,6 +1164,18 @@ public class ShenyuStreamableHttpServerTransportProvider
implements McpServerTra
this.lastSentMessage = null;
this.responseReady = false;
}
+
+ /**
+ * Clears the response captured for the given message id, leaving the
+ * last-sent-message slot used by the initialization handshake intact.
+ *
+ * @param messageId the JSON-RPC message id whose captured response
should be removed
+ */
+ public void resetCapturedMessage(final Object messageId) {
+ if (Objects.nonNull(messageId)) {
+ this.messageResponses.remove(String.valueOf(messageId));
+ }
+ }
}
}