diff --git a/src/mcp/server/streamable_http.py b/src/mcp/server/streamable_http.py index 416dd9e2b4..dc8ce1bae9 100644 --- a/src/mcp/server/streamable_http.py +++ b/src/mcp/server/streamable_http.py @@ -560,6 +560,7 @@ async def _handle_post_request(self, scope: Scope, request: Request, receive: Re writer = self._read_stream_writer if writer is None: # pragma: no cover raise ValueError("No read stream writer available. Ensure connect() is called first.") + response_sent = False try: # Validate Accept header if not await self._validate_accept_header(request, scope, send): @@ -623,6 +624,7 @@ async def _handle_post_request(self, scope: Scope, request: Request, receive: Re HTTPStatus.ACCEPTED, ) await response(scope, receive, send) + response_sent = True # Process the message after sending the response session_message = SessionMessage(message, metadata=self._message_metadata(request)) @@ -716,13 +718,19 @@ async def _handle_post_request(self, scope: Scope, request: Request, receive: Re except Exception as err: logger.exception("Error handling POST request") - response = self._create_error_response( - "Error handling POST request", - HTTPStatus.INTERNAL_SERVER_ERROR, - INTERNAL_ERROR, - ) - await response(scope, receive, send) - await writer.send(Exception(err)) + if response_sent: + logger.debug("Not sending error response: POST response already sent") + else: + response = self._create_error_response( + "Error handling POST request", + HTTPStatus.INTERNAL_SERVER_ERROR, + INTERNAL_ERROR, + ) + await response(scope, receive, send) + try: + await writer.send(Exception(err)) + except (anyio.ClosedResourceError, anyio.BrokenResourceError): + logger.debug("Writer closed while forwarding POST error; dropping exception") return async def _handle_get_request(self, request: Request, send: Send) -> None: diff --git a/tests/server/test_streamable_http_router.py b/tests/server/test_streamable_http_router.py index 0c5796c1f5..7efda2837b 100644 --- a/tests/server/test_streamable_http_router.py +++ b/tests/server/test_streamable_http_router.py @@ -1,5 +1,7 @@ """Regression coverage for the StreamableHTTP per-session response router.""" +import logging + import anyio import pytest from mcp_types import JSONRPCMessage, JSONRPCResponse @@ -157,3 +159,39 @@ async def test_terminated_transport_answers_404() -> None: assert post.sent[0]["type"] == "http.response.start" assert post.sent[0]["status"] == 404 + + +@pytest.mark.anyio +async def test_closed_writer_notification_post_sends_single_202( + caplog: pytest.LogCaptureFixture, +) -> None: + """A notification POST answered 202 before writer.send fails must not send a second response (#3651). + + The 202 completes the ASGI response, so the closed-writer error is logged + and dropped instead of answered with a 500 that Uvicorn rejects. + """ + transport = StreamableHTTPServerTransport(mcp_session_id="repro-session") + async with transport.connect(): + pass + post = _AsgiPost( + b'{"jsonrpc": "2.0", "method": "notifications/initialized"}', + [ + (b"accept", b"application/json, text/event-stream"), + (b"content-type", b"application/json"), + (b"mcp-session-id", b"repro-session"), + ], + ) + sent: list[Message] = [] + + async def strict_send(message: Message) -> None: + if message["type"] == "http.response.start" and any(m["type"] == "http.response.start" for m in sent): + raise RuntimeError("Unexpected ASGI message after completed") + sent.append(message) + + with caplog.at_level(logging.ERROR, logger="mcp.server.streamable_http"): + await transport.handle_request(post.scope, post.receive, strict_send) + + starts = [m for m in sent if m["type"] == "http.response.start"] + assert len(starts) == 1 + assert starts[0]["status"] == 202 + assert "Error handling POST request" in caplog.text