Skip to content

Commit dfbe6ba

Browse files
author
OpenSourceMaintenance-Luccc-grok-4.5
committed
fix(server): terminate streamable HTTP sessions on manager shutdown
Why: #2150 — active SSE sessions were cancelled without terminate(), and terminate() did not close _sse_stream_writers, so ASGI/EventSourceResponse could hang on shutdown. Local Atlas work only; no public PR until operator approval.
1 parent 3a6f299 commit dfbe6ba

4 files changed

Lines changed: 110 additions & 0 deletions

File tree

src/mcp/server/streamable_http.py

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -774,11 +774,19 @@ async def terminate(self) -> None:
774774
"""Terminate the current session, closing all streams.
775775
776776
Once terminated, all requests with this session ID will receive 404 Not Found.
777+
778+
Active SSE writers are closed first so EventSourceResponse / ASGI callables can
779+
complete instead of hanging until the task group is cancelled (see #2150).
777780
"""
778781

779782
self._terminated = True
780783
logger.info(f"Terminating session: {self.mcp_session_id}")
781784

785+
# Close SSE stream writers first so long-lived GET/POST SSE responses finish.
786+
# Copy keys: close_sse_stream mutates the dict (includes GET_STREAM_KEY).
787+
for request_id in list(self._sse_stream_writers.keys()):
788+
self.close_sse_stream(request_id)
789+
782790
# We need a copy of the keys to avoid modification during iteration
783791
request_stream_keys = list(self._request_streams.keys())
784792

src/mcp/server/streamable_http_manager.py

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -153,6 +153,18 @@ async def lifespan(app: Starlette) -> AsyncIterator[None]:
153153
yield # Let the application run
154154
finally:
155155
logger.info("StreamableHTTP session manager shutting down")
156+
# Terminate active transports before cancelling the task group so
157+
# in-flight SSE responses can complete cleanly (issue #2150).
158+
active_transports = list(self._server_instances.values())
159+
for transport in active_transports:
160+
if not transport.is_terminated: # pragma: no branch
161+
try:
162+
await transport.terminate()
163+
except Exception: # pragma: no cover
164+
logger.exception(
165+
"Error terminating streamable HTTP session %s during shutdown",
166+
transport.mcp_session_id,
167+
)
156168
# Cancel task group to stop all spawned tasks
157169
tg.cancel_scope.cancel()
158170
self._task_group = None
Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,42 @@
1+
# Issue #2150 local verification notes
2+
3+
Upstream: https://github.com/modelcontextprotocol/python-sdk/issues/2150
4+
SHA baseline: 3a6f2996cdd8358957479791e8b26198c07d6a75
5+
6+
## Bug (still present on main at scout time)
7+
8+
1. `StreamableHTTPSessionManager.run()` finally only:
9+
- `tg.cancel_scope.cancel()`
10+
- `_server_instances.clear()`
11+
without calling `transport.terminate()` on active sessions.
12+
13+
2. `StreamableHTTPServerTransport.terminate()` closed request/read/write streams but
14+
did **not** close `_sse_stream_writers`, leaving EventSourceResponse hung.
15+
16+
## Fix (local branch `atlas/fix-2150-shutdown-sessions`)
17+
18+
1. Manager shutdown: terminate each non-terminated transport before cancel.
19+
2. Transport.terminate: close all SSE writers via close_sse_stream / close_standalone_sse_stream first.
20+
21+
## Tests added
22+
23+
- `test_terminate_closes_active_sse_stream_writers`
24+
- `test_manager_shutdown_terminates_active_sessions`
25+
26+
in `tests/server/test_streamable_http_manager.py`.
27+
28+
## Local env note
29+
30+
This Atlas sandbox lacked `pip`/`uv`; tests were not executed here. Run upstream:
31+
32+
```bash
33+
uv sync
34+
uv run pytest tests/server/test_streamable_http_manager.py -k 2150 -q
35+
# or by test name:
36+
uv run pytest tests/server/test_streamable_http_manager.py::test_terminate_closes_active_sse_stream_writers -q
37+
uv run pytest tests/server/test_streamable_http_manager.py::test_manager_shutdown_terminates_active_sessions -q
38+
```
39+
40+
## Publication
41+
42+
Operator approval: none. Do not open PR until approved.

tests/server/test_streamable_http_manager.py

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -746,3 +746,51 @@ async def test_anonymous_session_accepts_anonymous_requests(
746746
session_id = await _open_session(manager, None)
747747

748748
assert await _request_session(manager, session_id, None) != 404
749+
750+
751+
@pytest.mark.anyio
752+
async def test_terminate_closes_active_sse_stream_writers():
753+
"""Regression for #2150: terminate must close SSE writers so ASGI can finish.
754+
755+
Without this, manager shutdown cancels the task group while EventSourceResponse
756+
is still open and uvicorn logs "ASGI callable returned without completing response".
757+
"""
758+
transport = StreamableHTTPServerTransport(mcp_session_id="test-session-2150")
759+
send_stream, receive_stream = anyio.create_memory_object_stream[object](1)
760+
transport._sse_stream_writers["req-1"] = send_stream # type: ignore[assignment]
761+
762+
await transport.terminate()
763+
764+
assert transport.is_terminated
765+
assert "req-1" not in transport._sse_stream_writers
766+
with pytest.raises(anyio.ClosedResourceError):
767+
await send_stream.send(object()) # type: ignore[arg-type]
768+
await receive_stream.aclose()
769+
770+
771+
@pytest.mark.anyio
772+
async def test_manager_shutdown_terminates_active_sessions():
773+
"""Regression for #2150: run() finally should terminate tracked transports."""
774+
app = Server("test-shutdown-terminate")
775+
manager = StreamableHTTPSessionManager(app=app)
776+
transport = StreamableHTTPServerTransport(mcp_session_id="shutdown-session")
777+
# Inject a live session as if a client still held an SSE connection.
778+
manager._server_instances[transport.mcp_session_id] = transport # type: ignore[index]
779+
original_terminate = transport.terminate
780+
terminate_calls = 0
781+
782+
async def counting_terminate() -> None:
783+
nonlocal terminate_calls
784+
terminate_calls += 1
785+
await original_terminate()
786+
787+
transport.terminate = counting_terminate # type: ignore[method-assign]
788+
789+
async with manager.run():
790+
assert transport.mcp_session_id in manager._server_instances
791+
# Exit context -> shutdown path should terminate then clear.
792+
793+
assert terminate_calls == 1
794+
assert transport.is_terminated
795+
assert transport.mcp_session_id not in manager._server_instances
796+
assert not manager._server_instances

0 commit comments

Comments
 (0)