Class: Hutch::Broker
Constant Summary collapse
- DEFAULT_AMQP_PORT =
case RUBY_ENGINE when "jruby" then com.rabbitmq.client.ConnectionFactory::DEFAULT_AMQP_PORT when "ruby" then AMQ::Protocol::DEFAULT_PORT end
- DEFAULT_AMQPS_PORT =
case RUBY_ENGINE when "jruby" then com.rabbitmq.client.ConnectionFactory::DEFAULT_AMQP_OVER_SSL_PORT when "ruby" then AMQ::Protocol::TLS_PORT end
Instance Attribute Summary collapse
-
#api_client ⇒ Object
Returns the value of attribute api_client.
-
#channel ⇒ Object
Returns the value of attribute channel.
-
#connection ⇒ Object
Returns the value of attribute connection.
-
#exchange ⇒ Object
Returns the value of attribute exchange.
Instance Method Summary collapse
- #ack(delivery_tag, channel: self.channel) ⇒ Object
-
#bind_queue(queue, routing_keys) ⇒ Object
Bind a queue to the broker's exchange on the routing keys provided.
-
#bindings ⇒ Object
Return a mapping of queue names to the routing keys they're bound to.
- #confirm_select(*args) ⇒ Object
-
#connect(options = {}) ⇒ Object
Connect to broker.
- #declare_exchange(ch = channel) ⇒ Object
- #declare_exchange!(*args) ⇒ Object
- #declare_publisher! ⇒ Object
- #disconnect ⇒ Object
- #http_api_use_enabled? ⇒ Boolean
-
#initialize(config = nil) ⇒ Broker
constructor
A new instance of Broker.
- #nack(delivery_tag, channel: self.channel) ⇒ Object
-
#namespaced_queue_name(name) ⇒ Object
Apply the configured namespace prefix to a queue name.
- #open_channel ⇒ Object
- #open_channel! ⇒ Object
- #open_connection ⇒ Object
- #open_connection! ⇒ Object
- #publish(*args) ⇒ Object
-
#queue(name, options = {}) ⇒ Object
Create / get a durable queue.
- #queue_exists?(name) ⇒ Boolean
- #reject(delivery_tag, requeue = false, channel: self.channel) ⇒ Object
-
#replace_channel! ⇒ Object
Closing the old channel makes RabbitMQ forget the consumers on it and requeue their unacknowledged deliveries.
-
#requeue(delivery_tag, channel: self.channel) ⇒ Object
Delivery tags are scoped to their channel: a tag from a replaced one is dropped, the server has requeued that delivery anyway.
-
#set_up_amqp_connection ⇒ Object
Connect to RabbitMQ via AMQP.
-
#set_up_api_connection ⇒ Object
Set up the connection to the RabbitMQ management API.
- #stop ⇒ Object
- #tracing_enabled? ⇒ Boolean
-
#unbind_redundant_bindings(queue, routing_keys) ⇒ Object
Find the existing bindings, and unbind any redundant bindings.
-
#using_publisher_confirmations? ⇒ Boolean
True if channel is set up to use publisher confirmations.
- #wait_for_confirms ⇒ Object
Methods included from Logging
logger, #logger, logger=, setup_logger
Constructor Details
Instance Attribute Details
#api_client ⇒ Object
Returns the value of attribute api_client.
13 14 15 |
# File 'lib/hutch/broker.rb', line 13 def api_client @api_client end |
#channel ⇒ Object
Returns the value of attribute channel.
13 14 15 |
# File 'lib/hutch/broker.rb', line 13 def channel @channel end |
#connection ⇒ Object
Returns the value of attribute connection.
13 14 15 |
# File 'lib/hutch/broker.rb', line 13 def connection @connection end |
#exchange ⇒ Object
Returns the value of attribute exchange.
13 14 15 |
# File 'lib/hutch/broker.rb', line 13 def exchange @exchange end |
Instance Method Details
#ack(delivery_tag, channel: self.channel) ⇒ Object
261 262 263 |
# File 'lib/hutch/broker.rb', line 261 def ack(delivery_tag, channel: self.channel) channel.ack(delivery_tag, false) if current_channel?(channel) end |
#bind_queue(queue, routing_keys) ⇒ Object
Bind a queue to the broker's exchange on the routing keys provided. Any existing bindings on the queue that aren't present in the array of routing keys will be unbound.
233 234 235 236 237 238 239 240 241 |
# File 'lib/hutch/broker.rb', line 233 def bind_queue(queue, routing_keys) unbind_redundant_bindings(queue, routing_keys) # Ensure all the desired bindings are present routing_keys.each do |routing_key| logger.debug "creating binding #{queue.name} <--> #{routing_key}" queue.bind(exchange, routing_key: routing_key) end end |
#bindings ⇒ Object
Return a mapping of queue names to the routing keys they're bound to.
203 204 205 206 207 208 209 210 211 212 213 214 215 |
# File 'lib/hutch/broker.rb', line 203 def bindings results = Hash.new { |hash, key| hash[key] = [] } filtered = api_client.bindings. reject { |b| b['destination'] == b['routing_key'] }. select { |b| b['source'] == @config[:mq_exchange] && b['vhost'] == vhost } filtered.each do |binding| results[binding['destination']] << binding['routing_key'] end results end |
#confirm_select(*args) ⇒ Object
273 274 275 |
# File 'lib/hutch/broker.rb', line 273 def confirm_select(*args) channel.confirm_select(*args) end |
#connect(options = {}) ⇒ Object
Connect to broker
47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 |
# File 'lib/hutch/broker.rb', line 47 def connect( = {}) @options = set_up_amqp_connection if http_api_use_enabled? logger.info "HTTP API use is enabled" set_up_api_connection else logger.info "HTTP API use is disabled" end if tracing_enabled? logger.info "tracing is enabled using #{@config[:tracer]}" else logger.info "tracing is disabled" end if block_given? begin yield ensure disconnect end end end |
#declare_exchange(ch = channel) ⇒ Object
135 136 137 138 139 140 141 142 143 144 |
# File 'lib/hutch/broker.rb', line 135 def declare_exchange(ch = channel) exchange_name = @config[:mq_exchange] exchange_type = @config[:mq_exchange_type] = { durable: true }.merge(@config[:mq_exchange_options]) logger.info "using topic exchange '#{exchange_name}'" with_bunny_precondition_handler('exchange') do Adapter.new_exchange(ch, exchange_type, exchange_name, ) end end |
#declare_exchange!(*args) ⇒ Object
146 147 148 |
# File 'lib/hutch/broker.rb', line 146 def declare_exchange!(*args) @exchange = declare_exchange(*args) end |
#declare_publisher! ⇒ Object
150 151 152 |
# File 'lib/hutch/broker.rb', line 150 def declare_publisher! @publisher = Hutch::Publisher.new(connection, channel, exchange, @config) end |
#disconnect ⇒ Object
72 73 74 75 76 77 78 79 |
# File 'lib/hutch/broker.rb', line 72 def disconnect @channel.close if @channel @connection.close if @connection @channel = nil @connection = nil @exchange = nil @api_client = nil end |
#http_api_use_enabled? ⇒ Boolean
170 171 172 173 174 175 176 177 178 179 |
# File 'lib/hutch/broker.rb', line 170 def http_api_use_enabled? op = @options.fetch(:enable_http_api_use, true) cf = if @config[:enable_http_api_use].nil? true else @config[:enable_http_api_use] end op && cf end |
#nack(delivery_tag, channel: self.channel) ⇒ Object
265 266 267 |
# File 'lib/hutch/broker.rb', line 265 def nack(delivery_tag, channel: self.channel) channel.nack(delivery_tag, false, false) if current_channel?(channel) end |
#namespaced_queue_name(name) ⇒ Object
Apply the configured namespace prefix to a queue name.
197 198 199 200 |
# File 'lib/hutch/broker.rb', line 197 def namespaced_queue_name(name) namespace = @config[:namespace].to_s.downcase.gsub(/[^-_:\.\w]/, "") namespace.present? ? "#{namespace}:#{name}" : name end |
#open_channel ⇒ Object
109 110 111 112 113 114 115 116 117 118 119 120 |
# File 'lib/hutch/broker.rb', line 109 def open_channel logger.info "opening rabbitmq channel with pool size #{consumer_pool_size}, abort on exception #{consumer_pool_abort_on_exception}" connection.create_channel(nil, consumer_pool_size, consumer_pool_abort_on_exception).tap do |ch| connection.prefetch_channel(ch, @config[:channel_prefetch]) if @config[:publisher_confirms] || @config[:force_publisher_confirms] logger.info 'enabling publisher confirms' ch.confirm_select end connection.install_channel_recovery(ch) end end |
#open_channel! ⇒ Object
122 123 124 |
# File 'lib/hutch/broker.rb', line 122 def open_channel! @channel = open_channel end |
#open_connection ⇒ Object
92 93 94 95 96 97 98 99 100 101 102 103 |
# File 'lib/hutch/broker.rb', line 92 def open_connection logger.info "connecting to rabbitmq (#{sanitized_uri})" connection = Hutch::Adapter.new(connection_params) with_bunny_connection_handler(sanitized_uri) do connection.start end logger.info "connected to RabbitMQ at #{connection_params[:host]} as #{connection_params[:username]}" connection end |
#open_connection! ⇒ Object
105 106 107 |
# File 'lib/hutch/broker.rb', line 105 def open_connection! @connection = open_connection end |
#publish(*args) ⇒ Object
269 270 271 |
# File 'lib/hutch/broker.rb', line 269 def publish(*args) @publisher.publish(*args) end |
#queue(name, options = {}) ⇒ Object
Create / get a durable queue.
186 187 188 189 190 |
# File 'lib/hutch/broker.rb', line 186 def queue(name, = {}) with_bunny_precondition_handler('queue') do channel.queue(name, **) end end |
#queue_exists?(name) ⇒ Boolean
192 193 194 |
# File 'lib/hutch/broker.rb', line 192 def queue_exists?(name) connection.queue_exists?(name) end |
#reject(delivery_tag, requeue = false, channel: self.channel) ⇒ Object
257 258 259 |
# File 'lib/hutch/broker.rb', line 257 def reject(delivery_tag, requeue = false, channel: self.channel) channel.reject(delivery_tag, requeue) if current_channel?(channel) end |
#replace_channel! ⇒ Object
Closing the old channel makes RabbitMQ forget the consumers on it and requeue their unacknowledged deliveries.
128 129 130 131 132 133 |
# File 'lib/hutch/broker.rb', line 128 def replace_channel! close_consumer_channel open_channel! declare_exchange! declare_publisher! end |
#requeue(delivery_tag, channel: self.channel) ⇒ Object
Delivery tags are scoped to their channel: a tag from a replaced one is dropped, the server has requeued that delivery anyway.
253 254 255 |
# File 'lib/hutch/broker.rb', line 253 def requeue(delivery_tag, channel: self.channel) channel.reject(delivery_tag, true) if current_channel?(channel) end |
#set_up_amqp_connection ⇒ Object
Connect to RabbitMQ via AMQP
This sets up the main connection and channel we use for talking to RabbitMQ. It also ensures the existence of the exchange we'll be using.
85 86 87 88 89 90 |
# File 'lib/hutch/broker.rb', line 85 def set_up_amqp_connection open_connection! open_channel! declare_exchange! declare_publisher! end |
#set_up_api_connection ⇒ Object
Set up the connection to the RabbitMQ management API. Unfortunately, this is necessary to do a few things that are impossible over AMQP. E.g. listing queues and bindings.
157 158 159 160 161 162 163 164 165 166 167 168 |
# File 'lib/hutch/broker.rb', line 157 def set_up_api_connection logger.info "connecting to rabbitmq HTTP API (#{api_config.sanitized_uri})" with_authentication_error_handler do with_connection_error_handler do @api_client = CarrotTop.new(host: api_config.host, port: api_config.port, user: api_config.username, password: api_config.password, ssl: api_config.ssl) @api_client.exchanges end end end |
#stop ⇒ Object
243 244 245 246 247 248 249 |
# File 'lib/hutch/broker.rb', line 243 def stop if defined?(JRUBY_VERSION) channel.close else drain_consumer_work_pool end end |
#tracing_enabled? ⇒ Boolean
181 182 183 |
# File 'lib/hutch/broker.rb', line 181 def tracing_enabled? @config[:tracer] && @config[:tracer] != Hutch::Tracers::NullTracer end |
#unbind_redundant_bindings(queue, routing_keys) ⇒ Object
Find the existing bindings, and unbind any redundant bindings
218 219 220 221 222 223 224 225 226 227 228 |
# File 'lib/hutch/broker.rb', line 218 def unbind_redundant_bindings(queue, routing_keys) return unless http_api_use_enabled? filtered = bindings.select { |dest, keys| dest == queue.name } filtered.each do |dest, keys| keys.reject { |key| routing_keys.include?(key) }.each do |key| logger.debug "removing redundant binding #{queue.name} <--> #{key}" queue.unbind(exchange, routing_key: key) end end end |
#using_publisher_confirmations? ⇒ Boolean
Returns True if channel is set up to use publisher confirmations.
282 283 284 |
# File 'lib/hutch/broker.rb', line 282 def using_publisher_confirmations? channel.using_publisher_confirmations? end |
#wait_for_confirms ⇒ Object
277 278 279 |
# File 'lib/hutch/broker.rb', line 277 def wait_for_confirms channel.wait_for_confirms end |