Class: Farce::TimerQueue

Inherits:
Abstract::TimerQueue
  • Object
show all
Includes:
Shareable::Unfreezable, Internal::ManagedQueue
Defined in:
lib/farce/timer_queue.rb

Overview

A shareable timer queue with transfer modes for unshareable values.

Instance Method Summary collapse

Methods included from Shareable

#ractor_shareable?

Constructor Details

#initialize(capacity: nil, mode: :copy, track_age: false) ⇒ BasicObject

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:

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

    the maximum number of values, or nil for an unbounded queue

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

    the transfer mode for unshareable values



14
15
16
17
# File 'lib/farce/timer_queue.rb', line 14

def initialize(capacity: nil, mode: :copy, track_age: false)
  @manager = ModeManager.new(mode:)
  super(capacity:, track_age:)
end

Instance Method Details

#peek { ... } ⇒ BasicObject?

Note:

Peeking at a moved value claims it for this Ractor while leaving it queued.

Return the earliest value without removing it, whether or not it is due.

Yields:

  • called when the queue is empty

Returns:

  • (BasicObject, nil) —

    the value or the fallback result



59
60
61
62
63
64
# File 'lib/farce/timer_queue.rb', line 59

def peek
  empty = false
  value = @queue.peek { empty = true }
  return @manager.unwrap(value) unless empty
  yield if block_given?
end

#pop(non_block = false, timeout: nil) { ... } ⇒ BasicObject?

Remove the earliest value, waiting until its timestamp is reached.

Parameters:

  • timeout (Numeric, nil) (defaults to: nil) —

    maximum number of seconds to wait

Yields:

  • called when the timeout expires first

Returns:

  • (BasicObject, nil) —

    the value or the fallback result



37
38
39
40
41
42
43
44
45
46
47
# File 'lib/farce/timer_queue.rb', line 37

def pop(non_block = false, timeout: nil) # rubocop:disable Style/OptionalBooleanParameter
  return try_pop { raise ThreadError, "queue empty" } if non_block
  deadline = timeout_at(timeout) unless timeout.nil?

  while true
    empty = false
    value = @queue.pop_before(Clock.now) { empty = true }
    return @manager.unwrap(value) unless empty
    return block_given? ? yield : nil if UNDEFINED.equal?(wait_for_timestamp(deadline))
  end
end

#push(value, non_block = false, mode: nil, timeout: nil, **time_options) ⇒ Boolean

Add a value to become available at the given time. 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:

  • value (BasicObject) —

    the value to add

  • non_block (Boolean) (defaults to: false) —

    whether to raise an exception when the queue is at capacity

  • timeout (Numeric, nil) (defaults to: nil) —

    maximum number of seconds to wait for capacity

  • time_options (Hash{Symbol => Object}) —

    a scheduling option accepted by Clock.parse: at:, time:, timeout_at:, delay:, in:, offset:, wait:, or clock:. With no scheduling option, the value is available immediately. Since timeout: controls the capacity wait, use delay: or wait: to schedule a relative offset.

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

    the transfer mode for unshareable values

Returns:

  • (Boolean) —

    whether the value was added

Raises:

  • (ThreadError) —

    when the queue is at capacity and non_block is true



22
23
24
25
# File 'lib/farce/timer_queue.rb', line 22

def push(value, non_block = false, mode: nil, timeout: nil, **time_options) # rubocop:disable Style/OptionalBooleanParameter
  at = Clock.parse(time_options)
  push_to_storage(at, non_block, @manager.wrap(value, mode:), timeout:)
end

#try_pop { ... } ⇒ BasicObject?

Try to remove the earliest value only if its timestamp has been reached.

Yields:

  • called when no value is ready

Returns:

  • (BasicObject, nil) —

    the value or the fallback result



50
51
52
53
54
55
# File 'lib/farce/timer_queue.rb', line 50

def try_pop
  empty = false
  value = @queue.pop_before(Clock.now) { empty = true }
  return @manager.unwrap(value) unless empty
  yield if block_given?
end

#try_push(value, mode: nil, **time_options) { ... } ⇒ Boolean, BasicObject

Try to add a value without waiting for capacity. 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:

  • value (BasicObject) —

    the value to add

  • time_options (Hash{Symbol => Object}) —

    a scheduling option accepted by Clock.parse: at:, time:, timeout_at:, delay:, in:, offset:, timeout:, wait:, or clock:. With no scheduling option, the value is available immediately.

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

    the transfer mode for unshareable values

Yields:

  • called when the queue is at capacity

Returns:

  • (Boolean, BasicObject) —

    true, or the fallback result when full



30
31
32
33
34
# File 'lib/farce/timer_queue.rb', line 30

def try_push(value, mode: nil, **time_options)
  at = Clock.parse(time_options)
  return true if @queue.push(at, @manager.wrap(value, mode:))
  block_given? ? yield : false
end