Class: RosettaQueue::Gateway::StompAdapter
- Inherits:
-
BaseAdapter
- Object
- BaseAdapter
- RosettaQueue::Gateway::StompAdapter
- Defined in:
- lib/rosetta_queue/adapters/stomp.rb
Instance Method Summary collapse
- #ack(msg) ⇒ Object
- #disconnect(message_handler) ⇒ Object
-
#initialize(adapter_settings = {}) ⇒ StompAdapter
constructor
A new instance of StompAdapter.
- #receive(options) ⇒ Object
- #receive_once(destination, opts) ⇒ Object
- #receive_with(message_handler) ⇒ Object
- #send_message(destination, message, options) ⇒ Object
- #subscribe(destination, options) ⇒ Object
- #unsubscribe(destination) ⇒ Object
Constructor Details
#initialize(adapter_settings = {}) ⇒ StompAdapter
Returns a new instance of StompAdapter.
13 14 15 16 17 18 19 20 |
# File 'lib/rosetta_queue/adapters/stomp.rb', line 13 def initialize(adapter_settings = {}) raise AdapterException, "Missing adapter settings" if adapter_settings.empty? @conn = Stomp::Connection.open(adapter_settings[:user], adapter_settings[:password], adapter_settings[:host], adapter_settings[:port], true) end |
Instance Method Details
#ack(msg) ⇒ Object
8 9 10 11 |
# File 'lib/rosetta_queue/adapters/stomp.rb', line 8 def ack(msg) raise AdapterException, "Unable to ack client because message-id is blank. Are your message handler options correct? (i.e., :ack => 'client')" if msg.headers["message-id"].nil? @conn.ack(msg.headers["message-id"]) end |
#disconnect(message_handler) ⇒ Object
22 23 24 25 |
# File 'lib/rosetta_queue/adapters/stomp.rb', line 22 def disconnect() unsubscribe(destination_for()) @conn.disconnect end |
#receive(options) ⇒ Object
27 28 29 30 31 |
# File 'lib/rosetta_queue/adapters/stomp.rb', line 27 def receive() msg = @conn.receive ack(msg) unless [:ack].nil? msg end |
#receive_once(destination, opts) ⇒ Object
33 34 35 36 37 38 39 |
# File 'lib/rosetta_queue/adapters/stomp.rb', line 33 def receive_once(destination, opts) subscribe(destination, opts) msg = receive(opts).body unsubscribe(destination) RosettaQueue.logger.info("Receiving from #{destination} :: #{msg}") filter_receiving(msg) end |
#receive_with(message_handler) ⇒ Object
41 42 43 44 45 46 47 48 49 50 51 52 53 |
# File 'lib/rosetta_queue/adapters/stomp.rb', line 41 def receive_with() = () destination = destination_for() @conn.subscribe(destination, ) running do msg = receive().body Thread.current[:processing] = true RosettaQueue.logger.info("Receiving from #{destination} :: #{msg}") .(filter_receiving(msg)) Thread.current[:processing] = false end end |
#send_message(destination, message, options) ⇒ Object
55 56 57 58 |
# File 'lib/rosetta_queue/adapters/stomp.rb', line 55 def (destination, , ) RosettaQueue.logger.info("Publishing to #{destination} :: #{}") @conn.send(destination, , ) end |
#subscribe(destination, options) ⇒ Object
60 61 62 |
# File 'lib/rosetta_queue/adapters/stomp.rb', line 60 def subscribe(destination, ) @conn.subscribe(destination, ) end |
#unsubscribe(destination) ⇒ Object
64 65 66 |
# File 'lib/rosetta_queue/adapters/stomp.rb', line 64 def unsubscribe(destination) @conn.unsubscribe(destination) end |