Class: Farce::Scheduler

Inherits:
Abstract::Scheduler show all
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.

Examples:

Submit arguments to a worker

# Runs a scheduler in a background Ractor/Thread
scheduler = Farce::Scheduler.create
scheduler.schedule("hello") { |message| puts message }

Using Farce::Scheduler as a Fiber scheduler

scheduler = Farce::Scheduler.new
Fiber.set_scheduler(scheduler)
Fiber.schedule { puts "Hello from the scheduler!" }

Class Method Summary collapse

Instance Method Summary collapse

Methods included from Shareable

#ractor_shareable?

Methods inherited from Abstract::Scheduler

#execute, #ractor_safe?

Methods included from Internal::Noncopyable

#duplicable?

Constructor Details

#initialize(capacity: 1024, backend: Farce.config.io_backend) ⇒ Scheduler #initialize(capacity: 1024, &block) ⇒ Scheduler

Creates a scheduler without starting a worker. Use this instead of create when you want to launch it later with #launch, or install it on the current thread with Fiber.set_scheduler(scheduler).

Overloads:

  • #initialize(capacity: 1024, backend: Farce.config.io_backend) ⇒ Scheduler

    Uses the configured fiber scheduler when launched or installed.

    Parameters:

    • capacity (Integer, nil) (defaults to: 1024) —

      Pending-task capacity, nil creates an unbounded queue.

    • backend (Symbol) (defaults to: Farce.config.io_backend) —

      IO backend for the built-in scheduler, defaults to Farce configuration.

  • #initialize(capacity: 1024, &block) ⇒ Scheduler

    Uses another library's fiber scheduler while retaining Farce's ability to accept tasks from other threads and Ractors.

    Supply a block that constructs and returns that scheduler. The block runs once when this handle is launched or installed, in the thread that will run the tasks, rather than during new. It must be Ractor-shareable. Farce installs its result. Submit tasks separately with #schedule. Configure the supplied scheduler in the block. Farce's backend option is ignored for this overload.

    Examples:

    Choose Async, then launch the worker later

    require "async"
    require "farce"
    
    scheduler = Farce::Scheduler.new { Async::Scheduler.new }
    worker = scheduler.launch_thread
    scheduler.schedule { puts "Running on Async" }

    Parameters:

    • capacity (Integer, nil) (defaults to: 1024) —

      Pending-task capacity, nil creates an unbounded queue.

    Yield Returns:

    • (Object) —

      The fiber scheduler to install when launched or installed.



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.

Overloads:

  • .create(executor = nil, name: nil, priority: nil, capacity: 1024, backend: Farce.config.io_backend) ⇒ Scheduler

    Runs tasks on the configured fiber scheduler. The backend option selects the IO driver when using Farce's built-in scheduler.

    Examples:

    scheduler = Farce::Scheduler.create
    scheduler.schedule { puts "Running on the built-in scheduler" }

    Parameters:

    • executor (Class, nil) (defaults to: nil) —

      Ractor or Thread, nil selects the runtime default.

    • name (String, nil) (defaults to: nil) —

      Optional worker name.

    • priority (Integer, nil) (defaults to: nil) —

      Optional worker thread priority.

    • capacity (Integer, nil) (defaults to: 1024) —

      Pending-task capacity, nil creates an unbounded queue.

    • backend (Symbol) (defaults to: Farce.config.io_backend) —

      IO backend for the built-in scheduler, defaults to Farce configuration.

    Returns:

  • .create(executor = nil, name: nil, priority: nil, capacity: 1024, &block) ⇒ Scheduler

    Runs tasks on a fiber scheduler supplied by your block. Use this to send work through Farce to another library's event loop, such as Async.

    Construct and return the scheduler inside the block: Farce calls it once in the worker, where that event loop will run, and installs the result. The block configures the scheduler. Submit actual tasks with #schedule. Configure any IO backend on the returned scheduler itself. Farce's backend option is ignored when a block is supplied.

    The block must be Ractor-shareable. Require the library before calling create and avoid capturing a scheduler created on another thread or Ractor.

    Examples:

    Run tasks on Async in a worker thread

    require "async"
    require "farce"
    
    # Using Thread as an executor as at the moment Async cannot run on a non-main Ractor
    scheduler = Farce::Scheduler.create(Thread) { Async::Scheduler.new }
    scheduler.schedule { puts "Running on Async" }

    Run tasks on CarbonFiber in a worker ractor

    require "carbon_fiber"
    require "farce"
    
    scheduler = Farce::Scheduler.create(Ractor) { CarbonFiber::Scheduler.new }
    scheduler.schedule { puts "Running on CarbonFiber" }

    Parameters:

    • executor (Class, nil) (defaults to: nil) —

      Ractor or Thread. nil selects the runtime default. Choose an executor supported by the supplied scheduler and its dependencies.

    • name (String, nil) (defaults to: nil) —

      Optional worker name.

    • priority (Integer, nil) (defaults to: nil) —

      Optional worker thread priority.

    • capacity (Integer, nil) (defaults to: 1024) —

      Pending-task capacity, nil creates an unbounded queue.

    Yield Returns:

    • (Object) —

      The fiber scheduler to install in the worker.

    Returns:

