Qrush
v3.0.1Lightweight job queue & task scheduler
Production-ready background jobs for Rust on Redis + Tokio. The core is web-framework agnostic; the optional dashboard works with either Actix Web or Axum. Integrated or separate worker modes, cron scheduling, delayed jobs, retries with backoff, and a built-in metrics UI. Lean by default — the CLI and worker binaries live behind a `cli` feature, so embedding qrush as a library pulls in no argument parser or log subscriber.
cargo add qrush// documentation
Qrush docs
Overview
A lightweight, production-ready job queue and task scheduler for Rust applications built on Redis and Tokio. The core is web-framework agnostic, and the optional built-in dashboard works with either Actix Web or Axum. Qrush provides both integrated and separate process modes, making it suitable for everything from simple background tasks to large-scale distributed systems.
Features
- Dual Deployment Modes: integrated (single process) or separate worker process
- Framework Choice: optional dashboard for Actix Web or Axum; the queue/worker core needs neither
- High Performance: built on Redis and Tokio for maximum throughput
- Cron Scheduling: full cron expression support for recurring tasks
- Delayed Jobs: schedule jobs to run after a specified delay
- Built-in Metrics UI: real-time dashboard for monitoring queues, jobs, and workers
- Security: optional Basic Auth for metrics endpoints
- Type-Safe: leverages Rust’s type system for safe job handling
- Graceful Shutdown: clean worker shutdown with configurable grace periods
- Scalable: support for multiple queues with different priorities and concurrency levels
Feature Flags
The built-in dashboard is optional and works with either Actix or Axum — pick the one that matches your app.
| Feature | Default | Description |
|---|---|---|
| dashboard-actix | off | Metrics dashboard served with Actix Web (qrush::routes::metrics_route). Pulls in Actix Web, Tera, and the web stack. |
| dashboard-axum | off | Metrics dashboard served with Axum (qrush::routes::axum_route). Pulls in Axum, Tera, and the web stack. |
| dashboard | off | Back-compat alias for dashboard-actix. |
Library-only usage (default). No dashboard framework is enabled by default, so a plain dependency gives you enqueue + workers with no web stack:
[dependencies]
qrush = "2.1.0"To mount the dashboard, opt into one framework:
# Actix
qrush = { version = "2.1.0", features = ["dashboard-actix"] }
# Axum
qrush = { version = "2.1.0", features = ["dashboard-axum"] }Migrating from 1.x to 2.0
In 1.x the dashboard was Actix-only and enabled by default. In 2.0 it is framework-selectable and off by default. Nothing else changed — the route-wiring function and all queue/worker/cron APIs are the same.
| 1.x | 2.0 | |
|---|---|---|
| Dashboard default | on (Actix) | off |
| Enable Actix dashboard | (default) | features = ["dashboard-actix"] |
| Enable Axum dashboard | not available | features = ["dashboard-axum"] |
# 1.x
qrush = "1.0.1"
# 2.0 — Actix (equivalent to the old default; no code changes needed)
qrush = { version = "2.1.0", features = ["dashboard-actix"] }If you only used enqueue + workers (no dashboard), a plain qrush = "2.1.0" now pulls in less — the web stack is no longer compiled by default. See the CHANGELOG for the full list of changes.
Installation
Add to your Cargo.toml:
[dependencies]
qrush = "2.1.0"
tokio = { version = "1", features = ["rt-multi-thread", "macros"] }
serde = { version = "1", features = ["derive"] }
async-trait = "0.1"
anyhow = "1"
futures = "0.3"Basic Usage (Integrated Mode)
use qrush::job::Job;
use qrush::queue::{enqueue, enqueue_in};
use qrush::config::QueueConfig;
use qrush::registry::register_job;
use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use futures::future::BoxFuture;
use anyhow::Result;
#[derive(Clone, Serialize, Deserialize)]
pub struct EmailJob {
pub to: String,
pub subject: String,
}
#[async_trait]
impl Job for EmailJob {
async fn perform(&self) -> Result<()> {
println!("Sending email to {}: {}", self.to, self.subject);
// Your email sending logic here
Ok(())
}
fn name(&self) -> &'static str { "EmailJob" }
fn queue(&self) -> &'static str { "default" }
}
impl EmailJob {
pub fn name() -> &'static str { "EmailJob" }
pub fn handler(payload: String) -> BoxFuture<'static, Result<Box<dyn Job>>> {
Box::pin(async move {
let job: EmailJob = serde_json::from_str(&payload)?;
Ok(Box::new(job) as Box<dyn Job>)
})
}
}
#[tokio::main]
async fn main() -> Result<()> {
// Set Redis URL
std::env::set_var("REDIS_URL", "redis://127.0.0.1:6379");
// Register job
register_job(EmailJob::name(), EmailJob::handler);
// Initialize queues
let queues = vec![
QueueConfig::new("default", 5, 0),
];
QueueConfig::initialize(
"redis://127.0.0.1:6379".to_string(),
queues
).await?;
// Enqueue a job
enqueue(EmailJob {
to: "user@example.com".to_string(),
subject: "Hello!".to_string(),
}).await?;
// Keep running
tokio::signal::ctrl_c().await?;
Ok(())
}Recommended Project Layout (qrushes/ module)
The snippets above inline everything into main to stay short. In a real app you’ll want main to stay minimal and keep all qrush wiring in one place. The convention used by the reference demos is a self-contained qrushes/ module: main only calls qrushes::initiate::initiate(), and every job, cron, and piece of configuration lives under qrushes/.
src/
├── main.rs # calls qrushes::initiate::initiate() — nothing else qrush-related
└── qrushes/
├── mod.rs # pub mod crons; pub mod initiate; pub mod jobs;
├── initiate.rs # ALL wiring: Redis URL, auth, register jobs+crons, init queues
├── jobs/
│ ├── mod.rs
│ └── send_email_job.rs
└── crons/
├── mod.rs
├── interval_1minutes_notify_slack_cron.rs
└── interval_2minutes_notify_slack_cron.rsmain.rs — minimal. Everything qrush-specific collapses to a single call. The only framework-specific line left in main is mounting the dashboard route.
Actix (features = ["dashboard-actix"]):
mod qrushes;
use actix_web::{web, App, HttpServer};
use qrush::routes::metrics_route::qrush_metrics_routes;
#[actix_web::main]
async fn main() -> std::io::Result<()> {
dotenvy::dotenv().ok();
// All qrush wiring (Redis, dashboard auth, jobs, crons, queues) lives here.
qrushes::initiate::initiate().await.expect("qrush init failed");
HttpServer::new(|| {
App::new().service(web::scope("/qrush").configure(qrush_metrics_routes))
})
.bind("0.0.0.0:8080")?
.run()
.await
}Axum (features = ["dashboard-axum"]):
mod qrushes;
use axum::Router;
use qrush::routes::axum_route::qrush_metrics_router;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
dotenvy::dotenv().ok();
// All qrush wiring (Redis, dashboard auth, jobs, crons, queues) lives here.
qrushes::initiate::initiate().await?;
let app = Router::new().nest("/qrush", qrush_metrics_router());
let listener = tokio::net::TcpListener::bind("0.0.0.0:8080").await?;
axum::serve(listener, app).await?;
Ok(())
}qrushes/initiate.rs — the single entry point. initiate() owns the four-step boot sequence (set Redis URL → register → register crons → initialize) plus the optional dashboard auth. It is framework-agnostic — the same file works under Actix and Axum.
use qrush::config::{set_basic_auth, set_redis_url, QrushBasicAuthConfig, QueueConfig};
use qrush::cron::cron_scheduler::CronScheduler;
use qrush::queue::{enqueue, enqueue_in};
use qrush::registry::register_job;
use crate::qrushes::crons::interval_1minutes_notify_slack_cron::Interval1MinutesNotifySlackCron;
use crate::qrushes::crons::interval_2minutes_notify_slack_cron::Interval2MinutesNotifySlackCron;
use crate::qrushes::jobs::send_email_job::SendEmailJob;
/// Parse a `user:password` pair for the optional `QRUSH_BASIC_AUTH` gate.
fn parse_user_pass(value: Option<&str>) -> Option<(String, String)> {
let (user, pass) = value?.split_once(':')?;
if user.is_empty() { return None; }
Some((user.to_string(), pass.to_string()))
}
/// Configure and start qrush. Reads `REDIS_URL` and the optional
/// `QRUSH_BASIC_AUTH`, registers jobs + crons, initializes the queues, and
/// seeds a couple of demo jobs so the dashboard has data.
pub async fn initiate() -> anyhow::Result<()> {
let redis_url = std::env::var("REDIS_URL")
.unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string());
set_redis_url(redis_url.clone())?; // 1
// Optional dashboard auth: QRUSH_BASIC_AUTH=user:password protects /qrush.
if let Some((username, password)) =
parse_user_pass(std::env::var("QRUSH_BASIC_AUTH").ok().as_deref())
{
set_basic_auth(Some(QrushBasicAuthConfig { username, password }));
}
register_job(SendEmailJob::type_name(), SendEmailJob::handler); // 2
register_job(
Interval1MinutesNotifySlackCron::type_name(),
Interval1MinutesNotifySlackCron::handler,
);
register_job(
Interval2MinutesNotifySlackCron::type_name(),
Interval2MinutesNotifySlackCron::handler,
);
// Register the schedules. Restart-safe: re-registering an existing cron_id
// returns an error we log and treat as a no-op instead of aborting startup.
if let Err(e) = CronScheduler::register_cron_job(Interval1MinutesNotifySlackCron {
label: "minutely slack notify".into(),
}).await {
println!("cron job already registered: {e}"); // 3
}
if let Err(e) = CronScheduler::register_cron_job(Interval2MinutesNotifySlackCron {
label: "2-minutely slack notify".into(),
}).await {
println!("cron job already registered: {e}");
}
let queues = vec![QueueConfig::new("default", 5, 0)]; // 4
QueueConfig::initialize(redis_url, queues).await?;
// Optional: seed a job or two so the dashboard isn't empty on first boot.
let _ = enqueue(SendEmailJob {
to: "user@example.com".into(),
subject: "Immediate hello".into(),
}).await;
let _ = enqueue_in(
SendEmailJob { to: "user@example.com".into(), subject: "Delayed hello".into() },
60,
).await;
Ok(())
}qrushes/jobs/send_email_job.rs — one job per file. Each job is a plain Job plus a type_name()/handler() pair so a worker can rebuild it from Redis. type_name() must match name().
use async_trait::async_trait;
use futures::future::BoxFuture;
use serde::{Deserialize, Serialize};
use qrush::job::Job;
#[derive(Serialize, Deserialize)]
pub struct SendEmailJob {
pub to: String,
pub subject: String,
}
#[async_trait]
impl Job for SendEmailJob {
async fn perform(&self) -> anyhow::Result<()> {
println!("Sending email to {} -> {}", self.to, self.subject);
Ok(())
}
fn name(&self) -> &'static str { "SendEmailJob" }
fn queue(&self) -> &'static str { "default" }
}
impl SendEmailJob {
pub fn type_name() -> &'static str { "SendEmailJob" }
pub fn handler(payload: String) -> BoxFuture<'static, anyhow::Result<Box<dyn Job>>> {
Box::pin(async move {
let job: SendEmailJob = serde_json::from_str(&payload)?;
Ok(Box::new(job) as Box<dyn Job>)
})
}
}qrushes/crons/interval_1minutes_notify_slack_cron.rs — one cron per file. A cron file is the same as a job file plus a CronJob impl (a cron_expression + a unique cron_id). Here perform() does real work — POSTing to a Slack incoming webhook (the HTTP equivalent of a curl POST with a JSON {"text": "…"} body):
use async_trait::async_trait;
use futures::future::BoxFuture;
use serde::{Deserialize, Serialize};
use serde_json::json;
use qrush::cron::cron_job::CronJob;
use qrush::job::Job;
#[derive(Serialize, Deserialize)]
pub struct Interval1MinutesNotifySlackCron {
pub label: String,
}
#[async_trait]
impl Job for Interval1MinutesNotifySlackCron {
async fn perform(&self) -> anyhow::Result<()> {
let webhook = std::env::var("SLACK_WEBHOOK_URL")?;
let resp = reqwest::Client::new()
.post(&webhook)
.json(&json!({ "text": format!("Hello, World! ({})", self.label) }))
.send()
.await?;
if !resp.status().is_success() {
anyhow::bail!("slack webhook returned {}", resp.status()); // -> retry
}
Ok(())
}
fn name(&self) -> &'static str { "Interval1MinutesNotifySlackCron" }
fn queue(&self) -> &'static str { "default" }
}
impl Interval1MinutesNotifySlackCron {
pub fn type_name() -> &'static str { "Interval1MinutesNotifySlackCron" }
pub fn handler(payload: String) -> BoxFuture<'static, anyhow::Result<Box<dyn Job>>> {
Box::pin(async move {
let job: Interval1MinutesNotifySlackCron = serde_json::from_str(&payload)?;
Ok(Box::new(job) as Box<dyn Job>)
})
}
}
#[async_trait]
impl CronJob for Interval1MinutesNotifySlackCron {
fn cron_expression(&self) -> &'static str { "0 * * * * *" } // every minute
fn cron_id(&self) -> &'static str { "interval_1min_notify_slack" } // unique per cron
}The 2-minute variant is identical apart from cron_expression ("0 */2 * * * *") and a distinct cron_id — each CronJob needs its own cron_id, or the second registration collides with the first in Redis.
Architecture
QRush supports two deployment modes.
Integrated Mode — workers run in the same process as your application. Perfect for small to medium applications.
┌─────────────────────┐
│ Application │
│ (Single Process) │
│ │
│ • HTTP Server │
│ • Enqueue Jobs │
│ • Process Jobs │ ← Workers here
└─────────────────────┘Separate Process Mode — workers run in a dedicated process. Recommended for production environments.
┌─────────────────────┐ ┌─────────────────────┐
│ Web Server │ │ qrush-engine │
│ (cargo run) │ │ (separate process) │
│ │ │ │
│ • HTTP Server │ │ • Worker Pools │
│ • Enqueue Jobs ─────┼──Redis──┼─▶ Process Jobs │
│ • Serve Routes │ │ • Cron Scheduler │
└─────────────────────┘ └─────────────────────┘API Reference
Core Traits:
- Job: implement this trait for your job types. Only perform, name, and queue are required; the before/after/on_error/always lifecycle hooks are optional overrides.
- CronJob: implement for recurring scheduled jobs.
Core Functions:
- enqueue(job) -> QrushResult<String>: enqueue a job immediately; returns the job ID
- enqueue_in(job, delay_secs) -> QrushResult<String>: enqueue a job with a delay; returns the job ID
- register_job(name, handler): register a job handler
- QueueConfig::initialize(redis_url, queues): start worker pools and the cron scheduler
- set_basic_auth(Some(QrushBasicAuthConfig { .. })): protect the dashboard with HTTP Basic Auth
Failed jobs are retried automatically with exponential backoff and moved to a dead-letter queue after MAX_RETRIES (3).
Cron Scheduling — all under qrush::cron::cron_scheduler::CronScheduler:
- register_cron_job(job) -> Result<()>: persist a schedule to Redis
- list_cron_jobs() -> Result<Vec<CronJobMeta>>: list registered cron jobs
- run_now(cron_id) -> Result<String>: enqueue a cron job immediately
- toggle_cron_job(cron_id, enabled) -> Result<()>: pause / resume a schedule
- delete_cron_job(cron_id) -> Result<()>: remove a schedule
Errors. The public API returns QrushResult<T> (Result<T, QrushError>). QrushError distinguishes Redis, Serialization, and Config failures, and implements std::error::Error, so it still propagates through ? in anyhow-based code.
Engine Runtime:
- qrush::engine::run_engine(redis_url, queues, shutdown_grace_secs): run worker process
- qrush::engine::parse_queues(spec): parse queue specification string
Command-Line Interface. The crate also ships reference binaries — qrush (a management CLI with start/stop/status/stats/queues/jobs subcommands) and qrush-engine (the worker process) — that you can adapt for your own app. See src/bin/cli.md for the full CLI guide.
Runnable Dashboard Examples
The repo ships a complete, runnable dashboard example for each framework. With a Redis instance available (REDIS_URL, defaults to redis://127.0.0.1:6379):
# Actix — serves http://127.0.0.1:8080/qrush/metrics
cargo run --example actix_dashboard --features dashboard-actix
# Axum — serves http://127.0.0.1:8080/qrush/metrics
cargo run --example axum_dashboard --features dashboard-axumJob Lifecycle Hooks
Beyond perform, the Job trait exposes optional hooks that wrap each execution. All are async and have default no-op implementations, so you only override the ones you need:
| Hook | When it runs | Signature | Notes |
|---|---|---|---|
| before | Before perform | async fn before(&self) -> Result<()> | Return Err to skip the job — it is marked skipped (a terminal, non-failure state) and perform never runs. |
| perform | The actual work | async fn perform(&self) -> Result<()> | Return Err to trigger retry / dead-letter. |
| after | After a successful perform | async fn after(&self) | Skipped if perform errored. |
| on_error | After a failed perform | async fn on_error(&self, err: &anyhow::Error) | Runs before the retry is scheduled. Good for logging/alerting. |
| always | After every attempt that ran perform | async fn always(&self) | Runs on both success and failure (but not when before skipped the job). |
use qrush::job::Job;
use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use anyhow::{bail, Result};
#[derive(Clone, Serialize, Deserialize)]
pub struct ChargeCard {
pub user_id: String,
pub amount_cents: u64,
}
#[async_trait]
impl Job for ChargeCard {
// Guard: bail out early (job is marked `skipped`, not `failed`).
async fn before(&self) -> Result<()> {
if self.amount_cents == 0 {
bail!("nothing to charge — skipping");
}
Ok(())
}
async fn perform(&self) -> Result<()> {
println!("Charging {} cents to {}", self.amount_cents, self.user_id);
// ... call payment gateway; return Err to retry ...
Ok(())
}
async fn after(&self) {
println!("charge succeeded — sending receipt");
}
async fn on_error(&self, err: &anyhow::Error) {
eprintln!("charge failed, will retry: {err}");
}
async fn always(&self) {
println!("charge attempt finished (success or failure)");
}
fn name(&self) -> &'static str { "ChargeCard" }
fn queue(&self) -> &'static str { "critical" }
}Delayed Jobs
enqueue_in(job, delay_secs) runs a job after a delay instead of immediately. It returns the job ID and is otherwise identical to enqueue — same job type, same worker, same retry semantics.
use qrush::queue::{enqueue, enqueue_in};
// Run now.
let id = enqueue(EmailJob {
to: "user@example.com".into(),
subject: "Welcome!".into(),
}).await?;
// Run in 10 minutes (600 seconds).
let id = enqueue_in(EmailJob {
to: "user@example.com".into(),
subject: "Don't forget to verify your email".into(),
}, 600).await?;Delayed jobs sit in a Redis sorted set keyed by their run-at timestamp; a dedicated delayed-worker pool (started by QueueConfig::initialize) promotes them onto their queue once due. Precision is bounded by the poll interval, so treat the delay as "at least N seconds", not an exact wall-clock alarm.
Retries & Dead-Letter Queue
When perform returns Err, QRush retries the job automatically — you don’t schedule retries yourself:
- 1. on_error is called, and the error string is stored on the job.
- 2. The job’s retry counter increments. While it’s <= 3 (MAX_RETRIES), the job is re-queued with exponential backoff plus jitter (10s * 2^retries, jittered to avoid thundering-herd retries) and its status becomes retrying.
- 3. After the 3rd retry is exhausted, the job moves to the dead-letter queue (status = dead) instead of being dropped. Inspect and requeue dead jobs from the dashboard at /qrush/metrics/extras/dead (or the dead-jobs view).
Job status values you’ll see in Redis / on the dashboard:
| Status | Meaning |
|---|---|
| pending | Enqueued, waiting for a worker |
| delayed | Scheduled via enqueue_in, not yet due |
| retrying | Failed once or more; waiting for its backoff to elapse |
| skipped | before() returned Err; terminal, treated as a non-failure |
| success | perform() completed successfully |
| dead | Retries exhausted; parked in the dead-letter queue |
| failed | Could not run at all (e.g. no handler registered for the job name) |
Cron Jobs
A cron job is a regular Job that runs on a schedule instead of being enqueued by hand. The work still lives in Job::perform; CronJob only adds when to run it. Follow these three steps.
Step 1 — Define the job and its perform(). This is identical to any other QRush job: implement Job (the work + a handler so a worker can rebuild it from Redis).
use qrush::job::Job;
use qrush::registry::register_job;
use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use futures::future::BoxFuture;
use anyhow::Result;
#[derive(Clone, Serialize, Deserialize)]
pub struct EmailJob {
pub to: String,
pub subject: String,
}
#[async_trait]
impl Job for EmailJob {
// This is the execution — it runs every time the schedule fires.
async fn perform(&self) -> Result<()> {
println!("Sending email to {}: {}", self.to, self.subject);
// Your recurring work goes here.
Ok(())
}
fn name(&self) -> &'static str { "EmailJob" }
fn queue(&self) -> &'static str { "default" }
}
impl EmailJob {
pub fn name() -> &'static str { "EmailJob" }
// Lets a worker rebuild the job from its stored payload.
pub fn handler(payload: String) -> BoxFuture<'static, Result<Box<dyn Job>>> {
Box::pin(async move {
let job: EmailJob = serde_json::from_str(&payload)?;
Ok(Box::new(job) as Box<dyn Job>)
})
}
}Step 2 — Add the schedule (CronJob). Attach a cron expression and a unique id to the same type.
use qrush::cron::cron_job::CronJob;
#[async_trait]
impl CronJob for EmailJob {
// 6-field: sec min hour day month weekday. See "Cron Expressions" below.
fn cron_expression(&self) -> &'static str { "0 0 * * * *" } // every hour
fn cron_id(&self) -> &'static str { "hourly_email" } // must be unique
}Step 3 — Register and start it in main. The job above is framework-agnostic; only main differs. In both frameworks the order is the same:
- 1. set_redis_url(...) — required before any Redis call.
- 2. register_job(...) — so a worker can run the job.
- 3. CronScheduler::register_cron_job(...) — saves the schedule to Redis.
- 4. QueueConfig::initialize(...) — starts the workers and the cron scheduler.
Actix (features = ["dashboard-actix"]):
use qrush::config::{set_redis_url, QueueConfig};
use qrush::registry::register_job;
use qrush::cron::cron_scheduler::CronScheduler;
use qrush::routes::metrics_route::qrush_metrics_routes;
use actix_web::{web, App, HttpServer};
#[actix_web::main]
async fn main() -> anyhow::Result<()> {
let redis_url = std::env::var("REDIS_URL")
.unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string());
set_redis_url(redis_url.clone())?; // 1
register_job(EmailJob::name(), EmailJob::handler); // 2
let job = EmailJob { // 3
to: "user@example.com".into(),
subject: "Hourly report".into(),
};
CronScheduler::register_cron_job(job).await?;
let queues = vec![QueueConfig::new("default", 5, 0)]; // 4
QueueConfig::initialize(redis_url, queues).await?;
// Serve the dashboard at http://127.0.0.1:8080/qrush/metrics
HttpServer::new(|| {
App::new().service(web::scope("/qrush").configure(qrush_metrics_routes))
})
.bind("0.0.0.0:8080")?
.run()
.await?;
Ok(())
}Axum (features = ["dashboard-axum"]):
use qrush::config::{set_redis_url, QueueConfig};
use qrush::registry::register_job;
use qrush::cron::cron_scheduler::CronScheduler;
use qrush::routes::axum_route::qrush_metrics_router;
use axum::Router;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
let redis_url = std::env::var("REDIS_URL")
.unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string());
set_redis_url(redis_url.clone())?; // 1
register_job(EmailJob::name(), EmailJob::handler); // 2
let job = EmailJob { // 3
to: "user@example.com".into(),
subject: "Hourly report".into(),
};
CronScheduler::register_cron_job(job).await?;
let queues = vec![QueueConfig::new("default", 5, 0)]; // 4
QueueConfig::initialize(redis_url, queues).await?;
// Serve the dashboard at http://127.0.0.1:8080/qrush/metrics
let app = Router::new().nest("/qrush", qrush_metrics_router());
let listener = tokio::net::TcpListener::bind("0.0.0.0:8080").await?;
axum::serve(listener, app).await?;
Ok(())
}That’s it. When the schedule fires, QRush enqueues the job onto its queue() and a worker runs perform(). Watch it (and manage schedules) on the dashboard at /qrush/metrics/extras/cron.
Managing Cron Jobs
Manage schedules from the dashboard at /qrush/metrics/extras/cron, or programmatically via CronScheduler:
use qrush::cron::cron_scheduler::CronScheduler;
CronScheduler::list_cron_jobs().await?; // -> Vec<CronJobMeta>
CronScheduler::run_now("hourly_email").await?; // enqueue once, right now
CronScheduler::toggle_cron_job("hourly_email", false).await?; // pause
CronScheduler::toggle_cron_job("hourly_email", true).await?; // resume
CronScheduler::delete_cron_job("hourly_email").await?; // remove entirelyTo register a job that starts paused, override enabled() on the CronJob impl (it defaults to true); enable it later from the dashboard or with toggle_cron_job:
fn enabled(&self) -> bool { false }A disabled job stays registered but is skipped and removed from the run schedule until re-enabled.
Multiple Queues
let queues = vec![
QueueConfig::new("default", 5, 0), // 5 workers, priority 0
QueueConfig::new("critical", 10, 0), // 10 workers, priority 0
QueueConfig::new("low", 2, 1), // 2 workers, priority 1
];Metrics UI
Access the built-in metrics dashboard at /qrush/metrics:
- Queue statistics and job counts
- Worker status and health
- Cron job management
- Job retry and deletion
- CSV export
Requirements
- Rust 1.89.0 or later
- Redis 6.0 or later
- Tokio runtime (multi-threaded)
Environment Variables
QRush itself only reads REDIS_URL (and only where you pass it — most APIs take the URL explicitly). The other variables below are conventions used by the example binaries; your code decides whether to read them.
# Read by qrush where a Redis URL is expected
REDIS_URL=redis://127.0.0.1:6379
# Conventions (you read these yourself — see the sections linked)
QRUSH_BASIC_AUTH=admin:password # dashboard auth — you parse it and call set_basic_auth()
RUST_LOG=info,qrush=info # tracing filter, honored by tracing_subscriberIntegrated Mode (Detailed)
Use this mode when: you want a simple setup with workers running in the same process as your web server.
1. Add Dependencies
[dependencies]
# Pick the dashboard framework you use: "dashboard-actix" or "dashboard-axum"
qrush = { version = "2.1.0", features = ["dashboard-actix"] }
actix-web = "4" # or: axum = "0.8"
tokio = { version = "1", features = ["rt-multi-thread", "macros"] }
serde = { version = "1", features = ["derive"] }
async-trait = "0.1"
anyhow = "1"
futures = "0.3"2. Define a Job
use qrush::job::Job;
use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use futures::future::BoxFuture;
use anyhow::Result;
#[derive(Clone, Serialize, Deserialize)]
pub struct NotifyUser {
pub user_id: String,
pub message: String,
}
#[async_trait]
impl Job for NotifyUser {
async fn perform(&self) -> Result<()> {
println!("Notify {} -> {}", self.user_id, self.message);
Ok(())
}
fn name(&self) -> &'static str { "NotifyUser" }
fn queue(&self) -> &'static str { "default" }
}
impl NotifyUser {
pub fn name() -> &'static str { "NotifyUser" }
pub fn handler(payload: String) -> BoxFuture<'static, Result<Box<dyn Job>>> {
Box::pin(async move {
let job: NotifyUser = serde_json::from_str(&payload)?;
Ok(Box::new(job) as Box<dyn Job>)
})
}
}3. Initialize QRush. The queue/worker setup is identical for both frameworks — only the dashboard wiring differs. The dashboard mounts at /qrush/metrics/... in both cases.
Actix (features = ["dashboard-actix"]):
use qrush::config::{QueueConfig, set_redis_url};
use qrush::registry::register_job;
use qrush::routes::metrics_route::qrush_metrics_routes;
use actix_web::{web, App, HttpServer};
#[actix_web::main]
async fn main() -> std::io::Result<()> {
// Set Redis URL
let redis_url = std::env::var("REDIS_URL")
.unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string());
set_redis_url(redis_url.clone())?;
// Register jobs
register_job(NotifyUser::name(), NotifyUser::handler);
// Initialize queues
let queues = vec![
QueueConfig::new("default", 5, 0),
];
QueueConfig::initialize(redis_url, queues).await?;
// Start web server
HttpServer::new(|| {
App::new()
.service(web::scope("/qrush").configure(qrush_metrics_routes))
})
.bind("0.0.0.0:8080")?
.run()
.await
}Axum (features = ["dashboard-axum"]):
use qrush::config::{QueueConfig, set_redis_url};
use qrush::registry::register_job;
use qrush::routes::axum_route::qrush_metrics_router;
use axum::Router;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
// Set Redis URL
let redis_url = std::env::var("REDIS_URL")
.unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string());
set_redis_url(redis_url.clone())?;
// Register jobs
register_job(NotifyUser::name(), NotifyUser::handler);
// Initialize queues
let queues = vec![
QueueConfig::new("default", 5, 0),
];
QueueConfig::initialize(redis_url, queues).await?;
// Mount the dashboard under /qrush
let app = Router::new().nest("/qrush", qrush_metrics_router());
let listener = tokio::net::TcpListener::bind("0.0.0.0:8080").await?;
axum::serve(listener, app).await?;
Ok(())
}4. Enqueue Jobs
use qrush::queue::{enqueue, enqueue_in};
// Immediate
enqueue(NotifyUser {
user_id: "123".to_string(),
message: "Hello".to_string(),
}).await?;
// Delayed (300 seconds)
enqueue_in(NotifyUser {
user_id: "123".to_string(),
message: "Reminder".to_string(),
}, 300).await?;Separate Process Mode (Detailed)
Use this mode when: you want production-ready separation with workers in a dedicated process.
Recommended Project Layout (qrushes_engines/ module). Integrated mode keeps all wiring in qrushes/. Separate process mode uses the same idea in a self-contained qrushes_engines/ module — the engine binary’s main only calls qrushes_engines::initiate::initiate(), and every job, cron, and piece of engine configuration lives under qrushes_engines/.
There are only two differences from the integrated qrushes/ layout:
- 1. initiate() ends with run_engine(...) instead of QueueConfig::initialize(...). run_engine starts the worker pools, the delayed-job handler, and the cron scheduler, then blocks until SIGINT/SIGTERM — so it is the last thing initiate() does, not a call it returns from.
- 2. Jobs and crons live in your crate’s library (src/lib.rs), because the engine process and the web server are two binaries that both need the same job/cron types. Put the module in the lib and both can use your_app::qrushes_engines::….
src/
├── lib.rs # pub mod qrushes_engines;
├── main.rs # web server: enqueue + dashboard, NO workers
├── bin/
│ └── qrush_engine.rs # worker process: calls qrushes_engines::initiate::initiate()
└── qrushes_engines/
├── mod.rs # pub mod initiate; pub mod jobs; pub mod crons;
├── initiate.rs # shared registry + two entry points: initiate_web() / initiate_engine()
├── jobs/
│ ├── mod.rs # pub mod send_email_job;
│ └── send_email_job.rs
└── crons/
├── mod.rs # pub mod interval_1minutes_notify_slack_cron;
└── interval_1minutes_notify_slack_cron.rssrc/lib.rs — expose the module to both binaries
pub mod qrushes_engines;src/bin/qrush_engine.rs — minimal worker process. Everything engine-specific collapses to a single call, exactly like main.rs does in integrated mode. Replace your_app with your crate’s name (the name under [package] in Cargo.toml).
use your_app::qrushes_engines;
#[tokio::main(flavor = "multi_thread")]
async fn main() -> anyhow::Result<()> {
dotenvy::dotenv().ok();
tracing_subscriber::fmt::init();
// All engine wiring (Redis, jobs, crons, queues) lives here; this call blocks
// until shutdown because run_engine() blocks.
qrushes_engines::initiate::initiate_engine().await
}src/qrushes_engines/mod.rs
pub mod initiate;
pub mod jobs;
pub mod crons;src/qrushes_engines/initiate.rs — two entry points, one registry. Both processes must register the same job/cron handlers — the web server to enqueue/serialize them, the engine to run them. So the register_job(...) list lives here once, in a shared register_all(), and two thin entry points build on it:
- initiate_web() — registers the handlers and returns. The web server calls this; it does not start workers.
- initiate_engine() — registers the handlers, registers the cron schedules, then calls run_engine(...), which starts the worker pools + delayed handler + cron scheduler and blocks until SIGINT/SIGTERM.
Keeping registration in one function means adding a job is a one-line change that both processes pick up — you can’t forget to register it in one of them.
use qrush::config::set_redis_url;
use qrush::cron::cron_scheduler::CronScheduler;
use qrush::engine::{parse_queues, run_engine};
use qrush::registry::register_job;
use crate::qrushes_engines::crons::interval_1minutes_notify_slack_cron::Interval1MinutesNotifySlackCron;
use crate::qrushes_engines::jobs::send_email_job::SendEmailJob;
/// Single source of truth for the type registry: set the Redis URL and register
/// every job + cron handler. Shared by both processes. Returns the Redis URL so
/// the engine can hand it to `run_engine`.
fn register_all() -> anyhow::Result<String> {
let redis_url = std::env::var("REDIS_URL")
.unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string());
set_redis_url(redis_url.clone())?;
register_job(SendEmailJob::type_name(), SendEmailJob::handler);
register_job(
Interval1MinutesNotifySlackCron::type_name(),
Interval1MinutesNotifySlackCron::handler,
);
Ok(redis_url)
}
/// Web-server entry point: register handlers so jobs can be enqueued, but do NOT
/// start workers — the engine process owns those.
pub async fn initiate_web() -> anyhow::Result<()> {
register_all()?;
Ok(())
}
/// Engine entry point: register handlers + cron schedules, then run the workers.
/// `run_engine` owns `QueueConfig::initialize` internally and BLOCKS until
/// shutdown (with a 5s graceful-shutdown grace period).
pub async fn initiate_engine() -> anyhow::Result<()> {
let redis_url = register_all()?;
// Restart-safe: re-registering an existing cron_id returns an error we log
// and treat as a no-op instead of aborting startup.
if let Err(e) = CronScheduler::register_cron_job(Interval1MinutesNotifySlackCron {
label: "minutely slack notify".into(),
}).await {
println!("cron job already registered: {e}");
}
let queues = parse_queues("default:5:0");
run_engine(redis_url, queues, 5).await
}These are the same job/cron types as integrated mode — a plain Job (plus a CronJob impl for crons) with a type_name()/handler() pair — so a job enqueued by the web server deserializes and runs in the engine process. They just live in the lib under qrushes_engines/ instead of qrushes/.
src/qrushes_engines/jobs/send_email_job.rs — one job per file
use async_trait::async_trait;
use futures::future::BoxFuture;
use serde::{Deserialize, Serialize};
use qrush::job::Job;
#[derive(Serialize, Deserialize)]
pub struct SendEmailJob {
pub to: String,
pub subject: String,
}
#[async_trait]
impl Job for SendEmailJob {
async fn perform(&self) -> anyhow::Result<()> {
println!("Sending email to {} -> {}", self.to, self.subject);
Ok(())
}
fn name(&self) -> &'static str { "SendEmailJob" }
fn queue(&self) -> &'static str { "default" }
}
impl SendEmailJob {
pub fn type_name() -> &'static str { "SendEmailJob" }
pub fn handler(payload: String) -> BoxFuture<'static, anyhow::Result<Box<dyn Job>>> {
Box::pin(async move {
let job: SendEmailJob = serde_json::from_str(&payload)?;
Ok(Box::new(job) as Box<dyn Job>)
})
}
}src/qrushes_engines/jobs/mod.rs just re-exports it:
pub mod send_email_job;src/qrushes_engines/crons/interval_1minutes_notify_slack_cron.rs — one cron per file
use async_trait::async_trait;
use futures::future::BoxFuture;
use serde::{Deserialize, Serialize};
use serde_json::json;
use qrush::cron::cron_job::CronJob;
use qrush::job::Job;
#[derive(Serialize, Deserialize)]
pub struct Interval1MinutesNotifySlackCron {
pub label: String,
}
#[async_trait]
impl Job for Interval1MinutesNotifySlackCron {
async fn perform(&self) -> anyhow::Result<()> {
let webhook = std::env::var("SLACK_WEBHOOK_URL")?;
let resp = reqwest::Client::new()
.post(&webhook)
.json(&json!({ "text": format!("Hello, World! ({})", self.label) }))
.send()
.await?;
if !resp.status().is_success() {
anyhow::bail!("slack webhook returned {}", resp.status()); // -> retry
}
Ok(())
}
fn name(&self) -> &'static str { "Interval1MinutesNotifySlackCron" }
fn queue(&self) -> &'static str { "default" }
}
impl Interval1MinutesNotifySlackCron {
pub fn type_name() -> &'static str { "Interval1MinutesNotifySlackCron" }
pub fn handler(payload: String) -> BoxFuture<'static, anyhow::Result<Box<dyn Job>>> {
Box::pin(async move {
let job: Interval1MinutesNotifySlackCron = serde_json::from_str(&payload)?;
Ok(Box::new(job) as Box<dyn Job>)
})
}
}
#[async_trait]
impl CronJob for Interval1MinutesNotifySlackCron {
fn cron_expression(&self) -> &'static str { "0 * * * * *" } // every minute
fn cron_id(&self) -> &'static str { "interval_1min_notify_slack" } // unique per cron
}src/qrushes_engines/crons/mod.rs just re-exports it:
pub mod interval_1minutes_notify_slack_cron;The web server (src/main.rs). The web server reuses the same module but does not start workers — its main calls qrushes_engines::initiate::initiate_web() (register handlers only), then mounts the dashboard. See Web Server (No Workers) below for the full Actix/Axum main.rs. The engine binary, its initiate.rs, and the job/cron files were all defined above. The two steps below just wire the two processes together — you do not create qrush_engine.rs again.
1. Register the engine binary in Cargo.toml. The engine binary above uses tracing_subscriber for logging and dotenvy to load .env, so add them alongside the [[bin]] entry. Adding this second binary makes a bare cargo run ambiguous (error: could not determine which binary to run), so set default-run to your web binary — then cargo run starts the web server and cargo run --bin qrush_engine starts the worker:
[package]
name = "your_app"
# ...
default-run = "your_app" # so a bare `cargo run` picks the web server, not the engine
[dependencies]
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
dotenvy = "0.15"
[[bin]]
name = "qrush_engine"
path = "src/bin/qrush_engine.rs"2. Web Server (No Workers). The web server’s main.rs calls initiate_web() — the same registry as the engine, minus the workers — then mounts the dashboard. Because initiate_web() never calls run_engine/QueueConfig::initialize, it returns immediately and the HTTP server starts. Just like integrated mode, main stays minimal.
Actix (features = ["dashboard-actix"]):
use your_app::qrushes_engines;
use qrush::routes::metrics_route::qrush_metrics_routes;
use actix_web::{web, App, HttpServer};
#[actix_web::main]
async fn main() -> std::io::Result<()> {
dotenvy::dotenv().ok();
// Set Redis + register the same job/cron handlers as the engine — but no
// workers. initiate_web() returns immediately.
qrushes_engines::initiate::initiate_web().await.expect("qrush init failed");
HttpServer::new(|| {
App::new().service(web::scope("/qrush").configure(qrush_metrics_routes))
})
.bind("0.0.0.0:8080")?
.run()
.await
}Axum (features = ["dashboard-axum"]):
use your_app::qrushes_engines;
use qrush::routes::axum_route::qrush_metrics_router;
use axum::Router;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
dotenvy::dotenv().ok();
// Set Redis + register the same job/cron handlers as the engine — but no
// workers. initiate_web() returns immediately.
qrushes_engines::initiate::initiate_web().await?;
let app = Router::new().nest("/qrush", qrush_metrics_router());
let listener = tokio::net::TcpListener::bind("0.0.0.0:8080").await?;
axum::serve(listener, app).await?;
Ok(())
}3. Run Both Processes
Terminal 1 - Web Server: (needs default-run from step 1; otherwise use cargo run --bin your_app)
export REDIS_URL=redis://127.0.0.1:6379
cargo runTerminal 2 - Worker Engine:
export REDIS_URL=redis://127.0.0.1:6379
cargo run --bin qrush_engineCron Expressions
QRush accepts both 6-field (sec min hour day month weekday) and 5-field (min hour day month weekday) expressions. A 5-field expression defaults seconds to 0, so */5 * * * * and 0 */5 * * * * are equivalent.
Common examples:
| Expression | Meaning |
|---|---|
| "* * * * *" | Every minute (5-field) |
| "0 * * * * *" | Every minute (6-field) |
| "0 */5 * * * *" | Every 5 minutes |
| "0 0 * * * *" | Every hour |
| "0 0 0 * * *" | Daily at midnight |
| "0 30 9 * * *" | Daily at 09:30 |
| "0 0 0 * * 1" | Every Monday at midnight |
| "0 0 9 * * MON-FRI" | Weekdays at 09:00 |
| "0 0 0 1 * *" | First day of every month |
| "0 0 12 1 JAN *" | Jan 1st at noon |
Each field supports the usual operators:
- * — any value
- a — an exact value
- a,b,c — a list
- a-b — an inclusive range
- */n — a step over the whole range (e.g. */15 in minutes)
- a-b/n — a step within a range
- Names: months JAN–DEC, weekdays SUN–SAT (case-insensitive). For the weekday field, both 0 and 7 mean Sunday.
Timezone. Expressions evaluate in UTC by default. Override per job with a timezone() method on the CronJob impl that returns any IANA name (e.g. "Asia/Kolkata", "America/New_York") — so "0 0 9 * * *" fires at 09:00 in that zone, DST included.
#[async_trait]
impl CronJob for EmailJob {
fn cron_expression(&self) -> &'static str { "0 0 9 * * *" } // 9 AM…
fn cron_id(&self) -> &'static str { "morning_email" }
fn timezone(&self) -> &'static str { "Asia/Kolkata" } // …IST
}Precision & missed runs. The scheduler ticks every ~5 seconds, so a job fires within a few seconds of its scheduled time (don’t rely on sub-5s precision). If the scheduler was down when a run was due, that run fires once on the next tick and is then re-anchored to its next future slot — missed cycles are not backfilled one-per-cycle. Claiming is atomic in Redis, so running multiple engine processes will not double-fire the same job.
Metrics Endpoints
Paths assume the dashboard is mounted at /qrush (as in the examples). The Actix and Axum adapters expose the same routes:
| Method & Path | Purpose |
|---|---|
| GET /qrush/metrics | Dashboard overview |
| GET /qrush/metrics/health | Health check (returns healthy) |
| GET /qrush/metrics/queues/{queue} | Per-queue details |
| GET /qrush/metrics/queues/{queue}/export | Export a queue’s jobs as CSV |
| GET /qrush/metrics/extras/summary | Aggregate metrics summary |
| GET /qrush/metrics/extras/delayed | Delayed (scheduled-later) jobs |
| GET /qrush/metrics/extras/scheduled | Scheduled jobs |
| GET /qrush/metrics/extras/retry | Jobs waiting to retry |
| GET /qrush/metrics/extras/failed | Failed jobs |
| GET /qrush/metrics/extras/dead | Dead-letter queue |
| GET /qrush/metrics/extras/cron | Cron job management |
| POST /qrush/metrics/jobs/action | Job actions (retry / delete) |
| POST /qrush/metrics/cron/action | Cron actions (run-now / toggle / delete) |
Securing the Dashboard (Basic Auth)
The dashboard is open by default. To require HTTP Basic Auth, register credentials with set_basic_auth before you start the web server. Once credentials are set, the built-in middleware (already wired into both the Actix and Axum routers) enforces them on every /qrush/metrics/... request using a constant-time credential comparison.
use qrush::config::{set_basic_auth, QrushBasicAuthConfig};
// Read from the environment (recommended) — the crate does NOT do this for you.
if let Ok(raw) = std::env::var("QRUSH_BASIC_AUTH") {
if let Some((username, password)) = raw.split_once(':') {
set_basic_auth(Some(QrushBasicAuthConfig {
username: username.to_string(),
password: password.to_string(),
}));
}
}
// ...then mount the dashboard and start the server as usual.- Call set_basic_auth once, during startup, before serving requests.
- Passing None (or never calling it) leaves the dashboard open.
- There’s no env-var auto-wiring: QRUSH_BASIC_AUTH is only a naming convention — you read it and call set_basic_auth yourself, as above.
- Basic Auth sends credentials base64-encoded, not encrypted. Terminate TLS in front of the dashboard (reverse proxy) for anything internet-facing.
Production Tips
- Use separate process mode for production
- Protect the dashboard with Basic Auth (and put TLS in front of it)
- Configure appropriate queue concurrency based on your workload
- Monitor Redis memory usage
- Use graceful shutdown for zero-downtime deployments
- Scale workers horizontally by running multiple engine processes
License & Support
This project is licensed under the MIT License — see the LICENSE file for details. Contributions are welcome; feel free to submit a Pull Request.
- Documentation: docs.rs/qrush
- Issues: github.com/srotas-space/qrush/issues
- Discussions: github.com/srotas-space/qrush/discussions
Ready to try Qrush?
cargo add qrush// more