pub(crate) mod backend;
pub mod create_rapina_jobs;
mod model;
pub(crate) mod retry;
pub(crate) mod worker;
pub(crate) use create_rapina_jobs::RapinaJobs;
pub use model::{JobRow, JobStatus};
pub use retry::RetryPolicy;
pub use worker::JobConfig;
use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use sea_orm::{ConnectionTrait, DatabaseConnection, DbBackend};
use uuid::Uuid;
use crate::state::AppState;
pub type JobId = Uuid;
pub type JobResult = Result<(), crate::error::Error>;
pub struct JobRequest {
pub job_type: &'static str,
pub payload: serde_json::Value,
pub queue: &'static str,
pub max_retries: i32,
}
#[doc(hidden)]
pub type JobHandlerFn =
fn(serde_json::Value, Arc<AppState>) -> Pin<Box<dyn Future<Output = JobResult> + Send>>;
pub struct JobDescriptor {
pub job_type: &'static str,
#[doc(hidden)]
pub handle: JobHandlerFn,
#[doc(hidden)]
pub retry_policy: &'static str,
#[doc(hidden)]
pub retry_delay_secs: f64,
}
inventory::collect!(JobDescriptor);
#[derive(Debug, Clone)]
pub struct Jobs {
pool: DatabaseConnection,
pub(crate) trace_id: Option<String>,
}
impl Jobs {
pub fn new(pool: DatabaseConnection, trace_id: Option<String>) -> Self {
Self { pool, trace_id }
}
pub async fn enqueue(&self, req: impl Into<JobRequest>) -> crate::error::Result<JobId> {
insert_job(&self.pool, req.into(), self.trace_id.as_deref()).await
}
pub async fn enqueue_with<C>(
&self,
conn: &C,
req: impl Into<JobRequest>,
) -> crate::error::Result<JobId>
where
C: ConnectionTrait,
{
insert_job(conn, req.into(), self.trace_id.as_deref()).await
}
}
async fn insert_job<C>(
conn: &C,
req: JobRequest,
trace_id: Option<&str>,
) -> crate::error::Result<JobId>
where
C: ConnectionTrait,
{
let id = Uuid::new_v4();
let stmt = match conn.get_database_backend() {
DbBackend::Postgres => backend::Postgres::build_insert_stmt(req, trace_id, id),
DbBackend::MySql => backend::Mysql::build_insert_stmt(req, trace_id, id),
DbBackend::Sqlite => backend::Sqlite::build_insert_stmt(req, trace_id, id),
};
conn.execute(stmt)
.await
.map_err(|e| crate::error::Error::internal(format!("failed to enqueue job: {e}")))?;
Ok(id)
}
#[cfg(test)]
mod tests {
use super::*;
fn sample_req() -> JobRequest {
JobRequest {
job_type: "send_email",
payload: serde_json::json!({ "to": "test@example.com" }),
queue: "emails",
max_retries: 5,
}
}
#[test]
fn job_request_fields() {
let req = sample_req();
assert_eq!(req.job_type, "send_email");
assert_eq!(req.queue, "emails");
assert_eq!(req.max_retries, 5);
assert_eq!(req.payload["to"], "test@example.com");
}
#[test]
fn default_convention() {
let req = JobRequest {
job_type: "process_event",
payload: serde_json::Value::Null,
queue: "default",
max_retries: 3,
};
assert_eq!(req.queue, "default");
assert_eq!(req.max_retries, 3);
}
#[test]
fn max_retries_is_i32() {
let req = JobRequest {
job_type: "t",
payload: serde_json::Value::Null,
queue: "default",
max_retries: i32::MAX,
};
assert_eq!(req.max_retries, i32::MAX);
}
#[test]
fn postgres_insert_stmt_uses_dollar_params() {
let req = sample_req();
let stmt = backend::Postgres::build_insert_stmt(req, None, Uuid::new_v4());
assert!(stmt.sql.contains("$1::uuid"), "Postgres uses $N params");
assert_eq!(stmt.db_backend, DbBackend::Postgres);
}
#[test]
fn postgres_insert_stmt_uses_string_uuid() {
let req = sample_req();
let stmt = backend::Postgres::build_insert_stmt(req, Some("trace-1"), Uuid::new_v4());
let vals = stmt.values.as_ref().unwrap();
assert!(matches!(&vals.0[0], sea_orm::Value::String(Some(_))));
}
#[test]
fn mysql_insert_stmt_has_current_timestamp() {
let req = sample_req();
let stmt = backend::Mysql::build_insert_stmt(req, None, Uuid::new_v4());
assert!(stmt.sql.contains("CURRENT_TIMESTAMP"));
assert_eq!(stmt.db_backend, DbBackend::MySql);
}
#[test]
fn mysql_insert_stmt_uses_uuid_value() {
let req = sample_req();
let stmt = backend::Mysql::build_insert_stmt(req, None, Uuid::new_v4());
let vals = stmt.values.as_ref().unwrap();
assert!(matches!(&vals.0[0], sea_orm::Value::Uuid(Some(_))));
}
#[test]
fn sqlite_insert_stmt_uses_datetime_now() {
let req = sample_req();
let stmt = backend::Sqlite::build_insert_stmt(req, None, Uuid::new_v4());
assert!(stmt.sql.contains("datetime('now')"));
assert_eq!(stmt.db_backend, DbBackend::Sqlite);
}
#[test]
fn sqlite_insert_stmt_uses_string_uuid() {
let req = sample_req();
let stmt = backend::Sqlite::build_insert_stmt(req, None, Uuid::new_v4());
let vals = stmt.values.as_ref().unwrap();
assert!(matches!(&vals.0[0], sea_orm::Value::String(Some(_))));
}
#[test]
fn mysql_retry_stmt_uses_microsecond() {
let stmt = backend::Mysql::build_retry_stmt(Uuid::new_v4(), "err", 1.7);
assert!(
stmt.sql.contains("MICROSECOND"),
"MySQL retry stmt should use MICROSECOND: {}",
stmt.sql
);
let vals = stmt.values.as_ref().unwrap();
assert!(matches!(&vals.0[1], sea_orm::Value::BigInt(Some(_))));
}
#[test]
fn mysql_build_retry_stmt_uses_uuid_value() {
let stmt = backend::Mysql::build_retry_stmt(Uuid::new_v4(), "err", 0.5);
let vals = stmt.values.as_ref().unwrap();
assert!(matches!(&vals.0[2], sea_orm::Value::Uuid(Some(_))));
}
#[test]
fn sqlite_build_retry_stmt_uses_string_uuid() {
let stmt = backend::Sqlite::build_retry_stmt(Uuid::new_v4(), "err", 0.5);
let vals = stmt.values.as_ref().unwrap();
assert!(matches!(&vals.0[2], sea_orm::Value::String(Some(_))));
}
#[test]
fn postgres_build_retry_stmt_uses_string_uuid() {
let stmt = backend::Postgres::build_retry_stmt(Uuid::new_v4(), "err", 0.5);
let vals = stmt.values.as_ref().unwrap();
assert!(matches!(&vals.0[2], sea_orm::Value::String(Some(_))));
}
}