Clean up cancelled SSE wait tasks
This commit is contained in:
parent
ca3d34053f
commit
3f33994cdc
2 changed files with 50 additions and 14 deletions
|
|
@ -169,20 +169,23 @@ async def render_stream(
|
|||
return
|
||||
queue_task = asyncio.create_task(queue.get())
|
||||
shutdown_task = asyncio.create_task(shutdown_event.wait())
|
||||
done, pending = await asyncio.wait(
|
||||
{queue_task, shutdown_task},
|
||||
return_when=asyncio.FIRST_COMPLETED,
|
||||
)
|
||||
for task in pending:
|
||||
task.cancel()
|
||||
for task in pending:
|
||||
with suppress(asyncio.CancelledError):
|
||||
await task
|
||||
if shutdown_task in done:
|
||||
with suppress(asyncio.CancelledError):
|
||||
await queue_task
|
||||
return
|
||||
event_name = queue_task.result()
|
||||
try:
|
||||
done, _pending = await asyncio.wait(
|
||||
{queue_task, shutdown_task},
|
||||
return_when=asyncio.FIRST_COMPLETED,
|
||||
)
|
||||
if shutdown_task in done:
|
||||
return
|
||||
event_name = queue_task.result()
|
||||
finally:
|
||||
for task in (queue_task, shutdown_task):
|
||||
if not task.done():
|
||||
task.cancel()
|
||||
for task in (queue_task, shutdown_task):
|
||||
if task.done() and not task.cancelled():
|
||||
continue
|
||||
with suppress(asyncio.CancelledError):
|
||||
await task
|
||||
last_event_id, event = await render_sse_event(
|
||||
render,
|
||||
last_event_id=last_event_id,
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue