From dd16fe06c4e2d127ad1a58c183f79624e852d35f Mon Sep 17 00:00:00 2001 From: Max Isbey <224885523+maxisbey@users.noreply.github.com> Date: Wed, 19 Aug 2026 14:51:55 +0000 Subject: [PATCH] Lift httpx2's default SSE event size cap on every client SSE reader httpx2 2.10 caps a single server-sent event at 1 MiB by default, on both `client.sse()` and a bare `EventSource(response)`, and raises `SSEError` past it. A JSON-RPC message is one event and MCP sets no message size limit, so any tool result or server notification over ~1 MiB delivered over SSE failed with `SSE stream ended without a response` (the cause was logged at DEBUG only), and with an event store the client also replayed the same oversized event on each reconnect attempt. Pass `max_event_size=None` at all five sites - the POST response stream, the standalone GET stream, resumption, reconnection, and the legacy SSE transport - which restores the pre-2.10 behaviour and matches the unbounded `application/json` response path. The floor moves to `httpx2>=2.10.0`, the first version that accepts the argument. Fixes #3332 --- pyproject.toml | 2 +- src/mcp/client/sse.py | 4 +- src/mcp/client/streamable_http.py | 10 +-- src/mcp/shared/_httpx_utils.py | 7 +- tests/interaction/_requirements.py | 19 +++++ .../transports/test_client_transport_http.py | 73 ++++++++++++++++++- tests/shared/test_sse.py | 25 +++++++ 7 files changed, 130 insertions(+), 10 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index d12cb7485e..15c7d3511e 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -130,7 +130,7 @@ dependencies = [ # stderr (agronholm/anyio#816, fixed in 4.10). "anyio>=4.10; python_version >= '3.14'", "anyio>=4.9; python_version < '3.14'", - "httpx2>=2.5.0", + "httpx2>=2.10.0", "mcp-types=={{ version }}", "pydantic>=2.12.0", "starlette>=0.48.0; python_version >= '3.14'", diff --git a/src/mcp/client/sse.py b/src/mcp/client/sse.py index 31d0f35391..ca1f5de81e 100644 --- a/src/mcp/client/sse.py +++ b/src/mcp/client/sse.py @@ -12,7 +12,7 @@ from mcp.shared._compat import resync_tracer from mcp.shared._context_streams import create_context_streams -from mcp.shared._httpx_utils import McpHttpClientFactory, create_mcp_http_client +from mcp.shared._httpx_utils import MCP_SSE_MAX_EVENT_SIZE, McpHttpClientFactory, create_mcp_http_client from mcp.shared.message import SessionMessage logger = logging.getLogger(__name__) @@ -55,7 +55,7 @@ async def sse_client( async with httpx_client_factory( headers=headers, auth=auth, timeout=httpx2.Timeout(timeout, read=sse_read_timeout) ) as client: - async with client.sse(url) as event_source: + async with client.sse(url, max_event_size=MCP_SSE_MAX_EVENT_SIZE) as event_source: event_source.response.raise_for_status() logger.debug("SSE connection established") diff --git a/src/mcp/client/streamable_http.py b/src/mcp/client/streamable_http.py index 226b0fecf9..42d1392ce2 100644 --- a/src/mcp/client/streamable_http.py +++ b/src/mcp/client/streamable_http.py @@ -33,7 +33,7 @@ from mcp.client._transport import TransportStreams from mcp.shared._compat import resync_tracer from mcp.shared._context_streams import ContextReceiveStream, ContextSendStream, create_context_streams -from mcp.shared._httpx_utils import create_mcp_http_client +from mcp.shared._httpx_utils import MCP_SSE_MAX_EVENT_SIZE, create_mcp_http_client from mcp.shared.inbound import MCP_PROTOCOL_VERSION_HEADER from mcp.shared.jsonrpc_dispatcher import cancelled_request_id_from_params from mcp.shared.message import ClientMessageMetadata, SessionMessage @@ -210,7 +210,7 @@ async def handle_get_stream(self, client: httpx2.AsyncClient, read_stream_writer if last_event_id: headers[LAST_EVENT_ID] = last_event_id - async with client.sse(self.url, headers=headers) as event_source: + async with client.sse(self.url, headers=headers, max_event_size=MCP_SSE_MAX_EVENT_SIZE) as event_source: event_source.response.raise_for_status() logger.debug("GET SSE connection established") @@ -253,7 +253,7 @@ async def _handle_resumption_request(self, ctx: RequestContext) -> None: if isinstance(ctx.session_message.message, JSONRPCRequest): # pragma: no branch original_request_id = ctx.session_message.message.id - async with ctx.client.sse(self.url, headers=headers) as event_source: + async with ctx.client.sse(self.url, headers=headers, max_event_size=MCP_SSE_MAX_EVENT_SIZE) as event_source: event_source.response.raise_for_status() logger.debug("Resumption GET SSE connection established") @@ -423,7 +423,7 @@ async def _handle_sse_response( original_request_id = ctx.session_message.message.id try: - event_source = EventSource(response) + event_source = EventSource(response, max_event_size=MCP_SSE_MAX_EVENT_SIZE) async for sse in event_source: # pragma: no branch # Track last event ID for potential reconnection if sse.id: @@ -501,7 +501,7 @@ async def _handle_reconnection( headers[LAST_EVENT_ID] = last_event_id try: - async with ctx.client.sse(self.url, headers=headers) as event_source: + async with ctx.client.sse(self.url, headers=headers, max_event_size=MCP_SSE_MAX_EVENT_SIZE) as event_source: event_source.response.raise_for_status() logger.info("Reconnected to SSE stream") diff --git a/src/mcp/shared/_httpx_utils.py b/src/mcp/shared/_httpx_utils.py index 6bb638886a..ba880afcef 100644 --- a/src/mcp/shared/_httpx_utils.py +++ b/src/mcp/shared/_httpx_utils.py @@ -4,12 +4,17 @@ import httpx2 -__all__ = ["create_mcp_http_client", "MCP_DEFAULT_TIMEOUT", "MCP_DEFAULT_SSE_READ_TIMEOUT"] +__all__ = ["create_mcp_http_client", "MCP_DEFAULT_TIMEOUT", "MCP_DEFAULT_SSE_READ_TIMEOUT", "MCP_SSE_MAX_EVENT_SIZE"] # Default MCP timeout configuration MCP_DEFAULT_TIMEOUT = 30.0 # General operations (seconds) MCP_DEFAULT_SSE_READ_TIMEOUT = 300.0 # SSE streams - 5 minutes (seconds) +# httpx2 >= 2.10 caps a single SSE event at 1 MiB by default. One JSON-RPC message +# is one event and MCP sets no message size limit (the application/json response +# path is unbounded too), so every SSE reader the SDK opens passes this instead. +MCP_SSE_MAX_EVENT_SIZE: int | None = None + class McpHttpClientFactory(Protocol): # pragma: no branch def __call__( # pragma: no branch diff --git a/tests/interaction/_requirements.py b/tests/interaction/_requirements.py index 86725bcb4f..4cd87894ca 100644 --- a/tests/interaction/_requirements.py +++ b/tests/interaction/_requirements.py @@ -3474,6 +3474,25 @@ def __post_init__(self) -> None: transports=("streamable-http",), note="Only observable over HTTP: per-request SSE streams are HTTP-specific.", ), + "client-transport:http:post-stream-large-event": Requirement( + source="sdk", + behavior=( + "A JSON-RPC message larger than httpx2's default 1 MiB SSE event cap is delivered intact over the " + "per-request POST stream; the transport lifts the cap because MCP sets no message size limit." + ), + transports=("streamable-http",), + note="Only observable over HTTP: SSE event framing is HTTP-specific.", + ), + "client-transport:http:get-stream-large-event": Requirement( + source="sdk", + behavior=( + "A server-initiated message larger than httpx2's default 1 MiB SSE event cap is delivered intact " + "over the standalone GET stream; the transport lifts the cap because MCP sets no message size limit." + ), + transports=("streamable-http",), + removed_in="2026-07-28", + note="removed in 2026-07-28 (SEP-2575); the standalone GET stream is replaced by subscriptions/listen.", + ), "client-transport:http:custom-client": Requirement( source="sdk", behavior=( diff --git a/tests/interaction/transports/test_client_transport_http.py b/tests/interaction/transports/test_client_transport_http.py index 625fcaad8d..3131b034fd 100644 --- a/tests/interaction/transports/test_client_transport_http.py +++ b/tests/interaction/transports/test_client_transport_http.py @@ -14,13 +14,23 @@ import mcp_types as types import pytest from inline_snapshot import snapshot -from mcp_types import INVALID_REQUEST, CallToolResult, ErrorData, ListToolsResult, TextContent, Tool +from mcp_types import ( + INVALID_REQUEST, + CallToolResult, + ErrorData, + ListToolsResult, + LoggingMessageNotification, + LoggingMessageNotificationParams, + TextContent, + Tool, +) from starlette.types import Receive, Scope, Send from mcp import MCPError from mcp.client.client import Client from mcp.client.streamable_http import streamable_http_client from mcp.server import Server, ServerRequestContext +from mcp.server.mcpserver import Context, MCPServer from tests.interaction._connect import BASE_URL, NO_DNS_REBINDING_PROTECTION, client_via_http, mounted_app from tests.interaction._requirements import requirement from tests.interaction.transports._bridge import StreamingASGITransport @@ -158,6 +168,67 @@ async def call(n: int) -> None: assert len(tools_call_posts) == 3 +# One byte past the 1 MiB that httpx2 >= 2.10 allows a single SSE event by default. +_OVERSIZED_TEXT = "x" * (1024 * 1024 + 1) + + +@requirement("client-transport:http:post-stream-large-event") +async def test_a_post_stream_delivers_a_tool_result_larger_than_one_mebibyte() -> None: + """A tool result bigger than httpx2's default per-event SSE cap arrives intact over the request's + POST stream. SDK-defined: MCP sets no message size limit, so the transport lifts the cap (#3332).""" + mcp = MCPServer("bulky") + + @mcp.tool() + def bulk() -> str: + """Return more than one SSE event may carry by default.""" + return _OVERSIZED_TEXT + + async with mounted_app(mcp) as (http, _), client_via_http(http) as client: + with anyio.fail_after(5): + result = await client.call_tool("bulk", {}) + + assert result.content == [TextContent(text=_OVERSIZED_TEXT)] + + +@requirement("client-transport:http:get-stream-large-event") +async def test_the_standalone_get_stream_delivers_a_notification_larger_than_one_mebibyte() -> None: + """A server-initiated notification bigger than httpx2's default per-event SSE cap arrives intact + over the standalone GET stream, which the transport opens with the same lifted cap (#3332).""" + mcp = MCPServer("bulky") + + @mcp.tool() + async def shout(ctx: Context) -> str: + """Emit one unrelated notification, which the server routes to the standalone stream.""" + params = LoggingMessageNotificationParams(level="info", data=_OVERSIZED_TEXT) + await ctx.session.send_notification(LoggingMessageNotification(params=params)) + return "sent" + + get_stream_open = anyio.Event() + + async def on_response(response: httpx2.Response) -> None: + if response.request.method == "GET": + get_stream_open.set() + + received: list[object] = [] + delivered = anyio.Event() + + async def collect(params: LoggingMessageNotificationParams) -> None: + received.append(params.data) + delivered.set() + + async with ( + mounted_app(mcp, on_response=on_response) as (http, _), + client_via_http(http, logging_callback=collect) as client, + ): + with anyio.fail_after(5): + # The server drops standalone messages emitted before the GET stream is established. + await get_stream_open.wait() + await client.call_tool("shout", {}) + await delivered.wait() + + assert received == [_OVERSIZED_TEXT] + + @requirement("client-transport:http:sse-405-tolerated") @requirement("client-transport:http:terminate-405-ok") async def test_client_tolerates_405_on_get_and_delete() -> None: diff --git a/tests/shared/test_sse.py b/tests/shared/test_sse.py index c27dd69db3..427ede9eea 100644 --- a/tests/shared/test_sse.py +++ b/tests/shared/test_sse.py @@ -224,6 +224,31 @@ async def test_sse_client_exception_handling( await session.read_resource(uri="xxx://will-not-work") +@pytest.mark.anyio +async def test_sse_client_delivers_a_result_larger_than_one_mebibyte() -> None: + """A resource read bigger than httpx2's default 1 MiB per-event SSE cap arrives intact. SDK-defined: + MCP sets no message size limit, and the legacy transport carries every server message on one + event stream, so the cap is lifted there too (#3332).""" + oversized = "x" * (1024 * 1024 + 1) + + async def read_resource(ctx: ServerRequestContext, params: ReadResourceRequestParams) -> ReadResourceResult: + return ReadResourceResult( + contents=[TextResourceContents(uri=str(params.uri), text=oversized, mime_type="text/plain")] + ) + + factory = in_process_client_factory(make_app(Server(SERVER_NAME, on_read_resource=read_resource))) + with anyio.fail_after(5): + async with ( + sse_client(f"{BASE_URL}/sse", httpx_client_factory=factory) as streams, + ClientSession(*streams) as session, + ): + await session.initialize() + response = await session.read_resource(uri="foobar://bulk") + + assert isinstance(response.contents[0], TextResourceContents) + assert response.contents[0].text == oversized + + @pytest.mark.anyio async def test_sse_client_basic_connection_mounted_app() -> None: """The SSE transport works unchanged when its app is mounted under a sub-path."""