use std::sync::Arc;
use std::sync::atomic::{AtomicU32, Ordering};
use std::time::Duration;
use async_trait::async_trait;
use serde_json::{Value, json};
use tokio::sync::Notify;
use tracing::error;
use crate::config::JobsConfig;
use crate::sqlite::db::Database;
use crate::sqlite::job::Job;
use crate::sqlite::nonce::now_secs;
pub mod registry;
pub mod runner;
pub mod sweep;
pub use registry::JobRegistry;
pub use runner::{spawn_runner, spawn_runner_watching};
pub use sweep::SweepJob;
#[derive(Debug)]
pub enum JobOutcome {
Done,
Retry(String),
Failed(String),
Reschedule(Duration),
}
#[derive(Debug, Clone)]
pub struct JobSpec {
pub kind: &'static str,
pub key: String,
pub payload: Value,
pub run_at: i64,
pub deadline: Option<i64>,
pub max_attempts: Option<u32>,
}
impl JobSpec {
#[must_use]
pub fn now(kind: &'static str, key: impl Into<String>) -> Self {
Self {
kind,
key: key.into(),
payload: json!({}),
run_at: now_secs(),
deadline: None,
max_attempts: None,
}
}
#[must_use]
pub fn with_payload(mut self, payload: Value) -> Self {
self.payload = payload;
self
}
#[must_use]
pub fn with_deadline(mut self, deadline: Option<i64>) -> Self {
self.deadline = deadline;
self
}
#[must_use]
pub fn with_delay(mut self, delay: Duration) -> Self {
self.run_at = now_secs().saturating_add(seconds(delay));
self
}
}
#[async_trait]
pub trait JobHandler: Send + Sync {
fn kind(&self) -> &'static str;
async fn run(&self, job: &Job) -> JobOutcome;
fn lease(&self) -> Option<Duration> {
None
}
async fn abandon(&self, _job: &Job, _reason: &str) {}
async fn recover(&self, _queue: &JobQueue) {}
}
#[derive(Clone)]
pub struct JobQueue {
database: Arc<Database>,
notify: Arc<Notify>,
default_max_attempts: Arc<AtomicU32>,
}
impl JobQueue {
#[must_use]
pub fn new(database: Arc<Database>, config: &JobsConfig) -> Self {
Self {
database,
notify: Arc::new(Notify::new()),
default_max_attempts: Arc::new(AtomicU32::new(config.max_attempts)),
}
}
pub fn set_max_attempts(&self, max_attempts: u32) {
self.default_max_attempts
.store(max_attempts, Ordering::Relaxed);
}
pub async fn enqueue(&self, spec: JobSpec) -> Result<bool, sqlx::Error> {
let max_attempts = i64::from(
spec.max_attempts
.unwrap_or_else(|| self.default_max_attempts.load(Ordering::Relaxed)),
);
let queued = Job::enqueue(
crate::sqlite::job::NewJob {
id: &uuid::Uuid::new_v4().to_string(),
kind: spec.kind,
dedup_key: &spec.key,
payload: &spec.payload,
run_at: spec.run_at,
deadline: spec.deadline,
max_attempts,
},
&self.database,
)
.await?;
if queued {
self.notify.notify_one();
}
Ok(queued)
}
pub async fn enqueue_or_log(&self, spec: JobSpec) -> bool {
let kind = spec.kind;
let key = spec.key.clone();
match self.enqueue(spec).await {
Ok(queued) => queued,
Err(error) => {
error!(
event = "job_enqueue_failed",
outcome = "failure",
job_kind = %kind,
dedup_key = %key,
error = %error,
);
false
}
}
}
#[must_use]
pub fn database(&self) -> &Arc<Database> {
&self.database
}
}
pub(crate) fn seconds(duration: Duration) -> i64 {
i64::try_from(duration.as_secs()).unwrap_or(i64::MAX)
}
#[cfg(test)]
mod tests {
use super::*;
async fn queue() -> JobQueue {
let database = Arc::new(Database::connect_in_memory().await.unwrap());
JobQueue::new(database, &JobsConfig::default())
}
#[tokio::test]
async fn enqueue_reports_whether_it_took_the_identity() {
let queue = queue().await;
assert!(queue.enqueue(JobSpec::now("test", "k")).await.unwrap());
assert!(
!queue.enqueue(JobSpec::now("test", "k")).await.unwrap(),
"a second live job for one key is refused, not duplicated"
);
assert!(queue.enqueue(JobSpec::now("test", "other")).await.unwrap());
}
#[tokio::test]
async fn a_spec_carries_its_payload_deadline_and_attempt_budget() {
let queue = queue().await;
let deadline = now_secs() + 60;
let spec = JobSpec {
max_attempts: Some(9),
..JobSpec::now("test", "k")
.with_payload(json!({"order_id": "ord-1"}))
.with_deadline(Some(deadline))
};
assert!(queue.enqueue(spec).await.unwrap());
let job = Job::find_live("test", "k", queue.database())
.await
.unwrap()
.unwrap();
assert_eq!(job.payload, json!({"order_id": "ord-1"}));
assert_eq!(job.deadline, Some(deadline));
assert_eq!(job.max_attempts, 9);
}
#[tokio::test]
async fn a_spec_with_no_budget_of_its_own_takes_the_configured_one() {
let queue = queue().await;
assert!(queue.enqueue(JobSpec::now("test", "k")).await.unwrap());
let job = Job::find_live("test", "k", queue.database())
.await
.unwrap()
.unwrap();
assert_eq!(
job.max_attempts,
i64::from(JobsConfig::default().max_attempts)
);
}
#[tokio::test]
async fn a_reloaded_max_attempts_reaches_the_clones_but_not_the_backlog() {
let queue = queue().await;
let held = queue.clone();
queue.enqueue(JobSpec::now("test", "before")).await.unwrap();
queue.set_max_attempts(11);
held.enqueue(JobSpec::now("test", "after")).await.unwrap();
let queued = |key: &'static str| {
let database = queue.database().clone();
async move {
Job::find_live("test", key, &database)
.await
.unwrap()
.unwrap()
.max_attempts
}
};
assert_eq!(
queued("before").await,
i64::from(JobsConfig::default().max_attempts),
"a row already queued keeps what it was queued under",
);
assert_eq!(
queued("after").await,
11,
"a clone taken before the change still enqueues under the new value",
);
}
#[tokio::test]
async fn with_delay_pushes_the_run_time_out() {
let queue = queue().await;
let spec = JobSpec::now("test", "k").with_delay(Duration::from_secs(600));
assert!(spec.run_at >= now_secs() + 599);
assert!(queue.enqueue(spec).await.unwrap());
}
#[tokio::test]
async fn enqueue_or_log_swallows_a_database_failure() {
let queue = queue().await;
queue.database().pool.close().await;
assert!(!queue.enqueue_or_log(JobSpec::now("test", "k")).await);
}
#[test]
fn seconds_saturates_rather_than_wrapping() {
assert_eq!(seconds(Duration::from_secs(90)), 90);
assert_eq!(seconds(Duration::MAX), i64::MAX);
}
}