Skip to main content

Crate later

Crate later 

Source
Expand description

§later

A distributed background job manager and runner for Rust.

Later stores job state and delivers work through SQLite or PostgreSQL. Every server using the same database and namespace joins the same worker pool. Renewable leases coordinate job execution across workers and server processes.

§Documentation index

§Retained logs

Enable retained-log with sqlite or postgres to append immutable bytes and replay them by partition offset. Consumer groups commit offsets and use 30-second, epoch-fenced leases. This is a database-backed log API, not a Kafka protocol implementation: there is no Kafka wire protocol, broker replication, or cross-region quorum.

§Set up

§1. Import later and required dependencies

later = { version = "0.0.27", features = ["sqlite"] }
serde = { version = "1.0", features = ["derive"] }

§2. Define some types to use as a payload to the background jobs

use serde::{Deserialize, Serialize};

#[derive(Serialize, Deserialize)] // <- Required derives
pub struct SendEmail {
    pub address: String,
    pub body: String,
}

// ... more as required

§3. Generate the stub

later::background_job! {
    struct Jobs {
        send_email: SendEmail,
    }
}

This generates two types

  • JobsBuilder - used to bootstrap the background job server - which can be used to enqueue jobs,
  • JobContext<T> - used to pass application context (T) in the handler as well as enqueue jobs,

§4. Use the generated code to bootstrap the background job server

For struct Jobs a type JobsBuilder will be generated. Use this to bootstrap the server.

ⓘ
use later::{backend::SqliteBackend, storage::Sqlite, BackgroundJobServer, Config};

// bootstrap the server
let ctx = MyContext{ /*..*/ };                  // Any context to pass onto the handlers
let storage = Sqlite::new("sqlite://later.db").await?;
let backend = SqliteBackend::new("fnf-example", storage)?;
let ctx = JobsBuilder::new(
    later::Config::builder()
        .context(ctx)                       // Pass the context here
        .backend(Box::new(backend))          // Configure state and delivery
        // ...
        .build()
    )
    // for each payload defined in the `struct Jobs` above
    // the generated fn name uses the pattern "with_[name]_handler"
    .with_send_email_handler(handle_send_email)     // Pass the handler function
    // ..
    .build()
    .await?;

// use ctx.enqueue(SendEmail{ ... }) to enqueue jobs,
// or ctx.enqueue_continue(parent_job_id, SendEmail{ ... }) to chain jobs.
// this will only accept types defined inside the macro above
// define handler
async fn handle_send_email(
        ctx: JobsContext<MyContext>, // JobContext is generated wrapper
        payload: SendEmail,
    ) -> anyhow::Result<()> {
    // handle `payload`

    // ctx.app -> Access the MyContext passed during bootstrapping
    // ctx.enqueue(_).await to enqueue more jobs
    // ctx.enqueue_continue(_).await to chain jobs

    Ok(()) // or Err(_) to retry this message
}

This example uses SQLite for state and delivery.


§How workers run

A BackgroundJobServer starts six worker tasks by default. Set Config::worker_count to change that number. The count applies to each server instance, so three server processes configured with four workers can run up to twelve job handlers at once. Every worker owns one delivery consumer and processes one command at a time. A slow job occupies one worker while the other workers continue receiving commands. BackgroundJobServer::add_worker starts another local worker at runtime. BackgroundJobServer::remove_worker waits for one worker’s current job before stopping it. At least one worker remains. For process shutdown, call BackgroundJobServer::shutdown with a deadline. It stops new claims, lets active handlers commit, and then deregisters a sequential-processing owner. If the deadline expires, remaining durable work recovers through its normal lease.

Server startup waits until every worker consumer is ready. Workers using the same backend database and namespace compete for work, including workers in other processes or on other hosts. SQLite can coordinate several processes only when they can all open the same database file. PostgreSQL can coordinate processes across hosts.

For each job delivery, a worker:

  1. claims a renewable execution lease so another worker does not run the same job at the same time;
  2. loads the durable job and changes its stage from enqueued to running;
  3. calls the generated handler for the stored payload type;
  4. stores success, failure, or the next retry; and
  5. acknowledges the delivery. A delivery error is reported to the backend so it can be made available again.

