Skip to content
42 changes: 32 additions & 10 deletions src/mcp/client/sse.py
Original file line number Diff line number Diff line change
Expand Up @@ -120,17 +120,39 @@
async with write_stream_reader, write_stream:

async def _send_message(session_message: SessionMessage) -> None:
# A POST failure must not raise: the post_writer handler below
# would swallow it, hanging the waiting caller forever and killing
# the write loop (#2110). Mirror the streamable-HTTP transport
# instead: resolve the waiter with an error correlated to its
# request id, keeping the session usable.
logger.debug(f"Sending client message: {session_message}")
response = await client.post(
endpoint_url,
json=session_message.message.model_dump(
by_alias=True,
mode="json",
exclude_unset=True,
),
)
response.raise_for_status()
logger.debug(f"Client message sent successfully: {response.status_code}")
message = session_message.message
try:
response = await client.post(
endpoint_url,
json=message.model_dump(
by_alias=True,
mode="json",
exclude_unset=True,
),
)
except httpx2.HTTPError as exc:
logger.exception("Error POSTing message")
error = types.ErrorData(
code=types.CONNECTION_CLOSED, message=f"Failed to send message: {exc}"
)

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

View check run for this annotation

Claude / Claude Code Review

SSE POST: SDK OAuth-flow exceptions bypass the new httpx2.HTTPError catch and still take the swallowed-hang path

The new `except httpx2.HTTPError` around `client.post()` doesn't catch the SDK's own OAuth exceptions (`OAuthFlowError`/`OAuthTokenError`/`OAuthRegistrationError` subclass `Exception`, not `HTTPError`), which `OAuthClientProvider.async_auth_flow` raises from inside the same `client.post()` call when a mid-session re-auth fails. Those still escape into `post_writer`'s catch-all and take the exact swallowed-hang path this comment says the code prevents — widening the catch to `except (httpx2.HTTPE
Comment thread
claude[bot] marked this conversation as resolved.
Comment thread
claude[bot] marked this conversation as resolved.
else:
if response.is_success:
logger.debug(f"Client message sent successfully: {response.status_code}")
return
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

error = types.ErrorData(
code=types.INTERNAL_ERROR, message="Server returned an error response"
)
# A notification has no waiter to resolve, so its failure is only logged.

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

View check run for this annotation

Claude / Claude Code Review

SSE POST 404 (session expiry) collapsed into generic INTERNAL_ERROR instead of the Session-terminated mapping the PR standardized elsewhere

On the SSE transport, a 404 to the message POST — the SDK server's own session-expiry signal (`SseServerTransport.handle_post_message` answers exactly 404 "Could not find session") — is collapsed into the generic `ErrorData(INTERNAL_ERROR, "Server returned an error response")`, while this same PR extended the 404 → `INVALID_REQUEST`/"Session terminated" mapping to the streamable resumption GET precisely so reconnect logic keyed on that shape works across paths. Consider a 404 branch before the g
Comment thread
claude[bot] marked this conversation as resolved.
if isinstance(message, types.JSONRPCRequest):
reply = types.JSONRPCError(jsonrpc="2.0", id=message.id, error=error)
await read_stream_writer.send(SessionMessage(reply))

async for session_message in write_stream_reader:
Comment thread
claude[bot] marked this conversation as resolved.
sender_ctx = write_stream_reader.last_context
Comment thread
claude[bot] marked this conversation as resolved.
Outdated
Expand Down
66 changes: 47 additions & 19 deletions src/mcp/client/streamable_http.py
Original file line number Diff line number Diff line change
Expand Up @@ -248,25 +248,51 @@
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: post_writer dispatches here on message type as well as
# 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

async with ctx.client.sse(self.url, headers=headers) as event_source:
event_source.response.raise_for_status()
logger.debug("Resumption GET SSE connection established")
try:
async with ctx.client.sse(self.url, headers=headers) as event_source:
if not event_source.response.is_success:
# 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).
if event_source.response.status_code == 404 and self.session_id is not None:
# The GET carried our Mcp-Session-Id, so a 404 is the session-expiry
# signal reconnect logic keys on - same mapping as the POST path.

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

