Skip to content
18 changes: 16 additions & 2 deletions src/mcp/client/sse.py
Original file line number Diff line number Diff line change
Expand Up @@ -118,19 +118,33 @@
async def post_writer(endpoint_url: str):
try:
async with write_stream_reader, write_stream:

async def _send_message(session_message: SessionMessage) -> None:
logger.debug(f"Sending client message: {session_message}")
message = session_message.message
response = await client.post(
endpoint_url,
json=session_message.message.model_dump(
json=message.model_dump(
by_alias=True,
mode="json",
exclude_unset=True,
),
)
response.raise_for_status()
if response.status_code >= 400:
# Resolve the waiting caller with an error correlated to its
# request id, mirroring the streamable-HTTP transport: raising
# here would be swallowed by the post_writer handler below and
# the caller would hang forever (#2110). A notification has no
# waiter to resolve, so the failure is only logged.
logger.error(f"Message POST returned HTTP status {response.status_code}")
Comment on lines +150 to +153

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Under the default client factory (create_mcp_http_client sets follow_redirects=True), a Location-bearing 301/302/303 on the SSE message POST is followed by httpx2, which rewrites the POST to a bodyless GET — so the JSON-RPC message is silently dropped, and a 2xx at the redirect target (e.g. an SSO login page) makes response.is_success pass, logging 'sent successfully' while the waiting caller hangs forever: the residual #2110 hang the is_success widening does not close, since the PR's 302 tests use Location-less responses that httpx cannot follow. Consider treating a method-rewriting redirect as a delivery failure (e.g. response.history non-empty and final request method != POST) and resolving the waiter via the same correlated path — method-preserving 307/308 keep working.

Extended reasoning...

What the bug is. The new success check at src/mcp/client/sse.py:150 — if response.is_success: → return — only ever sees the final response of httpx2's redirect-following. The default factory create_mcp_http_client (src/mcp/shared/_httpx_utils.py:79) hardcodes follow_redirects=True (its docstring: "Always enables follow_redirects"), and sse_client uses it by default. httpx2's redirect handling rewrites POST to GET on 301/302/303 (_redirect_method, browser semantics) and drops the request body when the method changes (_redirect_stream returns no stream). So a Location-bearing 3xx on the message POST silently discards the JSON-RPC message, GETs the redirect target with no body, and if that target answers any 2xx, is_success is True: the transport logs "Client message sent successfully" at debug level and returns without resolving the waiter.\n\nWhy this hangs the caller. On this transport, unlike streamable HTTP, real responses only ever arrive on the SSE stream — the POST response body is never inspected. With the message never delivered, nothing will arrive on the SSE stream for that request id, and _send_message reported success, so the correlated-error path this PR built is never taken. The caller of session.call_tool() hangs until its own timeout — the exact #2110 symptom, indistinguishable from a slow server, with no log above debug.\n\nWhy the PR's hardening doesn't cover it. The is_success widening and the [302] parametrization of test_sse_client_request_post_http_error_reaches_caller_and_session_survives cover only unfollowed redirects: make_app_rejecting_posts returns Response(status_code=302) with no Location header, which httpx cannot follow (has_redirect_location is False), so the 3xx reaches the is_success check — the test docstring itself says "An unfollowed redirect counts." With a Location present (the realistic production shape), the 3xx never reaches the check; only the redirect target's 2xx does. The in-process test factory explicitly enables follow_redirects=True "to match create_mcp_http_client," confirming the followed case is simply untested.\n\nWhy only this leg. The sibling paths self-heal in the same scenario: on the streamable message POST, a followed redirect to a 200 HTML page falls into _handle_post_request's unexpected-content-type branch and resolves the waiter; on the resumption GET and SSE initial GET, a redirect landing on non-text/event-stream content makes httpx2's EventSource raise SSEError, which the new containment resolves. The SSE message POST is the one leg that never inspects the response, so a followed-redirect-to-2xx uniquely reads as success.\n\nStep-by-step proof. (1) A client connects through a corporate gateway; the SSE GET stream is established while auth is valid. (2) Auth expires mid-session; the gateway answers the next message POST with 302 Location: https://sso.example/login. (3) httpx2 follows: _redirect_method rewrites POST→GET, _redirect_stream drops the JSON-RPC body, and the login page returns 200 text/html. (4) response.is_success is True; _send_message logs "sent successfully" (debug) and returns. (5) The waiter is never resolved; session.call_tool() hangs into its timeout, and every subsequent request repeats the cycle — no MCPError, no visible log.\n\nSeverity and fix. This is behaviorally pre-existing — v1's raise_for_status() on the followed final 200 passed identically, so merging this PR causes no new failure, and the PR strictly improves the path (the migration table's non-2xx rows aren't literally violated: a followed redirect terminates in a 2xx). That's why this is a nit, not blocking. But since the PR rewrites this exact line, reasons explicitly about redirects on this path, and documents/tests a redirects-resolve-the-caller contract, it's worth closing here or in a follow-up: after the POST, treat a method-rewriting redirect as a delivery failure — e.g. if response.history and response.request.method != "POST": resolve the waiter via the same status_error_data/correlated path (using the first redirect's status, or a generic delivery-failure error). Method-preserving 307/308 redirects keep the POST and body, so they genuinely deliver and continue to work. All three verifiers confirmed this chain end-to-end; none refuted it.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Real residual case — but behaviorally pre-existing per your own analysis: v1 passed identically on the followed 200, so this PR neither introduces nor worsens it. Treating method-rewriting redirects (301/302/303 turning the POST into a GET) as delivery failures is a deliberate behavior change, and it belongs in the grouped follow-up rather than another round here. Adding it to that follow-up's list alongside the reconnection-loop raise_for_status(), the uncontained POST-path reply sends, and the resumption/reconnection GET header drops. Declining for this PR.


Generated by Claude Code

if isinstance(message, types.JSONRPCRequest):
error_data = types.ErrorData(
code=types.INTERNAL_ERROR, message="Server returned an error response"
)
reply = types.JSONRPCError(jsonrpc="2.0", id=message.id, error=error_data)
await read_stream_writer.send(SessionMessage(reply))
return
logger.debug(f"Client message sent successfully: {response.status_code}")

Check warning on line 147 in src/mcp/client/sse.py

View check run for this annotation

Claude / Claude Code Review

3xx responses now treated as success: raise_for_status covered all non-2xx, replacement checks only >= 400

The old `raise_for_status()` raised on ANY non-2xx (httpx treats 1xx and 3xx as errors too), but the replacement `status_code >= 400` check silently reclassifies 3xx as success — on this SSE POST path a stray redirect is logged as "sent successfully" at debug level while the waiting caller hangs, and the same gap exists in the resumption GET check in `src/mcp/client/streamable_http.py:257`. Widening both checks to `if not response.is_success:` would match the semantics being replaced (and the "n
Comment thread
claude[bot] marked this conversation as resolved.
Outdated

async for session_message in write_stream_reader:

Check notice on line 149 in src/mcp/client/sse.py

View check run for this annotation

Claude / Claude Code Review

SSE message POST network errors still hang the caller forever

Pre-existing issue, deliberately scoped out by the PR body — but the parity rationale for the scope-out doesn't hold on this transport: a network-level error (`ConnectError`, `ReadError`, `ReadTimeout`) from `await client.post(...)` still lands in `post_writer`'s catch-all `except Exception`, so the waiting caller hangs forever — the exact #2110 symptom this PR's title addresses, left in place for what is in practice the more common failure class than a non-2xx. Consider catching `httpx2.HTTPErr
Comment thread
claude[bot] marked this conversation as resolved.
sender_ctx = write_stream_reader.last_context
Expand Down
17 changes: 12 additions & 5 deletions src/mcp/client/streamable_http.py
Original file line number Diff line number Diff line change
Expand Up @@ -248,13 +248,20 @@
else:
raise ResumptionError("Resumption request requires a resumption token") # pragma: no cover

# Extract original request ID to map responses
original_request_id = None
if isinstance(ctx.session_message.message, JSONRPCRequest): # pragma: no branch
original_request_id = ctx.session_message.message.id
# Only requests resume: a resumption token is only ever attached by a
# request's metadata, so the original id is always available to map responses.
assert isinstance(ctx.session_message.message, JSONRPCRequest)
original_request_id = ctx.session_message.message.id

Check warning on line 254 in src/mcp/client/streamable_http.py

View check run for this annotation

Claude / Claude Code Review

New assert in _handle_resumption_request crashes on a notification carrying a resumption token

The new `assert isinstance(ctx.session_message.message, JSONRPCRequest)` can crash `post_writer` for a notification carrying a resumption token: the resumption path is selected on metadata alone (`is_resumption = bool(metadata and metadata.resumption_token)`), not message type, so the invariant isn't structurally enforced here the way it is for `_handle_sse_response`/`_handle_reconnection`. Consider guarding the dispatch instead, e.g. `is_resumption = bool(metadata and metadata.resumption_token)
Comment thread
claude[bot] marked this conversation as resolved.
Outdated

async with ctx.client.sse(self.url, headers=headers) as event_source:
event_source.response.raise_for_status()
if event_source.response.status_code >= 400:
# Resolve the waiting caller with an error correlated to its request,
# mirroring `_handle_post_request`: an escaping `HTTPStatusError` would
# tear down the transport's task group and every stream with it (#2110).
error_data = ErrorData(code=INTERNAL_ERROR, message="Server returned an error response")
reply = JSONRPCError(jsonrpc="2.0", id=original_request_id, error=error_data)
await ctx.read_stream_writer.send(SessionMessage(reply))
return

Check warning on line 264 in src/mcp/client/streamable_http.py

View check run for this annotation

Claude / Claude Code Review

[quality] Resumption-GET error branch re-implements _resolve_abandoned_request, losing its closed-stream guard

The new >=400 branch in `_handle_resumption_request` hand-builds the `ErrorData` + `JSONRPCError` shape that `_resolve_abandoned_request` already produces, and the raw `ctx.read_stream_writer.send()` lacks that helper's `BrokenResourceError`/`ClosedResourceError` containment for the teardown race. Replacing the three lines with `await self._resolve_abandoned_request(ctx.read_stream_writer, original_request_id, "Server returned an error response", code=INTERNAL_ERROR)` is wire-identical and stric

Check notice on line 264 in src/mcp/client/streamable_http.py

View check run for this annotation

Claude / Claude Code Review

Resumption GET stream that dies mid-read still tears down the transport

Pre-existing issue (the loop itself is unchanged by this PR): the `async for sse in event_source:` loop that follows the new status check in `_handle_resumption_request` still has neither of the safeguards its siblings have — a mid-stream `ReadError`/`RemoteProtocolError` propagates into the transport task group and tears down every stream (the exact ExceptionGroup symptom from #2110, one packet after the case this PR fixes), and a resumed stream that ends cleanly without delivering a response r
Comment thread
claude[bot] marked this conversation as resolved.
Outdated
Comment thread
claude[bot] marked this conversation as resolved.
Outdated
Comment thread
claude[bot] marked this conversation as resolved.
Outdated
logger.debug("Resumption GET SSE connection established")

async for sse in event_source: # pragma: no branch
Expand Down
42 changes: 42 additions & 0 deletions tests/client/test_streamable_http.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
CLIENT_CAPABILITIES_META_KEY,
CLIENT_INFO_META_KEY,
CONNECTION_CLOSED,
INTERNAL_ERROR,
INVALID_REQUEST,
METHOD_NOT_FOUND,
PROTOCOL_VERSION_META_KEY,
Expand All @@ -31,6 +32,7 @@
from starlette.types import Receive, Scope, Send

from mcp.client.streamable_http import (
LAST_EVENT_ID,
MAX_RECONNECTION_ATTEMPTS,
RequestContext,
StreamableHTTPTransport,
Expand Down Expand Up @@ -132,6 +134,46 @@ def handler(request: httpx2.Request) -> httpx2.Response:
assert reply.message.error.code == METHOD_NOT_FOUND


@pytest.mark.anyio
@pytest.mark.parametrize("status", [401, 403, 500])
async def test_resumption_get_http_error_resolves_caller_and_transport_survives(status: int) -> None:
"""A non-2xx on the resumption GET resolves the waiting request with a JSON-RPC error
correlated to its id, and the transport stays usable for follow-up requests (SDK-defined;
#2110 — the status error used to escape into the task group and tear down every stream).
"""

def handler(request: httpx2.Request) -> httpx2.Response:
if request.method == "GET" and LAST_EVENT_ID in request.headers:
return httpx2.Response(status)
body = json.loads(request.content)
return httpx2.Response(200, json={"jsonrpc": "2.0", "id": body["id"], "result": {}})

with anyio.fail_after(5):
async with (
httpx2.AsyncClient(transport=httpx2.MockTransport(handler)) as http,
streamable_http_client("http://test/mcp", http_client=http) as (read, write),
):
await write.send(
SessionMessage(
message=JSONRPCRequest(jsonrpc="2.0", id=1, method="tools/call", params={}),
metadata=ClientMessageMetadata(resumption_token="token-1"),
)
)
reply = await read.receive()
assert isinstance(reply, SessionMessage)
assert isinstance(reply.message, JSONRPCError)
assert reply.message.id == 1
assert reply.message.error.code == INTERNAL_ERROR
assert reply.message.error.message == snapshot("Server returned an error response")

# The transport survived: a plain follow-up request still round-trips.
await write.send(SessionMessage(JSONRPCRequest(jsonrpc="2.0", id=2, method="tools/list", params={})))
follow_up = await read.receive()
assert isinstance(follow_up, SessionMessage)
assert isinstance(follow_up.message, JSONRPCResponse)
assert follow_up.message.id == 2


@pytest.mark.anyio
async def test_initialize_post_clears_cached_pv_header_and_unstamped_posts_read_it() -> None:
"""``initialize`` discards the cached protocol-version header; every other POST reads it.
Expand Down
76 changes: 76 additions & 0 deletions tests/shared/test_sse.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
from starlette.requests import Request
from starlette.responses import Response
from starlette.routing import Mount, Route
from starlette.types import Receive, Scope, Send

import mcp.client.sse
from mcp.client.session import ClientSession
Expand Down Expand Up @@ -108,6 +109,39 @@
return make_app(Server(SERVER_NAME, on_read_resource=_handle_read_resource))


def make_app_rejecting_posts(reject: dict[str, int]) -> Starlette:
"""Like `make_server_app`, but the message POST is answered with a bare HTTP error
for JSON-RPC messages whose method appears in `reject` (they never reach the server)."""
server = Server(SERVER_NAME, on_read_resource=_handle_read_resource)
sse = SseServerTransport(
"/messages/", security_settings=TransportSecuritySettings(enable_dns_rebinding_protection=False)
)

async def handle_sse(request: Request) -> Response:
async with sse.connect_sse(request.scope, request.receive, request._send) as (read_stream, write_stream):
await server.run(read_stream, write_stream, server.create_initialization_options())
return Response()

async def handle_post(scope: Scope, receive: Receive, send: Send) -> None:
body = await Request(scope, receive).body()
status = reject.get(json.loads(body).get("method"))
if status is not None:
await Response(status_code=status)(scope, receive, send)
return

async def replay() -> dict[str, Any]:
return {"type": "http.request", "body": body, "more_body": False}

await sse.handle_post_message(scope, replay, send)

return Starlette(
routes=[
Route("/sse", endpoint=handle_sse),
Mount("/messages/", app=handle_post),
]

Check warning on line 141 in tests/shared/test_sse.py

View check run for this annotation

Claude / Claude Code Review

[quality] make_app_rejecting_posts duplicates make_app's wiring

make_app_rejecting_posts duplicates the SseServerTransport construction, the handle_sse closure, and the Starlette route assembly verbatim from make_app in this same file, differing only in the POST endpoint — two copies of the SSE wiring that must now be kept in sync. Consider extending make_app with an optional POST-app wrapper (e.g. wrap_post: Callable[[ASGIApp], ASGIApp] applied around sse.handle_post_message) so this helper only supplies its rejecting handle_post.
Comment thread
claude[bot] marked this conversation as resolved.
Outdated
)


@pytest.mark.anyio
async def test_raw_sse_connection() -> None:
"""The SSE GET responds 200 with an event-stream content type, announcing the session
Expand Down Expand Up @@ -224,6 +258,48 @@
await session.read_resource(uri="xxx://will-not-work")


@pytest.mark.anyio
@pytest.mark.parametrize("status_code", [401, 403, 500])
async def test_sse_client_request_post_http_error_reaches_caller_and_session_survives(status_code: int) -> None:
"""A non-2xx on a request's message POST reaches the waiting caller promptly as a JSON-RPC
error correlated to the request, and the session stays usable (SDK-defined; #2110 — the
status error used to be swallowed inside post_writer, hanging the caller forever).
"""
factory = in_process_client_factory(make_app_rejecting_posts({"resources/read": status_code}))
with anyio.fail_after(5):
async with sse_client(f"{BASE_URL}/sse", httpx_client_factory=factory) as streams:
async with ClientSession(*streams) as session:
await session.initialize()

with pytest.raises(MCPError) as exc_info:
await session.read_resource(uri="foobar://should-work")
assert exc_info.value.error.code == types.INTERNAL_ERROR
assert exc_info.value.error.message == snapshot("Server returned an error response")

# The session survived the failed POST: the next request round-trips.
assert isinstance(await session.send_ping(), EmptyResult)


@pytest.mark.anyio
async def test_sse_client_notification_post_http_error_leaves_session_usable() -> None:
"""A non-2xx on a notification's message POST resolves no caller (a notification has no
waiter) and leaves the session usable for subsequent requests (SDK-defined; #2110)."""
factory = in_process_client_factory(make_app_rejecting_posts({"notifications/cancelled": 500}))
with anyio.fail_after(5):
async with sse_client(f"{BASE_URL}/sse", httpx_client_factory=factory) as streams:
async with ClientSession(*streams) as session:
await session.initialize()

# Fire-and-forget: the rejected POST must neither raise nor stall the writer.
await session.send_notification(
types.CancelledNotification(params=types.CancelledNotificationParams(request_id=999))
)

# The write loop is serialized, so this request's POST happens strictly after
# the rejected one; its success proves the failure was contained.
assert isinstance(await session.send_ping(), EmptyResult)


@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."""
Expand Down
Loading