fix(mcp): widen #59349 handshake bound to HTTP transports + cancel abandoned start() task
Sibling sites of the same bug class as the salvaged stdio fix: - SSE, streamable-HTTP (new + deprecated API) initialize() calls are now bounded by the same connect_timeout, so an endpoint that accepts the connection but never answers the handshake cannot park the run() task forever. - start() now cancels its ensure_future'd run() task when the caller's connect timeout cancels start() itself — the orphaned-task leak was the root mechanism behind #59349, and this closes the class for any future pre-ready hang.
This commit is contained in:
parent
1f6836cd81
commit
4638f3b433
1 changed files with 28 additions and 4 deletions
|
|
@ -2312,7 +2312,13 @@ class MCPServerTask:
|
|||
async with ClientSession(
|
||||
read_stream, write_stream, **sampling_kwargs
|
||||
) as session:
|
||||
self.initialize_result = await session.initialize()
|
||||
# Bound the handshake — same orphaned-task hang as the
|
||||
# stdio path (#59349): an endpoint that accepts the
|
||||
# connection but never answers ``initialize`` parks this
|
||||
# coroutine forever on the background loop.
|
||||
self.initialize_result = await asyncio.wait_for(
|
||||
session.initialize(), timeout=float(connect_timeout)
|
||||
)
|
||||
self.session = session
|
||||
await self._discover_tools()
|
||||
self._ready.set()
|
||||
|
|
@ -2366,7 +2372,10 @@ class MCPServerTask:
|
|||
read_stream, write_stream, _get_session_id,
|
||||
):
|
||||
async with ClientSession(read_stream, write_stream, **sampling_kwargs) as session:
|
||||
self.initialize_result = await session.initialize()
|
||||
# Bound the handshake (#59349) — see stdio path.
|
||||
self.initialize_result = await asyncio.wait_for(
|
||||
session.initialize(), timeout=float(connect_timeout)
|
||||
)
|
||||
self.session = session
|
||||
await self._discover_tools()
|
||||
self._ready.set()
|
||||
|
|
@ -2394,7 +2403,10 @@ class MCPServerTask:
|
|||
read_stream, write_stream, _get_session_id,
|
||||
):
|
||||
async with ClientSession(read_stream, write_stream, **sampling_kwargs) as session:
|
||||
self.initialize_result = await session.initialize()
|
||||
# Bound the handshake (#59349) — see stdio path.
|
||||
self.initialize_result = await asyncio.wait_for(
|
||||
session.initialize(), timeout=float(connect_timeout)
|
||||
)
|
||||
self.session = session
|
||||
await self._discover_tools()
|
||||
self._ready.set()
|
||||
|
|
@ -2715,7 +2727,19 @@ class MCPServerTask:
|
|||
async def start(self, config: dict):
|
||||
"""Create the background Task and wait until ready (or failed)."""
|
||||
self._task = asyncio.ensure_future(self.run(config))
|
||||
await self._ready.wait()
|
||||
try:
|
||||
await self._ready.wait()
|
||||
except asyncio.CancelledError:
|
||||
# The caller's connect timeout (discover_mcp_tools wraps start()
|
||||
# in asyncio.wait_for) cancels *this* coroutine, but the
|
||||
# ensure_future'd run() task is independent and would otherwise
|
||||
# keep running detached — parked on a hung transport with no
|
||||
# owner to reap it (#59349). Propagate the cancellation so the
|
||||
# transport context managers unwind and their finally blocks
|
||||
# release the child process / FDs.
|
||||
if self._task and not self._task.done():
|
||||
self._task.cancel()
|
||||
raise
|
||||
if self._error:
|
||||
raise self._error
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue