Skip to content

Commit 58e28e3

Browse files
committed
Backport abandoned-request resolution from 2.x to the 1.x Streamable HTTP client
On the 1.x line, a pending request stays unresolved forever when the server answers its POST with 202 Accepted, when the per-request SSE stream ends without carrying a response event, or when reconnection attempts are exhausted: the caller only learns of the failure when its own deadline fires, and it surfaces as a timeout rather than a disconnect. Port the _resolve_abandoned_request mechanism from 2.x (PR #3047) to the 1.x StreamableHTTPTransport: resolve the pending request with a synthesized JSONRPCError (CONNECTION_CLOSED, or INVALID_REQUEST for the 202 case) at the three sites where its response can never arrive. Notifications are not affected. Fixes #3441
1 parent 3eed7ce commit 58e28e3

2 files changed

Lines changed: 226 additions & 7 deletions

File tree

src/mcp/client/streamable_http.py

Lines changed: 55 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -28,8 +28,10 @@
2828
)
2929
from mcp.shared.message import ClientMessageMetadata, SessionMessage
3030
from mcp.types import (
31+
CONNECTION_CLOSED,
3132
ErrorData,
3233
InitializeResult,
34+
INVALID_REQUEST,
3335
JSONRPCError,
3436
JSONRPCMessage,
3537
JSONRPCNotification,
@@ -345,6 +347,15 @@ async def _handle_post_request(self, ctx: RequestContext) -> None:
345347
) as response:
346348
if response.status_code == 202:
347349
logger.debug("Received 202 Accepted")
350+
if isinstance(message.root, JSONRPCRequest):
351+
# A request's response arrives on this POST's body; 202 says
352+
# none will follow. Resolve rather than park the caller forever.
353+
await self._resolve_abandoned_request(
354+
ctx.read_stream_writer,
355+
message.root.id,
356+
"server answered a request with 202 Accepted",
357+
code=INVALID_REQUEST,
358+
)
348359
return
349360

