Rustango docs
← Guides

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.

Background jobs in Rustango: a handler dispatches a Job onto a queue, worker tasks run it, retryable failures back off and retry, fatal ones go to a dead-letter handler

New to a term here? queue, worker, retry/backoff, dead-letter — see the glossary.

Source: rustango::jobs (Job, JobQueue, InMemoryJobQueue, JobError, JobDeadLetter) and rustango::jobs::DatabaseJobQueue (the table-backed queue, alias of pg::PgJobQueue) — behind the jobs feature (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 by jobs_sqlite_live.rs (cargo test -p rustango --features sqlite,jobs-postgres --test jobs_sqlite_live).

Table of contents


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 after run(). Until #1409 this example called queue.shutdown() on the line after run(), 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 — but on_shutdown is 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 JobQueue trait is not object-safe (its register/dispatch are generic), so you can't hold Arc<dyn JobQueue> in state — keep the concrete Arc<InMemoryJobQueue> (or Arc<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:

ModeWhere rustango_jobs livesWhy a worker stays scoped
databaseinside the tenant's own database / SQLite filea different tenant's rows are not in the table the worker queries — cross-tenant reads aren't expressible
schema (PG)the tenant's schemascoped_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(&registry_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

BackendUse forSurvives restart?
InMemoryJobQueuesingle process, dev, testsno
DatabaseJobQueue (PgJobQueue)production; multi-worker / multi-replicayes (rustango_jobs table)

JobError

VariantEffect
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 + a Duration; 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.