Class: AmqpAdapter

Inherits:
Object
  • Object
show all
Defined in:
lib/orq/adapters/amqp_adapter.rb

Overview

Adapter providing access to AMQP based message systems from ORQ

Instance Method Summary collapse

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