Class: AmqpAdapter
- Inherits:
-
Object
- Object
- AmqpAdapter
- Defined in:
- lib/orq/adapters/amqp_adapter.rb
Overview
Adapter providing access to AMQP based message systems from ORQ
Instance Method Summary collapse
- #fire(impulseUri, content) ⇒ Object
-
#initialize(config) ⇒ AmqpAdapter
constructor
A new instance of AmqpAdapter.
- #start ⇒ Object
- #subscribe(impulseType, &block) ⇒ Object
Constructor Details
#initialize(config) ⇒ AmqpAdapter
Returns a new instance of AmqpAdapter.
5 6 7 |
# File 'lib/orq/adapters/amqp_adapter.rb', line 5 def initialize(config) @config = config end |
Instance Method Details
#fire(impulseUri, content) ⇒ Object
28 29 30 31 |
# File 'lib/orq/adapters/amqp_adapter.rb', line 28 def fire(impulseUri, content) exchange = MQ::Exchange.new @channel, :direct, impulseUri exchange.publish content, :content_type => 'application/json' end |
#start ⇒ Object
9 10 11 12 13 |
# File 'lib/orq/adapters/amqp_adapter.rb', line 9 def start # Do connection @connection = AMQP.connect :host => (@config['host'] || 'localhost'), :port => (@config['port'] || 5672).to_i @channel = MQ.new @connection end |
#subscribe(impulseType, &block) ⇒ Object
15 16 17 18 19 20 21 22 23 24 25 26 |
# File 'lib/orq/adapters/amqp_adapter.rb', line 15 def subscribe(impulseType, &block) raise "No block provided - a block must be provided to the subscribe method" unless block_given? queue = MQ::Queue.new @channel, impulseType.uri exchange = MQ::Exchange.new @channel, :direct, impulseType.uri queue.bind(exchange) queue.subscribe :ack => true do |headers, msg| yield impulseType.load(msg) headers.ack end end |