use std::collections::HashMap;
use std::sync::Arc;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use async_trait::async_trait;
use sz_orm_core::{DbError, Pool, PoolError, Value};
use thiserror::Error;
pub const JOBS_TABLE: &str = "sz_jobs";
const SCHEMA_SQL: &str = "CREATE TABLE IF NOT EXISTS sz_jobs (
id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
kind VARCHAR(64) NOT NULL,
payload TEXT NOT NULL,
status VARCHAR(16) NOT NULL DEFAULT 'pending',
attempts INT NOT NULL DEFAULT 0,
run_after BIGINT NOT NULL,
locked_until BIGINT NULL,
last_error TEXT NULL,
dedupe_key VARCHAR(255) NULL,
created_at BIGINT NOT NULL,
updated_at BIGINT NOT NULL,
UNIQUE KEY uq_sz_jobs_dedupe (kind, dedupe_key),
KEY idx_sz_jobs_status_run_after (status, run_after)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4";
const STATUS_PENDING: &str = "pending";
const STATUS_RUNNING: &str = "running";
const STATUS_SUCCEEDED: &str = "succeeded";
const STATUS_DEAD: &str = "dead";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum JobStatus {
Pending,
Running,
Succeeded,
Dead,
}
impl JobStatus {
pub fn as_str(&self) -> &'static str {
match self {
JobStatus::Pending => STATUS_PENDING,
JobStatus::Running => STATUS_RUNNING,
JobStatus::Succeeded => STATUS_SUCCEEDED,
JobStatus::Dead => STATUS_DEAD,
}
}
pub fn parse_status(s: &str) -> Option<JobStatus> {
match s {
STATUS_PENDING => Some(JobStatus::Pending),
STATUS_RUNNING => Some(JobStatus::Running),
STATUS_SUCCEEDED => Some(JobStatus::Succeeded),
STATUS_DEAD => Some(JobStatus::Dead),
_ => None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum JobErrorKind {
Temporary,
Permanent,
}
#[derive(Debug, Error)]
pub enum JobError {
#[error("temporary job failure: {0}")]
Temporary(String),
#[error("permanent job failure: {0}")]
Permanent(String),
}
impl JobError {
pub fn kind(&self) -> JobErrorKind {
match self {
JobError::Temporary(_) => JobErrorKind::Temporary,
JobError::Permanent(_) => JobErrorKind::Permanent,
}
}
}
#[async_trait]
pub trait TaskHandler: Send + Sync + 'static {
async fn handle(&self, payload: &serde_json::Value) -> Result<(), JobError>;
}
#[derive(Debug, Error)]
pub enum JobQueueError {
#[error("database error: {0}")]
Db(#[from] DbError),
#[error("pool error: {0}")]
Pool(#[from] PoolError),
#[error("invalid job row: {0}")]
InvalidRow(String),
#[error("json error: {0}")]
Json(#[from] serde_json::Error),
}
#[derive(Debug, Clone)]
pub struct JobQueueConfig {
pub batch_size: u32,
pub poll_interval: Duration,
pub max_attempts: u32,
pub backoff_base_secs: u64,
pub backoff_cap_secs: u64,
pub lease_seconds: u64,
pub jitter_ratio: f64,
pub handler_timeout: Duration,
}
impl Default for JobQueueConfig {
fn default() -> Self {
Self {
batch_size: 10,
poll_interval: Duration::from_secs(1),
max_attempts: 8,
backoff_base_secs: 1,
backoff_cap_secs: 64,
lease_seconds: 60,
jitter_ratio: 0.3,
handler_timeout: Duration::from_secs(30),
}
}
}
#[derive(Debug, Clone)]
pub struct Job {
pub id: u64,
pub kind: String,
pub payload: serde_json::Value,
pub status: JobStatus,
pub attempts: u32,
pub run_after: i64,
pub last_error: Option<String>,
pub dedupe_key: Option<String>,
pub created_at: i64,
}
#[derive(Debug, Clone, Copy, Default)]
pub struct QueueSnapshot {
pub pending: u64,
pub running: u64,
pub dead: u64,
pub succeeded: u64,
pub oldest_pending_seconds: u64,
}
#[derive(Clone)]
pub struct JobQueue {
pool: Arc<Pool>,
}
impl JobQueue {
pub fn new(pool: Arc<Pool>) -> Self {
Self { pool }
}
pub fn pool(&self) -> &Arc<Pool> {
&self.pool
}
pub async fn init_schema(&self) -> Result<(), JobQueueError> {
let mut conn = self.pool.acquire().await?;
conn.execute(SCHEMA_SQL).await?;
Ok(())
}
pub async fn enqueue(
&self,
kind: &str,
payload: serde_json::Value,
dedupe_key: Option<&str>,
) -> Result<u64, JobQueueError> {
self.enqueue_at(kind, payload, dedupe_key, now_ms()).await
}
pub async fn enqueue_delayed(
&self,
kind: &str,
payload: serde_json::Value,
dedupe_key: Option<&str>,
delay: Duration,
) -> Result<u64, JobQueueError> {
self.enqueue_at(
kind,
payload,
dedupe_key,
now_ms() + delay.as_millis() as i64,
)
.await
}
async fn enqueue_at(
&self,
kind: &str,
payload: serde_json::Value,
dedupe_key: Option<&str>,
run_after: i64,
) -> Result<u64, JobQueueError> {
let payload_str = serde_json::to_string(&payload)?;
let now = now_ms();
let mut conn = self.pool.acquire().await?;
conn.execute_with_params(
"INSERT INTO sz_jobs (kind, payload, status, attempts, run_after, dedupe_key, created_at, updated_at) \
VALUES (?, ?, ?, 0, ?, ?, ?, ?) \
ON DUPLICATE KEY UPDATE id = LAST_INSERT_ID(id)",
&[
Value::String(kind.into()),
Value::String(payload_str),
Value::String(STATUS_PENDING.into()),
Value::I64(run_after),
dedupe_key.map_or(Value::Null, |k| Value::String(k.into())),
Value::I64(now),
Value::I64(now),
],
)
.await?;
let rows = conn
.query_with_params("SELECT LAST_INSERT_ID() AS id", &[])
.await?;
rows.first()
.and_then(|r| r.get("id"))
.and_then(Value::as_i64)
.map(|v| v as u64)
.ok_or_else(|| JobQueueError::InvalidRow("LAST_INSERT_ID() 返回空".into()))
}
pub async fn retry_dead(&self, job_id: u64) -> Result<(), JobQueueError> {
let now = now_ms();
let mut conn = self.pool.acquire().await?;
conn.execute_with_params(
"UPDATE sz_jobs SET status = ?, run_after = ?, locked_until = NULL, updated_at = ? WHERE id = ? AND status = ?",
&[
Value::String(STATUS_PENDING.into()),
Value::I64(now),
Value::I64(now),
Value::I64(job_id as i64),
Value::String(STATUS_DEAD.into()),
],
)
.await?;
Ok(())
}
pub async fn queue_snapshot(&self) -> Result<QueueSnapshot, JobQueueError> {
let mut conn = self.pool.acquire().await?;
let rows = conn
.query("SELECT status, COUNT(*) AS cnt FROM sz_jobs GROUP BY status")
.await?;
let mut snap = QueueSnapshot::default();
for row in rows {
let status = row.get("status").and_then(Value::as_str).unwrap_or("");
let cnt = row.get("cnt").and_then(Value::as_i64).unwrap_or(0).max(0) as u64;
match status {
STATUS_PENDING => snap.pending = cnt,
STATUS_RUNNING => snap.running = cnt,
STATUS_SUCCEEDED => snap.succeeded = cnt,
STATUS_DEAD => snap.dead = cnt,
_ => {}
}
}
let rows = conn
.query("SELECT MIN(run_after) AS oldest FROM sz_jobs WHERE status = 'pending'")
.await?;
if let Some(oldest) = rows
.first()
.and_then(|r| r.get("oldest"))
.and_then(Value::as_i64)
{
snap.oldest_pending_seconds = ((now_ms() - oldest).max(0) / 1000) as u64;
}
Ok(snap)
}
pub async fn run_worker(
&self,
handlers: HashMap<String, Arc<dyn TaskHandler>>,
config: JobQueueConfig,
shutdown: tokio::sync::watch::Receiver<bool>,
) -> Result<(), JobQueueError> {
let mut interval = tokio::time::interval(config.poll_interval);
loop {
interval.tick().await;
if *shutdown.borrow() {
tracing::info!(target: "sz_orm::jobs", "job worker shutting down");
return Ok(());
}
if let Err(e) = self.reclaim_stale(&config).await {
tracing::error!(target: "sz_orm::jobs", "reclaim stale jobs failed: {e}");
continue;
}
let jobs = match self.claim_batch(&config).await {
Ok(jobs) => jobs,
Err(e) => {
tracing::error!(target: "sz_orm::jobs", "claim jobs failed: {e}");
continue;
}
};
if jobs.is_empty() {
if let Ok(snap) = self.queue_snapshot().await {
tracing::debug!(
target: "sz_orm::jobs",
"queue snapshot: pending={}, running={}, dead={}, oldest_pending_secs={}",
snap.pending, snap.running, snap.dead, snap.oldest_pending_seconds
);
}
continue;
}
for job in jobs {
let handler = handlers.get(&job.kind);
let outcome = match handler {
Some(h) => {
match tokio::time::timeout(config.handler_timeout, h.handle(&job.payload))
.await
{
Ok(Ok(())) => Ok(()),
Ok(Err(e)) => Err(e),
Err(_) => Err(JobError::Temporary(format!(
"handler timeout after {:?}",
config.handler_timeout
))),
}
}
None => Err(JobError::Permanent(format!(
"no handler registered for kind '{}'",
job.kind
))),
};
match outcome {
Ok(()) => {
self.mark_succeeded(job.id).await?;
tracing::debug!(target: "sz_orm::jobs", "job {} (kind={}) succeeded", job.id, job.kind);
}
Err(e) => {
self.handle_failure(
job.id,
job.attempts,
&e.to_string(),
e.kind(),
&config,
)
.await?;
tracing::warn!(
target: "sz_orm::jobs",
"job {} (kind={}) failed: {} (kind={:?}), attempts={}",
job.id, job.kind, e, e.kind(), job.attempts
);
}
}
}
}
}
async fn reclaim_stale(&self, config: &JobQueueConfig) -> Result<(), JobQueueError> {
let now = now_ms();
let lease_deadline = now - config.lease_seconds as i64 * 1000;
let mut conn = self.pool.acquire().await?;
conn.execute_with_params(
"UPDATE sz_jobs SET status = ?, locked_until = NULL, updated_at = ? \
WHERE status = ? AND locked_until < ?",
&[
Value::String(STATUS_PENDING.into()),
Value::I64(now),
Value::String(STATUS_RUNNING.into()),
Value::I64(lease_deadline),
],
)
.await?;
Ok(())
}
async fn claim_batch(&self, config: &JobQueueConfig) -> Result<Vec<Job>, JobQueueError> {
let now = now_ms();
let locked_until = now + config.lease_seconds as i64 * 1000;
let mut conn = self.pool.acquire().await?;
conn.begin_transaction().await?;
let rows = conn
.query_with_params(
"SELECT id FROM sz_jobs WHERE status = ? AND run_after <= ? \
ORDER BY created_at LIMIT ? FOR UPDATE SKIP LOCKED",
&[
Value::String(STATUS_PENDING.into()),
Value::I64(now),
Value::I64(config.batch_size as i64),
],
)
.await?;
let ids: Vec<Value> = rows
.iter()
.filter_map(|r| r.get("id").and_then(Value::as_i64).map(Value::I64))
.collect();
if !ids.is_empty() {
let placeholders = vec!["?"; ids.len()].join(",");
let mut params = Vec::with_capacity(ids.len() + 3);
params.push(Value::String(STATUS_RUNNING.into()));
params.push(Value::I64(locked_until));
params.push(Value::I64(now));
params.extend(ids);
conn.execute_with_params(
&format!(
"UPDATE sz_jobs SET status = ?, locked_until = ?, attempts = attempts + 1, updated_at = ? \
WHERE id IN ({placeholders})"
),
¶ms,
)
.await?;
}
conn.commit().await?;
let rows = conn
.query_with_params(
"SELECT id, kind, payload, status, attempts, run_after, last_error, dedupe_key, created_at \
FROM sz_jobs WHERE status = ? AND locked_until = ? ORDER BY created_at",
&[Value::String(STATUS_RUNNING.into()), Value::I64(locked_until)],
)
.await?;
rows.into_iter().map(row_to_job).collect()
}
async fn mark_succeeded(&self, job_id: u64) -> Result<(), JobQueueError> {
let mut conn = self.pool.acquire().await?;
conn.execute_with_params(
"UPDATE sz_jobs SET status = ?, locked_until = NULL, updated_at = ? WHERE id = ?",
&[
Value::String(STATUS_SUCCEEDED.into()),
Value::I64(now_ms()),
Value::I64(job_id as i64),
],
)
.await?;
Ok(())
}
async fn handle_failure(
&self,
job_id: u64,
attempts: u32,
error: &str,
kind: JobErrorKind,
config: &JobQueueConfig,
) -> Result<(), JobQueueError> {
let now = now_ms();
let (status, run_after) = match kind {
JobErrorKind::Temporary if attempts <= config.max_attempts => (
STATUS_PENDING,
now + backoff_delay_ms(config, attempts) as i64,
),
_ => (STATUS_DEAD, now),
};
let mut conn = self.pool.acquire().await?;
conn.execute_with_params(
"UPDATE sz_jobs SET status = ?, run_after = ?, locked_until = NULL, last_error = ?, updated_at = ? WHERE id = ?",
&[
Value::String(status.into()),
Value::I64(run_after),
Value::String(error.into()),
Value::I64(now),
Value::I64(job_id as i64),
],
)
.await?;
Ok(())
}
}
pub fn backoff_delay_ms(config: &JobQueueConfig, attempts: u32) -> u64 {
let exp = (attempts as i32).min(6);
let delay_secs =
(config.backoff_base_secs as f64 * 2f64.powi(exp)).min(config.backoff_cap_secs as f64);
let jitter = delay_secs * config.jitter_ratio.clamp(0.0, 1.0) * rand::random::<f64>();
((delay_secs + jitter) * 1000.0) as u64
}
pub fn now_ms() -> i64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_millis() as i64)
.unwrap_or(0)
}
fn row_to_job(row: HashMap<String, Value>) -> Result<Job, JobQueueError> {
let id = row
.get("id")
.and_then(Value::as_i64)
.ok_or_else(|| JobQueueError::InvalidRow("id".into()))?;
let kind = row
.get("kind")
.and_then(Value::as_str)
.ok_or_else(|| JobQueueError::InvalidRow("kind".into()))?
.to_string();
let payload_str = row
.get("payload")
.and_then(Value::as_str)
.ok_or_else(|| JobQueueError::InvalidRow("payload".into()))?;
let payload = serde_json::from_str(payload_str)?;
let status = row
.get("status")
.and_then(Value::as_str)
.and_then(JobStatus::parse_status)
.ok_or_else(|| JobQueueError::InvalidRow("status".into()))?;
let attempts = row
.get("attempts")
.and_then(Value::as_i64)
.unwrap_or(0)
.max(0) as u32;
let run_after = row.get("run_after").and_then(Value::as_i64).unwrap_or(0);
let last_error = row
.get("last_error")
.and_then(Value::as_str)
.map(str::to_string);
let dedupe_key = row
.get("dedupe_key")
.and_then(Value::as_str)
.map(str::to_string);
let created_at = row.get("created_at").and_then(Value::as_i64).unwrap_or(0);
Ok(Job {
id: id as u64,
kind,
payload,
status,
attempts,
run_after,
last_error,
dedupe_key,
created_at,
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_job_status_roundtrip() {
for s in [
JobStatus::Pending,
JobStatus::Running,
JobStatus::Succeeded,
JobStatus::Dead,
] {
assert_eq!(JobStatus::parse_status(s.as_str()), Some(s));
}
assert_eq!(JobStatus::parse_status("unknown"), None);
}
#[test]
fn test_backoff_delay_respects_cap() {
let config = JobQueueConfig {
max_attempts: 10,
backoff_base_secs: 1,
backoff_cap_secs: 64,
jitter_ratio: 0.0,
..JobQueueConfig::default()
};
let delay = backoff_delay_ms(&config, 10);
assert!((delay as f64 - 64_000.0).abs() < 1.0, "delay={delay}");
}
#[test]
fn test_backoff_delay_jitter_range() {
let config = JobQueueConfig {
max_attempts: 8,
backoff_base_secs: 1,
backoff_cap_secs: 64,
jitter_ratio: 0.5,
..JobQueueConfig::default()
};
for _ in 0..50 {
let delay = backoff_delay_ms(&config, 2);
assert!(delay >= 4_000, "delay={delay}");
assert!(delay <= 6_000, "delay={delay}");
}
}
#[test]
fn test_backoff_delay_escalation() {
let config = JobQueueConfig {
max_attempts: 8,
backoff_base_secs: 1,
backoff_cap_secs: 64,
jitter_ratio: 0.0,
..JobQueueConfig::default()
};
let delay = backoff_delay_ms(&config, 3);
assert!((delay as f64 - 8_000.0).abs() < 1.0, "delay={delay}");
}
#[test]
fn test_now_ms_monotonic() {
let a = now_ms();
std::thread::sleep(Duration::from_millis(5));
let b = now_ms();
assert!(b > a);
}
#[test]
fn test_job_error_kind() {
let t = JobError::Temporary("downstream 503".into());
assert_eq!(t.kind(), JobErrorKind::Temporary);
let p = JobError::Permanent("user not found".into());
assert_eq!(p.kind(), JobErrorKind::Permanent);
}
fn make_row() -> HashMap<String, Value> {
let mut row = HashMap::new();
row.insert("id".into(), Value::I64(42));
row.insert("kind".into(), Value::String("email".into()));
row.insert("payload".into(), Value::String(r#"{"to":"a@b"}"#.into()));
row.insert("status".into(), Value::String("pending".into()));
row.insert("attempts".into(), Value::I64(3));
row.insert("run_after".into(), Value::I64(1000));
row.insert("last_error".into(), Value::String("timeout".into()));
row.insert("dedupe_key".into(), Value::String("k1".into()));
row.insert("created_at".into(), Value::I64(500));
row
}
#[test]
fn test_row_to_job_success() {
let job = row_to_job(make_row()).unwrap();
assert_eq!(job.id, 42);
assert_eq!(job.kind, "email");
assert_eq!(job.status, JobStatus::Pending);
assert_eq!(job.attempts, 3);
assert_eq!(job.run_after, 1000);
assert_eq!(job.last_error.as_deref(), Some("timeout"));
assert_eq!(job.dedupe_key.as_deref(), Some("k1"));
assert_eq!(job.created_at, 500);
}
#[test]
fn test_row_to_job_missing_id() {
let mut row = make_row();
row.remove("id");
let err = row_to_job(row).unwrap_err();
assert!(matches!(err, JobQueueError::InvalidRow(_)));
}
#[test]
fn test_row_to_job_missing_kind() {
let mut row = make_row();
row.remove("kind");
let err = row_to_job(row).unwrap_err();
assert!(matches!(err, JobQueueError::InvalidRow(_)));
}
#[test]
fn test_row_to_job_missing_payload() {
let mut row = make_row();
row.remove("payload");
let err = row_to_job(row).unwrap_err();
assert!(matches!(err, JobQueueError::InvalidRow(_)));
}
#[test]
fn test_row_to_job_invalid_status() {
let mut row = make_row();
row.insert("status".into(), Value::String("unknown".into()));
let err = row_to_job(row).unwrap_err();
assert!(matches!(err, JobQueueError::InvalidRow(_)));
}
#[test]
fn test_row_to_job_invalid_payload_json() {
let mut row = make_row();
row.insert("payload".into(), Value::String("{bad json".into()));
let err = row_to_job(row).unwrap_err();
assert!(matches!(err, JobQueueError::Json(_)));
}
#[test]
fn test_row_to_job_optional_fields_default() {
let mut row = make_row();
row.remove("last_error");
row.remove("dedupe_key");
row.remove("attempts");
row.remove("run_after");
row.remove("created_at");
let job = row_to_job(row).unwrap();
assert_eq!(job.attempts, 0);
assert_eq!(job.run_after, 0);
assert!(job.last_error.is_none());
assert!(job.dedupe_key.is_none());
assert_eq!(job.created_at, 0);
}
#[test]
fn test_job_queue_config_default() {
let config = JobQueueConfig::default();
assert_eq!(config.batch_size, 10);
assert_eq!(config.max_attempts, 8);
assert_eq!(config.backoff_base_secs, 1);
assert_eq!(config.backoff_cap_secs, 64);
assert_eq!(config.lease_seconds, 60);
}
use std::future::Future;
use std::pin::Pin;
use sz_orm_core::{Connection, ConnectionFactory, PoolConfig, QueryRows};
struct MockConnection;
impl Connection for MockConnection {
fn execute<'a>(
&'a mut self,
_sql: &'a str,
) -> Pin<Box<dyn Future<Output = Result<u64, DbError>> + Send + 'a>> {
Box::pin(async { Ok(0) })
}
fn query<'a>(
&'a mut self,
_sql: &'a str,
) -> Pin<Box<dyn Future<Output = Result<QueryRows, DbError>> + Send + 'a>> {
Box::pin(async { Ok(vec![]) })
}
fn execute_with_params<'a>(
&'a mut self,
_sql: &'a str,
_params: &'a [Value],
) -> Pin<Box<dyn Future<Output = Result<u64, DbError>> + Send + 'a>> {
Box::pin(async { Ok(1) })
}
fn query_with_params<'a>(
&'a mut self,
sql: &'a str,
_params: &'a [Value],
) -> Pin<Box<dyn Future<Output = Result<QueryRows, DbError>> + Send + 'a>> {
Box::pin(async move {
if sql.contains("LAST_INSERT_ID") {
let mut row = HashMap::new();
row.insert("id".into(), Value::I64(1));
Ok(vec![row])
} else {
Ok(vec![])
}
})
}
fn begin_transaction<'a>(
&'a mut self,
) -> Pin<Box<dyn Future<Output = Result<(), DbError>> + Send + 'a>> {
Box::pin(async { Ok(()) })
}
fn commit<'a>(
&'a mut self,
) -> Pin<Box<dyn Future<Output = Result<(), DbError>> + Send + 'a>> {
Box::pin(async { Ok(()) })
}
fn rollback<'a>(
&'a mut self,
) -> Pin<Box<dyn Future<Output = Result<(), DbError>> + Send + 'a>> {
Box::pin(async { Ok(()) })
}
fn is_connected(&self) -> bool {
true
}
fn ping<'a>(&'a mut self) -> Pin<Box<dyn Future<Output = bool> + Send + 'a>> {
Box::pin(async { true })
}
fn close<'a>(
&'a mut self,
) -> Pin<Box<dyn Future<Output = Result<(), DbError>> + Send + 'a>> {
Box::pin(async { Ok(()) })
}
}
struct MockConnectionFactory;
#[async_trait]
impl ConnectionFactory for MockConnectionFactory {
async fn create(&self) -> Result<Box<dyn Connection>, DbError> {
Ok(Box::new(MockConnection))
}
}
fn make_mock_pool() -> Arc<Pool> {
let config = PoolConfig::default();
let factory: Arc<dyn ConnectionFactory> = Arc::new(MockConnectionFactory);
Arc::new(Pool::new(config, factory).expect("mock pool creation should not fail"))
}
#[test]
fn test_job_queue_new_and_pool() {
let pool = make_mock_pool();
let queue = JobQueue::new(pool.clone());
assert!(Arc::ptr_eq(queue.pool(), &pool));
}
#[tokio::test]
async fn test_job_queue_init_schema() {
let queue = JobQueue::new(make_mock_pool());
let result = queue.init_schema().await;
assert!(result.is_ok());
}
#[tokio::test]
async fn test_job_queue_enqueue() {
let queue = JobQueue::new(make_mock_pool());
let id = queue
.enqueue("email", serde_json::json!({"to": "a@b"}), None)
.await
.unwrap();
assert_eq!(id, 1);
}
#[tokio::test]
async fn test_job_queue_enqueue_with_dedupe() {
let queue = JobQueue::new(make_mock_pool());
let id = queue
.enqueue("email", serde_json::json!({}), Some("k1"))
.await
.unwrap();
assert_eq!(id, 1);
}
#[tokio::test]
async fn test_job_queue_enqueue_delayed() {
let queue = JobQueue::new(make_mock_pool());
let id = queue
.enqueue_delayed(
"email",
serde_json::json!({}),
None,
Duration::from_secs(60),
)
.await
.unwrap();
assert_eq!(id, 1);
}
#[tokio::test]
async fn test_job_queue_retry_dead() {
let queue = JobQueue::new(make_mock_pool());
let result = queue.retry_dead(42).await;
assert!(result.is_ok());
}
#[tokio::test]
async fn test_job_queue_snapshot() {
let queue = JobQueue::new(make_mock_pool());
let snap = queue.queue_snapshot().await.unwrap();
assert_eq!(snap.pending, 0);
assert_eq!(snap.running, 0);
}
}