Class: Farce::Pool

Inherits:
Abstract::Scheduler
  • Object
show all
Includes:
Shareable::Unfreezable, Internal::Inspect, Internal::MarshalSupport::Reject
Defined in:
lib/farce/pool.rb

Overview

A dynamically sized pool of Ractor schedulers.

Tasks wait in one shared queue. An empty pool starts its first worker immediately when work arrives. A task that remains queued for grow_after causes the pool to add a worker until max_size is reached. Workers above min_size retire after shrink_after without queued or running work. By default, the pool starts with no workers and can shrink back to zero.

Most importantly, it exposes a #schedule method compatible with Scheduler#schedule.

Examples:

# Create a new pool
pool = Farce::Pool.new

# Schedule some work on the pool.
pool.schedule do
  # The work runs in a non-blocking fiber, multiple tasks can run concurrently even on the same worker.
  loop do
    sleep 5
    MyClass.recurring_work
  end
end

Custom Fiber scheduler

require "farce"
require "carbon_fiber"

# At least two Ractors, up to four. Launch a new one if a task waits longer than 100 milliseconds.
# Use the fiber scheduler from the carbon_fiber gem.
pool = Farce::Pool.new(min_size: 2, max_size: 4, grow_after: 0.1) { CarbonFiber::Scheduler.new }

# Schedule some work
pool.schedule { MyClass.do_something }

Instance Attribute Summary collapse

Instance Method Summary collapse

Methods included from Shareable

#ractor_shareable?

Constructor Details

#initialize(min_size: 0, max_size: 4, max_inflight: 64, grow_after: 0.005, shrink_after: 30, capacity: 1024, backend: Internal::FROZEN_CONFIG.io_backend, &constructor) ⇒ Pool

Create and start an elastic pool using the configured fiber scheduler. An explicit constructor block overrides the configured default.

max_inflight is a soft limit. A worker may exceed it when its scheduler has no runnable fibers or when the pool has reached max_size.

Parameters:

  • min_size (Integer) (defaults to: 0) —

    Workers started immediately and retained while open.

  • max_size (Integer) (defaults to: 4) —

    Maximum workers allowed.

  • max_inflight (Integer, nil) (defaults to: 64) —

    Soft task limit for each worker.

  • grow_after (Numeric) (defaults to: 0.005) —

    Queue wait before adding another worker.

  • shrink_after (Numeric, nil) (defaults to: 30) —

    Idle time before an extra worker retires.

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

    Pending-task capacity.

  • backend (Symbol) (defaults to: Internal::FROZEN_CONFIG.io_backend) —

    IO backend for the built-in scheduler.

Yield Returns:

  • (Object) —

    Fiber scheduler constructed in each worker.



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
# File 'lib/farce/pool.rb', line 73

def initialize(min_size: 0, max_size: 4, max_inflight: 64,
               grow_after: 0.005, shrink_after: 30, capacity: 1024,
               backend: Internal::FROZEN_CONFIG.io_backend, &constructor)
  constructor ||= Internal::FROZEN_CONFIG.fiber_scheduler_constructor
  @min_size     = Integer(min_size)
  @max_size     = Integer(max_size)
  @max_inflight = max_inflight && Integer(max_inflight)
  @grow_after   = duration(grow_after, :grow_after)
  @shrink_after = duration(shrink_after, :shrink_after, nil_allowed: true)
  validate_sizes

  @backend      = backend
  @constructor  = constructor
  @constructor  = Ractor.shareable_proc(&constructor) if constructor && !Ractor.shareable?(constructor)
  @queue        = Internal::Queue.new(capacity: capacity && Integer(capacity), track_age: true)
  @state        = Internal::Atom.new(:running)
  @error        = Internal::Atom.new
  @worker_count = Internal::Atom.new(0)
  @pressure     = Internal::Atom.new
  @admission    = Internal::Signal.new
  super()
  begin
    @min_size.times { start_worker }
  rescue Exception # rubocop:disable Lint/RescueException
    close
    raise
  end
