<div align="center">
<img src="https://srotasspace.s3.ap-south-1.amazonaws.com/srotas.svg" alt="QRush" width="100" />
# Qrush
**Lightweight, production-ready job queue & task scheduler for Rust β Redis + Tokio, with an optional Actix / Axum dashboard.**
[](https://crates.io/crates/qrush)
[](https://docs.rs/qrush)
[](LICENSE)
[](https://crates.io/crates/qrush)
[](https://www.youtube.com/@srotas-space)
πΊ **[Watch the video walkthrough β](https://www.youtube.com/@srotas-space)**
</div>
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.
<details>
<summary><b>π Table of Contents</b></summary>
- [Features](#features)
- [Feature Flags](#feature-flags)
- [Quick Start](#quick-start)
- [Recommended Project Layout (`qrushes/`)](#recommended-project-layout-qrushes-module)
- [Architecture](#architecture)
- [API Reference](#api-reference)
- [Examples](#examples)
- [Job Lifecycle Hooks](#job-lifecycle-hooks)
- [Delayed Jobs](#delayed-jobs)
- [Retries & Dead-Letter Queue](#retries--dead-letter-queue)
- [Cron Jobs](#cron-jobs)
- [Integrated Mode (Detailed)](#integrated-mode-detailed)
- [Separate Process Mode (Detailed)](#separate-process-mode-detailed)
- [Cron Expressions](#cron-expressions)
- [Metrics Endpoints](#metrics-endpoints)
- [Securing the Dashboard](#securing-the-dashboard-basic-auth)
- [Production Tips](#production-tips)
</details>
## 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.
| `dashboard-actix` | β | Metrics dashboard served with Actix Web (`qrush::routes::metrics_route`). Pulls in Actix Web, Tera, and the web stack. |
| `dashboard-axum` | β | Metrics dashboard served with Axum (`qrush::routes::axum_route`). Pulls in Axum, Tera, and the web stack. |
| `dashboard` | β | 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:
```toml
[dependencies]
qrush = "3.0.0"
```
To mount the dashboard, opt into one framework:
```toml
# Actix
qrush = { version = "3.0.0", features = ["dashboard-actix"] }
# Axum
qrush = { version = "3.0.0", features = ["dashboard-axum"] }
```
### Migrating from 1.x to 2.0
<details>
<summary><b>π¦ Show the 1.x β 2.0 migration guide</b></summary>
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.
| Dashboard default | on (Actix) | off |
| Enable Actix dashboard | (default) | `features = ["dashboard-actix"]` |
| Enable Axum dashboard | not available | `features = ["dashboard-axum"]` |
```toml
# 1.x
qrush = "1.0.1"
# 2.0 β Actix (equivalent to the old default; no code changes needed)
qrush = { version = "3.0.0", features = ["dashboard-actix"] }
```
If you only used `enqueue` + workers (no dashboard), a plain `qrush = "3.0.0"`
now pulls in **less** β the web stack is no longer compiled by default. See the
[CHANGELOG](CHANGELOG.md) for the full list of changes.
</details>
## Quick Start
### Installation
Add to your `Cargo.toml`:
```toml
[dependencies]
qrush = "3.0.0"
tokio = { version = "1", features = ["rt-multi-thread", "macros"] }
serde = { version = "1", features = ["derive"] }
async-trait = "0.1"
anyhow = "1"
futures = "0.3"
```
> `qrush` bundles its own Redis client (with cluster support), so you don't need
> to depend on `redis` directly unless you use it yourself.
### Basic Usage (Integrated Mode)
```rust
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(())
}
```
---
## Integrated Process Mode
**Use this mode when:** You want production-ready in the same server process
### 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.rs
```
### `main.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"]`):
```rust
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"]`):
```rust
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(())
}
```
#### `src/qrushes/mod.rs`
```rust
pub mod initiate;
pub mod jobs;
pub mod crons;
```
### `qrushes/initiate.rs` β the single entry point
`initiate()` owns the four-step boot sequence ([set Redis URL β register β
register crons β initialize](#step-3--register-and-start-it-in-main)) plus the
optional dashboard auth. It is framework-agnostic β the same file works under
Actix and Axum.
```rust
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(())
}
```
> Because the cron registration is wrapped (log-and-continue instead of `?`),
> `initiate()` is **restart-safe** β see the [restart note](#step-3--register-and-start-it-in-main).
> If your app already has a `user:password` parser elsewhere, import that instead
> of the small `parse_user_pass` shown here.
`src/qrushes/jobs/mod.rs` just re-exports it:
```rust
pub mod send_email_job;
```
### `qrushes/jobs/send_email_job.rs` β one job per file
Each job is a plain [`Job`](#core-traits) plus a `type_name()`/`handler()` pair so
a worker can rebuild it from Redis. `type_name()` must match `name()`.
```rust
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/crons/crons/mod.rs` just re-exports it:
```rust
pub mod interval_1minutes_notify_slack_cron;
```
### `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`](#step-2--add-the-schedule-cronjob)
impl (a `cron_expression` + a **unique** `cron_id`). Here `perform()` does real
work β POSTing to a Slack incoming webhook, the HTTP equivalent of
`curl -X POST -H 'Content-type: application/json' --data '{"text":"β¦"}' <webhook>`:
```rust
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.
> Returning `Err` from a cron's `perform()` (e.g. a non-2xx webhook response)
> triggers the same [retry / dead-letter](#retries--dead-letter-queue) path as any
> other job. The Slack webhook needs an HTTP client β the demos use
> `reqwest = { version = "0.12", default-features = false, features = ["rustls-tls", "json"] }`.
## 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 β
βββββββββββββββββββββββ βββββββββββββββββββββββ
```
## Documentation
### Integrated Mode
See [Part 1: Integrated Mode](#integrated-mode-detailed) below for complete setup instructions.
### Separate Process Mode
See [Part 2: Separate Process Mode](#separate-process-mode-detailed) below for production deployment.
## 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](#job-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](#delayed-jobs); 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](#securing-the-dashboard-basic-auth) with HTTP Basic Auth
Failed jobs are [retried automatically](#retries--dead-letter-queue) with
exponential backoff and moved to a dead-letter queue after `MAX_RETRIES` (3).
### Cron Scheduling
All under `qrush::cron::cron_scheduler::CronScheduler` (see [Cron Jobs](#cron-jobs)):
- `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. They live behind the
`cli` feature so that library users don't compile an argument parser and a log
subscriber they never call:
```bash
cargo install qrush --features cli
```
See [`src/bin/cli.md`](src/bin/cli.md) for the full CLI guide.
## Examples
### 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`):
```sh
# 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-axum
```
### Job Lifecycle Hooks
Beyond `perform`, the [`Job`](#core-traits) 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:
<details>
<summary><b>πͺ Hook reference table</b></summary>
| `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](#retries--dead-letter-queue). |
| `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). |
</details>
```rust
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.
```rust
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:
| `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) |
> Retries and the dead-letter queue are handled by the **worker** process, so
> they apply wherever `QueueConfig::initialize` runs β the app in integrated
> mode, or the engine binary in [separate process mode](#separate-process-mode-detailed).
### Cron Jobs
A cron job is a regular [`Job`](#basic-usage-integrated-mode) 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).
```rust
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.
```rust
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.
> β οΈ The cron scheduler only runs where `QueueConfig::initialize` is called. In
> [Separate Process Mode](#separate-process-mode-detailed) that's the engine
> binary, not the web server β put steps 2β4 there.
**Actix** (`features = ["dashboard-actix"]`):
```rust
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"]`):
```rust
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`.
> **On restart:** `register_cron_job` errors if a job with the same `cron_id`
> already exists in Redis, so the `?` above would abort a second boot. The
> schedule already survives restarts, so either skip re-registering, treat the
> duplicate as non-fatal (log and continue instead of `?`), or call
> `CronScheduler::delete_cron_job("hourly_email")` first to re-seed it.
#### `register_job` vs `register_cron_job` β why a cron needs both
They answer two different questions β one is **how** to run a job type, the other
is **when** to run it:
- **`register_job(name, handler)`** β the *how*. It maps a job's name to a handler
that rebuilds the struct from its stored JSON so a **worker** can run
`perform()`. It lives in an in-memory registry, so it must be called on **every
startup, in every process that runs workers**. Miss it and the job fires but the
worker can't reconstruct it (status `failed`).
- **`CronScheduler::register_cron_job(job)`** β the *when*. It persists the
**schedule** (cron expression, timezone, next-run, payload) to Redis, so the
scheduler enqueues the job when it's due. Miss it and the handler exists but
nothing ever triggers it.
When a schedule fires, the stored payload is enqueued **tagged with the job name**,
and the worker resolves the handler via the registry `register_job` populated β
which is why a cron calls both.
| Stores | in-memory `HashMap` | Redis (durable) |
| Answers | *how* to rebuild & run | *when* to fire |
| Call frequency | every startup, every worker process | once (persists across restarts) |
| Duplicate handling | overwrites silently | **errors** if `cron_id` already exists |
| Needed by | **every** job a worker runs | **only** cron / scheduled jobs |
| Sync / async | sync | async (writes to Redis) |
**Rule of thumb:** a plain job = `register_job` + `enqueue(...)`; a cron job =
`register_job` **and** `register_cron_job`. In [separate process
mode](#separate-process-mode-detailed), `register_cron_job` must run where the
scheduler runs β the **engine**, not the web server.
#### Managing cron jobs
Manage schedules from the dashboard at `/qrush/metrics/extras/cron`, or
programmatically via `CronScheduler`:
```rust
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 entirely
```
To 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`:
```rust
fn enabled(&self) -> bool { false }
```
A disabled job stays registered but is skipped and removed from the run schedule
until re-enabled.
### Multiple Queues
```rust
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
> Requires a dashboard feature β `dashboard-actix` or `dashboard-axum` (not
> enabled by default). See [Feature Flags](#feature-flags).
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.
```bash
# 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_subscriber
```
> β οΈ Setting `QRUSH_BASIC_AUTH` alone does **nothing** β the crate never reads
> it. Dashboard auth is configured programmatically; see
> [Securing the Dashboard](#securing-the-dashboard-basic-auth).
# Detailed Documentation
## Integrated 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
```toml
[dependencies]
# Pick the dashboard framework you use: "dashboard-actix" or "dashboard-axum"
qrush = { version = "3.0.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
```rust
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"]`):
```rust
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"]`):
```rust
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
```rust
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/`](#recommended-project-layout-qrushes-module).
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.rs
```
#### `src/lib.rs` β expose the module to both binaries
```rust
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`).
```rust
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`
```rust
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.
```rust
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}");
}
// Some jobs to test
// enqueue instant
let _ = enqueue(SendEmailJob {
to: "user@example.com".into(),
subject: "Immediate hello".into(),
})
.await;
// enqueue later after 60 seconds
let _ = enqueue_in(
SendEmailJob {
to: "user@example.com".into(),
subject: "Delayed hello (30s)".into(),
},
60,
)
.await;
let queues = parse_queues("default:5:0");
run_engine(redis_url, queues, 5).await
}
```
> β οΈ The cron scheduler runs **only** where `run_engine` runs β the engine
> process, via `initiate_engine()`. `initiate_web()` deliberately skips both the
> cron registration and the workers.
These are the **same** job/cron types as integrated mode β a plain
[`Job`](#core-traits) (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
```rust
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:
```rust
pub mod send_email_job;
```
#### `src/qrushes_engines/crons/interval_1minutes_notify_slack_cron.rs` β one cron per file
```rust
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:
```rust
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)](#2-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:
```toml
[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"
```
> The web-server binary is named after your package (`src/main.rs` β the `name`
> under `[package]`). Without `default-run` you must always disambiguate:
> `cargo run --bin your_app`.
### 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"]`):
```rust
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"]`):
```rust
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`)
```bash
export REDIS_URL=redis://127.0.0.1:6379
cargo run
```
**Terminal 2 - Worker Engine:**
```bash
export REDIS_URL=redis://127.0.0.1:6379
cargo run --bin qrush_engine
```
---
## Cron 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.
<details>
<summary><b>ποΈ Common examples & field operators</b></summary>
| `"* * * * *"` | 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.
</details>
**Timezone.** Expressions evaluate in UTC by default. Override per job with
`fn timezone(&self) -> &'static str` on the `CronJob` impl, returning any IANA
name (e.g. `"Asia/Kolkata"`, `"America/New_York"`) β so `"0 0 9 * * *"` fires at
09:00 in that zone, DST included.
```rust
#[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:
<details>
<summary><b>π Full endpoint reference</b></summary>
| `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) |
</details>
## 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.
```rust
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](#securing-the-dashboard-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
This project is licensed under the MIT License - see the [LICENSE](LICENSE) file for details.
## Contributing
Contributions are welcome! Please feel free to submit a Pull Request.
## Support
- **Video walkthrough**: [@srotas-space on YouTube](https://www.youtube.com/@srotas-space)
- **Documentation**: [docs.rs/qrush](https://docs.rs/qrush)
- **Issues**: [GitHub Issues](https://github.com/srotas-space/qrush/issues)
- **Discussions**: [GitHub Discussions](https://github.com/srotas-space/qrush/discussions)
---
Made with β€οΈ by [Srotas Space](https://open-source.srotas.space)
---
## π₯ Contributors
- **[Sandeep Maurya](https://github.com/srotas-space)** - Creator & Lead Developer
<img src="https://srotasspace.s3.ap-south-1.amazonaws.com/snm.png" alt="Sandeep Maurya" width="80" height="80" style="border-radius: 50%;">
[LinkedIn](https://www.linkedin.com/in/snmmaurya/)
---
[](https://github.com/srotas-space/qrush)