From bb7c12e29ff9f8b1e43fdada589f14c4409b9b56 Mon Sep 17 00:00:00 2001 From: Ahmed ElNahas Date: Wed, 12 Aug 2026 16:46:25 -0400 Subject: [PATCH] add DeleteACLs to Manager DeleteACLs mirrors CreateACLs by delegating the request to the underlying kadm.Client.DeleteACLs. --- kafka/manager.go | 15 +++++++++ kafka/manager_test.go | 72 +++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 87 insertions(+) diff --git a/kafka/manager.go b/kafka/manager.go index 6bdce133..c6b6a1a5 100644 --- a/kafka/manager.go +++ b/kafka/manager.go @@ -368,6 +368,21 @@ func (m *Manager) CreateACLs(ctx context.Context, acls *kadm.ACLBuilder) error { return errors.Join(errs...) } +// DeleteACLs deletes the specified ACLs from the Kafka cluster. +func (m *Manager) DeleteACLs(ctx context.Context, acls *kadm.ACLBuilder) error { + res, err := m.adminClient.DeleteACLs(ctx, acls) + if err != nil { + return fmt.Errorf("failed to delete ACLs: %w", err) + } + var errs []error + for _, r := range res { + if r.Err != nil { + errs = append(errs, r.Err) + } + } + return errors.Join(errs...) +} + // ListTopics returns all topics that begin with prefix in lexicographical order from the Kafka broker. func (m *Manager) ListTopics(ctx context.Context, prefix string) ([]string, error) { details, err := m.adminClient.ListTopics(ctx) diff --git a/kafka/manager_test.go b/kafka/manager_test.go index 106f47c8..c605606d 100644 --- a/kafka/manager_test.go +++ b/kafka/manager_test.go @@ -771,6 +771,78 @@ func TestManagerCreateACLs(t *testing.T) { }) } +func TestManagerDeleteACLs(t *testing.T) { + t.Run("Success", func(t *testing.T) { + cluster, commonConfig := newFakeCluster(t) + cluster.ControlKey(kmsg.DeleteACLs.Int16(), func(req kmsg.Request) (kmsg.Response, error, bool) { + return &kmsg.DeleteACLsResponse{ + Version: req.GetVersion(), + Results: []kmsg.DeleteACLsResponseResult{ + {}, // Empty result means success + }, + }, nil, true + }) + m, err := NewManager(ManagerConfig{CommonConfig: commonConfig}) + require.NoError(t, err) + t.Cleanup(func() { m.Close() }) + + // For delete filters, Allow() must be paired with AllowHosts(). + // Explicit "*" matches only ACLs stored with host "*". + acls := kadm.NewACLs(). + Allow("User:*"). + AllowHosts("*"). + Topics("topic"). + Operations(kadm.OpRead). + ResourcePatternType(kadm.ACLPatternPrefixed) + + cluster.ControlKey(kmsg.ApiVersions.Int16(), func(req kmsg.Request) (kmsg.Response, error, bool) { + return &kmsg.ApiVersionsResponse{ + Version: req.GetVersion(), + ApiKeys: []kmsg.ApiVersionsResponseApiKey{ + {ApiKey: kmsg.DeleteACLs.Int16(), MaxVersion: 3}, + }, + }, nil, true + }) + + err = m.DeleteACLs(context.Background(), acls) + assert.NoError(t, err) + }) + t.Run("Partial Failure", func(t *testing.T) { + cluster, commonConfig := newFakeCluster(t) + respErr := kerr.InvalidPrincipalType + cluster.ControlKey(kmsg.DeleteACLs.Int16(), func(req kmsg.Request) (kmsg.Response, error, bool) { + return &kmsg.DeleteACLsResponse{ + Version: req.GetVersion(), + Results: []kmsg.DeleteACLsResponseResult{ + {ErrorCode: respErr.Code, ErrorMessage: &respErr.Message}, + }, + }, nil, true + }) + m, err := NewManager(ManagerConfig{CommonConfig: commonConfig}) + require.NoError(t, err) + t.Cleanup(func() { m.Close() }) + + acls := kadm.NewACLs(). + Allow("User:*"). + AllowHosts(). + Topics("topic"). + Operations(kadm.OpRead). + ResourcePatternType(kadm.ACLPatternPrefixed) + + cluster.ControlKey(kmsg.ApiVersions.Int16(), func(req kmsg.Request) (kmsg.Response, error, bool) { + return &kmsg.ApiVersionsResponse{ + Version: req.GetVersion(), + ApiKeys: []kmsg.ApiVersionsResponseApiKey{ + {ApiKey: kmsg.DeleteACLs.Int16(), MaxVersion: 3}, + }, + }, nil, true + }) + + err = m.DeleteACLs(context.Background(), acls) + assert.EqualError(t, err, respErr.Error()) + }) +} + func TestListTopics(t *testing.T) { cluster, commonConfig := newFakeCluster(t) m, err := NewManager(ManagerConfig{CommonConfig: commonConfig})