Skip to content

Commit 9e33664

Browse files
author
Jianke LIN
committed
fix(streamable-http): bound empty SSE reconnect loops
1 parent 2abc3c1 commit 9e33664

2 files changed

Lines changed: 63 additions & 3 deletions

File tree

src/mcp/client/streamable_http.py

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -518,12 +518,15 @@ async def _handle_reconnection(
518518
# Track for potential further reconnection
519519
reconnect_last_event_id: str = last_event_id
520520
reconnect_retry_ms = retry_interval_ms
521+
made_progress = False
521522

522523
async for sse in event_source.aiter_sse():
523524
if sse.id: # pragma: no branch
524525
reconnect_last_event_id = sse.id
525526
if sse.retry is not None:
526527
reconnect_retry_ms = sse.retry
528+
if sse.event == "message" and bool(sse.data):
529+
made_progress = True
527530

528531
is_complete = await self._handle_sse_event(
529532
sse,
@@ -535,9 +538,11 @@ async def _handle_reconnection(
535538
await event_source.response.aclose()
536539
return
537540

538-
# Stream ended again without response - reconnect again (reset attempt counter)
541+
# Stream ended again without response - reconnect again. Only reset
542+
# the retry counter when the resumed stream delivered real data.
539543
logger.info("SSE stream disconnected, reconnecting...")
540-
await self._handle_reconnection(ctx, reconnect_last_event_id, reconnect_retry_ms, 0)
544+
next_attempt = 0 if made_progress else attempt + 1
545+
await self._handle_reconnection(ctx, reconnect_last_event_id, reconnect_retry_ms, next_attempt)
541546
except Exception as e:
542547
logger.debug(f"Reconnection failed: {e}")
543548
# Try to reconnect again if we still have an event ID

tests/client/test_streamable_http.py

Lines changed: 56 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -98,14 +98,69 @@ async def test_sse_response_disconnect_before_any_event_id_fails_request() -> No
9898

9999
async with read_stream_writer, read_stream:
100100
await transport._handle_sse_response(response, ctx)
101-
message = await read_stream.receive()
101+
with anyio.fail_after(5):
102+
message = await read_stream.receive()
102103

103104
assert isinstance(message, SessionMessage)
104105
assert isinstance(message.message, JSONRPCError)
105106
assert message.message.id == 1
106107
assert message.message.error.code == CONNECTION_CLOSED
107108

108109

110+
@pytest.mark.anyio
111+
async def test_reconnection_empty_streams_count_toward_max_attempts(monkeypatch: pytest.MonkeyPatch) -> None:
112+
class PrimingOnlyEventSource:
113+
def __init__(self) -> None:
114+
self.response = httpx.Response(200)
115+
116+
async def __aenter__(self) -> "PrimingOnlyEventSource":
117+
nonlocal reconnect_attempts
118+
reconnect_attempts += 1
119+
return self
120+
121+
async def __aexit__(self, *args: object) -> None:
122+
return None
123+
124+
async def aiter_sse(self) -> object:
125+
yield type(
126+
"SSE",
127+
(),
128+
{"event": "message", "data": "", "id": f"event-{reconnect_attempts}", "retry": 0},
129+
)()
130+
131+
def connect_sse(*args: object, **kwargs: object) -> PrimingOnlyEventSource:
132+
return PrimingOnlyEventSource()
133+
134+
reconnect_attempts = 0
135+
monkeypatch.setattr(
136+
"mcp.client.streamable_http.aconnect_sse",
137+
connect_sse,
138+
)
139+
140+
transport = StreamableHTTPTransport("http://example.com/mcp")
141+
async with httpx.AsyncClient() as client:
142+
read_stream_writer, read_stream = create_context_streams[SessionMessage | Exception](1)
143+
request = JSONRPCRequest(jsonrpc="2.0", id=1, method="tools/call", params={"name": "noop", "arguments": {}})
144+
ctx = RequestContext(
145+
client=client,
146+
session_id=None,
147+
session_message=SessionMessage(request),
148+
metadata=None,
149+
read_stream_writer=read_stream_writer,
150+
)
151+
152+
async with read_stream_writer, read_stream:
153+
with anyio.fail_after(5):
154+
await transport._handle_reconnection(ctx, "event-1", retry_interval_ms=0)
155+
message = await read_stream.receive()
156+
157+
assert reconnect_attempts == 2
158+
assert isinstance(message, SessionMessage)
159+
assert isinstance(message.message, JSONRPCError)
160+
assert message.message.id == 1
161+
assert message.message.error.code == CONNECTION_CLOSED
162+
163+
109164
@pytest.mark.anyio
110165
async def test_sse_response_disconnect_ignores_closed_read_stream() -> None:
111166
transport = StreamableHTTPTransport("http://example.com/mcp")

0 commit comments

Comments
 (0)