Skip to content

Commit 7417623

Browse files
committed
fix(acl): Improve error handling
Signed-off-by: Viacheslav Vasilyev <avoidik@gmail.com>
1 parent 128ffcd commit 7417623

3 files changed

Lines changed: 125 additions & 7 deletions

File tree

Makefile

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -224,6 +224,7 @@ kind-kafka-setup: $(HELM) $(KIND) $(KUBECTL)
224224
@$(KUBECTL) wait --for=condition=ready -n kafka-cluster kafka/dev --timeout=300s
225225
@$(KUBECTL) wait --for=condition=ready -n kafka-cluster kafkauser/user --timeout=300s
226226
@$(KUBECTL) wait --for=condition=ready -n kafka-cluster kafkatopic/pre-existing --timeout=300s
227+
@$(KUBECTL) annotate kafkatopic --all -n kafka-cluster strimzi.io/pause-reconciliation=true --overwrite
227228
@$(INFO) Waiting for Kafka broker to be fully initialized...
228229
@for i in $$(seq 1 30); do $(KUBECTL) -n kafka-cluster exec kafka-dev-0 -- bash -c 'echo "test" | kafka-console-producer.sh --broker-list localhost:9092 --topic pre-existing' > /dev/null 2>&1 && break || sleep 2; done
229230
@$(INFO) Getting service IP and port

internal/clients/kafka/acl/acl.go

Lines changed: 35 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,36 @@ 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 {
112+
return errors.New("no create response for acl")
113+
}
114+
if resp[0].Err != nil {
115+
return fmt.Errorf("create ACL failed: %w", resp[0].Err)
116+
}
117+
if len(resp[0].Principal) == 0 {
98118
return errors.New("no create response for acl")
99119
}
100120

101121
return nil
102122
}
103123

104124
// Delete deletes an ACL from the Kafka side
105-
func Delete(ctx context.Context, cl *kadm.Client, accessControlList *AccessControlList) error {
125+
func Delete(ctx context.Context, cl adminClient, accessControlList *AccessControlList) error {
106126
ab, err := buildACLBuilder(accessControlList)
107127
if err != nil {
108128
return err
109129
}
110130

111-
_, err = cl.DeleteACLs(ctx, ab)
112-
return err
131+
resp, err := cl.DeleteACLs(ctx, ab)
132+
if err != nil {
133+
return err
134+
}
135+
for _, r := range resp {
136+
if r.Err != nil {
137+
return fmt.Errorf("delete ACL failed: %w", r.Err)
138+
}
139+
}
140+
return nil
113141
}
114142

115143
// ConvertToJSON performs a json marshalling for ACLs

internal/clients/kafka/acl/acl_test.go

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

1010
"github.com/google/go-cmp/cmp"
11+
"github.com/stretchr/testify/assert"
12+
"github.com/stretchr/testify/require"
13+
"github.com/twmb/franz-go/pkg/kadm"
14+
"github.com/twmb/franz-go/pkg/kerr"
1115
"k8s.io/apimachinery/pkg/util/json"
1216

1317
"github.com/crossplane-contrib/provider-kafka/apis/v1alpha1"
1418
"github.com/crossplane-contrib/provider-kafka/internal/clients/kafka"
1519
)
1620

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

1945
var baseACL = AccessControlList{
@@ -477,6 +503,69 @@ func TestList(t *testing.T) {
477503
}
478504
}
479505

506+
// --- Unit tests for broker-level error propagation (no real Kafka needed) ---
507+
508+
func TestCreateBrokerError(t *testing.T) {
509+
t.Parallel()
510+
cl := &fakeACLAdmin{
511+
createResults: kadm.CreateACLsResults{
512+
{Principal: "User:alice", Err: kerr.ClusterAuthorizationFailed},
513+
},
514+
}
515+
err := Create(context.Background(), cl, &baseACL)
516+
require.Error(t, err)
517+
}
518+
519+
func TestCreateEmptyResponse(t *testing.T) {
520+
t.Parallel()
521+
cl := &fakeACLAdmin{createResults: kadm.CreateACLsResults{}}
522+
err := Create(context.Background(), cl, &baseACL)
523+
require.Error(t, err)
524+
}
525+
526+
func TestListBrokerError(t *testing.T) {
527+
t.Parallel()
528+
cl := &fakeACLAdmin{
529+
describeResults: kadm.DescribeACLsResults{
530+
{Err: kerr.ClusterAuthorizationFailed},
531+
},
532+
}
533+
got, err := List(context.Background(), cl, &baseACL)
534+
require.Error(t, err)
535+
assert.Nil(t, got)
536+
}
537+
538+
func TestListNotFound(t *testing.T) {
539+
t.Parallel()
540+
cl := &fakeACLAdmin{
541+
describeResults: kadm.DescribeACLsResults{
542+
{Described: kadm.DescribedACLs{}},
543+
},
544+
}
545+
got, err := List(context.Background(), cl, &baseACL)
546+
require.NoError(t, err)
547+
assert.Nil(t, got)
548+
}
549+
550+
func TestListEmptyResponse(t *testing.T) {
551+
t.Parallel()
552+
cl := &fakeACLAdmin{describeResults: kadm.DescribeACLsResults{}}
553+
got, err := List(context.Background(), cl, &baseACL)
554+
require.NoError(t, err)
555+
assert.Nil(t, got)
556+
}
557+
558+
func TestDeleteBrokerError(t *testing.T) {
559+
t.Parallel()
560+
cl := &fakeACLAdmin{
561+
deleteResults: kadm.DeleteACLsResults{
562+
{Err: kerr.ClusterAuthorizationFailed},
563+
},
564+
}
565+
err := Delete(context.Background(), cl, &baseACL)
566+
require.Error(t, err)
567+
}
568+
480569
// TestListAtProviderNotFound verifies that List returns nil when the ACL does not exist.
481570
func TestListAtProviderNotFound(t *testing.T) {
482571
if len(dataTesting) == 0 {

0 commit comments

Comments
 (0)