Class: Bunny::Queue
- Inherits:
-
Object
- Object
- Bunny::Queue
- Defined in:
- lib/bunny/queue.rb
Overview
Represents AMQP 0.9.1 queue.
Defined Under Namespace
Instance Attribute Summary collapse
-
#channel ⇒ Bunny::Channel
readonly
Channel this queue uses.
-
#name ⇒ String
readonly
Queue name.
-
#options ⇒ Hash
readonly
Options this queue was created with.
Instance Method Summary collapse
-
#arguments ⇒ Hash
Additional optional arguments (typically used by RabbitMQ extensions and plugins).
-
#auto_delete? ⇒ Boolean
True if this queue was declared as automatically deleted (deleted as soon as last consumer unbinds).
-
#bind(exchange, opts = {}) ⇒ Object
Binds queue to an exchange.
-
#consumer_count ⇒ Integer
How many active consumers the queue has.
- #declare! ⇒ Object
-
#delete(opts = {}) ⇒ Object
Deletes the queue.
-
#durable? ⇒ Boolean
True if this queue was declared as durable (will survive broker restart).
-
#exclusive? ⇒ Boolean
True if this queue was declared as exclusive (limited to just one consumer).
-
#initialize(channel, name = AMQ::Protocol::EMPTY_STRING, opts = {}) ⇒ Queue
constructor
A new instance of Queue.
- #inspect ⇒ Object
-
#message_count ⇒ Integer
How many messages the queue has ready (e.g. not delivered but not unacknowledged).
-
#pop(opts = {:manual_ack => false}, &block) ⇒ Array
(also: #get)
Triple of delivery info, message properties and message content.
-
#publish(payload, opts = {}) ⇒ Object
Publishes a message to the queue via default exchange.
-
#purge(opts = {}) ⇒ Object
Purges a queue (removes all messages from it).
- #recover_bindings ⇒ Object
- #recover_from_network_failure ⇒ Object
-
#server_named? ⇒ Boolean
True if this queue was declared as server named.
-
#status ⇒ Hash
A hash with information about the number of queue messages and consumers.
-
#subscribe(opts = { :consumer_tag => @channel.generate_consumer_tag, :manual_ack => false, :exclusive => false, :block => false, :on_cancellation => nil }, &block) ⇒ Object
Adds a consumer to the queue (subscribes for message deliveries).
-
#subscribe_with(consumer, opts = {:block => false}) ⇒ Object
Adds a consumer object to the queue (subscribes for message deliveries).
- #to_s ⇒ Object
-
#unbind(exchange, opts = {}) ⇒ Object
Unbinds queue from an exchange.
Constructor Details
#initialize(channel, name = AMQ::Protocol::EMPTY_STRING, opts = {}) ⇒ Queue
Returns a new instance of Queue.
52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 |
# File 'lib/bunny/queue.rb', line 52 def initialize(channel, name = AMQ::Protocol::EMPTY_STRING, opts = {}) # old Bunny versions pass a connection here. In that case, # we just use default channel from it. MK. @channel = channel @name = name @options = self.class.(name, opts) @durable = @options[:durable] @exclusive = @options[:exclusive] @server_named = @name.empty? @auto_delete = @options[:auto_delete] @type = @options[:type] @arguments = if @type and !@type.empty? then (@options[:arguments] || {}).merge({XArgs::QUEUE_TYPE => @type}) else @options[:arguments] end verify_type!(@arguments) # reassigns updated and verified arguments because Bunny::Channel#declare_queue # accepts a map of options @options[:arguments] = @arguments @bindings = Array.new @default_consumer = nil declare! unless opts[:no_declare] @channel.register_queue(self) end |
Instance Attribute Details
#channel ⇒ Bunny::Channel (readonly)
Returns Channel this queue uses.
32 33 34 |
# File 'lib/bunny/queue.rb', line 32 def channel @channel end |
#name ⇒ String (readonly)
Returns Queue name.
34 35 36 |
# File 'lib/bunny/queue.rb', line 34 def name @name end |
#options ⇒ Hash (readonly)
Returns Options this queue was created with.
36 37 38 |
# File 'lib/bunny/queue.rb', line 36 def @options end |
Instance Method Details
#arguments ⇒ Hash
Returns Additional optional arguments (typically used by RabbitMQ extensions and plugins).
114 115 116 |
# File 'lib/bunny/queue.rb', line 114 def arguments @arguments end |
#auto_delete? ⇒ Boolean
Returns true if this queue was declared as automatically deleted (deleted as soon as last consumer unbinds).
101 102 103 |
# File 'lib/bunny/queue.rb', line 101 def auto_delete? @auto_delete end |
#bind(exchange, opts = {}) ⇒ Object
Binds queue to an exchange
138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 |
# File 'lib/bunny/queue.rb', line 138 def bind(exchange, opts = {}) @channel.queue_bind(@name, exchange, opts) exchange_name = if exchange.respond_to?(:name) exchange.name else exchange end # store bindings for automatic recovery. We need to be very careful to # not cause an infinite rebinding loop here when we recover. MK. binding = { :exchange => exchange_name, :routing_key => (opts[:routing_key] || opts[:key]), :arguments => opts[:arguments] } @bindings.push(binding) unless @bindings.include?(binding) self end |
#consumer_count ⇒ Integer
Returns How many active consumers the queue has.
350 351 352 353 |
# File 'lib/bunny/queue.rb', line 350 def consumer_count s = self.status s[:consumer_count] end |
#declare! ⇒ Object
396 397 398 399 |
# File 'lib/bunny/queue.rb', line 396 def declare! queue_declare_ok = @channel.queue_declare(@name, @options) @name = queue_declare_ok.queue end |
#delete(opts = {}) ⇒ Object
Deletes the queue
320 321 322 323 |
# File 'lib/bunny/queue.rb', line 320 def delete(opts = {}) @channel.deregister_queue(self) @channel.queue_delete(@name, opts) end |
#durable? ⇒ Boolean
Returns true if this queue was declared as durable (will survive broker restart).
87 88 89 |
# File 'lib/bunny/queue.rb', line 87 def durable? @durable end |
#exclusive? ⇒ Boolean
Returns true if this queue was declared as exclusive (limited to just one consumer).
94 95 96 |
# File 'lib/bunny/queue.rb', line 94 def exclusive? @exclusive end |
#inspect ⇒ Object
123 124 125 |
# File 'lib/bunny/queue.rb', line 123 def inspect to_s end |
#message_count ⇒ Integer
Returns How many messages the queue has ready (e.g. not delivered but not unacknowledged).
344 345 346 347 |
# File 'lib/bunny/queue.rb', line 344 def s = self.status s[:message_count] end |
#pop(opts = {:manual_ack => false}, &block) ⇒ Array Also known as: get
Returns Triple of delivery info, message properties and message content. If the queue is empty, all three will be nils.
269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 |
# File 'lib/bunny/queue.rb', line 269 def pop(opts = {:manual_ack => false}, &block) unless opts[:ack].nil? warn "[DEPRECATION] `:ack` is deprecated. Please use `:manual_ack` instead." opts[:manual_ack] = opts[:ack] end get_response, properties, content = @channel.basic_get(@name, opts) if block if properties di = GetResponse.new(get_response, @channel) mp = MessageProperties.new(properties) block.call(di, mp, content) else block.call(nil, nil, nil) end else if properties di = GetResponse.new(get_response, @channel) mp = MessageProperties.new(properties) [di, mp, content] else [nil, nil, nil] end end end |
#publish(payload, opts = {}) ⇒ Object
Publishes a message to the queue via default exchange. Takes the same arguments as Exchange#publish
304 305 306 307 308 |
# File 'lib/bunny/queue.rb', line 304 def publish(payload, opts = {}) @channel.default_exchange.publish(payload, opts.merge(:routing_key => @name)) self end |
#purge(opts = {}) ⇒ Object
Purges a queue (removes all messages from it)
328 329 330 331 332 |
# File 'lib/bunny/queue.rb', line 328 def purge(opts = {}) @channel.queue_purge(@name, opts) self end |
#recover_bindings ⇒ Object
382 383 384 385 386 387 388 |
# File 'lib/bunny/queue.rb', line 382 def recover_bindings @bindings.each do |b| # TODO: inject and use logger # puts "Recovering binding #{b.inspect}" self.bind(b[:exchange], b) end end |
#recover_from_network_failure ⇒ Object
360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 |
# File 'lib/bunny/queue.rb', line 360 def recover_from_network_failure if self.server_named? old_name = @name.dup @name = AMQ::Protocol::EMPTY_STRING @channel.deregister_queue_named(old_name) end # TODO: inject and use logger # puts "Recovering queue #{@name}" begin declare! unless @options[:no_declare] @channel.register_queue(self) rescue Exception => e # TODO: inject and use logger puts "Caught #{e.inspect} while redeclaring and registering #{@name}!" end recover_bindings end |
#server_named? ⇒ Boolean
Returns true if this queue was declared as server named.
108 109 110 |
# File 'lib/bunny/queue.rb', line 108 def server_named? @server_named end |
#status ⇒ Hash
Returns A hash with information about the number of queue messages and consumers.
337 338 339 340 341 |
# File 'lib/bunny/queue.rb', line 337 def status queue_declare_ok = @channel.queue_declare(@name, @options.merge(:passive => true)) {:message_count => queue_declare_ok., :consumer_count => queue_declare_ok.consumer_count} end |
#subscribe(opts = { :consumer_tag => @channel.generate_consumer_tag, :manual_ack => false, :exclusive => false, :block => false, :on_cancellation => nil }, &block) ⇒ Object
Adds a consumer to the queue (subscribes for message deliveries).
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 |
# File 'lib/bunny/queue.rb', line 195 def subscribe(opts = { :consumer_tag => @channel.generate_consumer_tag, :manual_ack => false, :exclusive => false, :block => false, :on_cancellation => nil }, &block) unless opts[:ack].nil? warn "[DEPRECATION] `:ack` is deprecated. Please use `:manual_ack` instead." opts[:manual_ack] = opts[:ack] end ctag = opts.fetch(:consumer_tag, @channel.generate_consumer_tag) consumer = Consumer.new(@channel, self, ctag, !opts[:manual_ack], opts[:exclusive], opts[:arguments]) consumer.on_delivery(&block) consumer.on_cancellation(&opts[:on_cancellation]) if opts[:on_cancellation] @channel.basic_consume_with(consumer) if opts[:block] # joins current thread with the consumers pool, will block # the current thread for as long as the consumer pool is active @channel.work_pool.join end consumer end |
#subscribe_with(consumer, opts = {:block => false}) ⇒ Object
Adds a consumer object to the queue (subscribes for message deliveries).
238 239 240 241 242 243 |
# File 'lib/bunny/queue.rb', line 238 def subscribe_with(consumer, opts = {:block => false}) @channel.basic_consume_with(consumer) @channel.work_pool.join if opts[:block] consumer end |
#to_s ⇒ Object
118 119 120 121 |
# File 'lib/bunny/queue.rb', line 118 def to_s oid = ("0x%x" % (self.object_id << 1)) "<#{self.class.name}:#{oid} @name=\"#{name}\" channel=#{@channel.to_s} @durable=#{@durable} @auto_delete=#{@auto_delete} @exclusive=#{@exclusive} @arguments=#{@arguments}>" end |
#unbind(exchange, opts = {}) ⇒ Object
Unbinds queue from an exchange
167 168 169 170 171 172 173 174 175 176 177 178 179 180 |
# File 'lib/bunny/queue.rb', line 167 def unbind(exchange, opts = {}) @channel.queue_unbind(@name, exchange, opts) exchange_name = if exchange.respond_to?(:name) exchange.name else exchange end @bindings.delete_if { |b| b[:exchange] == exchange_name && b[:routing_key] == (opts[:routing_key] || opts[:key]) && b[:arguments] == opts[:arguments] } self end |