Class: Farce::Abstract::Queue Abstract

Inherits:
Object
  • Object
show all
Includes:
Internal::Inspect, Internal::MarshalSupport::Reject
Defined in:
lib/farce/abstract/queue.rb,
lib/farce/integrations/active_support/duplicable.rb

Overview

This class is abstract.

A shared super class for all queues defined by Farce.

A queue is a collection of items that can be added to and removed from in a thread-safe manner. The order is typically FIFO (first-in, first-out), but other orderings may be applied on top.

Queues have a blocking and non-blocking interface.

Instance Attribute Summary collapse

ActiveSupport Integration collapse

Instance Method Summary collapse

Constructor Details

#initialize(capacity: 1024, track_age: false) ⇒ BasicObject

This method is abstract.

Subclasses may add additional parameters to this method.

Parameters:

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

    The maximum number of items the queue can hold. If nil, the queue is unbounded.

    The default may vary for subclasses. Most notably, priority and timer queues default to nil.

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

    Whether to track enqueue age and queue generations.



147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
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/abstract/queue.rb', line 147

class Queue
  include Internal::MarshalSupport::Reject
  include Internal::Inspect

  # The default mode used to transfer values between Ractors.
  # @return [Symbol]
  def mode = :raise

  # Alias for {#pop} to match the interface of Ruby's Queue class.
  # @return (see #pop)
  def deq(...) = pop(...)
  alias shift deq

  # Alias for {#push} to match the interface of Ruby's Queue class.
  # @return (see #push)
  def enq(...) = push(...)
  alias << enq

  # Alias for {#size} to match the interface of Ruby's Queue class.
  # @return (see #size)
  def length = size

  # The maximum number of items the queue can hold.
  # @return [Integer, Float] the maximum number of items, or `Float::INFINITY` if the queue is unbounded.
  def max = capacity || Float::INFINITY

  # Whether the queue is at capacity.
  # @return [Boolean] `true` if the queue is at capacity, `false` otherwise.
  def full? = capacity && size >= capacity

  def capacity = internal_queue.capacity

  def clear
    internal_queue.clear
    self
  end

  def close
    internal_queue.close
    self
  end

  # Stop new pushes and close after the final queued item is removed.
  # @return [self]
  def seal
    internal_queue.seal
    self
  end

  def closed? = internal_queue.closed?
  def sealed? = internal_queue.sealed?

  # @return [Boolean] Whether enqueue-age tracking is enabled.
  def age_tracking? = internal_queue.age_tracking?

  # @return [Integer, nil] The mutation generation, or `nil` when tracking is disabled.
  def generation = internal_queue.generation

  # The timestamp is monotonic. Its origin is runtime-specific.
  # @return [Float, nil] The monotonic enqueue time of the oldest item.
  def oldest_enqueued_at = internal_queue.oldest_enqueued_at

  # @return [Float, nil] Seconds since the oldest item was enqueued.
  def oldest_age = internal_queue.oldest_age

  def empty? = size.zero?

  def num_waiting = internal_queue.num_waiting

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

  def push(value, non_block = false, timeout: nil) # rubocop:disable Style/OptionalBooleanParameter
    if non_block
      return true if internal_queue.try_push(value)
      raise ThreadError, "queue full"
    end
    timeout.nil? ? internal_queue.push(value) : internal_queue.push(value, timeout:)
  end

  def try_pop(&) = internal_queue.try_pop(&)

  def try_push(value)
    return true if internal_queue.try_push(value)
    block_given? ? yield : false
  end

  def size = internal_queue.size

  def wait_pop(timeout: nil) = internal_queue.wait_pop(timeout:)

  def wait_push(timeout: nil) = internal_queue.wait_push(timeout:)

  # @api private
  def inspect_with(inspector)
    super do
      if closed?
        inspector.breakable " "
        inspector.text "closed"
      else
        inspector.attributes(inspect_info)
      end
    end
  end

  private

  def initialize_copy(_other)
    raise TypeError, "queues cannot be copied"
  end

  def inspect_info
    info = { size:, capacity: }.compact
    info[:num_waiting] = num_waiting if num_waiting.positive?
    info
  end

  private def internal_queue = @queue
end

Instance Attribute Details

#capacity ⇒ Integer? (readonly)

This method is abstract.

Returns The maximum number of items the queue can hold.

Returns:

  • (Integer, nil) —

    The maximum number of items the queue can hold.



147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
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/abstract/queue.rb', line 147

class Queue
  include Internal::MarshalSupport::Reject
  include Internal::Inspect

  # The default mode used to transfer values between Ractors.
  # @return [Symbol]
  def mode = :raise

  # Alias for {#pop} to match the interface of Ruby's Queue class.
  # @return (see #pop)
  def deq(...) = pop(...)
  alias shift deq

  # Alias for {#push} to match the interface of Ruby's Queue class.
  # @return (see #push)
  def enq(...) = push(...)
  alias << enq

  # Alias for {#size} to match the interface of Ruby's Queue class.
  # @return (see #size)
  def length = size

  # The maximum number of items the queue can hold.
  # @return [Integer, Float] the maximum number of items, or `Float::INFINITY` if the queue is unbounded.
  def max = capacity || Float::INFINITY

  # Whether the queue is at capacity.
  # @return [Boolean] `true` if the queue is at capacity, `false` otherwise.
  def full? = capacity && size >= capacity

  def capacity = internal_queue.capacity

  def clear
    internal_queue.clear
    self
  end

  def close
    internal_queue.close
    self
  end

  # Stop new pushes and close after the final queued item is removed.
  # @return [self]
  def seal
    internal_queue.seal
    self
  end

  def closed? = internal_queue.closed?
  def sealed? = internal_queue.sealed?

  # @return [Boolean] Whether enqueue-age tracking is enabled.
  def age_tracking? = internal_queue.age_tracking?

  # @return [Integer, nil] The mutation generation, or `nil` when tracking is disabled.
  def generation = internal_queue.generation

  # The timestamp is monotonic. Its origin is runtime-specific.
  # @return [Float, nil] The monotonic enqueue time of the oldest item.
  def oldest_enqueued_at = internal_queue.oldest_enqueued_at

  # @return [Float, nil] Seconds since the oldest item was enqueued.
  def oldest_age = internal_queue.oldest_age

  def empty? = size.zero?

  def num_waiting = internal_queue.num_waiting

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

  def push(value, non_block = false, timeout: nil) # rubocop:disable Style/OptionalBooleanParameter
    if non_block
      return true if internal_queue.try_push(value)
      raise ThreadError, "queue full"
    end
    timeout.nil? ? internal_queue.push(value) : internal_queue.push(value, timeout:)
  end

  def try_pop(&) = internal_queue.try_pop(&)

  def try_push(value)
    return true if internal_queue.try_push(value)
    block_given? ? yield : false
  end

  def size = internal_queue.size

  def wait_pop(timeout: nil) = internal_queue.wait_pop(timeout:)

  def wait_push(timeout: nil) = internal_queue.wait_push(timeout:)

  # @api private
  def inspect_with(inspector)
    super do
      if closed?
        inspector.breakable " "
        inspector.text "closed"
      else
        inspector.attributes(inspect_info)
      end
    end
  end

  private

  def initialize_copy(_other)
    raise TypeError, "queues cannot be copied"
  end

  def inspect_info
    info = { size:, capacity: }.compact
    info[:num_waiting] = num_waiting if num_waiting.positive?
    info
  end

  private def internal_queue = @queue
end

#mode ⇒ Symbol (readonly)

The default mode used to transfer values between Ractors.

Returns:

  • (Symbol)


147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
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/abstract/queue.rb', line 147

class Queue
  include Internal::MarshalSupport::Reject
  include Internal::Inspect

  # The default mode used to transfer values between Ractors.
  # @return [Symbol]
  def mode = :raise

  # Alias for {#pop} to match the interface of Ruby's Queue class.
  # @return (see #pop)
  def deq(...) = pop(...)
  alias shift deq

  # Alias for {#push} to match the interface of Ruby's Queue class.
  # @return (see #push)
  def enq(...) = push(...)
  alias << enq

  # Alias for {#size} to match the interface of Ruby's Queue class.
  # @return (see #size)
  def length = size

  # The maximum number of items the queue can hold.
  # @return [Integer, Float] the maximum number of items, or `Float::INFINITY` if the queue is unbounded.
  def max = capacity || Float::INFINITY

  # Whether the queue is at capacity.
  # @return [Boolean] `true` if the queue is at capacity, `false` otherwise.
  def full? = capacity && size >= capacity

  def capacity = internal_queue.capacity

  def clear
    internal_queue.clear
    self
  end

  def close
    internal_queue.close
    self
  end

  # Stop new pushes and close after the final queued item is removed.
  # @return [self]
  def seal
    internal_queue.seal
    self
  end

  def closed? = internal_queue.closed?
  def sealed? = internal_queue.sealed?

  # @return [Boolean] Whether enqueue-age tracking is enabled.
  def age_tracking? = internal_queue.age_tracking?

  # @return [Integer, nil] The mutation generation, or `nil` when tracking is disabled.
  def generation = internal_queue.generation

  # The timestamp is monotonic. Its origin is runtime-specific.
  # @return [Float, nil] The monotonic enqueue time of the oldest item.
  def oldest_enqueued_at = internal_queue.oldest_enqueued_at

  # @return [Float, nil] Seconds since the oldest item was enqueued.
  def oldest_age = internal_queue.oldest_age

  def empty? = size.zero?

  def num_waiting = internal_queue.num_waiting

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

  def push(value, non_block = false, timeout: nil) # rubocop:disable Style/OptionalBooleanParameter
    if non_block
      return true if internal_queue.try_push(value)
      raise ThreadError, "queue full"
    end
    timeout.nil? ? internal_queue.push(value) : internal_queue.push(value, timeout:)
  end

  def try_pop(&) = internal_queue.try_pop(&)

  def try_push(value)
    return true if internal_queue.try_push(value)
    block_given? ? yield : false
  end

  def size = internal_queue.size

  def wait_pop(timeout: nil) = internal_queue.wait_pop(timeout:)

  def wait_push(timeout: nil) = internal_queue.wait_push(timeout:)

  # @api private
  def inspect_with(inspector)
    super do
      if closed?
        inspector.breakable " "
        inspector.text "closed"
      else
        inspector.attributes(inspect_info)
      end
    end
  end

  private

  def initialize_copy(_other)
    raise TypeError, "queues cannot be copied"
  end

  def inspect_info
    info = { size:, capacity: }.compact
    info[:num_waiting] = num_waiting if num_waiting.positive?
    info
  end

  private def internal_queue = @queue
end

#num_waiting ⇒ Integer (readonly)

This method is abstract.

The number of threads/fibers currently waiting on the queue. This is a momentary snapshot and may be an approximation. The number may also include threads waiting for a push to succeed, as well as threads waiting without modifying the queue.

Returns:

  • (Integer) —

    Number of threads/fibers waiting on the queue.



147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
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/abstract/queue.rb', line 147

class Queue
  include Internal::MarshalSupport::Reject
  include Internal::Inspect

  # The default mode used to transfer values between Ractors.
  # @return [Symbol]
  def mode = :raise

  # Alias for {#pop} to match the interface of Ruby's Queue class.
  # @return (see #pop)
  def deq(...) = pop(...)
  alias shift deq

  # Alias for {#push} to match the interface of Ruby's Queue class.
  # @return (see #push)
  def enq(...) = push(...)
  alias << enq

  # Alias for {#size} to match the interface of Ruby's Queue class.
  # @return (see #size)
  def length = size

  # The maximum number of items the queue can hold.
  # @return [Integer, Float] the maximum number of items, or `Float::INFINITY` if the queue is unbounded.
  def max = capacity || Float::INFINITY

  # Whether the queue is at capacity.
  # @return [Boolean] `true` if the queue is at capacity, `false` otherwise.
  def full? = capacity && size >= capacity

  def capacity = internal_queue.capacity

  def clear
    internal_queue.clear
    self
  end

  def close
    internal_queue.close
    self
  end

  # Stop new pushes and close after the final queued item is removed.
  # @return [self]
  def seal
    internal_queue.seal
    self
  end

  def closed? = internal_queue.closed?
  def sealed? = internal_queue.sealed?

  # @return [Boolean] Whether enqueue-age tracking is enabled.
  def age_tracking? = internal_queue.age_tracking?

  # @return [Integer, nil] The mutation generation, or `nil` when tracking is disabled.
  def generation = internal_queue.generation

  # The timestamp is monotonic. Its origin is runtime-specific.
  # @return [Float, nil] The monotonic enqueue time of the oldest item.
  def oldest_enqueued_at = internal_queue.oldest_enqueued_at

  # @return [Float, nil] Seconds since the oldest item was enqueued.
  def oldest_age = internal_queue.oldest_age

  def empty? = size.zero?

  def num_waiting = internal_queue.num_waiting

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

  def push(value, non_block = false, timeout: nil) # rubocop:disable Style/OptionalBooleanParameter
    if non_block
      return true if internal_queue.try_push(value)
      raise ThreadError, "queue full"
    end
    timeout.nil? ? internal_queue.push(value) : internal_queue.push(value, timeout:)
  end

  def try_pop(&) = internal_queue.try_pop(&)

  def try_push(value)
    return true if internal_queue.try_push(value)
    block_given? ? yield : false
  end

  def size = internal_queue.size

  def wait_pop(timeout: nil) = internal_queue.wait_pop(timeout:)

  def wait_push(timeout: nil) = internal_queue.wait_push(timeout:)

  # @api private
  def inspect_with(inspector)
    super do
      if closed?
        inspector.breakable " "
        inspector.text "closed"
      else
        inspector.attributes(inspect_info)
      end
    end
  end

  private

  def initialize_copy(_other)
    raise TypeError, "queues cannot be copied"
  end

  def inspect_info
    info = { size:, capacity: }.compact
    info[:num_waiting] = num_waiting if num_waiting.positive?
    info
  end

  private def internal_queue = @queue
end

#size ⇒ Integer (readonly)

This method is abstract.

Returns The current number of items in the queue.

Returns:

  • (Integer) —

    The current number of items in the queue.



147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
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/abstract/queue.rb', line 147

class Queue
  include Internal::MarshalSupport::Reject
  include Internal::Inspect

  # The default mode used to transfer values between Ractors.
  # @return [Symbol]
  def mode = :raise

  # Alias for {#pop} to match the interface of Ruby's Queue class.
  # @return (see #pop)
  def deq(...) = pop(...)
  alias shift deq

  # Alias for {#push} to match the interface of Ruby's Queue class.
  # @return (see #push)
  def enq(...) = push(...)
  alias << enq

  # Alias for {#size} to match the interface of Ruby's Queue class.
  # @return (see #size)
  def length = size

  # The maximum number of items the queue can hold.
  # @return [Integer, Float] the maximum number of items, or `Float::INFINITY` if the queue is unbounded.
  def max = capacity || Float::INFINITY

  # Whether the queue is at capacity.
  # @return [Boolean] `true` if the queue is at capacity, `false` otherwise.
  def full? = capacity && size >= capacity

  def capacity = internal_queue.capacity

  def clear
    internal_queue.clear
    self
  end

  def close
    internal_queue.close
    self
  end

  # Stop new pushes and close after the final queued item is removed.
  # @return [self]
  def seal
    internal_queue.seal
    self
  end

  def closed? = internal_queue.closed?
  def sealed? = internal_queue.sealed?

  # @return [Boolean] Whether enqueue-age tracking is enabled.
  def age_tracking? = internal_queue.age_tracking?

  # @return [Integer, nil] The mutation generation, or `nil` when tracking is disabled.
  def generation = internal_queue.generation

  # The timestamp is monotonic. Its origin is runtime-specific.
  # @return [Float, nil] The monotonic enqueue time of the oldest item.
  def oldest_enqueued_at = internal_queue.oldest_enqueued_at

  # @return [Float, nil] Seconds since the oldest item was enqueued.
  def oldest_age = internal_queue.oldest_age

  def empty? = size.zero?

  def num_waiting = internal_queue.num_waiting

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

  def push(value, non_block = false, timeout: nil) # rubocop:disable Style/OptionalBooleanParameter
    if non_block
      return true if internal_queue.try_push(value)
      raise ThreadError, "queue full"
    end
    timeout.nil? ? internal_queue.push(value) : internal_queue.push(value, timeout:)
  end

  def try_pop(&) = internal_queue.try_pop(&)

  def try_push(value)
    return true if internal_queue.try_push(value)
    block_given? ? yield : false
  end

  def size = internal_queue.size

  def wait_pop(timeout: nil) = internal_queue.wait_pop(timeout:)

  def wait_push(timeout: nil) = internal_queue.wait_push(timeout:)

  # @api private
  def inspect_with(inspector)
    super do
      if closed?
        inspector.breakable " "
        inspector.text "closed"
      else
        inspector.attributes(inspect_info)
      end
    end
  end

  private

  def initialize_copy(_other)
    raise TypeError, "queues cannot be copied"
  end

  def inspect_info
    info = { size:, capacity: }.compact
    info[:num_waiting] = num_waiting if num_waiting.positive?
    info
  end

  private def internal_queue = @queue
end

Instance Method Details

#age_tracking? ⇒ Boolean

Returns Whether enqueue-age tracking is enabled.

Returns:

  • (Boolean) —

    Whether enqueue-age tracking is enabled.



200
# File 'lib/farce/abstract/queue.rb', line 200

def age_tracking? = internal_queue.age_tracking?

#clear ⇒ self

This method is abstract.

Removes all items from the queue.

Returns:

  • (self)


147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
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/abstract/queue.rb', line 147

class Queue
  include Internal::MarshalSupport::Reject
  include Internal::Inspect

  # The default mode used to transfer values between Ractors.
  # @return [Symbol]
  def mode = :raise

  # Alias for {#pop} to match the interface of Ruby's Queue class.
  # @return (see #pop)
  def deq(...) = pop(...)
  alias shift deq

  # Alias for {#push} to match the interface of Ruby's Queue class.
  # @return (see #push)
  def enq(...) = push(...)
  alias << enq

  # Alias for {#size} to match the interface of Ruby's Queue class.
  # @return (see #size)
  def length = size

  # The maximum number of items the queue can hold.
  # @return [Integer, Float] the maximum number of items, or `Float::INFINITY` if the queue is unbounded.
  def max = capacity || Float::INFINITY

  # Whether the queue is at capacity.
  # @return [Boolean] `true` if the queue is at capacity, `false` otherwise.
  def full? = capacity && size >= capacity

  def capacity = internal_queue.capacity

  def clear
    internal_queue.clear
    self
  end

  def close
    internal_queue.close
    self
  end

  # Stop new pushes and close after the final queued item is removed.
  # @return [self]
  def seal
    internal_queue.seal
    self
  end

  def closed? = internal_queue.closed?
  def sealed? = internal_queue.sealed?

  # @return [Boolean] Whether enqueue-age tracking is enabled.
  def age_tracking? = internal_queue.age_tracking?

  # @return [Integer, nil] The mutation generation, or `nil` when tracking is disabled.
  def generation = internal_queue.generation

  # The timestamp is monotonic. Its origin is runtime-specific.
  # @return [Float, nil] The monotonic enqueue time of the oldest item.
  def oldest_enqueued_at = internal_queue.oldest_enqueued_at

  # @return [Float, nil] Seconds since the oldest item was enqueued.
  def oldest_age = internal_queue.oldest_age

  def empty? = size.zero?

  def num_waiting = internal_queue.num_waiting

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

  def push(value, non_block = false, timeout: nil) # rubocop:disable Style/OptionalBooleanParameter
    if non_block
      return true if internal_queue.try_push(value)
      raise ThreadError, "queue full"
    end
    timeout.nil? ? internal_queue.push(value) : internal_queue.push(value, timeout:)
  end

  def try_pop(&) = internal_queue.try_pop(&)

  def try_push(value)
    return true if internal_queue.try_push(value)
    block_given? ? yield : false
  end

  def size = internal_queue.size

  def wait_pop(timeout: nil) = internal_queue.wait_pop(timeout:)

  def wait_push(timeout: nil) = internal_queue.wait_push(timeout:)

  # @api private
  def inspect_with(inspector)
    super do
      if closed?
        inspector.breakable " "
        inspector.text "closed"
      else
        inspector.attributes(inspect_info)
      end
    end
  end

  private

  def initialize_copy(_other)
    raise TypeError, "queues cannot be copied"
  end

  def inspect_info
    info = { size:, capacity: }.compact
    info[:num_waiting] = num_waiting if num_waiting.positive?
    info
  end

  private def internal_queue = @queue
end

#close ⇒ self

This method is abstract.

Closes the queue, preventing any further items from being added.

Returns:

  • (self)


147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
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/abstract/queue.rb', line 147

class Queue
  include Internal::MarshalSupport::Reject
  include Internal::Inspect

  # The default mode used to transfer values between Ractors.
  # @return [Symbol]
  def mode = :raise

  # Alias for {#pop} to match the interface of Ruby's Queue class.
  # @return (see #pop)
  def deq(...) = pop(...)
  alias shift deq

  # Alias for {#push} to match the interface of Ruby's Queue class.
  # @return (see #push)
  def enq(...) = push(...)
  alias << enq

  # Alias for {#size} to match the interface of Ruby's Queue class.
  # @return (see #size)
  def length = size

  # The maximum number of items the queue can hold.
  # @return [Integer, Float] the maximum number of items, or `Float::INFINITY` if the queue is unbounded.
  def max = capacity || Float::INFINITY

  # Whether the queue is at capacity.
  # @return [Boolean] `true` if the queue is at capacity, `false` otherwise.
  def full? = capacity && size >= capacity

  def capacity = internal_queue.capacity

  def clear
    internal_queue.clear
    self
  end

  def close
    internal_queue.close
    self
  end

  # Stop new pushes and close after the final queued item is removed.
  # @return [self]
  def seal
    internal_queue.seal
    self
  end

  def closed? = internal_queue.closed?
  def sealed? = internal_queue.sealed?

  # @return [Boolean] Whether enqueue-age tracking is enabled.
  def age_tracking? = internal_queue.age_tracking?

  # @return [Integer, nil] The mutation generation, or `nil` when tracking is disabled.
  def generation = internal_queue.generation

  # The timestamp is monotonic. Its origin is runtime-specific.
  # @return [Float, nil] The monotonic enqueue time of the oldest item.
  def oldest_enqueued_at = internal_queue.oldest_enqueued_at

  # @return [Float, nil] Seconds since the oldest item was enqueued.
  def oldest_age = internal_queue.oldest_age

  def empty? = size.zero?

  def num_waiting = internal_queue.num_waiting

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

  def push(value, non_block = false, timeout: nil) # rubocop:disable Style/OptionalBooleanParameter
    if non_block
      return true if internal_queue.try_push(value)
      raise ThreadError, "queue full"
    end
    timeout.nil? ? internal_queue.push(value) : internal_queue.push(value, timeout:)
  end

  def try_pop(&) = internal_queue.try_pop(&)

  def try_push(value)
    return true if internal_queue.try_push(value)
    block_given? ? yield : false
  end

  def size = internal_queue.size

  def wait_pop(timeout: nil) = internal_queue.wait_pop(timeout:)

  def wait_push(timeout: nil) = internal_queue.wait_push(timeout:)

  # @api private
  def inspect_with(inspector)
    super do
      if closed?
        inspector.breakable " "
        inspector.text "closed"
      else
        inspector.attributes(inspect_info)
      end
    end
  end

  private

  def initialize_copy(_other)
    raise TypeError, "queues cannot be copied"
  end

  def inspect_info
    info = { size:, capacity: }.compact
    info[:num_waiting] = num_waiting if num_waiting.positive?
    info
  end

  private def internal_queue = @queue
end

#closed? ⇒ Boolean

This method is abstract.

Checks whether the queue is closed.

Returns:

  • (Boolean) —

    true if the queue is closed, false otherwise.



147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
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/abstract/queue.rb', line 147

class Queue
  include Internal::MarshalSupport::Reject
  include Internal::Inspect

  # The default mode used to transfer values between Ractors.
  # @return [Symbol]
  def mode = :raise

  # Alias for {#pop} to match the interface of Ruby's Queue class.
  # @return (see #pop)
  def deq(...) = pop(...)
  alias shift deq

  # Alias for {#push} to match the interface of Ruby's Queue class.
  # @return (see #push)
  def enq(...) = push(...)
  alias << enq

  # Alias for {#size} to match the interface of Ruby's Queue class.
  # @return (see #size)
  def length = size

  # The maximum number of items the queue can hold.
  # @return [Integer, Float] the maximum number of items, or `Float::INFINITY` if the queue is unbounded.
  def max = capacity || Float::INFINITY

  # Whether the queue is at capacity.
  # @return [Boolean] `true` if the queue is at capacity, `false` otherwise.
  def full? = capacity && size >= capacity

  def capacity = internal_queue.capacity

  def clear
    internal_queue.clear
    self
  end

  def close
    internal_queue.close
    self
  end

  # Stop new pushes and close after the final queued item is removed.
  # @return [self]
  def seal
    internal_queue.seal
    self
  end

  def closed? = internal_queue.closed?
  def sealed? = internal_queue.sealed?

  # @return [Boolean] Whether enqueue-age tracking is enabled.
  def age_tracking? = internal_queue.age_tracking?

  # @return [Integer, nil] The mutation generation, or `nil` when tracking is disabled.
  def generation = internal_queue.generation

  # The timestamp is monotonic. Its origin is runtime-specific.
  # @return [Float, nil] The monotonic enqueue time of the oldest item.
  def oldest_enqueued_at = internal_queue.oldest_enqueued_at

  # @return [Float, nil] Seconds since the oldest item was enqueued.
  def oldest_age = internal_queue.oldest_age

  def empty? = size.zero?

  def num_waiting = internal_queue.num_waiting

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

  def push(value, non_block = false, timeout: nil) # rubocop:disable Style/OptionalBooleanParameter
    if non_block
      return true if internal_queue.try_push(value)
      raise ThreadError, "queue full"
    end
    timeout.nil? ? internal_queue.push(value) : internal_queue.push(value, timeout:)
  end

  def try_pop(&) = internal_queue.try_pop(&)

  def try_push(value)
    return true if internal_queue.try_push(value)
    block_given? ? yield : false
  end

  def size = internal_queue.size

  def wait_pop(timeout: nil) = internal_queue.wait_pop(timeout:)

  def wait_push(timeout: nil) = internal_queue.wait_push(timeout:)

  # @api private
  def inspect_with(inspector)
    super do
      if closed?
        inspector.breakable " "
        inspector.text "closed"
      else
        inspector.attributes(inspect_info)
      end
    end
  end

  private

  def initialize_copy(_other)
    raise TypeError, "queues cannot be copied"
  end

  def inspect_info
    info = { size:, capacity: }.compact
    info[:num_waiting] = num_waiting if num_waiting.positive?
    info
  end

  private def internal_queue = @queue
end

#deq ⇒ BasicObject Also known as: shift

Alias for #pop to match the interface of Ruby's Queue class.



157
# File 'lib/farce/abstract/queue.rb', line 157

def deq(...) = pop(...)

#duplicable? ⇒ Boolean

Note:

This methods is only available if ActiveSupport has been loaded.

Returns false.

Returns:

  • (Boolean) —

    false



37
# File 'lib/farce/integrations/active_support/duplicable.rb', line 37

def duplicable? = false

#empty? ⇒ Boolean

This method is abstract.

Checks whether the queue is empty.

Returns:

  • (Boolean) —

    true if the queue is empty, false otherwise.



147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
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/abstract/queue.rb', line 147

class Queue
  include Internal::MarshalSupport::Reject
  include Internal::Inspect

  # The default mode used to transfer values between Ractors.
  # @return [Symbol]
  def mode = :raise

  # Alias for {#pop} to match the interface of Ruby's Queue class.
  # @return (see #pop)
  def deq(...) = pop(...)
  alias shift deq

  # Alias for {#push} to match the interface of Ruby's Queue class.
  # @return (see #push)
  def enq(...) = push(...)
  alias << enq

  # Alias for {#size} to match the interface of Ruby's Queue class.
  # @return (see #size)
  def length = size

  # The maximum number of items the queue can hold.
  # @return [Integer, Float] the maximum number of items, or `Float::INFINITY` if the queue is unbounded.
  def max = capacity || Float::INFINITY

  # Whether the queue is at capacity.
  # @return [Boolean] `true` if the queue is at capacity, `false` otherwise.
  def full? = capacity && size >= capacity

  def capacity = internal_queue.capacity

  def clear
    internal_queue.clear
    self
  end

  def close
    internal_queue.close
    self
  end

  # Stop new pushes and close after the final queued item is removed.
  # @return [self]
  def seal
    internal_queue.seal
    self
  end

  def closed? = internal_queue.closed?
  def sealed? = internal_queue.sealed?

  # @return [Boolean] Whether enqueue-age tracking is enabled.
  def age_tracking? = internal_queue.age_tracking?

  # @return [Integer, nil] The mutation generation, or `nil` when tracking is disabled.
  def generation = internal_queue.generation

  # The timestamp is monotonic. Its origin is runtime-specific.
  # @return [Float, nil] The monotonic enqueue time of the oldest item.
  def oldest_enqueued_at = internal_queue.oldest_enqueued_at

  # @return [Float, nil] Seconds since the oldest item was enqueued.
  def oldest_age = internal_queue.oldest_age

  def empty? = size.zero?

  def num_waiting = internal_queue.num_waiting

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

  def push(value, non_block = false, timeout: nil) # rubocop:disable Style/OptionalBooleanParameter
    if non_block
      return true if internal_queue.try_push(value)
      raise ThreadError, "queue full"
    end
    timeout.nil? ? internal_queue.push(value) : internal_queue.push(value, timeout:)
  end

  def try_pop(&) = internal_queue.try_pop(&)

  def try_push(value)
    return true if internal_queue.try_push(value)
    block_given? ? yield : false
  end

  def size = internal_queue.size

  def wait_pop(timeout: nil) = internal_queue.wait_pop(timeout:)

  def wait_push(timeout: nil) = internal_queue.wait_push(timeout:)

  # @api private
  def inspect_with(inspector)
    super do
      if closed?
        inspector.breakable " "
        inspector.text "closed"
      else
        inspector.attributes(inspect_info)
      end
    end
  end

  private

  def initialize_copy(_other)
    raise TypeError, "queues cannot be copied"
  end

  def inspect_info
    info = { size:, capacity: }.compact
    info[:num_waiting] = num_waiting if num_waiting.positive?
    info
  end

  private def internal_queue = @queue
end

#enq ⇒ BasicObject Also known as: <<

Alias for #push to match the interface of Ruby's Queue class.



162
# File 'lib/farce/abstract/queue.rb', line 162

def enq(...) = push(...)

#full? ⇒ Boolean

Whether the queue is at capacity.

Returns:

  • (Boolean) —

    true if the queue is at capacity, false otherwise.



175
# File 'lib/farce/abstract/queue.rb', line 175

def full? = capacity && size >= capacity

#generation ⇒ Integer?

Returns The mutation generation, or nil when tracking is disabled.

Returns:

  • (Integer, nil) —

    The mutation generation, or nil when tracking is disabled.



203
# File 'lib/farce/abstract/queue.rb', line 203

def generation = internal_queue.generation

#length ⇒ Integer

Alias for #size to match the interface of Ruby's Queue class.

Returns:

  • (Integer) —

    The current number of items in the queue.



167
# File 'lib/farce/abstract/queue.rb', line 167

def length = size

#max ⇒ Integer, Float

The maximum number of items the queue can hold.

Returns:

  • (Integer, Float) —

    the maximum number of items, or Float::INFINITY if the queue is unbounded.



171
# File 'lib/farce/abstract/queue.rb', line 171

def max = capacity || Float::INFINITY

#oldest_age ⇒ Float?

Returns Seconds since the oldest item was enqueued.

Returns:

  • (Float, nil) —

    Seconds since the oldest item was enqueued.



210
# File 'lib/farce/abstract/queue.rb', line 210

def oldest_age = internal_queue.oldest_age

#oldest_enqueued_at ⇒ Float?

The timestamp is monotonic. Its origin is runtime-specific.

Returns:

  • (Float, nil) —

    The monotonic enqueue time of the oldest item.



207
# File 'lib/farce/abstract/queue.rb', line 207

def oldest_enqueued_at = internal_queue.oldest_enqueued_at

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

rubocop:disable Style/OptionalBooleanParameter



147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
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/abstract/queue.rb', line 147

class Queue
  include Internal::MarshalSupport::Reject
  include Internal::Inspect

  # The default mode used to transfer values between Ractors.
  # @return [Symbol]
  def mode = :raise

  # Alias for {#pop} to match the interface of Ruby's Queue class.
  # @return (see #pop)
  def deq(...) = pop(...)
  alias shift deq

  # Alias for {#push} to match the interface of Ruby's Queue class.
  # @return (see #push)
  def enq(...) = push(...)
  alias << enq

  # Alias for {#size} to match the interface of Ruby's Queue class.
  # @return (see #size)
  def length = size

  # The maximum number of items the queue can hold.
  # @return [Integer, Float] the maximum number of items, or `Float::INFINITY` if the queue is unbounded.
  def max = capacity || Float::INFINITY

  # Whether the queue is at capacity.
  # @return [Boolean] `true` if the queue is at capacity, `false` otherwise.
  def full? = capacity && size >= capacity

  def capacity = internal_queue.capacity

  def clear
    internal_queue.clear
    self
  end

  def close
    internal_queue.close
    self
  end

  # Stop new pushes and close after the final queued item is removed.
  # @return [self]
  def seal
    internal_queue.seal
    self
  end

  def closed? = internal_queue.closed?
  def sealed? = internal_queue.sealed?

  # @return [Boolean] Whether enqueue-age tracking is enabled.
  def age_tracking? = internal_queue.age_tracking?

  # @return [Integer, nil] The mutation generation, or `nil` when tracking is disabled.
  def generation = internal_queue.generation

  # The timestamp is monotonic. Its origin is runtime-specific.
  # @return [Float, nil] The monotonic enqueue time of the oldest item.
  def oldest_enqueued_at = internal_queue.oldest_enqueued_at

  # @return [Float, nil] Seconds since the oldest item was enqueued.
  def oldest_age = internal_queue.oldest_age

  def empty? = size.zero?

  def num_waiting = internal_queue.num_waiting

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

  def push(value, non_block = false, timeout: nil) # rubocop:disable Style/OptionalBooleanParameter
    if non_block
      return true if internal_queue.try_push(value)
      raise ThreadError, "queue full"
    end
    timeout.nil? ? internal_queue.push(value) : internal_queue.push(value, timeout:)
  end

  def try_pop(&) = internal_queue.try_pop(&)

  def try_push(value)
    return true if internal_queue.try_push(value)
    block_given? ? yield : false
  end

  def size = internal_queue.size

  def wait_pop(timeout: nil) = internal_queue.wait_pop(timeout:)

  def wait_push(timeout: nil) = internal_queue.wait_push(timeout:)

  # @api private
  def inspect_with(inspector)
    super do
      if closed?
        inspector.breakable " "
        inspector.text "closed"
      else
        inspector.attributes(inspect_info)
      end
    end
  end

  private

  def initialize_copy(_other)
    raise TypeError, "queues cannot be copied"
  end

  def inspect_info
    info = { size:, capacity: }.compact
    info[:num_waiting] = num_waiting if num_waiting.positive?
    info
  end

  private def internal_queue = @queue
end

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

rubocop:disable Style/OptionalBooleanParameter



147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
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/abstract/queue.rb', line 147

class Queue
  include Internal::MarshalSupport::Reject
  include Internal::Inspect

  # The default mode used to transfer values between Ractors.
  # @return [Symbol]
  def mode = :raise

  # Alias for {#pop} to match the interface of Ruby's Queue class.
  # @return (see #pop)
  def deq(...) = pop(...)
  alias shift deq

  # Alias for {#push} to match the interface of Ruby's Queue class.
  # @return (see #push)
  def enq(...) = push(...)
  alias << enq

  # Alias for {#size} to match the interface of Ruby's Queue class.
  # @return (see #size)
  def length = size

  # The maximum number of items the queue can hold.
  # @return [Integer, Float] the maximum number of items, or `Float::INFINITY` if the queue is unbounded.
  def max = capacity || Float::INFINITY

  # Whether the queue is at capacity.
  # @return [Boolean] `true` if the queue is at capacity, `false` otherwise.
  def full? = capacity && size >= capacity

  def capacity = internal_queue.capacity

  def clear
    internal_queue.clear
    self
  end

  def close
    internal_queue.close
    self
  end

  # Stop new pushes and close after the final queued item is removed.
  # @return [self]
  def seal
    internal_queue.seal
    self
  end

  def closed? = internal_queue.closed?
  def sealed? = internal_queue.sealed?

  # @return [Boolean] Whether enqueue-age tracking is enabled.
  def age_tracking? = internal_queue.age_tracking?

  # @return [Integer, nil] The mutation generation, or `nil` when tracking is disabled.
  def generation = internal_queue.generation

  # The timestamp is monotonic. Its origin is runtime-specific.
  # @return [Float, nil] The monotonic enqueue time of the oldest item.
  def oldest_enqueued_at = internal_queue.oldest_enqueued_at

  # @return [Float, nil] Seconds since the oldest item was enqueued.
  def oldest_age = internal_queue.oldest_age

  def empty? = size.zero?

  def num_waiting = internal_queue.num_waiting

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

  def push(value, non_block = false, timeout: nil) # rubocop:disable Style/OptionalBooleanParameter
    if non_block
      return true if internal_queue.try_push(value)
      raise ThreadError, "queue full"
    end
    timeout.nil? ? internal_queue.push(value) : internal_queue.push(value, timeout:)
  end

  def try_pop(&) = internal_queue.try_pop(&)

  def try_push(value)
    return true if internal_queue.try_push(value)
    block_given? ? yield : false
  end

  def size = internal_queue.size

  def wait_pop(timeout: nil) = internal_queue.wait_pop(timeout:)

  def wait_push(timeout: nil) = internal_queue.wait_push(timeout:)

  # @api private
  def inspect_with(inspector)
    super do
      if closed?
        inspector.breakable " "
        inspector.text "closed"
      else
        inspector.attributes(inspect_info)
      end
    end
  end

  private

  def initialize_copy(_other)
    raise TypeError, "queues cannot be copied"
  end

  def inspect_info
    info = { size:, capacity: }.compact
    info[:num_waiting] = num_waiting if num_waiting.positive?
    info
  end

  private def internal_queue = @queue
end

#seal ⇒ self

Stop new pushes and close after the final queued item is removed.

Returns:

  • (self)


147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
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/abstract/queue.rb', line 147

class Queue
  include Internal::MarshalSupport::Reject
  include Internal::Inspect

  # The default mode used to transfer values between Ractors.
  # @return [Symbol]
  def mode = :raise

  # Alias for {#pop} to match the interface of Ruby's Queue class.
  # @return (see #pop)
  def deq(...) = pop(...)
  alias shift deq

  # Alias for {#push} to match the interface of Ruby's Queue class.
  # @return (see #push)
  def enq(...) = push(...)
  alias << enq

  # Alias for {#size} to match the interface of Ruby's Queue class.
  # @return (see #size)
  def length = size

  # The maximum number of items the queue can hold.
  # @return [Integer, Float] the maximum number of items, or `Float::INFINITY` if the queue is unbounded.
  def max = capacity || Float::INFINITY

  # Whether the queue is at capacity.
  # @return [Boolean] `true` if the queue is at capacity, `false` otherwise.
  def full? = capacity && size >= capacity

  def capacity = internal_queue.capacity

  def clear
    internal_queue.clear
    self
  end

  def close
    internal_queue.close
    self
  end

  # Stop new pushes and close after the final queued item is removed.
  # @return [self]
  def seal
    internal_queue.seal
    self
  end

  def closed? = internal_queue.closed?
  def sealed? = internal_queue.sealed?

  # @return [Boolean] Whether enqueue-age tracking is enabled.
  def age_tracking? = internal_queue.age_tracking?

  # @return [Integer, nil] The mutation generation, or `nil` when tracking is disabled.
  def generation = internal_queue.generation

  # The timestamp is monotonic. Its origin is runtime-specific.
  # @return [Float, nil] The monotonic enqueue time of the oldest item.
  def oldest_enqueued_at = internal_queue.oldest_enqueued_at

  # @return [Float, nil] Seconds since the oldest item was enqueued.
  def oldest_age = internal_queue.oldest_age

  def empty? = size.zero?

  def num_waiting = internal_queue.num_waiting

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

  def push(value, non_block = false, timeout: nil) # rubocop:disable Style/OptionalBooleanParameter
    if non_block
      return true if internal_queue.try_push(value)
      raise ThreadError, "queue full"
    end
    timeout.nil? ? internal_queue.push(value) : internal_queue.push(value, timeout:)
  end

  def try_pop(&) = internal_queue.try_pop(&)

  def try_push(value)
    return true if internal_queue.try_push(value)
    block_given? ? yield : false
  end

  def size = internal_queue.size

  def wait_pop(timeout: nil) = internal_queue.wait_pop(timeout:)

  def wait_push(timeout: nil) = internal_queue.wait_push(timeout:)

  # @api private
  def inspect_with(inspector)
    super do
      if closed?
        inspector.breakable " "
        inspector.text "closed"
      else
        inspector.attributes(inspect_info)
      end
    end
  end

  private

  def initialize_copy(_other)
    raise TypeError, "queues cannot be copied"
  end

  def inspect_info
    info = { size:, capacity: }.compact
    info[:num_waiting] = num_waiting if num_waiting.positive?
    info
  end

  private def internal_queue = @queue
end

#sealed? ⇒ Boolean

Returns Whether the queue rejects new pushes.

Returns:

  • (Boolean) —

    Whether the queue rejects new pushes.



147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
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/abstract/queue.rb', line 147

class Queue
  include Internal::MarshalSupport::Reject
  include Internal::Inspect

  # The default mode used to transfer values between Ractors.
  # @return [Symbol]
  def mode = :raise

  # Alias for {#pop} to match the interface of Ruby's Queue class.
  # @return (see #pop)
  def deq(...) = pop(...)
  alias shift deq

  # Alias for {#push} to match the interface of Ruby's Queue class.
  # @return (see #push)
  def enq(...) = push(...)
  alias << enq

  # Alias for {#size} to match the interface of Ruby's Queue class.
  # @return (see #size)
  def length = size

  # The maximum number of items the queue can hold.
  # @return [Integer, Float] the maximum number of items, or `Float::INFINITY` if the queue is unbounded.
  def max = capacity || Float::INFINITY

  # Whether the queue is at capacity.
  # @return [Boolean] `true` if the queue is at capacity, `false` otherwise.
  def full? = capacity && size >= capacity

  def capacity = internal_queue.capacity

  def clear
    internal_queue.clear
    self
  end

  def close
    internal_queue.close
    self
  end

  # Stop new pushes and close after the final queued item is removed.
  # @return [self]
  def seal
    internal_queue.seal
    self
  end

  def closed? = internal_queue.closed?
  def sealed? = internal_queue.sealed?

  # @return [Boolean] Whether enqueue-age tracking is enabled.
  def age_tracking? = internal_queue.age_tracking?

  # @return [Integer, nil] The mutation generation, or `nil` when tracking is disabled.
  def generation = internal_queue.generation

  # The timestamp is monotonic. Its origin is runtime-specific.
  # @return [Float, nil] The monotonic enqueue time of the oldest item.
  def oldest_enqueued_at = internal_queue.oldest_enqueued_at

  # @return [Float, nil] Seconds since the oldest item was enqueued.
  def oldest_age = internal_queue.oldest_age

  def empty? = size.zero?

  def num_waiting = internal_queue.num_waiting

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

  def push(value, non_block = false, timeout: nil) # rubocop:disable Style/OptionalBooleanParameter
    if non_block
      return true if internal_queue.try_push(value)
      raise ThreadError, "queue full"
    end
    timeout.nil? ? internal_queue.push(value) : internal_queue.push(value, timeout:)
  end

  def try_pop(&) = internal_queue.try_pop(&)

  def try_push(value)
    return true if internal_queue.try_push(value)
    block_given? ? yield : false
  end

  def size = internal_queue.size

  def wait_pop(timeout: nil) = internal_queue.wait_pop(timeout:)

  def wait_push(timeout: nil) = internal_queue.wait_push(timeout:)

  # @api private
  def inspect_with(inspector)
    super do
      if closed?
        inspector.breakable " "
        inspector.text "closed"
      else
        inspector.attributes(inspect_info)
      end
    end
  end

  private

  def initialize_copy(_other)
    raise TypeError, "queues cannot be copied"
  end

  def inspect_info
    info = { size:, capacity: }.compact
    info[:num_waiting] = num_waiting if num_waiting.positive?
    info
  end

  private def internal_queue = @queue
end

#try_pop { ... } ⇒ BasicObject?

This method is abstract.

Subclasses may add additional parameters to this method.

Tries to take an item from the queue without blocking. If no item is available, the method will call the block if given, or return nil if no block is given.

Yields:

  • Block called if no item is available.

Returns:

  • (BasicObject, nil) —

    The item taken from the queue, or the return value of the block or nil if no item was available.



147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
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/abstract/queue.rb', line 147

class Queue
  include Internal::MarshalSupport::Reject
  include Internal::Inspect

  # The default mode used to transfer values between Ractors.
  # @return [Symbol]
  def mode = :raise

  # Alias for {#pop} to match the interface of Ruby's Queue class.
  # @return (see #pop)
  def deq(...) = pop(...)
  alias shift deq

  # Alias for {#push} to match the interface of Ruby's Queue class.
  # @return (see #push)
  def enq(...) = push(...)
  alias << enq

  # Alias for {#size} to match the interface of Ruby's Queue class.
  # @return (see #size)
  def length = size

  # The maximum number of items the queue can hold.
  # @return [Integer, Float] the maximum number of items, or `Float::INFINITY` if the queue is unbounded.
  def max = capacity || Float::INFINITY

  # Whether the queue is at capacity.
  # @return [Boolean] `true` if the queue is at capacity, `false` otherwise.
  def full? = capacity && size >= capacity

  def capacity = internal_queue.capacity

  def clear
    internal_queue.clear
    self
  end

  def close
    internal_queue.close
    self
  end

  # Stop new pushes and close after the final queued item is removed.
  # @return [self]
  def seal
    internal_queue.seal
    self
  end

  def closed? = internal_queue.closed?
  def sealed? = internal_queue.sealed?

  # @return [Boolean] Whether enqueue-age tracking is enabled.
  def age_tracking? = internal_queue.age_tracking?

  # @return [Integer, nil] The mutation generation, or `nil` when tracking is disabled.
  def generation = internal_queue.generation

  # The timestamp is monotonic. Its origin is runtime-specific.
  # @return [Float, nil] The monotonic enqueue time of the oldest item.
  def oldest_enqueued_at = internal_queue.oldest_enqueued_at

  # @return [Float, nil] Seconds since the oldest item was enqueued.
  def oldest_age = internal_queue.oldest_age

  def empty? = size.zero?

  def num_waiting = internal_queue.num_waiting

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

  def push(value, non_block = false, timeout: nil) # rubocop:disable Style/OptionalBooleanParameter
    if non_block
      return true if internal_queue.try_push(value)
      raise ThreadError, "queue full"
    end
    timeout.nil? ? internal_queue.push(value) : internal_queue.push(value, timeout:)
  end

  def try_pop(&) = internal_queue.try_pop(&)

  def try_push(value)
    return true if internal_queue.try_push(value)
    block_given? ? yield : false
  end

  def size = internal_queue.size

  def wait_pop(timeout: nil) = internal_queue.wait_pop(timeout:)

  def wait_push(timeout: nil) = internal_queue.wait_push(timeout:)

  # @api private
  def inspect_with(inspector)
    super do
      if closed?
        inspector.breakable " "
        inspector.text "closed"
      else
        inspector.attributes(inspect_info)
      end
    end
  end

  private

  def initialize_copy(_other)
    raise TypeError, "queues cannot be copied"
  end

  def inspect_info
    info = { size:, capacity: }.compact
    info[:num_waiting] = num_waiting if num_waiting.positive?
    info
  end

  private def internal_queue = @queue
end

#try_push(value) { ... } ⇒ Boolean

This method is abstract.

Subclasses may add additional parameters to this method.

Tries to push an item onto the queue without blocking. If the queue is at capacity, the method will call the block if given, or return false if no block is given.

Parameters:

  • value (BasicObject) —

    The item to add to the queue.

Yields:

  • Block called if the queue is at capacity.

Returns:

  • (Boolean) —

    true if the item was added to the queue, false if the queue was at capacity and no block was given.

Raises:



147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
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/abstract/queue.rb', line 147

class Queue
  include Internal::MarshalSupport::Reject
  include Internal::Inspect

  # The default mode used to transfer values between Ractors.
  # @return [Symbol]
  def mode = :raise

  # Alias for {#pop} to match the interface of Ruby's Queue class.
  # @return (see #pop)
  def deq(...) = pop(...)
  alias shift deq

  # Alias for {#push} to match the interface of Ruby's Queue class.
  # @return (see #push)
  def enq(...) = push(...)
  alias << enq

  # Alias for {#size} to match the interface of Ruby's Queue class.
  # @return (see #size)
  def length = size

  # The maximum number of items the queue can hold.
  # @return [Integer, Float] the maximum number of items, or `Float::INFINITY` if the queue is unbounded.
  def max = capacity || Float::INFINITY

  # Whether the queue is at capacity.
  # @return [Boolean] `true` if the queue is at capacity, `false` otherwise.
  def full? = capacity && size >= capacity

  def capacity = internal_queue.capacity

  def clear
    internal_queue.clear
    self
  end

  def close
    internal_queue.close
    self
  end

  # Stop new pushes and close after the final queued item is removed.
  # @return [self]
  def seal
    internal_queue.seal
    self
  end

  def closed? = internal_queue.closed?
  def sealed? = internal_queue.sealed?

  # @return [Boolean] Whether enqueue-age tracking is enabled.
  def age_tracking? = internal_queue.age_tracking?

  # @return [Integer, nil] The mutation generation, or `nil` when tracking is disabled.
  def generation = internal_queue.generation

  # The timestamp is monotonic. Its origin is runtime-specific.
  # @return [Float, nil] The monotonic enqueue time of the oldest item.
  def oldest_enqueued_at = internal_queue.oldest_enqueued_at

  # @return [Float, nil] Seconds since the oldest item was enqueued.
  def oldest_age = internal_queue.oldest_age

  def empty? = size.zero?

  def num_waiting = internal_queue.num_waiting

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

  def push(value, non_block = false, timeout: nil) # rubocop:disable Style/OptionalBooleanParameter
    if non_block
      return true if internal_queue.try_push(value)
      raise ThreadError, "queue full"
    end
    timeout.nil? ? internal_queue.push(value) : internal_queue.push(value, timeout:)
  end

  def try_pop(&) = internal_queue.try_pop(&)

  def try_push(value)
    return true if internal_queue.try_push(value)
    block_given? ? yield : false
  end

  def size = internal_queue.size

  def wait_pop(timeout: nil) = internal_queue.wait_pop(timeout:)

  def wait_push(timeout: nil) = internal_queue.wait_push(timeout:)

  # @api private
  def inspect_with(inspector)
    super do
      if closed?
        inspector.breakable " "
        inspector.text "closed"
      else
        inspector.attributes(inspect_info)
      end
    end
  end

  private

  def initialize_copy(_other)
    raise TypeError, "queues cannot be copied"
  end

  def inspect_info
    info = { size:, capacity: }.compact
    info[:num_waiting] = num_waiting if num_waiting.positive?
    info
  end

  private def internal_queue = @queue
end

#wait_pop(timeout: nil) ⇒ Boolean

This method is abstract.

Subclasses may add additional parameters to this method.

Waits until an item is available to pop from the queue, or the timeout expires. Does not remove the item from the queue. Note that this method does not guarantee a subsequent call to #pop or #try_pop will succeed, as another thread may remove the item first.

Parameters:

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

    The maximum time to wait for an item to be available. If nil, the method will wait indefinitely. If 0, the method will not wait at all.

Returns:

  • (Boolean) —

    true if an item is available, false if the timeout expired.

Raises:



147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
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/abstract/queue.rb', line 147

class Queue
  include Internal::MarshalSupport::Reject
  include Internal::Inspect

  # The default mode used to transfer values between Ractors.
  # @return [Symbol]
  def mode = :raise

  # Alias for {#pop} to match the interface of Ruby's Queue class.
  # @return (see #pop)
  def deq(...) = pop(...)
  alias shift deq

  # Alias for {#push} to match the interface of Ruby's Queue class.
  # @return (see #push)
  def enq(...) = push(...)
  alias << enq

  # Alias for {#size} to match the interface of Ruby's Queue class.
  # @return (see #size)
  def length = size

  # The maximum number of items the queue can hold.
  # @return [Integer, Float] the maximum number of items, or `Float::INFINITY` if the queue is unbounded.
  def max = capacity || Float::INFINITY

  # Whether the queue is at capacity.
  # @return [Boolean] `true` if the queue is at capacity, `false` otherwise.
  def full? = capacity && size >= capacity

  def capacity = internal_queue.capacity

  def clear
    internal_queue.clear
    self
  end

  def close
    internal_queue.close
    self
  end

  # Stop new pushes and close after the final queued item is removed.
  # @return [self]
  def seal
    internal_queue.seal
    self
  end

  def closed? = internal_queue.closed?
  def sealed? = internal_queue.sealed?

  # @return [Boolean] Whether enqueue-age tracking is enabled.
  def age_tracking? = internal_queue.age_tracking?

  # @return [Integer, nil] The mutation generation, or `nil` when tracking is disabled.
  def generation = internal_queue.generation

  # The timestamp is monotonic. Its origin is runtime-specific.
  # @return [Float, nil] The monotonic enqueue time of the oldest item.
  def oldest_enqueued_at = internal_queue.oldest_enqueued_at

  # @return [Float, nil] Seconds since the oldest item was enqueued.
  def oldest_age = internal_queue.oldest_age

  def empty? = size.zero?

  def num_waiting = internal_queue.num_waiting

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

  def push(value, non_block = false, timeout: nil) # rubocop:disable Style/OptionalBooleanParameter
    if non_block
      return true if internal_queue.try_push(value)
      raise ThreadError, "queue full"
    end
    timeout.nil? ? internal_queue.push(value) : internal_queue.push(value, timeout:)
  end

  def try_pop(&) = internal_queue.try_pop(&)

  def try_push(value)
    return true if internal_queue.try_push(value)
    block_given? ? yield : false
  end

  def size = internal_queue.size

  def wait_pop(timeout: nil) = internal_queue.wait_pop(timeout:)

  def wait_push(timeout: nil) = internal_queue.wait_push(timeout:)

  # @api private
  def inspect_with(inspector)
    super do
      if closed?
        inspector.breakable " "
        inspector.text "closed"
      else
        inspector.attributes(inspect_info)
      end
    end
  end

  private

  def initialize_copy(_other)
    raise TypeError, "queues cannot be copied"
  end

  def inspect_info
    info = { size:, capacity: }.compact
    info[:num_waiting] = num_waiting if num_waiting.positive?
    info
  end

  private def internal_queue = @queue
end

#wait_push(timeout: nil) ⇒ Boolean

This method is abstract.

Subclasses may add additional parameters to this method.

Waits until there is space to push an item onto the queue, or the timeout expires. Note that this method does not guarantee a subsequent call to #push or #try_push will succeed, as another thread may add an item first.

Parameters:

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

    The maximum time to wait for space to be available. If nil, the method will wait indefinitely. If 0, the method will not wait at all.

Returns:

  • (Boolean) —

    true if space is available, false if the timeout expired.

Raises:



147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
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/abstract/queue.rb', line 147

class Queue
  include Internal::MarshalSupport::Reject
  include Internal::Inspect

  # The default mode used to transfer values between Ractors.
  # @return [Symbol]
  def mode = :raise

  # Alias for {#pop} to match the interface of Ruby's Queue class.
  # @return (see #pop)
  def deq(...) = pop(...)
  alias shift deq

  # Alias for {#push} to match the interface of Ruby's Queue class.
  # @return (see #push)
  def enq(...) = push(...)
  alias << enq

  # Alias for {#size} to match the interface of Ruby's Queue class.
  # @return (see #size)
  def length = size

  # The maximum number of items the queue can hold.
  # @return [Integer, Float] the maximum number of items, or `Float::INFINITY` if the queue is unbounded.
  def max = capacity || Float::INFINITY

  # Whether the queue is at capacity.
  # @return [Boolean] `true` if the queue is at capacity, `false` otherwise.
  def full? = capacity && size >= capacity

  def capacity = internal_queue.capacity

  def clear
    internal_queue.clear
    self
  end

  def close
    internal_queue.close
    self
  end

  # Stop new pushes and close after the final queued item is removed.
  # @return [self]
  def seal
    internal_queue.seal
    self
  end

  def closed? = internal_queue.closed?
  def sealed? = internal_queue.sealed?

  # @return [Boolean] Whether enqueue-age tracking is enabled.
  def age_tracking? = internal_queue.age_tracking?

  # @return [Integer, nil] The mutation generation, or `nil` when tracking is disabled.
  def generation = internal_queue.generation

  # The timestamp is monotonic. Its origin is runtime-specific.
  # @return [Float, nil] The monotonic enqueue time of the oldest item.
  def oldest_enqueued_at = internal_queue.oldest_enqueued_at

  # @return [Float, nil] Seconds since the oldest item was enqueued.
  def oldest_age = internal_queue.oldest_age

  def empty? = size.zero?

  def num_waiting = internal_queue.num_waiting

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

  def push(value, non_block = false, timeout: nil) # rubocop:disable Style/OptionalBooleanParameter
    if non_block
      return true if internal_queue.try_push(value)
      raise ThreadError, "queue full"
    end
    timeout.nil? ? internal_queue.push(value) : internal_queue.push(value, timeout:)
  end

  def try_pop(&) = internal_queue.try_pop(&)

  def try_push(value)
    return true if internal_queue.try_push(value)
    block_given? ? yield : false
  end

  def size = internal_queue.size

  def wait_pop(timeout: nil) = internal_queue.wait_pop(timeout:)

  def wait_push(timeout: nil) = internal_queue.wait_push(timeout:)

  # @api private
  def inspect_with(inspector)
    super do
      if closed?
        inspector.breakable " "
        inspector.text "closed"
      else
        inspector.attributes(inspect_info)
      end
    end
  end

  private

  def initialize_copy(_other)
    raise TypeError, "queues cannot be copied"
  end

  def inspect_info
    info = { size:, capacity: }.compact
    info[:num_waiting] = num_waiting if num_waiting.positive?
    info
  end

  private def internal_queue = @queue
end