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 e3cc44e298 fix: return 202 empty body for JSON-RPC notifications and 
reset captured response in Streamable HTTP transport (#6985)
e3cc44e298 is described below

commit e3cc44e2980bfea38eec261ed57ca43a0bd9ccc3
Author: wy471x <[email protected]>
AuthorDate: Fri Sep 4 12:06:07 2026 +0800

    fix: return 202 empty body for JSON-RPC notifications and reset captured 
response in Streamable HTTP transport (#6985)
    
    * fix: return 202 empty body for JSON-RPC notifications and reset captured 
response in Streamable HTTP transport
    
    - Detect JSONRPCNotification in processWithExistingSession and acknowledge
      with HTTP 202 empty body instead of waiting for a transport response that
      could replay the previous request's stale response or a fabricated one
    - Map null response body to an empty HTTP response in handleUnifiedEndpoint
    - Reset the captured transport message after each completed response so
      subsequent messages cannot observe stale state
    - Add regression tests covering notifications without a prior request and
      notifications after a business request
    
    Co-Authored-By: Claude <[email protected]>
    
    * fix: mcp streamable http notification returns empty 202 response in 
plugin path
    
    ---------
    
    Co-authored-by: Claude <[email protected]>
    Co-authored-by: aias00 <[email protected]>
---
 .../shenyu/plugin/mcp/server/McpServerPlugin.java  |  10 +-
 ...henyuStreamableHttpServerTransportProvider.java |  31 +++++
 .../plugin/mcp/server/McpServerPluginTest.java     |  52 ++++++++
 ...uStreamableHttpServerTransportProviderTest.java | 131 +++++++++++++++++++++
 4 files changed, 223 insertions(+), 1 deletion(-)

diff --git 
a/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/McpServerPlugin.java
 
b/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/McpServerPlugin.java
index 0c268eb361..ef62e517b2 100644
--- 
a/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/McpServerPlugin.java
+++ 
b/shenyu-plugin/shenyu-plugin-mcp-server/src/main/java/org/apache/shenyu/plugin/mcp/server/McpServerPlugin.java
@@ -474,7 +474,16 @@ public class McpServerPlugin extends AbstractShenyuPlugin {
         // Configure response
         configureStreamableHttpResponse(exchange, result);
 
+        // Notifications are acknowledged with 202 Accepted and an empty body
+        // per the Streamable HTTP spec: a null response body means no body
+        // should be written.
+        if (Objects.isNull(result.getResponseBody())) {
+            LOG.debug("Response body is null, completing response without 
body");
+            return exchange.getResponse().setComplete();
+        }
+
         // Write response body
+        exchange.getResponse().getHeaders().set("Content-Type", 
"application/json");
         final String responseBodyJson = result.getResponseBodyAsJson();
         final byte[] responseBytes = 
responseBodyJson.getBytes(StandardCharsets.UTF_8);
 
@@ -511,7 +520,6 @@ public class McpServerPlugin extends AbstractShenyuPlugin {
 
         // Set standard headers
         setCorsHeaders(exchange);
-        exchange.getResponse().getHeaders().set("Content-Type", 
"application/json");
 
         // Clean up potentially conflicting headers
         exchange.getResponse().getHeaders().remove("Transfer-Encoding");
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 540d834853..3ba4371c63 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
@@ -247,6 +247,11 @@ public class ShenyuStreamableHttpServerTransportProvider 
implements McpServerTra
                 if (Objects.nonNull(result.getSessionId())) {
                     builder.header(SESSION_ID_HEADER, result.getSessionId());
                 }
+                // Notifications are acknowledged with 202 Accepted and an 
empty body
+                // per the Streamable HTTP spec.
+                if (Objects.isNull(result.getResponseBody())) {
+                    return builder.build();
+                }
                 builder.contentType(MediaType.APPLICATION_JSON);
                 return builder.bodyValue(result.getResponseBodyAsJson());
             });
@@ -513,11 +518,37 @@ public class ShenyuStreamableHttpServerTransportProvider 
implements McpServerTra
                     sessionId, System.identityHashCode(verifyExchange));
         }
         final StreamableHttpSessionTransport transport = 
getSessionTransport(sessionId);
+        // JSON-RPC notifications (messages without an id, e.g. 
notifications/initialized,
+        // notifications/cancelled) must be acknowledged with HTTP 202 and an 
empty body
+        // per the Streamable HTTP spec. They are still dispatched to the MCP 
framework so
+        // that session state is updated, but they never produce a JSON-RPC 
response, so we
+        // must not wait for a captured transport response here. Otherwise a 
stale response
+        // from a previous request could be replayed with the wrong id.
+        if (message instanceof McpSchema.JSONRPCNotification) {
+            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))
+                    .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));
+                    });
+        }
         // Let MCP framework handle the message - framework will send response 
