Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 15 additions & 0 deletions kafka/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
72 changes: 72 additions & 0 deletions kafka/manager_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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})
Expand Down