resque 0.3.0

rust resque clone
Documentation
use std::thread;
use std::time::Duration;
use std::ops::Deref;
use std::borrow::Cow;
use chrono;
use redis;
use serde_json;
use error_chain::ChainedError;
use wild_thread_pool::{NewWorker, Shutdown, Worker};
use {redis_keys, Error, ErrorKind, JobItem, JobServices, Result, WithRedis};

pub(crate) struct NewResqueWorker<R> {
    host: String,
    with_redis: R,
    job_services: JobServices,
    tick_sleep: Duration,
    fetch_timeout: u64,
}

impl<R> NewResqueWorker<R> {
    pub fn new(
        host: String,
        with_redis: R,
        job_services: JobServices,
        tick_sleep: Duration,
        fetch_timeout: u64,
    ) -> NewResqueWorker<R> {
        NewResqueWorker {
            host,
            with_redis,
            job_services,
            tick_sleep,
            fetch_timeout,
        }
    }
}

impl<R> NewWorker for NewResqueWorker<R>
where
    R: WithRedis,
{
    type Worker = ResqueWorker<R>;

    fn new_worker(&self, worker_id: u64) -> Self::Worker {
        ResqueWorker::new(
            &self.host,
            worker_id,
            self.with_redis.clone(),
            self.job_services.clone(),
            self.tick_sleep,
            self.fetch_timeout,
        )
    }
}

#[derive(Clone)]
pub(crate) struct ResqueWorker<R> {
    worker_id: u64,
    worker_name: String,
    with_redis: R,
    job_services: JobServices,
    tick_sleep: Duration,
    fetch_cmd: redis::Cmd,
}

impl<R> ResqueWorker<R>
where
    R: WithRedis,
{
    fn new(
        host: &str,
        worker_id: u64,
        with_redis: R,
        job_services: JobServices,
        tick_sleep: Duration,
        fetch_timeout: u64,
    ) -> ResqueWorker<R> {
        let worker_name = format!("{}:{}:*", host, worker_id);

        let mut fetch_cmd = redis::cmd("BLPOP");
        for job_service_pair in job_services.values() {
            fetch_cmd.arg(redis_keys::queue(job_service_pair.queue()));
        }
        fetch_cmd.arg(fetch_timeout);

        ResqueWorker {
            worker_id,
            worker_name,
            with_redis,
            job_services,
            tick_sleep,
            fetch_cmd,
        }
    }

    fn reg_worker(with_redis: &R, worker_name: &str) -> Result<()> {
        with_redis.with_redis(|redis_conn| {
            let now = chrono::Utc::now().to_rfc3339();

            let _: (i64, String) = redis::pipe()
                .cmd("SADD")
                .arg(redis_keys::workers())
                .arg(worker_name)
                .cmd("SET")
                .arg(redis_keys::worker_start_time_key(worker_name))
                .arg(now)
                .query(redis_conn.deref())?;
            Ok(())
        })
    }

    fn unreg_worker(with_redis: &R, worker_name: &str) -> Result<()> {
        with_redis.with_redis(|redis_conn| {
            let _: (i64, i64) = redis::pipe()
                .cmd("SREM")
                .arg(redis_keys::workers())
                .arg(worker_name)
                .cmd("DEL")
                .arg(redis_keys::worker_start_time_key(worker_name))
                .query(redis_conn.deref())?;
            Ok(())
        })
    }

    // Result<Option<(Queue, JobItem)>>
    fn fetch_one(&self) -> Result<Option<JobItem>> {
        let job_item: Option<(String, String)> = self.with_redis.with_redis(|redis_conn| {
            self.fetch_cmd
                .query(redis_conn.deref())
                .map_err(|err| err.into())
        })?;

        let job_item = match job_item {
            Some(v) => v.1,
            None => return Ok(None),
        };

        let job_item: JobItem = match serde_json::from_str(&job_item) {
            Ok(job_item) => job_item,
            Err(err) => {
                error!(
                    "worker-{} fetch job: {:?}, err: \n{:?}",
                    self.worker_id,
                    &job_item,
                    err
                );
                return Err(err.into());
            }
        };

        Ok(Some(job_item))
    }

    fn working_on(&self, queue: &str, job_item: &JobItem) -> Result<()> {
        let working_on = WorkingOn {
            queue: Cow::from(queue),
            run_at: chrono::Utc::now(),
            payload: Cow::Borrowed(job_item),
        };
        let working_on = serde_json::to_string(&working_on)?;

        self.with_redis.with_redis(|redis_conn| {
            let _: bool = redis::cmd("SET")
                .arg(redis_keys::working_on(&self.worker_name))
                .arg(working_on)
                .query(redis_conn.deref())?;
            Ok(())
        })
    }

    fn report_failed(with_redis: &R, worker_id: u64, worker_name: &str, err: &Error) -> Result<()> {
        error!("worker-{} report_failed: {:?}", worker_id, err);

        let err_msg = format!("{:?}", err);

        with_redis.with_redis(|redis_conn| {
            let working_on: Option<String> = redis::cmd("GET")
                .arg(redis_keys::working_on(worker_name))
                .query(redis_conn.deref())?;

            let working_on = match working_on {
                Some(working_on) => working_on,
                None => return Ok(()),
            };
            let working_on: WorkingOn = serde_json::from_str(&working_on)?;

            let failed_job = FailedJob {
                failed_at: chrono::Utc::now(),
                payload: Cow::from(working_on.payload),
                exception: Cow::Borrowed(&err_msg),
                error: Some(Cow::Borrowed(&err_msg)),
                backtrace: err.iter().map(|e| format!("{:?}", e)).collect(),
                worker: Cow::from(worker_name),
                queue: Cow::from(working_on.queue),
                retried_at: None,
            };
            let failed_job = serde_json::to_string(&failed_job)?;

            let _: i64 = redis::cmd("RPUSH")
                .arg(redis_keys::failed())
                .arg(failed_job)
                .query(redis_conn.deref())?;

            Ok(())
        })
    }

    fn done_working(&self) -> Result<()> {
        self.with_redis.with_redis(|redis_conn| {
            let _: (i64, i64) = redis::pipe()
                .cmd("DEL")
                .arg(redis_keys::working_on(&self.worker_name))
                .cmd("INCR")
                .arg(redis_keys::processed())
                .query(redis_conn.deref())?;

            Ok(())
        })
    }
}

