FrameProcessor: wait_for_task is now deprecated

This commit is contained in:
Aleix Conchillo Flaqué
2025-08-21 21:17:47 -07:00
parent 256ecf4d71
commit bc51e7abc6
12 changed files with 32 additions and 63 deletions

View File

@@ -48,6 +48,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
- Updated `daily-python` to 0.19.7. - Updated `daily-python` to 0.19.7.
### Deprecated
- `FrameProcessor.wait_for_task()` is deprecated. Use `await task` or `await
asyncio.wait_for(task, timeout)` instead.
### Removed ### Removed
- Watchdog timers have been removed. They were introduced in 0.0.72 to help - Watchdog timers have been removed. They were introduced in 0.0.72 to help

View File

@@ -28,7 +28,7 @@ SPEAKING_THRESHOLD = 20
def create_default_resampler(**kwargs) -> BaseAudioResampler: def create_default_resampler(**kwargs) -> BaseAudioResampler:
"""Create a default audio resampler instance. """Create a default audio resampler instance.
. deprecated:: 0.0.74 .. deprecated:: 0.0.74
This function is deprecated and will be removed in a future version. This function is deprecated and will be removed in a future version.
Use `create_stream_resampler` for real-time processing scenarios or Use `create_stream_resampler` for real-time processing scenarios or
`create_file_resampler` for batch processing of complete audio files. `create_file_resampler` for batch processing of complete audio files.

View File

@@ -363,9 +363,9 @@ class PipelineTask(BasePipelineTask):
# Create all main tasks and wait of the main push task. This is the # Create all main tasks and wait of the main push task. This is the
# task that pushes frames to the very beginning of our pipeline (our # task that pushes frames to the very beginning of our pipeline (our
# controlled PipelineTaskSource processor). # controlled source processor).
push_task = await self._create_tasks() push_task = await self._create_tasks()
await self._task_manager.wait_for_task(push_task) await push_task
# We have already cleaned up the pipeline inside the task. # We have already cleaned up the pipeline inside the task.
cleanup_pipeline = False cleanup_pipeline = False

View File

@@ -985,10 +985,6 @@ class LLMAssistantContextAggregator(LLMContextResponseAggregator):
def _context_updated_task_finished(self, task: asyncio.Task): def _context_updated_task_finished(self, task: asyncio.Task):
self._context_updated_tasks.discard(task) self._context_updated_tasks.discard(task)
# The task is finished so this should exit immediately. We need to do
# this because otherwise the task manager would report a dangling task
# if we don't remove it.
asyncio.run_coroutine_threadsafe(self.wait_for_task(task), self.get_event_loop())
class LLMUserResponseAggregator(LLMUserContextAggregator): class LLMUserResponseAggregator(LLMUserContextAggregator):

View File

