Skip to content

Commit b2fe38d

Browse files
committed
Expose an SSE event size limit in Streamable HTTP clients
1 parent 06d1d1e commit b2fe38d

11 files changed

Lines changed: 381 additions & 31 deletions

File tree

‎docs/client/transports.md‎

Lines changed: 18 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -46,16 +46,31 @@ 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 16 MiB per event, measured in bytes before the event is parsed. The limit applies to
58+
POST responses, the GET stream, and resumed streams. If an event exceeds it, the request fails with
59+
an error naming the limit. Set `max_sse_event_size=None` to disable the cap when you trust the server
60+
and need larger events. JSON responses are unaffected. If you use `ClientSessionGroup`, set the same
61+
option on `StreamableHttpParameters`.
62+
4963
!!! warning
5064
`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
65+
its parameters are `url`, `http_client`, `terminate_on_close`, and `max_sse_event_size`. Reach for `headers=` out
5266
of habit and you get:
5367

5468
```text
5569
TypeError: streamable_http_client() got an unexpected keyword argument 'headers'
5670
```
5771

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

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

133148
* `Client("http://.../mcp")` (a URL) connects over Streamable HTTP, the production transport.
134149
* Headers, auth, proxies and timeouts belong on an `httpx2.AsyncClient` you pass to `streamable_http_client(url, http_client=...)`. There is no `headers=` keyword.
150+
* Use `streamable_http_client(url, max_sse_event_size=...)` to change the byte limit for each SSE event.
135151
* 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.
136152
* stdio is `Client(StdioServerParameters(...))`. Wrap it in `stdio_client(...)` yourself only to redirect the child's stderr.
137153
* 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=16 * 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: 47 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 = 16 * 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}")
@@ -252,6 +259,9 @@ async def handle_get_stream(self, client: httpx2.AsyncClient, read_stream_writer
252259
# Stream ended normally (server closed) - reset attempt counter
253260
attempt = 0
254261

262+
except httpx2.SSEError:
263+
logger.exception("GET SSE stream failed")
264+
return
255265
except Exception:
256266
logger.debug("GET stream error", exc_info=True)
257267
attempt += 1
@@ -278,7 +288,9 @@ async def _handle_resumption_request(self, ctx: RequestContext) -> None:
278288
if isinstance(ctx.session_message.message, JSONRPCRequest): # pragma: no branch
279289
original_request_id = ctx.session_message.message.id
280290

281-
async with sse_within_origin(ctx.client, self.url, headers=headers) as event_source:
291+
async with sse_within_origin(
292+
ctx.client, self.url, headers=headers, max_event_size=self.max_sse_event_size
293+
) as event_source:
282294
if (redirect := _unfollowed_redirect(event_source.response)) is not None:
283295
logger.warning(redirect)
284296
assert original_request_id is not None
@@ -289,16 +301,22 @@ async def _handle_resumption_request(self, ctx: RequestContext) -> None:
289301
event_source.response.raise_for_status()
290302
logger.debug("Resumption GET SSE connection established")
291303

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,
304+
try:
305+
async for sse in event_source: # pragma: no branch
306+
is_complete = await self._handle_sse_event(
307+
sse,
308+
ctx.read_stream_writer,
309+
original_request_id,
310+
ctx.metadata.on_resumption_token_update if ctx.metadata else None,
311+
)
312+
if is_complete:
313+
await event_source.response.aclose()
314+
break
315+
except httpx2.SSEError as exc:
316+
assert original_request_id is not None
317+
await self._resolve_abandoned_request(
318+
ctx.read_stream_writer, original_request_id, f"SSE stream failed: {exc}"
298319
)
299-
if is_complete:
300-
await event_source.response.aclose()
301-
break
302320

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

466484
try:
467-
event_source = EventSource(response)
485+
event_source = EventSource(response, max_event_size=self.max_sse_event_size)
468486
async for sse in event_source: # pragma: no branch
469487
# Track last event ID for potential reconnection
470488
if sse.id:
@@ -485,6 +503,11 @@ async def _handle_sse_response(
485503
if is_complete:
486504
await response.aclose()
487505
return # Normal completion, no reconnect needed
506+
except httpx2.SSEError as exc:
507+
await self._resolve_abandoned_request(
508+
ctx.read_stream_writer, original_request_id, f"SSE stream failed: {exc}"
509+
)
510+
return
488511
except Exception:
489512
logger.debug("SSE stream ended", exc_info=True) # pragma: lax no cover
490513

@@ -542,7 +565,9 @@ async def _handle_reconnection(
542565
headers[LAST_EVENT_ID] = last_event_id
543566

544567
try:
545-
async with sse_within_origin(ctx.client, self.url, headers=headers) as event_source:
568+
async with sse_within_origin(
569+
ctx.client, self.url, headers=headers, max_event_size=self.max_sse_event_size
570+
) as event_source:
546571
event_source.response.raise_for_status()
547572
logger.info("Reconnected to SSE stream")
548573

@@ -569,6 +594,10 @@ 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+
await self._resolve_abandoned_request(
599+
ctx.read_stream_writer, original_request_id, f"SSE stream failed: {exc}"
600+
)
572601
except Exception as e: # pragma: no cover
573602
logger.debug(f"Reconnection failed: {e}")
574603
# Try to reconnect again if we still have an event ID
@@ -683,6 +712,7 @@ async def streamable_http_client(
683712
*,
684713
http_client: httpx2.AsyncClient | None = None,
685714
terminate_on_close: bool = True,
715+
max_sse_event_size: int | None = DEFAULT_MAX_SSE_EVENT_SIZE,
686716
) -> AsyncGenerator[TransportStreams, None]:
687717
"""Client transport for StreamableHTTP.
688718
@@ -699,6 +729,8 @@ async def streamable_http_client(
699729
client's `follow_redirects` setting is not consulted; the SDK's OAuth providers apply the
700730
same rule to the requests they make.
701731
terminate_on_close: If True, send a DELETE request to terminate the session when the context exits.
732+
max_sse_event_size: Maximum bytes buffered for one SSE event. None disables the limit.
733+
JSON responses are not affected.
702734
703735
Yields:
704736
Tuple containing:
@@ -716,7 +748,7 @@ async def streamable_http_client(
716748
# Create default client with recommended MCP timeouts
717749
client = create_mcp_http_client()
718750

719-
transport = StreamableHTTPTransport(url)
751+
transport = StreamableHTTPTransport(url, max_sse_event_size=max_sse_event_size)
720752

721753
logger.debug(f"Connecting to StreamableHTTP endpoint: {url}")
722754

‎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)