Class: Farce::Abstract::Scheduler Abstract

Inherits:
Object
  • Object
show all
Includes:
Internal::Noncopyable
Defined in:
lib/farce/abstract/scheduler.rb

Overview

This class is abstract.

A shared super class for all scheduler implementations.

Direct Known Subclasses

Pool, Scheduler, ThreadScheduler

Instance Method Summary collapse

Methods included from Internal::Noncopyable

#duplicable?

Instance Method Details

#close ⇒ self

This method is abstract.

This method must be implemented by concrete scheduler subclasses.

Closes the scheduler, releasing any resources it holds. Once closed, the scheduler cannot be used to schedule new tasks. Will drain any pending tasks before fully closing. If tasks are running indefinitely, it may block.

Returns:

  • (self)


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
# File 'lib/farce/abstract/scheduler.rb', line 66

class Scheduler
  include Internal::Noncopyable

  # An error that occurred during the scheduler's operation, if any.
  # This is primarily for scheduling errors, not execution errors.
  # You can check this if the {#state} is `:error`.
  def error = nil

  # @param wait [Boolean] Whether to wait for the owner to be non-nil before checking if the scheduler is local.
  # @return [Boolean] Whether the scheduler is local for the current Ractor.
  def local?(wait: true) = false # rubocop:disable Lint/UnusedMethodArgument

  # Like {#schedule}, with two differences:
  # 1. If {#local?} returns true, and it would run in local mode (either due to `auto_local` or explicitly set
  #    `mode: :local`), it executes the block directly without enqueuing it.
  # 2. It blocks until the block has been executed, either immediately in local mode or after being scheduled.
  #
  # @param (see #schedule)
  # @return [self]
  def execute(*, mode: :copy, auto_local: true, &callback)
    raise LocalJumpError, "Cannot yield without a block" unless block_given?
    raise SchedulerClosedError, "cannot execute task on a closed scheduler" if closed?

    if (auto_local || mode == :local) && local?
      raise SchedulerClosedError, "cannot execute task on a closed scheduler" if closed?
      yield(*)
      return self
    end

    ran      = Farce::Flag.new(false)
    signal   = Signal.new
    callback = Ractor.shareable_proc(&callback) unless Ractor.shareable?(callback)

    schedule(callback, ran, signal, *, mode:, auto_local:) do |callback, ran, signal, *args|
      callback.call(*args)
    ensure
      ran.set
      signal.broadcast
    end

    signal.wait_until { ran.value }
    self
  end
end

#closed? ⇒ Boolean

This method is abstract.

This method must be implemented by concrete scheduler subclasses.

Returns Whether the scheduler is closed.

Returns:

  • (Boolean) —

    Whether the scheduler is closed.



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
# File 'lib/farce/abstract/scheduler.rb', line 66

class Scheduler
  include Internal::Noncopyable

  # An error that occurred during the scheduler's operation, if any.
  # This is primarily for scheduling errors, not execution errors.
  # You can check this if the {#state} is `:error`.
  def error = nil

  # @param wait [Boolean] Whether to wait for the owner to be non-nil before checking if the scheduler is local.
  # @return [Boolean] Whether the scheduler is local for the current Ractor.
  def local?(wait: true) = false # rubocop:disable Lint/UnusedMethodArgument

  # Like {#schedule}, with two differences:
  # 1. If {#local?} returns true, and it would run in local mode (either due to `auto_local` or explicitly set
  #    `mode: :local`), it executes the block directly without enqueuing it.
  # 2. It blocks until the block has been executed, either immediately in local mode or after being scheduled.
  #
  # @param (see #schedule)
  # @return [self]
  def execute(*, mode: :copy, auto_local: true, &callback)
    raise LocalJumpError, "Cannot yield without a block" unless block_given?
    raise SchedulerClosedError, "cannot execute task on a closed scheduler" if closed?

    if (auto_local || mode == :local) && local?
      raise SchedulerClosedError, "cannot execute task on a closed scheduler" if closed?
      yield(*)
      return self
    end

    ran      = Farce::Flag.new(false)
    signal   = Signal.new
    callback = Ractor.shareable_proc(&callback) unless Ractor.shareable?(callback)

    schedule(callback, ran, signal, *, mode:, auto_local:) do |callback, ran, signal, *args|
      callback.call(*args)
    ensure
      ran.set
      signal.broadcast
    end

    signal.wait_until { ran.value }
    self
  end
end

#error ⇒ BasicObject