A delivery backend may make a command available more than once. The execution lease prevents concurrent handler calls, and the stored job stage makes workers ignore duplicate commands after a job has moved on from the enqueued stage. A regular (non-partitioned) job whose handler is interrupted after the running stage is stored - its worker crashed or was killed mid-handler - is reclaimed by the periodic PollStuckJobs maintenance command once its execution lease has been expired for a while: it’s retried (or marked failed, once retries are exhausted) the same way a handler error would be. Partitioned/sequential jobs don’t need this - a stuck partition head already recovers through its own partition-lease expiry.

Delayed jobs and retries do not occupy a sleeping worker. They remain in durable storage until a periodic command finds that they are ready, then they rejoin the delivery queue. Polling and queue load can add delay, so a scheduled time or retry delay is a lower bound rather than an exact start time.

A separate periodic command, PollExpiredStorage, physically deletes already-expired rows in small batches. Expiry (Storage::expire, and the TTLs later sets internally on terminal jobs and dashboard bookkeeping) only ever makes a row invisible to reads - something still has to delete it, or storage only ever grows. This runs automatically; there is nothing to configure.

Keep the BackgroundJobServer handle alive for as long as this process should accept work. Dropping it aborts its local worker and maintenance tasks, including the two periodic commands described above.

With the dashboard feature, each stage change updates job state and its dashboard projection in the same atomic storage operation. Dashboard work therefore cannot form a second delivery backlog during producer load.

§Fire and forget jobs

Fire and forget jobs are made available to a worker immediately. A failed handler is retried according to its resolved retry policy.

ctx.enqueue(SendEmail{
    address: "hello@rust-lang.org".to_string(),
    body: "You rock!".to_string()
}).await?;

§Continuations

One or many jobs can be chained into a workflow. A child becomes available only after its parent succeeds.

let email_welcome = ctx.enqueue(SendEmail{
    address: "customer@example.com".to_string(),
    body: "Creating your account!".to_string()
}).await?;

let create_account = ctx.enqueue_continue(email_welcome, CreateAccount { id: "accout-1".to_string() }).await?;

let email_confirmation = ctx.enqueue_continue(create_account, SendEmail{
    address: "customer@example.com".to_string(),
    body: "Your account has been created!".to_string()
}).await;

§Delayed jobs

A delayed job becomes available at or after a chosen interval or timestamp.

// delay
ctx.enqueue_delayed(SendEmail{
    address: "hello@rust-lang.org".to_string(),
    body: "You rock!".to_string()
}, std::time::Duration::from_secs(60)).await?;

// specific time
let run_job_at : chrono::DateTime<chrono::Utc> = todo!();
ctx.enqueue_delayed_at(SendEmail{
    address: "hello@rust-lang.org".to_string(),
    body: "You rock!".to_string()
}, run_job_at).await?;

§Recurring jobs

Run a recurring job on a cron schedule. BackgroundJobServerPublisher::enqueue_recurring registers a payload and cron expression under an identifier; calling it again with the same identifier and cron updates the existing definition in place (upsert) rather than creating a second schedule, and returns None rather than enqueuing another occurrence - safe to call on every process startup to recreate the schedule, which is the intended pattern; calling it unconditionally on every startup was, once, not safe, and would enqueue one more permanently-queued occurrence per restart. Registering with a different cron is treated as a real change and does enqueue a fresh occurrence, returning Some. A background poller proactively enqueues each subsequent due occurrence, independent of whether a previous occurrence’s handler is still running or ever ran — a crashed process cannot silently stop the schedule.

ctx.enqueue_recurring("send-newsletter-1".to_string(),
    SendNewsletter{
        address: "hello@rust-lang.org".to_string(),
    },
    "0 6 1 * * *".to_string() // 6am, 1st day of every month
).await?;

By default occurrences run on the plain unordered path and may overlap, the same as any other job. Use BackgroundJobServerPublisher::enqueue_recurring_sequential instead when at most one instance of a recurring job — including any jobs chained off an occurrence with BackgroundJobServerPublisher::enqueue_recurring_continue — may ever be running at once. It reuses the sequential partition mode described below: every occurrence, and every job chained off one, is enqueued into the same topic partition (keyed by identifier), so they run in strict order and never overlap. Requires Config::recurring_sequential_partitions to be set.

