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
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
|
module EventMachine
# A simple async resource pool based on a resource and work queue. Resources
# are enqueued and work waits for resources to become available.
#
# @example
# require 'em-http-request'
#
# EM.run do
# pool = EM::Pool.new
# spawn = lambda { pool.add EM::HttpRequest.new('http://example.org') }
# 10.times { spawn[] }
# done, scheduled = 0, 0
#
# check = lambda do
# done += 1
# if done >= scheduled
# EM.stop
# end
# end
#
# pool.on_error { |conn| spawn[] }
#
# 100.times do |i|
# scheduled += 1
# pool.perform do |conn|
# req = conn.get :path => '/', :keepalive => true
#
# req.callback do
# p [:success, conn.object_id, i, req.response.size]
# check[]
# end
#
# req.errback { check[] }
#
# req
# end
# end
# end
#
# Resources are expected to be controlled by an object responding to a
# deferrable/completion style API with callback and errback blocks.
#
class Pool
def initialize
@resources = EM::Queue.new
@removed = []
@contents = []
@on_error = nil
end
def add resource
@contents << resource
requeue resource
end
def remove resource
@contents.delete resource
@removed << resource
end
# Returns a list for introspection purposes only. You should *NEVER* call
# modification or work oriented methods on objects in this list. A good
# example use case is periodic statistics collection against a set of
# connection resources.
#
# @example
# pool.contents.inject(0) { |sum, connection| connection.num_bytes }
def contents
@contents.dup
end
# Define a default catch-all for when the deferrables returned by work
# blocks enter a failed state. By default all that happens is that the
# resource is returned to the pool. If on_error is defined, this block is
# responsible for re-adding the resource to the pool if it is still usable.
# In other words, it is generally assumed that on_error blocks explicitly
# handle the rest of the lifetime of the resource.
def on_error *a, &b
@on_error = EM::Callback(*a, &b)
end
# Perform a given #call-able object or block. The callable object will be
# called with a resource from the pool as soon as one is available, and is
# expected to return a deferrable.
#
# The deferrable will have callback and errback added such that when the
# deferrable enters a finished state, the object is returned to the pool.
#
# If on_error is defined, then objects are not automatically returned to the
# pool.
def perform(*a, &b)
work = EM::Callback(*a, &b)
@resources.pop do |resource|
if removed? resource
@removed.delete resource
reschedule work
else
process work, resource
end
end
end
alias reschedule perform
# A peek at the number of enqueued jobs waiting for resources
def num_waiting
@resources.num_waiting
end
# Removed will show resources in a partial pruned state. Resources in the
# removed list may not appear in the contents list if they are currently in
# use.
def removed? resource
@removed.include? resource
end
protected
def requeue resource
@resources.push resource
end
def failure resource
if @on_error
@contents.delete resource
@on_error.call resource
# Prevent users from calling a leak.
@removed.delete resource
else
requeue resource
end
end
def completion deferrable, resource
deferrable.callback { requeue resource }
deferrable.errback { failure resource }
end
def process work, resource
deferrable = work.call resource
if deferrable.kind_of?(EM::Deferrable)
completion deferrable, resource
else
raise ArgumentError, "deferrable expected from work"
end
rescue
failure resource
raise
end
end
end
|