An error that occurred during the scheduler's operation, if any. This is primarily for scheduling errors, not execution errors. You can check this if the #state is :error.



72
# File 'lib/farce/abstract/scheduler.rb', line 72

def error = nil

#execute(mode: :copy, auto_local: true, &callback) ⇒ self

Like #schedule, with two differences:

  1. If #local? returns true, and it would run in local mode (either due to auto_local or explicitly set mode: :local), it executes the block directly without enqueuing it.
  2. It blocks until the block has been executed, either immediately in local mode or after being scheduled.

Parameters:

  • args (Array) —

    Arguments to pass to the scheduled block. May be handled according to mode.

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

    The mode in which to schedule the task. Defaults to :copy. May be ignored by some schedulers.

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

    Whether to automatically use a local mode when possible. Defaults to true. May be ignored by some schedulers.

Returns:

  • (self)

Raises:

  • (LocalJumpError)


85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
# File 'lib/farce/abstract/scheduler.rb', line 85

def execute(*, mode: :copy, auto_local: true, &callback)
  raise LocalJumpError, "Cannot yield without a block" unless block_given?
  raise SchedulerClosedError, "cannot execute task on a closed scheduler" if closed?

  if (auto_local || mode == :local) && local?
    raise SchedulerClosedError, "cannot execute task on a closed scheduler" if closed?
    yield(*)
    return self
  end

  ran      = Farce::Flag.new(false)
  signal   = Signal.new
  callback = Ractor.shareable_proc(&callback) unless Ractor.shareable?(callback)

  schedule(callback, ran, signal, *, mode:, auto_local:) do |callback, ran, signal, *args|
    callback.call(*args)
  ensure
    ran.set
    signal.broadcast
  end

  signal.wait_until { ran.value }
  self
end

#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.



76
# File 'lib/farce/abstract/scheduler.rb', line 76

def local?(wait: true) = false # rubocop:disable Lint/UnusedMethodArgument

#ractor_safe? ⇒ Boolean

This method is abstract.

Schedulers should include Shareable or Unshareable, which implement this method.

Returns Whether the scheduler is safe to use across Ractors.

Returns:

  • (Boolean) —

    Whether the scheduler is safe to use across Ractors.



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
# File 'lib/farce/abstract/scheduler.rb', line 66

class Scheduler
  include Internal::Noncopyable

  # An error that occurred during the scheduler's operation, if any.
  # This is primarily for scheduling errors, not execution errors.
  # You can check this if the {#state} is `:error`.
  def error = nil

  # @param wait [Boolean] Whether to wait for the owner to be non-nil before checking if the scheduler is local.
  # @return [Boolean] Whether the scheduler is local for the current Ractor.
  def local?(wait: true) = false # rubocop:disable Lint/UnusedMethodArgument

  # Like {#schedule}, with two differences:
  # 1. If {#local?} returns true, and it would run in local mode (either due to `auto_local` or explicitly set
  #    `mode: :local`), it executes the block directly without enqueuing it.
  # 2. It blocks until the block has been executed, either immediately in local mode or after being scheduled.
  #
  # @param (see #schedule)
  # @return [self]
  def execute(*, mode: :copy, auto_local: true, &callback)
    raise LocalJumpError, "Cannot yield without a block" unless block_given?
    raise SchedulerClosedError, "cannot execute task on a closed scheduler" if closed?

    if (auto_local || mode == :local) && local?
      raise SchedulerClosedError, "cannot execute task on a closed scheduler" if closed?
      yield(*)
      return self
    end

    ran      = Farce::Flag.new(false)
    signal   = Signal.new
    callback = Ractor.shareable_proc(&callback) unless Ractor.shareable?(callback)

    schedule(callback, ran, signal, *, mode:, auto_local:) do |callback, ran, signal, *args|
      callback.call(*args)
    ensure
      ran.set
      signal.broadcast
    end

    signal.wait_until { ran.value }
    self
  end
end

#schedule(*args, mode: :copy, auto_local: true) {|*args| ... } ⇒ self

This method is abstract.

This method must be implemented by concrete scheduler subclasses.

Schedules a block for execution. What actually happens depends on the concrete scheduler implementation. Some schedulers may ignore the auto_local option. Some might treat every mode as :local.

The scheduler may choose to block until the task has been executed, but most will add it to a queue of pending tasks, only blocking when that queue's capacity has been reached.

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) —

    Arguments to pass to the scheduled block. May be handled according to mode.

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

    The mode in which to schedule the task. Defaults to :copy. May be ignored by some schedulers.

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

    Whether to automatically use a local mode when possible. Defaults to true. May be ignored by some schedulers.

