Module: Async::Background::Runtime
- Defined in:
- lib/async/background/runtime.rb,
lib/async/background/runtime/task.rb,
lib/async/background/runtime/semaphore.rb,
lib/async/background/runtime/task_group.rb,
lib/async/background/runtime/notification.rb
Defined Under Namespace
Classes: Cancel, Deadline, Error, Notification, SchedulerRequired, Semaphore, Task, TaskGroup, TimeoutError
Constant Summary
collapse
- UNSET =
Object.new
- CURRENT_TASK_KEY =
:async_background_current_task
- WAITER_KEY =
:async_background_waiter
- DEADLINE_MESSAGE =
'execution expired'
- NO_SCHEDULER_MESSAGE =
<<~MESSAGE
Async::Background requires an active Fiber scheduler.
Install one in the host process before calling this, for example:
require "async/background/scheduler"
Async::Background::Scheduler.run { runner.run }
or install one yourself:
Fiber.set_scheduler(Itsi::Scheduler.new) # itsi-scheduler
Async { runner.run } # async / falcon
MESSAGE
- MISSING_TIMEOUT_HOOK_WARNING =
<<~MESSAGE
Async::Background: %s does not implement #timeout_after.
Falling back to stdlib Timeout, which uses Thread#raise and can deliver
the timeout to an unrelated fiber. Job timeouts are therefore not safe
on this scheduler. Use a scheduler that implements #timeout_after
(async, itsi-scheduler) or run jobs with `timeout: nil`.
MESSAGE
Class Method Summary
collapse
-
.arm_deadline(seconds, &block) ⇒ Object
-
.build_waiter(blocker) ⇒ Object
-
.clear_fiber_local(key, expected) ⇒ Object
-
.current_task ⇒ Object
-
.current_task=(task) ⇒ Object
-
.deadline_for(timeout) ⇒ Object
-
.error_handler ⇒ Object
-
.error_handler=(handler) ⇒ Object
-
.fiber_local(key) ⇒ Object
-
.fiber_local_fallback ⇒ Object
-
.monotonic_now ⇒ Object
-
.native_timeouts?(target = Fiber.scheduler) ⇒ Boolean
-
.park(blocker, waiter, deadline = nil) ⇒ Object
-
.report_error(task, error, handler = UNSET) ⇒ Object
-
.scheduler ⇒ Object
-
.scheduler! ⇒ Object
-
.set_fiber_local(key, value) ⇒ Object
-
.spawn(name: nil, on_error: UNSET, &block) ⇒ Object
-
.wake(waiter, blocker) ⇒ Object
-
.wake_dequeued(waiter, blocker) ⇒ Object
-
.warn_missing_timeout_hook(target) ⇒ Object
-
.with_error_handler(handler) ⇒ Object
-
.with_timeout(duration, on_timeout: UNSET) ⇒ Object
-
.with_waiter(blocker, waiters) ⇒ Object
Class Method Details
.arm_deadline(seconds, &block) ⇒ Object
220
221
222
223
224
225
226
227
|
# File 'lib/async/background/runtime.rb', line 220
def arm_deadline(seconds, &block)
target = scheduler!
return target.timeout_after(seconds, Deadline, DEADLINE_MESSAGE, &block) if native_timeouts?(target)
warn_missing_timeout_hook(target)
::Timeout.timeout(seconds, Deadline, DEADLINE_MESSAGE, &block)
end
|
.build_waiter(blocker) ⇒ Object
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
|
# File 'lib/async/background/runtime.rb', line 75
def build_waiter(blocker)
existing = fiber_local(WAITER_KEY)
if existing&.[](:blocker)
raise Error, "waiter already parked on #{existing[:blocker].class}; " \
'a fiber may only park in one place at a time'
end
waiter = {
fiber: Fiber.current,
scheduler: scheduler!,
ready: false,
queued: false,
blocker: blocker
}
set_fiber_local(WAITER_KEY, waiter)
waiter
end
|
.clear_fiber_local(key, expected) ⇒ Object
194
195
196
197
|
# File 'lib/async/background/runtime.rb', line 194
def clear_fiber_local(key, expected)
current = fiber_local(key)
set_fiber_local(key, nil) if current.equal?(expected)
end
|
.current_task ⇒ Object
67
68
69
|
# File 'lib/async/background/runtime.rb', line 67
def current_task
fiber_local(CURRENT_TASK_KEY)
end
|
.current_task=(task) ⇒ Object
71
72
73
|
# File 'lib/async/background/runtime.rb', line 71
def current_task=(task)
set_fiber_local(CURRENT_TASK_KEY, task)
end
|
.deadline_for(timeout) ⇒ Object
147
148
149
150
151
152
153
154
|
# File 'lib/async/background/runtime.rb', line 147
def deadline_for(timeout)
return nil if timeout.nil?
seconds = Float(timeout)
raise ArgumentError, 'timeout must be non-negative and finite' unless seconds.finite? && seconds >= 0
monotonic_now + seconds
end
|
.error_handler ⇒ Object
160
161
162
|
# File 'lib/async/background/runtime.rb', line 160
def error_handler
@error_handler
end
|
.error_handler=(handler) ⇒ Object
156
157
158
|
# File 'lib/async/background/runtime.rb', line 156
def error_handler=(handler)
@error_handler = handler
end
|
.fiber_local(key) ⇒ Object
182
183
184
185
186
|
# File 'lib/async/background/runtime.rb', line 182
def fiber_local(key)
Fiber[key]
rescue ArgumentError, FiberError
fiber_local_fallback[key]
end
|
.fiber_local_fallback ⇒ Object
199
200
201
202
203
|
# File 'lib/async/background/runtime.rb', line 199
def fiber_local_fallback
fiber = Fiber.current
store = fiber.instance_variable_get(:@async_background_locals)
store || fiber.instance_variable_set(:@async_background_locals, {})
end
|
.monotonic_now ⇒ Object
143
144
145
|
# File 'lib/async/background/runtime.rb', line 143
def monotonic_now
Process.clock_gettime(Process::CLOCK_MONOTONIC)
end
|
.native_timeouts?(target = Fiber.scheduler) ⇒ Boolean
63
64
65
|
# File 'lib/async/background/runtime.rb', line 63
def native_timeouts?(target = Fiber.scheduler)
!target.nil? && target.respond_to?(:timeout_after)
end
|
.park(blocker, waiter, deadline = nil) ⇒ Object
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
|
# File 'lib/async/background/runtime.rb', line 93
def park(blocker, waiter, deadline = nil)
scheduler = waiter[:scheduler]
task = current_task
task&.enter_block(waiter)
until waiter[:ready] || yield
if deadline
remaining = deadline - monotonic_now
return false if remaining <= 0
scheduler.block(blocker, remaining)
else
scheduler.block(blocker, nil)
end
task&.raise_if_cancelled!
end
true
ensure
waiter[:blocker] = nil
clear_fiber_local(WAITER_KEY, waiter)
task&.exit_block
end
|
.report_error(task, error, handler = UNSET) ⇒ Object
172
173
174
175
176
177
178
179
180
|
# File 'lib/async/background/runtime.rb', line 172
def report_error(task, error, handler = UNSET)
handler = @error_handler if UNSET.equal?(handler)
return false unless handler
handler.call(task, error)
true
rescue StandardError
false
end
|
.scheduler ⇒ Object
55
56
57
|
# File 'lib/async/background/runtime.rb', line 55
def scheduler
Fiber.scheduler
end
|
.set_fiber_local(key, value) ⇒ Object
188
189
190
191
192
|
# File 'lib/async/background/runtime.rb', line 188
def set_fiber_local(key, value)
Fiber[key] = value
rescue ArgumentError, FiberError
fiber_local_fallback[key] = value
end
|
.spawn(name: nil, on_error: UNSET, &block) ⇒ Object
51
52
53
|
# File 'lib/async/background/runtime.rb', line 51
def spawn(name: nil, on_error: UNSET, &block)
Task.spawn(name: name, on_error: on_error, &block)
end
|
.wake(waiter, blocker) ⇒ Object
132
133
134
135
136
137
138
139
140
141
|
# File 'lib/async/background/runtime.rb', line 132
def wake(waiter, blocker)
return false if waiter[:ready]
waiter[:ready] = true
fiber = waiter[:fiber]
return false unless fiber.alive?
waiter[:scheduler].unblock(blocker, fiber)
true
end
|
.wake_dequeued(waiter, blocker) ⇒ Object
127
128
129
130
|
# File 'lib/async/background/runtime.rb', line 127
def wake_dequeued(waiter, blocker)
waiter[:queued] = false
wake(waiter, blocker)
end
|
.warn_missing_timeout_hook(target) ⇒ Object
229
230
231
232
233
234
235
|
# File 'lib/async/background/runtime.rb', line 229
def warn_missing_timeout_hook(target)
key = target.class
return if @warned_schedulers[key]
@warned_schedulers[key] = true
warn(format(MISSING_TIMEOUT_HOOK_WARNING, key))
end
|
.with_error_handler(handler) ⇒ Object
164
165
166
167
168
169
170
|
# File 'lib/async/background/runtime.rb', line 164
def with_error_handler(handler)
previous = @error_handler
@error_handler = handler
yield
ensure
@error_handler = previous
end
|
.with_timeout(duration, on_timeout: UNSET) ⇒ Object
205
206
207
208
209
210
211
212
213
214
215
216
217
218
|
# File 'lib/async/background/runtime.rb', line 205
def with_timeout(duration, on_timeout: UNSET)
return yield if duration.nil?
seconds = Float(duration)
raise ArgumentError, 'timeout must be positive and finite' unless seconds.finite? && seconds.positive?
begin
arm_deadline(seconds) { yield }
rescue Deadline => e
raise TimeoutError, e.message if UNSET.equal?(on_timeout)
on_timeout
end
end
|
.with_waiter(blocker, waiters) ⇒ Object
118
119
120
121
122
123
124
125
|
# File 'lib/async/background/runtime.rb', line 118
def with_waiter(blocker, waiters)
waiter = build_waiter(blocker)
waiter[:queued] = true
waiters << waiter
yield waiter
ensure
waiters.delete(waiter) if waiter[:queued]
end
|