Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 19 additions & 2 deletions docs/client/transports.md
Original file line number Diff line number Diff line change
Expand Up @@ -46,16 +46,32 @@ environment variables or pass an explicit `verify=ssl_context` to your `httpx2.A
(background in
[`httpx` and `httpx-sse` replaced by `httpx2`](../migration.md#httpx-and-httpx-sse-replaced-by-httpx2)).

### Larger SSE events

Pass `max_sse_event_size` when a server sends a large tool result or notification in one SSE event:

```python title="client.py" hl_lines="6-9"
--8<-- "docs_src/client_transports/tutorial005.py"
```

The default is 1 MiB per event, measured in bytes before the event is parsed. The limit applies to
POST responses, the GET stream, and resumed streams. An oversized event in a POST response or resumed
stream fails that request with an SSE error. On the background GET stream, the client logs
the error and retries the stream. Set `max_sse_event_size=None` to disable the cap when you trust the
server and need larger events. JSON responses are unaffected. If you use `ClientSessionGroup`, set the
same option on `StreamableHttpParameters`.

!!! warning
`streamable_http_client` used to take `headers=` and `timeout=` directly. It does not any more:
its only parameters are `url`, `http_client` and `terminate_on_close`. Reach for `headers=` out
its parameters are `url`, `http_client`, `terminate_on_close`, and `max_sse_event_size`. Reach for `headers=` out
of habit and you get:

```text
TypeError: streamable_http_client() got an unexpected keyword argument 'headers'
```

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

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

* `Client("http://.../mcp")` (a URL) connects over Streamable HTTP, the production transport.
* Headers, auth, proxies and timeouts belong on an `httpx2.AsyncClient` you pass to `streamable_http_client(url, http_client=...)`. There is no `headers=` keyword.
* Use `streamable_http_client(url, max_sse_event_size=...)` to change the byte limit for each SSE event.
* 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.
* stdio is `Client(StdioServerParameters(...))`. Wrap it in `stdio_client(...)` yourself only to redirect the child's stderr.
* The subprocess gets an allow-listed environment, not yours; `env=` adds to it.
Expand Down
2 changes: 1 addition & 1 deletion docs/migration.md
Original file line number Diff line number Diff line change
Expand Up @@ -2104,7 +2104,7 @@ async with http_client:

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.

`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:
`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:

- `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.
- `httpx_client_factory`: gone with no replacement — call your factory yourself and pass the result as `http_client`.
Expand Down
12 changes: 12 additions & 0 deletions docs_src/client_transports/tutorial005.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
from mcp import Client
from mcp.client.streamable_http import streamable_http_client


async def main() -> None:
transport = streamable_http_client(
"http://localhost:8000/mcp",
max_sse_event_size=32 * 1024 * 1024,
)
async with Client(transport) as client:
result = await client.list_tools()
print([tool.name for tool in result.tools])
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -131,7 +131,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",
Comment thread
Kludex marked this conversation as resolved.
"mcp-types=={{ version }}",
"pydantic>=2.12.0",
"starlette>=0.48.0; python_version >= '3.14'",
Expand Down
6 changes: 5 additions & 1 deletion src/mcp/client/session_group.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@
from mcp.client.session import ElicitationFnT, ListRootsFnT, LoggingFnT, MessageHandlerFnT, SamplingFnT
from mcp.client.sse import sse_client
from mcp.client.stdio import StdioServerParameters
from mcp.client.streamable_http import streamable_http_client
from mcp.client.streamable_http import DEFAULT_MAX_SSE_EVENT_SIZE, streamable_http_client
from mcp.shared._httpx_utils import create_mcp_http_client
from mcp.shared.dispatcher import ProgressFnT
from mcp.shared.exceptions import MCPError
Expand Down Expand Up @@ -63,6 +63,9 @@ class StreamableHttpParameters(BaseModel):
# Close the client session when the transport closes.
terminate_on_close: bool = True

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


ServerParameters: TypeAlias = StdioServerParameters | SseServerParameters | StreamableHttpParameters

Expand Down Expand Up @@ -335,6 +338,7 @@ async def _establish_session(
url=server_params.url,
http_client=httpx_client,
terminate_on_close=server_params.terminate_on_close,
max_sse_event_size=server_params.max_sse_event_size,
)
read, write = await session_stack.enter_async_context(client)

Expand Down
65 changes: 50 additions & 15 deletions src/mcp/client/streamable_http.py
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,7 @@
# Reconnection defaults
DEFAULT_RECONNECTION_DELAY_MS = 1000 # 1 second fallback when server doesn't provide retry
MAX_RECONNECTION_ATTEMPTS = 2 # Max retry attempts before giving up
DEFAULT_MAX_SSE_EVENT_SIZE = 1024 * 1024


class StreamableHTTPError(Exception):
Expand Down Expand Up @@ -110,13 +111,17 @@ class _InFlightPost:
class StreamableHTTPTransport:
"""StreamableHTTP client transport implementation."""

