Tasks

Flaxon provides a task queue system for background job processing with retries, scheduling, and result storage.

Defining Tasks

from flaxon.tasks import task

@task(name="send_email")
async def send_email(to: str, subject: str, body: str):
    # Send email
    return {"sent": True, "to": to}

@task(name="process_image")
def process_image(image_path: str):
    # CPU-intensive processing
    return {"processed": True}

Running Tasks

from flaxon.tasks import TaskQueue, TaskRegistry

registry = TaskRegistry()
queue = TaskQueue()

registry.register("send_email", send_email)
registry.register("process_image", process_image)

task = Task("send_email", send_email, args=["user@example.com", "Hello"])
await queue.push(task)

Starting Workers

# Start a worker
flaxon worker app:app --concurrency 4

# With specific queue
flaxon worker app:app --queue email --concurrency 2

Scheduling Tasks

from flaxon.tasks import Scheduler

scheduler = Scheduler(queue)

# Schedule a one-time task
scheduler.schedule(
    task=Task("send_email", send_email, args=["user@example.com", "Hello"]),
    delay=60,  # 60 seconds from now
)

# Schedule recurring task
scheduler.schedule(
    task=Task("process_image", process_image),
    interval=300,  # Every 5 minutes
)

# Scheduled decorator
from flaxon.tasks import scheduled_task

@scheduled_task(interval=60)
async def cleanup_tokens():
    await db.execute("DELETE FROM expired_tokens")

Task Results

task = Task("send_email", send_email, args=["user@example.com", "Hello"])
await queue.push(task)

result = await queue.get_result(task.id)
print(result.result)  # {"sent": True, "to": "user@example.com"}

Retry Policies

from flaxon.tasks import RetryPolicy

retry_policy = RetryPolicy(
    max_retries=5,
    delay=1.0,
    backoff=2.0,
    max_delay=60.0,
    random_jitter=0.1,
)

@task(name="retry_task", retry_policy=retry_policy)
async def retry_task():
    # Will retry up to 5 times with exponential backoff
    return {"success": True}

Task Timeouts

@task(name="timeout_task", timeout=30)
async def long_running_task():
    # Will be cancelled after 30 seconds
    return {"done": True}
Tip

Tasks support priority queues. Higher priority tasks are processed first.