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
65type 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}