Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
The table of contents is too big for display.
Diff view
Diff view
  •  
  •  
  •  
23 changes: 0 additions & 23 deletions actions/actionsfakes/fake_iactions.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1,570 changes: 1,161 additions & 409 deletions backends/awskinesis/kinesisfakes/fake_kinesis.go

Large diffs are not rendered by default.

45 changes: 32 additions & 13 deletions backends/awskinesis/read.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import (
"time"

"github.com/aws/aws-sdk-go/aws"
"github.com/aws/aws-sdk-go/aws/arn"
"github.com/aws/aws-sdk-go/service/kinesis"
"github.com/pkg/errors"
uuid "github.com/satori/go.uuid"
Expand All @@ -22,10 +23,15 @@ import (

// getShardIterator creates a shard interator based on read arguments
// lastSequenceNumber argument is used when creating anther iterator after a recoverable read failure
func (k *Kinesis) getShardIterator(args *args.AWSKinesisReadArgs, shardID string, lastSequenceNumber *string) (*string, error) {
func (k *Kinesis) getShardIterator(ctx context.Context, args *args.AWSKinesisReadArgs, shardID string, lastSequenceNumber *string) (*string, error) {
getOpts := &kinesis.GetShardIteratorInput{
ShardId: aws.String(shardID),
StreamName: aws.String(args.Stream),
ShardId: aws.String(shardID),
}

if arn.IsARN(args.Stream) {
getOpts.StreamARN = aws.String(args.Stream)
} else {
getOpts.StreamName = aws.String(args.Stream)
}

if args.ReadTrimHorizon {
Expand All @@ -48,7 +54,7 @@ func (k *Kinesis) getShardIterator(args *args.AWSKinesisReadArgs, shardID string
}

//Get shard iterator
shardResp, err := k.client.GetShardIterator(getOpts)
shardResp, err := k.client.GetShardIteratorWithContext(ctx, getOpts)
if err != nil {
return nil, err
}
Expand All @@ -59,10 +65,15 @@ func (k *Kinesis) getShardIterator(args *args.AWSKinesisReadArgs, shardID string
// getShards returns a slice of all shards that exist for a stream
// This method is used when no --stream argument is provided. In this case, we will
// start a read for each available shard
func (k *Kinesis) getShards(args *args.AWSKinesisReadArgs) ([]string, error) {
shardResp, err := k.client.DescribeStream(&kinesis.DescribeStreamInput{
StreamName: aws.String(args.Stream),
})
func (k *Kinesis) getShards(ctx context.Context, args *args.AWSKinesisReadArgs) ([]string, error) {
describeOpts := &kinesis.DescribeStreamInput{}
if arn.IsARN(args.Stream) {
describeOpts.StreamARN = aws.String(args.Stream)
} else {
describeOpts.StreamName = aws.String(args.Stream)
}

shardResp, err := k.client.DescribeStreamWithContext(ctx, describeOpts)
if err != nil {
return nil, errors.Wrap(err, "unable to get shards")
}
Expand All @@ -88,7 +99,7 @@ func (k *Kinesis) Read(ctx context.Context, readOpts *opts.ReadOptions, resultsC
}

// No shard specified. Get all shards and launch a goroutine to read each shard
shards, err := k.getShards(readOpts.AwsKinesis.Args)
shards, err := k.getShards(ctx, readOpts.AwsKinesis.Args)
if err != nil {
return errors.Wrap(err, "unable to get shards")
}
Expand Down Expand Up @@ -120,18 +131,26 @@ func (k *Kinesis) readShard(ctx context.Context, shardID string, readOpts *opts.
// is reading by specific sequence number. Otherwise this value will be ignored
var lastSequenceNumber *string

shardIterator, err := k.getShardIterator(readOpts.AwsKinesis.Args, shardID, lastSequenceNumber)
shardIterator, err := k.getShardIterator(ctx, readOpts.AwsKinesis.Args, shardID, lastSequenceNumber)
if err != nil {
return errors.Wrap(err, "unable to create shard iterator")
}

k.log.Infof("Waiting for messages for shard '%s'...", shardID)

for {
resp, err := k.client.GetRecordsWithContext(ctx, &kinesis.GetRecordsInput{
getOpts := &kinesis.GetRecordsInput{
ShardIterator: shardIterator,
Limit: aws.Int64(readOpts.AwsKinesis.Args.MaxRecords),
})
}

if arn.IsARN(readOpts.AwsKinesis.Args.Stream) {
getOpts.StreamARN = aws.String(readOpts.AwsKinesis.Args.Stream)
} else {
// stream name isn't required in this case
}

resp, err := k.client.GetRecordsWithContext(ctx, getOpts)
if err != nil {
// Some errors are recoverable
if !canRetry(err) {
Expand All @@ -145,7 +164,7 @@ func (k *Kinesis) readShard(ctx context.Context, shardID string, readOpts *opts.

// Get new iterator. lastSequenceNumber is specified in the event that the user
// started reading from a specific sequence number. We want to start where we left off
shardIterator, err = k.getShardIterator(readOpts.AwsKinesis.Args, shardID, lastSequenceNumber)
shardIterator, err = k.getShardIterator(ctx, readOpts.AwsKinesis.Args, shardID, lastSequenceNumber)
if err != nil {
return errors.Wrap(err, "unable to create shard iterator")
}
Expand Down
50 changes: 47 additions & 3 deletions backends/awskinesis/read_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -80,7 +80,9 @@ var _ = Describe("AWS Kinesis Backend", func() {
resultsCh := make(chan *records.ReadRecord, 1)

fakeKinesis := &kinesisfakes.FakeKinesisAPI{}
fakeKinesis.GetRecordsWithContextStub = func(context.Context, *kinesis.GetRecordsInput, ...request.Option) (*kinesis.GetRecordsOutput, error) {
fakeKinesis.GetRecordsWithContextStub = func(_ context.Context, input *kinesis.GetRecordsInput, _ ...request.Option) (*kinesis.GetRecordsOutput, error) {
Expect(input.StreamARN).To(BeNil())

return &kinesis.GetRecordsOutput{
Records: []*kinesis.Record{
{
Expand All @@ -89,7 +91,10 @@ var _ = Describe("AWS Kinesis Backend", func() {
},
}, nil
}
fakeKinesis.GetShardIteratorStub = func(*kinesis.GetShardIteratorInput) (*kinesis.GetShardIteratorOutput, error) {
fakeKinesis.GetShardIteratorWithContextStub = func(_ context.Context, input *kinesis.GetShardIteratorInput, _ ...request.Option) (*kinesis.GetShardIteratorOutput, error) {
Expect(input.StreamName).ToNot(BeNil())
Expect(input.StreamARN).To(BeNil())

return &kinesis.GetShardIteratorOutput{
ShardIterator: aws.String("test"),
}, nil
Expand All @@ -104,7 +109,46 @@ var _ = Describe("AWS Kinesis Backend", func() {

Expect(err).ToNot(HaveOccurred())
Expect(fakeKinesis.GetRecordsWithContextCallCount()).To(Equal(1))
Expect(fakeKinesis.GetShardIteratorCallCount()).To(Equal(1))
Expect(fakeKinesis.GetShardIteratorWithContextCallCount()).To(Equal(1))
Expect(resultsCh).Should(Receive())
})
It("reads from Kinesis using an ARN", func() {
errorsCh := make(chan *records.ErrorRecord, 1)
resultsCh := make(chan *records.ReadRecord, 1)

fakeKinesis := &kinesisfakes.FakeKinesisAPI{}
fakeKinesis.GetRecordsWithContextStub = func(_ context.Context, input *kinesis.GetRecordsInput, _ ...request.Option) (*kinesis.GetRecordsOutput, error) {
Expect(input.StreamARN).ToNot(BeNil())

return &kinesis.GetRecordsOutput{
Records: []*kinesis.Record{
{
Data: []byte("test"),
},
},
}, nil
}
fakeKinesis.GetShardIteratorWithContextStub = func(_ context.Context, input *kinesis.GetShardIteratorInput, _ ...request.Option) (*kinesis.GetShardIteratorOutput, error) {
Expect(input.StreamARN).ToNot(BeNil())
Expect(input.StreamName).To(BeNil())

return &kinesis.GetShardIteratorOutput{
ShardIterator: aws.String("test"),
}, nil
}

a := &Kinesis{
client: fakeKinesis,
log: logrus.NewEntry(&logrus.Logger{Out: ioutil.Discard}),
}

readOpts.AwsKinesis.Args.Stream = "arn:aws:kinesis:us-east-1:123456789012:stream/test"

err := a.Read(context.Background(), readOpts, resultsCh, errorsCh)

Expect(err).ToNot(HaveOccurred())
Expect(fakeKinesis.GetRecordsWithContextCallCount()).To(Equal(1))
Expect(fakeKinesis.GetShardIteratorWithContextCallCount()).To(Equal(1))
Expect(resultsCh).Should(Receive())
})
})
Expand Down
Loading