pipeline(task): cleanup processors only if we need to
This commit is contained in:
@@ -19,6 +19,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
|||||||
|
|
||||||
- Fixed issues with Ctrl-C program termination.
|
- Fixed issues with Ctrl-C program termination.
|
||||||
|
|
||||||
|
- Fixed an issue that was causing `StopTaskFrame` to actually not exit the
|
||||||
|
`PipelineTask`.
|
||||||
|
|
||||||
## [0.0.16] - 2024-05-16
|
## [0.0.16] - 2024-05-16
|
||||||
|
|
||||||
### Fixed
|
### Fixed
|
||||||
|
|||||||
@@ -61,8 +61,6 @@ class PipelineTask:
|
|||||||
self._process_up_task = asyncio.create_task(self._process_up_queue())
|
self._process_up_task = asyncio.create_task(self._process_up_queue())
|
||||||
self._process_down_task = asyncio.create_task(self._process_down_queue())
|
self._process_down_task = asyncio.create_task(self._process_down_queue())
|
||||||
await asyncio.gather(self._process_up_task, self._process_down_task)
|
await asyncio.gather(self._process_up_task, self._process_down_task)
|
||||||
await self._source.cleanup()
|
|
||||||
await self._pipeline.cleanup()
|
|
||||||
|
|
||||||
async def queue_frame(self, frame: Frame):
|
async def queue_frame(self, frame: Frame):
|
||||||
await self._down_queue.put(frame)
|
await self._down_queue.put(frame)
|
||||||
@@ -81,14 +79,20 @@ class PipelineTask:
|
|||||||
await self._source.process_frame(
|
await self._source.process_frame(
|
||||||
StartFrame(allow_interruptions=self._allow_interruptions), FrameDirection.DOWNSTREAM)
|
StartFrame(allow_interruptions=self._allow_interruptions), FrameDirection.DOWNSTREAM)
|
||||||
running = True
|
running = True
|
||||||
|
should_cleanup = True
|
||||||
while running:
|
while running:
|
||||||
try:
|
try:
|
||||||
frame = await self._down_queue.get()
|
frame = await self._down_queue.get()
|
||||||
await self._source.process_frame(frame, FrameDirection.DOWNSTREAM)
|
await self._source.process_frame(frame, FrameDirection.DOWNSTREAM)
|
||||||
self._down_queue.task_done()
|
|
||||||
running = not (isinstance(frame, StopTaskFrame) or isinstance(frame, EndFrame))
|
running = not (isinstance(frame, StopTaskFrame) or isinstance(frame, EndFrame))
|
||||||
|
should_cleanup = not isinstance(frame, StopTaskFrame)
|
||||||
|
self._down_queue.task_done()
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
break
|
break
|
||||||
|
# Cleanup only if we need to.
|
||||||
|
if should_cleanup:
|
||||||
|
await self._source.cleanup()
|
||||||
|
await self._pipeline.cleanup()
|
||||||
# We just enqueue None to terminate the task gracefully.
|
# We just enqueue None to terminate the task gracefully.
|
||||||
self._process_up_task.cancel()
|
self._process_up_task.cancel()
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user