use std::borrow::Cow;
use std::ops::Deref;
use std::sync::atomic::Ordering;
use std::thread;
use std::time::Duration;
use redis;
use chrono;
use serde_json;
use config::ScheduleConfig;
use {enqueue, redis_keys, Result, WithRedis};
pub fn enqueue_at<C>(
redis_conn: &C,
run_at: i64,
queue: &str,
class: &str,
args: &[String],
) -> Result<()>
where
C: Deref<Target = redis::Connection>,
{
let item = ScheduleJobItem {
queue: Cow::from(queue),
class: Cow::from(class),
args: Cow::from(args),
};
let job_json = serde_json::to_string(&item)?;
let _: i64 = redis::cmd("zadd")
.arg(redis_keys::schedules())
.arg(run_at - 1)
.arg(job_json)
.query(redis_conn.deref())?;
Ok(())
}
pub fn start_schedule<R>(c: ScheduleConfig<R>) -> Result<()>
where
R: WithRedis,
{
let ScheduleConfig {
with_redis,
tick_sleep,
closed,
} = c;
assert_ne!(tick_sleep, 0);
let tick_sleep = Duration::from_secs(tick_sleep);
info!("resque schedule running...");
loop {
if closed.load(Ordering::Relaxed) {
break;
}
let any_job = match with_redis.with_redis(|redis_conn| {
let now_i = chrono::Utc::now().timestamp();
try_schedule(redis_conn, now_i + 1)
}) {
Ok(job_size) => {
if job_size > 0 {
debug!("schedule {} jobs", job_size);
}
job_size > 0
}
Err(err) => {
error!("schedule jobs err: {:?}", err);
false
}
};
if !any_job {
thread::sleep(tick_sleep);
}
}
Ok(())
}
fn try_schedule<C>(redis_conn: &C, now_i: i64) -> Result<i64>
where
C: Deref<Target = redis::Connection>,
{
let (jobs, c): (Vec<String>, i64) = redis::pipe()
.atomic()
.cmd("zrangebyscore")
.arg(redis_keys::schedules())
.arg("-inf")
.arg(now_i)
.cmd("zremrangebyscore")
.arg(redis_keys::schedules())
.arg("-inf")
.arg(now_i)
.query(redis_conn.deref())?;
if c <= 0 {
return Ok(0);
}
debug!("scheduled {} jobs", c);
for job_s in jobs {
let job: ScheduleJobItem = serde_json::from_str(&job_s)?;
enqueue(
redis_conn,
job.queue.as_ref(),
job.class.as_ref(),
job.args.as_ref(),
)?;
}
Ok(c)
}
#[derive(Debug, Serialize, Deserialize)]
struct ScheduleJobItem<'a> {
queue: Cow<'a, str>,
class: Cow<'a, str>,
args: Cow<'a, [String]>,
}