Module: Smash::CloudPowers::Synapse::Queue
- Included in:
- Smash::CloudPowers::SelfAwareness, Smash::CloudPowers::Synapse
- Defined in:
- lib/cloud_powers/synapse/queue.rb
Defined Under Namespace
Classes: Board
Instance Method Summary collapse
- #board_name(url) ⇒ Object
-
#create_queue(name) ⇒ Object
def board_name(url) # TODO: figure out a way to not have this and :name in Board # gets the name from the url if url =~ URI.regexp url = URI.parse(url) url.path.split('/').last.split('_').last else env(url) end end.
- #delete_queue_message(queue, opts = {}) ⇒ Object
- #get_count(board) ⇒ Object
-
#pluck_message(board) ⇒ Object
Params: board
returns a message and deletes it from its origin. - #poll(board, opts = {}) ⇒ Object
- #poller(board_name) ⇒ Object
- #queue_exists?(name) ⇒ Boolean
- #send_queue_message(message, *board_info) ⇒ Object
- #sqs ⇒ Object
Instance Method Details
#board_name(url) ⇒ Object
32 33 34 |
# File 'lib/cloud_powers/synapse/queue.rb', line 32 def board_name(url) url.to_s.split('/').last end |
#create_queue(name) ⇒ Object
def board_name(url)
TODO: figure out a way to not have this and :name in Board
gets the name from the url
if url =~ URI.regexp url = URI.parse(url) url.path.split('/').last.split('_').last else env(url) end end
47 48 49 |
# File 'lib/cloud_powers/synapse/queue.rb', line 47 def create_queue(name) sqs.create_queue(queue_name: to_camel(name)) end |
#delete_queue_message(queue, opts = {}) ⇒ Object
51 52 53 54 55 56 |
# File 'lib/cloud_powers/synapse/queue.rb', line 51 def (queue, opts = {}) poll(queue, opts) do |msg, stats| poller(queue).(msg) throw :stop_polling end end |
#get_count(board) ⇒ Object
58 59 60 61 62 63 |
# File 'lib/cloud_powers/synapse/queue.rb', line 58 def get_count(board) sqs.get_queue_attributes( queue_url: board_name(board), attribute_names: ['ApproximateNumberOfMessages'] ).attributes['ApproximateNumberOfMessages'].to_f end |
#pluck_message(board) ⇒ Object
Params: board
67 68 69 70 71 72 |
# File 'lib/cloud_powers/synapse/queue.rb', line 67 def (board) poll(board) do |msg, poller| poller.(msg) return msg end end |
#poll(board, opts = {}) ⇒ Object
74 75 76 77 78 79 |
# File 'lib/cloud_powers/synapse/queue.rb', line 74 def poll(board, opts = {}) this_poller = poller(board) this_poller.poll(opts) do |msg| yield msg, this_poller if block_given? end end |
#poller(board_name) ⇒ Object
81 82 83 84 85 86 87 88 89 90 91 |
# File 'lib/cloud_powers/synapse/queue.rb', line 81 def poller(board_name) board = Board.new(board_name) unless instance_variable_defined?(board.i_var) instance_variable_set( board.i_var, Aws::SQS::QueuePoller.new(board.address) ) end instance_variable_get(board.i_var) end |
#queue_exists?(name) ⇒ Boolean
93 94 95 |
# File 'lib/cloud_powers/synapse/queue.rb', line 93 def queue_exists?(name) sqs.list_queues(queue_name_prefix: name) end |
#send_queue_message(message, *board_info) ⇒ Object
97 98 99 100 101 102 103 104 |
# File 'lib/cloud_powers/synapse/queue.rb', line 97 def (, *board_info) board = board_info.first.kind_of?(Board) ? board_info.first : Board.new(*board_info) = .to_json unless .kind_of? String sqs.( queue_url: board.address, message_body: ) end |