2023-07-28 07:53:51 -04:00
|
|
|
# frozen_string_literal: true
|
|
|
|
|
|
|
|
require "monitor"
|
|
|
|
|
|
|
|
module WorkQueue
|
|
|
|
class WorkQueueFull < StandardError
|
|
|
|
end
|
|
|
|
|
|
|
|
class ThreadSafeWrapper
|
|
|
|
include MonitorMixin
|
|
|
|
|
|
|
|
def initialize(queue)
|
|
|
|
mon_initialize
|
|
|
|
|
|
|
|
@queue = queue
|
|
|
|
@has_items = new_cond
|
|
|
|
end
|
|
|
|
|
|
|
|
def push(task, force:)
|
|
|
|
synchronize do
|
|
|
|
previously_empty = @queue.empty?
|
|
|
|
@queue.push(task, force: force)
|
|
|
|
|
|
|
|
@has_items.signal if previously_empty
|
|
|
|
end
|
|
|
|
end
|
|
|
|
|
|
|
|
def shift(block:)
|
|
|
|
synchronize do
|
|
|
|
loop do
|
|
|
|
if task = @queue.shift
|
|
|
|
break task
|
|
|
|
elsif block
|
|
|
|
@has_items.wait
|
|
|
|
else
|
|
|
|
break nil
|
|
|
|
end
|
|
|
|
end
|
|
|
|
end
|
|
|
|
end
|
|
|
|
|
|
|
|
def empty?
|
|
|
|
synchronize { @queue.empty? }
|
|
|
|
end
|
|
|
|
|
|
|
|
def size
|
|
|
|
synchronize { @queue.size }
|
|
|
|
end
|
|
|
|
end
|
|
|
|
|
|
|
|
class FairQueue
|
|
|
|
attr_reader :size
|
|
|
|
|
2024-02-07 14:47:50 -05:00
|
|
|
def initialize(key, limit, &blk)
|
2023-07-28 07:53:51 -04:00
|
|
|
@limit = limit
|
|
|
|
@size = 0
|
2024-02-07 14:47:50 -05:00
|
|
|
@key = key
|
2023-07-28 07:53:51 -04:00
|
|
|
@elements = Hash.new { |h, k| h[k] = blk.call }
|
|
|
|
end
|
|
|
|
|
|
|
|
def push(task, force:)
|
|
|
|
raise WorkQueueFull if !force && @size >= @limit
|
2024-02-07 14:47:50 -05:00
|
|
|
key = task[@key]
|
2023-07-28 07:53:51 -04:00
|
|
|
@elements[key].push(task, force: force)
|
|
|
|
@size += 1
|
|
|
|
nil
|
|
|
|
end
|
|
|
|
|
|
|
|
def shift
|
|
|
|
unless @elements.empty?
|
|
|
|
key, queue = @elements.shift
|
|
|
|
|
|
|
|
task = queue.shift
|
|
|
|
|
|
|
|
@elements[key] = queue unless queue.empty?
|
|
|
|
@size -= 1
|
2024-02-07 14:47:50 -05:00
|
|
|
task
|
2023-07-28 07:53:51 -04:00
|
|
|
end
|
|
|
|
end
|
|
|
|
|
|
|
|
def empty?
|
|
|
|
@elements.empty?
|
|
|
|
end
|
|
|
|
end
|
|
|
|
|
|
|
|
class BoundedQueue
|
|
|
|
def initialize(limit)
|
|
|
|
@limit = limit
|
|
|
|
@elements = []
|
|
|
|
end
|
|
|
|
|
|
|
|
def push(task, force:)
|
|
|
|
raise WorkQueueFull if !force && @elements.size >= @limit
|
|
|
|
@elements << task
|
|
|
|
nil
|
|
|
|
end
|
|
|
|
|
|
|
|
def shift
|
|
|
|
@elements.shift
|
|
|
|
end
|
|
|
|
|
|
|
|
def empty?
|
|
|
|
@elements.empty?
|
|
|
|
end
|
|
|
|
|
|
|
|
def size
|
|
|
|
@elements.size
|
|
|
|
end
|
|
|
|
end
|
|
|
|
end
|