Merge pull request #505 from pipecat-ai/aleix/init-task-variables
initialize task variables and add minor description
This commit is contained in:
@@ -5,7 +5,6 @@
|
|||||||
#
|
#
|
||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
import time
|
|
||||||
|
|
||||||
from enum import Enum
|
from enum import Enum
|
||||||
|
|
||||||
|
|||||||
@@ -423,18 +423,26 @@ class RTVIProcessor(FrameProcessor):
|
|||||||
await self._maybe_send_bot_ready()
|
await self._maybe_send_bot_ready()
|
||||||
|
|
||||||
async def _stop(self, frame: EndFrame):
|
async def _stop(self, frame: EndFrame):
|
||||||
self._action_task.cancel()
|
if self._action_task:
|
||||||
await self._action_task
|
self._action_task.cancel()
|
||||||
|
await self._action_task
|
||||||
|
self._action_task = None
|
||||||
|
|
||||||
self._message_task.cancel()
|
if self._message_task:
|
||||||
await self._message_task
|
self._message_task.cancel()
|
||||||
|
await self._message_task
|
||||||
|
self._message_task = None
|
||||||
|
|
||||||
async def _cancel(self, frame: CancelFrame):
|
async def _cancel(self, frame: CancelFrame):
|
||||||
self._action_task.cancel()
|
if self._action_task:
|
||||||
await self._action_task
|
self._action_task.cancel()
|
||||||
|
await self._action_task
|
||||||
|
self._action_task = None
|
||||||
|
|
||||||
self._message_task.cancel()
|
if self._message_task:
|
||||||
await self._message_task
|
self._message_task.cancel()
|
||||||
|
await self._message_task
|
||||||
|
self._message_task = None
|
||||||
|
|
||||||
async def _push_transport_message(self, model: BaseModel, exclude_none: bool = True):
|
async def _push_transport_message(self, model: BaseModel, exclude_none: bool = True):
|
||||||
frame = TransportMessageFrame(
|
frame = TransportMessageFrame(
|
||||||
|
|||||||
@@ -350,6 +350,7 @@ class AsyncWordTTSService(AsyncTTSService):
|
|||||||
if self._words_task:
|
if self._words_task:
|
||||||
self._words_task.cancel()
|
self._words_task.cancel()
|
||||||
await self._words_task
|
await self._words_task
|
||||||
|
self._words_task = None
|
||||||
|
|
||||||
async def _words_task_handler(self):
|
async def _words_task_handler(self):
|
||||||
while True:
|
while True:
|
||||||
|
|||||||
@@ -37,6 +37,10 @@ class BaseInputTransport(FrameProcessor):
|
|||||||
|
|
||||||
self._executor = ThreadPoolExecutor(max_workers=5)
|
self._executor = ThreadPoolExecutor(max_workers=5)
|
||||||
|
|
||||||
|
# Task to process incoming audio (VAD) and push audio frames downstream
|
||||||
|
# if passthrough is enabled.
|
||||||
|
self._audio_task = None
|
||||||
|
|
||||||
async def start(self, frame: StartFrame):
|
async def start(self, frame: StartFrame):
|
||||||
# Create audio input queue and task if needed.
|
# Create audio input queue and task if needed.
|
||||||
if self._params.audio_in_enabled or self._params.vad_enabled:
|
if self._params.audio_in_enabled or self._params.vad_enabled:
|
||||||
@@ -45,16 +49,17 @@ class BaseInputTransport(FrameProcessor):
|
|||||||
|
|
||||||
async def stop(self, frame: EndFrame):
|
async def stop(self, frame: EndFrame):
|
||||||
# Cancel and wait for the audio input task to finish.
|
# Cancel and wait for the audio input task to finish.
|
||||||
if self._params.audio_in_enabled or self._params.vad_enabled:
|
if self._audio_task and (self._params.audio_in_enabled or self._params.vad_enabled):
|
||||||
self._audio_task.cancel()
|
self._audio_task.cancel()
|
||||||
await self._audio_task
|
await self._audio_task
|
||||||
|
self._audio_task = None
|
||||||
|
|
||||||
async def cancel(self, frame: CancelFrame):
|
async def cancel(self, frame: CancelFrame):
|
||||||
# Cancel all the tasks and wait for them to finish.
|
# Cancel and wait for the audio input task to finish.
|
||||||
|
if self._audio_task and (self._params.audio_in_enabled or self._params.vad_enabled):
|
||||||
if self._params.audio_in_enabled or self._params.vad_enabled:
|
|
||||||
self._audio_task.cancel()
|
self._audio_task.cancel()
|
||||||
await self._audio_task
|
await self._audio_task
|
||||||
|
self._audio_task = None
|
||||||
|
|
||||||
def vad_analyzer(self) -> VADAnalyzer | None:
|
def vad_analyzer(self) -> VADAnalyzer | None:
|
||||||
return self._params.vad_analyzer
|
return self._params.vad_analyzer
|
||||||
|
|||||||
@@ -47,6 +47,18 @@ class BaseOutputTransport(FrameProcessor):
|
|||||||
|
|
||||||
self._params = params
|
self._params = params
|
||||||
|
|
||||||
|
# Task to process incoming frames so we don't block upstream elements.
|
||||||
|
self._sink_task = None
|
||||||
|
|
||||||
|
# Task to process incoming frames using a clock.
|
||||||
|
self._sink_clock_task = None
|
||||||
|
|
||||||
|
# Task to write/send audio frames.
|
||||||
|
self._audio_out_task = None
|
||||||
|
|
||||||
|
# Task to write/send image frames.
|
||||||
|
self._camera_out_task = None
|
||||||
|
|
||||||
# These are the images that we should send to the camera at our desired
|
# These are the images that we should send to the camera at our desired
|
||||||
# framerate.
|
# framerate.
|
||||||
self._camera_images = None
|
self._camera_images = None
|
||||||
@@ -88,36 +100,53 @@ class BaseOutputTransport(FrameProcessor):
|
|||||||
# that EndFrame to be processed by the sink tasks. We also need to wait
|
# that EndFrame to be processed by the sink tasks. We also need to wait
|
||||||
# for these tasks before cancelling the camera and audio tasks below
|
# for these tasks before cancelling the camera and audio tasks below
|
||||||
# because they might be still rendering.
|
# because they might be still rendering.
|
||||||
await self._sink_task
|
if self._sink_task:
|
||||||
await self._sink_clock_task
|
await self._sink_task
|
||||||
|
if self._sink_clock_task:
|
||||||
|
await self._sink_clock_task
|
||||||
|
|
||||||
# Cancel and wait for the camera output task to finish.
|
# Cancel and wait for the camera output task to finish.
|
||||||
if self._params.camera_out_enabled:
|
if self._camera_out_task and self._params.camera_out_enabled:
|
||||||
self._camera_out_task.cancel()
|
self._camera_out_task.cancel()
|
||||||
await self._camera_out_task
|
await self._camera_out_task
|
||||||
|
self._camera_out_task = None
|
||||||
|
|
||||||
# Cancel and wait for the audio output task to finish.
|
# Cancel and wait for the audio output task to finish.
|
||||||
if self._params.audio_out_enabled and self._params.audio_out_is_live:
|
if (
|
||||||
|
self._audio_out_task
|
||||||
|
and self._params.audio_out_enabled
|
||||||
|
and self._params.audio_out_is_live
|
||||||
|
):
|
||||||
self._audio_out_task.cancel()
|
self._audio_out_task.cancel()
|
||||||
await self._audio_out_task
|
await self._audio_out_task
|
||||||
|
self._audio_out_task = None
|
||||||
|
|
||||||
async def cancel(self, frame: CancelFrame):
|
async def cancel(self, frame: CancelFrame):
|
||||||
# Since we are cancelling everything it doesn't matter if we cancel sink
|
# Since we are cancelling everything it doesn't matter if we cancel sink
|
||||||
# tasks first or not.
|
# tasks first or not.
|
||||||
self._sink_task.cancel()
|
if self._sink_task:
|
||||||
self._sink_clock_task.cancel()
|
self._sink_task.cancel()
|
||||||
await self._sink_task
|
await self._sink_task
|
||||||
await self._sink_clock_task
|
self._sink_task = None
|
||||||
|
|
||||||
|
if self._sink_clock_task:
|
||||||
|
self._sink_clock_task.cancel()
|
||||||
|
await self._sink_clock_task
|
||||||
|
self._sink_clock_task = None
|
||||||
|
|
||||||
# Cancel and wait for the camera output task to finish.
|
# Cancel and wait for the camera output task to finish.
|
||||||
if self._params.camera_out_enabled:
|
if self._camera_out_task and self._params.camera_out_enabled:
|
||||||
self._camera_out_task.cancel()
|
self._camera_out_task.cancel()
|
||||||
await self._camera_out_task
|
await self._camera_out_task
|
||||||
|
self._camera_out_task = None
|
||||||
|
|
||||||
# Cancel and wait for the audio output task to finish.
|
# Cancel and wait for the audio output task to finish.
|
||||||
if self._params.audio_out_enabled and self._params.audio_out_is_live:
|
if self._audio_out_task and (
|
||||||
|
self._params.audio_out_enabled and self._params.audio_out_is_live
|
||||||
|
):
|
||||||
self._audio_out_task.cancel()
|
self._audio_out_task.cancel()
|
||||||
await self._audio_out_task
|
await self._audio_out_task
|
||||||
|
self._audio_out_task = None
|
||||||
|
|
||||||
async def send_message(self, frame: TransportMessageFrame):
|
async def send_message(self, frame: TransportMessageFrame):
|
||||||
pass
|
pass
|
||||||
@@ -183,11 +212,13 @@ class BaseOutputTransport(FrameProcessor):
|
|||||||
|
|
||||||
if isinstance(frame, StartInterruptionFrame):
|
if isinstance(frame, StartInterruptionFrame):
|
||||||
# Stop sink tasks.
|
# Stop sink tasks.
|
||||||
self._sink_task.cancel()
|
if self._sink_task:
|
||||||
await self._sink_task
|
self._sink_task.cancel()
|
||||||
|
await self._sink_task
|
||||||
# Stop sink clock tasks.
|
# Stop sink clock tasks.
|
||||||
self._sink_clock_task.cancel()
|
if self._sink_clock_task:
|
||||||
await self._sink_clock_task
|
self._sink_clock_task.cancel()
|
||||||
|
await self._sink_clock_task
|
||||||
# Create sink tasks.
|
# Create sink tasks.
|
||||||
self._create_sink_tasks()
|
self._create_sink_tasks()
|
||||||
# Let's send a bot stopped speaking if we have to.
|
# Let's send a bot stopped speaking if we have to.
|
||||||
|
|||||||
@@ -575,6 +575,9 @@ class DailyInputTransport(BaseInputTransport):
|
|||||||
self._client = client
|
self._client = client
|
||||||
|
|
||||||
self._video_renderers = {}
|
self._video_renderers = {}
|
||||||
|
|
||||||
|
# Task that gets audio data from a device or the network and queues it
|
||||||
|
# internally to be processed.
|
||||||
self._audio_in_task = None
|
self._audio_in_task = None
|
||||||
|
|
||||||
self._vad_analyzer: VADAnalyzer | None = params.vad_analyzer
|
self._vad_analyzer: VADAnalyzer | None = params.vad_analyzer
|
||||||
@@ -603,6 +606,7 @@ class DailyInputTransport(BaseInputTransport):
|
|||||||
if self._audio_in_task and (self._params.audio_in_enabled or self._params.vad_enabled):
|
if self._audio_in_task and (self._params.audio_in_enabled or self._params.vad_enabled):
|
||||||
self._audio_in_task.cancel()
|
self._audio_in_task.cancel()
|
||||||
await self._audio_in_task
|
await self._audio_in_task
|
||||||
|
self._audio_in_task = None
|
||||||
|
|
||||||
async def cancel(self, frame: CancelFrame):
|
async def cancel(self, frame: CancelFrame):
|
||||||
# Parent stop.
|
# Parent stop.
|
||||||
@@ -613,6 +617,7 @@ class DailyInputTransport(BaseInputTransport):
|
|||||||
if self._audio_in_task and (self._params.audio_in_enabled or self._params.vad_enabled):
|
if self._audio_in_task and (self._params.audio_in_enabled or self._params.vad_enabled):
|
||||||
self._audio_in_task.cancel()
|
self._audio_in_task.cancel()
|
||||||
await self._audio_in_task
|
await self._audio_in_task
|
||||||
|
self._audio_in_task = None
|
||||||
|
|
||||||
async def cleanup(self):
|
async def cleanup(self):
|
||||||
await super().cleanup()
|
await super().cleanup()
|
||||||
|
|||||||
Reference in New Issue
Block a user