Yields:

  • (*args) —

    The block to be executed by the scheduler.

Yield Parameters:

  • args (Array) —

    The arguments passed to the block.

Returns:

  • (self)

Yield Receiver:

  • (BasicObject, nil) —

    Local mode will not touch the block's binding (self stays the same), but other modes may change it to nil.



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
# File 'lib/farce/abstract/scheduler.rb', line 66

class Scheduler
  include Internal::Noncopyable

  # An error that occurred during the scheduler's operation, if any.
  # This is primarily for scheduling errors, not execution errors.
  # You can check this if the {#state} is `:error`.
  def error = nil

  # @param wait [Boolean] Whether to wait for the owner to be non-nil before checking if the scheduler is local.
  # @return [Boolean] Whether the scheduler is local for the current Ractor.
  def local?(wait: true) = false # rubocop:disable Lint/UnusedMethodArgument

  # Like {#schedule}, with two differences:
  # 1. If {#local?} returns true, and it would run in local mode (either due to `auto_local` or explicitly set
  #    `mode: :local`), it executes the block directly without enqueuing it.
  # 2. It blocks until the block has been executed, either immediately in local mode or after being scheduled.
  #
  # @param (see #schedule)
  # @return [self]
  def execute(*, mode: :copy, auto_local: true, &callback)
    raise LocalJumpError, "Cannot yield without a block" unless block_given?
    raise SchedulerClosedError, "cannot execute task on a closed scheduler" if closed?

    if (auto_local || mode == :local) && local?
      raise SchedulerClosedError, "cannot execute task on a closed scheduler" if closed?
      yield(*)
      return self
    end

    ran      = Farce::Flag.new(false)
    signal   = Signal.new
    callback = Ractor.shareable_proc(&callback) unless Ractor.shareable?(callback)

    schedule(callback, ran, signal, *, mode:, auto_local:) do |callback, ran, signal, *args|
      callback.call(*args)
    ensure
      ran.set
      signal.broadcast
    end

    signal.wait_until { ran.value }
    self
  end
end

#state ⇒ Symbol

This method is abstract.

This method must be implemented by concrete scheduler subclasses.

Returns The current state of the scheduler.

Common states include:

  • :initialized, :launching and :setup for the setup phase.
  • :running for the execution phase.
  • :closing, :closed, and :error for the termination phase.

Not all schedulers will implement every state.

Returns:

  • (Symbol) —

    The current state of the scheduler.

    Common states include:

    • :initialized, :launching and :setup for the setup phase.
    • :running for the execution phase.
    • :closing, :closed, and :error for the termination phase.

    Not all schedulers will implement every state.



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
# File 'lib/farce/abstract/scheduler.rb', line 66

class Scheduler
  include Internal::Noncopyable

  # An error that occurred during the scheduler's operation, if any.
  # This is primarily for scheduling errors, not execution errors.
  # You can check this if the {#state} is `:error`.
  def error = nil

  # @param wait [Boolean] Whether to wait for the owner to be non-nil before checking if the scheduler is local.
  # @return [Boolean] Whether the scheduler is local for the current Ractor.
  def local?(wait: true) = false # rubocop:disable Lint/UnusedMethodArgument

  # Like {#schedule}, with two differences:
  # 1. If {#local?} returns true, and it would run in local mode (either due to `auto_local` or explicitly set
  #    `mode: :local`), it executes the block directly without enqueuing it.
  # 2. It blocks until the block has been executed, either immediately in local mode or after being scheduled.
  #
  # @param (see #schedule)
  # @return [self]
  def execute(*, mode: :copy, auto_local: true, &callback)
    raise LocalJumpError, "Cannot yield without a block" unless block_given?
    raise SchedulerClosedError, "cannot execute task on a closed scheduler" if closed?

    if (auto_local || mode == :local) && local?
      raise SchedulerClosedError, "cannot execute task on a closed scheduler" if closed?
      yield(*)
      return self
    end

    ran      = Farce::Flag.new(false)
    signal   = Signal.new
    callback = Ractor.shareable_proc(&callback) unless Ractor.shareable?(callback)

    schedule(callback, ran, signal, *, mode:, auto_local:) do |callback, ran, signal, *args|
      callback.call(*args)
    ensure
      ran.set
      signal.broadcast
    end

    signal.wait_until { ran.value }
    self
  end
end