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
|
# File 'lib/frenzy_bunnies/worker.rb', line 26
def start(context)
@jobs_stats = { :failed => Atomic.new(0), :passed => Atomic.new(0) }
@working_since = Time.now
@logger = context.logger
@queue_opts[:prefetch] ||= 10
@queue_opts[:durable] ||= false
@queue_opts[:timeout_job_after] = 5 if @queue_opts[:timeout_job_after].nil?
@queue_opts[:handler] ||= FrenzyBunnies::Handlers::Oneshot
if @queue_opts[:threads]
@thread_pool = MarchHare::ThreadPools.fixed_of_size(@queue_opts[:threads])
else
@thread_pool = MarchHare::ThreadPools.dynamically_growing
end
q = context.queue_factory.build_queue(@queue_name, @queue_opts)
say "#{@queue_opts[:threads]} with #{@queue_opts[:prefetch]} prefetch on <#{@queue_name}>."
q.subscribe(:ack => true, :blocking => false) do |, msg|
@thread_pool.submit do
worker = new
handler = @queue_opts[:handler].new(.channel, q, @logger, { queue_options: @queue_opts })
begin
Timeout::timeout(@queue_opts[:timeout_job_after]) do
if worker.work(msg)
handler.acknowledge(, msg)
incr! :passed
else
handler.reject(, msg)
incr! :failed
error "REJECTED", msg
end
end
rescue Timeout::Error
handler.timeout(, msg)
incr! :failed
error "TIMEOUT #{@queue_opts[:timeout_job_after]}s", msg
rescue
handler.reject(, msg)
incr! :failed
error "ERROR #{$!}", msg
end
end
end
say "workers up."
end
|