Skip to content

Commit 7fdc944

Browse files
authored
Expose an SSE event size limit in Streamable HTTP clients (#3600)
1 parent 772ecf2 commit 7fdc944

11 files changed

Lines changed: 445 additions & 54 deletions

File tree

‎docs/client/transports.md‎

Lines changed: 19 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -46,16 +46,32 @@ environment variables or pass an explicit `verify=ssl_context` to your `httpx2.A
4646
(background in
4747
[`httpx` and `httpx-sse` replaced by `httpx2`](../migration.md#httpx-and-httpx-sse-replaced-by-httpx2)).
4848

49+
### Larger SSE events
50+
51+
Pass `max_sse_event_size` when a server sends a large tool result or notification in one SSE event:
52+
53+
```python title="client.py" hl_lines="6-9"
54+
--8<-- "docs_src/client_transports/tutorial005.py"
55+
```
56+
57+
The default is 1 MiB per event, measured in bytes before the event is parsed. The limit applies to
58+
POST responses, the GET stream, and resumed streams. An oversized event in a POST response or resumed
59+
stream fails that request with an SSE error. On the background GET stream, the client logs
60+
the error and retries the stream. Set `max_sse_event_size=None` to disable the cap when you trust the
61+
server and need larger events. JSON responses are unaffected. If you use `ClientSessionGroup`, set the
62+
same option on `StreamableHttpParameters`.
63+
4964
!!! warning
5065
`streamable_http_client` used to take `headers=` and `timeout=` directly. It does not any more:
51-
its only parameters are `url`, `http_client` and `terminate_on_close`. Reach for `headers=` out
66+
its parameters are `url`, `http_client`, `terminate_on_close`, and `max_sse_event_size`. Reach for `headers=` out
5267
of habit and you get:
5368

5469
```text
5570
TypeError: streamable_http_client() got an unexpected keyword argument 'headers'
5671
```
5772

58-
Everything HTTP-shaped now lives on the one `httpx2.AsyncClient` you pass in.
73+
Headers, authentication, proxies, and timeouts live on the one `httpx2.AsyncClient` you pass in.
74+
`max_sse_event_size` applies to the MCP transport's SSE readers instead.
5975

6076
!!! info
6177
`httpx2` keeps the familiar `httpx` API, so if you know `httpx` you already know how to do auth,
@@ -132,6 +148,7 @@ A **transport** is any async context manager that yields a `(read, write)` pair
132148

133149
* `Client("http://.../mcp")` (a URL) connects over Streamable HTTP, the production transport.
134150
* Headers, auth, proxies and timeouts belong on an `httpx2.AsyncClient` you pass to `streamable_http_client(url, http_client=...)`. There is no `headers=` keyword.
151+
* Use `streamable_http_client(url, max_sse_event_size=...)` to change the byte limit for each SSE event.
135152
* Redirects are followed only within the URL's own origin (a trailing-slash `307`/`308`), plus `http`→`https` on the same host. Anything else fails with `Redirect to … not followed`; configure the final URL.
136153
* stdio is `Client(StdioServerParameters(...))`. Wrap it in `stdio_client(...)` yourself only to redirect the child's stderr.
137154
* The subprocess gets an allow-listed environment, not yours; `env=` adds to it.

‎docs/migration.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2104,7 +2104,7 @@ async with http_client:
21042104

21052105
v1's internal client set `follow_redirects=True`. You don't need it on your own client: the transport follows a method-preserving redirect within the endpoint's origin (a trailing-slash 307/308, say) itself, and does not follow one anywhere else, whatever the client is configured to do.
21062106

2107-
`streamable_http_client` itself keeps a small signature — `streamable_http_client(url, *, http_client=None, terminate_on_close=True)` — and now yields a 2-tuple (next section). The removed function's other parameters map onto the client you build:
2107+
`streamable_http_client` itself keeps a small signature — `streamable_http_client(url, *, http_client=None, terminate_on_close=True, max_sse_event_size=1024 * 1024)` — and now yields a 2-tuple (next section). The removed function's other parameters map onto the client you build:
21082108

21092109
- `headers`, `timeout`, `sse_read_timeout`, `auth`: set them on the `httpx2.AsyncClient` as above. `streamablehttp_client` defaulted to `httpx.Timeout(30, read=300)`; a bare `httpx2.AsyncClient()` falls back to httpx2's flat 5-second timeout, too short for the long-lived GET stream, so set `timeout=httpx2.Timeout(30, read=300)` (as shown) to keep v1's values. Omitting `http_client` still gives you a default client with those timeouts.
21102110
- `httpx_client_factory`: gone with no replacement — call your factory yourself and pass the result as `http_client`.
Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,12 @@
1+
from mcp import Client
2+
from mcp.client.streamable_http import streamable_http_client
3+
4+
5+
async def main() -> None:
6+
transport = streamable_http_client(
7+
"http://localhost:8000/mcp",
8+
max_sse_event_size=32 * 1024 * 1024,
9+
)
10+
async with Client(transport) as client:
11+
result = await client.list_tools()
12+
print([tool.name for tool in result.tools])

‎pyproject.toml‎

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

‎src/mcp/client/session_group.py‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@
2323
from mcp.client.session import ElicitationFnT, ListRootsFnT, LoggingFnT, MessageHandlerFnT, SamplingFnT
2424
from mcp.client.sse import sse_client
2525
from mcp.client.stdio import StdioServerParameters
26-
from mcp.client.streamable_http import streamable_http_client
26+
from mcp.client.streamable_http import DEFAULT_MAX_SSE_EVENT_SIZE, streamable_http_client
2727
from mcp.shared._httpx_utils import create_mcp_http_client
2828
from mcp.shared.dispatcher import ProgressFnT
2929
from mcp.shared.exceptions import MCPError
@@ -63,6 +63,9 @@ class StreamableHttpParameters(BaseModel):
6363
# Close the client session when the transport closes.
6464
terminate_on_close: bool = True
6565

66+
# Maximum bytes in one server-sent event. None disables the limit.
67+
max_sse_event_size: int | None = Field(default=DEFAULT_MAX_SSE_EVENT_SIZE, gt=0)
68+
6669

6770
ServerParameters: TypeAlias = StdioServerParameters | SseServerParameters | StreamableHttpParameters
6871

@@ -335,6 +338,7 @@ async def _establish_session(
335338
url=server_params.url,
336339
http_client=httpx_client,
337340
terminate_on_close=server_params.terminate_on_close,
341+
max_sse_event_size=server_params.max_sse_event_size,
338342
)
339343
read, write = await session_stack.enter_async_context(client)
340344

‎src/mcp/client/streamable_http.py‎

Lines changed: 50 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -58,6 +58,7 @@
5858
# Reconnection defaults
5959
DEFAULT_RECONNECTION_DELAY_MS = 1000 # 1 second fallback when server doesn't provide retry
6060
MAX_RECONNECTION_ATTEMPTS = 2 # Max retry attempts before giving up
61+
DEFAULT_MAX_SSE_EVENT_SIZE = 1024 * 1024
6162

6263

6364
class StreamableHTTPError(Exception):
@@ -110,13 +111,17 @@ class _InFlightPost:
110111
class StreamableHTTPTransport:
111112
"""StreamableHTTP client transport implementation."""
112113

113-
def __init__(self, url: str) -> None:
114+
def __init__(self, url: str, *, max_sse_event_size: int | None = DEFAULT_MAX_SSE_EVENT_SIZE) -> None:
114115
"""Initialize the StreamableHTTP transport.
115116
116117
Args:
117118
url: The endpoint URL.
119+
max_sse_event_size: Maximum bytes in one SSE event. None disables the limit.
118120
"""
121+
if max_sse_event_size is not None and max_sse_event_size <= 0:
122+
raise ValueError("max_sse_event_size must be positive or None")
119123
self.url = url
124+
self.max_sse_event_size = max_sse_event_size
120125
self.session_id: str | None = None
121126
# Captured from each stamped message's metadata, synchronously in the
122127
# post_writer loop so the cache always reflects wire order (a POST task's
@@ -231,7 +236,9 @@ async def handle_get_stream(self, client: httpx2.AsyncClient, read_stream_writer
231236
if last_event_id:
232237
headers[LAST_EVENT_ID] = last_event_id
233238

234-
async with sse_within_origin(client, self.url, headers=headers) as event_source:
239+
async with sse_within_origin(
240+
client, self.url, headers=headers, max_event_size=self.max_sse_event_size
241+
) as event_source:
235242
if (redirect := _unfollowed_redirect(event_source.response)) is not None:
236243
# The same GET would be redirected again, so retrying cannot help.
237244
logger.warning(f"GET stream not opened: {redirect}")
@@ -278,7 +285,9 @@ async def _handle_resumption_request(self, ctx: RequestContext) -> None:
278285
if isinstance(ctx.session_message.message, JSONRPCRequest): # pragma: no branch
279286
original_request_id = ctx.session_message.message.id
280287

281-
async with sse_within_origin(ctx.client, self.url, headers=headers) as event_source:
288+
async with sse_within_origin(
289+
ctx.client, self.url, headers=headers, max_event_size=self.max_sse_event_size
290+
) as event_source:
282291
if (redirect := _unfollowed_redirect(event_source.response)) is not None:
283292
logger.warning(redirect)
284293
assert original_request_id is not None
@@ -289,16 +298,22 @@ async def _handle_resumption_request(self, ctx: RequestContext) -> None:
289298
event_source.response.raise_for_status()
290299
logger.debug("Resumption GET SSE connection established")
291300

292-
async for sse in event_source: # pragma: no branch
293-
is_complete = await self._handle_sse_event(
294-
sse,
295-
ctx.read_stream_writer,
296-
original_request_id,
297-
ctx.metadata.on_resumption_token_update if ctx.metadata else None,
301+
try:
302+
async for sse in event_source: # pragma: no branch
303+
is_complete = await self._handle_sse_event(
304+
sse,
305+
ctx.read_stream_writer,
306+
original_request_id,
307+
ctx.metadata.on_resumption_token_update if ctx.metadata else None,
308+
)
309+
if is_complete:
310+
await event_source.response.aclose()
311+
break
312+
except httpx2.SSEError as exc:
313+
assert original_request_id is not None
314+
await self._resolve_abandoned_request(
315+
ctx.read_stream_writer, original_request_id, f"SSE stream failed: {exc}"
298316
)
299-
if is_complete:
300-
await event_source.response.aclose()
301-
break
302317

303318
def _consume_modern_cancellation(self, session_message: SessionMessage) -> bool:
304319
"""Translate an outbound `notifications/cancelled` at 2026; True means "do not POST".
@@ -464,7 +479,7 @@ async def _handle_sse_response(
464479
original_request_id = ctx.session_message.message.id
465480

466481
try:
467-
event_source = EventSource(response)
482+
event_source = EventSource(response, max_event_size=self.max_sse_event_size)
468483
async for sse in event_source: # pragma: no branch
469484
# Track last event ID for potential reconnection
470485
if sse.id:
@@ -485,6 +500,11 @@ async def _handle_sse_response(
485500
if is_complete:
486501
await response.aclose()
487502
return # Normal completion, no reconnect needed
503+
except httpx2.SSEError as exc:
504+
await self._resolve_abandoned_request(
505+
ctx.read_stream_writer, original_request_id, f"SSE stream failed: {exc}"
506+
)
507+
return
488508
except Exception:
489509
logger.debug("SSE stream ended", exc_info=True) # pragma: lax no cover
490510

@@ -541,9 +561,14 @@ async def _handle_reconnection(
541561
headers = self._prepare_headers()
542562
headers[LAST_EVENT_ID] = last_event_id
543563

564+
is_sse_response = False
544565
try:
545-
async with sse_within_origin(ctx.client, self.url, headers=headers) as event_source:
566+
async with sse_within_origin(
567+
ctx.client, self.url, headers=headers, max_event_size=self.max_sse_event_size
568+
) as event_source:
546569
event_source.response.raise_for_status()
570+
content_type = event_source.response.headers.get("content-type", "").partition(";")[0]
571+
is_sse_response = content_type.strip().lower() == "text/event-stream"
547572
logger.info("Reconnected to SSE stream")
548573

549574
# Track for potential further reconnection
@@ -569,6 +594,13 @@ async def _handle_reconnection(
569594
# Stream ended again without response - reconnect again (reset attempt counter)
570595
logger.info("SSE stream disconnected, reconnecting...")
571596
await self._handle_reconnection(ctx, reconnect_last_event_id, reconnect_retry_ms, 0)
597+
except httpx2.SSEError as exc:
598+
if is_sse_response:
599+
await self._resolve_abandoned_request(
600+
ctx.read_stream_writer, original_request_id, f"SSE stream failed: {exc}"
601+
)
602+
else:
603+
await self._handle_reconnection(ctx, last_event_id, retry_interval_ms, attempt + 1)
572604
except Exception as e: # pragma: no cover
573605
logger.debug(f"Reconnection failed: {e}")
574606
# Try to reconnect again if we still have an event ID
@@ -683,6 +715,7 @@ async def streamable_http_client(
683715
*,
684716
http_client: httpx2.AsyncClient | None = None,
685717
terminate_on_close: bool = True,
718+
max_sse_event_size: int | None = DEFAULT_MAX_SSE_EVENT_SIZE,
686719
) -> AsyncGenerator[TransportStreams, None]:
687720
"""Client transport for StreamableHTTP.
688721
@@ -699,6 +732,8 @@ async def streamable_http_client(
699732
client's `follow_redirects` setting is not consulted; the SDK's OAuth providers apply the
700733
same rule to the requests they make.
701734
terminate_on_close: If True, send a DELETE request to terminate the session when the context exits.
735+
max_sse_event_size: Maximum bytes buffered for one SSE event. None disables the limit.
736+
JSON responses are not affected.
702737
703738
Yields:
704739
Tuple containing:
@@ -716,7 +751,7 @@ async def streamable_http_client(
716751
# Create default client with recommended MCP timeouts
717752
client = create_mcp_http_client()
718753

719-
transport = StreamableHTTPTransport(url)
754+
transport = StreamableHTTPTransport(url, max_sse_event_size=max_sse_event_size)
720755

721756
logger.debug(f"Connecting to StreamableHTTP endpoint: {url}")
722757

‎src/mcp/shared/_httpx_utils.py‎

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -156,13 +156,17 @@ async def request_within_origin(
156156

157157
@asynccontextmanager
158158
async def sse_within_origin(
159-
client: httpx2.AsyncClient, url: httpx2.URL | str, *, headers: dict[str, str] | None = None
159+
client: httpx2.AsyncClient,
160+
url: httpx2.URL | str,
161+
*,
162+
headers: dict[str, str] | None = None,
163+
max_event_size: int | None = 1024 * 1024,
160164
) -> AsyncGenerator[httpx2.EventSource]:
161165
"""`client.sse(url)` with the redirect handling of `stream_within_origin`."""
162166
merged = httpx2.Headers(_SSE_HEADERS)
163167
merged.update(headers or {})
164168
async with stream_within_origin(client, "GET", url, headers=merged) as response:
165-
yield httpx2.EventSource(response)
169+
yield httpx2.EventSource(response, max_event_size=max_event_size)
166170

167171

168172
def redirect_location(response: httpx2.Response) -> httpx2.URL | None:

‎tests/client/test_session_group.py‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -311,7 +311,9 @@ async def test_client_session_group_disconnect_non_existent_server():
311311
"mcp.client.session_group.sse_client",
312312
), # url, headers, timeout, sse_read_timeout
313313
(
314-
StreamableHttpParameters(url="http://test.com/stream", terminate_on_close=False),
314+
StreamableHttpParameters(
315+
url="http://test.com/stream", terminate_on_close=False, max_sse_event_size=32 * 1024 * 1024
316+
),
315317
"streamablehttp",
316318
"mcp.client.session_group.streamable_http_client",
317319
), # url, headers, timeout, sse_read_timeout, terminate_on_close
@@ -380,6 +382,7 @@ async def test_client_session_group_establish_session_parameterized(
380382
call_args = mock_specific_client_func.call_args
381383
assert call_args.kwargs["url"] == server_params_instance.url
382384
assert call_args.kwargs["terminate_on_close"] == server_params_instance.terminate_on_close
385+
assert call_args.kwargs["max_sse_event_size"] == server_params_instance.max_sse_event_size
383386
assert isinstance(call_args.kwargs["http_client"], httpx2.AsyncClient)
384387

385388
mock_client_cm_instance.__aenter__.assert_awaited_once()

0 commit comments

Comments
 (0)