Module: Smash::CloudPowers::Synapse::Queue
Defined Under Namespace
Classes: Board
Instance Method Summary
collapse
#ec2, #image, #kinesis, #region, #s3, #sqs
Methods included from Zenv
#env_vars, #file_tree_search, #i_vars, #project_root, #project_root=, #system_vars, #zfind
Methods included from Helper
#attr_map!, #called_from, #create_logger, #errors, #format_error_message, #log_file, #logger, #smart_retry, #task_path, #task_require_path, #to_camel, #to_i_var, #to_pascal, #to_ruby_file_name, #to_snake, #update_message_body, #valid_json?
Methods included from Auth
creds
Instance Method Details
#board_name(url) ⇒ Object
62
63
64
|
# File 'lib/cloud_powers/synapse/queue.rb', line 62
def board_name(url)
url.to_s.split('/').last
end
|
#build_queue(name) ⇒ Object
75
76
77
|
# File 'lib/cloud_powers/synapse/queue.rb', line 75
def build_queue(name)
Board.new(sqs, to_camel(name))
end
|
#create_queue!(name) ⇒ Object
66
67
68
69
70
71
72
73
|
# File 'lib/cloud_powers/synapse/queue.rb', line 66
def create_queue!(name)
begin
Board.new(sqs, to_camel(name)).create_queue!
rescue Aws::SQS::Errors::QueueDeletedRecently => e
sleep 5
retry
end
end
|
#delete_queue_message(queue, opts = {}) ⇒ Object
79
80
81
82
83
84
|
# File 'lib/cloud_powers/synapse/queue.rb', line 79
def delete_queue_message(queue, opts = {})
poll(queue, opts) do |msg, stats|
poller(queue).delete_message(msg)
throw :stop_polling
end
end
|
#get_queue_message_count(board_url) ⇒ Object
86
87
88
89
90
91
|
# File 'lib/cloud_powers/synapse/queue.rb', line 86
def get_queue_message_count(board_url)
sqs.get_queue_attributes(
queue_url: board_url,
attribute_names: ['ApproximateNumberOfMessages']
).attributes['ApproximateNumberOfMessages'].to_f
end
|
#pluck_queue_message(board) ⇒ Object
@params: board<String|symbol>: The name of the board
95
96
97
98
99
100
|
# File 'lib/cloud_powers/synapse/queue.rb', line 95
def pluck_queue_message(board)
poll(board) do |msg, poller|
poller.delete_message(msg)
return valid_json?(msg.body) ? JSON.parse(msg.body) : msg.body
end
end
|
#poll(board, opts = {}) ⇒ Object
102
103
104
105
106
107
108
109
110
111
|
# File 'lib/cloud_powers/synapse/queue.rb', line 102
def poll(board, opts = {})
this_poller = poller(board)
results = nil
this_poller.poll(opts) do |msg|
results = yield msg, this_poller if block_given?
this_poller.delete_message(msg)
throw :stop_polling
end
results
end
|
#poller(board_name) ⇒ Object
113
114
115
116
117
118
119
120
121
122
123
|
# File 'lib/cloud_powers/synapse/queue.rb', line 113
def poller(board_name)
board = Board.new(sqs, 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
125
126
127
|
# File 'lib/cloud_powers/synapse/queue.rb', line 125
def queue_exists?(name)
!sqs.list_queues(queue_name_prefix: name).queue_urls.empty?
end
|
#queue_search(name) ⇒ Object
129
130
131
|
# File 'lib/cloud_powers/synapse/queue.rb', line 129
def queue_search(name)
sqs.list_queues(queue_name_prefix: name).queue_urls
end
|
#send_queue_message(address, message) ⇒ Object
133
134
135
136
137
138
|
# File 'lib/cloud_powers/synapse/queue.rb', line 133
def send_queue_message(address, message)
sqs.send_message(
queue_url: address,
message_body: message
)
end
|