forked from discourse/discourse
-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathwork_queue.rb
110 lines (88 loc) · 1.81 KB
/
work_queue.rb
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
# 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
def initialize(key, limit, &blk)
@limit = limit
@size = 0
@key = key
@elements = Hash.new { |h, k| h[k] = blk.call }
end
def push(task, force:)
raise WorkQueueFull if !force && @size >= @limit
key = task[@key]
@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
task
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