Skip to content

Commit 957ac18

Browse files
committed
chore(acl): Improve error handling
Signed-off-by: Viacheslav Vasilyev <avoidik@gmail.com>
1 parent 35a1ada commit 957ac18

2 files changed

Lines changed: 129 additions & 7 deletions

File tree

internal/clients/kafka/acl/acl.go

Lines changed: 30 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,14 @@ import (
1414
"github.com/crossplane-contrib/provider-kafka/internal/clients/kafka"
1515
)
1616

17+
// adminClient is the subset of kadm.Client methods used by this package.
18+
// *kadm.Client satisfies this interface without any changes to callers.
19+
type adminClient interface {
20+
CreateACLs(ctx context.Context, b *kadm.ACLBuilder) (kadm.CreateACLsResults, error)
21+
DeleteACLs(ctx context.Context, b *kadm.ACLBuilder) (kadm.DeleteACLsResults, error)
22+
DescribeACLs(ctx context.Context, b *kadm.ACLBuilder) (kadm.DescribeACLsResults, error)
23+
}
24+
1725
// AccessControlList is a holistic representation of a Kafka ACL with configurable
1826
// fields
1927
type AccessControlList struct {
@@ -58,7 +66,7 @@ func buildACLBuilder(accessControlList *AccessControlList) (*kadm.ACLBuilder, er
5866
}
5967

6068
// List lists all the ACLs in Kafka
61-
func List(ctx context.Context, cl *kadm.Client, accessControlList *AccessControlList) (*AccessControlList, error) {
69+
func List(ctx context.Context, cl adminClient, accessControlList *AccessControlList) (*AccessControlList, error) {
6270
ab, err := buildACLBuilder(accessControlList)
6371
if err != nil {
6472
return nil, err
@@ -68,7 +76,13 @@ func List(ctx context.Context, cl *kadm.Client, accessControlList *AccessControl
6876
if err != nil {
6977
return nil, fmt.Errorf("describe ACLs failed: %w", err)
7078
}
71-
if len(resp) == 0 || len(resp[0].Described) == 0 {
79+
if len(resp) == 0 {
80+
return nil, nil
81+
}
82+
if resp[0].Err != nil {
83+
return nil, fmt.Errorf("describe ACLs failed: %w", resp[0].Err)
84+
}
85+
if len(resp[0].Described) == 0 {
7286
return nil, nil
7387
}
7488

@@ -84,7 +98,7 @@ func List(ctx context.Context, cl *kadm.Client, accessControlList *AccessControl
8498
}
8599

86100
// Create creates an ACL from the Kafka side
87-
func Create(ctx context.Context, cl *kadm.Client, accessControlList *AccessControlList) error {
101+
func Create(ctx context.Context, cl adminClient, accessControlList *AccessControlList) error {
88102
ab, err := buildACLBuilder(accessControlList)
89103
if err != nil {
90104
return err
@@ -94,22 +108,31 @@ func Create(ctx context.Context, cl *kadm.Client, accessControlList *AccessContr
94108
if err != nil {
95109
return err
96110
}
97-
if len(resp) == 0 || len(resp[0].Principal) == 0 {
111+
if len(resp) == 0 {
98112
return errors.New("no create response for acl")
99113
}
114+
if resp[0].Err != nil {
115+
return fmt.Errorf("create ACL failed: %w", resp[0].Err)
116+
}
100117

101118
return nil
102119
}
103120

104121
// Delete deletes an ACL from the Kafka side
105-
func Delete(ctx context.Context, cl *kadm.Client, accessControlList *AccessControlList) error {
122+
func Delete(ctx context.Context, cl adminClient, accessControlList *AccessControlList) error {
106123
ab, err := buildACLBuilder(accessControlList)
107124
if err != nil {
108125
return err
109126
}
110127

111-
_, err = cl.DeleteACLs(ctx, ab)
112-
return err
128+
resp, err := cl.DeleteACLs(ctx, ab)
129+
if err != nil {
130+
return err
131+
}
132+
if len(resp) > 0 && resp[0].Err != nil {
133+
return fmt.Errorf("delete ACL failed: %w", resp[0].Err)
134+
}
135+
return nil
113136
}
114137

115138
// ConvertToJSON performs a json marshalling for ACLs

internal/clients/kafka/acl/acl_test.go

Lines changed: 99 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,12 +8,36 @@ import (
88
"time"
99

1010
"github.com/google/go-cmp/cmp"
11+
"github.com/twmb/franz-go/pkg/kadm"
12+
"github.com/twmb/franz-go/pkg/kerr"
1113
"k8s.io/apimachinery/pkg/util/json"
1214

1315
"github.com/crossplane-contrib/provider-kafka/apis/v1alpha1"
1416
"github.com/crossplane-contrib/provider-kafka/internal/clients/kafka"
1517
)
1618

19+
// fakeACLAdmin is an in-process implementation of adminClient for unit tests.
20+
type fakeACLAdmin struct {
21+
createResults kadm.CreateACLsResults
22+
createErr error
23+
describeResults kadm.DescribeACLsResults
24+
describeErr error
25+
deleteResults kadm.DeleteACLsResults
26+
deleteErr error
27+
}
28+
29+
func (f *fakeACLAdmin) CreateACLs(_ context.Context, _ *kadm.ACLBuilder) (kadm.CreateACLsResults, error) {
30+
return f.createResults, f.createErr
31+
}
32+
33+
func (f *fakeACLAdmin) DescribeACLs(_ context.Context, _ *kadm.ACLBuilder) (kadm.DescribeACLsResults, error) {
34+
return f.describeResults, f.describeErr
35+
}
36+
37+
func (f *fakeACLAdmin) DeleteACLs(_ context.Context, _ *kadm.ACLBuilder) (kadm.DeleteACLsResults, error) {
38+
return f.deleteResults, f.deleteErr
39+
}
40+
1741
var dataTesting = []byte(os.Getenv("KAFKA_CONFIG"))
1842

1943
var baseACL = AccessControlList{
@@ -477,6 +501,81 @@ func TestList(t *testing.T) {
477501
}
478502
}
479503

504+
// --- Unit tests for broker-level error propagation (no real Kafka needed) ---
505+
506+
func TestCreateBrokerError(t *testing.T) {
507+
cl := &fakeACLAdmin{
508+
createResults: kadm.CreateACLsResults{
509+
{Principal: "User:alice", Err: kerr.ClusterAuthorizationFailed},
510+
},
511+
}
512+
err := Create(context.Background(), cl, &baseACL)
513+
if err == nil {
514+
t.Fatal("Create() expected error for broker-level Err, got nil")
515+
}
516+
}
517+
518+
func TestCreateEmptyResponse(t *testing.T) {
519+
cl := &fakeACLAdmin{createResults: kadm.CreateACLsResults{}}
520+
err := Create(context.Background(), cl, &baseACL)
521+
if err == nil {
522+
t.Fatal("Create() expected error for empty response, got nil")
523+
}
524+
}
525+
526+
func TestListBrokerError(t *testing.T) {
527+
cl := &fakeACLAdmin{
528+
describeResults: kadm.DescribeACLsResults{
529+
{Err: kerr.ClusterAuthorizationFailed},
530+
},
531+
}
532+
got, err := List(context.Background(), cl, &baseACL)
533+
if err == nil {
534+
t.Fatal("List() expected error for broker-level Err, got nil")
535+
}
536+
if got != nil {
537+
t.Errorf("List() expected nil result on error, got %v", got)
538+
}
539+
}
540+
541+
func TestListNotFound(t *testing.T) {
542+
cl := &fakeACLAdmin{
543+
describeResults: kadm.DescribeACLsResults{
544+
{Described: kadm.DescribedACLs{}},
545+
},
546+
}
547+
got, err := List(context.Background(), cl, &baseACL)
548+
if err != nil {
549+
t.Fatalf("List() unexpected error: %v", err)
550+
}
551+
if got != nil {
552+
t.Errorf("List() expected nil for no matching ACL, got %v", got)
553+
}
554+
}
555+
556+
func TestListEmptyResponse(t *testing.T) {
557+
cl := &fakeACLAdmin{describeResults: kadm.DescribeACLsResults{}}
558+
got, err := List(context.Background(), cl, &baseACL)
559+
if err != nil {
560+
t.Fatalf("List() unexpected error: %v", err)
561+
}
562+
if got != nil {
563+
t.Errorf("List() expected nil for empty response, got %v", got)
564+
}
565+
}
566+
567+
func TestDeleteBrokerError(t *testing.T) {
568+
cl := &fakeACLAdmin{
569+
deleteResults: kadm.DeleteACLsResults{
570+
{Err: kerr.ClusterAuthorizationFailed},
571+
},
572+
}
573+
err := Delete(context.Background(), cl, &baseACL)
574+
if err == nil {
575+
t.Fatal("Delete() expected error for broker-level Err, got nil")
576+
}
577+
}
578+
480579
// TestListAtProviderNotFound verifies that List returns nil when the ACL does not exist.
481580
func TestListAtProviderNotFound(t *testing.T) {
482581
if len(dataTesting) == 0 {

0 commit comments

Comments
 (0)