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