refac
This commit is contained in:
@@ -1923,8 +1923,6 @@ async def chat_completion(
|
||||
# uses anyio task groups whose cancel scopes enforce
|
||||
# same-task exit. Do NOT wrap in asyncio.shield() or
|
||||
# asyncio.wait_for() — both create a new task.
|
||||
# MCPClient.disconnect() self-shields via
|
||||
# anyio.CancelScope(shield=True).
|
||||
try:
|
||||
if mcp_clients := metadata.get('mcp_clients'):
|
||||
for client in reversed(list(mcp_clients.values())):
|
||||
@@ -1932,14 +1930,29 @@ async def chat_completion(
|
||||
await client.disconnect()
|
||||
except Exception as e:
|
||||
log.debug(f'Error disconnecting MCP client: {e}')
|
||||
except asyncio.CancelledError:
|
||||
# Let the client close asynchronously by GC
|
||||
pass
|
||||
except Exception as e:
|
||||
log.debug(f'Error cleaning up MCP clients: {e}')
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
|
||||
try:
|
||||
if metadata.get('chat_id'):
|
||||
event_emitter = await get_event_emitter(metadata, update_db=False)
|
||||
if event_emitter:
|
||||
await event_emitter({'type': 'chat:active', 'data': {'active': False}})
|
||||
async def emit_inactive_event():
|
||||
try:
|
||||
event_emitter = await get_event_emitter(metadata, update_db=False)
|
||||
if event_emitter:
|
||||
await event_emitter({'type': 'chat:active', 'data': {'active': False}})
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
try:
|
||||
# Shield the event emission so it finishes even if the main task is cancelled
|
||||
await asyncio.shield(emit_inactive_event())
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
@@ -156,14 +156,14 @@ class MCPClient:
|
||||
|
||||
try:
|
||||
# IMPORTANT: Do NOT use asyncio.shield() or asyncio.wait_for()
|
||||
# here — both create a new asyncio task. The MCP SDK's
|
||||
# streamablehttp_client uses anyio task groups / cancel scopes
|
||||
# that MUST be exited in the same task they were entered in.
|
||||
# Using anyio.CancelScope(shield=True) protects from
|
||||
# CancelledError while staying in the current task.
|
||||
with anyio.CancelScope(shield=True):
|
||||
with anyio.fail_after(5.0):
|
||||
await exit_stack.aclose()
|
||||
# because they create a new asyncio task, which violates the MCP SDK's
|
||||
# requirement that its TaskGroup be exited in the exact same task.
|
||||
# ALSO do NOT use anyio.CancelScope(shield=True) or anyio.fail_after(),
|
||||
# because they push a new cancel scope onto the task, violating LIFO
|
||||
# order when aclose() attempts to exit the inner TaskGroup.
|
||||
# We simply call aclose() directly. If the task is cancelled, the
|
||||
# sockets will eventually be cleaned up by garbage collection.
|
||||
await exit_stack.aclose()
|
||||
except TimeoutError:
|
||||
log.warning('MCPClient.disconnect() timed out after 5 s')
|
||||
except RuntimeError as exc:
|
||||
|
||||
@@ -22,6 +22,7 @@ aiocache
|
||||
aiofiles
|
||||
starlette-compress==1.7.0
|
||||
Brotli==1.2.0
|
||||
brotlicffi==1.2.0.1
|
||||
httpx[socks,http2,zstd,cli,brotli]==0.28.1
|
||||
starsessions[redis]==2.2.1
|
||||
|
||||
|
||||
Reference in New Issue
Block a user