end

Instance Attribute Details

#grow_after ⇒ Float (readonly)

Returns Maximum queue wait before adding a worker.

Returns:

  • (Float) —

    Maximum queue wait before adding a worker.



45
46
47
# File 'lib/farce/pool.rb', line 45

def grow_after
  @grow_after
end

#max_inflight ⇒ Integer? (readonly)

Returns Soft task limit for each worker.

Returns:

  • (Integer, nil) —

    Soft task limit for each worker.



48
49
50
# File 'lib/farce/pool.rb', line 48

def max_inflight
  @max_inflight
end

#max_size ⇒ Integer (readonly)

Returns Maximum worker count.

Returns:

  • (Integer) —

    Maximum worker count.



51
52
53
# File 'lib/farce/pool.rb', line 51

def max_size
  @max_size
end

#min_size ⇒ Integer (readonly)

Returns Minimum worker count.

Returns:

  • (Integer) —

    Minimum worker count.



54
55
56
# File 'lib/farce/pool.rb', line 54

def min_size
  @min_size
end

#shrink_after ⇒ Float? (readonly)

Returns Idle time before an extra worker retires.

Returns:

  • (Float, nil) —

    Idle time before an extra worker retires.



57
58
59
# File 'lib/farce/pool.rb', line 57

def shrink_after
  @shrink_after
end

Instance Method Details

#at_max_size? ⇒ Boolean

Return whether all configured workers have been started.

Returns:

  • (Boolean)


164
# File 'lib/farce/pool.rb', line 164

def at_max_size? = size >= @max_size

#close ⇒ Pool

Stop accepting tasks and drain queued work.

Returns:



130
131
132
133
134
135
136
137
138
139
# File 'lib/farce/pool.rb', line 130

def close
  return self if closed?
  @state.compare_and_set(:running, :closing)
  @queue.seal
  admission_changed
  if size.zero?
    @queue.closed? ? finish_close : start_worker(drain: true)
  end
  self
end

#closed? ⇒ Boolean

Return whether the pool rejects new tasks.

Returns:

  • (Boolean)


154
# File 'lib/farce/pool.rb', line 154

def closed? = state != :running

#closing? ⇒ Boolean

Return whether the pool is draining queued work.

Returns:

  • (Boolean)


157
# File 'lib/farce/pool.rb', line 157

def closing? = state == :closing

#error ⇒ BasicObject

Return the first worker error, if any.



142
143
144
145
# File 'lib/farce/pool.rb', line 142

def error
  value = @error.value
  value.is_a?(Envelope) ? value.value : value
end

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

Enqueue a task for any worker in the pool.

Pools have no single owner, so auto_local does not select local transfer. Explicit local transfer is rejected because another Ractor may run the task.

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.

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

    Accepted for compatibility with Scheduler#schedule.

Yields:

  • (*args) —

    Task to execute.

Returns:

Raises:



114
115
116
117
118
119
120
121
122
123
124
125
# File 'lib/farce/pool.rb', line 114

def schedule(*args, mode: :copy, auto_local: true, &block) # rubocop:disable Lint/UnusedMethodArgument
  raise PoolClosedError, "cannot schedule task on a closed pool" if closed?
  raise ArgumentError, "local transfer is not supported by a pool" if mode == :local

  task = Internal::ScheduledTask.new(args, block, mode)
  @queue.push(task)
  arm_scaler
  self
rescue ClosedQueueError
  raise unless closed?
  raise PoolClosedError, "cannot schedule task on a closed pool"
end

#size ⇒ BasicObject

Return the current number of running and starting workers.



148
# File 'lib/farce/pool.rb', line 148

def size = @worker_count.value

#state ⇒ BasicObject

Return the current pool state.



151
# File 'lib/farce/pool.rb', line 151

def state = @state.value