forked from orionz/minion
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathminion.rb
More file actions
144 lines (118 loc) · 2.67 KB
/
Copy pathminion.rb
File metadata and controls
144 lines (118 loc) · 2.67 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
require 'uri'
require 'json' unless defined? ActiveSupport::JSON
require 'mq'
require 'bunny'
require 'minion/handler'
module Minion
extend self
def url=(url)
@@config_url = url
end
def enqueue(jobs, data = {})
raise "cannot enqueue a nil job" if jobs.nil?
raise "cannot enqueue an empty job" if jobs.empty?
## jobs can be one or more jobs
if jobs.respond_to? :shift
queue = jobs.shift
data["next_job"] = jobs unless jobs.empty?
else
queue = jobs
end
encoded = JSON.dump(data)
log "send: #{queue}:#{encoded}"
bunny.queue(queue, :durable => true, :auto_delete => false).publish(encoded)
end
def log(msg)
@@logger ||= proc { |m| puts "#{Time.now} :minion: #{m}" }
@@logger.call(msg)
end
def error(&blk)
@@error_handler = blk
end
def logger(&blk)
@@logger = blk
end
def job(queue, options = {}, &blk)
handler = Minion::Handler.new queue
handler.when = options[:when] if options[:when]
handler.unsub = lambda do
log "unsubscribing to #{queue}"
MQ.queue(queue, :durable => true, :auto_delete => false).unsubscribe
end
handler.sub = lambda do
log "subscribing to #{queue}"
MQ.queue(queue, :durable => true, :auto_delete => false).subscribe(:ack => true) do |h,m|
return if AMQP.closing?
begin
log "recv: #{queue}:#{m}"
args = decode_json(m)
result = yield(args)
next_job(args, result)
rescue Object => e
raise unless error_handler
error_handler.call(e,queue,m,h)
end
h.ack
check_all
end
end
@@handlers ||= []
at_exit { Minion.run } if @@handlers.size == 0
@@handlers << handler
end
def decode_json(string)
if defined? ActiveSupport::JSON
ActiveSupport::JSON.decode string
else
JSON.load string
end
end
def check_all
@@handlers.each { |h| h.check }
end
def run
log "Starting minion"
Signal.trap('INT') { AMQP.stop{ EM.stop } }
Signal.trap('TERM'){ AMQP.stop{ EM.stop } }
EM.run do
AMQP.start(amqp_config) do
MQ.prefetch(1)
check_all
end
end
end
def amqp_url
@@amqp_url ||= ENV["AMQP_URL"] || "amqp://guest:guest@localhost/"
end
def amqp_url=(url)
@@amqp_url = url
end
private
def amqp_config
uri = URI.parse(amqp_url)
{
:vhost => uri.path,
:host => uri.host,
:user => uri.user,
:port => (uri.port || 5672),
:pass => uri.password
}
rescue Object => e
raise "invalid AMQP_URL: #{uri.inspect} (#{e})"
end
def new_bunny
b = Bunny.new(amqp_config)
b.start
b
end
def bunny
@@bunny ||= new_bunny
end
def next_job(args, response)
queue = args.delete("next_job")
enqueue(queue,args.merge(response)) if queue and not queue.empty?
end
def error_handler
@@error_handler ||= nil
end
end