View check run for this annotation

Claude / Claude Code Review

User-visible error-behavior changes missing from docs/migration.md (required by AGENTS.md in the same PR)

This PR changes user-visible client error behavior — a non-2xx on the streamable-HTTP resumption GET now surfaces as a correlated `MCPError` instead of an escaping `HTTPStatusError`, and SSE-transport POST failures now raise `MCPError` from the waiting call instead of hanging — but touches no `docs/` file, which AGENTS.md requires in the same PR. It also makes the existing migration-guide section's closing claim (docs/migration.md ~2233: "Connect-level failures such as `httpx2.ConnectError` stil
Comment thread
claude[bot] marked this conversation as resolved.
Outdated
await self._resolve_abandoned_request(
ctx.read_stream_writer, original_request_id, "Session terminated", code=INVALID_REQUEST
)
return
await self._resolve_abandoned_request(
ctx.read_stream_writer,
original_request_id,
"Server returned an error response",
code=INTERNAL_ERROR,
)
return
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,
)
if is_complete:
await event_source.response.aclose()
break
async for sse in event_source:
is_complete = await self._handle_sse_event(
sse,
ctx.read_stream_writer,
original_request_id,
Comment on lines +261 to +274

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.

🟣 Pre-existing issue (not introduced by this PR): the automatic Last-Event-ID reconnection GET in _handle_reconnection still does a bare raise_for_status() inside its except Exception retry loop, so a deterministic session-expiry 404 is retried MAX_RECONNECTION_ATTEMPTS times (~2s of futile sleeps) and then surfaced as CONNECTION_CLOSED / "SSE stream ended and reconnection attempts were exhausted" — never the INVALID_REQUEST / "Session terminated" mapping this PR standardized on the resumption GET and both POST paths (handle_get_stream has the same pattern, with lower impact since it has no waiter). Now that status_error_data exists, the fix is a small mirror of this PR's own resumption-GET branch: check event_source.response.is_success before retrying and resolve the waiter via status_error_data(status, has_session=self.session_id is not None), since a non-2xx GET can never deliver the stream.

Extended reasoning...

What the bug is. This PR establishes a cross-path contract — a 404 while a session is held maps to MCPError(INVALID_REQUEST, "Session terminated"), now centralized in status_error_data() (src/mcp/client/_transport.py) and applied on the streamable message POST, the metadata-driven resumption GET, and the SSE message POST. But there is a third Last-Event-ID GET leg the mapping never reaches: _handle_reconnection (src/mcp/client/streamable_http.py:518-549), the automatic reconnect that fires from _handle_sse_response whenever a request's SSE response stream drops after carrying event ids. It still calls a bare event_source.response.raise_for_status() (line ~520) inside a try whose except Exception handler retries with attempt + 1 (lines ~546-549). handle_get_stream (line ~213) has the same bare-raise_for_status-inside-retry pattern for the server-initiated stream.\n\nThe code path. The trigger is realistic — a server restart drops the in-flight SSE stream and evicts the session:\n\n1. session.call_tool() over streamable HTTP; the server answers with an SSE stream that emits an event id, then the connection drops mid-stream.\n2. _handle_sse_response catches the read error and, because last_event_id is set, calls _handle_reconnection(ctx, last_event_id, ...).\n3. Meanwhile the session has expired. The SDK's own server answers any request bearing an unknown/expired Mcp-Session-Id with HTTP 404 (src/mcp/server/streamable_http_manager.py:361-371, "return 404 per MCP spec"), and the reconnection GET carries the stale Mcp-Session-Id via _prepare_headers().\n4. raise_for_status() raises HTTPStatusError; the except Exception handler never inspects the status, sleeps DEFAULT_RECONNECTION_DELAY_MS (1s, or the server-provided retry), and re-sends the doomed GET — MAX_RECONNECTION_ATTEMPTS times.\n5. On give-up it resolves the waiter via _resolve_abandoned_request with the default code=CONNECTION_CLOSED and "SSE stream ended and reconnection attempts were exhausted".\n\nWhy existing code doesn't prevent it. The retry handler treats every failure as transient — it sees only the exception, never the status, and the except branch carries a # pragma: no cover, so no test drives a non-2xx through this leg. The mapping this PR added lives in _handle_resumption_request, which handles only metadata-driven resumption (an explicit resumption_token stamped by the caller); the automatic mid-stream reconnect never routes through it.\n\nImpact. The caller's error shape for the identical server-side event — session expired, observed on a Last-Event-ID GET — depends on which internal leg observed it. On the metadata-driven resumption GET the caller gets MCPError(-32600, "Session terminated") promptly (pinned by test_resumption_get_404_with_session_reports_session_terminated). On the automatic reconnection GET the same expiry costs ~2s of futile retries and then arrives as MCPError(-32000, "SSE stream ended and reconnection attempts were exhausted"). A reconnect wrapper written per the migration guide's v2 pattern (exc.code == INVALID_REQUEST and exc.message == "Session terminated" → rebuild the connection) silently never matches on this leg, so "reconnecting would fix this" is indistinguishable from a dead network — the precise ambiguity the PR's 404 mapping exists to remove. In handle_get_stream the impact is only the futile retry budget plus a silently abandoned server-initiated stream (no waiter to mis-resolve).\n\nWhy pre-existing, not blocking. All three verifiers converged on this: _handle_reconnection and handle_get_stream are untouched by this PR's diff, and their behavior is unchanged from before — the error was already contained and retried pre-PR; nothing hangs and nothing tears down, so merging this PR breaks nothing that worked. The migration table's literal wording ("a request re-attached with a resumption token (Last-Event-ID)") is scoped to the metadata-driven resumption GET, so no documented claim is strictly falsified — only the PR body's "works everywhere" prose and cross-path consistency. This PR merely makes the residual gap salient (and cheap to close).\n\nHow to fix. Mirror this PR's own _handle_resumption_request branch in _handle_reconnection: before raise_for_status(), check event_source.response.is_success; on a non-2xx, resolve the waiter immediately via self._resolve_abandoned_request(ctx.read_stream_writer, original_request_id, error_data.message, code=error_data.code) with error_data = status_error_data(event_source.response.status_code, has_session=self.session_id is not None) and return — a non-2xx GET can never deliver the stream, so retrying is pointless. The same two-line guard fits handle_get_stream (there, just return on non-2xx instead of resolving a waiter). Fine as a follow-up PR.

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.

Agreed this is a real gap — but it's pre-existing, on a path this PR doesn't rewrite (_handle_reconnection's retry loop has its own containment and give-up resolution, so the failure mode is bounded retries rather than a hang or teardown), and as you note it's not blocking. Declining here to keep this PR reviewable; it belongs in a grouped follow-up together with the other two pre-existing findings from this pass (uncontained POST-path reply sends, dropped ctx.metadata.headers on the resumption/reconnection GETs).


Generated by Claude Code

ctx.metadata.on_resumption_token_update if ctx.metadata else None,
)
if is_complete:
await event_source.response.aclose()
return
except Exception:
logger.debug("Resumption stream ended", exc_info=True)

# Stream ended without a response, cleanly or mid-read: resolve the waiter,
# mirroring `_handle_sse_response`, else the caller would hang forever.
await self._resolve_abandoned_request(
ctx.read_stream_writer, original_request_id, "resumption stream ended without a response"
)

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 @@ -556,8 +582,10 @@
else None
)

