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

Reply via email to