350361
if response.status_code == 404: # pragma: no branch
@@ -404,6 +415,10 @@ async def _handle_sse_response(
404415
last_event_id: str | None = None
405416
retry_interval_ms: int | None = None
406417

418+
original_request_id = None
419+
if isinstance(ctx.session_message.message.root, JSONRPCRequest):
420+
original_request_id = ctx.session_message.message.root.id
421+
407422
try:
408423
event_source = EventSource(response)
409424
async for sse in event_source.aiter_sse(): # pragma: no branch
@@ -430,9 +445,34 @@ async def _handle_sse_response(
430445
logger.debug(f"SSE stream ended: {e}")
431446

432447
# Stream ended without response - reconnect if we received an event with ID
433-
if last_event_id is not None: # pragma: no branch
448+
if last_event_id is not None:
434449
logger.info("SSE stream disconnected, reconnecting...")
435450
await self._handle_reconnection(ctx, last_event_id, retry_interval_ms)
451+
else:
452+
# Not resumable: resolve the waiter, else the request would hang
453+
# forever instead of learning the connection is lost.
454+
await self._resolve_abandoned_request(
455+
ctx.read_stream_writer, original_request_id, "SSE stream ended without a response"
456+
)
457+
458+
async def _resolve_abandoned_request(
459+
self,
460+
read_stream_writer: StreamWriter,
461+
request_id: RequestId,
462+
message: str,
463+
*,
464+
code: int = CONNECTION_CLOSED,
465+
) -> None:
466+
"""Resolve a request whose response can never arrive with a synthesized error.
467+
468+
Best-effort: a closed read stream means the session is tearing down.
469+
"""
470+
error_data = ErrorData(code=code, message=message)
471+
session_message = SessionMessage(JSONRPCMessage(JSONRPCError(jsonrpc="2.0", id=request_id, error=error_data)))
472+
try:
473+
await read_stream_writer.send(session_message)
474+
except (anyio.BrokenResourceError, anyio.ClosedResourceError):
475+
logger.debug("read stream closed before request %r could be resolved", request_id)
436476

437477
async def _handle_reconnection(
438478
self,
@@ -442,9 +482,22 @@ async def _handle_reconnection(
442482
attempt: int = 0,
443483
) -> None:
444484
"""Reconnect with Last-Event-ID to resume stream after server disconnect."""
485+
# Extract original request ID to map responses
486+
original_request_id = None
487+
if isinstance(ctx.session_message.message.root, JSONRPCRequest):
488+
original_request_id = ctx.session_message.message.root.id
489+
445490
# Bail if max retries exceeded
446-
if attempt >= MAX_RECONNECTION_ATTEMPTS: # pragma: no cover
491+
if attempt >= MAX_RECONNECTION_ATTEMPTS:
447492
logger.debug(f"Max reconnection attempts ({MAX_RECONNECTION_ATTEMPTS}) exceeded")
493+
# Resolve on give-up: a request with no read timeout would otherwise
494+
# hang its caller forever.
495+
if original_request_id is not None:
496+
await self._resolve_abandoned_request(
497+
ctx.read_stream_writer,
498+
original_request_id,
499+
"SSE stream ended and reconnection attempts were exhausted",
500+
)
448501
return
449502

450503
# Always wait - use server value or default
@@ -454,11 +507,6 @@ async def _handle_reconnection(
454507
headers = self._prepare_headers()
455508
headers[LAST_EVENT_ID] = last_event_id
456509

457-
# Extract original request ID to map responses
458-
original_request_id = None
459-
if isinstance(ctx.session_message.message.root, JSONRPCRequest): # pragma: no branch
460-
original_request_id = ctx.session_message.message.root.id
461-
462510
try:
463511
async with aconnect_sse(
464512
ctx.client,
Lines changed: 171 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,171 @@
1+
"""Deterministic regression tests for abandoned-request resolution (#3441).
2+
3+
When a request's response can never arrive - the server answers the POST
4+
with 202 Accepted, the per-request SSE stream ends without a response
5+
event, or reconnection attempts are exhausted - the v1.x transport must
6+
resolve the pending request with a synthesized JSONRPCError instead of
7+
parking the caller until its own timeout fires.
8+
"""
9+
10+
from collections.abc import AsyncIterator
11+
12+
import anyio
13+
import httpx
14+
import pytest
15+
from mcp.shared.message import SessionMessage
16+
from mcp.types import (
17+
CONNECTION_CLOSED,
18+
INVALID_REQUEST,
19+
JSONRPCError,
20+
JSONRPCMessage,
21+
JSONRPCNotification,
22+
JSONRPCRequest,
23+
)
24+
25+
from mcp.client.streamable_http import (
26+
MAX_RECONNECTION_ATTEMPTS,
27+
RequestContext,
28+
StreamableHTTPTransport,
29+
)
30+
31+
32+
def _make_request_context(
33+
client: httpx.AsyncClient,
34+
message: JSONRPCMessage,
35+
read_stream_writer,
36+
) -> RequestContext:
37+
session_message = SessionMessage(message)
38+
return RequestContext(
39+
client=client,
40+
session_id=None,
41+
session_message=session_message,
42+
metadata=None,
43+
read_stream_writer=read_stream_writer,
44+
)
45+
46+
47+
def _request(id_: str) -> JSONRPCMessage:
48+
return JSONRPCMessage(JSONRPCRequest(jsonrpc="2.0", id=id_, method="tools/call", params={}))
49+
50+
51+
class _DyingSSEStream(httpx.AsyncByteStream):
52+
"""Emits one id-less comment then breaks - a non-resumable stream dropping."""
53+
54+
def __init__(self) -> None:
55+
self.opened = anyio.Event()
56+
57+
async def __aiter__(self) -> AsyncIterator[bytes]:
58+
self.opened.set()
59+
yield b": hello\n\n"
60+
raise httpx.ReadError("connection reset")
61+
62+
async def aclose(self) -> None:
63+
pass
64+
65+
66+
@pytest.mark.anyio
67+
async def test_non_resumable_sse_drop_resolves_request_with_error() -> None:
68+
"""A per-request SSE stream that dies having carried no event ids can never
69+
deliver its response; the transport resolves the waiter with CONNECTION_CLOSED
70+
instead of hanging forever."""
71+
transport = StreamableHTTPTransport("http://test/mcp")
72+
73+
def handler(request: httpx.Request) -> httpx.Response:
74+
return httpx.Response(200, headers={"content-type": "text/event-stream"}, stream=_DyingSSEStream())
75+
76+
streams = anyio.create_memory_object_stream(4)
77+
try:
78+
write_stream, read_stream = streams
79+
async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as client:
80+
ctx = _make_request_context(client, _request("req-sse"), write_stream)
81+
with anyio.fail_after(5):
82+
await transport._handle_post_request(ctx)
83+
reply = await read_stream.receive()
84+
finally:
85+
write_stream.close()
86+
read_stream.close()
87+
88+
assert isinstance(reply.message.root, JSONRPCError)
89+
assert reply.message.root.id == "req-sse"
90+
assert reply.message.root.error.code == CONNECTION_CLOSED
91+
92+
93+
@pytest.mark.anyio
94+
async def test_post_answered_with_202_resolves_request_with_error() -> None:
95+
"""A request answered with 202 Accepted will never receive a response body;
96+
the transport resolves the waiter with INVALID_REQUEST instead of hanging."""
97+
98+
def handler(request: httpx.Request) -> httpx.Response:
99+
return httpx.Response(202)
100+
101+
transport = StreamableHTTPTransport("http://test/mcp")
102+
streams = anyio.create_memory_object_stream(4)
103+
try:
104+
write_stream, read_stream = streams
105+
async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as client:
106+
ctx = _make_request_context(client, _request("req-202"), write_stream)
107+
with anyio.fail_after(5):
108+
await transport._handle_post_request(ctx)
109+
reply = await read_stream.receive()
110+
finally:
111+
write_stream.close()
112+
read_stream.close()
113+
114+
assert isinstance(reply.message.root, JSONRPCError)
115+
assert reply.message.root.id == "req-202"
116+
assert reply.message.root.error.code == INVALID_REQUEST
117+
118+
119+
@pytest.mark.anyio
120+
async def test_post_answered_with_202_does_not_resolve_notifications() -> None:
121+
"""Notifications have no waiter to resolve; a 202 answer must not inject an
122+
error into the read stream for them."""
123+
124+
def handler(request: httpx.Request) -> httpx.Response:
125+
return httpx.Response(202)
126+
127+
transport = StreamableHTTPTransport("http://test/mcp")
128+
notification = JSONRPCMessage(JSONRPCNotification(jsonrpc="2.0", method="notifications/initialized"))
129+
streams = anyio.create_memory_object_stream(4)
130+
try:
131+
write_stream, read_stream = streams
132+
async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as client:
133+
ctx = _make_request_context(client, notification, write_stream)
134+
with anyio.move_on_after(1):
135+
await transport._handle_post_request(ctx)
136+
await read_stream.receive()
137+
pytest.fail("a notification must not be resolved with an error")
138+
finally:
139+
write_stream.close()
140+
read_stream.close()
141+
142+
143+
@pytest.mark.anyio
144+
async def test_exhausted_reconnection_attempts_resolve_request_with_error() -> None:
145+
"""When reconnection attempts are exhausted for a request whose SSE stream
146+
keeps dying, the transport resolves the waiter with CONNECTION_CLOSED."""
147+
transport = StreamableHTTPTransport("http://test/mcp")
148+
149+
def handler(request: httpx.Request) -> httpx.Response: # pragma: no cover
150+
pytest.fail("exhausted reconnection must resolve without further HTTP calls")
151+
152+
streams = anyio.create_memory_object_stream(4)
153+
try:
154+
write_stream, read_stream = streams
155+
async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as client:
156+
ctx = _make_request_context(client, _request("req-ex"), write_stream)
157+
with anyio.fail_after(5):
158+
await transport._handle_reconnection(
159+
ctx,
160+
last_event_id="1",
161+
retry_interval_ms=1,
162+
attempt=MAX_RECONNECTION_ATTEMPTS,
163+
)
164+
reply = await read_stream.receive()
165+
finally:
166+
write_stream.close()
167+
read_stream.close()
168+
169+
assert isinstance(reply.message.root, JSONRPCError)
170+
assert reply.message.root.id == "req-ex"
171+
assert reply.message.root.error.code == CONNECTION_CLOSED

0 commit comments

Comments
 (0)