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 de6d66304b fix(mcp): reuse configured response json mapper (#7116)
de6d66304b is described below
commit de6d66304b7c04bf41681e52acb7b5a3cc8761a9
Author: Liming Deng <[email protected]>
AuthorDate: Tue Sep 22 08:44:52 2026 +0800
fix(mcp): reuse configured response json mapper (#7116)
Co-authored-by: aias00 <[email protected]>
---
.../server/transport/MessageHandlingResult.java | 10 ++--
...henyuStreamableHttpServerTransportProvider.java | 34 ++++++++------
.../transport/MessageHandlingResultTest.java | 53 ++++++++++++++++++++++
3 files changed, 79 insertions(+), 18 deletions(-)
diff --git
a/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/transport/MessageHandlingResult.java
b/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/transport/MessageHandlingResult.java
index 7ecb36399b..f63ad3dfa2 100644
---
a/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/transport/MessageHandlingResult.java
+++
b/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/transport/MessageHandlingResult.java
@@ -17,7 +17,7 @@
package org.apache.shenyu.plugin.mcp.server.transport;
-import com.fasterxml.jackson.databind.ObjectMapper;
+import io.modelcontextprotocol.json.McpJsonMapper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -36,17 +36,21 @@ public class MessageHandlingResult {
private final String sessionId;
+ private final McpJsonMapper jsonMapper;
+
/**
* Creates a new message handling result.
*
* @param statusCode the HTTP status code for the response
* @param responseBody the response body object
* @param sessionId the session identifier for correlation (nullable)
+ * @param jsonMapper the configured MCP JSON mapper
*/
- public MessageHandlingResult(final int statusCode, final Object
responseBody, final String sessionId) {
+ public MessageHandlingResult(final int statusCode, final Object
responseBody, final String sessionId, final McpJsonMapper jsonMapper) {
this.statusCode = statusCode;
this.responseBody = responseBody;
this.sessionId = sessionId;
+ this.jsonMapper = jsonMapper;
}
/**
@@ -90,7 +94,7 @@ public class MessageHandlingResult {
}
try {
- return new ObjectMapper().writeValueAsString(responseBody);
+ return jsonMapper.writeValueAsString(responseBody);
} catch (Exception e) {
LOGGER.error("Failed to serialize response body to JSON: {}",
e.getMessage());
return
"{\"jsonrpc\":\"2.0\",\"error\":{\"code\":-32603,\"message\":\"Internal
error\"}}";
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 8773277d3d..370e445fd4 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
@@ -166,6 +166,10 @@ public class ShenyuStreamableHttpServerTransportProvider
implements McpServerTra
return new StreamableHttpProviderBuilder();
}
+ private MessageHandlingResult createMessageHandlingResult(final int
statusCode, final Object responseBody, final String sessionId) {
+ return new MessageHandlingResult(statusCode, responseBody, sessionId,
jsonMapper);
+ }
+
@Override
public void setSessionFactory(final McpServerSession.Factory
sessionFactory) {
this.sessionFactory = sessionFactory;
@@ -289,11 +293,11 @@ public class ShenyuStreamableHttpServerTransportProvider
implements McpServerTra
} catch (IOException e) {
LOGGER.warn("Failed to parse JSON-RPC message: {}",
e.getMessage());
final Object errorResponse = createJsonRpcError(null,
-32700, "Parse error: Invalid JSON-RPC message");
- return Mono.just(new MessageHandlingResult(400,
errorResponse, null));
+ return Mono.just(createMessageHandlingResult(400,
errorResponse, null));
} catch (Exception e) {
LOGGER.error("Unexpected error handling message: {}",
e.getMessage(), e);
final Object errorResponse = createJsonRpcError(null,
-32603, "Internal error: " + e.getMessage());
- return Mono.just(new MessageHandlingResult(500,
errorResponse, null));
+ return Mono.just(createMessageHandlingResult(500,
errorResponse, null));
}
});
}
@@ -340,17 +344,17 @@ public class ShenyuStreamableHttpServerTransportProvider
implements McpServerTra
cleanupInvalidSession(newSessionId);
final Object errorResponse = createJsonRpcError(messageId,
-32600,
"Unsupported protocol version. Supported versions: " +
SUPPORTED_PROTOCOL_VERSIONS);
- return Mono.just(new MessageHandlingResult(400, errorResponse,
null));
+ return Mono.just(createMessageHandlingResult(400,
errorResponse, null));
}
// Create initialize response
final Object initializeResponse =
createInitializeResponse(messageId, clientProtocolVersion, newSessionId);
LOGGER.debug("Initialize request processed successfully for
session: {}", newSessionId);
- return Mono.just(new MessageHandlingResult(200,
initializeResponse, newSessionId));
+ return Mono.just(createMessageHandlingResult(200,
initializeResponse, newSessionId));
} catch (Exception e) {
LOGGER.error("Error handling initialize request: {}",
e.getMessage(), e);
final Object errorResponse =
createJsonRpcError(extractMessageId(message), -32603,
"Internal error during initialization: " + e.getMessage());
- return Mono.just(new MessageHandlingResult(500, errorResponse,
null));
+ return Mono.just(createMessageHandlingResult(500, errorResponse,
null));
}
}
@@ -395,7 +399,7 @@ public class ShenyuStreamableHttpServerTransportProvider
implements McpServerTra
.map(result -> {
if (!requestedSessionId.equals(result.getSessionId())) {
LOGGER.info("Returning actual session ID {} instead of
requested ID {}", result.getSessionId(), requestedSessionId);
- return new
MessageHandlingResult(result.getStatusCode(), result.getResponseBody(),
result.getSessionId());
+ return
createMessageHandlingResult(result.getStatusCode(), result.getResponseBody(),
result.getSessionId());
}
return result;
});
@@ -438,12 +442,12 @@ public class ShenyuStreamableHttpServerTransportProvider
implements McpServerTra
removeSession(actualSessionId);
ShenyuMcpExchangeHolder.remove(actualSessionId);
})
- .map(result -> new
MessageHandlingResult(result.getStatusCode(), result.getResponseBody(), null));
+ .map(result ->
createMessageHandlingResult(result.getStatusCode(), result.getResponseBody(),
null));
} catch (Exception e) {
LOGGER.error("Error creating temporary session: {}",
e.getMessage(), e);
final Object errorResponse = createJsonRpcError(messageId, -32603,
"Internal error creating temporary session: " +
e.getMessage());
- return Mono.just(new MessageHandlingResult(500, errorResponse,
null));
+ return Mono.just(createMessageHandlingResult(500, errorResponse,
null));
}
}
@@ -488,7 +492,7 @@ public class ShenyuStreamableHttpServerTransportProvider
implements McpServerTra
.map(result -> {
if (!actualSessionId.equals(requestedSessionId)) {
LOGGER.info("Returning actual session ID {}
instead of requested ID {}", actualSessionId, requestedSessionId);
- return new
MessageHandlingResult(result.getStatusCode(), result.getResponseBody(),
actualSessionId);
+ return
createMessageHandlingResult(result.getStatusCode(), result.getResponseBody(),
actualSessionId);
}
return result;
});
@@ -496,7 +500,7 @@ public class ShenyuStreamableHttpServerTransportProvider
implements McpServerTra
LOGGER.error("Error creating session with restored ID {}: {}",
requestedSessionId, e.getMessage(), e);
final Object errorResponse = createJsonRpcError(messageId, -32603,
"Internal error restoring session: " + e.getMessage());
- return Mono.just(new MessageHandlingResult(500, errorResponse,
requestedSessionId));
+ return Mono.just(createMessageHandlingResult(500, errorResponse,
requestedSessionId));
}
}
@@ -534,12 +538,12 @@ public class ShenyuStreamableHttpServerTransportProvider
implements McpServerTra
return session.handle(message)
.cast(Object.class)
.doOnSuccess(result -> LOGGER.debug("Successfully
processed notification for session: {}", sessionId))
- .thenReturn(new
MessageHandlingResult(HttpStatus.ACCEPTED.value(), null, sessionId))
+
.thenReturn(createMessageHandlingResult(HttpStatus.ACCEPTED.value(), null,
sessionId))
.onErrorResume(error -> {
LOGGER.error("Error processing notification for
session {}: {}", sessionId, error.getMessage(), error);
final Object errorResponse = createJsonRpcError(null,
-32603,
"Internal error: " + error.getMessage());
- return Mono.just(new MessageHandlingResult(500,
errorResponse, sessionId));
+ return Mono.just(createMessageHandlingResult(500,
errorResponse, sessionId));
});
}
// Let MCP framework handle the message - framework will send response
through transport
@@ -559,7 +563,7 @@ public class ShenyuStreamableHttpServerTransportProvider
implements McpServerTra
LOGGER.error("Error processing message for session {}:
{}", sessionId, error.getMessage(), error);
final Object errorResponse = createJsonRpcError(messageId,
-32603,
"Internal error: " + error.getMessage());
- return Mono.just(new MessageHandlingResult(500,
errorResponse, sessionId));
+ return Mono.just(createMessageHandlingResult(500,
errorResponse, sessionId));
});
}
@@ -833,11 +837,11 @@ public class ShenyuStreamableHttpServerTransportProvider
implements McpServerTra
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 new MessageHandlingResult(200, sentMessage, sessionId);
+ return createMessageHandlingResult(200, sentMessage,
sessionId);
} else {
LOGGER.debug("No response captured from transport, returning
default success for session: {}", sessionId);
final Object successResponse =
createJsonRpcResponse(messageId, new java.util.HashMap<>());
- return new MessageHandlingResult(200, successResponse,
sessionId);
+ return createMessageHandlingResult(200, successResponse,
sessionId);
}
});
}
diff --git
a/shenyu-plugin/shenyu-plugin-mcp-server/src/test/java/org/apache/shenyu/plugin/mcp/server/transport/MessageHandlingResultTest.java
b/shenyu-plugin/shenyu-plugin-mcp-server/src/test/java/org/apache/shenyu/plugin/mcp/server/transport/MessageHandlingResultTest.java
new file mode 100644
index 0000000000..e20e9be56f
--- /dev/null
+++
b/shenyu-plugin/shenyu-plugin-mcp-server/src/test/java/org/apache/shenyu/plugin/mcp/server/transport/MessageHandlingResultTest.java
@@ -0,0 +1,53 @@
+/*
+ * 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.transport;
+
+import io.modelcontextprotocol.json.McpJsonMapper;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
+import static org.mockito.Mockito.when;
+
+/**
+ * Test cases for MessageHandlingResult.
+ */
+public final class MessageHandlingResultTest {
+
+ @Test
+ public void testConfiguredMapperSerializesResponseBody() throws Exception {
+ McpJsonMapper jsonMapper = mock(McpJsonMapper.class);
+ Object responseBody = new Object();
+
when(jsonMapper.writeValueAsString(responseBody)).thenReturn("{\"jsonrpc\":\"2.0\"}");
+ MessageHandlingResult result = new MessageHandlingResult(200,
responseBody, "session", jsonMapper);
+
+ assertEquals("{\"jsonrpc\":\"2.0\"}", result.getResponseBodyAsJson());
+ verify(jsonMapper).writeValueAsString(responseBody);
+ }
+
+ @Test
+ public void testStringResponseBodyDoesNotRequireSerialization() {
+ McpJsonMapper jsonMapper = mock(McpJsonMapper.class);
+ MessageHandlingResult result = new MessageHandlingResult(200,
"response", "session", jsonMapper);
+
+ assertEquals("response", result.getResponseBodyAsJson());
+ verifyNoInteractions(jsonMapper);
+ }
+}