From 6b42707eb0853bd6987f8dba6fc3ed111dc5fc56 Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Wed, 12 Aug 2026 18:22:23 +0800 Subject: [PATCH 1/3] first commit --- .../topicmanager/kafka_topic_manager_test.go | 429 +++++++++--------- pkg/errors/error.go | 4 + pkg/errors/error_test.go | 5 + pkg/sink/kafka/admin.go | 26 +- pkg/sink/kafka/admin_test.go | 250 +++++++++- 5 files changed, 470 insertions(+), 244 deletions(-) diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go index 9e4b7e32b0..4f18d4e3ef 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,146 @@ 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 adminAuthorizationError(operation, resource string) error { + authorizationErr := errors.ErrKafkaAdminAuthorizationFailed.FastGenByArgs(operation, resource) + return errors.WrapError(errors.ErrKafkaAdminAPI, authorizationErr, operation, resource) } 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") - 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.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") + + require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) + require.Equal(t, "new-topic", createdTopic.Name) + }) } func TestCreateTopicValidatesReplicationFactor(t *testing.T) { @@ -191,18 +173,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 +188,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 +198,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 +249,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 +274,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, adminAuthorizationError("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 +302,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(adminAuthorizationError("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..e6445df9c4 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"), ) + ErrKafkaAdminAuthorizationFailed = errors.Normalize( + "kafka admin API %s authorization failed: %s", + errors.RFCCodeText("CDC:ErrKafkaAdminAuthorizationFailed"), + ) 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..17c5791a01 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: "ErrKafkaAdminAuthorizationFailed should return false", + err: ErrKafkaAdminAuthorizationFailed.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..f6e3fc866d 100644 --- a/pkg/sink/kafka/admin.go +++ b/pkg/sink/kafka/admin.go @@ -61,7 +61,7 @@ func (a *saramaAdminClient) GetAllBrokers() []Broker { func (a *saramaAdminClient) GetBrokerConfig(configName string) (string, bool, error) { _, controller, err := a.admin.DescribeCluster() if err != nil { - return "", false, errors.WrapError(errors.ErrKafkaAdminAPI, err, "describe-cluster", "cluster") + return "", false, wrapSaramaAdminError(err, "describe-cluster", "cluster") } configEntries, err := a.admin.DescribeConfig(sarama.ConfigResource{ @@ -70,7 +70,7 @@ func (a *saramaAdminClient) GetBrokerConfig(configName string) (string, bool, er ConfigNames: []string{configName}, }) if err != nil { - return "", false, errors.WrapError(errors.ErrKafkaAdminAPI, err, "describe-config", configName) + return "", false, wrapSaramaAdminError(err, "describe-config", configName) } // For compatibility with KOP, we checked all return values. @@ -91,7 +91,7 @@ func (a *saramaAdminClient) GetTopicConfig(topicName string, configName string) ConfigNames: []string{configName}, }) if err != nil { - return "", false, errors.WrapError(errors.ErrKafkaAdminAPI, err, "describe-config", topicName) + return "", false, wrapSaramaAdminError(err, "describe-config", topicName) } // For compatibility with KOP, we checked all return values. @@ -110,7 +110,7 @@ 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, ",")) + return nil, wrapSaramaAdminError(err, "describe-topics", strings.Join(topics, ",")) } for _, meta := range metaList { @@ -119,7 +119,7 @@ func (a *saramaAdminClient) GetTopicsMeta(topics []string, ignoreTopicError bool continue } if !ignoreTopicError { - return nil, errors.WrapError(errors.ErrKafkaAdminAPI, meta.Err, "describe-topic", meta.Name) + return nil, wrapSaramaAdminError(meta.Err, "describe-topic", meta.Name) } log.Warn("kafka topic metadata refresh failed", zap.String("keyspace", a.changefeed.Keyspace()), @@ -138,6 +138,18 @@ func (a *saramaAdminClient) GetTopicsMeta(topics []string, ignoreTopicError bool // IsAdminAuthorizationFailed checks whether err is an authorization failure from Kafka admin APIs. func IsAdminAuthorizationFailed(err error) bool { + return errors.Is(err, errors.ErrKafkaAdminAuthorizationFailed) +} + +func wrapSaramaAdminError(err error, operation, resource string) error { + if isSaramaAdminAuthorizationFailed(err) { + // Preserve the ErrKafkaAdminAPI RFC code and avoid adding a second stack. + err = errors.ErrKafkaAdminAuthorizationFailed.Wrap(err).FastGenByArgs(operation, resource) + } + return errors.WrapError(errors.ErrKafkaAdminAPI, err, operation, resource) +} + +func isSaramaAdminAuthorizationFailed(err error) bool { return errors.Is(err, sarama.ErrTopicAuthorizationFailed) || errors.Is(err, sarama.ErrClusterAuthorizationFailed) } @@ -147,7 +159,7 @@ func (a *saramaAdminClient) GetTopicsPartitionsNum(topics []string) (map[string] for _, topic := range topics { partition, err := a.client.Partitions(topic) if err != nil { - return nil, errors.WrapError(errors.ErrKafkaAdminAPI, err, "list-partitions", topic) + return nil, wrapSaramaAdminError(err, "list-partitions", topic) } result[topic] = int32(len(partition)) } @@ -164,7 +176,7 @@ 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()) { - return errors.WrapError(errors.ErrKafkaAdminAPI, err, "create-topic", detail.Name) + return wrapSaramaAdminError(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..ad6a5a10c6 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,221 @@ func TestGetBrokerConfig(t *testing.T) { }) } -func TestCreateTopic(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, IsAdminAuthorizationFailed(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, IsAdminAuthorizationFailed(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 TestSaramaAdminAuthorizationErrorMapping(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) + for _, cause := range []error{ + sarama.ErrTopicAuthorizationFailed, + sarama.ErrClusterAuthorizationFailed, + } { + cause := cause + t.Run(cause.Error(), func(t *testing.T) { + err := wrapSaramaAdminError(cause, "describe-topic", "test-topic") - client := &saramaAdminClient{ - changefeed: common.NewChangeFeedIDWithName("test", "default"), - admin: admin, + require.ErrorIs(t, err, errors.ErrKafkaAdminAuthorizationFailed) + require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) + require.ErrorIs(t, err, cause) + require.True(t, IsAdminAuthorizationFailed(err)) + code, ok := errors.RFCCode(err) + require.True(t, ok) + require.Equal(t, errors.ErrKafkaAdminAPI.RFCCode(), code) + }) } - err := client.CreateTopic(&TopicDetail{ - Name: "test-topic", - NumPartitions: 3, - ReplicationFactor: 2, + + t.Run("general admin error", func(t *testing.T) { + err := wrapSaramaAdminError(io.ErrUnexpectedEOF, "describe-topic", "test-topic") + + require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) + require.ErrorIs(t, err, io.ErrUnexpectedEOF) + require.NotErrorIs(t, err, errors.ErrKafkaAdminAuthorizationFailed) + require.False(t, IsAdminAuthorizationFailed(err)) }) - require.NoError(t, err) +} + +func TestCreateTopic(t *testing.T) { + t.Parallel() + + tests := []struct { + name string + adminErr error + expectedErr error + }{ + {name: "success"}, + {name: "topic already exists", adminErr: sarama.ErrTopicAlreadyExists}, + {name: "authorization error", adminErr: sarama.ErrClusterAuthorizationFailed, expectedErr: errors.ErrKafkaAdminAuthorizationFailed}, + {name: "general error", adminErr: sarama.ErrInvalidReplicationFactor, expectedErr: errors.ErrKafkaAdminAPI}, + } + + for _, test := range tests { + test := test + 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, + }) + + if test.expectedErr == nil { + require.NoError(t, err) + return + } + require.ErrorIs(t, err, test.expectedErr) + require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) + require.ErrorIs(t, err, test.adminErr) + }) + } } func TestAdminClientClose(t *testing.T) { From 0b7df2862db4f5c3297044a059ab4c5d20b2881c Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Wed, 12 Aug 2026 19:04:25 +0800 Subject: [PATCH 2/3] adjust code --- .../topicmanager/kafka_topic_manager_test.go | 9 ++------- pkg/sink/kafka/admin.go | 11 +++-------- pkg/sink/kafka/admin_test.go | 19 ++++++++++--------- 3 files changed, 15 insertions(+), 24 deletions(-) diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go index 4f18d4e3ef..f4411ce82f 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go @@ -26,11 +26,6 @@ import ( const kafkaTopicManagerTestTopic = "mock_topic" -func adminAuthorizationError(operation, resource string) error { - authorizationErr := errors.ErrKafkaAdminAuthorizationFailed.FastGenByArgs(operation, resource) - return errors.WrapError(errors.ErrKafkaAdminAPI, authorizationErr, operation, resource) -} - func TestCreateTopic(t *testing.T) { t.Parallel() @@ -277,7 +272,7 @@ func TestCreateTopicWithTopicDescribeDenied(t *testing.T) { 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, adminAuthorizationError("describe-topic", "default-topic")) + nil, errors.ErrKafkaAdminAuthorizationFailed.GenWithStackByArgs("describe-topic", "default-topic")) manager := newKafkaTopicManager( "default-topic", common.NewChangefeedID4Test("test", "test"), @@ -309,7 +304,7 @@ func TestCreateTopicWithCreateDenied(t *testing.T) { Name: "default-topic", NumPartitions: 2, ReplicationFactor: 1, - }).Return(adminAuthorizationError("create-topic", "default-topic")) + }).Return(errors.ErrKafkaAdminAuthorizationFailed.GenWithStackByArgs("create-topic", "default-topic")) manager := newKafkaTopicManager( "default-topic", common.NewChangefeedID4Test("test", "test"), diff --git a/pkg/sink/kafka/admin.go b/pkg/sink/kafka/admin.go index f6e3fc866d..6f56e58cdd 100644 --- a/pkg/sink/kafka/admin.go +++ b/pkg/sink/kafka/admin.go @@ -142,18 +142,13 @@ func IsAdminAuthorizationFailed(err error) bool { } func wrapSaramaAdminError(err error, operation, resource string) error { - if isSaramaAdminAuthorizationFailed(err) { - // Preserve the ErrKafkaAdminAPI RFC code and avoid adding a second stack. - err = errors.ErrKafkaAdminAuthorizationFailed.Wrap(err).FastGenByArgs(operation, resource) + if errors.Is(err, sarama.ErrTopicAuthorizationFailed) || + errors.Is(err, sarama.ErrClusterAuthorizationFailed) { + return errors.WrapError(errors.ErrKafkaAdminAuthorizationFailed, err, operation, resource) } return errors.WrapError(errors.ErrKafkaAdminAPI, err, operation, resource) } -func isSaramaAdminAuthorizationFailed(err error) bool { - return errors.Is(err, sarama.ErrTopicAuthorizationFailed) || - errors.Is(err, sarama.ErrClusterAuthorizationFailed) -} - func (a *saramaAdminClient) GetTopicsPartitionsNum(topics []string) (map[string]int32, error) { result := make(map[string]int32, len(topics)) for _, topic := range topics { diff --git a/pkg/sink/kafka/admin_test.go b/pkg/sink/kafka/admin_test.go index ad6a5a10c6..097f1fe4ca 100644 --- a/pkg/sink/kafka/admin_test.go +++ b/pkg/sink/kafka/admin_test.go @@ -234,17 +234,16 @@ func TestSaramaAdminAuthorizationErrorMapping(t *testing.T) { sarama.ErrTopicAuthorizationFailed, sarama.ErrClusterAuthorizationFailed, } { - cause := cause t.Run(cause.Error(), func(t *testing.T) { err := wrapSaramaAdminError(cause, "describe-topic", "test-topic") require.ErrorIs(t, err, errors.ErrKafkaAdminAuthorizationFailed) - require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) + require.NotErrorIs(t, err, errors.ErrKafkaAdminAPI) require.ErrorIs(t, err, cause) require.True(t, IsAdminAuthorizationFailed(err)) code, ok := errors.RFCCode(err) require.True(t, ok) - require.Equal(t, errors.ErrKafkaAdminAPI.RFCCode(), code) + require.Equal(t, errors.ErrKafkaAdminAuthorizationFailed.RFCCode(), code) }) } @@ -262,18 +261,18 @@ func TestCreateTopic(t *testing.T) { t.Parallel() tests := []struct { - name string - adminErr error - expectedErr error + 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.ErrKafkaAdminAuthorizationFailed}, + {name: "authorization error", adminErr: sarama.ErrClusterAuthorizationFailed, expectedErr: errors.ErrKafkaAdminAuthorizationFailed, authorization: true}, {name: "general error", adminErr: sarama.ErrInvalidReplicationFactor, expectedErr: errors.ErrKafkaAdminAPI}, } for _, test := range tests { - test := test t.Run(test.name, func(t *testing.T) { ctrl := gomock.NewController(t) admin := NewMocksaramaClusterAdmin(ctrl) @@ -297,8 +296,10 @@ func TestCreateTopic(t *testing.T) { return } require.ErrorIs(t, err, test.expectedErr) - require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) require.ErrorIs(t, err, test.adminErr) + if test.authorization { + require.NotErrorIs(t, err, errors.ErrKafkaAdminAPI) + } }) } } From 2b3c99ee64f25cecefdbfbad2d56320248558709 Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Thu, 13 Aug 2026 15:32:29 +0800 Subject: [PATCH 3/3] adjust code --- .../sink/topicmanager/kafka_topic_manager.go | 6 +- .../topicmanager/kafka_topic_manager_test.go | 4 +- pkg/errors/error.go | 6 +- pkg/errors/error_test.go | 4 +- pkg/sink/kafka/admin.go | 52 ++++++++---- pkg/sink/kafka/admin_test.go | 84 +++++++++++++------ 6 files changed, 101 insertions(+), 55 deletions(-) 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 f4411ce82f..e4dc2612eb 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go @@ -272,7 +272,7 @@ func TestCreateTopicWithTopicDescribeDenied(t *testing.T) { 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.ErrKafkaAdminAuthorizationFailed.GenWithStackByArgs("describe-topic", "default-topic")) + nil, errors.ErrKafkaAuthorizationFailed.GenWithStackByArgs("describe-topic", "default-topic")) manager := newKafkaTopicManager( "default-topic", common.NewChangefeedID4Test("test", "test"), @@ -304,7 +304,7 @@ func TestCreateTopicWithCreateDenied(t *testing.T) { Name: "default-topic", NumPartitions: 2, ReplicationFactor: 1, - }).Return(errors.ErrKafkaAdminAuthorizationFailed.GenWithStackByArgs("create-topic", "default-topic")) + }).Return(errors.ErrKafkaAuthorizationFailed.GenWithStackByArgs("create-topic", "default-topic")) manager := newKafkaTopicManager( "default-topic", common.NewChangefeedID4Test("test", "test"), diff --git a/pkg/errors/error.go b/pkg/errors/error.go index e6445df9c4..9956d39544 100644 --- a/pkg/errors/error.go +++ b/pkg/errors/error.go @@ -143,9 +143,9 @@ var ( "kafka admin API %s failed: %s", errors.RFCCodeText("CDC:ErrKafkaAdminAPI"), ) - ErrKafkaAdminAuthorizationFailed = errors.Normalize( - "kafka admin API %s authorization failed: %s", - errors.RFCCodeText("CDC:ErrKafkaAdminAuthorizationFailed"), + ErrKafkaAuthorizationFailed = errors.Normalize( + "kafka %s authorization failed: %s", + errors.RFCCodeText("CDC:ErrKafkaAuthorizationFailed"), ) ErrPulsarInvalidTopicExpression = errors.Normalize( "invalid topic expression", diff --git a/pkg/errors/error_test.go b/pkg/errors/error_test.go index 17c5791a01..31879157cb 100644 --- a/pkg/errors/error_test.go +++ b/pkg/errors/error_test.go @@ -109,8 +109,8 @@ func TestShouldFailChangefeed(t *testing.T) { expected: false, }, { - name: "ErrKafkaAdminAuthorizationFailed should return false", - err: ErrKafkaAdminAuthorizationFailed.GenWithStackByArgs("describe-topic", "test-topic"), + name: "ErrKafkaAuthorizationFailed should return false", + err: ErrKafkaAuthorizationFailed.GenWithStackByArgs("describe-topic", "test-topic"), expected: false, }, { diff --git a/pkg/sink/kafka/admin.go b/pkg/sink/kafka/admin.go index 6f56e58cdd..d45ec702d1 100644 --- a/pkg/sink/kafka/admin.go +++ b/pkg/sink/kafka/admin.go @@ -61,7 +61,10 @@ func (a *saramaAdminClient) GetAllBrokers() []Broker { func (a *saramaAdminClient) GetBrokerConfig(configName string) (string, bool, error) { _, controller, err := a.admin.DescribeCluster() if err != nil { - return "", false, wrapSaramaAdminError(err, "describe-cluster", "cluster") + if IsAuthorizationFailed(err) { + return "", false, errors.WrapError(errors.ErrKafkaAuthorizationFailed, err, "describe-cluster", "cluster") + } + return "", false, errors.WrapError(errors.ErrKafkaAdminAPI, err, "describe-cluster", "cluster") } configEntries, err := a.admin.DescribeConfig(sarama.ConfigResource{ @@ -70,7 +73,10 @@ func (a *saramaAdminClient) GetBrokerConfig(configName string) (string, bool, er ConfigNames: []string{configName}, }) if err != nil { - return "", false, wrapSaramaAdminError(err, "describe-config", configName) + if IsAuthorizationFailed(err) { + return "", false, errors.WrapError(errors.ErrKafkaAuthorizationFailed, err, "describe-config", configName) + } + return "", false, errors.WrapError(errors.ErrKafkaAdminAPI, err, "describe-config", configName) } // For compatibility with KOP, we checked all return values. @@ -91,7 +97,10 @@ func (a *saramaAdminClient) GetTopicConfig(topicName string, configName string) ConfigNames: []string{configName}, }) if err != nil { - return "", false, wrapSaramaAdminError(err, "describe-config", topicName) + if IsAuthorizationFailed(err) { + return "", false, errors.WrapError(errors.ErrKafkaAuthorizationFailed, err, "describe-config", topicName) + } + return "", false, errors.WrapError(errors.ErrKafkaAdminAPI, err, "describe-config", topicName) } // For compatibility with KOP, we checked all return values. @@ -110,7 +119,11 @@ func (a *saramaAdminClient) GetTopicsMeta(topics []string, ignoreTopicError bool metaList, err := a.admin.DescribeTopics(topics) if err != nil { - return nil, wrapSaramaAdminError(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,7 +132,10 @@ func (a *saramaAdminClient) GetTopicsMeta(topics []string, ignoreTopicError bool continue } if !ignoreTopicError { - return nil, wrapSaramaAdminError(meta.Err, "describe-topic", meta.Name) + 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", zap.String("keyspace", a.changefeed.Keyspace()), @@ -136,17 +152,11 @@ 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, errors.ErrKafkaAdminAuthorizationFailed) -} - -func wrapSaramaAdminError(err error, operation, resource string) error { - if errors.Is(err, sarama.ErrTopicAuthorizationFailed) || - errors.Is(err, sarama.ErrClusterAuthorizationFailed) { - return errors.WrapError(errors.ErrKafkaAdminAuthorizationFailed, err, operation, resource) - } - return errors.WrapError(errors.ErrKafkaAdminAPI, err, operation, resource) +// 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) } func (a *saramaAdminClient) GetTopicsPartitionsNum(topics []string) (map[string]int32, error) { @@ -154,7 +164,10 @@ func (a *saramaAdminClient) GetTopicsPartitionsNum(topics []string) (map[string] for _, topic := range topics { partition, err := a.client.Partitions(topic) if err != nil { - return nil, wrapSaramaAdminError(err, "list-partitions", topic) + 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)) } @@ -171,7 +184,10 @@ 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()) { - return wrapSaramaAdminError(err, "create-topic", detail.Name) + 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 097f1fe4ca..90fc3530dd 100644 --- a/pkg/sink/kafka/admin_test.go +++ b/pkg/sink/kafka/admin_test.go @@ -140,7 +140,7 @@ func TestGetTopicConfig(t *testing.T) { require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) require.ErrorIs(t, err, context.DeadlineExceeded) - require.False(t, IsAdminAuthorizationFailed(err)) + require.False(t, IsAuthorizationFailed(err)) }) } @@ -206,7 +206,46 @@ func TestGetTopicsMeta(t *testing.T) { require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) require.ErrorIs(t, err, sarama.ErrInvalidTopic) - require.False(t, IsAdminAuthorizationFailed(err)) + 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) { @@ -227,34 +266,25 @@ func TestGetTopicsMeta(t *testing.T) { }) } -func TestSaramaAdminAuthorizationErrorMapping(t *testing.T) { +func TestIsAuthorizationFailed(t *testing.T) { t.Parallel() - for _, cause := range []error{ - sarama.ErrTopicAuthorizationFailed, - sarama.ErrClusterAuthorizationFailed, - } { - t.Run(cause.Error(), func(t *testing.T) { - err := wrapSaramaAdminError(cause, "describe-topic", "test-topic") - - require.ErrorIs(t, err, errors.ErrKafkaAdminAuthorizationFailed) - require.NotErrorIs(t, err, errors.ErrKafkaAdminAPI) - require.ErrorIs(t, err, cause) - require.True(t, IsAdminAuthorizationFailed(err)) - code, ok := errors.RFCCode(err) - require.True(t, ok) - require.Equal(t, errors.ErrKafkaAdminAuthorizationFailed.RFCCode(), code) - }) + 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}, } - t.Run("general admin error", func(t *testing.T) { - err := wrapSaramaAdminError(io.ErrUnexpectedEOF, "describe-topic", "test-topic") - - require.ErrorIs(t, err, errors.ErrKafkaAdminAPI) - require.ErrorIs(t, err, io.ErrUnexpectedEOF) - require.NotErrorIs(t, err, errors.ErrKafkaAdminAuthorizationFailed) - require.False(t, IsAdminAuthorizationFailed(err)) - }) + 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) { @@ -268,7 +298,7 @@ func TestCreateTopic(t *testing.T) { }{ {name: "success"}, {name: "topic already exists", adminErr: sarama.ErrTopicAlreadyExists}, - {name: "authorization error", adminErr: sarama.ErrClusterAuthorizationFailed, expectedErr: errors.ErrKafkaAdminAuthorizationFailed, authorization: true}, + {name: "authorization error", adminErr: sarama.ErrClusterAuthorizationFailed, expectedErr: errors.ErrKafkaAuthorizationFailed, authorization: true}, {name: "general error", adminErr: sarama.ErrInvalidReplicationFactor, expectedErr: errors.ErrKafkaAdminAPI}, }