fix: Fix some swallowed CancelledError (#6625)
Fix some swallowed CancelledError Co-authored-by: Gabriel Luiz Freitas Almeida <gabriel@langflow.org>
This commit is contained in:
parent
c7e2e1eca8
commit
23244fe8c2
2 changed files with 8 additions and 16 deletions
|
|
@ -347,14 +347,11 @@ async def build_flow(
|
||||||
client_consumed_queue: asyncio.Queue,
|
client_consumed_queue: asyncio.Queue,
|
||||||
event_manager: EventManager,
|
event_manager: EventManager,
|
||||||
) -> None:
|
) -> None:
|
||||||
build_task = asyncio.create_task(_build_vertex(vertex_id, graph, event_manager))
|
|
||||||
try:
|
try:
|
||||||
await build_task
|
vertex_build_response: VertexBuildResponse = await _build_vertex(vertex_id, graph, event_manager)
|
||||||
vertex_build_response: VertexBuildResponse = build_task.result()
|
|
||||||
except asyncio.CancelledError as exc:
|
except asyncio.CancelledError as exc:
|
||||||
logger.exception(exc)
|
logger.exception(exc)
|
||||||
build_task.cancel()
|
raise
|
||||||
return
|
|
||||||
|
|
||||||
# send built event or error event
|
# send built event or error event
|
||||||
try:
|
try:
|
||||||
|
|
@ -370,12 +367,7 @@ async def build_flow(
|
||||||
for next_vertex_id in vertex_build_response.next_vertices_ids:
|
for next_vertex_id in vertex_build_response.next_vertices_ids:
|
||||||
task = asyncio.create_task(build_vertices(next_vertex_id, graph, client_consumed_queue, event_manager))
|
task = asyncio.create_task(build_vertices(next_vertex_id, graph, client_consumed_queue, event_manager))
|
||||||
tasks.append(task)
|
tasks.append(task)
|
||||||
try:
|
|
||||||
await asyncio.gather(*tasks)
|
await asyncio.gather(*tasks)
|
||||||
except asyncio.CancelledError:
|
|
||||||
for task in tasks:
|
|
||||||
task.cancel()
|
|
||||||
return
|
|
||||||
|
|
||||||
async def event_generator(event_manager: EventManager, client_consumed_queue: asyncio.Queue) -> None:
|
async def event_generator(event_manager: EventManager, client_consumed_queue: asyncio.Queue) -> None:
|
||||||
try:
|
try:
|
||||||
|
|
@ -397,9 +389,8 @@ async def build_flow(
|
||||||
await asyncio.gather(*tasks)
|
await asyncio.gather(*tasks)
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
background_tasks.add_task(graph.end_all_traces)
|
background_tasks.add_task(graph.end_all_traces)
|
||||||
for task in tasks:
|
raise
|
||||||
task.cancel()
|
|
||||||
return
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"Error building vertices: {e}")
|
logger.error(f"Error building vertices: {e}")
|
||||||
custom_component = graph.get_vertex(vertex_id).custom_component
|
custom_component = graph.get_vertex(vertex_id).custom_component
|
||||||
|
|
|
||||||
|
|
@ -3,7 +3,6 @@ import base64
|
||||||
import json
|
import json
|
||||||
import logging
|
import logging
|
||||||
import traceback
|
import traceback
|
||||||
from contextlib import suppress
|
|
||||||
from contextvars import ContextVar
|
from contextvars import ContextVar
|
||||||
from typing import Annotated
|
from typing import Annotated
|
||||||
from urllib.parse import quote, unquote, urlparse
|
from urllib.parse import quote, unquote, urlparse
|
||||||
|
|
@ -266,8 +265,9 @@ async def handle_call_tool(
|
||||||
return collected_results
|
return collected_results
|
||||||
finally:
|
finally:
|
||||||
progress_task.cancel()
|
progress_task.cancel()
|
||||||
with suppress(asyncio.CancelledError):
|
await asyncio.wait([progress_task])
|
||||||
await progress_task
|
if not progress_task.cancelled() and (exc := progress_task.exception()) is not None:
|
||||||
|
raise exc
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
msg = f"Error in async session: {e}"
|
msg = f"Error in async session: {e}"
|
||||||
logger.exception(msg)
|
logger.exception(msg)
|
||||||
|
|
@ -333,6 +333,7 @@ async def handle_sse(request: Request, current_user: Annotated[User, Depends(get
|
||||||
logger.info("Client disconnected from SSE connection")
|
logger.info("Client disconnected from SSE connection")
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
logger.info("SSE connection was cancelled")
|
logger.info("SSE connection was cancelled")
|
||||||
|
raise
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
msg = f"Error in MCP: {e!s}"
|
msg = f"Error in MCP: {e!s}"
|
||||||
logger.exception(msg)
|
logger.exception(msg)
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue