Class: Farce::Abstract::Queue Abstract
- Inherits:
-
Object
- Object
- Farce::Abstract::Queue
- Includes:
- Internal::Inspect, Internal::MarshalSupport::Reject
- Defined in:
- lib/farce/abstract/queue.rb,
lib/farce/integrations/active_support/duplicable.rb
Overview
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.
Direct Known Subclasses
PriorityQueue, TimerQueue, Local::Queue, Queue, Strict::Queue, Unshared::Queue
Instance Attribute Summary collapse
-
#capacity ⇒ Integer?
readonly
abstract
The maximum number of items the queue can hold.
-
#mode ⇒ Symbol
readonly
The default mode used to transfer values between Ractors.
-
#num_waiting ⇒ Integer
readonly
abstract
The number of threads/fibers currently waiting on the queue.
-
#size ⇒ Integer
readonly
abstract
The current number of items in the queue.
ActiveSupport Integration collapse
-
#duplicable? ⇒ Boolean
False.
Instance Method Summary collapse
-
#age_tracking? ⇒ Boolean
Whether enqueue-age tracking is enabled.
-
#clear ⇒ self
abstract
Removes all items from the queue.
-
#close ⇒ self
abstract
Closes the queue, preventing any further items from being added.
-
#closed? ⇒ Boolean
abstract
Checks whether the queue is closed.
-
#deq ⇒ BasicObject
(also: #shift)
Alias for #pop to match the interface of Ruby's Queue class.
-
#empty? ⇒ Boolean
abstract
Checks whether the queue is empty.
-
#enq ⇒ BasicObject
(also: #<<)
Alias for #push to match the interface of Ruby's Queue class.
-
#full? ⇒ Boolean
Whether the queue is at capacity.
-
#generation ⇒ Integer?
The mutation generation, or
nilwhen tracking is disabled. - #initialize(capacity: 1024, track_age: false) ⇒ BasicObject constructor abstract
-
#length ⇒ Integer
Alias for #size to match the interface of Ruby's Queue class.
-
#max ⇒ Integer, Float
The maximum number of items the queue can hold.
-
#oldest_age ⇒ Float?
Seconds since the oldest item was enqueued.
-
#oldest_enqueued_at ⇒ Float?
The timestamp is monotonic.
-
#pop(non_block = false, timeout: nil) ⇒ BasicObject
rubocop:disable Style/OptionalBooleanParameter.
-
#push(value, non_block = false, timeout: nil) ⇒ BasicObject
rubocop:disable Style/OptionalBooleanParameter.
-
#seal ⇒ self
Stop new pushes and close after the final queued item is removed.
-
#sealed? ⇒ Boolean
Whether the queue rejects new pushes.
-
#try_pop { ... } ⇒ BasicObject?
abstract
Tries to take an item from the queue without blocking.
-
#try_push(value) { ... } ⇒ Boolean
abstract
Tries to push an item onto the queue without blocking.
-
#wait_pop(timeout: nil) ⇒ Boolean
abstract
Waits until an item is available to pop from the queue, or the timeout expires.
-
#wait_push(timeout: nil) ⇒ Boolean
abstract
Waits until there is space to push an item onto the queue, or the timeout expires.
Constructor Details
#initialize(capacity: 1024, track_age: false) ⇒ BasicObject
Subclasses may add additional parameters to this method.
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)
Returns 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.
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)
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.
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)
Returns 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.
200 |
# File 'lib/farce/abstract/queue.rb', line 200 def age_tracking? = internal_queue.age_tracking? |
#clear ⇒ self
Removes all items from 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 |
#close ⇒ self
Closes the queue, preventing any further items from being added.
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
Checks whether the queue is closed.
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
This methods is only available if ActiveSupport has been loaded.
Returns false.
37 |
# File 'lib/farce/integrations/active_support/duplicable.rb', line 37 def duplicable? = false |
#empty? ⇒ Boolean
Checks whether the queue is empty.
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.
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.
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.
167 |
# File 'lib/farce/abstract/queue.rb', line 167 def length = size |
#max ⇒ Integer, Float
The maximum number of items the queue can hold.
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.
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.
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.
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.
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?
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.
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
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.
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
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.
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
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.
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 |