Skip to content
Merged
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
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 16 additions & 6 deletions kinsumer.go
Original file line number Diff line number Diff line change
Expand Up @@ -78,11 +78,18 @@ func NewWithSession(session *session.Session, streamName, applicationName, clien
k := kinesis.New(session)
d := dynamodb.New(session)

return NewWithInterfaces(k, d, streamName, applicationName, clientName, config)
return NewWithInterfaces(k, d, streamName, applicationName, clientName, "", config)
}

// NewWithInterfaces allows you to override the Kinesis and Dynamo instances for mocking or using a local set of servers
func NewWithInterfaces(kinesis kinesisiface.KinesisAPI, dynamodb dynamodbiface.DynamoDBAPI, streamName, applicationName, clientName string, config Config) (*Kinsumer, error) {
func NewWithInterfaces(
kinesis kinesisiface.KinesisAPI,
dynamodb dynamodbiface.DynamoDBAPI,
streamName,
applicationName,
clientName string,
clientID string,
config Config) (*Kinsumer, error) {
if kinesis == nil {
return nil, ErrNoKinesisInterface
}
Expand All @@ -95,6 +102,9 @@ func NewWithInterfaces(kinesis kinesisiface.KinesisAPI, dynamodb dynamodbiface.D
if applicationName == "" {
return nil, ErrNoApplicationName
}
if clientID == "" {
clientID = uuid.New().String()
}
if err := validateConfig(&config); err != nil {
return nil, err
}
Expand All @@ -115,7 +125,7 @@ func NewWithInterfaces(kinesis kinesisiface.KinesisAPI, dynamodb dynamodbiface.D
checkpointTableName: applicationName + "_checkpoints",
clientsTableName: applicationName + "_clients",
metadataTableName: applicationName + "_metadata",
clientID: uuid.New().String(),
clientID: clientID,
clientName: clientName,
config: config,
maxAgeForClientRecord: *config.clientRecordMaxAge,
Expand All @@ -126,7 +136,7 @@ func NewWithInterfaces(kinesis kinesisiface.KinesisAPI, dynamodb dynamodbiface.D

// refreshShards registers our client, refreshes the lists of clients and shards, checks if we
// have become/unbecome the leader, and returns whether the shards/clients changed.
//TODO: Write unit test - needs dynamo _and_ kinesis mocking
// TODO: Write unit test - needs dynamo _and_ kinesis mocking
func (k *Kinsumer) refreshShards() (bool, error) {
var shardIDs []string

Expand Down Expand Up @@ -351,7 +361,7 @@ func (k *Kinsumer) kinesisStreamReady() error {
// Run runs the main kinesis consumer process. This is a non-blocking call, use Stop() to force it to return.
// This goroutine is responsible for starting/stopping consumers, aggregating all consumers' records,
// updating checkpointers as records are consumed, and refreshing our shard/client list and leadership
//TODO: Can we unit test this at all?
// TODO: Can we unit test this at all?
func (k *Kinsumer) Run() error {
if err := k.dynamoTableReady(k.checkpointTableName); err != nil {
return err
Expand Down Expand Up @@ -467,7 +477,7 @@ func (k *Kinsumer) Run() error {
}

// Stop stops the consumption of kinesis events
//TODO: Can we unit test this at all?
// TODO: Can we unit test this at all?
func (k *Kinsumer) Stop() {
k.stoprequest <- true
k.mainWG.Wait()
Expand Down
24 changes: 12 additions & 12 deletions kinsumer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,27 +44,27 @@ func TestNewWithInterfaces(t *testing.T) {
d := dynamodb.New(s)

// No kinesis
_, err := NewWithInterfaces(nil, d, "stream", "app", "client", NewConfig())
_, err := NewWithInterfaces(nil, d, "stream", "app", "client", "", NewConfig())
assert.NotEqual(t, err, nil)

// No dynamodb
_, err = NewWithInterfaces(k, nil, "stream", "app", "client", NewConfig())
_, err = NewWithInterfaces(k, nil, "stream", "app", "client", "", NewConfig())
assert.NotEqual(t, err, nil)

// No streamName
_, err = NewWithInterfaces(k, d, "", "app", "client", NewConfig())
_, err = NewWithInterfaces(k, d, "", "app", "client", "", NewConfig())
assert.NotEqual(t, err, nil)

// No applicationName
_, err = NewWithInterfaces(k, d, "stream", "", "client", NewConfig())
_, err = NewWithInterfaces(k, d, "stream", "", "client", "", NewConfig())
assert.NotEqual(t, err, nil)

// Invalid config
_, err = NewWithInterfaces(k, d, "stream", "app", "client", Config{})
_, err = NewWithInterfaces(k, d, "stream", "app", "client", "", Config{})
assert.NotEqual(t, err, nil)

// All ok
kinsumer, err := NewWithInterfaces(k, d, "stream", "app", "client", NewConfig())
kinsumer, err := NewWithInterfaces(k, d, "stream", "app", "client", "", NewConfig())
assert.Equal(t, err, nil)
assert.NotEqual(t, kinsumer, nil)
}
Expand Down Expand Up @@ -115,7 +115,7 @@ func setupTestEnvironment(t *testing.T, k kinesisiface.KinesisAPI, d dynamodbifa
}

testConf := NewConfig().WithDynamoWaiterDelay(*resourceChangeTimeout)
client, clientErr := NewWithInterfaces(k, d, streamName, *applicationName, "N/A", testConf)
client, clientErr := NewWithInterfaces(k, d, streamName, *applicationName, "N/A", "", testConf)
if clientErr != nil {
return fmt.Errorf("Error creating new Kinsumer Client: %s", clientErr)
}
Expand Down Expand Up @@ -190,7 +190,7 @@ func cleanupTestEnvironment(t *testing.T, k kinesisiface.KinesisAPI, d dynamodbi
}

testConf := NewConfig().WithDynamoWaiterDelay(*resourceChangeTimeout)
client, clientErr := NewWithInterfaces(k, d, "N/A", *applicationName, "N/A", testConf)
client, clientErr := NewWithInterfaces(k, d, "N/A", *applicationName, "N/A", "", testConf)
if clientErr != nil {
return fmt.Errorf("Error creating new Kinsumer Client: %s", clientErr)
}
Expand Down Expand Up @@ -332,7 +332,7 @@ func TestKinsumer(t *testing.T) {
time.Sleep(50 * time.Millisecond) // Add the clients slowly
}

clients[i], err = NewWithInterfaces(k, d, streamName, *applicationName, fmt.Sprintf("test_%d", i), config)
clients[i], err = NewWithInterfaces(k, d, streamName, *applicationName, fmt.Sprintf("test_%d", i), "", config)
require.NoError(t, err, "NewWithInterfaces() failed")

err = clients[i].Run()
Expand Down Expand Up @@ -431,7 +431,7 @@ func TestLeader(t *testing.T) {
time.Sleep(50 * time.Millisecond) // Add the clients slowly
}

clients[i], err = NewWithInterfaces(k, d, streamName, *applicationName, fmt.Sprintf("test_%d", i), config)
clients[i], err = NewWithInterfaces(k, d, streamName, *applicationName, fmt.Sprintf("test_%d", i), "", config)
require.NoError(t, err, "NewWithInterfaces() failed")
clients[i].clientID = strconv.Itoa(i + 1)

Expand Down Expand Up @@ -475,7 +475,7 @@ func TestLeader(t *testing.T) {
assert.Equal(t, true, clients[0].isLeader, "First client is not leader")
assert.Equal(t, false, clients[1].isLeader, "Second leader is also leader")

c, err := NewWithInterfaces(k, d, streamName, *applicationName, fmt.Sprintf("_test_%d", numberOfClients), config)
c, err := NewWithInterfaces(k, d, streamName, *applicationName, fmt.Sprintf("_test_%d", numberOfClients), "", config)
require.NoError(t, err, "NewWithInterfaces() failed")
c.clientID = "0"

Expand Down Expand Up @@ -537,7 +537,7 @@ func TestSplit(t *testing.T) {
time.Sleep(50 * time.Millisecond) // Add the clients slowly
}

clients[i], err = NewWithInterfaces(k, d, streamName, *applicationName, fmt.Sprintf("test_%d", i), config)
clients[i], err = NewWithInterfaces(k, d, streamName, *applicationName, fmt.Sprintf("test_%d", i), "", config)
require.NoError(t, err, "NewWithInterfaces() failed")
clients[i].clientID = strconv.Itoa(i + 1)

Expand Down
26 changes: 13 additions & 13 deletions shard_consumer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ func TestShardConsumer(t *testing.T) {
config = config.WithShardCheckFrequency(500 * time.Millisecond)
config = config.WithLeaderActionFrequency(500 * time.Millisecond)

kinsumer1, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_1", config)
kinsumer1, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_1", "", config)

desc, err := k.DescribeStream(&kinesis.DescribeStreamInput{
StreamName: &streamName,
Expand Down Expand Up @@ -85,8 +85,8 @@ func TestForcefulOwnershipChange(t *testing.T) {
maxAge2 := 500 * time.Millisecond
config2 := config.WithClientRecordMaxAge(&maxAge2)

kinsumer1, err1 := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_1", config1)
kinsumer2, err2 := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_2", config2)
kinsumer1, err1 := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_1", "", config1)
kinsumer2, err2 := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_2", "", config2)
require.NoError(t, err1)
require.NoError(t, err2)

Expand Down Expand Up @@ -227,8 +227,8 @@ func TestPotentialLegitimateDuplicates(t *testing.T) {
maxAge2 := 500 * time.Millisecond
config2 := config.WithClientRecordMaxAge(&maxAge2)

kinsumer1, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_1", config1)
kinsumer2, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_2", config2)
kinsumer1, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_1", "", config1)
kinsumer2, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_2", "", config2)

desc, err := k.DescribeStream(&kinesis.DescribeStreamInput{
StreamName: &streamName,
Expand Down Expand Up @@ -351,9 +351,9 @@ func TestShardsMerged(t *testing.T) {
config = config.WithLeaderActionFrequency(500 * time.Millisecond)
config = config.WithCommitFrequency(100 * time.Millisecond)

kinsumer1, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_1", config)
kinsumer2, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_2", config)
kinsumer3, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_3", config)
kinsumer1, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_1", "", config)
kinsumer2, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_2", "", config)
kinsumer3, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_3", "", config)

desc, err := k.DescribeStream(&kinesis.DescribeStreamInput{
StreamName: &streamName,
Expand Down Expand Up @@ -524,7 +524,7 @@ func TestConsumerStopStart(t *testing.T) {
config = config.WithLeaderActionFrequency(500 * time.Millisecond)
config = config.WithCommitFrequency(50 * time.Millisecond)

kinsumer, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_1", config)
kinsumer, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_1", "", config)

desc, err := k.DescribeStream(&kinesis.DescribeStreamInput{
StreamName: &streamName,
Expand Down Expand Up @@ -596,9 +596,9 @@ func TestMultipleConsumerStopStart(t *testing.T) {
config = config.WithLeaderActionFrequency(500 * time.Millisecond)
config = config.WithCommitFrequency(50 * time.Millisecond)

kinsumer1, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_1", config)
kinsumer2, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_2", config)
kinsumer3, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_3", config)
kinsumer1, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_1", "", config)
kinsumer2, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_2", "", config)
kinsumer3, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_3", "", config)

desc, err := k.DescribeStream(&kinesis.DescribeStreamInput{
StreamName: &streamName,
Expand Down Expand Up @@ -719,7 +719,7 @@ func TestDelayedUpdateDuplicates(t *testing.T) {
config = config.WithLeaderActionFrequency(500 * time.Millisecond)
config = config.WithCommitFrequency(50 * time.Millisecond)

kinsumer, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_1", config)
kinsumer, err := NewWithInterfaces(k, dynamo, streamName, *applicationName, "client_1", "", config)

desc, err := k.DescribeStream(&kinesis.DescribeStreamInput{
StreamName: &streamName,
Expand Down
Loading