Class: CI::Queue::Redis::Base

Inherits:
Object
  • Object
show all
Includes:
Common
Defined in:
lib/ci/queue/redis/base.rb

Direct Known Subclasses

Supervisor, Worker

Constant Summary collapse

CONNECTION_ERRORS =
[
  ::Redis::BaseConnectionError,
  ::SocketError, # https://github.com/redis/redis-rb/pull/631
].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

Returns:

  • (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

Returns:

  • (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

Raises:



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