See Also:



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.

Examples:

Create a scheduler for Async, and send it to a different Ractor

require "async"
require "farce"

# No fiber scheduler set here
Farce::Scheduler.current # => nil

Async do
  # A scheduler that will hand tasks to Async::Scheduler
  scheduler = Farce::Scheduler.current # => #<Farce::Scheduler>

  # This looks like doing the same as nested Async { ... } calls or
  # Fiber.schedule { ... }, …
  scheduler.schedule { puts "Scheduled from the main Ractor" }

  # … but it can be called from other threads and ractors
  Ractor.new(scheduler) do |scheduler|
    scheduler.schedule { puts "Scheduled from a background Ractor" }
  end
end

Full circle

scheduler = Farce::Scheduler.new
Farce::Scheduler.current # => nil
Fiber.set_scheduler(scheduler)
Farce::Scheduler.current == scheduler # => true

Returns:



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.

Returns:

  • (Boolean) —

    Whether the dispatcher is in its :running state.

See Also:



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.

Returns:



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.

Returns:

  • (Boolean) —

    Whether new submissions are rejected, including while closing or after a dispatch error. Does not imply all tasks have finished.

See Also:



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.

Returns:

  • (Exception, nil) —

    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.

Returns:

  • (Object, nil) —

    The external scheduler, or nil before setup or when using the built-in scheduler.

Raises:



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.

Examples:

Using the return value to assign a ThreadGroup

scheduler = Farce::Scheduler.new
worker    = scheduler.launch_thread
group     = ThreadGroup.new

group.add(worker)
group.enclose

Using a custom worker

class MyWorker
  def initialize(...)
    # flip a coin whether we're using Ractor or Thread
    @worker = [Ractor, Thread].sample.new(...)
  end

  def join = @worker.join
end

scheduler = Farce::Scheduler.new
scheduler.launch(MyWorker)

Parameters:

  • executor (Class, nil) (defaults to: nil) —

    Ractor or Thread. Defaults to Ractor when native Ractors are available, otherwise Thread.

  • name (String, nil) (defaults to: nil) —

    Optional worker name.

  • priority (Integer, nil) (defaults to: nil) —

    Optional worker thread priority.

Returns:

  • (Ractor, Thread, Object) —

    The worker.

Raises:

  • (RuntimeError) —

    If this handle has already been launched, installed, or closed.



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

Launches a Ractor worker. Accepts the keyword arguments of #launch.

Parameters:

  • name (String, nil) (defaults to: nil) —

    Optional worker name.

  • priority (Integer, nil) (defaults to: nil) —

    Optional worker thread priority.

Returns:



200
# File 'lib/farce/scheduler.rb', line 200

def launch_ractor(...) = launch(Ractor, ...)

#launch_thread(name: nil, priority: nil) ⇒ Thread

Launches a Thread worker. Accepts the keyword arguments of #launch.

Parameters:

  • name (String, nil) (defaults to: nil) —

    Optional worker name.

  • priority (Integer, nil) (defaults to: nil) —

    Optional worker thread priority.

Returns:

  • (Thread) —

    The worker.



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.

Parameters:

  • wait (Boolean) (defaults to: true) —

    Whether to wait for the owner to be non-nil before checking if the scheduler is local.

Returns:

  • (Boolean) —

    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.

Returns:

  • (Ractor, nil) —

    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.

Examples:

scheduler = Farce::Scheduler.new

# Scheduler hasn't been launched, so it could block.
scheduler.owner = Ractor.current

# Now this doesn't block
scheduler.schedule { some_task }

# Can use #launch_thread, but #launch_ractor would fail now that the owner is set.
scheduler.launch_thread

Parameters:

Raises:

  • (ArgumentError) —

    If value is not a Ractor/nil or changes an assigned owner.



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 a Farce::Proxy that 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.

Parameters:

  • args (Array<Object>) —

    Positional arguments passed to the task block.

  • mode (Symbol) (defaults to: :copy) —

    Argument transfer mode: :copy, :move, :local, :make_shareable, :shareable_copy, or :raise. See ModeManager#wrap.

  • auto_local (Boolean) (defaults to: true) —

    Whether to override mode with :local in the owning Ractor.

Yields:

  • (*args) —

    The task to execute.

Returns:

Raises:

See Also:



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.

Returns:

  • (Symbol) —

    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.

Returns:

  • (Boolean) —

    Whether a fiber scheduler constructor was supplied.



300
# File 'lib/farce/scheduler.rb', line 300

def wraps_external? = @external