impl<R> Worker for ResqueWorker<R>
where
    R: WithRedis,
{
    type Shutdown = ResqueShutdown<R>;
    type Err = Error;

    #[inline]
    fn worker_id(&self) -> u64 {
        self.worker_id
    }

    fn on_start(&self) -> Result<()> {
        ResqueWorker::reg_worker(&self.with_redis, &self.worker_name)?;

        info!("worker-{} started", self.worker_id());
        Ok(())
    }

    fn on_stop(&self) -> Result<()> {
        ResqueWorker::unreg_worker(&self.with_redis, &self.worker_name)?;

        info!("worker-{} stoped", self.worker_id());
        Ok(())
    }

    fn on_tick_err(&self, err: &Error) {
        error!(
            "worker-{} on_tick_err: {:?}",
            self.worker_id(),
            err.display_chain()
        );

        if let Err(err) =
            ResqueWorker::report_failed(&self.with_redis, self.worker_id, &self.worker_name, err)
        {
            error!("worker-{} report_failed err: {:?}", self.worker_id, err);
        }
    }

    fn on_tick(&self) -> Result<()> {
        let job_item = match self.fetch_one()? {
            Some(job_item) => job_item,
            None => {
                debug!("worker-{} fetched nothing", self.worker_id);

                thread::sleep(self.tick_sleep);
                return Ok(());
            }
        };

        debug!("worker-{} fetched one job: {:?}", self.worker_id, &job_item);

        let job_service = match self.job_services.get(job_item.class.as_ref()) {
            Some(job_service) => job_service,
            None => return Err(ErrorKind::JobServiceNotFound(job_item.class.into_owned()).into()),
        };

        self.working_on(job_service.queue(), &job_item)?;

        job_service.run(&job_item.args)?;

        self.done_working()?;

        debug!(
            "worker-{} done working job: {:?}",
            self.worker_id,
            &job_item
        );

        Ok(())
    }

    fn new_shutdown(&self) -> Self::Shutdown {
        ResqueShutdown {
            worker_id: self.worker_id,
            worker_name: self.worker_name.clone(),
            with_redis: self.with_redis.clone(),
        }
    }
}

pub(crate) struct ResqueShutdown<R> {
    worker_id: u64,
    worker_name: String,
    with_redis: R,
}

impl<R> Shutdown for ResqueShutdown<R>
where
    R: WithRedis,
{
    #[inline]
    fn worker_id(&self) -> u64 {
        self.worker_id
    }

    fn on_stop_timeout(&self) {
        error!("worker-{} stop_timeout", self.worker_id());

        let err: Error = "DirtyExit".into();
        if let Err(err) =
            ResqueWorker::report_failed(&self.with_redis, self.worker_id, &self.worker_name, &err)
        {
            error!("worker-{} report_failed err: {:?}", self.worker_id, err);
        }

        if let Err(err) = ResqueWorker::unreg_worker(&self.with_redis, &self.worker_name) {
            error!("worker-{} stop_timeout err: {:?}", self.worker_id, err);
        }
    }
}

#[derive(Debug, Serialize, Deserialize)]
struct WorkingOn<'a> {
    queue: Cow<'a, str>,
    run_at: chrono::DateTime<chrono::Utc>,
    payload: Cow<'a, JobItem<'a>>,
}

#[derive(Debug, Serialize, Deserialize)]
struct FailedJob<'a> {
    failed_at: chrono::DateTime<chrono::Utc>,
    payload: Cow<'a, JobItem<'a>>,
    exception: Cow<'a, str>,
    error: Option<Cow<'a, str>>,
    backtrace: Vec<String>,
    worker: Cow<'a, str>,
    queue: Cow<'a, str>,
    retried_at: Option<chrono::DateTime<chrono::Utc>>,
}