diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go index 8e43a52df6..a0110217f3 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go @@ -277,7 +277,7 @@ func (m *kafkaTopicManager) CreateTopicAndWaitUntilVisible( // which means we should create the topic later. topicDetails, err := m.admin.GetTopicsMeta([]string{topicName}, true) if err != nil { - if kafka.IsAdminAuthorizationFailed(err) { + if kafka.IsAuthorizationFailed(err) { return m.useConfiguredPartitionNum(topicName, err), nil } return 0, err @@ -288,7 +288,7 @@ func (m *kafkaTopicManager) CreateTopicAndWaitUntilVisible( topicDetails, err = m.admin.GetTopicsMeta([]string{topicName}, false) if err != nil { - if kafka.IsAdminAuthorizationFailed(err) { + if kafka.IsAuthorizationFailed(err) { return m.useConfiguredPartitionNum(topicName, err), nil } } else if numPartition, ok := m.tryStoreTopicMeta(topicName, topicDetails); ok { @@ -298,7 +298,7 @@ func (m *kafkaTopicManager) CreateTopicAndWaitUntilVisible( start := time.Now() partitionNum, err := m.createTopic(ctx, topicName) if err != nil { - if kafka.IsAdminAuthorizationFailed(err) { + if kafka.IsAuthorizationFailed(err) { return m.useConfiguredPartitionNum(topicName, err), nil } return 0, err diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go index 9e4b7e32b0..e4dc2612eb 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go @@ -17,7 +17,6 @@ import ( "context" "testing" - "github.com/IBM/sarama" "github.com/golang/mock/gomock" "github.com/pingcap/ticdc/pkg/common" "github.com/pingcap/ticdc/pkg/errors" @@ -27,163 +26,141 @@ import ( const kafkaTopicManagerTestTopic = "mock_topic" -type mockAdminClientWithDeniedDescribe struct { - *kafka.MockAdminClient - createTopicCalled bool - describeCount int -} - -func (m *mockAdminClientWithDeniedDescribe) GetTopicsMeta( - topics []string, - ignoreTopicError bool, -) (map[string]kafka.TopicDetail, error) { - m.describeCount++ - if ignoreTopicError { - return map[string]kafka.TopicDetail{}, nil - } - return nil, errors.WrapError( - errors.ErrKafkaAdminAPI, - sarama.ErrTopicAuthorizationFailed, - "describe-topic", - topics[0], - ) -} - -func (m *mockAdminClientWithDeniedDescribe) CreateTopic( - detail *kafka.TopicDetail, -) error { - m.createTopicCalled = true - return nil -} - -type mockAdminClientWithDeniedCreate struct { - *kafka.MockAdminClient - createTopicCalled bool - describeCount int -} - -func (m *mockAdminClientWithDeniedCreate) GetTopicsMeta( - topics []string, - ignoreTopicError bool, -) (map[string]kafka.TopicDetail, error) { - m.describeCount++ - return map[string]kafka.TopicDetail{}, nil -} - -func (m *mockAdminClientWithDeniedCreate) CreateTopic( - detail *kafka.TopicDetail, -) error { - m.createTopicCalled = true - return errors.WrapError( - errors.ErrKafkaAdminAPI, - sarama.ErrClusterAuthorizationFailed, - "create-topic", - detail.Name, - ) -} - func TestCreateTopic(t *testing.T) { t.Parallel() - ctrl := gomock.NewController(t) - adminClient := kafka.NewMockAdminClient(ctrl) - cfg := &kafka.AutoCreateTopicConfig{ - AutoCreate: true, - PartitionNum: 2, - ReplicationFactor: 1, - RequiredAcks: kafka.WaitForAll, - } - changefeedID := common.NewChangefeedID4Test("test", "test") - ctx := context.Background() - var gotNewTopicDetail *kafka.TopicDetail - var gotFailedTopicDetail *kafka.TopicDetail - gomock.InOrder( + + t.Run("existing topic", func(t *testing.T) { + t.Parallel() + + ctrl := gomock.NewController(t) + adminClient := kafka.NewMockAdminClient(ctrl) adminClient.EXPECT().GetTopicsMeta([]string{kafkaTopicManagerTestTopic}, true).Return( map[string]kafka.TopicDetail{ kafkaTopicManagerTestTopic: { Name: kafkaTopicManagerTestTopic, NumPartitions: 2, }, - }, nil), - adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, true).Return( - map[string]kafka.TopicDetail{}, nil), - adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, false).Return( - map[string]kafka.TopicDetail{}, nil), + }, nil) + manager := newKafkaTopicManager( + kafkaTopicManagerTestTopic, + changefeedID, + adminClient, + &kafka.AutoCreateTopicConfig{PartitionNum: 2}, + ) + + partitionNum, err := manager.CreateTopicAndWaitUntilVisible(context.Background(), kafkaTopicManagerTestTopic) + + require.NoError(t, err) + require.Equal(t, int32(2), partitionNum) + }) + + t.Run("create missing topic", func(t *testing.T) { + t.Parallel() + + ctrl := gomock.NewController(t) + adminClient := kafka.NewMockAdminClient(ctrl) + var createdTopic *kafka.TopicDetail + adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, true).Return(map[string]kafka.TopicDetail{}, nil) + adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, false).DoAndReturn( + func([]string, bool) (map[string]kafka.TopicDetail, error) { + if createdTopic == nil { + return map[string]kafka.TopicDetail{}, nil + } + return map[string]kafka.TopicDetail{ + createdTopic.Name: { + Name: createdTopic.Name, + NumPartitions: createdTopic.NumPartitions, + }, + }, nil + }).Times(2) adminClient.EXPECT().CreateTopic(gomock.Any()).DoAndReturn( func(detail *kafka.TopicDetail) error { - gotNewTopicDetail = detail + copy := *detail + createdTopic = © return nil - }), - adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, false).Return( - map[string]kafka.TopicDetail{ - "new-topic": { - Name: "new-topic", - NumPartitions: 2, - }, - }, nil), - adminClient.EXPECT().GetTopicsMeta([]string{"new-topic2"}, true).Return( - map[string]kafka.TopicDetail{}, nil), - adminClient.EXPECT().GetTopicsMeta([]string{"new-topic2"}, false).Return( - map[string]kafka.TopicDetail{}, nil), - adminClient.EXPECT().GetTopicsMeta([]string{"new-topic-failed"}, true).Return( - map[string]kafka.TopicDetail{}, nil), - adminClient.EXPECT().GetTopicsMeta([]string{"new-topic-failed"}, false).Return( - map[string]kafka.TopicDetail{}, nil), - adminClient.EXPECT().CreateTopic(gomock.Any()).DoAndReturn( - func(detail *kafka.TopicDetail) error { - gotFailedTopicDetail = detail - return errors.WrapError(errors.ErrKafkaAdminAPI, sarama.ErrInvalidReplicationFactor, "create-topic", detail.Name) - }), - ) + }) + manager := newKafkaTopicManager( + kafkaTopicManagerTestTopic, + changefeedID, + adminClient, + &kafka.AutoCreateTopicConfig{ + AutoCreate: true, + PartitionNum: 2, + ReplicationFactor: 1, + RequiredAcks: kafka.WaitForLocal, + }, + ) - manager := newKafkaTopicManager(kafkaTopicManagerTestTopic, changefeedID, adminClient, cfg) - partitionNum, err := manager.CreateTopicAndWaitUntilVisible(ctx, kafkaTopicManagerTestTopic) - require.NoError(t, err) - require.Equal(t, int32(2), partitionNum) + partitionNum, err := manager.CreateTopicAndWaitUntilVisible(context.Background(), "new-topic") - cfg.RequiredAcks = kafka.WaitForLocal - partitionNum, err = manager.CreateTopicAndWaitUntilVisible(ctx, "new-topic") - require.NoError(t, err) - require.Equal(t, int32(2), partitionNum) - require.Equal(t, &kafka.TopicDetail{ - Name: "new-topic", - NumPartitions: 2, - ReplicationFactor: 1, - }, gotNewTopicDetail) - partitionsNum, err := manager.GetPartitionNum(ctx, "new-topic") - require.NoError(t, err) - require.Equal(t, int32(2), partitionsNum) + require.NoError(t, err) + require.Equal(t, int32(2), partitionNum) + require.Equal(t, &kafka.TopicDetail{ + Name: "new-topic", + NumPartitions: 2, + ReplicationFactor: 1, + }, createdTopic) + partitionsNum, err := manager.GetPartitionNum(context.Background(), "new-topic") + require.NoError(t, err) + require.Equal(t, int32(2), partitionsNum) + }) + + t.Run("auto create disabled", func(t *testing.T) { + t.Parallel() + + ctrl := gomock.NewController(t) + adminClient := kafka.NewMockAdminClient(ctrl) + adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, true).Return(map[string]kafka.TopicDetail{}, nil) + adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, false).Return(map[string]kafka.TopicDetail{}, nil) + manager := newKafkaTopicManager( + "new-topic", + changefeedID, + adminClient, + &kafka.AutoCreateTopicConfig{ + AutoCreate: false, + PartitionNum: 2, + ReplicationFactor: 1, + RequiredAcks: kafka.WaitForAll, + }, + ) - // Try to create a topic without auto create. - cfg = &kafka.AutoCreateTopicConfig{ - AutoCreate: false, - PartitionNum: 2, - ReplicationFactor: 1, - RequiredAcks: kafka.WaitForAll, - } - manager = newKafkaTopicManager("new-topic2", changefeedID, adminClient, cfg) - _, err = manager.CreateTopicAndWaitUntilVisible(ctx, "new-topic2") - require.Regexp( - t, - "`auto-create-topic` is false, and new-topic2 not found", - err, - ) + _, err := manager.CreateTopicAndWaitUntilVisible(context.Background(), "new-topic") + + require.ErrorContains(t, err, "`auto-create-topic` is false, and new-topic not found") + }) + + t.Run("create error", func(t *testing.T) { + t.Parallel() + + ctrl := gomock.NewController(t) + adminClient := kafka.NewMockAdminClient(ctrl) + adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, true).Return(map[string]kafka.TopicDetail{}, nil) + adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, false).Return(map[string]kafka.TopicDetail{}, nil) + var createdTopic *kafka.TopicDetail + adminClient.EXPECT().CreateTopic(gomock.Any()).DoAndReturn( + func(detail *kafka.TopicDetail) error { + copy := *detail + createdTopic = © + return errors.ErrKafkaAdminAPI.GenWithStackByArgs("create-topic", detail.Name) + }) + manager := newKafkaTopicManager( + "new-topic", + changefeedID, + adminClient, + &kafka.AutoCreateTopicConfig{ + AutoCreate: true, + PartitionNum: 2, + ReplicationFactor: 4, + }, + ) + + _, err := manager.CreateTopicAndWaitUntilVisible(context.Background(), "new-topic") - topic := "new-topic-failed" - // Invalid replication factor. - // It happens when replication-factor is greater than the number of brokers. - cfg = &kafka.AutoCreateTopicConfig{ - AutoCreate: true, - PartitionNum: 2, - ReplicationFactor: 4, - } - manager = newKafkaTopicManager(topic, changefeedID, adminClient, cfg) - _, err = manager.CreateTopicAndWaitUntilVisible(ctx, topic) - require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) - require.ErrorIs(t, err, sarama.ErrInvalidReplicationFactor) - require.NotNil(t, gotFailedTopicDetail) - require.Equal(t, "new-topic-failed", gotFailedTopicDetail.Name) + require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) + require.Equal(t, "new-topic", createdTopic.Name) + }) } func TestCreateTopicValidatesReplicationFactor(t *testing.T) { @@ -191,18 +168,11 @@ func TestCreateTopicValidatesReplicationFactor(t *testing.T) { ctrl := gomock.NewController(t) adminClient := kafka.NewMockAdminClient(ctrl) - topic := "new-topic" - gomock.InOrder( - adminClient.EXPECT().GetTopicsMeta([]string{topic}, true). - Return(map[string]kafka.TopicDetail{}, nil), - adminClient.EXPECT().GetTopicsMeta([]string{topic}, false). - Return(map[string]kafka.TopicDetail{}, nil), - adminClient.EXPECT().GetBrokerConfig(kafka.MinInsyncReplicasConfigName). - Return("2", true, nil), - ) - + adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, true).Return(map[string]kafka.TopicDetail{}, nil) + adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, false).Return(map[string]kafka.TopicDetail{}, nil) + adminClient.EXPECT().GetBrokerConfig(kafka.MinInsyncReplicasConfigName).Return("2", true, nil) manager := newKafkaTopicManager( - topic, + "new-topic", common.NewChangefeedID4Test("test", "test"), adminClient, &kafka.AutoCreateTopicConfig{ @@ -213,7 +183,8 @@ func TestCreateTopicValidatesReplicationFactor(t *testing.T) { }, ) - _, err := manager.CreateTopicAndWaitUntilVisible(context.Background(), topic) + _, err := manager.CreateTopicAndWaitUntilVisible(context.Background(), "new-topic") + require.ErrorContains(t, err, "`replication-factor` 1 is smaller than the `min.insync.replicas` 2 of broker") } @@ -222,42 +193,50 @@ func TestEnsureTopicExistsWaitsUntilVisible(t *testing.T) { ctrl := gomock.NewController(t) adminClient := kafka.NewMockAdminClient(ctrl) - cfg := &kafka.AutoCreateTopicConfig{ - AutoCreate: true, - PartitionNum: 2, - ReplicationFactor: 1, - } - - topic := "delayed-topic" - gomock.InOrder( - adminClient.EXPECT().GetTopicsMeta([]string{topic}, true).Return( - map[string]kafka.TopicDetail{}, nil), - adminClient.EXPECT().GetTopicsMeta([]string{topic}, false).Return( - map[string]kafka.TopicDetail{}, nil), - adminClient.EXPECT().CreateTopic(gomock.Any()).DoAndReturn( - func(detail *kafka.TopicDetail) error { - require.Equal(t, &kafka.TopicDetail{ - Name: topic, - NumPartitions: 2, - ReplicationFactor: 1, - }, detail) - return nil - }), - adminClient.EXPECT().GetTopicsMeta([]string{topic}, false).Return( - map[string]kafka.TopicDetail{}, nil), - adminClient.EXPECT().GetTopicsMeta([]string{topic}, false).Return( - map[string]kafka.TopicDetail{ - topic: { - Name: topic, + created := false + postCreateDescribeCount := 0 + adminClient.EXPECT().GetTopicsMeta([]string{"delayed-topic"}, true).Return(map[string]kafka.TopicDetail{}, nil) + adminClient.EXPECT().GetTopicsMeta([]string{"delayed-topic"}, false).DoAndReturn( + func([]string, bool) (map[string]kafka.TopicDetail, error) { + if !created { + return map[string]kafka.TopicDetail{}, nil + } + postCreateDescribeCount++ + if postCreateDescribeCount == 1 { + return map[string]kafka.TopicDetail{}, nil + } + return map[string]kafka.TopicDetail{ + "delayed-topic": { + Name: "delayed-topic", NumPartitions: 2, }, - }, nil), + }, nil + }).Times(3) + adminClient.EXPECT().CreateTopic(gomock.Any()).DoAndReturn( + func(detail *kafka.TopicDetail) error { + require.Equal(t, &kafka.TopicDetail{ + Name: "delayed-topic", + NumPartitions: 2, + ReplicationFactor: 1, + }, detail) + created = true + return nil + }) + + err := EnsureTopic( + context.Background(), + common.NewChangefeedID4Test("test", "test"), + "delayed-topic", + &kafka.AutoCreateTopicConfig{ + AutoCreate: true, + PartitionNum: 2, + ReplicationFactor: 1, + }, + adminClient, ) - ctx := context.Background() - changefeedID := common.NewChangefeedID4Test("test", "test") - err := EnsureTopic(ctx, changefeedID, topic, cfg, adminClient) require.NoError(t, err) + require.Equal(t, 2, postCreateDescribeCount) } func TestGetTopicManagerStartsBackgroundRefreshAfterTopicReady(t *testing.T) { @@ -265,23 +244,22 @@ func TestGetTopicManagerStartsBackgroundRefreshAfterTopicReady(t *testing.T) { ctrl := gomock.NewController(t) adminClient := kafka.NewMockAdminClient(ctrl) - topic := "existing-topic" - adminClient.EXPECT().GetTopicsMeta([]string{topic}, true).Return( + adminClient.EXPECT().GetTopicsMeta([]string{"existing-topic"}, true).Return( map[string]kafka.TopicDetail{ - topic: { - Name: topic, + "existing-topic": { + Name: "existing-topic", NumPartitions: 2, }, - }, nil, - ) + }, nil) manager, err := GetTopicManagerAndTryCreateTopic( t.Context(), common.NewChangefeedID4Test("test", "test"), - topic, + "existing-topic", &kafka.AutoCreateTopicConfig{PartitionNum: 2}, adminClient, ) + require.NoError(t, err) defer manager.Close() require.NotNil(t, manager.(*kafkaTopicManager).cancel) @@ -291,27 +269,26 @@ func TestCreateTopicWithTopicDescribeDenied(t *testing.T) { t.Parallel() ctrl := gomock.NewController(t) - adminClient := &mockAdminClientWithDeniedDescribe{ - MockAdminClient: kafka.NewMockAdminClient(ctrl), - } - cfg := &kafka.AutoCreateTopicConfig{ - AutoCreate: true, - PartitionNum: 2, - ReplicationFactor: 1, - } + adminClient := kafka.NewMockAdminClient(ctrl) + adminClient.EXPECT().GetTopicsMeta([]string{"default-topic"}, true).Return(map[string]kafka.TopicDetail{}, nil) + adminClient.EXPECT().GetTopicsMeta([]string{"default-topic"}, false).Return( + nil, errors.ErrKafkaAuthorizationFailed.GenWithStackByArgs("describe-topic", "default-topic")) + manager := newKafkaTopicManager( + "default-topic", + common.NewChangefeedID4Test("test", "test"), + adminClient, + &kafka.AutoCreateTopicConfig{ + AutoCreate: true, + PartitionNum: 2, + ReplicationFactor: 1, + }, + ) - changefeedID := common.NewChangefeedID4Test("test", "test") - ctx := context.Background() - defaultTopic := "default-topic" - manager := newKafkaTopicManager(defaultTopic, changefeedID, adminClient, cfg) + partitionNum, err := manager.CreateTopicAndWaitUntilVisible(context.Background(), "default-topic") - partitionNum, err := manager.CreateTopicAndWaitUntilVisible(ctx, defaultTopic) require.NoError(t, err) require.Equal(t, int32(2), partitionNum) - require.False(t, adminClient.createTopicCalled) - require.Equal(t, 2, adminClient.describeCount) - - partitions, ok := manager.topics.Load(defaultTopic) + partitions, ok := manager.topics.Load("default-topic") require.True(t, ok) require.Equal(t, int32(2), partitions) } @@ -320,27 +297,30 @@ func TestCreateTopicWithCreateDenied(t *testing.T) { t.Parallel() ctrl := gomock.NewController(t) - adminClient := &mockAdminClientWithDeniedCreate{ - MockAdminClient: kafka.NewMockAdminClient(ctrl), - } - cfg := &kafka.AutoCreateTopicConfig{ - AutoCreate: true, - PartitionNum: 2, + adminClient := kafka.NewMockAdminClient(ctrl) + adminClient.EXPECT().GetTopicsMeta([]string{"default-topic"}, true).Return(map[string]kafka.TopicDetail{}, nil) + adminClient.EXPECT().GetTopicsMeta([]string{"default-topic"}, false).Return(map[string]kafka.TopicDetail{}, nil) + adminClient.EXPECT().CreateTopic(&kafka.TopicDetail{ + Name: "default-topic", + NumPartitions: 2, ReplicationFactor: 1, - } + }).Return(errors.ErrKafkaAuthorizationFailed.GenWithStackByArgs("create-topic", "default-topic")) + manager := newKafkaTopicManager( + "default-topic", + common.NewChangefeedID4Test("test", "test"), + adminClient, + &kafka.AutoCreateTopicConfig{ + AutoCreate: true, + PartitionNum: 2, + ReplicationFactor: 1, + }, + ) - changefeedID := common.NewChangefeedID4Test("test", "test") - ctx := context.Background() - defaultTopic := "default-topic" - manager := newKafkaTopicManager(defaultTopic, changefeedID, adminClient, cfg) + partitionNum, err := manager.CreateTopicAndWaitUntilVisible(context.Background(), "default-topic") - partitionNum, err := manager.CreateTopicAndWaitUntilVisible(ctx, defaultTopic) require.NoError(t, err) require.Equal(t, int32(2), partitionNum) - require.True(t, adminClient.createTopicCalled) - require.Equal(t, 2, adminClient.describeCount) - - partitions, ok := manager.topics.Load(defaultTopic) + partitions, ok := manager.topics.Load("default-topic") require.True(t, ok) require.Equal(t, int32(2), partitions) } diff --git a/pkg/errors/error.go b/pkg/errors/error.go index 5dcb4160f3..9956d39544 100644 --- a/pkg/errors/error.go +++ b/pkg/errors/error.go @@ -143,6 +143,10 @@ var ( "kafka admin API %s failed: %s", errors.RFCCodeText("CDC:ErrKafkaAdminAPI"), ) + ErrKafkaAuthorizationFailed = errors.Normalize( + "kafka %s authorization failed: %s", + errors.RFCCodeText("CDC:ErrKafkaAuthorizationFailed"), + ) ErrPulsarInvalidTopicExpression = errors.Normalize( "invalid topic expression", errors.RFCCodeText("CDC:ErrPulsarTopicExprInvalid"), diff --git a/pkg/errors/error_test.go b/pkg/errors/error_test.go index 2c76ca8978..31879157cb 100644 --- a/pkg/errors/error_test.go +++ b/pkg/errors/error_test.go @@ -108,6 +108,11 @@ func TestShouldFailChangefeed(t *testing.T) { err: ErrKafkaAdminAPI.GenWithStackByArgs("describe-topic", "test-topic"), expected: false, }, + { + name: "ErrKafkaAuthorizationFailed should return false", + err: ErrKafkaAuthorizationFailed.GenWithStackByArgs("describe-topic", "test-topic"), + expected: false, + }, { name: "ErrKafkaSendMessage should return false", err: ErrKafkaSendMessage.GenWithStackByArgs(), diff --git a/pkg/sink/kafka/admin.go b/pkg/sink/kafka/admin.go index 8ad44cf5ea..d45ec702d1 100644 --- a/pkg/sink/kafka/admin.go +++ b/pkg/sink/kafka/admin.go @@ -61,6 +61,9 @@ func (a *saramaAdminClient) GetAllBrokers() []Broker { func (a *saramaAdminClient) GetBrokerConfig(configName string) (string, bool, error) { _, controller, err := a.admin.DescribeCluster() if err != nil { + if IsAuthorizationFailed(err) { + return "", false, errors.WrapError(errors.ErrKafkaAuthorizationFailed, err, "describe-cluster", "cluster") + } return "", false, errors.WrapError(errors.ErrKafkaAdminAPI, err, "describe-cluster", "cluster") } @@ -70,6 +73,9 @@ func (a *saramaAdminClient) GetBrokerConfig(configName string) (string, bool, er ConfigNames: []string{configName}, }) if err != nil { + if IsAuthorizationFailed(err) { + return "", false, errors.WrapError(errors.ErrKafkaAuthorizationFailed, err, "describe-config", configName) + } return "", false, errors.WrapError(errors.ErrKafkaAdminAPI, err, "describe-config", configName) } @@ -91,6 +97,9 @@ func (a *saramaAdminClient) GetTopicConfig(topicName string, configName string) ConfigNames: []string{configName}, }) if err != nil { + if IsAuthorizationFailed(err) { + return "", false, errors.WrapError(errors.ErrKafkaAuthorizationFailed, err, "describe-config", topicName) + } return "", false, errors.WrapError(errors.ErrKafkaAdminAPI, err, "describe-config", topicName) } @@ -110,7 +119,11 @@ func (a *saramaAdminClient) GetTopicsMeta(topics []string, ignoreTopicError bool metaList, err := a.admin.DescribeTopics(topics) if err != nil { - return nil, errors.WrapError(errors.ErrKafkaAdminAPI, err, "describe-topics", strings.Join(topics, ",")) + resource := strings.Join(topics, ",") + if IsAuthorizationFailed(err) { + return nil, errors.WrapError(errors.ErrKafkaAuthorizationFailed, err, "describe-topics", resource) + } + return nil, errors.WrapError(errors.ErrKafkaAdminAPI, err, "describe-topics", resource) } for _, meta := range metaList { @@ -119,6 +132,9 @@ func (a *saramaAdminClient) GetTopicsMeta(topics []string, ignoreTopicError bool continue } if !ignoreTopicError { + if IsAuthorizationFailed(meta.Err) { + return nil, errors.WrapError(errors.ErrKafkaAuthorizationFailed, meta.Err, "describe-topic", meta.Name) + } return nil, errors.WrapError(errors.ErrKafkaAdminAPI, meta.Err, "describe-topic", meta.Name) } log.Warn("kafka topic metadata refresh failed", @@ -136,9 +152,10 @@ func (a *saramaAdminClient) GetTopicsMeta(topics []string, ignoreTopicError bool return result, nil } -// IsAdminAuthorizationFailed checks whether err is an authorization failure from Kafka admin APIs. -func IsAdminAuthorizationFailed(err error) bool { - return errors.Is(err, sarama.ErrTopicAuthorizationFailed) || +// IsAuthorizationFailed checks whether err is a Kafka authorization failure. +func IsAuthorizationFailed(err error) bool { + return errors.Is(err, errors.ErrKafkaAuthorizationFailed) || + errors.Is(err, sarama.ErrTopicAuthorizationFailed) || errors.Is(err, sarama.ErrClusterAuthorizationFailed) } @@ -147,6 +164,9 @@ func (a *saramaAdminClient) GetTopicsPartitionsNum(topics []string) (map[string] for _, topic := range topics { partition, err := a.client.Partitions(topic) if err != nil { + if IsAuthorizationFailed(err) { + return nil, errors.WrapError(errors.ErrKafkaAuthorizationFailed, err, "list-partitions", topic) + } return nil, errors.WrapError(errors.ErrKafkaAdminAPI, err, "list-partitions", topic) } result[topic] = int32(len(partition)) @@ -164,6 +184,9 @@ func (a *saramaAdminClient) CreateTopic(detail *TopicDetail) error { err := a.admin.CreateTopic(detail.Name, request, false) // Ignore the already exists error because it's not harmful. if err != nil && !strings.Contains(err.Error(), sarama.ErrTopicAlreadyExists.Error()) { + if IsAuthorizationFailed(err) { + return errors.WrapError(errors.ErrKafkaAuthorizationFailed, err, "create-topic", detail.Name) + } return errors.WrapError(errors.ErrKafkaAdminAPI, err, "create-topic", detail.Name) } return nil diff --git a/pkg/sink/kafka/admin_test.go b/pkg/sink/kafka/admin_test.go index a79ad5b06a..90fc3530dd 100644 --- a/pkg/sink/kafka/admin_test.go +++ b/pkg/sink/kafka/admin_test.go @@ -14,6 +14,7 @@ package kafka import ( + "context" "io" "testing" @@ -27,6 +28,30 @@ import ( func TestGetBrokerConfig(t *testing.T) { t.Parallel() + t.Run("found", func(t *testing.T) { + ctrl := gomock.NewController(t) + admin := NewMocksaramaClusterAdmin(ctrl) + admin.EXPECT().DescribeCluster().Return(nil, int32(1), nil) + admin.EXPECT().DescribeConfig(sarama.ConfigResource{ + Type: sarama.BrokerResource, + Name: "1", + ConfigNames: []string{"message.max.bytes"}, + }).Return([]sarama.ConfigEntry{ + {Name: "unrelated", Value: "value"}, + {Name: "message.max.bytes", Value: "1048576"}, + }, nil) + + client := &saramaAdminClient{ + changefeed: common.NewChangeFeedIDWithName("test", "default"), + admin: admin, + } + value, found, err := client.GetBrokerConfig("message.max.bytes") + + require.NoError(t, err) + require.True(t, found) + require.Equal(t, "1048576", value) + }) + t.Run("not found", func(t *testing.T) { ctrl := gomock.NewController(t) admin := NewMocksaramaClusterAdmin(ctrl) @@ -61,26 +86,252 @@ func TestGetBrokerConfig(t *testing.T) { }) } +func TestGetTopicConfig(t *testing.T) { + t.Parallel() + + t.Run("found", func(t *testing.T) { + ctrl := gomock.NewController(t) + admin := NewMocksaramaClusterAdmin(ctrl) + admin.EXPECT().DescribeConfig(sarama.ConfigResource{ + Type: sarama.TopicResource, + Name: "test-topic", + ConfigNames: []string{"max.message.bytes"}, + }).Return([]sarama.ConfigEntry{ + {Name: "max.message.bytes", Value: "1048576"}, + }, nil) + client := &saramaAdminClient{ + changefeed: common.NewChangeFeedIDWithName("test", "default"), + admin: admin, + } + + value, found, err := client.GetTopicConfig("test-topic", "max.message.bytes") + + require.NoError(t, err) + require.True(t, found) + require.Equal(t, "1048576", value) + }) + + t.Run("not found", func(t *testing.T) { + ctrl := gomock.NewController(t) + admin := NewMocksaramaClusterAdmin(ctrl) + admin.EXPECT().DescribeConfig(gomock.Any()).Return([]sarama.ConfigEntry{}, nil) + client := &saramaAdminClient{ + changefeed: common.NewChangeFeedIDWithName("test", "default"), + admin: admin, + } + + value, found, err := client.GetTopicConfig("test-topic", "missing") + + require.NoError(t, err) + require.False(t, found) + require.Empty(t, value) + }) + + t.Run("admin error", func(t *testing.T) { + ctrl := gomock.NewController(t) + admin := NewMocksaramaClusterAdmin(ctrl) + admin.EXPECT().DescribeConfig(gomock.Any()).Return(nil, context.DeadlineExceeded) + client := &saramaAdminClient{ + changefeed: common.NewChangeFeedIDWithName("test", "default"), + admin: admin, + } + + _, _, err := client.GetTopicConfig("test-topic", "missing") + + require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) + require.ErrorIs(t, err, context.DeadlineExceeded) + require.False(t, IsAuthorizationFailed(err)) + }) +} + +func TestGetTopicsMeta(t *testing.T) { + t.Parallel() + + t.Run("returns valid topics and ignores unknown topics", func(t *testing.T) { + ctrl := gomock.NewController(t) + admin := NewMocksaramaClusterAdmin(ctrl) + admin.EXPECT().DescribeTopics([]string{"valid-topic", "missing-topic"}).Return([]*sarama.TopicMetadata{ + { + Name: "valid-topic", + Partitions: []*sarama.PartitionMetadata{{}, {}}, + }, + { + Name: "missing-topic", + Err: sarama.ErrUnknownTopicOrPartition, + }, + }, nil) + client := &saramaAdminClient{ + changefeed: common.NewChangeFeedIDWithName("test", "default"), + admin: admin, + } + + topics, err := client.GetTopicsMeta([]string{"valid-topic", "missing-topic"}, false) + + require.NoError(t, err) + require.Equal(t, map[string]TopicDetail{ + "valid-topic": { + Name: "valid-topic", + NumPartitions: 2, + }, + }, topics) + }) + + t.Run("missing response", func(t *testing.T) { + ctrl := gomock.NewController(t) + admin := NewMocksaramaClusterAdmin(ctrl) + admin.EXPECT().DescribeTopics([]string{"missing-topic"}).Return(nil, nil) + client := &saramaAdminClient{ + changefeed: common.NewChangeFeedIDWithName("test", "default"), + admin: admin, + } + + topics, err := client.GetTopicsMeta([]string{"missing-topic"}, false) + + require.NoError(t, err) + require.Empty(t, topics) + }) + + t.Run("topic error", func(t *testing.T) { + ctrl := gomock.NewController(t) + admin := NewMocksaramaClusterAdmin(ctrl) + admin.EXPECT().DescribeTopics([]string{"test-topic"}).Return([]*sarama.TopicMetadata{ + {Name: "test-topic", Err: sarama.ErrInvalidTopic}, + }, nil) + client := &saramaAdminClient{ + changefeed: common.NewChangeFeedIDWithName("test", "default"), + admin: admin, + } + + _, err := client.GetTopicsMeta([]string{"test-topic"}, false) + + require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) + require.ErrorIs(t, err, sarama.ErrInvalidTopic) + require.False(t, IsAuthorizationFailed(err)) + }) + + t.Run("topic authorization error", func(t *testing.T) { + ctrl := gomock.NewController(t) + admin := NewMocksaramaClusterAdmin(ctrl) + admin.EXPECT().DescribeTopics([]string{"test-topic"}).Return([]*sarama.TopicMetadata{ + {Name: "test-topic", Err: sarama.ErrTopicAuthorizationFailed}, + }, nil) + client := &saramaAdminClient{ + changefeed: common.NewChangeFeedIDWithName("test", "default"), + admin: admin, + } + + _, err := client.GetTopicsMeta([]string{"test-topic"}, false) + + require.ErrorIs(t, err, errors.ErrKafkaAuthorizationFailed) + require.NotErrorIs(t, err, errors.ErrKafkaAdminAPI) + require.ErrorIs(t, err, sarama.ErrTopicAuthorizationFailed) + require.True(t, IsAuthorizationFailed(err)) + code, ok := errors.RFCCode(err) + require.True(t, ok) + require.Equal(t, errors.ErrKafkaAuthorizationFailed.RFCCode(), code) + }) + + t.Run("cluster authorization error", func(t *testing.T) { + ctrl := gomock.NewController(t) + admin := NewMocksaramaClusterAdmin(ctrl) + admin.EXPECT().DescribeTopics([]string{"test-topic"}).Return(nil, sarama.ErrClusterAuthorizationFailed) + client := &saramaAdminClient{ + changefeed: common.NewChangeFeedIDWithName("test", "default"), + admin: admin, + } + + _, err := client.GetTopicsMeta([]string{"test-topic"}, false) + + require.ErrorIs(t, err, errors.ErrKafkaAuthorizationFailed) + require.NotErrorIs(t, err, errors.ErrKafkaAdminAPI) + require.ErrorIs(t, err, sarama.ErrClusterAuthorizationFailed) + require.True(t, IsAuthorizationFailed(err)) + }) + + t.Run("ignored topic error", func(t *testing.T) { + ctrl := gomock.NewController(t) + admin := NewMocksaramaClusterAdmin(ctrl) + admin.EXPECT().DescribeTopics([]string{"test-topic"}).Return([]*sarama.TopicMetadata{ + {Name: "test-topic", Err: sarama.ErrInvalidTopic}, + }, nil) + client := &saramaAdminClient{ + changefeed: common.NewChangeFeedIDWithName("test", "default"), + admin: admin, + } + + topics, err := client.GetTopicsMeta([]string{"test-topic"}, true) + + require.NoError(t, err) + require.Empty(t, topics) + }) +} + +func TestIsAuthorizationFailed(t *testing.T) { + t.Parallel() + + tests := []struct { + name string + err error + expected bool + }{ + {name: "TiCDC authorization error", err: errors.ErrKafkaAuthorizationFailed.GenWithStackByArgs("describe-topic", "test-topic"), expected: true}, + {name: "topic authorization error", err: sarama.ErrTopicAuthorizationFailed, expected: true}, + {name: "cluster authorization error", err: sarama.ErrClusterAuthorizationFailed, expected: true}, + {name: "general error", err: sarama.ErrInvalidTopic}, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + require.Equal(t, test.expected, IsAuthorizationFailed(test.err)) + }) + } +} + func TestCreateTopic(t *testing.T) { t.Parallel() - ctrl := gomock.NewController(t) - admin := NewMocksaramaClusterAdmin(ctrl) - admin.EXPECT().CreateTopic("test-topic", &sarama.TopicDetail{ - NumPartitions: 3, - ReplicationFactor: 2, - }, false).Return(nil) + tests := []struct { + name string + adminErr error + expectedErr error + authorization bool + }{ + {name: "success"}, + {name: "topic already exists", adminErr: sarama.ErrTopicAlreadyExists}, + {name: "authorization error", adminErr: sarama.ErrClusterAuthorizationFailed, expectedErr: errors.ErrKafkaAuthorizationFailed, authorization: true}, + {name: "general error", adminErr: sarama.ErrInvalidReplicationFactor, expectedErr: errors.ErrKafkaAdminAPI}, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + ctrl := gomock.NewController(t) + admin := NewMocksaramaClusterAdmin(ctrl) + admin.EXPECT().CreateTopic("test-topic", &sarama.TopicDetail{ + NumPartitions: 3, + ReplicationFactor: 2, + }, false).Return(test.adminErr) + client := &saramaAdminClient{ + changefeed: common.NewChangeFeedIDWithName("test", "default"), + admin: admin, + } + + err := client.CreateTopic(&TopicDetail{ + Name: "test-topic", + NumPartitions: 3, + ReplicationFactor: 2, + }) - client := &saramaAdminClient{ - changefeed: common.NewChangeFeedIDWithName("test", "default"), - admin: admin, + if test.expectedErr == nil { + require.NoError(t, err) + return + } + require.ErrorIs(t, err, test.expectedErr) + require.ErrorIs(t, err, test.adminErr) + if test.authorization { + require.NotErrorIs(t, err, errors.ErrKafkaAdminAPI) + } + }) } - err := client.CreateTopic(&TopicDetail{ - Name: "test-topic", - NumPartitions: 3, - ReplicationFactor: 2, - }) - require.NoError(t, err) } func TestAdminClientClose(t *testing.T) {