diff --git a/mcp-core/src/main/java/io/modelcontextprotocol/server/transport/HttpServletSseServerTransportProvider.java b/mcp-core/src/main/java/io/modelcontextprotocol/server/transport/HttpServletSseServerTransportProvider.java index a14271afc..81f8c4a2c 100644 --- a/mcp-core/src/main/java/io/modelcontextprotocol/server/transport/HttpServletSseServerTransportProvider.java +++ b/mcp-core/src/main/java/io/modelcontextprotocol/server/transport/HttpServletSseServerTransportProvider.java @@ -304,6 +304,7 @@ protected void doGet(HttpServletRequest request, HttpServletResponse response) response.setCharacterEncoding(UTF_8); response.setHeader("Cache-Control", "no-cache"); response.setHeader("Connection", "keep-alive"); + response.setHeader("X-Accel-Buffering", "no"); String sessionId = UUID.randomUUID().toString(); AsyncContext asyncContext = request.startAsync(); diff --git a/mcp-core/src/main/java/io/modelcontextprotocol/server/transport/HttpServletStreamableServerTransportProvider.java b/mcp-core/src/main/java/io/modelcontextprotocol/server/transport/HttpServletStreamableServerTransportProvider.java index 5f883cb98..03a9ad904 100644 --- a/mcp-core/src/main/java/io/modelcontextprotocol/server/transport/HttpServletStreamableServerTransportProvider.java +++ b/mcp-core/src/main/java/io/modelcontextprotocol/server/transport/HttpServletStreamableServerTransportProvider.java @@ -430,6 +430,7 @@ protected void doGet(HttpServletRequest request, HttpServletResponse response) response.setCharacterEncoding(UTF_8); response.setHeader("Cache-Control", "no-cache"); response.setHeader("Connection", "keep-alive"); + response.setHeader("X-Accel-Buffering", "no"); AsyncContext asyncContext = request.startAsync(); asyncContext.setTimeout(0); @@ -597,6 +598,7 @@ else if (message instanceof McpSchema.JSONRPCRequest jsonrpcRequest) { response.setCharacterEncoding(UTF_8); response.setHeader("Cache-Control", "no-cache"); response.setHeader("Connection", "keep-alive"); + response.setHeader("X-Accel-Buffering", "no"); AsyncContext asyncContext = request.startAsync(); asyncContext.setTimeout(0); diff --git a/mcp-test/src/test/java/io/modelcontextprotocol/server/HttpServletSseIntegrationTests.java b/mcp-test/src/test/java/io/modelcontextprotocol/server/HttpServletSseIntegrationTests.java index c4e334f4e..c91c6615a 100644 --- a/mcp-test/src/test/java/io/modelcontextprotocol/server/HttpServletSseIntegrationTests.java +++ b/mcp-test/src/test/java/io/modelcontextprotocol/server/HttpServletSseIntegrationTests.java @@ -121,6 +121,29 @@ public void after() { protected void prepareClients(int port, String mcpEndpoint) { } + @Test + void sseResponseIncludesXAccelBufferingHeader() throws Exception { + // https://github.com/modelcontextprotocol/java-sdk/issues/293 - without this + // header, proxies like Nginx buffer the SSE response, breaking real-time + // streaming. + prepareAsyncServerBuilder().build(); + + var httpClient = HttpClient.newHttpClient(); + var sseRequest = HttpRequest.newBuilder() + .uri(URI.create("http://localhost:" + PORT + CUSTOM_SSE_ENDPOINT)) + .header("Accept", "text/event-stream") + .GET() + .build(); + + var response = httpClient.send(sseRequest, HttpResponse.BodyHandlers.ofInputStream()); + try { + assertThat(response.headers().firstValue("X-Accel-Buffering")).contains("no"); + } + finally { + response.body().close(); + } + } + @Test void rejectsWhenBodyBytesExceedLimitWithoutContentLengthHeader() throws Exception { var httpClient = HttpClient.newHttpClient(); diff --git a/mcp-test/src/test/java/io/modelcontextprotocol/server/HttpServletStreamableIntegrationTests.java b/mcp-test/src/test/java/io/modelcontextprotocol/server/HttpServletStreamableIntegrationTests.java index c6796ce3f..a2a492d9e 100644 --- a/mcp-test/src/test/java/io/modelcontextprotocol/server/HttpServletStreamableIntegrationTests.java +++ b/mcp-test/src/test/java/io/modelcontextprotocol/server/HttpServletStreamableIntegrationTests.java @@ -14,6 +14,7 @@ import java.net.http.HttpResponse; import java.nio.ByteBuffer; import java.time.Duration; +import java.util.List; import java.util.Map; import java.util.Queue; import java.util.concurrent.CompletableFuture; @@ -204,6 +205,69 @@ void testMissingHandlerReturnsMethodNotFoundError() { } } + /** + * https://github.com/modelcontextprotocol/java-sdk/issues/293 - without this header, + * proxies like Nginx buffer the SSE response, breaking real-time streaming. + */ + @Test + void listeningStreamIncludesXAccelBufferingHeader() throws Exception { + prepareAsyncServerBuilder().serverInfo("test-server", "1.0.0").build(); + var sessionId = initializeSession(httpClient); + + var get = HttpRequest.newBuilder() + .uri(URI.create("http://localhost:" + PORT + MESSAGE_ENDPOINT)) + .header("Accept", "text/event-stream") + .header(HttpHeaders.MCP_SESSION_ID, sessionId) + .GET() + .build(); + + var response = httpClient.send(get, HttpResponse.BodyHandlers.ofInputStream()); + try { + assertThat(response.headers().firstValue("X-Accel-Buffering")).contains("no"); + } + finally { + response.body().close(); + } + } + + /** + * https://github.com/modelcontextprotocol/java-sdk/issues/293 - same header is also + * required on the SSE response a POST request opens for a streamed tool call. + */ + @Test + void toolCallSseResponseIncludesXAccelBufferingHeader() throws Exception { + prepareAsyncServerBuilder().serverInfo("test-server", "1.0.0") + .capabilities(McpSchema.ServerCapabilities.builder().tools(true).build()) + .tools(McpServerFeatures.AsyncToolSpecification.builder() + .tool(McpSchema.Tool.builder("echo", EMPTY_JSON_SCHEMA).description("returns immediately").build()) + .callHandler((exchange, + request) -> Mono.just(McpSchema.CallToolResult.builder() + .content(List.of(McpSchema.TextContent.builder("ok").build())) + .isError(false) + .build())) + .build()) + .build(); + var sessionId = initializeSession(httpClient); + + var post = HttpRequest.newBuilder() + .uri(URI.create("http://localhost:" + PORT + MESSAGE_ENDPOINT)) + .header("Content-Type", "application/json") + .header("Accept", "text/event-stream, application/json") + .header(HttpHeaders.MCP_SESSION_ID, sessionId) + .POST(HttpRequest.BodyPublishers + .ofString("{\"jsonrpc\":\"2.0\",\"id\":\"call-1\",\"method\":\"tools/call\"," + + "\"params\":{\"name\":\"echo\",\"arguments\":{}}}")) + .build(); + + var response = httpClient.send(post, HttpResponse.BodyHandlers.ofInputStream()); + try { + assertThat(response.headers().firstValue("X-Accel-Buffering")).contains("no"); + } + finally { + response.body().close(); + } + } + @Override protected void prepareClients(int port, String mcpEndpoint) { }