Background jobs
Some work shouldn't happen during a request — sending a welcome email, resizing an upload, syncing a third-party API. Doing it inline makes the user wait and couples the response to a flaky external call. A background job moves that work onto a queue: the handler returns immediately, and a pool of workers runs the job moments later, with automatic retries and a dead-letter path for failures. This is Celery or Laravel queues, in Rust.
New to a term here? queue, worker, retry/backoff, dead-letter — see the glossary.
Source:
rustango::jobs(Job,JobQueue,InMemoryJobQueue,JobError,JobDeadLetter) andrustango::jobs::DatabaseJobQueue(the table-backed queue, alias ofpg::PgJobQueue) — behind thejobsfeature (on by default).Runnable version: the in-memory, retry, and dead-letter snippets are copied from
jobs_doc.rs(cargo test -p rustango --test jobs_doc); the persistent queue is dogfooded on SQLite byjobs_sqlite_live.rs(cargo test -p rustango --features sqlite,jobs-postgres --test jobs_sqlite_live).
Table of contents
- Step 1 — Define a job
- Step 2 — Start a queue
- Step 3 — Dispatch from a handler
- Step 4 — Wire it into your app — the full example
- Running the workers (CLI + production)
- Retries and backoff
- The dead-letter handler
- The persistent queue (production)
- Jobs under multi-tenancy
- Scheduled sweeps under multi-tenancy
- Reference
- See also
Step 1 — Define a job
A job is a serializable struct that implements Job. The payload is what gets
queued (serialized to JSON); run() is the work. NAME routes a queued payload
back to its handler, so it must be unique.
Scaffold a skeleton with the CLI — cargo run -- make:job WelcomeEmail — or
write it by hand:
use rustango::jobs::{Job, JobError};
use serde::{Deserialize, Serialize};
#[derive(Serialize, Deserialize)]
struct WelcomeEmail {
user_id: i64,
}
#[async_trait::async_trait] // add `async-trait` to your Cargo.toml
impl Job for WelcomeEmail {
const NAME: &'static str = "welcome_email";
async fn run(&self) -> Result<(), JobError> {
// ... look up the user and send the email
Ok(())
}
}
run() returns:
Ok(())— done.Err(JobError::Retryable(msg))— transient; the worker retries with backoff.Err(JobError::Fatal(msg))— permanent; skip retries, dead-letter it now.
Override const MAX_ATTEMPTS: u32 = 3; on the impl to change the ceiling on
total attempts — not retries. The default of 5 is one initial run plus four
retries; MAX_ATTEMPTS = 3 gives two retries.
Step 2 — Start a queue
For a single process (and for dev and tests), InMemoryJobQueue runs jobs on
tokio worker tasks — no database. Register every job type, then start() the
workers (jobs aren't picked up until you do):
use rustango::jobs::{InMemoryJobQueue, JobQueue};
use std::sync::Arc;
let queue = Arc::new(InMemoryJobQueue::with_workers(4)); // 4 worker tasks
queue.register::<WelcomeEmail>().await; // before start/dispatch
queue.start().await;
// ... app runs ...
queue.shutdown().await; // on shutdown: drain in-flight jobs, then stop
shutdown() gives running jobs a grace period (shutdown_grace, default 5s),
then aborts and re-queues them. Queued jobs and parked retries stay queued, and
start() again picks them up.
Jobs are at-least-once. No stop signal reaches a running job, and an aborted job runs again from the start (the DB queue releases its row even if the job had just finished). The abort spends an attempt. Make every job idempotent.
Keep the Arc<InMemoryJobQueue> in your app state so handlers can reach it.
In-memory means in-memory. Jobs queued or in-flight are lost on restart. For anything you can't afford to drop, use the persistent queue.
Step 3 — Dispatch from a handler
Once the queue is running, dispatch is one call — it enqueues and returns immediately, so the request doesn't wait for the work:
queue.dispatch(&WelcomeEmail { user_id: 42 }).await?;
A worker picks it up and runs run() moments later. Verified end to end: three
dispatched jobs all run on the workers.
for user_id in 1..=3 {
queue.dispatch(&WelcomeEmail { user_id }).await.unwrap();
}
// → all three WelcomeEmail::run() calls execute on the worker pool
Step 4 — Wire it into your app
Putting Steps 1–3 together: build the queue and register every job once at
boot, start the workers, hand the queue to your routes so handlers can reach
it, then drain on shutdown. The queue lives in your main.rs:
// src/main.rs
use std::sync::Arc;
use rustango::jobs::{InMemoryJobQueue, JobQueue};
use rustango::manage::Cli;
#[rustango::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
// 1. Build the queue + register every job type BEFORE starting workers.
let queue = Arc::new(InMemoryJobQueue::with_workers(4));
queue.register::<WelcomeEmail>().await;
queue
.on_dead_letter(|dl| async move {
tracing::error!(job = dl.name, attempts = dl.attempts, error = %dl.error,
"job dead-lettered");
})
.await;
// 2. Start the workers — tokio tasks that run inside THIS process.
queue.start().await;
// 3. Share the queue with handlers (axum extension/state).
let app = urls::api().layer(axum::Extension(queue.clone()));
// 4. Boot the server, and say what to do when it stops. The hook
// runs after the server drains, on SIGINT *and* SIGTERM.
let draining = Arc::clone(&queue);
Cli::new()
.api(app)
.with_health()
.on_shutdown(move || async move { draining.shutdown().await })
.run()
.await?;
Ok(())
}
Put the drain in
on_shutdown, not afterrun(). Until #1409 this example calledqueue.shutdown()on the line afterrun(), and that line could never execute: nothing handled SIGTERM, so an orchestrator's stop killed the process outright.run()now returns on signal, so code after it does run — buton_shutdownis still the right place, because it also runs on the tenancy path and orders correctly against the server's own drain.
A handler dispatches by reading the queue back out of the request:
use axum::Extension;
use std::sync::Arc;
use rustango::jobs::{InMemoryJobQueue, JobQueue};
async fn signup(Extension(queue): Extension<Arc<InMemoryJobQueue>>) -> StatusCode {
// ... create the user ...
queue.dispatch(&WelcomeEmail { user_id: 42 }).await.ok(); // returns instantly
StatusCode::CREATED
}
Store the concrete queue type. The
JobQueuetrait is not object-safe (itsregister/dispatchare generic), so you can't holdArc<dyn JobQueue>in state — keep the concreteArc<InMemoryJobQueue>(orArc<DatabaseJobQueue>).
Running the workers (CLI + production)
start() spawns the workers as tokio tasks inside the current process. So
when you launch your app with the CLI — cargo run (which runs the server) —
the workers run right alongside it. For most apps that's all you need: one
process serves requests and drains the queue; there's no separate worker
command.
cargo run # starts the server AND the in-process workers
cargo run -- make:job WelcomeEmail # scaffold a new job type
A dedicated worker process (production)
At scale you often want workers separate from the web tier — so a traffic
spike can't starve jobs, and you scale each independently. With the
persistent queue every process pulls from the
same rustango_jobs table, so just run a second, server-less binary that builds
the queue, starts it, and blocks until a signal:
// src/bin/worker.rs — run with `cargo run --bin worker`
use std::sync::Arc;
use rustango::jobs::{DatabaseJobQueue, JobQueue};
#[rustango::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let pool = /* connect your tri-dialect Pool */;
DatabaseJobQueue::ensure_table_pool(&pool).await?;
let queue = Arc::new(DatabaseJobQueue::with_workers_pool(pool, 8));
queue.register::<WelcomeEmail>().await; // register the SAME job types
queue.start().await;
// Blocks until Ctrl-C (SIGINT) *or* SIGTERM. `tokio::signal::ctrl_c()`
// alone is SIGINT-only on Unix, so a worker waiting on it is killed
// outright by `docker stop` / a pod eviction and drains nothing.
rustango::shutdown::shutdown_signal().await;
queue.shutdown().await; // drain in-flight, then exit
Ok(())
}
Deploy it as its own container/service and scale to N replicas — they all
pull from the shared table safely. The web process then only needs to
dispatch (it doesn't have to start() workers). Pair the worker with a
periodic reclaim_stuck_jobs_pool sweep to
recover jobs from a crashed worker.
Retries and backoff
A job that returns Retryable is re-queued with exponential backoff (1s, 2s,
4s, 8s, …, capped at 1024s) until MAX_ATTEMPTS total attempts are spent. Use it
for transient failures — a timeout, a rate-limited API, a deadlock:
#[async_trait::async_trait]
impl Job for FlakyImport {
const NAME: &'static str = "flaky_import";
async fn run(&self) -> Result<(), JobError> {
match call_external_api().await {
Ok(_) => Ok(()),
Err(e) => Err(JobError::Retryable(e.to_string())), // try again later
}
}
}
The backing test dispatches a job that fails once then succeeds, and asserts it ran more than once and eventually succeeded — the retry actually happens.
The dead-letter handler
When a job exhausts its retries — or returns Fatal immediately — it's handed
to the dead-letter callback instead of vanishing. Register one (before
start) to log, alert, or persist the failure:
queue
.on_dead_letter(|dl| async move {
// dl: JobDeadLetter { name, payload, attempts, error }
tracing::error!(job = dl.name, attempts = dl.attempts, error = %dl.error,
"job dead-lettered");
})
.await;
A Fatal error skips retries and lands here on the first attempt:
async fn run(&self) -> Result<(), JobError> {
Err(JobError::Fatal("unprocessable payload".into())) // → dead-letter now
}
The test confirms the callback fires exactly once for a Fatal job, with the
job's name and error intact.
The persistent queue (production)
InMemoryJobQueue is single-process and forgets on restart. For real
deployments — multiple workers, multiple replicas, jobs that must survive a
crash — use DatabaseJobQueue (the same type as PgJobQueue). Jobs live in
a rustango_jobs table; workers pick them up with a transaction-bounded
UPDATE … RETURNING, so it's safe across processes. It's tri-dialect
(PostgreSQL, MySQL, SQLite). The Job definitions from Step 1 are unchanged.
use rustango::jobs::DatabaseJobQueue; // = rustango::jobs::pg::PgJobQueue
use rustango::jobs::JobQueue;
use std::sync::Arc;
use std::time::Duration;
// 1. Create the jobs table once at boot (idempotent, tri-dialect).
DatabaseJobQueue::ensure_table_pool(&pool).await?;
// 2. Build the queue over your pool + N workers, and start it.
let queue = Arc::new(
DatabaseJobQueue::with_workers_pool(pool.clone(), 4)
.poll_interval(Duration::from_millis(50)),
);
queue.register::<WelcomeEmail>().await;
queue.start().await;
// 3. Dispatch exactly as before — now it's durable.
queue.dispatch(&WelcomeEmail { user_id: 42 }).await?;
Recovering crashed workers. If a worker dies mid-job, the row stays locked. Run a periodic sweep to release locks older than a threshold so another worker picks the job back up:
// e.g. from a scheduled task every few minutes:
let reclaimed = DatabaseJobQueue::reclaim_stuck_jobs_pool(&pool, Duration::from_secs(300)).await?;
A running job refreshes its lock every heartbeat_interval (10 s by default), so
keep the threshold well above it. Each pickup spends an attempt, so a job that
crashes its worker is dead-lettered once MAX_ATTEMPTS runs are used. A job that
panics is a Retryable failure; the worker keeps running.
Both the dispatch-and-run flow and reclaim_stuck_jobs_pool are dogfooded
against SQLite in jobs_sqlite_live.rs.
Jobs under multi-tenancy
A worker is not a request. Nothing resolved a tenant for it — no host, no
header, no middleware ran — so a job carries no tenant context:
run(&self) receives only the deserialized payload. Isolation comes from
the pool the queue was built on, and the rule follows from that:
One queue per tenant pool. A tenant id in the payload is routing, not isolation.
That single decision is what makes each mode safe:
| Mode | Where rustango_jobs lives | Why a worker stays scoped |
|---|---|---|
database | inside the tenant's own database / SQLite file | a different tenant's rows are not in the table the worker queries — cross-tenant reads aren't expressible |
schema (PG) | the tenant's schema | scoped_pool_dyn returns a pool with search_path baked into its connect options, not a per-request SET — so every checkout for the pool's whole lifetime is scoped, including a worker that lives for days |
Build the queues at boot, one per active tenant:
use rustango::core::Column as _;
use rustango::jobs::{DatabaseJobQueue, JobQueue};
use rustango::sql::FetcherPool as _;
use rustango::tenancy::{Org, TenantPools};
use std::sync::Arc;
use std::time::Duration;
let pools = Arc::new(TenantPools::new(registry));
let registry_pool = pools.registry_pool();
// The registry holds the tenant list; each `Org` names its own
// database (or PG schema).
let orgs: Vec<Org> = Org::objects()
.where_(Org::active.eq(true))
.fetch(®istry_pool)
.await?;
let mut queues = Vec::new();
for org in orgs {
let pool = pools.scoped_pool_dyn(&org).await?; // scoped for its lifetime
DatabaseJobQueue::ensure_table_pool(&pool).await?;
let queue = Arc::new(
DatabaseJobQueue::with_workers_pool(pool, 2)
.poll_interval(Duration::from_secs(1)),
);
queue.register::<WelcomeEmail>().await;
queue.start().await;
queues.push((org.slug, queue));
}
Dispatch then goes to that tenant's queue — in a request handler, the one
keyed by the Org the resolver already produced.
What not to do. One queue on the registry pool with an org_id field in
the payload puts every tenant's jobs in one shared table, and leaves each
run() to remember to re-scope. One forgotten re-scope is a cross-tenant
write. If you need a shared queue for operational reasons, resolve the tenant
pool as the first thing run() does and never touch the registry pool
afterwards.
Known limits. The framework does not do the fan-out for you: there is no
per-tenant worker supervisor, no Job::run(&ctx) with the tenant attached,
and no tenant-aware reclaim_stuck_jobs_pool sweep — you loop over tenants
yourself, as above. New tenants provisioned after boot get no workers until
the process restarts. Tracked in
#1223.
Ambient context: audit source and timezone only. InMemoryJobQueue
captures them at dispatch and the scheduler at every(), and each reinstalls
them around the run. PgJobQueue does not yet: its jobs run as
AuditSource::System with the default timezone, so carry the actor in the
payload and re-enter the scope inside run() with audit::with_source. No
queue carries a session or a tenant. Tracked in
#1229.
Scheduled sweeps under multi-tenancy
The same "a worker is not a request" problem hits time-based work, and the
framework's own sweep helpers are where it bites: MediaManager::purge_orphans,
audit::cleanup_older_than_pool and prunable::prune_all each take one
pool, and every table they touch is per-tenant. One pool means one tenant —
and a registry pool in schema mode means only public, while every tenant's
rows accumulate and the sweep still reports success.
Fan out with for_each_tenant, which resolves each active tenant's own
pool and keeps going when one tenant fails:
use rustango::tenancy::for_each_tenant;
let opts = &opts;
let sweep = for_each_tenant(&pools, move |_org, pool| async move {
rustango::prunable::prune_all(&pool, opts).await
})
.await?;
tracing::info!(ok = sweep.succeeded(), failed = sweep.failed(), "nightly prune");
for (slug, err) in sweep.errors() {
tracing::warn!(%slug, %err, "tenant prune failed");
}
A broken tenant — rotated credential, unreachable database — is recorded in the
report rather than aborting the run, so it cannot starve the tenants after it.
Inactive orgs are skipped. purge_orphans is the one sweep that reaches outside
the database (it deletes storage objects), so it has a
purge_orphans_dry_run that runs the same query and deletes nothing — worth a
look before you wire the real thing.
If you guard the sweep so only one replica runs it, scope the lock per tenant:
use rustango::distributed_lock::DistributedLock;
let lock = DistributedLock::new(cache.clone()).for_tenant(&org.slug);
lock.with_lock("nightly_prune", ttl, || async { /* … */ }).await;
Unscoped, every tenant contends for one lock:nightly_prune: the first tenant
wins and the rest are skipped for the whole TTL, silently — a refused acquire is
the expected outcome, so nothing is logged. Tracked in
#1226 and
#1228.
Reference
Backends
| Backend | Use for | Survives restart? |
|---|---|---|
InMemoryJobQueue | single process, dev, tests | no |
DatabaseJobQueue (PgJobQueue) | production; multi-worker / multi-replica | yes (rustango_jobs table) |
JobError
| Variant | Effect |
|---|---|
Retryable(String) | re-queue with exponential backoff, up to MAX_ATTEMPTS |
Fatal(String) | skip retries → dead-letter immediately |
Queue(String) | internal queue error (serialization/registration) |
JobQueue methods: register::<T>() · dispatch(&payload) · start() ·
shutdown() · pending_count(). Both queues have shutdown_grace. The DatabaseJobQueue adds
ensure_table_pool, with_workers_pool, poll_interval, and
reclaim_stuck_jobs_pool.
See also
- Scheduler — for time-based recurring work at a fixed interval
(
Scheduler::every+ aDuration; there are no cron expressions), as opposed to on-demand jobs. - Email — the canonical "do it in a job" workload.
- Caching — the other way to keep request handlers fast.
- Signals — fire-and-forget hooks that often dispatch a job.
- Tenancy commands — provisioning the tenants a per-tenant queue fans out over.