// Chain a follow-up step off the currently-running occurrence itself,
// using its own job id. Because this must happen before the handler
// returns, the continuation is always ordered ahead of the next
// scheduled occurrence.
ctx.enqueue_recurring_continue(ctx.job_id().clone(), RunReport).await?;
ctx.enqueue_recurring_sequential("nightly-report".to_string(),
    RunReport,
    "0 0 3 * * *".to_string() // 3am daily
).await?;

§Sequential partition mode

Named topics with a fixed number of Kafka-like partitions give strict, in-order execution within one (topic, partition), while different partitions and topics still run concurrently. Jobs enqueued without a topic keep the default unordered, competing-worker behaviour described above; nothing changes for them.

A message type opts in by implementing topic::JobPartition and being declared with #[topic("name")] inside background_job!. The macro then requires that type to implement topic::JobPartition and generates the topic::JobTopic marker the topic-aware enqueue methods need. topic::JobPartition::partition_key returns an arbitrary topic::PartitionKey — a customer ID, an order ID, anything that identifies “things that must stay in order relative to each other” — and Later hashes it to a partition with topic::partition_for_key, the way a Kafka producer key is hashed by the partitioner. The caller never names a partition number directly. Register every topic’s fixed partition count on Config::topics before starting the server; an unregistered topic is rejected before anything is written.

use later::topic::{JobPartition, PartitionKey};

#[derive(serde::Serialize, serde::Deserialize)]
pub struct ProcessOrder {
    pub customer_id: u32,
    pub description: String,
}

// Every order from the same customer must apply in order; different
// customers can usually be processed at the same time.
impl JobPartition for ProcessOrder {
    fn partition_key(&self) -> PartitionKey {
        PartitionKey::from(self.customer_id)
    }
}

later::background_job! {
    struct Jobs {
        #[topic("orders")]
        process_order: ProcessOrder,
    }
}

Register the topic when building the server, and use the partition-aware enqueue methods for that message type:

use later::topic::TopicConfig;

let config = later::Config::builder()
    .context(context)
    .backend(backend)
    .topics(vec![TopicConfig::new("orders", 8)?])
    .build();
// Only a type declared with a topic can use the partition-aware methods.
ctx.enqueue_to_partition(ProcessOrder {
    customer_id: 42,
    description: "add item".to_string(),
}).await?;

BackgroundJobServerPublisher::enqueue_to_partition, BackgroundJobServerPublisher::enqueue_to_partition_delayed, and BackgroundJobServerPublisher::enqueue_continue_to_partition mirror BackgroundJobServerPublisher::enqueue, BackgroundJobServerPublisher::enqueue_delayed, and BackgroundJobServerPublisher::enqueue_continue, but assign the job a monotonic sequence within its (topic, partition) as part of the same atomic write. The SQLite and PostgreSQL backends also expose enqueue_to_partition_in for writing application data, the job, and its sequence in one caller-owned transaction, matching enqueue_in.

Ordering rules:

  • The oldest unfinished job in a partition blocks every later job in that partition, including a retry (which keeps its position), a delayed head (which blocks until its scheduled time), and a waiting continuation (which blocks until its parent succeeds, even when the parent is not itself partitioned or lives in a different partition).
  • A terminal success or failure releases the next sequence.
  • Different partitions, and different topics, never block each other, even when they share the same partition ID.

Every server with configured topics registers a heartbeat and runs a bounded number of partition poller tasks (one per worker) that lease and run the oldest ready, unowned head. Which server owns a given (topic, partition) is decided by rendezvous hashing over every currently live server (heartbeats expire after 15 seconds), so a small membership change only moves the partitions whose winner changes rather than reshuffling every assignment. A lease keeps one worker’s claim exclusive across every process sharing the same database and namespace; losing a race for a lease is normal and just means another worker already owns that partition’s head. The lease renews itself while a handler keeps running, so a slow handler does not let another worker reclaim and re-run the same head concurrently. As with the default delivery path, a worker crash can still run the current head job again after its lease expires: order is preserved, but execution is at-least-once rather than exactly-once.

