resque 0.3.0

rust resque clone
Documentation

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(())
}

// Result<scheduled_job_size>
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]>,
}