From b947f350efab9bbf7a9b196f5dbbcf4fd69251ec Mon Sep 17 00:00:00 2001 From: colmsnowplow Date: Wed, 12 Feb 2025 12:44:42 +0000 Subject: [PATCH] Allow configuration of client ID --- kinsumer.go | 22 ++++++++++++++++------ kinsumer_test.go | 24 ++++++++++++------------ shard_consumer_test.go | 26 +++++++++++++------------- 3 files changed, 41 insertions(+), 31 deletions(-) diff --git a/kinsumer.go b/kinsumer.go index c8d2655..0a44ad2 100644 --- a/kinsumer.go +++ b/kinsumer.go @@ -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 } @@ -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 } @@ -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, @@ -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 @@ -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 @@ -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() diff --git a/kinsumer_test.go b/kinsumer_test.go index df7a727..4f926bf 100644 --- a/kinsumer_test.go +++ b/kinsumer_test.go @@ -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) } @@ -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) } @@ -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) } @@ -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() @@ -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) @@ -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" @@ -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) diff --git a/shard_consumer_test.go b/shard_consumer_test.go index 13a7aa0..0dce81a 100644 --- a/shard_consumer_test.go +++ b/shard_consumer_test.go @@ -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, @@ -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) @@ -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, @@ -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, @@ -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, @@ -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, @@ -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,