through transport
         return session.handle(message)
                 .cast(Object.class)
                 .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.
+                    if (Objects.nonNull(transport)) {
+                        transport.resetCapturedMessage();
+                    }
+                })
                 .onErrorResume(error -> {
                     LOGGER.error("Error processing message for session {}: 
{}", sessionId, error.getMessage(), error);
                     final Object errorResponse = createJsonRpcError(messageId, 
-32603,
diff --git 
a/shenyu-plugin/shenyu-plugin-mcp-server/src/test/java/org/apache/shenyu/plugin/mcp/server/McpServerPluginTest.java
 
b/shenyu-plugin/shenyu-plugin-mcp-server/src/test/java/org/apache/shenyu/plugin/mcp/server/McpServerPluginTest.java
index ea27e72bda..603e3ec008 100644
--- 
a/shenyu-plugin/shenyu-plugin-mcp-server/src/test/java/org/apache/shenyu/plugin/mcp/server/McpServerPluginTest.java
+++ 
b/shenyu-plugin/shenyu-plugin-mcp-server/src/test/java/org/apache/shenyu/plugin/mcp/server/McpServerPluginTest.java
@@ -24,7 +24,9 @@ import org.apache.shenyu.common.enums.PluginEnum;
 import org.apache.shenyu.common.enums.RpcTypeEnum;
 import org.apache.shenyu.plugin.api.ShenyuPluginChain;
 import org.apache.shenyu.plugin.api.context.ShenyuContext;
+import org.apache.shenyu.plugin.mcp.server.holder.ShenyuMcpExchangeHolder;
 import org.apache.shenyu.plugin.mcp.server.manager.ShenyuMcpServerManager;
+import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.extension.ExtendWith;
@@ -46,6 +48,7 @@ import java.util.Locale;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.ArgumentMatchers.anyString;
 import static org.mockito.Mockito.when;
@@ -87,6 +90,11 @@ class McpServerPluginTest {
         mcpServerPlugin = new McpServerPlugin(shenyuMcpServerManager, 
messageReaders);
     }
 
+    @AfterEach
+    void tearDown() {
+        ShenyuMcpExchangeHolder.clear();
+    }
+
     @Test
     void testNamed() {
         assertEquals(PluginEnum.MCP_SERVER.getName(), mcpServerPlugin.named());
@@ -202,4 +210,48 @@ class McpServerPluginTest {
         final String allowHeaders = 
webExchange.getResponse().getHeaders().getFirst("Access-Control-Allow-Headers");
         assertTrue(allowHeaders.toLowerCase(Locale.ROOT).contains("xrequest"));
     }
+
+    /**
+     * A JSON-RPC notification routed through the real plugin path must be
+     * acknowledged with HTTP 202 and an empty body, not a serialized null
+     * response body.
+     */
+    @Test
+    void testStreamableHttpNotificationReturnsEmptyBodyViaPluginPath() {
+        final ShenyuMcpServerManager manager = new ShenyuMcpServerManager();
+        manager.getOrCreateStreamableHttpTransport("/mcp/streamablehttp");
+        final McpServerPlugin plugin = new McpServerPlugin(manager,
+                HandlerStrategies.withDefaults().messageReaders());
+
+        // Initialize a session through the plugin path.
+        final MockServerWebExchange initExchange = 
MockServerWebExchange.from(MockServerHttpRequest
+                .post("/mcp/streamablehttp")
+                .header("Content-Type", "application/json")
+                
.body("{\"jsonrpc\":\"2.0\",\"id\":\"init-1\",\"method\":\"initialize\","
+                        + 
"\"params\":{\"protocolVersion\":\"2025-03-26\",\"capabilities\":{},"
+                        + 
"\"clientInfo\":{\"name\":\"test-client\",\"version\":\"1.0.0\"}}}"));
+        initExchange.getAttributes().put(Constants.CONTEXT, new 
ShenyuContext());
+
+        StepVerifier.create(plugin.doExecute(initExchange, chain, selector, 
rule))
+                .verifyComplete();
+
+        assertEquals(HttpStatus.OK, 
initExchange.getResponse().getStatusCode());
+        final String sessionId = 
initExchange.getResponse().getHeaders().getFirst("Mcp-Session-Id");
+        assertNotNull(sessionId);
+
+        // Send a notification on the same session through the plugin path.
+        final MockServerWebExchange notificationExchange = 
MockServerWebExchange.from(MockServerHttpRequest
+                .post("/mcp/streamablehttp")
+                .header("Content-Type", "application/json")
+                .header("Mcp-Session-Id", sessionId)
+                
.body("{\"jsonrpc\":\"2.0\",\"method\":\"notifications/initialized\",\"params\":{}}"));
+        notificationExchange.getAttributes().put(Constants.CONTEXT, new 
ShenyuContext());
+
+        StepVerifier.create(plugin.doExecute(notificationExchange, chain, 
selector, rule))
+                .verifyComplete();
+
+        assertEquals(HttpStatus.ACCEPTED, 
notificationExchange.getResponse().getStatusCode());
+        assertEquals(sessionId, 
notificationExchange.getResponse().getHeaders().getFirst("Mcp-Session-Id"));
+        assertEquals("", 
notificationExchange.getResponse().getBodyAsString().block());
+    }
 }
diff --git 
a/shenyu-plugin/shenyu-plugin-mcp-server/src/test/java/org/apache/shenyu/plugin/mcp/server/transport/ShenyuStreamableHttpServerTransportProviderTest.java
 
b/shenyu-plugin/shenyu-plugin-mcp-server/src/test/java/org/apache/shenyu/plugin/mcp/server/transport/ShenyuStreamableHttpServerTransportProviderTest.java
index d10af3ccb1..7d3bddb70b 100644
--- 
a/shenyu-plugin/shenyu-plugin-mcp-server/src/test/java/org/apache/shenyu/plugin/mcp/server/transport/ShenyuStreamableHttpServerTransportProviderTest.java
+++ 
b/shenyu-plugin/shenyu-plugin-mcp-server/src/test/java/org/apache/shenyu/plugin/mcp/server/transport/ShenyuStreamableHttpServerTransportProviderTest.java
@@ -18,16 +18,31 @@
 package org.apache.shenyu.plugin.mcp.server.transport;
 
 import com.fasterxml.jackson.databind.ObjectMapper;
+import io.modelcontextprotocol.spec.McpSchema;
+import io.modelcontextprotocol.spec.McpServerSession;
+import io.modelcontextprotocol.server.McpRequestHandler;
+import org.apache.shenyu.plugin.mcp.server.holder.ShenyuMcpExchangeHolder;
+import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.Test;
 import org.springframework.http.HttpStatus;
+import org.springframework.http.codec.HttpMessageWriter;
 import org.springframework.mock.http.server.reactive.MockServerHttpRequest;
+import org.springframework.mock.http.server.reactive.MockServerHttpResponse;
 import org.springframework.mock.web.server.MockServerWebExchange;
 import org.springframework.web.reactive.function.server.HandlerStrategies;
 import org.springframework.web.reactive.function.server.ServerRequest;
 import org.springframework.web.reactive.function.server.ServerResponse;
+import org.springframework.web.reactive.result.view.ViewResolver;
+import reactor.core.publisher.Mono;
 import reactor.test.StepVerifier;
 
+import java.time.Duration;
+import java.util.Collections;
+import java.util.List;
 import java.util.Locale;
+import java.util.Map;
+import java.util.Objects;
+import java.util.UUID;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
@@ -38,6 +53,38 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
  */
 class ShenyuStreamableHttpServerTransportProviderTest {
 
+    private static final String SESSION_ID_HEADER = "Mcp-Session-Id";
+
+    private static final String INITIALIZE_REQUEST_BODY = 
"{\"jsonrpc\":\"2.0\",\"id\":\"init-1\",\"method\":\"initialize\","
+            + 
"\"params\":{\"protocolVersion\":\"2025-03-26\",\"capabilities\":{},"
+            + 
"\"clientInfo\":{\"name\":\"test-client\",\"version\":\"1.0.0\"}}}";
+
+    private static final String TOOLS_LIST_REQUEST_BODY = 
"{\"jsonrpc\":\"2.0\",\"id\":\"req-1\",\"method\":\"tools/list\","
+            + "\"params\":{}}";
+
+    private static final String INITIALIZED_NOTIFICATION_BODY = 
"{\"jsonrpc\":\"2.0\",\"method\":\"notifications/initialized\","
+            + "\"params\":{}}";
+
+    private static final String CANCELLED_NOTIFICATION_BODY = 
"{\"jsonrpc\":\"2.0\",\"method\":\"notifications/cancelled\","
+            + "\"params\":{\"requestId\":\"req-1\",\"reason\":\"test\"}}";
+
+    private static final ServerResponse.Context RESPONSE_CONTEXT = new 
ServerResponse.Context() {
+        @Override
+        public List<HttpMessageWriter<?>> messageWriters() {
+            return HandlerStrategies.withDefaults().messageWriters();
+        }
+
+        @Override
+        public List<ViewResolver> viewResolvers() {
+            return Collections.emptyList();
+        }
+    };
+
+    @AfterEach
+    void tearDown() {
+        ShenyuMcpExchangeHolder.clear();
+    }
+
     @Test
     void testPreflightUsesConfiguredHeadersAndMethods() {
         ShenyuStreamableHttpServerTransportProvider provider =
@@ -78,6 +125,90 @@ class ShenyuStreamableHttpServerTransportProviderTest {
         assertTrue(allowHeaders.toLowerCase(Locale.ROOT).contains("xrequest"));
     }
 
+    /**
+     * JSON-RPC notifications on an existing session must be acknowledged with
+     * HTTP 202 and an empty body instead of a fabricated JSON-RPC response.
+     */
+    @Test
+    void testNotificationWithExistingSessionReturnsAcceptedEmptyBody() {
+        ShenyuStreamableHttpServerTransportProvider provider = 
providerWithRealSessions();
+        MockServerHttpResponse initResponse = performRequest(provider, 
postRequest(INITIALIZE_REQUEST_BODY, null));
+        String sessionId = 
initResponse.getHeaders().getFirst(SESSION_ID_HEADER);
+        assertNotNull(sessionId);
+
+        MockServerHttpResponse notificationResponse = performRequest(provider, 
postRequest(CANCELLED_NOTIFICATION_BODY, sessionId));
+        assertEquals(HttpStatus.ACCEPTED, 
notificationResponse.getStatusCode());
+        assertEquals(sessionId, 
notificationResponse.getHeaders().getFirst(SESSION_ID_HEADER));
+        assertEquals("", notificationResponse.getBodyAsString().block());
+    }
+
+    /**
+     * A notification sent after a business request must not replay the stale
+     * response captured by the previous request on the same session.
+     */
+    @Test
+    void testNotificationAfterRequestDoesNotReplayStaleResponse() {
+        ShenyuStreamableHttpServerTransportProvider provider = 
providerWithRealSessions();
+        MockServerHttpResponse initResponse = performRequest(provider, 
postRequest(INITIALIZE_REQUEST_BODY, null));
+        String sessionId = 
initResponse.getHeaders().getFirst(SESSION_ID_HEADER);
+        assertNotNull(sessionId);
+
+        // Complete the handshake state so the business request handler can 
run.
+        MockServerHttpResponse initializedResponse = performRequest(provider, 
postRequest(INITIALIZED_NOTIFICATION_BODY, sessionId));
+        assertEquals(HttpStatus.ACCEPTED, initializedResponse.getStatusCode());
+        assertEquals("", initializedResponse.getBodyAsString().block());
+
+        // Business request populates the captured transport response.
+        MockServerHttpResponse toolsResponse = performRequest(provider, 
postRequest(TOOLS_LIST_REQUEST_BODY, sessionId));
+        assertEquals(HttpStatus.OK, toolsResponse.getStatusCode());
+        String toolsBody = toolsResponse.getBodyAsString().block();
+        assertTrue(toolsBody.contains("\"id\":\"req-1\""));
+        assertTrue(toolsBody.contains("\"tools\""));
+
+        // Notification on the same session must be 202 with an empty body
+        // instead of replaying the stale tools/list response.
+        MockServerHttpResponse notificationResponse = performRequest(provider, 
postRequest(CANCELLED_NOTIFICATION_BODY, sessionId));
+        assertEquals(HttpStatus.ACCEPTED, 
notificationResponse.getStatusCode());
+        assertEquals("", notificationResponse.getBodyAsString().block());
+    }
+
+    private ShenyuStreamableHttpServerTransportProvider 
providerWithRealSessions() {
+        ShenyuStreamableHttpServerTransportProvider provider =
+                new ShenyuStreamableHttpServerTransportProvider(new 
ObjectMapper(), "/mcp/streamablehttp");
+        provider.setSessionFactory(transport -> new McpServerSession(
+                UUID.randomUUID().toString(),
+                Duration.ofSeconds(10),
+                transport,
+                initializeRequest -> Mono.just(new McpSchema.InitializeResult(
+                        "2025-03-26",
+                        McpSchema.ServerCapabilities.builder().build(),
+                        new McpSchema.Implementation("ShenyuMcpServer", 
"1.0.0"),
+                        "test")),
+                Map.<String, McpRequestHandler<?>>of("tools/list",
+                        (McpRequestHandler<Map<String, Object>>) (exchange, 
params) -> Mono.just(Map.<String, Object>of("tools", List.of()))),
+                Map.of()));
+        return provider;
+    }
+
+    private MockServerHttpRequest postRequest(final String body, final String 
sessionId) {
+        MockServerHttpRequest.BodyBuilder builder = 
MockServerHttpRequest.post("/mcp/streamablehttp")
+                .header("Content-Type", "application/json");
+        if (Objects.nonNull(sessionId)) {
+            builder.header(SESSION_ID_HEADER, sessionId);
+        }
+        return builder.body(body);
+    }
+
+    private MockServerHttpResponse performRequest(final 
ShenyuStreamableHttpServerTransportProvider provider,
+                                                  final MockServerHttpRequest 
request) {
+        MockServerWebExchange exchange = MockServerWebExchange.from(request);
+        ServerRequest serverRequest = ServerRequest.create(exchange, 
HandlerStrategies.withDefaults().messageReaders());
+        ServerResponse response = 
provider.handleUnifiedEndpoint(serverRequest).block();
+        assertNotNull(response);
+        response.writeTo(exchange, RESPONSE_CONTEXT).block();
+        return exchange.getResponse();
+    }
+
     private ServerRequest createRequest(final MockServerHttpRequest request) {
         return ServerRequest.create(
                 MockServerWebExchange.from(request),

Reply via email to