ParallelPipeline: wait for CancelFrame in all branches
This commit is contained in:
@@ -79,6 +79,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
|||||||
|
|
||||||
### Changed
|
### Changed
|
||||||
|
|
||||||
|
- `ParallelPipeline` now waits for `CancelFrame` to finish in all branches
|
||||||
|
before pushing it downstream.
|
||||||
|
|
||||||
- Added `sip_codecs` to the `DailyRoomSipParams`.
|
- Added `sip_codecs` to the `DailyRoomSipParams`.
|
||||||
|
|
||||||
- Updated the `configure()` function in `pipecat.runner.daily` to include new
|
- Updated the `configure()` function in `pipecat.runner.daily` to include new
|
||||||
|
|||||||
@@ -16,7 +16,7 @@ from typing import Dict, List
|
|||||||
|
|
||||||
from loguru import logger
|
from loguru import logger
|
||||||
|
|
||||||
from pipecat.frames.frames import EndFrame, Frame, StartFrame
|
from pipecat.frames.frames import CancelFrame, EndFrame, Frame, StartFrame
|
||||||
from pipecat.pipeline.base_pipeline import BasePipeline
|
from pipecat.pipeline.base_pipeline import BasePipeline
|
||||||
from pipecat.pipeline.pipeline import Pipeline, PipelineSink, PipelineSource
|
from pipecat.pipeline.pipeline import Pipeline, PipelineSink, PipelineSource
|
||||||
from pipecat.processors.frame_processor import FrameDirection, FrameProcessor, FrameProcessorSetup
|
from pipecat.processors.frame_processor import FrameDirection, FrameProcessor, FrameProcessorSetup
|
||||||
@@ -141,7 +141,7 @@ class ParallelPipeline(BasePipeline):
|
|||||||
await super().process_frame(frame, direction)
|
await super().process_frame(frame, direction)
|
||||||
|
|
||||||
# Parallel pipeline synchronized frames.
|
# Parallel pipeline synchronized frames.
|
||||||
if isinstance(frame, (StartFrame, EndFrame)):
|
if isinstance(frame, (StartFrame, EndFrame, CancelFrame)):
|
||||||
self._frame_counter[frame.id] = len(self._pipelines)
|
self._frame_counter[frame.id] = len(self._pipelines)
|
||||||
await self.pause_processing_system_frames()
|
await self.pause_processing_system_frames()
|
||||||
await self.pause_processing_frames()
|
await self.pause_processing_frames()
|
||||||
@@ -158,7 +158,7 @@ class ParallelPipeline(BasePipeline):
|
|||||||
|
|
||||||
async def _pipeline_sink_push_frame(self, frame: Frame, direction: FrameDirection):
|
async def _pipeline_sink_push_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
# Parallel pipeline synchronized frames.
|
# Parallel pipeline synchronized frames.
|
||||||
if isinstance(frame, (StartFrame, EndFrame)):
|
if isinstance(frame, (StartFrame, EndFrame, CancelFrame)):
|
||||||
# Decrement counter.
|
# Decrement counter.
|
||||||
frame_counter = self._frame_counter.get(frame.id, 0)
|
frame_counter = self._frame_counter.get(frame.id, 0)
|
||||||
if frame_counter > 0:
|
if frame_counter > 0:
|
||||||
|
|||||||
Reference in New Issue
Block a user