Class: Pwwka::Receiver
- Inherits:
-
Object
- Object
- Pwwka::Receiver
- Extended by:
- Logging
- Defined in:
- lib/pwwka/receiver.rb
Constant Summary
Constants included from Logging
Instance Attribute Summary collapse
-
#channel ⇒ Object
readonly
Returns the value of attribute channel.
-
#channel_connector ⇒ Object
readonly
Returns the value of attribute channel_connector.
-
#queue_name ⇒ Object
readonly
Returns the value of attribute queue_name.
-
#routing_key ⇒ Object
readonly
Returns the value of attribute routing_key.
-
#topic_exchange ⇒ Object
readonly
Returns the value of attribute topic_exchange.
Class Method Summary collapse
Instance Method Summary collapse
- #ack(delivery_tag) ⇒ Object
- #drop_queue ⇒ Object
-
#initialize(queue_name, routing_key, prefetch: Pwwka.configuration.default_prefetch) ⇒ Receiver
constructor
A new instance of Receiver.
- #nack(delivery_tag) ⇒ Object
- #nack_requeue(delivery_tag) ⇒ Object
- #test_teardown ⇒ Object
- #topic_queue ⇒ Object
Methods included from Logging
Constructor Details
#initialize(queue_name, routing_key, prefetch: Pwwka.configuration.default_prefetch) ⇒ Receiver
Returns a new instance of Receiver.
12 13 14 15 16 17 18 |
# File 'lib/pwwka/receiver.rb', line 12 def initialize(queue_name, routing_key, prefetch: Pwwka.configuration.default_prefetch) @queue_name = queue_name @routing_key = routing_key @channel_connector = ChannelConnector.new(prefetch: prefetch, connection_name: "c: #{Pwwka.configuration.app_id} #{Pwwka.configuration.process_name}".strip) @channel = @channel_connector.channel @topic_exchange = @channel_connector.topic_exchange end |
Instance Attribute Details
#channel ⇒ Object (readonly)
Returns the value of attribute channel.
7 8 9 |
# File 'lib/pwwka/receiver.rb', line 7 def channel @channel end |
#channel_connector ⇒ Object (readonly)
Returns the value of attribute channel_connector.
6 7 8 |
# File 'lib/pwwka/receiver.rb', line 6 def channel_connector @channel_connector end |
#queue_name ⇒ Object (readonly)
Returns the value of attribute queue_name.
9 10 11 |
# File 'lib/pwwka/receiver.rb', line 9 def queue_name @queue_name end |
#routing_key ⇒ Object (readonly)
Returns the value of attribute routing_key.
10 11 12 |
# File 'lib/pwwka/receiver.rb', line 10 def routing_key @routing_key end |
#topic_exchange ⇒ Object (readonly)
Returns the value of attribute topic_exchange.
8 9 10 |
# File 'lib/pwwka/receiver.rb', line 8 def topic_exchange @topic_exchange end |
Class Method Details
.subscribe(handler_klass, queue_name, routing_key: "#.#", block: true, prefetch: Pwwka.configuration.default_prefetch, payload_parser: Pwwka.configuration.payload_parser) ⇒ Object
20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 |
# File 'lib/pwwka/receiver.rb', line 20 def self.subscribe(handler_klass, queue_name, routing_key: "#.#", block: true, prefetch: Pwwka.configuration.default_prefetch, payload_parser: Pwwka.configuration.payload_parser) raise "#{handler_klass.name} must respond to `handle!`" unless handler_klass.respond_to?(:handle!) receiver = new(queue_name, routing_key, prefetch: prefetch) begin info "Receiving on #{queue_name}" receiver.topic_queue.subscribe(manual_ack: true, block: block) do |delivery_info, properties, payload| begin payload = payload_parser.(payload) handler_klass.handle!(delivery_info, properties, payload) receiver.ack(delivery_info.delivery_tag) logf "Processed Message on %{queue_name} -> %{payload}, %{routing_key}", queue_name: queue_name, payload: payload, routing_key: delivery_info.routing_key rescue => exception Pwwka::ErrorHandlers::Chain.new( Pwwka.configuration.error_handling_chain ).handle_error( handler_klass, receiver, queue_name, payload, delivery_info, exception) end end rescue Interrupt => _ # TODO: trap TERM within channel.work_pool info "Interrupting queue #{queue_name} subscriber safely" ensure receiver.channel_connector.connection_close end return receiver end |
Instance Method Details
#ack(delivery_tag) ⇒ Object
64 65 66 |
# File 'lib/pwwka/receiver.rb', line 64 def ack(delivery_tag) channel.acknowledge(delivery_tag, false) end |
#drop_queue ⇒ Object
76 77 78 79 |
# File 'lib/pwwka/receiver.rb', line 76 def drop_queue topic_queue.purge topic_queue.delete end |
#nack(delivery_tag) ⇒ Object
68 69 70 |
# File 'lib/pwwka/receiver.rb', line 68 def nack(delivery_tag) channel.nack(delivery_tag, false, false) end |
#nack_requeue(delivery_tag) ⇒ Object
72 73 74 |
# File 'lib/pwwka/receiver.rb', line 72 def nack_requeue(delivery_tag) channel.nack(delivery_tag, false, true) end |
#test_teardown ⇒ Object
81 82 83 84 85 |
# File 'lib/pwwka/receiver.rb', line 81 def test_teardown drop_queue topic_exchange.delete channel_connector.connection_close end |
#topic_queue ⇒ Object
56 57 58 59 60 61 62 |
# File 'lib/pwwka/receiver.rb', line 56 def topic_queue @topic_queue ||= begin queue = channel.queue(queue_name, durable: true, arguments: {}) routing_key.split(',').each { |k| queue.bind(topic_exchange, routing_key: k) } queue end end |