Each partition claim also has a fencing epoch. If renewal shows that a lease was lost, Later cancels the handler future and rejects that claim’s retry or completion write. Only the owner holding the current epoch can advance the stored partition head. This protects Later’s own state during a process pause or a reclaimed lease. It cannot undo an external side effect that began before cancellation, so handlers that call external systems must still be idempotent, normally using the job ID as their key.

Adding a worker lets it start claiming newly idle or unassigned partitions within a couple of membership refresh cycles, without waiting for any lease to expire. A job already running when its partition’s assignment changes finishes under its current owner; the new owner only takes over once that worker stops claiming the partition’s next head. Dropping a BackgroundJobServer deregisters its worker immediately (best-effort) so other workers stop considering it right away; an ungraceful removal (a crash) is only detected once its heartbeat expires.

SQL is always the source of order and ownership. A poller also reacts immediately (instead of waiting out its idle backoff) to a local enqueue or continuation. This wake-up is a pure latency optimization; nothing depends on it arriving, since every poller still scans SQL on its own.

If a head job would otherwise block its partition indefinitely (a “poison” job whose retries never succeed, or one you simply need unblocked sooner), BackgroundJobServerPublisher::force_fail_partition_head moves it straight to a terminal failed state and frees the partition, without waiting out its retry policy and without needing to hold its lease first.

See <https://github.com/mustakimali/later/blob/main/plans/SEQUENTIAL_PARTITION_MODE.md> for what is planned next.

§Backends

StateDefault deliveryScaleFeatures
SQLiteSQLiteSeveral processes sharing one local filesqlite
PostgreSQLPostgreSQLProcesses on many hostspostgres
Give every backend a namespace. Servers that use the same database and
namespace share work; different namespaces stay isolated. The SQL backends
also accept an application-owned SQLx pool, which allows application data
and a job to be written in the same transaction with enqueue_in.

§Retry behavior

The default server policy retries six times with exponential backoff and bounded jitter. Override it for the whole server through Config:

use later::{retry::RetryPolicy, Config};
use std::time::Duration;

let config = Config::builder()
    .context(context)
    .backend(backend)
    .default_retry_policy(RetryPolicy::fixed(3, Duration::from_secs(5)))
    .build();

A message type can take priority over that setting:

use later::retry::{JobRetryPolicy, RetryPolicy};

struct AppContext;
struct SendEmail;

#[later::async_trait::async_trait]
impl JobRetryPolicy for SendEmail {
    type Context = AppContext;

    async fn retry_policy(&self, _context: &AppContext) -> anyhow::Result<RetryPolicy> {
        Ok(RetryPolicy::no_retries())
    }
}

§Dashboard

Enable feature dashboard to enable the experimental dashboard. The host application is responsible for protecting its route. Dashboard metadata is retained for one day and does not include job payload bytes. Mount the response returned by BackgroundJobServerPublisher::get_dashboard in the HTTP framework used by the application.

The bundled page also shows stage-transition throughput from the last 60 seconds and workers with a heartbeat in the last 15 seconds, plus average and maximum queue-dispatch wait time over that same window - how long a job sat enqueued/delayed/requeued before a worker picked it up, tracked separately for regular jobs and sequential/partitioned ones (a sequential job’s wait includes time blocked behind earlier jobs in its partition, so the two are not directly comparable). This is the overhead later itself adds before a handler starts, not time spent inside the handler. These stats are stored only when the dashboard feature is enabled - no prometheus feature needed. Each job response includes total elapsed, handler-processing, and waiting durations computed from its stage history. Job details link to a continuation’s parent and to at most 25 jobs that directly continue from it. The jobs table shows only the last 6 characters of each ID (the full ID is still in a hover tooltip and used for the detail link) - later’s IDs are lexicographically sortable ULIDs, so the trailing characters are what actually varies between IDs minted close together.

Every topic in Config::topics also gets a badge in the topic menu showing its queue depth (not-yet-completed jobs across every partition), refreshed on the same cycle as everything else on the page, plus a per-partition breakdown in the partition picker once a topic is selected - no separate opt-in needed. Set Config::retained_log_lag (needs the retained-log feature too) to add each configured retained-log topic’s consumer lag to the same badges - see RetainedLogLag.

