From 128c06f832caf343382b60319da396dc98f4ac8d Mon Sep 17 00:00:00 2001 From: Ramesh Reddy Adutla Date: Mon, 28 Sep 2026 13:37:32 +0100 Subject: [PATCH] Avoid opening SSE stream after non-success POST response markInitialized() and reconnect() ran before the status check, so a 405 from a server without Streamable HTTP still triggered a GET and left a duplicate SSE session during transport fallback. Only do this on 2xx. Fixes #773. --- .../HttpClientStreamableHttpTransport.java | 17 +++++++++-------- ...treamableHttpTransportErrorHandlingTest.java | 16 ++++++++++++++++ 2 files changed, 25 insertions(+), 8 deletions(-) diff --git a/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransport.java b/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransport.java index 5517823b6..65c5931bc 100644 --- a/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransport.java +++ b/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransport.java @@ -574,17 +574,18 @@ public Mono sendMessage(McpSchema.JSONRPCMessage sentMessage) { "Authorization error when sending message", requestSnapshot, responseEvent.responseInfo())); } - if (transportSession.markInitialized( - responseEvent.responseInfo().headers().firstValue("mcp-session-id").orElseGet(() -> null))) { - // Once we have a session, we try to open an async stream for - // the server to send notifications and requests out-of-band. - - reconnect(null).contextWrite(deliveredSink.contextView()).subscribe(); - } - String sessionRepresentation = sessionIdOrPlaceholder(transportSession); if (statusCode >= 200 && statusCode < 300) { + if (transportSession.markInitialized(responseEvent.responseInfo() + .headers() + .firstValue("mcp-session-id") + .orElseGet(() -> null))) { + // Once we have a session, we try to open an async stream + // for the server to send notifications and requests + // out-of-band. + reconnect(null).contextWrite(deliveredSink.contextView()).subscribe(); + } String contentType = responseEvent.responseInfo() .headers() diff --git a/mcp-test/src/test/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransportErrorHandlingTest.java b/mcp-test/src/test/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransportErrorHandlingTest.java index 0d3b69661..41d1242ee 100644 --- a/mcp-test/src/test/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransportErrorHandlingTest.java +++ b/mcp-test/src/test/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransportErrorHandlingTest.java @@ -389,6 +389,22 @@ void test405OnConnectReturnsEmptyFlux() { StepVerifier.create(transport.closeGracefully()).verifyComplete(); } + @Test + void test405OnSendMessageDoesNotOpenSseConnection() { + serverResponseStatus.set(405); + currentServerSessionId.set("ignored-session-id"); + + StepVerifier.create(transport.sendMessage(createTestRequestMessage())) + .expectErrorMatches( + error -> error instanceof McpTransportException && error.getMessage().contains("Status code: 405")) + .verify(); + + Awaitility.await() + .during(Duration.ofMillis(300)) + .atMost(Duration.ofSeconds(1)) + .untilAsserted(() -> assertThat(processedSseConnectCount.get()).isZero()); + } + @Nested class AuthorizationError {