@@ -442,11 +442,27 @@ class FrameProcessor(BaseObject):
async def wait_for_task(self, task: asyncio.Task, timeout: Optional[float] = None): async def wait_for_task(self, task: asyncio.Task, timeout: Optional[float] = None):
"""Wait for a task to complete. """Wait for a task to complete.
.. deprecated:: 0.0.81
This function is deprecated, use `await task` or
`await asyncio.wait_for(task, timeout) instead.
Args: Args:
task: The task to wait for. task: The task to wait for.
timeout: Optional timeout for waiting. timeout: Optional timeout for waiting.
""" """
await self.task_manager.wait_for_task(task, timeout) import warnings
warnings.warn(
"`FrameProcessor.wait_for_task()` is deprecated. "
"Use `await task` or `await asyncio.wait_for(task, timeout)` instead.",
DeprecationWarning,
stacklevel=2,
)
if timeout:
await asyncio.wait_for(task, timeout)
else:
await task
async def setup(self, setup: FrameProcessorSetup): async def setup(self, setup: FrameProcessorSetup):
"""Set up the processor with required components. """Set up the processor with required components.

View File

@@ -63,7 +63,7 @@ class SentryMetrics(FrameProcessorMetrics):
await super().cleanup() await super().cleanup()
if self._sentry_task: if self._sentry_task:
await self._sentry_queue.put(None) await self._sentry_queue.put(None)
await self.task_manager.wait_for_task(self._sentry_task) await self._sentry_task
self._sentry_task = None self._sentry_task = None
logger.trace(f"{self} Flushing Sentry metrics") logger.trace(f"{self} Flushing Sentry metrics")
sentry_sdk.flush(timeout=5.0) sentry_sdk.flush(timeout=5.0)

View File

@@ -487,7 +487,7 @@ class LLMService(AIService):
self._function_call_tasks[task] = runner_item self._function_call_tasks[task] = runner_item
# Since we run tasks sequentially we don't need to call # Since we run tasks sequentially we don't need to call
# task.add_done_callback(self._function_call_task_finished). # task.add_done_callback(self._function_call_task_finished).
await self.wait_for_task(task) await task
del self._function_call_tasks[task] del self._function_call_tasks[task]
async def _run_function_call(self, runner_item: FunctionCallRunnerItem): async def _run_function_call(self, runner_item: FunctionCallRunnerItem):
@@ -616,7 +616,3 @@ class LLMService(AIService):
def _function_call_task_finished(self, task: asyncio.Task): def _function_call_task_finished(self, task: asyncio.Task):
if task in self._function_call_tasks: if task in self._function_call_tasks:
del self._function_call_tasks[task] del self._function_call_tasks[task]
# The task is finished so this should exit immediately. We need to
# do this because otherwise the task manager would report a dangling
# task if we don't remove it.
asyncio.run_coroutine_threadsafe(self.wait_for_task(task), self.get_event_loop())

View File

@@ -208,7 +208,7 @@ class SonioxSTTService(STTService):
if self._receive_task: if self._receive_task:
# Task cannot cancel itself. If task called _cleanup() we expect it to cancel itself. # Task cannot cancel itself. If task called _cleanup() we expect it to cancel itself.
if self._receive_task != asyncio.current_task(): if self._receive_task != asyncio.current_task():
await self.wait_for_task(self._receive_task) await self._receive_task
self._receive_task = None self._receive_task = None
async def stop(self, frame: EndFrame): async def stop(self, frame: EndFrame):

View File

@@ -794,7 +794,7 @@ class AudioContextWordTTSService(WebsocketWordTTSService):
# Indicate no more audio contexts are available. this will end the # Indicate no more audio contexts are available. this will end the
# task cleanly after all contexts have been processed. # task cleanly after all contexts have been processed.
await self._contexts_queue.put(None) await self._contexts_queue.put(None)
await self.wait_for_task(self._audio_context_task) await self._audio_context_task
self._audio_context_task = None self._audio_context_task = None
async def cancel(self, frame: CancelFrame): async def cancel(self, frame: CancelFrame):

View File

@@ -435,9 +435,9 @@ class BaseOutputTransport(FrameProcessor):
# also need to wait for these tasks before cancelling the video task # also need to wait for these tasks before cancelling the video task
# because it might be still rendering. # because it might be still rendering.
if self._audio_task: if self._audio_task:
await self._transport.wait_for_task(self._audio_task) await self._audio_task
if self._clock_task: if self._clock_task:
await self._transport.wait_for_task(self._clock_task) await self._clock_task
# Stop audio mixer. # Stop audio mixer.
if self._mixer: if self._mixer:

View File

@@ -154,7 +154,7 @@ class WebsocketServerInputTransport(BaseInputTransport):
await self.cancel_task(self._monitor_task) await self.cancel_task(self._monitor_task)
self._monitor_task = None self._monitor_task = None
if self._server_task: if self._server_task:
await self.wait_for_task(self._server_task) await self._server_task
self._server_task = None self._server_task = None
async def cancel(self, frame: CancelFrame): async def cancel(self, frame: CancelFrame):

View File

@@ -69,21 +69,6 @@ class BaseTaskManager(ABC):
""" """
pass pass
@abstractmethod
async def wait_for_task(self, task: asyncio.Task, timeout: Optional[float] = None):
"""Wait for an asyncio.Task to complete with optional timeout handling.
This function awaits the specified asyncio.Task and handles scenarios for
timeouts, cancellations, and other exceptions. It also ensures that the task
is removed from the set of registered tasks upon completion or failure.
Args:
task: The asyncio Task to wait for.
timeout: The maximum number of seconds to wait for the task to complete.
If None, waits indefinitely.
"""
pass
@abstractmethod @abstractmethod
async def cancel_task(self, task: asyncio.Task, timeout: Optional[float] = None): async def cancel_task(self, task: asyncio.Task, timeout: Optional[float] = None):
"""Cancels the given asyncio Task and awaits its completion with an optional timeout. """Cancels the given asyncio Task and awaits its completion with an optional timeout.
@@ -189,35 +174,6 @@ class TaskManager(BaseTaskManager):
logger.trace(f"{name}: task created") logger.trace(f"{name}: task created")
return task return task
async def wait_for_task(self, task: asyncio.Task, timeout: Optional[float] = None):
"""Wait for an asyncio.Task to complete with optional timeout handling.
This function awaits the specified asyncio.Task and handles scenarios for
timeouts, cancellations, and other exceptions. It also ensures that the task
is removed from the set of registered tasks upon completion or failure.
Args:
task: The asyncio Task to wait for.
timeout: The maximum number of seconds to wait for the task to complete.
If None, waits indefinitely.
"""
name = task.get_name()
try:
if timeout:
await asyncio.wait_for(task, timeout=timeout)
else:
await task
except asyncio.TimeoutError:
logger.warning(f"{name}: timed out waiting for task to finish")
except asyncio.CancelledError:
logger.trace(f"{name}: unexpected task cancellation (maybe Ctrl-C?)")
raise
except Exception as e:
logger.exception(f"{name}: unexpected exception while stopping task: {e}")
except BaseException as e:
logger.critical(f"{name}: fatal base exception while stopping task: {e}")
raise
async def cancel_task(self, task: asyncio.Task, timeout: Optional[float] = None): async def cancel_task(self, task: asyncio.Task, timeout: Optional[float] = None):
"""Cancels the given asyncio Task and awaits its completion with an optional timeout. """Cancels the given asyncio Task and awaits its completion with an optional timeout.