summaryrefslogtreecommitdiff
path: root/lib/custodian/queue.rb
blob: 8d48575b41be96de3b27732c592d80b85a096e6c (plain)
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
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
#
# We don't necessarily expect that both libraries will be present,
# so long as one is that'll allow things to work.
#
%w( redis beanstalk-client ).each do |library|
  begin
    require library
  rescue LoadError
    ENV["DEBUG"] && puts( "Failed to load the library: #{library}" )
  end
end



module Custodian


  #
  # An abstraction layer for our queue.
  #
  class QueueType

    #
    # Class-Factory
    #
    def self.create type
      case type
      when "redis"
        RedisQueueType.new
      when "beanstalk"
        BeanstalkQueueType.new
      else
        raise "Bad queue-type: #{type}"
      end
    end


    #
    # Retrieve a job from the queue.
    #
    def fetch(timeout)
      raise "Subclasses must implement this method!"
    end


    #
    # Add a new job to the queue.
    #
    def add(job_string)
      raise "Subclasses must implement this method!"
    end


    #
    # Get the size of the queue
    #
    def size?
      raise "Subclasses must implement this method!"
    end

  end




  #
  # This is a simple FIFO queue which uses Redis for storage.
  #
  class RedisQueueType < QueueType


    #
    #  Connect to the server on localhost
    #
    def initialize
      host = ENV["REDIS"] || "127.0.0.1"
      @redis = Redis.new( :host => host )
    end


    #
    #  Fetch a job from the queue - the timeout parameter is ignored.
    #
    def fetch(timeout)
      job = false
      while( ! job )
        job = @redis.lpop( "queue" )
      end
      return( job )
    end

    #
    #  Add a new job to the queue.
    #
    def add(job_string)
      @redis.rpush( "queue", job_string )
    end


    #
    #  How many jobs in the queue?
    #
    def size?
      @redis.llen( "queue" )
    end
  end



  #
  #  Use the beanstalkd-queue for its intended purpose
  #
  class BeanstalkQueueType < QueueType

    #
    #  Connect to the server on localhost
    #
    def initialize
      host = ENV["QUEUE"] || "127.0.0.1:11300"
      @queue = Beanstalk::Pool.new([host] )
    end

    #
    #  Here we fetch a value from the queue, and delete it at the same time.
    #
    #  The timeout is used to specify the period we wait for a new job.
    #
    def fetch(timeout)
      begin
        j = @queue.reserve(timeout)
        if ( j ) then
          b = j.body
          j.delete
          return b
        else
          raise "ERRROR"
        end
      rescue Beanstalk::TimedOut => ex
        return nil
      end
    end


    #
    #  Add a new job to the queue.
    #
    def add(job_string)
      @queue.put(job_string)
    end


    #
    #  Get the size of the queue
    #
    def size?
      stats = @queue.stats()
      ( stats['current-jobs-ready'] || 0 )
    end

  end

end