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
|
# File 'lib/resque/unique_by_arity/modulizer.rb', line 4
def self.to_mod(configuration)
Module.new do
if configuration.unique_in_queue || configuration.unique_at_runtime || configuration.unique_across_queues
define_method(:redis_unique_hash) do |payload, arity_for_uniqueness|
arity_for_uniqueness ||= configuration.arity_for_uniqueness
payload = Resque.decode(Resque.encode(payload))
Resque::UniqueByArity.debug("payload is #{payload.inspect}")
job = payload["class"]
args = payload["args"] || []
args.map! do |arg|
arg.is_a?(Hash) ? arg.sort : arg
end
uniqueness_args = if arity_for_uniqueness.zero?
[]
else
args[0..(arity_for_uniqueness - 1)]
end
args = {class: job, args: uniqueness_args}
[Digest::MD5.hexdigest(Resque.encode(args)), uniqueness_args]
end
end
if configuration.unique_in_queue || configuration.unique_across_queues
define_method(:unique_in_queue_redis_key_prefix) do
"unique_job:#{self}"
end
define_method(:unique_in_queue_redis_key) do |queue, payload|
arity_for_uniqueness = if configuration.unique_in_queue
configuration.arity_for_uniqueness_in_queue
elsif configuration.unique_across_queues
configuration.arity_for_uniqueness_across_queues
end
unique_hash, args_for_uniqueness = redis_unique_hash(payload, arity_for_uniqueness)
key = "#{unique_in_queue_key_namespace(queue)}:#{unique_in_queue_redis_key_prefix}:#{unique_hash}"
Resque::UniqueByArity.debug("#{self}.unique_in_queue_redis_key for #{args_for_uniqueness} is: #{key}")
key
end
define_method(:purge_unique_queued_redis_keys) do
key_match = "#{unique_in_queue_key_namespace(instance_variable_get(:@queue))}:#{unique_in_queue_redis_key_prefix}:*"
keys = Resque.redis.keys(key_match)
Resque::UniqueByArity.log("#{Resque::UniqueByArity::PLUGIN_TAG}#{Resque::UniqueInQueue::PLUGIN_TAG} Purging #{keys.length} keys from #{key_match}")
Resque.redis.del keys unless keys.empty?
end
if configuration.unique_in_queue
define_method(:unique_in_queue_key_namespace) do |queue = nil|
"#{unique_in_queue_key_base}:queue:#{queue}:job"
end
elsif configuration.unique_across_queues
define_method(:unique_in_queue_key_namespace) do |_queue = nil|
"#{unique_in_queue_key_base}:across_queues:job"
end
end
end
if configuration.unique_at_runtime
define_method(:runtime_key_namespace) do
"#{unique_at_runtime_key_base}:#{self}"
end
define_method(:unique_at_runtime_redis_key) do |*args|
unique_hash, args_for_uniqueness = redis_unique_hash({"class" => to_s, "args" => args}, configuration.arity_for_uniqueness_at_runtime)
key = "#{runtime_key_namespace}:#{unique_hash}"
Resque::UniqueByArity.debug("#{Resque::UniqueAtRuntime::PLUGIN_TAG} #{self}.unique_at_runtime_redis_key for #{args_for_uniqueness} is: #{key}")
key
end
define_method(:purge_unique_at_runtime_redis_keys) do
key_match = "#{runtime_key_namespace}:*"
keys = Resque.redis.keys(key_match)
Resque::UniqueByArity.log("#{Resque::UniqueByArity::PLUGIN_TAG}#{Resque::UniqueAtRuntime::PLUGIN_TAG} Purging #{keys.length} keys from #{key_match}")
Resque.redis.del keys unless keys.empty?
end
end
end
end
|