2015-06-19 12:51:54 -04:00
|
|
|
module Fog
|
|
|
|
module AWS
|
|
|
|
class Kinesis < Fog::Service
|
|
|
|
extend Fog::AWS::CredentialFetcher::ServiceMethods
|
|
|
|
|
2015-06-24 11:54:30 -04:00
|
|
|
class LimitExceeded < Fog::Errors::Error; end
|
|
|
|
class ResourceNotFound < Fog::Errors::Error; end
|
|
|
|
class ExpiredIterator < Fog::Errors::Error; end
|
|
|
|
class InvalidArgument < Fog::Errors::Error; end
|
|
|
|
class ProvisionedThroughputExceeded < Fog::Errors::Error; end
|
|
|
|
|
2015-06-19 12:51:54 -04:00
|
|
|
requires :aws_access_key_id, :aws_secret_access_key
|
|
|
|
recognizes :region, :host, :path, :port, :scheme, :persistent, :use_iam_profile, :aws_session_token, :aws_credentials_expire_at, :instrumentor, :instrumentor_name
|
|
|
|
|
|
|
|
request_path 'fog/aws/requests/kinesis'
|
|
|
|
|
|
|
|
request :list_streams
|
2015-06-19 16:55:32 -04:00
|
|
|
request :describe_stream
|
|
|
|
request :create_stream
|
|
|
|
request :delete_stream
|
2015-06-19 17:15:33 -04:00
|
|
|
request :get_shard_iterator
|
2015-06-19 17:30:10 -04:00
|
|
|
request :put_records
|
2015-07-01 13:13:23 -04:00
|
|
|
request :put_record
|
2015-06-23 11:03:25 -04:00
|
|
|
request :get_records
|
2015-07-02 15:23:25 -04:00
|
|
|
request :split_shard
|
2015-07-02 15:57:29 -04:00
|
|
|
request :merge_shards
|
2015-07-02 16:12:09 -04:00
|
|
|
request :add_tags_to_stream
|
2015-07-02 16:21:18 -04:00
|
|
|
request :list_tags_for_stream
|
2015-07-02 16:52:41 -04:00
|
|
|
request :remove_tags_from_stream
|
2015-06-19 12:51:54 -04:00
|
|
|
|
|
|
|
class Real
|
|
|
|
include Fog::AWS::CredentialFetcher::ConnectionMethods
|
|
|
|
|
|
|
|
def initialize(options={})
|
|
|
|
@use_iam_profile = options[:use_iam_profile]
|
|
|
|
|
|
|
|
@connection_options = options[:connection_options] || {}
|
|
|
|
|
|
|
|
@instrumentor = options[:instrumentor]
|
|
|
|
@instrumentor_name = options[:instrumentor_name] || 'fog.aws.kinesis'
|
|
|
|
|
|
|
|
options[:region] ||= 'us-east-1'
|
|
|
|
@region = options[:region]
|
|
|
|
@host = options[:host] || "kinesis.#{options[:region]}.amazonaws.com"
|
|
|
|
@path = options[:path] || '/'
|
|
|
|
@persistent = options[:persistent] || false
|
|
|
|
@port = options[:port] || 443
|
|
|
|
@scheme = options[:scheme] || 'https'
|
|
|
|
@connection = Fog::XML::Connection.new("#{@scheme}://#{@host}:#{@port}#{@path}", @persistent, @connection_options)
|
2015-06-24 16:41:21 -04:00
|
|
|
@version = "20131202"
|
2015-06-19 12:51:54 -04:00
|
|
|
|
|
|
|
setup_credentials(options)
|
|
|
|
end
|
|
|
|
|
|
|
|
private
|
|
|
|
|
|
|
|
def setup_credentials(options)
|
|
|
|
@aws_access_key_id = options[:aws_access_key_id]
|
|
|
|
@aws_secret_access_key = options[:aws_secret_access_key]
|
|
|
|
@aws_session_token = options[:aws_session_token]
|
|
|
|
@aws_credentials_expire_at = options[:aws_credentials_expire_at]
|
|
|
|
|
|
|
|
@signer = Fog::AWS::SignatureV4.new( @aws_access_key_id, @aws_secret_access_key, @region, 'kinesis')
|
|
|
|
end
|
|
|
|
|
|
|
|
def request(params)
|
|
|
|
refresh_credentials_if_expired
|
|
|
|
idempotent = params.delete(:idempotent)
|
|
|
|
parser = params.delete(:parser)
|
|
|
|
|
|
|
|
date = Fog::Time.now
|
|
|
|
headers = {
|
2015-06-25 16:22:39 -04:00
|
|
|
'X-Amz-Target' => params['X-Amz-Target'],
|
|
|
|
'Content-Type' => 'application/x-amz-json-1.1',
|
|
|
|
'Host' => @host,
|
|
|
|
'x-amz-date' => date.to_iso8601_basic
|
|
|
|
}
|
2015-06-19 12:51:54 -04:00
|
|
|
headers['x-amz-security-token'] = @aws_session_token if @aws_session_token
|
|
|
|
body = MultiJson.dump(params[:body])
|
2015-07-01 11:16:30 -04:00
|
|
|
headers['Authorization'] = @signer.sign({:method => "POST", :headers => headers, :body => body, :query => {}, :path => @path}, date)
|
2015-06-19 12:51:54 -04:00
|
|
|
|
|
|
|
if @instrumentor
|
|
|
|
@instrumentor.instrument("#{@instrumentor_name}.request", params) do
|
|
|
|
_request(body, headers, idempotent, parser)
|
|
|
|
end
|
|
|
|
else
|
|
|
|
_request(body, headers, idempotent, parser)
|
|
|
|
end
|
|
|
|
end
|
|
|
|
|
|
|
|
def _request(body, headers, idempotent, parser)
|
|
|
|
@connection.request({
|
2015-06-25 16:22:39 -04:00
|
|
|
:body => body,
|
|
|
|
:expects => 200,
|
|
|
|
:headers => headers,
|
|
|
|
:idempotent => idempotent,
|
|
|
|
:method => 'POST',
|
|
|
|
:parser => parser
|
2015-06-24 11:54:30 -04:00
|
|
|
})
|
|
|
|
rescue Excon::Errors::HTTPStatusError => error
|
|
|
|
match = Fog::AWS::Errors.match_error(error)
|
|
|
|
raise if match.empty?
|
|
|
|
raise case match[:code]
|
|
|
|
when 'LimitExceededException'
|
|
|
|
Fog::AWS::Kinesis::LimitExceeded.slurp(error, match[:message])
|
|
|
|
when 'ResourceNotFoundException'
|
|
|
|
Fog::AWS::Kinesis::ResourceNotFound.slurp(error, match[:message])
|
|
|
|
when 'ExpiredIteratorException'
|
|
|
|
Fog::AWS::Kinesis::ExpiredIterator.slurp(error, match[:message])
|
|
|
|
when 'InvalidArgumentException'
|
|
|
|
Fog::AWS::Kinesis::InvalidArgument.slurp(error, match[:message])
|
|
|
|
when 'ProvisionedThroughputExceededException'
|
|
|
|
Fog::AWS::Kinesis::ProvisionedThroughputExceeded.slurp(error, match[:message])
|
|
|
|
else
|
|
|
|
Fog::AWS::Kinesis::Error.slurp(error, "#{match[:code]} => #{match[:message]}")
|
|
|
|
end
|
2015-06-19 12:51:54 -04:00
|
|
|
end
|
|
|
|
|
|
|
|
end
|
|
|
|
|
|
|
|
class Mock
|
2015-07-02 15:23:25 -04:00
|
|
|
def self.mutex
|
|
|
|
@mutex ||= Mutex.new
|
|
|
|
end
|
|
|
|
def mutex; self.class.mutex; end
|
|
|
|
|
2015-06-30 18:15:14 -04:00
|
|
|
def self.data
|
|
|
|
@data ||= Hash.new do |hash, region|
|
|
|
|
hash[region] = Hash.new do |region_hash, key|
|
|
|
|
region_hash[key] = {
|
|
|
|
:kinesis_streams => {}
|
|
|
|
}
|
|
|
|
end
|
|
|
|
end
|
|
|
|
end
|
|
|
|
|
|
|
|
def self.reset
|
|
|
|
@data = nil
|
|
|
|
end
|
|
|
|
|
2015-06-19 12:51:54 -04:00
|
|
|
def initialize(options={})
|
2015-06-30 18:15:14 -04:00
|
|
|
@account_id = Fog::AWS::Mock.owner_id
|
|
|
|
@aws_access_key_id = options[:aws_access_key_id]
|
|
|
|
@region = options[:region] || 'us-east-1'
|
|
|
|
|
|
|
|
unless ['ap-northeast-1', 'ap-southeast-1', 'ap-southeast-2', 'eu-central-1', 'eu-west-1', 'sa-east-1', 'us-east-1', 'us-west-1', 'us-west-2'].include?(@region)
|
|
|
|
raise ArgumentError, "Unknown region: #{@region.inspect}"
|
|
|
|
end
|
|
|
|
end
|
|
|
|
|
|
|
|
def data
|
|
|
|
self.class.data[@region][@aws_access_key_id]
|
|
|
|
end
|
|
|
|
|
|
|
|
def reset_data
|
|
|
|
self.class.data[@region].delete(@aws_access_key_id)
|
|
|
|
end
|
|
|
|
|
2015-07-02 15:23:25 -04:00
|
|
|
def self.next_sequence_number
|
|
|
|
mutex.synchronize do
|
|
|
|
@sequence_number ||= -1
|
|
|
|
@sequence_number += 1
|
|
|
|
@sequence_number.to_s
|
|
|
|
end
|
2015-06-19 12:51:54 -04:00
|
|
|
end
|
2015-07-02 15:23:25 -04:00
|
|
|
def next_sequence_number; self.class.next_sequence_number; end
|
|
|
|
|
|
|
|
def self.next_shard_id
|
|
|
|
mutex.synchronize do
|
|
|
|
@shard_id ||= -1
|
|
|
|
@shard_id += 1
|
|
|
|
"shardId-#{@shard_id.to_s.rjust(12, "0")}"
|
|
|
|
end
|
|
|
|
end
|
|
|
|
def next_shard_id; self.class.next_shard_id; end
|
|
|
|
|
2015-06-19 12:51:54 -04:00
|
|
|
end
|
|
|
|
|
|
|
|
end
|
|
|
|
end
|
|
|
|
end
|