utils/tasks: added new documentation
This commit is contained in:
@@ -11,6 +11,18 @@ from pipecat.frames.frames import Frame
|
|||||||
|
|
||||||
|
|
||||||
class BaseTask(ABC):
|
class BaseTask(ABC):
|
||||||
|
@property
|
||||||
|
@abstractmethod
|
||||||
|
def id(self) -> int:
|
||||||
|
"""Returns the unique indetifier for this task."""
|
||||||
|
pass
|
||||||
|
|
||||||
|
@property
|
||||||
|
@abstractmethod
|
||||||
|
def name(self) -> str:
|
||||||
|
"""Returns the name of this task."""
|
||||||
|
pass
|
||||||
|
|
||||||
@abstractmethod
|
@abstractmethod
|
||||||
def has_finished(self) -> bool:
|
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
|
||||||
|
|||||||
@@ -125,10 +125,12 @@ class PipelineTask(BaseTask):
|
|||||||
|
|
||||||
@property
|
@property
|
||||||
def id(self) -> int:
|
def id(self) -> int:
|
||||||
|
"""Returns the unique indetifier for this task."""
|
||||||
return self._id
|
return self._id
|
||||||
|
|
||||||
@property
|
@property
|
||||||
def name(self) -> str:
|
def name(self) -> str:
|
||||||
|
"""Returns the name of this task."""
|
||||||
return self._name
|
return self._name
|
||||||
|
|
||||||
def has_finished(self) -> bool:
|
def has_finished(self) -> bool:
|
||||||
|
|||||||
@@ -44,6 +44,20 @@ def obj_count(obj) -> int:
|
|||||||
|
|
||||||
|
|
||||||
def create_task(loop: asyncio.AbstractEventLoop, coroutine: Coroutine, name: str) -> asyncio.Task:
|
def create_task(loop: asyncio.AbstractEventLoop, coroutine: Coroutine, name: str) -> asyncio.Task:
|
||||||
|
"""
|
||||||
|
Creates and schedules a new asyncio Task that runs the given coroutine.
|
||||||
|
|
||||||
|
The task is added to a global set of created tasks.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
loop (asyncio.AbstractEventLoop): The event loop to use for creating the task.
|
||||||
|
coroutine (Coroutine): The coroutine to be executed within the task.
|
||||||
|
name (str): The name to assign to the task for identification.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
asyncio.Task: The created task object.
|
||||||
|
"""
|
||||||
|
|
||||||
async def run_coroutine():
|
async def run_coroutine():
|
||||||
try:
|
try:
|
||||||
await coroutine
|
await coroutine
|
||||||
@@ -62,6 +76,18 @@ def create_task(loop: asyncio.AbstractEventLoop, coroutine: Coroutine, name: str
|
|||||||
|
|
||||||
|
|
||||||
async def wait_for_task(task: asyncio.Task, timeout: Optional[float] = None):
|
async def wait_for_task(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 (asyncio.Task): The asyncio Task to wait for.
|
||||||
|
timeout (Optional[float], optional): The maximum number of seconds
|
||||||
|
to wait for the task to complete. If None, waits indefinitely.
|
||||||
|
Defaults to None.
|
||||||
|
"""
|
||||||
name = task.get_name()
|
name = task.get_name()
|
||||||
try:
|
try:
|
||||||
if timeout:
|
if timeout:
|
||||||
@@ -82,6 +108,17 @@ async def wait_for_task(task: asyncio.Task, timeout: Optional[float] = None):
|
|||||||
|
|
||||||
|
|
||||||
async def cancel_task(task: asyncio.Task, timeout: Optional[float] = None):
|
async def cancel_task(task: asyncio.Task, timeout: Optional[float] = None):
|
||||||
|
"""Cancels the given asyncio Task and awaits its completion with an
|
||||||
|
optional timeout.
|
||||||
|
|
||||||
|
This function removes the task from the set of registered tasks upon
|
||||||
|
completion or failure.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
task (asyncio.Task): The task to be cancelled.
|
||||||
|
timeout (Optional[float]): The optional timeout in seconds to wait for the task to cancel.
|
||||||
|
|
||||||
|
"""
|
||||||
name = task.get_name()
|
name = task.get_name()
|
||||||
task.cancel()
|
task.cancel()
|
||||||
try:
|
try:
|
||||||
@@ -104,4 +141,5 @@ async def cancel_task(task: asyncio.Task, timeout: Optional[float] = None):
|
|||||||
|
|
||||||
|
|
||||||
def current_tasks() -> Set[asyncio.Task]:
|
def current_tasks() -> Set[asyncio.Task]:
|
||||||
|
"""Returns the list of currently created/registered tasks."""
|
||||||
return _TASKS
|
return _TASKS
|
||||||
|
|||||||
Reference in New Issue
Block a user