Class: CI::Queue::Redis::Base
- Inherits:
-
Object
- Object
- CI::Queue::Redis::Base
show all
- Includes:
- Common
- Defined in:
- lib/ci/queue/redis/base.rb
Constant Summary
collapse
- CONNECTION_ERRORS =
[
::Redis::BaseConnectionError,
::SocketError,
].freeze
Instance Attribute Summary
Attributes included from Common
#config
Instance Method Summary
collapse
Methods included from Common
#flaky?, #report_failure!, #report_success!, #rescue_connection_errors
Constructor Details
#initialize(redis_url, config) ⇒ Base
Returns a new instance of Base.
12
13
14
15
16
|
# File 'lib/ci/queue/redis/base.rb', line 12
def initialize(redis_url, config)
@redis_url = redis_url
@redis = ::Redis.new(url: redis_url)
@config = config
end
|
Instance Method Details
#exhausted? ⇒ Boolean
18
19
20
|
# File 'lib/ci/queue/redis/base.rb', line 18
def exhausted?
queue_initialized? && size == 0
end
|
#progress ⇒ Object
36
37
38
|
# File 'lib/ci/queue/redis/base.rb', line 36
def progress
total - size
end
|
#queue_initialized? ⇒ Boolean
56
57
58
59
60
61
|
# File 'lib/ci/queue/redis/base.rb', line 56
def queue_initialized?
@queue_initialized ||= begin
status = master_status
status == 'ready' || status == 'finished'
end
end
|
#size ⇒ Object
22
23
24
25
26
27
|
# File 'lib/ci/queue/redis/base.rb', line 22
def size
redis.multi do
redis.llen(key('queue'))
redis.zcard(key('running'))
end.inject(:+)
end
|
#to_a ⇒ Object
29
30
31
32
33
34
|
# File 'lib/ci/queue/redis/base.rb', line 29
def to_a
redis.multi do
redis.lrange(key('queue'), 0, -1)
redis.zrange(key('running'), 0, -1)
end.flatten.reverse.map { |k| index.fetch(k) }
end
|
#wait_for_master(timeout: 10) ⇒ Object
40
41
42
43
44
45
46
47
48
49
50
|
# File 'lib/ci/queue/redis/base.rb', line 40
def wait_for_master(timeout: 10)
return true if master?
(timeout * 10 + 1).to_i.times do
if queue_initialized?
return true
else
sleep 0.1
end
end
raise LostMaster, "The master worker is still `#{master_status}` after 10 seconds waiting."
end
|
#workers_count ⇒ Object
52
53
54
|
# File 'lib/ci/queue/redis/base.rb', line 52
def workers_count
redis.scard(key('workers'))
end
|