Class: Farce::PriorityQueue

Inherits:
Abstract::PriorityQueue show all
Includes:
Shareable::Unfreezable, Internal::ManagedQueue
Defined in:
lib/farce/priority_queue.rb

Overview

A shareable priority queue with transfer modes for unshareable values.

Instance Attribute Summary

Attributes inherited from Abstract::PriorityQueue

#default_priority, #order

Attributes inherited from Abstract::Queue

#capacity, #mode, #num_waiting, #size

Instance Method Summary collapse

Methods included from Shareable

#ractor_shareable?

Methods inherited from Abstract::PriorityQueue

#delete, #first_priority, #last_priority, #wait_pop

Methods inherited from Abstract::Queue

#age_tracking?, #clear, #close, #closed?, #deq, #duplicable?, #empty?, #enq, #full?, #generation, #length, #max, #oldest_age, #oldest_enqueued_at, #seal, #sealed?, #wait_pop, #wait_push

Constructor Details

#initialize(capacity: nil, default_priority: 0, order: :ascending, 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

  • default_priority (BasicObject) (defaults to: 0) —

    the priority used when push or try_push is called without an explicit priority

  • order (:ascending, :descending) (defaults to: :ascending) —

    the priority order

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

    the default mode for sending unshareable values



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

def initialize(capacity: nil, default_priority: 0, order: :ascending, mode: :copy, track_age: false)
  @manager = ModeManager.new(mode:)
  super(capacity:, default_priority:, order:, 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 next value without removing it.

Yields:

  • called when the queue is empty

Returns:

  • (BasicObject, nil) —

    the next value or the fallback result



60
61
62
63
64
65
66
67
68
# File 'lib/farce/priority_queue.rb', line 60

def peek
  empty  = false
  result = if @reverse_order
             @queue.peek_last { empty = true }
           else
             @queue.peek { empty = true }
           end
  empty ? (block_given? ? yield : nil) : @manager.unwrap(result)
end

#pop(non_block = false, timeout: nil) ⇒ BasicObject

Remove the oldest value at the lowest or highest priority, depending on Abstract::PriorityQueue#order, waiting when empty.



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

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
    result = @reverse_order ? @queue.pop_last { empty = true } : @queue.pop { empty = true }

    return @manager.unwrap(result) unless empty
    remaining = remaining_timeout(deadline)
    return block_given? ? yield : nil unless wait_pop(timeout: remaining)
  end
end

#push(value, non_block = false, priority: default_priority, timeout: nil, mode: nil) ⇒ Boolean

Add a value, waiting for capacity when the queue is bounded and full. 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

  • priority (BasicObject) (defaults to: default_priority) —

    the value used to order the entry

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

    maximum number of seconds to wait

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

    the default mode for sending 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
# File 'lib/farce/priority_queue.rb', line 22

def push(value, non_block = false, priority: default_priority, timeout: nil, mode: nil) # rubocop:disable Style/OptionalBooleanParameter
  push_to_storage(priority, non_block, @manager.wrap(value, mode:), timeout:)
end

#try_pop { ... } ⇒ BasicObject?

Try to remove the oldest value at the lowest priority without waiting.

Yields:

  • called when the queue is empty

Returns:

  • (BasicObject, nil) —

    the next value or the fallback result



51
52
53
54
55
56
# File 'lib/farce/priority_queue.rb', line 51

def try_pop
  empty  = false
  result = @reverse_order ? @queue.pop_last { empty = true } : @queue.pop { empty = true }
  return @manager.unwrap(result) unless empty
  yield if block_given?
end

#try_push(value, priority: default_priority, mode: nil) { ... } ⇒ Boolean, BasicObject

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

  • priority (BasicObject) (defaults to: default_priority) —

    the value used to order the entry

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

    the default mode for sending unshareable values

Yields:

  • called when the queue is at capacity

Returns:

  • (Boolean, BasicObject) —

    true, or the fallback result when full



29
30
31
32
33
# File 'lib/farce/priority_queue.rb', line 29

def try_push(value, priority: default_priority, mode: nil)
  return true if @queue.push(priority, @manager.wrap(value, mode:))

  block_given? ? yield : false
end