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

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 delete_queue_message(queue, opts = {})
  poll(queue, opts) do |msg, stats|
    poller(queue).delete_message(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 returns a message and deletes it from its origin



67
68
69
70
71
72
# File 'lib/cloud_powers/synapse/queue.rb', line 67

def pluck_message(board)
  poll(board) do |msg, poller|
    poller.delete_message(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

Returns:

  • (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 send_queue_message(message, *board_info)
  board = board_info.first.kind_of?(Board) ? board_info.first : Board.new(*board_info)
  message = message.to_json unless message.kind_of? String
  sqs.send_message(
    queue_url: board.address,
    message_body: message
  )
end

#sqs ⇒ Object



106
107
108
# File 'lib/cloud_powers/synapse/queue.rb', line 106

def sqs
  @sqs ||= Aws::SQS::Client.new(credentials: Auth.creds)
end