Skip to content

Commit dd16fe0

Browse files
committed
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
1 parent 0d92192 commit dd16fe0

7 files changed

Lines changed: 130 additions & 10 deletions

File tree

pyproject.toml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -130,7 +130,7 @@ dependencies = [
130130
# stderr (agronholm/anyio#816, fixed in 4.10).
131131
"anyio>=4.10; python_version >= '3.14'",
132132
"anyio>=4.9; python_version < '3.14'",
133-
"httpx2>=2.5.0",
133+
"httpx2>=2.10.0",
134134
"mcp-types=={{ version }}",
135135
"pydantic>=2.12.0",
136136
"starlette>=0.48.0; python_version >= '3.14'",

src/mcp/client/sse.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@
1212

1313
from mcp.shared._compat import resync_tracer
1414
from mcp.shared._context_streams import create_context_streams
15-
from mcp.shared._httpx_utils import McpHttpClientFactory, create_mcp_http_client
15+
from mcp.shared._httpx_utils import MCP_SSE_MAX_EVENT_SIZE, McpHttpClientFactory, create_mcp_http_client
1616
from mcp.shared.message import SessionMessage
1717

1818
logger = logging.getLogger(__name__)
@@ -55,7 +55,7 @@ async def sse_client(
5555
async with httpx_client_factory(
5656
headers=headers, auth=auth, timeout=httpx2.Timeout(timeout, read=sse_read_timeout)
5757
) as client:
58-
async with client.sse(url) as event_source:
58+
async with client.sse(url, max_event_size=MCP_SSE_MAX_EVENT_SIZE) as event_source:
5959
event_source.response.raise_for_status()
6060
logger.debug("SSE connection established")
6161

src/mcp/client/streamable_http.py

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,7 @@
3333
from mcp.client._transport import TransportStreams
3434
from mcp.shared._compat import resync_tracer
3535
from mcp.shared._context_streams import ContextReceiveStream, ContextSendStream, create_context_streams
36-
from mcp.shared._httpx_utils import create_mcp_http_client
36+
from mcp.shared._httpx_utils import MCP_SSE_MAX_EVENT_SIZE, create_mcp_http_client
3737
from mcp.shared.inbound import MCP_PROTOCOL_VERSION_HEADER
3838
from mcp.shared.jsonrpc_dispatcher import cancelled_request_id_from_params
3939
from mcp.shared.message import ClientMessageMetadata, SessionMessage
@@ -210,7 +210,7 @@ async def handle_get_stream(self, client: httpx2.AsyncClient, read_stream_writer
210210
if last_event_id:
211211
headers[LAST_EVENT_ID] = last_event_id
212212

213-
async with client.sse(self.url, headers=headers) as event_source:
213+
async with client.sse(self.url, headers=headers, max_event_size=MCP_SSE_MAX_EVENT_SIZE) as event_source:
214214
event_source.response.raise_for_status()
215215
logger.debug("GET SSE connection established")
216216

@@ -253,7 +253,7 @@ async def _handle_resumption_request(self, ctx: RequestContext) -> None:
253253
if isinstance(ctx.session_message.message, JSONRPCRequest): # pragma: no branch
254254
original_request_id = ctx.session_message.message.id
255255

256-
async with ctx.client.sse(self.url, headers=headers) as event_source:
256+
async with ctx.client.sse(self.url, headers=headers, max_event_size=MCP_SSE_MAX_EVENT_SIZE) as event_source:
257257
event_source.response.raise_for_status()
258258
logger.debug("Resumption GET SSE connection established")
259259

@@ -423,7 +423,7 @@ async def _handle_sse_response(
423423
original_request_id = ctx.session_message.message.id
424424

425425
try:
426-
event_source = EventSource(response)
426+
event_source = EventSource(response, max_event_size=MCP_SSE_MAX_EVENT_SIZE)
427427
async for sse in event_source: # pragma: no branch
428428
# Track last event ID for potential reconnection
429429
if sse.id:
@@ -501,7 +501,7 @@ async def _handle_reconnection(
501501
headers[LAST_EVENT_ID] = last_event_id
502502

503503
try:
504-
async with ctx.client.sse(self.url, headers=headers) as event_source:
504+
async with ctx.client.sse(self.url, headers=headers, max_event_size=MCP_SSE_MAX_EVENT_SIZE) as event_source:
505505
event_source.response.raise_for_status()
506506
logger.info("Reconnected to SSE stream")
507507

src/mcp/shared/_httpx_utils.py

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,12 +4,17 @@
44

55
import httpx2
66

7-
__all__ = ["create_mcp_http_client", "MCP_DEFAULT_TIMEOUT", "MCP_DEFAULT_SSE_READ_TIMEOUT"]
7+
__all__ = ["create_mcp_http_client", "MCP_DEFAULT_TIMEOUT", "MCP_DEFAULT_SSE_READ_TIMEOUT", "MCP_SSE_MAX_EVENT_SIZE"]
88

99
# Default MCP timeout configuration
1010
MCP_DEFAULT_TIMEOUT = 30.0 # General operations (seconds)
1111
MCP_DEFAULT_SSE_READ_TIMEOUT = 300.0 # SSE streams - 5 minutes (seconds)
1212

13+
# httpx2 >= 2.10 caps a single SSE event at 1 MiB by default. One JSON-RPC message
14+
# is one event and MCP sets no message size limit (the application/json response
15+
# path is unbounded too), so every SSE reader the SDK opens passes this instead.
16+
MCP_SSE_MAX_EVENT_SIZE: int | None = None
17+
1318

1419
class McpHttpClientFactory(Protocol): # pragma: no branch
1520
def __call__( # pragma: no branch

tests/interaction/_requirements.py

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3474,6 +3474,25 @@ def __post_init__(self) -> None:
34743474
transports=("streamable-http",),
34753475
note="Only observable over HTTP: per-request SSE streams are HTTP-specific.",
34763476
),
3477+
"client-transport:http:post-stream-large-event": Requirement(
3478+
source="sdk",
3479+
behavior=(
3480+
"A JSON-RPC message larger than httpx2's default 1 MiB SSE event cap is delivered intact over the "
3481+
"per-request POST stream; the transport lifts the cap because MCP sets no message size limit."
3482+
),
3483+
transports=("streamable-http",),
3484+
note="Only observable over HTTP: SSE event framing is HTTP-specific.",
3485+
),
3486+
"client-transport:http:get-stream-large-event": Requirement(
3487+
source="sdk",
3488+
behavior=(
3489+
"A server-initiated message larger than httpx2's default 1 MiB SSE event cap is delivered intact "
3490+
"over the standalone GET stream; the transport lifts the cap because MCP sets no message size limit."
3491+
),
3492+
transports=("streamable-http",),
3493+
removed_in="2026-07-28",
3494+
note="removed in 2026-07-28 (SEP-2575); the standalone GET stream is replaced by subscriptions/listen.",
3495+
),
34773496
"client-transport:http:custom-client": Requirement(
34783497
source="sdk",
34793498
behavior=(

tests/interaction/transports/test_client_transport_http.py

Lines changed: 72 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,13 +14,23 @@
1414
import mcp_types as types
1515
import pytest
1616
from inline_snapshot import snapshot
17-
from mcp_types import INVALID_REQUEST, CallToolResult, ErrorData, ListToolsResult, TextContent, Tool
17+
from mcp_types import (
18+
INVALID_REQUEST,
19+
CallToolResult,
20+
ErrorData,
21+
ListToolsResult,
22+
LoggingMessageNotification,
23+
LoggingMessageNotificationParams,
24+
TextContent,
25+
Tool,
26+
)
1827
from starlette.types import Receive, Scope, Send
1928

2029
from mcp import MCPError
2130
from mcp.client.client import Client
2231
from mcp.client.streamable_http import streamable_http_client
2332
from mcp.server import Server, ServerRequestContext
33+
from mcp.server.mcpserver import Context, MCPServer
2434
from tests.interaction._connect import BASE_URL, NO_DNS_REBINDING_PROTECTION, client_via_http, mounted_app
2535
from tests.interaction._requirements import requirement
2636
from tests.interaction.transports._bridge import StreamingASGITransport
@@ -158,6 +168,67 @@ async def call(n: int) -> None:
158168
assert len(tools_call_posts) == 3
159169

160170

171+
# One byte past the 1 MiB that httpx2 >= 2.10 allows a single SSE event by default.
172+
_OVERSIZED_TEXT = "x" * (1024 * 1024 + 1)
173+
174+
175+
@requirement("client-transport:http:post-stream-large-event")
176+
async def test_a_post_stream_delivers_a_tool_result_larger_than_one_mebibyte() -> None:
177+
"""A tool result bigger than httpx2's default per-event SSE cap arrives intact over the request's
178+
POST stream. SDK-defined: MCP sets no message size limit, so the transport lifts the cap (#3332)."""
179+
mcp = MCPServer("bulky")
180+
181+
@mcp.tool()
182+
def bulk() -> str:
183+
"""Return more than one SSE event may carry by default."""
184+
return _OVERSIZED_TEXT
185+
186+
async with mounted_app(mcp) as (http, _), client_via_http(http) as client:
187+
with anyio.fail_after(5):
188+
result = await client.call_tool("bulk", {})
189+
190+
assert result.content == [TextContent(text=_OVERSIZED_TEXT)]
191+
192+
193+
@requirement("client-transport:http:get-stream-large-event")
194+
async def test_the_standalone_get_stream_delivers_a_notification_larger_than_one_mebibyte() -> None:
195+
"""A server-initiated notification bigger than httpx2's default per-event SSE cap arrives intact
196+
over the standalone GET stream, which the transport opens with the same lifted cap (#3332)."""
197+
mcp = MCPServer("bulky")
198+
199+
@mcp.tool()
200+
async def shout(ctx: Context) -> str:
201+
"""Emit one unrelated notification, which the server routes to the standalone stream."""
202+
params = LoggingMessageNotificationParams(level="info", data=_OVERSIZED_TEXT)
203+
await ctx.session.send_notification(LoggingMessageNotification(params=params))
204+
return "sent"
205+
206+
get_stream_open = anyio.Event()
207+
208+
async def on_response(response: httpx2.Response) -> None:
209+
if response.request.method == "GET":
210+
get_stream_open.set()
211+
212+
received: list[object] = []
213+
delivered = anyio.Event()
214+
215+
async def collect(params: LoggingMessageNotificationParams) -> None:
216+
received.append(params.data)
217+
delivered.set()
218+
219+
async with (
220+
mounted_app(mcp, on_response=on_response) as (http, _),
221+
client_via_http(http, logging_callback=collect) as client,
222+
):
223+
with anyio.fail_after(5):
224+
# The server drops standalone messages emitted before the GET stream is established.
225+
await get_stream_open.wait()
226+
await client.call_tool("shout", {})
227+
await delivered.wait()
228+
229+
assert received == [_OVERSIZED_TEXT]
230+
231+
161232
@requirement("client-transport:http:sse-405-tolerated")
162233
@requirement("client-transport:http:terminate-405-ok")
163234
async def test_client_tolerates_405_on_get_and_delete() -> None:

tests/shared/test_sse.py

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -224,6 +224,31 @@ async def test_sse_client_exception_handling(
224224
await session.read_resource(uri="xxx://will-not-work")
225225

226226

227+
@pytest.mark.anyio
228+
async def test_sse_client_delivers_a_result_larger_than_one_mebibyte() -> None:
229+
"""A resource read bigger than httpx2's default 1 MiB per-event SSE cap arrives intact. SDK-defined:
230+
MCP sets no message size limit, and the legacy transport carries every server message on one
231+
event stream, so the cap is lifted there too (#3332)."""
232+
oversized = "x" * (1024 * 1024 + 1)
233+
234+
async def read_resource(ctx: ServerRequestContext, params: ReadResourceRequestParams) -> ReadResourceResult:
235+
return ReadResourceResult(
236+
contents=[TextResourceContents(uri=str(params.uri), text=oversized, mime_type="text/plain")]
237+
)
238+
239+
factory = in_process_client_factory(make_app(Server(SERVER_NAME, on_read_resource=read_resource)))
240+
with anyio.fail_after(5):
241+
async with (
242+
sse_client(f"{BASE_URL}/sse", httpx_client_factory=factory) as streams,
243+
ClientSession(*streams) as session,
244+
):
245+
await session.initialize()
246+
response = await session.read_resource(uri="foobar://bulk")
247+
248+
assert isinstance(response.contents[0], TextResourceContents)
249+
assert response.contents[0].text == oversized
250+
251+
227252
@pytest.mark.anyio
228253
async def test_sse_client_basic_connection_mounted_app() -> None:
229254
"""The SSE transport works unchanged when its app is mounted under a sub-path."""

0 commit comments

Comments
 (0)