Module: Smash::CloudPowers::Synapse::Queue
- Includes:
- AwsResources, Helper
- Included in:
- Smash::CloudPowers::SelfAwareness, Smash::CloudPowers::Synapse, Board
- Defined in:
- lib/cloud_powers/synapse/queue/board.rb,
lib/cloud_powers/synapse/queue/queue.rb
Defined Under Namespace
Instance Method Summary collapse
-
#board_name(url) ⇒ Object
This method can be used to parse a queue name from its address.
-
#build_queue(name) ⇒ Object
This method builds a Queue::Board object for you to use but doesn't invoke the #create!() method, so no API call is made to create the Queue on SQS.
-
#create_queue!(name) ⇒ Object
This method allows you to create a queue on SQS without explicitly creating a Board object.
-
#delete_queue_message(queue, opts = {}) ⇒ Object
Deletes a queue message without caring about reading/interacting with the message.
-
#get_queue_message_count(board_url) ⇒ Object
This method is used to get the approximate count of messages in a given queue.
-
#pluck_queue_message(board) ⇒ Object
Get a message from a Queue.
-
#poll(board_name, opts = {}) ⇒ Object
Polls the given board with the given options hash and a block that interacts with the message that is retrieved from the queue.
-
#queue_exists?(name) ⇒ Boolean
Checks SQS for the existence of this queue using the #queue_search() method.
-
#queue_poller(board_name) ⇒ Object
This method can be used to gain a SQS::QueuePoller.
-
#queue_search(name) ⇒ Object
Searches for a queue based on the name.
-
#send_queue_message(address, message, this_sqs = sqs) ⇒ Object
Sends a given message to a given queue.
Methods included from Helper
#attr_map!, #available_resources, #called_from, #create_logger, #deep_modify_keys_with, #format_error_message, #log_file, #logger, #modify_keys_with, #smart_retry, #task_path, #task_require_path, #to_camel, #to_hyph, #to_i_var, #to_pascal, #to_ruby_file_name, #to_snake, #update_message_body, #valid_json?, #valid_url?
Methods included from AwsResources
#ec2, #image, #kinesis, #region, #s3, #sns, #sqs
Methods included from Zenv
#env_vars, #file_tree_search, #i_vars, #project_root, #project_root=, #system_vars, #zfind
Methods included from Auth
Instance Method Details
#board_name(url) ⇒ Object
This method can be used to parse a queue name from its address. It can be handy if you need the name of a queue but you don't want the overhead of creating a Board object.
Parameters
- url
String
Returns
String
Example board_name('https://sqs.us-west-53.amazonaws.com/001101010010/fooBar') => fooBar
62 63 64 |
# File 'lib/cloud_powers/synapse/queue/queue.rb', line 62 def board_name(url) url.to_s.split('/').last end |
#build_queue(name) ⇒ Object
This method builds a Queue::Board object for you to use but doesn't invoke the #create!() method, so no API call is made to create the Queue on SQS. This can be used if the Board and/or Queue already exists.
Parameters
- name
String- name of the Queue you want to interact with
Returns Queue::Board
Example queue_object = build_queue('exampleQueue') queue_object.address => https://sqs.us-west-2.amazonaws.com/81234567/exampleQueue
80 81 82 |
# File 'lib/cloud_powers/synapse/queue/queue.rb', line 80 def build_queue(name) Smash::CloudPowers::Queue::Board.build(to_camel(name), sqs) end |
#create_queue!(name) ⇒ Object
This method allows you to create a queue on SQS without explicitly creating a Board object
Parameters
- name
String- The name of the Queue to be created
Returns Queue::Board
Example create_queue('exampleQueue') get_queue_message_count
95 96 97 98 99 100 101 102 |
# File 'lib/cloud_powers/synapse/queue/queue.rb', line 95 def create_queue!(name) begin Smash::CloudPowers::Queue::Board.create!(to_camel(name), sqs) rescue Aws::SQS::Errors::QueueDeletedRecently => e sleep 5 retry end end |
#delete_queue_message(queue, opts = {}) ⇒ Object
Deletes a queue message without caring about reading/interacting with the message. This is usually used for progress tracking, ie; a Neuron takes a message from the Backlog, moves it to WIP and deletes it from Backlog. Then repeats these steps for the remaining States in the Workflow
Parameters
- queue
String- queue is the name of theQueueto interact with - opts
Hash(optional) - a configurationHashfor theSQS::QueuePoller
Notes
- throws :stop_polling after the message is deleted
Example get_queue_message_count('exampleQueue')
=> n
delete_queue_message('exampleQueue') get_queue_message_count('exampleQueue')
=> n-1
121 122 123 124 125 126 |
# File 'lib/cloud_powers/synapse/queue/queue.rb', line 121 def (queue, opts = {}) poll(queue, opts) do |msg, stats| poller(queue).(msg) throw :stop_polling end end |
#get_queue_message_count(board_url) ⇒ Object
This method is used to get the approximate count of messages in a given queue
Parameters
- board_url
String- the URL for the board you need to get a count from
Returns
Float
Example get_queue_message_count('exampleQueue')
=> n
delete_queue_message('exampleQueue') get_queue_message_count('exampleQueue')
=> n-1
142 143 144 145 146 147 |
# File 'lib/cloud_powers/synapse/queue/queue.rb', line 142 def (board_url) sqs.get_queue_attributes( queue_url: board_url, attribute_names: ['ApproximateNumberOfMessages'] ).attributes['ApproximateNumberOfMessages'].to_f end |
#pluck_queue_message(board) ⇒ Object
Get a message from a Queue
Parameters
- board<String|symbol>: The name of the board
Returns
Stringifmsg.bodyis not valid JSONHashifmsg.bodyis valid JSON
Example
msg.body == 'Hey' # String
pluck_queue_message('exampleQueue')
=> 'Hey' # String
# msg.body == "\{"tally":"ho"\}" # +JSON+
('exampleQueue')
# => { 'tally' => 'ho' } # +Hash+
166 167 168 169 170 171 |
# File 'lib/cloud_powers/synapse/queue/queue.rb', line 166 def (board) poll(board) do |msg, poller| poller.(msg) return valid_json?(msg.body) ? JSON.parse(msg.body) : msg.body end end |
#poll(board_name, opts = {}) ⇒ Object
Polls the given board with the given options hash and a block that interacts with the message that is retrieved from the queue
Parameters
boardString- the name of the queue you want to polloptsHash- costomizes the Aws::SQS::QueuePoller's #poll(opts) method and can have anyAWS::SQS:QueuePollerpolling configuration option(s)blockis the block that is used to interact with the message that was retrieved
Returns
the results from the message and the block that interacts with the message(s)
Example
continuously run jobs from messages in the Queue and leaves the message in the queue
using the :skip_delete parameter
poll(:backlog, :skip_delete) do |msg| demo_job = Job.new(msg.body) demo_job.run end
192 193 194 195 196 197 198 199 200 201 |
# File 'lib/cloud_powers/synapse/queue/queue.rb', line 192 def poll(board_name, opts = {}) this_poller = queue_poller(board_name) results = nil this_poller.poll(opts) do |msg| results = yield msg, this_poller if block_given? this_poller.(msg) throw :stop_polling end results end |
#queue_exists?(name) ⇒ Boolean
Checks SQS for the existence of this queue using the #queue_search() method
Parameters
- name
String
Returns Boolean
Notes:
* see <tt>#queue_search()</tt>
242 243 244 |
# File 'lib/cloud_powers/synapse/queue/queue.rb', line 242 def queue_exists?(name) !queue_search(name).empty? end |
#queue_poller(board_name) ⇒ Object
This method can be used to gain a SQS::QueuePoller. It creates a Board object, the Board then sends the API call to SQS to create the queue and sets an instance variable, using the board's name, to the Board object itself
Parameters
- board_name
String- name of the Queue you want to gain a QueuePoller for
Returns board_name:Queue::Board
Notes
- An instance variable is set with this method, if one doesn't exist for the board The instance variable that is created/used is named the same name that was given as a parameter.
Example
these are equivalent after @exp_queue_poller is set but before it is set,
exp_queue_poller
queue_poller('exampleQueue').poll { |msg| Job.new(msg.body).run } @example_queue.poll { |msg| Job.new(msg.body).run }
223 224 225 226 227 228 229 230 |
# File 'lib/cloud_powers/synapse/queue/queue.rb', line 223 def queue_poller(board_name) board = Smash::CloudPowers::Queue::Board.create!(board_name, sqs) unless instance_variable_defined?(board.i_var) instance_variable_set(board.i_var, board) end instance_variable_get(board.i_var).poller end |
#queue_search(name) ⇒ Object
Searches for a queue based on the name
Parameters
name String
Returns
queue_urls String
Example results = queue_search('exampleQueue') # returns related URLs results.first =~ /exampleQueue/ # regex match against the URL
257 258 259 |
# File 'lib/cloud_powers/synapse/queue/queue.rb', line 257 def queue_search(name) sqs.list_queues(queue_name_prefix: name).queue_urls end |
#send_queue_message(address, message, this_sqs = sqs) ⇒ Object
Sends a given message to a given queue
Parameters
- address
String- address of the Queue you want to interact with - message
String- message to be sent
Returns
Array
Example legit_address = 'https://sqs.us-west-2.amazonaws.com/12345678/exampleQueue' random_message = 'Wowza, this is pretty easy.' resp = send_queue_message(legit_address, random_message)) resp.message_id => 'some message id'
276 277 278 |
# File 'lib/cloud_powers/synapse/queue/queue.rb', line 276 def (address, , this_sqs = sqs) this_sqs.(queue_url: address, message_body: ) end |