# Check if this is a resumption request
is_resumption = bool(metadata and metadata.resumption_token)
# Only a request resumes: the token names an interrupted request's
# stream, and `_handle_resumption_request` needs the id to correlate
# its outcome. A notification stamped with one is POSTed as usual.
is_resumption = bool(metadata and metadata.resumption_token) and isinstance(message, JSONRPCRequest)

logger.debug(f"Sending client message: {message}")

Expand Down
185 changes: 185 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,7 +32,9 @@
from starlette.types import Receive, Scope, Send

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


@pytest.mark.anyio
@pytest.mark.parametrize("status", [302, 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).
An unfollowed redirect counts: its body is no event stream, so no response can arrive.
"""

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_resumption_get_404_with_session_reports_session_terminated() -> None:
"""A 404 on the resumption GET while a session id is held reports "Session terminated"
(INVALID_REQUEST) to the waiter, the same session-expiry mapping as the POST path, so
reconnect logic keyed on that error works across both (SDK-defined)."""

def handler(request: httpx2.Request) -> httpx2.Response:
if request.method == "GET" and LAST_EVENT_ID in request.headers:
return httpx2.Response(404)
if request.method == "DELETE": # session termination on close
return httpx2.Response(200)
body = json.loads(request.content)
return httpx2.Response(
200, json={"jsonrpc": "2.0", "id": body["id"], "result": {}}, headers={MCP_SESSION_ID: "sess-1"}
)

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),
):
# An initialize round-trip stores the session id the server stamps on its response.
await write.send(SessionMessage(JSONRPCRequest(jsonrpc="2.0", id=1, method="initialize", params={})))
assert isinstance(await read.receive(), SessionMessage)

await write.send(
SessionMessage(
message=JSONRPCRequest(jsonrpc="2.0", id=2, 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 == 2
assert reply.message.error.code == INVALID_REQUEST
assert reply.message.error.message == snapshot("Session terminated")


@pytest.mark.anyio
async def test_notification_with_resumption_token_is_posted_not_resumed() -> None:
"""A notification stamped with a resumption token is POSTed like any notification, and the
write loop survives to serve the next request (SDK-defined: the token names an interrupted
request's stream, so resumption applies to requests only)."""
recorded: list[httpx2.Request] = []

def handler(request: httpx2.Request) -> httpx2.Response:
recorded.append(request)
body = json.loads(request.content)
if "id" not in body:
return httpx2.Response(202)
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=JSONRPCNotification(jsonrpc="2.0", method="notifications/foo", params={}),
metadata=ClientMessageMetadata(resumption_token="token-1"),
)
)
await write.send(SessionMessage(JSONRPCRequest(jsonrpc="2.0", id=1, method="tools/list", params={})))
reply = await read.receive()
assert isinstance(reply, SessionMessage)
assert isinstance(reply.message, JSONRPCResponse)
assert reply.message.id == 1
# The stamped notification went out as a plain POST, not a resumption GET.
assert [r.method for r in recorded] == ["POST", "POST"]


