Skip to main content

resque/
lib.rs

1extern crate chrono;
2#[macro_use]
3extern crate error_chain;
4#[macro_use]
5extern crate log;
6extern crate redis;
7extern crate serde;
8#[macro_use]
9extern crate serde_derive;
10extern crate serde_json;
11extern crate wild_thread_pool;
12
13mod config;
14mod redis_keys;
15mod schedule;
16mod worker;
17
18pub use config::{Config, ScheduleConfig};
19pub use schedule::{enqueue_at, start_schedule};
20
21use std::collections::HashMap;
22use std::borrow::Cow;
23use std::time::Duration;
24use std::ops::Deref;
25use std::clone::Clone;
26use wild_thread_pool::{ErrorExt, ThreadPool};
27
28error_chain!{
29    foreign_links {
30        RedisError(::redis::RedisError);
31        SerdeJsonError(::serde_json::Error);
32    }
33
34    errors {
35        RegDupJobService(job_type: String){
36            description("reg dup job service")
37            display("reg dup job_service({})", job_type)
38        }
39
40        JobServiceNotFound(job_type: String){
41            description("job service not found")
42            display("job_service({}) not found", job_type)
43        }
44    }
45}
46
47pub trait WithRedis: Send + Clone + 'static {
48    type Conn: Deref<Target = redis::Connection>;
49
50    fn with_redis<F, R>(&self, f: F) -> Result<R>
51    where
52        F: FnOnce(&Self::Conn) -> Result<R>;
53}
54
55pub trait JobService: Send {
56    fn run(&self, args: &[String]) -> Result<()>;
57
58    fn queue(&self) -> &str;
59
60    fn job_type(&self) -> &str;
61
62    fn box_clone(&self) -> Box<JobService>;
63}
64
65// HashMap<JobType, Box<JobService>>
66type JobServices = HashMap<String, Box<JobService>>;
67
68impl Clone for Box<JobService> {
69    fn clone(&self) -> Box<JobService> {
70        self.box_clone()
71    }
72}
73
74pub fn start<R>(c: Config<R>) -> Result<()>
75where
76    R: WithRedis,
77{
78    assert_ne!(c.fetch_timeout, 0);
79    assert!(!c.job_services.is_empty());
80    assert_ne!(c.tick_sleep, 0);
81    assert_ne!(c.watch_dead_thread, 0);
82    assert_ne!(c.workers, 0);
83
84    let new_worker = worker::NewResqueWorker::new(
85        c.host,
86        c.with_redis,
87        c.job_services,
88        Duration::from_secs(c.tick_sleep),
89        c.fetch_timeout,
90    );
91
92    let thread_pool = ThreadPool::new(
93        c.workers,
94        Duration::from_secs(c.stop_timeout),
95        Duration::from_secs(c.watch_dead_thread),
96        c.closed,
97    );
98
99    thread_pool.run(new_worker).wait()
100}
101
102pub fn enqueue<C>(redis_conn: &C, queue: &str, class: &str, args: &[String]) -> Result<()>
103where
104    C: Deref<Target = redis::Connection>,
105{
106    let item = JobItem {
107        class: Cow::from(class),
108        args: Cow::from(args),
109    };
110
111    let json = serde_json::to_string(&item)?;
112
113    let _: (i64, i64) = redis::pipe()
114        .cmd("SADD")
115        .arg(redis_keys::queues())
116        .arg(queue)
117        .cmd("RPUSH")
118        .arg(redis_keys::queue(queue))
119        .arg(json)
120        .query(redis_conn.deref())?;
121
122    Ok(())
123}
124
125#[derive(Debug, Serialize, Deserialize, Clone)]
126struct JobItem<'a> {
127    class: Cow<'a, str>,
128    args: Cow<'a, [String]>,
129}
130
131impl ErrorExt for Error {
132    fn new_err(s: Cow<'static, str>) -> Self {
133        match s {
134            Cow::Borrowed(s) => s.into(),
135            Cow::Owned(s) => s.into(),
136        }
137    }
138}