Class: Purplelight::ByteQueue
- Inherits:
-
Object
- Object
- Purplelight::ByteQueue
- Defined in:
- lib/purplelight/queue.rb
Overview
Sized queue that tracks bytes to apply backpressure.
Instance Method Summary collapse
- #close ⇒ Object
-
#initialize(max_bytes: 128 * 1024 * 1024) ⇒ ByteQueue
constructor
A new instance of ByteQueue.
- #pop ⇒ Object
- #push(item, bytes:) ⇒ Object
- #size_bytes ⇒ Object
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
#close ⇒ Object
46 47 48 49 50 51 |
# File 'lib/purplelight/queue.rb', line 46 def close @mutex.synchronize do @closed = true @cv.broadcast end end |
#pop ⇒ Object
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_bytes ⇒ Object
53 54 55 |
# File 'lib/purplelight/queue.rb', line 53 def size_bytes @mutex.synchronize { @bytes } end |