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

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_taskObject



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

Raises:

  • (ArgumentError)


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_handlerObject



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_fallbackObject



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_nowObject



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

Returns:

  • (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

.schedulerObject



55
56
57
# File 'lib/async/background/runtime.rb', line 55

def scheduler
  Fiber.scheduler
end

.scheduler!Object



59
60
61
# File 'lib/async/background/runtime.rb', line 59

def scheduler!
  Fiber.scheduler or raise SchedulerRequired, NO_SCHEDULER_MESSAGE
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

Raises:

  • (ArgumentError)


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