Files
aiomql/docs/core/task_queue.md
T
Ichinga Samuel 2a55f49f78 v4
2024-11-13 04:29:37 +01:00

11 KiB

TaskQueue and QueueItem

Table of Contents

QueueItem

class QueueItem

A task to be executed by the TaskQueue. The task can be any coroutine callable. The task is wrapped as a QueueItem object, which is then added to the TaskQueue for execution. The arguments and keyword arguments are passed to the task when it is executed.

Attributes:

Name Type Description
task_item Callable | Coroutine A coroutine function to be executed by the TaskQueue
args tuple[Any, ...] Positional arguments to be passed to the task_item
kwargs dict[str, Any] Keyword arguments to be passed to the task_item
must_complete bool If True, the item must be completed even if the queue is stopped.
time float The time the item was added to the queue. For sorting priority queues

_init_

def __init__(self, task: Callable | Coroutine, *args, **kwargs):

Parameters:

Name Type Description
task Callable | Coroutine A coroutine to be executed by the TaskQueue
args Any Positional arguments to be passed to the task when it is executed
kwargs Any Keyword arguments to be passed to the task when it is executed

run

def run(self)

Run the task. If the task is a coroutine, it is awaited. If the task is a callable, it is called.

TaskQueue

class TaskQueue

A perpetual task queue that processes QueueItem objects. The TaskQueue runs indefinitely, processing QueueItem objects as they are added to the queue. The TaskQueue is a wrapper around an asyncio.Queue that can be passed in as an argument or defaults to an asyncio.PriorityQueue. It is added to the bot executor of the Bot class on a separate thread.

Attributes:

Name Type Description
queue asyncio.Queue An asyncio.Queue queue of QueueItem objects to be executed by the TaskQueue. If not provided during instantiation, an asyncio.PriorityQueue is used
stop bool A flag to stop the task_queue instance.
workers int The number of workers to process the queue items. Defaults to 10.
timeout int The maximum time to wait for the queue to complete. Default is None. If timeout is provided the queue is joined using asyncio.wait_for with the timeout
on_exit Literal["cancel", "complete_priority"] The action to take when the queue is stopped. If "cancel" the queue is cancelled and the remaining items are not processed. If "complete_priority" the queue is completed with the priority items. Default is "cancel"
mode Literal["finite", "infinite"] The mode of the queue. If finite the queue will stop when all tasks are completed. If infinite the queue will continue to run until stopped.
worker_timeout int The time to wait for a task to be added to the queue before stopping the worker or adding a dummy sleep task to the queue.
tasks List[Task] A list of the worker tasks running concurrently, including the main task that joins the queue.
priority_tasks set[QueueItem] A set to store the QueueItems that must complete before the queue stops

_init_

def __init__(self, queue: asyncio.Queue = None, workers: int = 10, timeout: int = None, size: int = None,
             on_exit: Literal["cancel", "complete_priority"] = "cancel",
             mode: Literal["finite", "infinite"] = "infinite", worker_timeout: int = 60)

Create a new TaskQueue instance.

Parameters:

Name Type Description Default
queue asyncio.Queue An asyncio.Queue queue instance None
workers int The number of workers to process the queue items. 10
timeout int The maximum time to wait for the queue to complete. None
size int The maximum size of the queue. None
on_exit Literal["cancel", "complete_priority"] The action to take when the queue is stopped. "complete_priority"
mode Literal["finite", "infinite"] The mode of the queue. "infinite"
worker_timeout int The time to wait for a task to be added to the queue before stopping the worker or adding a dummy sleep task to the queue. 60

add

def add(*, item: QueueItem, priority: int = 3, must_complete_false: bool = False) -> None

Add a QueueItem to the TaskQueue queue.

Parameters:

Name Type Description
item QueueItem A QueueItem to be added to the queue
priority int The priority of the item. The lower the number, the higher the priority. Default is 3.
must_complete_false bool If True, the item must be completed even if the queue is stopped. Default is False.

add_task

def add_task(self, task: Callable | Awaitable, *args, **kwargs)

Create a QueueItem from the task and add it to the TaskQueue queue. The task can be a callable or an awaitable. The arguments and keyword arguments are passed to the QueueItem.

Parameters:

Name Type Description
task Callable | Awaitable A callable or awaitable task to be executed by the TaskQueue
args Any Positional arguments to be passed to the task when it is executed
kwargs Any Keyword arguments to be passed to the task when it is executed

worker

async def worker()

A worker that processes the QueueItem objects in the TaskQueue queue.

run

async def run(timeout: int = None)

Start the TaskQueue instance. If a timeout is provided, the queue is joined using asyncio.wait_for with the timeout. This is the main entry point for the TaskQueue instance. It is added to the bot executor of the Bot class on a separate thread.

Parameters:

Name Type Description
timeout int The maximum time to wait for the queue to complete. Default is None.

stop_queue

def stop_queue()

Stop the TaskQueue instance. This sets the stop attribute to True, changes the on_exit attribute to "cancel", and cancels the queue.

clean_up

async def clean_up()

Clean up the TaskQueue instance. This is called when the queue is stopped. It cancels the queue and processes the remaining priority items based on the on_exit attribute.

cancel

def cancel()

Cancel all remaining tasks.