Class: Purplelight::ByteQueue

Inherits:
Object
  • Object
show all
Defined in:
lib/purplelight/queue.rb

Overview

Sized queue that tracks bytes to apply backpressure.

Instance Method Summary collapse

Constructor Details

#initialize(max_bytes: 128 * 1024 * 1024) ⇒ ByteQueue

Returns a new instance of ByteQueue.



6
7
8
9
10
11
12
13
14
# File 'lib/purplelight/queue.rb', line 6

def initialize(max_bytes: 128 * 1024 * 1024)
  @max_bytes = max_bytes
  @queue = []
  @sizes = []
  @bytes = 0
  @closed = false
  @mutex = Mutex.new
  @cv = ConditionVariable.new
end

Instance Method Details

#closeObject



46
47
48
49
50
51
# File 'lib/purplelight/queue.rb', line 46

def close
  @mutex.synchronize do
    @closed = true
    @cv.broadcast
  end
end

#popObject



31
32
33
34
35
36
37
38
39
40
41
42
43
44
# File 'lib/purplelight/queue.rb', line 31

def pop
  @mutex.synchronize do
    while @queue.empty?
      return nil if @closed

      @cv.wait(@mutex)
    end
    item = @queue.shift
    bytes = @sizes.shift
    @bytes -= bytes
    @cv.broadcast
    item
  end
end

#push(item, bytes:) ⇒ Object



16
17
18
19
20
21
22
23
24
25
26
27
28
29
# File 'lib/purplelight/queue.rb', line 16

def push(item, bytes:)
  @mutex.synchronize do
    raise 'queue closed' if @closed

    while !@queue.empty? && (@bytes + bytes) > @max_bytes
      @cv.wait(@mutex)
      raise 'queue closed' if @closed
    end
    @queue << item
    @sizes << bytes
    @bytes += bytes
    @cv.broadcast
  end
end

#size_bytesObject



53
54
55
# File 'lib/purplelight/queue.rb', line 53

def size_bytes
  @mutex.synchronize { @bytes }
end