| name | litestar-saq |
| description | Auto-activate for litestar_saq, SAQPlugin, SAQConfig, QueueConfig, TaskQueues, CronJob, litestar workers run, background jobs, schedules, or SAQ web UI. Not for Celery/RQ/Dramatiq. |
litestar-saq
litestar-saq is the first-party plugin that integrates SAQ (Simple Async Queue) with Litestar. It provides:
SAQPlugin โ registers queues, workers, lifespan management, and DI for TaskQueues
SAQConfig / QueueConfig โ declarative plugin, queue, worker, broker, shutdown, polling, and OpenTelemetry configuration
litestar workers run โ CLI to start worker processes, optionally filtered by queue
- Optional web UI mounted under the Litestar app
- DI injection of
TaskQueues into route handlers for ergonomic enqueueing
Code Style Rules
- Use PEP 604 unions:
T | None, never Optional[T]
- Consumer Litestar app modules MAY use
from __future__ import annotations โ canonical Litestar apps do.
- Async all I/O โ task bodies and enqueue calls are
async def.
- First positional arg of every task is
ctx: dict (the SAQ context dict).
- Task params after
* are keyword-only.
- Use
NamedDependency[TaskQueues] for handler injection. TaskQueues is registered under the task_queues dependency key, and Litestar 2.24 deprecates implicit DI.
Quick Reference
Plugin Setup (canonical pattern)
The canonical pattern from litestar-fullstack (src/py/app/server/plugins.py) uses lazy initialization and use_server_lifespan=True so worker child processes start and stop with the Litestar server lifespan:
from litestar_saq import SAQConfig, SAQPlugin, QueueConfig, CronJob
from app.lib.settings import get_settings
def create_saq_plugin() -> SAQPlugin:
settings = get_settings()
return SAQPlugin(
config=SAQConfig(
use_server_lifespan=True,
web_enabled=settings.saq.web_enabled,
enable_otel=None,
queue_configs=[
QueueConfig(
name="default",
dsn=settings.redis.url,
tasks=["app.domain.system.tasks.send_email"],
scheduled_tasks=[
CronJob(
function="app.domain.system.tasks.cleanup_sessions",
cron="*/15 * * * *",
timeout=120,
),
],
),
],
),
)
saq_plugin = create_saq_plugin()
def create_saq_plugin_pg() -> SAQPlugin:
settings = get_settings()
return SAQPlugin(
config=SAQConfig(
use_server_lifespan=True,
web_enabled=settings.saq.web_enabled,
queue_configs=[
QueueConfig(
name="default",
dsn=settings.database.url,
tasks=["app.domain.system.tasks.send_email"],
),
],
),
)
Pick Branch A (SAQ + Redis) when: Redis is already in-stack (cache, sessions, Channels), you need multi-queue fanout across many workers, want the SAQ web UI, or prioritize queue throughput.
Pick Branch B (SAQ + PostgreSQL) when: the deployment is PostgreSQL-only, avoiding Redis matters, or SQL inspection of retained SAQ jobs is useful.
Pick Branch C (sidecar worker) when: you need a project-owned transactional outbox with business data, a project-owned job schema, frontend progress updates through channels, or multi-target execution routing (local / cloudrun / immediate) โ and you're willing to own the TaskService + Worker + WorkerSidecar + WorkerPlugin stack. See below.
Anti-pattern: hard-coding dsn=settings.redis.url in a PG-only project just because Redis is the "default" example. Match the broker to the stack.
Each QueueConfig accepts exactly one connection source: a redis://... or postgresql://... dsn, or a supported broker_instance. Supplying both or neither raises ImproperlyConfiguredException. Redis can use litestar-saq[hiredis]; PostgreSQL requires litestar-saq[psycopg]. Configure queue behavior with broker_options and connection/client construction with broker_instance_options.
Branch C โ Sidecar Worker Pattern
Some projects need a project-owned PostgreSQL worker stack. This wins when you want to write business data and an outbox job in one project-owned transaction, own the job table, publish frontend progress through channels, or route execution across local, cloudrun, and immediate targets.
In this pattern, TaskService owns the SQL state transitions, Worker owns claim/execution/retry flow, WorkerSidecar owns LISTEN/NOTIFY plus batched heartbeat/progress/channel publishing on a dedicated asyncpg connection, and WorkerPlugin only wires task discovery, schedule sync, DI, and optional in-process startup into Litestar.
See references/postgres-native-sidecar-worker.md for the full pattern.
from app.lib.worker.jobs import task
@task(cron="0 2 * * *", timeout=120)
async def nightly_cleanup() -> None:
"""Purge soft-deleted records every night at 02:00 UTC."""
...
@task(priority=5, retries=1, timeout=300, execution_target="cloudrun")
async def generate_report(*, report_id: int) -> None:
"""Export report โ runs on Cloud Run for isolation."""
...
await generate_report.enqueue(execution_target="cloudrun", report_id=42)
Wire into Litestar
from litestar import Litestar
from app.server.plugins import saq_plugin
app = Litestar(
route_handlers=[...],
plugins=[saq_plugin],
)
Define a Task
async def send_email(ctx: dict, *, recipient: str, subject: str, body: str) -> None:
"""Send an email as a background job.
Args:
ctx: SAQ context dict populated by worker hooks.
recipient: To address.
subject: Email subject.
body: Email body.
"""
email_service = ctx["email_service"]
await email_service.send(recipient, subject, body)
For long-running work, set the job's heartbeat stale threshold and decorate the task with monitored_job() so the plugin signals its batched HeartbeatManager while the task runs. heartbeat is not an update interval: SAQ marks an active job stuck when its last touch is older than that threshold.
from litestar_saq import monitored_job
@monitored_job()
async def rebuild_index(ctx: dict, *, index_name: str) -> dict[str, str]:
await run_rebuild(index_name)
return {"status": "complete"}
Enqueue from a Handler (DI of TaskQueues)
from litestar import Controller, post
from litestar.di import NamedDependency
from litestar_saq import TaskQueues
class NotificationController(Controller):
path = "/api/notifications"
@post("/")
async def queue_notification(
self,
data: NotificationCreate,
task_queues: NamedDependency[TaskQueues],
) -> dict[str, str]:
queue = task_queues.get("default")
job = await queue.enqueue(
"send_email",
recipient=data.email,
subject=data.subject,
body=data.body,
timeout=30,
retries=2,
key=f"notify-{data.email}",
)
return {"status": "queued" if job is not None else "duplicate"}
CLI
litestar --app app:app workers run
litestar --app app:app workers run --workers 4
litestar --app app:app workers run --queues emails --queues reports
litestar --app app:app workers status
Web UI
When web_enabled=True, the SAQ web UI is mounted under the Litestar app for queue introspection and job retry.
Job Options
| Option | Default | Use |
|---|
timeout | 10 | Always set explicitly โ SAQ's default is usually too low or too high for real jobs |
retries | 1 | Retry count on exception |
retry_delay | 0.0 | Seconds to wait before retrying |
retry_backoff | False | True for exponential backoff with jitter, or a numeric maximum delay |
ttl | 600 | Seconds to retain result after completion |
key | generated | Stable uniqueness key; enqueue returns None while that key already exists |
heartbeat | 0 | Maximum seconds an active job may go without a touch; 0 disables stale detection |
scheduled | 0 | Unix timestamp to delay start |
Workflow
Step 1: Install
pip install litestar-saq
pip install "litestar-saq[psycopg]"
pip install "litestar-saq[otel]"
Step 2: Define Queues
Build QueueConfig instances for each logical queue ("default", "emails", "reports"). Put exactly one of dsn or broker_instance on each QueueConfig. Reference task functions by dotted path or callable; the plugin imports dotted paths at startup.
Step 3: Configure Plugin
Wrap QueueConfigs in SAQConfig. Pick Redis (redis://...) when Redis is already in the stack; pick PostgreSQL (postgresql://...) when the project is PG-only. Set use_server_lifespan=True when the web process should own worker child processes. Toggle web_enabled / web_guards for the introspection UI and enable_otel for tracing.
Step 4: Define Tasks
Place task functions in app/domain/<domain>/tasks.py. First arg ctx: dict, rest keyword-only. Add shared resources (DB, HTTP client, email service) in QueueConfig.startup / before_process hooks and read them from ctx.
Step 5: Schedule Cron Work
Add CronJob entries to QueueConfig.scheduled_tasks for recurring work. Always set timeout. Do not use external cron tools for work that belongs in the queue.
Step 6: Enqueue from Handlers
Inject TaskQueues into route handlers. Use task_queues.get("name") then await queue.enqueue("task_name", ...). The result is the new Job, or None when its unique key already exists; handle both outcomes. Use key= for deduplication.
Step 7: Publish to Channels (optional)
For real-time updates after a job completes, publish to Litestar Channels from inside the task. See ../litestar-realtime/references/websockets.md.
Step 8: Run
For dev: litestar run can start worker child processes when use_server_lifespan=True.
For production: run litestar workers run --workers N as a separate service/process from litestar run. There is no --process flag in litestar-saq 0.8.0.
For portable multi-process workers, configure each queue with dsn. Under spawn and forkserver, litestar-saq removes cached live broker objects from the child configuration and rebuilds them from that DSN.
Guardrails
- Use
litestar-saq, not raw SAQ, in Litestar apps โ the plugin handles DI, lifespan, CLI, and the web UI. Raw SAQ misses all of that.
- Always set
timeout on tasks and CronJobs โ SAQ defaults to 10s, which is rarely the correct production value.
- Pair
heartbeat with monitored_job() for long-running jobs โ heartbeat defines when a job is stale; it does not emit updates. Keep the stale threshold longer than the decorator signal interval and the heartbeat manager's flush cadence.
- Inject
TaskQueues via DI โ don't import a global queue inside handlers. The plugin owns the queue lifecycle.
- Use
CronJob for scheduled work โ not external cron. CronJobs participate in retries, timeouts, and observability.
- Handle
queue.enqueue() returning None โ an existing unique key prevents insertion, so do not report every enqueue attempt as newly queued.
- Use
key= for deduplication โ same logical job (per-user sync, per-resource refresh) should not stack. The key remains occupied until the stored job expires or is removed, including after terminal completion.
use_server_lifespan=True for dev and small-to-mid apps that should start worker child processes with the web server. For high-throughput production, run litestar workers run --workers N as a separate service.
- Use
dsn for portable multi-process workers โ forkserver/spawn workers rebuild brokers from QueueConfig.dsn. A broker_instance-only queue works in the parent and under fork, but spawn preparation rejects it.
- Set graceful shutdown controls for long jobs โ use
shutdown_grace_period_s and, when needed, cancellation_hard_deadline_s on QueueConfig.
- Publish to Litestar Channels from tasks when the job result must update connected websocket clients. See
../litestar-realtime/references/websockets.md.
- Pull shared resources from
ctx populated by QueueConfig hooks, not module-level globals โ keeps tests deterministic and supports per-worker init.
- Reach for the sidecar worker pattern when you need a project-owned transactional outbox, job schema, batched heartbeats for many running jobs, frontend updates through channels, or execution-target routing across Cloud Run / local. Keep the stack explicit:
TaskService for fenced SQL transitions, Worker for execution, WorkerSidecar for wakeups/batched heartbeats/channel publishing, and WorkerPlugin for Litestar lifecycle wiring. For normal PG-backed queueing, SAQ+PG is the simpler default. See references/postgres-native-sidecar-worker.md.
Validation Checkpoint
Before delivering Litestar + SAQ code, verify:
Example
Task: A Litestar app with a default queue, an email task, a cleanup CronJob, and a handler that enqueues notifications. This example uses Redis as the SAQ broker; swap dsn=settings.redis.url for dsn=settings.database.url if the project is PG-only โ see Quick Reference above for both patterns.
from litestar_saq import SAQConfig, SAQPlugin, QueueConfig, CronJob
from app.lib.settings import get_settings
def create_saq_plugin() -> SAQPlugin:
settings = get_settings()
return SAQPlugin(
config=SAQConfig(
use_server_lifespan=True,
web_enabled=settings.saq.web_enabled,
queue_configs=[
QueueConfig(
name="default",
dsn=settings.redis.url,
startup="app.domain.system.tasks.worker_startup",
shutdown="app.domain.system.tasks.worker_shutdown",
tasks=[
"app.domain.system.tasks.send_email",
"app.domain.system.tasks.cleanup_sessions",
],
scheduled_tasks=[
CronJob(
function="app.domain.system.tasks.cleanup_sessions",
cron="*/15 * * * *",
timeout=120,
),
],
),
],
),
)
saq_plugin = create_saq_plugin()
async def worker_startup(ctx: dict) -> None:
"""Initialize shared resources for this worker."""
ctx["email_service"] = create_email_service()
ctx["db"] = create_database_client()
async def worker_shutdown(ctx: dict) -> None:
"""Dispose shared worker resources."""
await ctx["db"].close()
async def send_email(ctx: dict, *, recipient: str, subject: str, body: str) -> None:
"""Send an email as a background job."""
email = ctx["email_service"]
await email.send(recipient, subject, body)
async def cleanup_sessions(ctx: dict) -> None:
"""Purge expired sessions every 15 minutes."""
db = ctx["db"]
await db.execute("DELETE FROM session WHERE expires_at < now()")
from litestar import Controller, post
from litestar.di import NamedDependency
from litestar_saq import TaskQueues
from app.domain.notifications.schemas import NotificationCreate
class NotificationController(Controller):
path = "/api/notifications"
tags = ["Notifications"]
@post("/")
async def queue_notification(
self,
data: NotificationCreate,
task_queues: NamedDependency[TaskQueues],
) -> dict[str, str]:
queue = task_queues.get("default")
job = await queue.enqueue(
"send_email",
recipient=data.email,
subject=data.subject,
body=data.body,
timeout=30,
retries=2,
key=f"notify-{data.email}",
)
return {"status": "queued" if job is not None else "duplicate"}
from litestar import Litestar
from app.domain.notifications.controllers import NotificationController
from app.server.plugins import saq_plugin
app = Litestar(
route_handlers=[NotificationController],
plugins=[saq_plugin],
)
litestar --app app:app run
litestar --app app:app workers run --workers 4
References Index
- Advanced Patterns โ Heartbeat tuning, dead-letter handling, job chaining, queue priorities, worker lifecycle hooks, Postgres backend.
- Sidecar Worker Pattern โ TaskService + Worker + WorkerSidecar + WorkerPlugin pattern for a project-owned transactional outbox, job schema, sidecar batched heartbeats/wakeups, channel publish-back to the frontend,
@task decorator + ScheduleConfig cron registry, and execution_target routing (local / cloudrun / immediate).
Cross-References
Official References
Shared Styleguide Baseline
- Use shared styleguides for generic language/framework rules to reduce duplication in this skill.
- General Principles
- Python
- Litestar
- Keep this skill focused on tool-specific workflows, edge cases, and integration details.