def __init__(self, url: str) -> None:
def __init__(self, url: str, *, max_sse_event_size: int | None = DEFAULT_MAX_SSE_EVENT_SIZE) -> None:
"""Initialize the StreamableHTTP transport.
Args:
url: The endpoint URL.
max_sse_event_size: Maximum bytes in one SSE event. None disables the limit.
"""
if max_sse_event_size is not None and max_sse_event_size <= 0:
raise ValueError("max_sse_event_size must be positive or None")
self.url = url
self.max_sse_event_size = max_sse_event_size
self.session_id: str | None = None
# Captured from each stamped message's metadata, synchronously in the
# post_writer loop so the cache always reflects wire order (a POST task's
Expand Down Expand Up @@ -231,7 +236,9 @@ 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 sse_within_origin(client, self.url, headers=headers) as event_source:
async with sse_within_origin(
client, self.url, headers=headers, max_event_size=self.max_sse_event_size
) as event_source:
if (redirect := _unfollowed_redirect(event_source.response)) is not None:
# The same GET would be redirected again, so retrying cannot help.
logger.warning(f"GET stream not opened: {redirect}")
Expand Down Expand Up @@ -278,7 +285,9 @@ 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 sse_within_origin(ctx.client, self.url, headers=headers) as event_source:
async with sse_within_origin(
ctx.client, self.url, headers=headers, max_event_size=self.max_sse_event_size
) as event_source:
if (redirect := _unfollowed_redirect(event_source.response)) is not None:
logger.warning(redirect)
assert original_request_id is not None
Expand All @@ -289,16 +298,22 @@ async def _handle_resumption_request(self, ctx: RequestContext) -> None:
event_source.response.raise_for_status()
logger.debug("Resumption GET SSE connection established")

async for sse in event_source: # pragma: no branch
is_complete = await self._handle_sse_event(
sse,
ctx.read_stream_writer,
original_request_id,
ctx.metadata.on_resumption_token_update if ctx.metadata else None,
try:
async for sse in event_source: # pragma: no branch
is_complete = await self._handle_sse_event(
sse,
ctx.read_stream_writer,
original_request_id,
ctx.metadata.on_resumption_token_update if ctx.metadata else None,
)
if is_complete:
await event_source.response.aclose()
break
except httpx2.SSEError as exc:
assert original_request_id is not None
await self._resolve_abandoned_request(
ctx.read_stream_writer, original_request_id, f"SSE stream failed: {exc}"
)
Comment thread
Kludex marked this conversation as resolved.
if is_complete:
await event_source.response.aclose()
break

def _consume_modern_cancellation(self, session_message: SessionMessage) -> bool:
"""Translate an outbound `notifications/cancelled` at 2026; True means "do not POST".
Expand Down Expand Up @@ -464,7 +479,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=self.max_sse_event_size)
async for sse in event_source: # pragma: no branch
# Track last event ID for potential reconnection
if sse.id:
Expand All @@ -485,6 +500,11 @@ async def _handle_sse_response(
if is_complete:
await response.aclose()
return # Normal completion, no reconnect needed
except httpx2.SSEError as exc:
await self._resolve_abandoned_request(
ctx.read_stream_writer, original_request_id, f"SSE stream failed: {exc}"
)
return
except Exception:
logger.debug("SSE stream ended", exc_info=True) # pragma: lax no cover

Expand Down Expand Up @@ -541,9 +561,14 @@ async def _handle_reconnection(
headers = self._prepare_headers()
headers[LAST_EVENT_ID] = last_event_id

is_sse_response = False
try:
async with sse_within_origin(ctx.client, self.url, headers=headers) as event_source:
async with sse_within_origin(
ctx.client, self.url, headers=headers, max_event_size=self.max_sse_event_size
) as event_source:
event_source.response.raise_for_status()
content_type = event_source.response.headers.get("content-type", "").partition(";")[0]
is_sse_response = content_type.strip().lower() == "text/event-stream"
logger.info("Reconnected to SSE stream")

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

transport = StreamableHTTPTransport(url)
transport = StreamableHTTPTransport(url, max_sse_event_size=max_sse_event_size)

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

Expand Down
8 changes: 6 additions & 2 deletions src/mcp/shared/_httpx_utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -156,13 +156,17 @@ async def request_within_origin(

@asynccontextmanager
async def sse_within_origin(
client: httpx2.AsyncClient, url: httpx2.URL | str, *, headers: dict[str, str] | None = None
client: httpx2.AsyncClient,
url: httpx2.URL | str,
*,
headers: dict[str, str] | None = None,
max_event_size: int | None = 1024 * 1024,
Comment thread
Kludex marked this conversation as resolved.
) -> AsyncGenerator[httpx2.EventSource]:
"""`client.sse(url)` with the redirect handling of `stream_within_origin`."""
merged = httpx2.Headers(_SSE_HEADERS)
merged.update(headers or {})
async with stream_within_origin(client, "GET", url, headers=merged) as response:
yield httpx2.EventSource(response)
yield httpx2.EventSource(response, max_event_size=max_event_size)


def redirect_location(response: httpx2.Response) -> httpx2.URL | None:
Expand Down
5 changes: 4 additions & 1 deletion tests/client/test_session_group.py
Original file line number Diff line number Diff line change
Expand Up @@ -311,7 +311,9 @@ async def test_client_session_group_disconnect_non_existent_server():
"mcp.client.session_group.sse_client",
), # url, headers, timeout, sse_read_timeout
(
StreamableHttpParameters(url="http://test.com/stream", terminate_on_close=False),
StreamableHttpParameters(
url="http://test.com/stream", terminate_on_close=False, max_sse_event_size=32 * 1024 * 1024
),
"streamablehttp",
"mcp.client.session_group.streamable_http_client",
), # url, headers, timeout, sse_read_timeout, terminate_on_close
Expand Down Expand Up @@ -380,6 +382,7 @@ async def test_client_session_group_establish_session_parameterized(
call_args = mock_specific_client_func.call_args
assert call_args.kwargs["url"] == server_params_instance.url
assert call_args.kwargs["terminate_on_close"] == server_params_instance.terminate_on_close
assert call_args.kwargs["max_sse_event_size"] == server_params_instance.max_sse_event_size
assert isinstance(call_args.kwargs["http_client"], httpx2.AsyncClient)

mock_client_cm_instance.__aenter__.assert_awaited_once()
Expand Down
Loading
Loading