The performance panel also shows the entire database’s on-disk size (not just Later’s own tables) when the storage backend can report one - SQLite and PostgreSQL both can; a backend without a single meaningful database size (Redis, in-memory) shows a dash instead. This is the whole database an application shares with Later, useful for keeping an eye on total footprint even when most of it is application data.

A partitioned topic’s dashboard history (the list backing its per-partition “partition jobs” view) is a full, append-only record by design and grows forever unless capped. Set Config::partition_history_retention to a duration (for example seven days) to have new entries self-expire instead - or call [BackgroundJobServerPublisher::truncate_partition_history] to trim an already-accumulated history down to its most recent N entries in one shot.

§Prometheus metrics

Enable prometheus to add process-local worker metrics. Expose the text returned by BackgroundJobServerPublisher::get_metrics on an HTTP route that your Prometheus server can scrape. The feature records:

  • later_job_transitions_total for durable stage changes;
  • later_job_handler_duration_seconds for handler latency and outcome per worker;
  • later_job_duration_seconds and later_job_wait_duration_seconds for completed job workflow latency;
  • later_queue_jobs for durable ready, leased, delayed, retry-waiting, and dependency-waiting job counts;
  • later_worker_commands_total for successful, failed, and discarded delivery commands per worker;
  • later_workers for each active and busy worker task; and
  • for servers with configured topics: later_partition_claims_total (claim attempts, whether or not a head was found), later_partition_heads_claimed_total and later_partitions_active by topic, later_partition_lease_renewals_total by topic and outcome, later_partition_workers_live for this process’s last-observed live worker count, later_partition_backlog / later_partition_oldest_head_age_seconds by topic for the ready backlog and how long its oldest head has been waiting, and later_partition_queue_depth by topic and partition for not-yet-completed jobs currently queued there; and
  • with retained-log also enabled and Config::retained_log_lag set: later_retained_log_consumer_lag by topic and consumer group, for records appended but not yet committed by that group.

Worker metrics use a numeric worker_id that is local to one server. Combine it with Prometheus’s scrape instance label to distinguish workers running in different processes. Metrics never use a job ID, error message, or partition worker owner ID as a label — only the configured topic name, which stays bounded the way queue and job_type already are. The one exception is later_retained_log_consumer_lag’s group label: it is whatever name the application commits offsets under, not a fixed set later controls, so keep the number of distinct consumer group names small the same way you would for any other Prometheus label. Values cover only the current process. Prometheus combines several server processes when a query sums their scraped series. Queue depth is refreshed from shared storage every five seconds. Each server reports the same shared value, so use max rather than sum across Prometheus scrape instances.

Re-exports§

pub use anyhow;
pub use async_trait;
pub use futures;

Modules§

backend
Valid state and delivery backend combinations.
core
Traits used by job payloads and generated job handlers.
encoder
MessagePack encoding helpers used by Later’s wire format. Encoding helpers for payloads and internal records.
instrument
Attach a span to a std::future::Future.
mq
Low-level job delivery interfaces and built-in delivery implementations.
retry
Retry limits and delay strategies for failed jobs. Retry policies for failed job handlers.
storage
Low-level state storage interfaces and built-in implementations.
topic
Named topics and Kafka-like partitions for strictly ordered jobs. Named topics and Kafka-like partitions for strictly ordered jobs.

Macros§

background_job
Defines the payload types handled by one Later server.

Structs§

BackgroundJobServer
A running job server and its worker tasks.
BackgroundJobServerPublisher
Enqueues jobs and exposes dashboard and metrics data.
Config
Settings used by a generated job server builder.
JobId
The stable identifier assigned to one job execution.
RecurringJobId
The caller-selected identifier for a recurring job definition.
RetainedLogLag
Wires an optional [retained_log::LogLagSource] into a Config, gated by the retained-log feature so Config itself needs no feature flag. Build with [RetainedLogLag::new], or use Default::default to leave lag reporting disabled.
ServerConfig
Runtime settings consumed by BackgroundJobServer::start.

Functions§

generate_id
Generates a lexicographically sortable ULID string.

Type Aliases§

UtcDateTime
A UTC timestamp used in persisted job records.

Attribute Macros§

instrument
Instruments a function to create and enter a tracing span every time the function is called.