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