Class: WorkQueue

Inherits:
Object
  • Object
show all
Defined in:
lib/workqueue.rb,
lib/workqueue/version.rb

Defined Under Namespace

Classes: ThreadsafeCounter

Instance Attribute Summary collapse

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(init_queue = [], opts = {}, &job) ⇒ WorkQueue

Returns a new instance of WorkQueue.



17
18
19
20
21
22
23
24
25
# File 'lib/workqueue.rb', line 17

def initialize(init_queue=[], opts={}, &job)
  @job = job
  @queue = Queue.new
  @mutex = Mutex.new

  opts.each { |k, v| send(:"#{k}=", v) }

  concat(init_queue)
end

Instance Attribute Details

#job ⇒ Object (readonly)

Returns the value of attribute job.



16
17
18
# File 'lib/workqueue.rb', line 16

def job
  @job
end

#queue ⇒ Object (readonly)

Returns the value of attribute queue.



15
16
17
# File 'lib/workqueue.rb', line 15

def queue
  @queue
end

#size ⇒ Object



28
29
30
# File 'lib/workqueue.rb', line 28

def size
  @size ||= 2
end

Class Method Details

.version ⇒ Object



2
3
4
# File 'lib/workqueue/version.rb', line 2

def self.version
  '0.1.2'
end

Instance Method Details

#abort! ⇒ Object



55
56
57
# File 'lib/workqueue.rb', line 55

def abort!
  @mutex.synchronize { @aborted = true }
end

#concat(arr) ⇒ Object



74
75
76
# File 'lib/workqueue.rb', line 74

def concat(arr)
  arr.each { |x| push(x) }
end

#join ⇒ Object



78
79
80
81
82
83
84
# File 'lib/workqueue.rb', line 78

def join
  @mutex.synchronize { @joined = true }

  workers.each(&:join)

  self
end

#push(e) ⇒ Object Also known as: <<



67
68
69
70
71
# File 'lib/workqueue.rb', line 67

def push(e)
  queue.push([e, cursor.incr])

  self
end

#results ⇒ Object



86
87
88
89
# File 'lib/workqueue.rb', line 86

def results
  join
  aggregate
end

#run ⇒ Object



59
60
61
62
63
64
65
# File 'lib/workqueue.rb', line 59

def run
  @workers = (1..size).map do
    Thread.new { work! }
  end

  self
end

#work! ⇒ Object



32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
# File 'lib/workqueue.rb', line 32

def work!
  loop do
    begin
      # initialize these here so they can be set in
      # the synchronize block
      payload, index = []

      @mutex.synchronize {
        unless @aborted or (@joined and queue.empty?)
          payload, index = queue.shift
        end
      }

      break if payload.nil?

      aggregate[index] = job.call(payload)
    rescue Exception
      abort!
      raise
    end
  end
end