@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 Expand Up @@ -632,6 +749,74 @@ def handler(request: httpx2.Request) -> httpx2.Response:
assert reply.message.error.code == CONNECTION_CLOSED


@pytest.mark.anyio
async def test_resumption_stream_dying_mid_read_resolves_caller_and_transport_survives() -> None:
"""A resumption GET stream that dies mid-read resolves the waiter with CONNECTION_CLOSED
and the transport stays usable for follow-up requests (SDK-defined; #2110 — the read error
used to escape into the task group and tear down every stream)."""
dying = _DyingSSEStream()

def handler(request: httpx2.Request) -> httpx2.Response:
if request.method == "GET" and LAST_EVENT_ID in request.headers:
return httpx2.Response(200, headers={"content-type": "text/event-stream"}, stream=dying)
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 == CONNECTION_CLOSED
assert reply.message.error.message == snapshot("resumption stream ended without a 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_resumption_stream_clean_end_without_response_resolves_caller() -> None:
"""A resumption GET stream that closes cleanly without delivering a response (e.g. the
server no longer holds the resumed request's events) resolves the waiter with an error
instead of hanging it forever (SDK-defined; #2110)."""

def handler(request: httpx2.Request) -> httpx2.Response:
assert request.method == "GET" and LAST_EVENT_ID in request.headers
return httpx2.Response(200, headers={"content-type": "text/event-stream"}, content=b": nothing to replay\n\n")

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 == CONNECTION_CLOSED
assert reply.message.error.message == snapshot("resumption stream ended without a response")


class _DeliverOnCommandSSEStream(httpx2.AsyncByteStream):
"""Parks after opening, then delivers one JSON-RPC response when told."""

Expand Down
Loading
Loading