Class: Farce::Scheduler
- Inherits:
-
Abstract::Scheduler
- Object
- Abstract::Scheduler
- Farce::Scheduler
- Includes:
- Farce::Shareable::Unfreezable, Internal::Inspect
- Defined in:
- lib/farce/scheduler.rb
Overview
A Ractor-shareable task scheduler.
Tasks run as fibers using a fiber scheduler, or as threads when the Ruby implementation doesn't support Fiber schedulers (ie TruffleRuby).
If supports any existing Fiber scheduler, including Async and CarbonFiber, but can also be used as a Fiber scheduler itself.
Use Scheduler.create to start a worker, or construct an instance and install it with
Fiber.set_scheduler on the current thread. A scheduler can only be launched
or installed once.
It will keep the fiber scheduler alive until it is explicitly closed.
Closing a scheduler stops new submissions through #schedule. The underlying
fiber scheduler has its own lifecycle: draining fibers may still spawn child
fibers with Fiber.schedule while it is closing.
Class Method Summary collapse
-
.create(executor = nil, name: nil, priority: nil) ⇒ BasicObject
Creates a new scheduler and starts it in a background ractor, thread, or custom executor.
-
.current ⇒ Scheduler?
Returns a Farce scheduler for the current thread's installed fiber scheduler.
Instance Method Summary collapse
-
#alive? ⇒ Boolean
Whether the dispatcher is in its
:runningstate. -
#close ⇒ Scheduler
Requests shutdown and rejects further submissions.
-
#closed? ⇒ Boolean
Whether new submissions are rejected, including while closing or after a dispatch error.
-
#error ⇒ Exception?
The error encountered by the scheduler, if any.
-
#external_scheduler ⇒ Object?
Retrieves the constructed external scheduler from its owning Ractor.
-
#initialize(capacity: 1024, backend: Internal::FROZEN_CONFIG.io_backend, queue: nil, external: false, pool_worker: nil, queue_owner: true, &constructor) ⇒ Scheduler
constructor
Creates a scheduler without starting a worker.
-
#launch(executor = nil, name: nil, priority: nil) ⇒ Ractor, ...
Starts dispatching tasks in a new worker and waits for setup to leave its launching/setup states before returning.
-
#launch_ractor(name: nil, priority: nil) ⇒ Ractor
Launches a Ractor worker.
-
#launch_thread(name: nil, priority: nil) ⇒ Thread
Launches a Thread worker.
-
#local?(wait: true) ⇒ Boolean
Whether the scheduler is local for the current Ractor.
-
#owner ⇒ Ractor?
The owning Ractor, or nil before ownership is assigned.
-
#owner=(value) ⇒ BasicObject
Assigns ownership during setup.
-
#schedule(*args, mode: :copy, auto_local: true) {|*args| ... } ⇒ Scheduler
Enqueues a task, waiting for queue space if necessary.
-
#state ⇒ Symbol
The current state of the scheduler.
-
#wraps_external? ⇒ Boolean
Whether a fiber scheduler constructor was supplied.
Methods included from Shareable
Methods inherited from Abstract::Scheduler
Methods included from Internal::Noncopyable
Constructor Details
#initialize(capacity: 1024, backend: Farce.config.io_backend) ⇒ Scheduler #initialize(capacity: 1024, &block) ⇒ Scheduler
173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 |
# File 'lib/farce/scheduler.rb', line 173 def initialize(capacity: 1024, backend: Internal::FROZEN_CONFIG.io_backend, queue: nil, external: false, pool_worker: nil, queue_owner: true, &constructor) constructor ||= Internal::FROZEN_CONFIG.fiber_scheduler_constructor @capacity = capacity ? Integer(capacity) : nil @backend = backend @constructor = Ractor.shareable?(constructor) ? constructor : Ractor.shareable_proc(&constructor) if constructor @external = external || !constructor.nil? @owner = Internal::Atom.new(nil) @state = Internal::Atom.new(:initialized) @error = Internal::Atom.new @queue = queue || Internal::Queue.new(capacity: @capacity) @pool_worker = pool_worker @queue_owner = queue_owner super() end |
Class Method Details
.create(executor = nil, name: nil, priority: nil, capacity: 1024, backend: Farce.config.io_backend) ⇒ Scheduler .create(executor = nil, name: nil, priority: nil, capacity: 1024, &block) ⇒ Scheduler
Creates a new scheduler and starts it in a background ractor, thread, or custom executor.
138 139 140 |
# File 'lib/farce/scheduler.rb', line 138 def self.create(executor = nil, name: nil, priority: nil, **, &) new(**, &).tap { it.launch(executor, name:, priority:) } end |
.current ⇒ Scheduler?
Returns a Farce scheduler for the current thread's installed fiber scheduler.
Returns nil when no fiber scheduler is installed.
Creates and installs a wrapper when necessary, starting its dispatcher on
that scheduler. Returns nil when no fiber scheduler is installed.
Additional keyword arguments are forwarded to #initialize when wrapping.
74 75 76 77 78 79 80 81 82 |
# File 'lib/farce/scheduler.rb', line 74 def self.current(**) return unless scheduler = Fiber.scheduler register = Internal::Storage.store_if_absent(self) { Internal::Storage.new } register.store_if_absent(scheduler) do wrapper = new(**, external: true) wrapper.__send__(:set_scheduler!, scheduler:) wrapper end end |
Instance Method Details
#alive? ⇒ Boolean
Returns Whether the dispatcher is in its :running state.
402 |
# File 'lib/farce/scheduler.rb', line 402 def alive? = @state.value == :running |
#close ⇒ Scheduler
Requests shutdown and rejects further submissions. Repeated calls are harmless. The dispatcher drains queued tasks before closing the queue. This method does not wait for dispatch or running tasks to finish.
This wakes a dispatcher blocked on an empty queue. It does not directly close the underlying fiber scheduler.
371 372 373 374 375 376 377 378 379 380 |
# File 'lib/farce/scheduler.rb', line 371 def close until closed? state = @state.value @state.compare_and_set(state, :closing) unless CLOSED_STATES.include?(state) end # if the scheduler never acquired an owner, #schedule might be blocked @owner.compare_and_set(nil, Ractor.current) wake_dispatcher self end |
#closed? ⇒ Boolean
Returns Whether new submissions are rejected, including while closing or after a dispatch error. Does not imply all tasks have finished.
398 |
# File 'lib/farce/scheduler.rb', line 398 def closed? = CLOSED_STATES.include?(@state.value) |
#error ⇒ Exception?
Returns The error encountered by the scheduler, if any.
190 191 192 193 |
# File 'lib/farce/scheduler.rb', line 190 def error value = @error.value value.is_a?(Envelope) ? value.value : value end |
#external_scheduler ⇒ Object?
Retrieves the constructed external scheduler from its owning Ractor.
307 308 309 310 311 |
# File 'lib/farce/scheduler.rb', line 307 def external_scheduler return unless wraps_external? return Internal::Storage[self] if owner.nil? || owner == Ractor.current raise Ractor::IsolationError, "external_scheduler is owned by a different Ractor" end |
#launch(executor = nil, name: nil, priority: nil) ⇒ Ractor, ...
Starts dispatching tasks in a new worker and waits for setup to leave its launching/setup states before returning. Returns the created worker instance.
239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 |
# File 'lib/farce/scheduler.rb', line 239 def launch(executor = nil, name: nil, priority: nil) raise "Scheduler may not be reused" unless @state.compare_and_set(:initialized, :launching) executor ||= Internal.native_ractors? ? Ractor : Thread priority &&= Integer(priority) name &&= -String(name) callback = Ractor.shareable_proc(self: self) { _1.__send__(:launch!, _3, _2) } return executor.new(self, priority, name, &callback) unless executor <= Ractor # Requiring an autoloaded implementation from a new Ractor can deadlock # while its launcher waits for setup to finish. Internal.const_get(:FiberScheduler, false) Internal.const_get(:Storage, false) Internal.const_get(:ThreadPool, false) executor.new(self, priority, name, name:, &callback) rescue Exception => e # rubocop:disable Lint/RescueException @state.value = :error @error.value = Envelope.new(e) raise ensure # if we return the thread too early, someone could kill it while another threads blocks on a schedule call if System.windows? && executor <= Ractor sleep 0.001 while @state.value == :launching sleep 0.001 while @state.value == :setup else @state.wait_until_changed(:launching) @state.wait_until_changed(:setup) end end |
#launch_ractor(name: nil, priority: nil) ⇒ Ractor
200 |
# File 'lib/farce/scheduler.rb', line 200 def launch_ractor(...) = launch(Ractor, ...) |
#launch_thread(name: nil, priority: nil) ⇒ Thread
207 |
# File 'lib/farce/scheduler.rb', line 207 def launch_thread(...) = launch(Thread, ...) |
#local?(wait: true) ⇒ Boolean
Returns Whether the scheduler is local for the current Ractor.
359 360 361 362 |
# File 'lib/farce/scheduler.rb', line 359 def local?(wait: true) owner = wait ? @owner.wait_until_non_nil : @owner.value owner == Ractor.current end |
#owner ⇒ Ractor?
Returns The owning Ractor, or nil before ownership is assigned.
270 |
# File 'lib/farce/scheduler.rb', line 270 def owner = @owner.value |
#owner=(value) ⇒ BasicObject
Assigns ownership during setup. An existing owner cannot be changed. This is used by schedule to determine if a task is scheduled locally or remotely.
If #schedule is called without setting auto_local to false, it will block until the owner can be determined.
Setting the owner explicitly avoids blocking on ownership determination.
292 293 294 295 296 297 |
# File 'lib/farce/scheduler.rb', line 292 def owner=(value) raise ArgumentError, "owner must be a Ractor" unless value.nil? || value.is_a?(Ractor) return if @owner.compare_and_set(nil, value) return if owner == value raise ArgumentError, "owner has already been set" end |
#schedule(*args, mode: :copy, auto_local: true) {|*args| ... } ⇒ Scheduler
Enqueues a task, waiting for queue space if necessary. Returns after submission. It does not wait for task completion or return the task's value.
With auto_local enabled, waits for ownership to be assigned and uses :local
mode whenever the caller and owner are in the same Ractor, including different
threads in that Ractor.
Local tasks retain their block and arguments.
Otherwise, the block must be convertible to a Ractor-shareable proc. Pass task data as arguments so the selected transfer mode can be applied to it.
Valid modes are:
:copy- The value will be copied between Ractors. This is the default mode.:make_shareable- The value will be made Ractor-shareable using Ractor.make_shareable.:move- The value will be moved between Ractors. This saves memory compared to copying, and supports values that can't be copied but moved (like IO objects). However, the value will no longer be accessible on the Ractor that pushed it.:mutable- A Mutable instance will be created for the value. This isn't done recursively and thus will fail for nested unshareable values.:local- The value will be kept local to the Ractor that pushed it. Another ractor trying to receive it will get an error. Useful for usage contained within a single Ractor.:proxy- The value will be wrapped in aFarce::Proxythat executes calls in the original Ractor.:raise- An error will be raised if the value is not Ractor-shareable. Useful for enforcing shareability.:dedup- The value will be deduplicated using Farce.dedup, then made Ractor-shareable. This may update and freeze the original. Already-shareable values pass through unchanged.:shareable_copy- The value will be copied and the copy will be made Ractor-shareable.
336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 |
# File 'lib/farce/scheduler.rb', line 336 def schedule(*args, mode: :copy, auto_local: true, &block) raise SchedulerClosedError, "cannot schedule task on a closed scheduler" if closed? if auto_local owner = @owner.wait_until_non_nil raise SchedulerClosedError, "cannot schedule task on a closed scheduler" if closed? mode = :local if owner == Ractor.current elsif mode == :local && (owner = @owner.value) && owner != Ractor.current raise Ractor::IsolationError, "cannot schedule local task outside the owning Ractor" end task = mode == :local ? Envelope::Local.new(args.any? ? -> { block.call(*args) } : block) : Internal::ScheduledTask.new(args, block, mode) @queue.push(task) self rescue ClosedQueueError raise unless closed? raise SchedulerClosedError, "cannot schedule task on a closed scheduler" end |
#state ⇒ Symbol
Returns The current state of the scheduler. One of:
:initialized: The scheduler has been created but not yet launched.:launching: The scheduler is in the process of being launched on a background worker.:setup: The scheduler is in the process of being setting up a fiber scheduler.:running: The scheduler is actively dispatching tasks.:closing: The scheduler has been requested to shut down and is draining its queue.:closed: The scheduler has finished shutting down.:error: The scheduler encountered an error during dispatch.
391 |
# File 'lib/farce/scheduler.rb', line 391 def state = @state.value |
#wraps_external? ⇒ Boolean
Returns Whether a fiber scheduler constructor was supplied.
300 |
# File 'lib/farce/scheduler.rb', line 300 def wraps_external? = @external |