Module: Smash::CloudPowers::Synapse::Queue

Includes:
AwsResources, Helper
Included in:
Smash::CloudPowers::SelfAwareness, Smash::CloudPowers::Synapse, Board
Defined in:
lib/cloud_powers/synapse/queue.rb

Defined Under Namespace

Classes: Board

Instance Method Summary collapse

Methods included from AwsResources

#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

Returns:

  • (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