Class: PgPipeline::BoundedQueue

Inherits:
Object
  • Object
show all
Defined in:
lib/pg_pipeline/bounded_queue.rb

Instance Method Summary collapse

Constructor Details

#initialize(limit) ⇒ BoundedQueue

Returns a new instance of BoundedQueue.



9
10
11
12
13
14
15
16
17
18
19
20
# File 'lib/pg_pipeline/bounded_queue.rb', line 9

def initialize(limit)
  @limit = Integer(limit)
  raise ArgumentError, "limit must be >= 1" if @limit < 1
rescue ArgumentError, TypeError
  raise ArgumentError, "limit must be an integer >= 1"
else
  @items = []
  @consumers = []
  @producers = []
  @closed = false
  @close_error = nil
end

Instance Method Details

#close(error) ⇒ Object



50
51
52
53
54
55
56
57
58
# File 'lib/pg_pipeline/bounded_queue.rb', line 50

def close(error)
  return if @closed

  @closed = true
  @close_error = error
  wake_all(@consumers)
  wake_all(@producers)
  nil
end

#dequeueObject



36
37
38
39
40
41
42
43
44
45
46
47
48
# File 'lib/pg_pipeline/bounded_queue.rb', line 36

def dequeue
  loop do
    unless @items.empty?
      item = @items.shift
      wake_one(@producers)
      return item
    end

    return nil if @closed

    wait_on(@consumers)
  end
end

#drainObject



65
66
67
68
69
70
# File 'lib/pg_pipeline/bounded_queue.rb', line 65

def drain
  items = @items
  @items = []
  wake_all(@producers)
  items
end

#empty?Boolean

Returns:

  • (Boolean)


60
# File 'lib/pg_pipeline/bounded_queue.rb', line 60

def empty? = @items.empty?

#enqueue(item) ⇒ Object



22
23
24
25
26
27
28
29
30
31
32
33
34
# File 'lib/pg_pipeline/bounded_queue.rb', line 22

def enqueue(item)
  loop do
    raise_close_error if @closed

    if @items.size < @limit
      @items << item
      wake_one(@consumers)
      return item
    end

    wait_on(@producers)
  end
end

#sizeObject



61
# File 'lib/pg_pipeline/bounded_queue.rb', line 61

def size = @items.size

#waiting_consumersObject



63
# File 'lib/pg_pipeline/bounded_queue.rb', line 63

def waiting_consumers = @consumers.size

#waiting_producersObject



62
# File 'lib/pg_pipeline/bounded_queue.rb', line 62

def waiting_producers = @producers.size