use std::collections::HashMap;
use async_trait::async_trait;
use redis::AsyncCommands;
use serde::{Deserialize, Serialize};
use qrush::config::set_redis_url;
use qrush::cron::cron_job::CronJob;
use qrush::cron::cron_scheduler::CronScheduler;
use qrush::job::Job;
use qrush::utils::rdconfig::get_redis_connection;
#[derive(Clone, Serialize, Deserialize)]
struct CronTestJob {
message: String,
}
#[async_trait]
impl Job for CronTestJob {
async fn perform(&self) -> anyhow::Result<()> {
Ok(())
}
fn name(&self) -> &'static str {
"CronTestJob"
}
fn queue(&self) -> &'static str {
"cron_test_queue"
}
}
#[async_trait]
impl CronJob for CronTestJob {
fn cron_expression(&self) -> &'static str {
"0 0 * * * *"
}
fn cron_id(&self) -> &'static str {
"cron_enqueue_test_job"
}
}
#[tokio::test]
async fn run_now_enqueues_a_dispatchable_job() {
let redis_url =
std::env::var("REDIS_URL").unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string());
set_redis_url(redis_url).expect("failed to set Redis URL");
let mut conn = match get_redis_connection().await {
Ok(c) => c,
Err(e) => {
eprintln!("skipping cron_enqueue test: Redis unavailable ({e})");
return;
}
};
let cron_id = CronTestJob {
message: String::new(),
}
.cron_id();
CronScheduler::delete_cron_job(cron_id)
.await
.expect("cleanup delete failed");
let job = CronTestJob {
message: "hello from cron".into(),
};
CronScheduler::register_cron_job(job)
.await
.expect("register_cron_job failed");
let enqueued_id = CronScheduler::run_now(cron_id)
.await
.expect("run_now failed");
let job_key = format!("snm:job:{enqueued_id}");
let job_hash: HashMap<String, String> =
conn.hgetall(&job_key).await.expect("hgetall failed");
assert_eq!(
job_hash.get("job_name").map(String::as_str),
Some("CronTestJob"),
"enqueued cron job is missing the job_name the worker dispatches on"
);
assert_eq!(
job_hash.get("queue").map(String::as_str),
Some("cron_test_queue"),
"enqueued cron job landed on the wrong queue"
);
let payload = job_hash.get("payload").expect("missing payload");
let rebuilt: CronTestJob =
serde_json::from_str(payload).expect("payload is not a valid CronTestJob");
assert_eq!(rebuilt.message, "hello from cron");
CronScheduler::delete_cron_job(cron_id)
.await
.expect("cleanup delete failed");
let _: () = conn.del(&job_key).await.unwrap_or_default();
}