task: added BaseTask
This commit is contained in:
56
src/pipecat/pipeline/base_task.py
Normal file
56
src/pipecat/pipeline/base_task.py
Normal file
@@ -0,0 +1,56 @@
|
|||||||
|
#
|
||||||
|
# Copyright (c) 2024–2025, Daily
|
||||||
|
#
|
||||||
|
# SPDX-License-Identifier: BSD 2-Clause License
|
||||||
|
#
|
||||||
|
|
||||||
|
from abc import ABC, abstractmethod
|
||||||
|
from typing import AsyncIterable, Iterable
|
||||||
|
|
||||||
|
from pipecat.frames.frames import Frame
|
||||||
|
|
||||||
|
|
||||||
|
class BaseTask(ABC):
|
||||||
|
@abstractmethod
|
||||||
|
def has_finished(self) -> bool:
|
||||||
|
"""Indicates whether the tasks has finished. That is, all processors
|
||||||
|
have stopped.
|
||||||
|
|
||||||
|
"""
|
||||||
|
pass
|
||||||
|
|
||||||
|
@abstractmethod
|
||||||
|
async def stop_when_done(self):
|
||||||
|
"""This is a helper function that sends an EndFrame to the pipeline in
|
||||||
|
order to stop the task after everything in it has been processed.
|
||||||
|
|
||||||
|
"""
|
||||||
|
pass
|
||||||
|
|
||||||
|
@abstractmethod
|
||||||
|
async def cancel(self):
|
||||||
|
"""
|
||||||
|
Stops the running pipeline immediately.
|
||||||
|
"""
|
||||||
|
pass
|
||||||
|
|
||||||
|
@abstractmethod
|
||||||
|
async def run(self):
|
||||||
|
"""
|
||||||
|
Starts running the given pipeline.
|
||||||
|
"""
|
||||||
|
pass
|
||||||
|
|
||||||
|
@abstractmethod
|
||||||
|
async def queue_frame(self, frame: Frame):
|
||||||
|
"""
|
||||||
|
Queue a frame to be pushed down the pipeline.
|
||||||
|
"""
|
||||||
|
pass
|
||||||
|
|
||||||
|
@abstractmethod
|
||||||
|
async def queue_frames(self, frames: Iterable[Frame] | AsyncIterable[Frame]):
|
||||||
|
"""
|
||||||
|
Queues multiple frames to be pushed down the pipeline.
|
||||||
|
"""
|
||||||
|
pass
|
||||||
@@ -27,6 +27,7 @@ from pipecat.frames.frames import (
|
|||||||
from pipecat.metrics.metrics import ProcessingMetricsData, TTFBMetricsData
|
from pipecat.metrics.metrics import ProcessingMetricsData, TTFBMetricsData
|
||||||
from pipecat.observers.base_observer import BaseObserver
|
from pipecat.observers.base_observer import BaseObserver
|
||||||
from pipecat.pipeline.base_pipeline import BasePipeline
|
from pipecat.pipeline.base_pipeline import BasePipeline
|
||||||
|
from pipecat.pipeline.base_task import BaseTask
|
||||||
from pipecat.pipeline.task_observer import TaskObserver
|
from pipecat.pipeline.task_observer import TaskObserver
|
||||||
from pipecat.processors.frame_processor import FrameDirection, FrameProcessor
|
from pipecat.processors.frame_processor import FrameDirection, FrameProcessor
|
||||||
from pipecat.utils.utils import obj_count, obj_id
|
from pipecat.utils.utils import obj_count, obj_id
|
||||||
@@ -86,7 +87,7 @@ class Sink(FrameProcessor):
|
|||||||
await self._down_queue.put(frame)
|
await self._down_queue.put(frame)
|
||||||
|
|
||||||
|
|
||||||
class PipelineTask:
|
class PipelineTask(BaseTask):
|
||||||
def __init__(
|
def __init__(
|
||||||
self,
|
self,
|
||||||
pipeline: BasePipeline,
|
pipeline: BasePipeline,
|
||||||
@@ -122,7 +123,7 @@ class PipelineTask:
|
|||||||
|
|
||||||
self._observer = TaskObserver(params.observers)
|
self._observer = TaskObserver(params.observers)
|
||||||
|
|
||||||
def has_finished(self):
|
def has_finished(self) -> bool:
|
||||||
"""Indicates whether the tasks has finished. That is, all processors
|
"""Indicates whether the tasks has finished. That is, all processors
|
||||||
have stopped.
|
have stopped.
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user