This is an automated email from the ASF dual-hosted git repository. squakez pushed a commit to branch main in repository https://gitbox.apache.org/repos/asf/camel-k.git
commit 7820173acdab8288f16399d9739d5909516cd112 Author: harbinresearcher <[email protected]> AuthorDate: Thu Oct 1 13:18:30 2026 +0800 chore: expand Strimzi binding test coverage Signed-off-by: harbinresearcher <[email protected]> --- pkg/util/bindings/strimzi.go | 65 +++++++------ pkg/util/bindings/strimzi_test.go | 196 +++++++++++++++----------------------- 2 files changed, 111 insertions(+), 150 deletions(-) diff --git a/pkg/util/bindings/strimzi.go b/pkg/util/bindings/strimzi.go index 3ddeb491b..bf51296fd 100644 --- a/pkg/util/bindings/strimzi.go +++ b/pkg/util/bindings/strimzi.go @@ -23,13 +23,11 @@ import ( camelv1 "github.com/apache/camel-k/v2/pkg/apis/camel/v1" strimziv1 "github.com/apache/camel-k/v2/pkg/apis/duck/strimzi/v1" - "github.com/apache/camel-k/v2/pkg/util/kubernetes" + "github.com/apache/camel-k/v2/pkg/client/strimzi/clientset/internalclientset" "github.com/apache/camel-k/v2/pkg/util/uri" k8serrors "k8s.io/apimachinery/pkg/api/errors" - "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" - "k8s.io/apimachinery/pkg/runtime" + v1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime/schema" - "sigs.k8s.io/controller-runtime/pkg/client" ) func init() { @@ -43,7 +41,9 @@ type camelKafka struct { } // StrimziBindingProvider allows to connect to a Kafka topic via Binding. -type StrimziBindingProvider struct{} +type StrimziBindingProvider struct { + Client internalclientset.Interface +} func (s StrimziBindingProvider) ID() string { return "strimzi" @@ -98,7 +98,16 @@ func (s StrimziBindingProvider) fromKafkaToCamel(ctx BindingContext, endpoint ca } topicName := props["topic"] delete(props, "topic") + //nolint:nestif if props["brokers"] == "" { + // build the client if needed + if s.Client == nil { + kafkaClient, err := internalclientset.NewForConfig(ctx.Client.GetConfig()) + if err != nil { + return nil, err + } + s.Client = kafkaClient + } namespace := endpoint.Ref.Namespace if namespace == "" { namespace = ctx.Namespace @@ -142,6 +151,15 @@ func (s StrimziBindingProvider) fromKafkaTopicToCamel(ctx BindingContext, endpoi } func (s StrimziBindingProvider) lookupBootstrapServers(ctx BindingContext, endpoint camelv1.Endpoint) (string, error) { + // build the client if needed + if s.Client == nil { + kafkaClient, err := internalclientset.NewForConfig(ctx.Client.GetConfig()) + if err != nil { + return "", err + } + s.Client = kafkaClient + } + topic, err := s.lookupTopic(ctx, endpoint) if err != nil { return "", err @@ -164,16 +182,10 @@ func (s StrimziBindingProvider) lookupBootstrapServers(ctx BindingContext, endpo } func (s StrimziBindingProvider) getBootstrapServers(ctx BindingContext, clusterName, namespace string) (string, error) { - // Read without the manager cache to preserve live, cross-namespace lookups without requiring watch permissions. - object, err := kubernetes.GetUnstructured(ctx.Ctx, ctx.Client, - strimziv1.SchemeGroupVersion.WithKind(strimziv1.StrimziKindKafkaCluster), clusterName, namespace) + cluster, err := s.Client.KafkaV1().Kafkas(namespace).Get(ctx.Ctx, clusterName, v1.GetOptions{}) if err != nil { return "", err } - cluster := strimziv1.Kafka{} - if err := runtime.DefaultUnstructuredConverter.FromUnstructured(object.Object, &cluster); err != nil { - return "", err - } for _, l := range cluster.Status.Listeners { if l.Name == strimziv1.StrimziListenerNamePlain { @@ -194,40 +206,27 @@ func (s StrimziBindingProvider) lookupTopic(ctx BindingContext, endpoint camelv1 namespace = ctx.Namespace } // first check by KafkaTopic name - object, err := kubernetes.GetUnstructured(ctx.Ctx, ctx.Client, - strimziv1.SchemeGroupVersion.WithKind(strimziv1.StrimziKindTopic), endpoint.Ref.Name, namespace) + topic, err := s.Client.KafkaV1().KafkaTopics(namespace).Get(ctx.Ctx, endpoint.Ref.Name, v1.GetOptions{}) if err != nil && !k8serrors.IsNotFound(err) { return nil, err } if err == nil { - topic := strimziv1.KafkaTopic{} - if err := runtime.DefaultUnstructuredConverter.FromUnstructured(object.Object, &topic); err != nil { - return nil, err - } - - return &topic, nil + return topic, nil } // if not found, then, look at the .status.topicName (it may be autogenerated) - // Unstructured lists also bypass the cache, without introducing a custom field index. - objects := unstructured.UnstructuredList{} - objects.SetGroupVersionKind(strimziv1.SchemeGroupVersion.WithKind("KafkaTopicList")) - err = ctx.Client.List(ctx.Ctx, &objects, client.InNamespace(namespace)) + topics, err := s.Client.KafkaV1().KafkaTopics(namespace).List(ctx.Ctx, v1.ListOptions{ + FieldSelector: "status.topicName=" + endpoint.Ref.Name, + }) if err != nil { return nil, fmt.Errorf("couldn't find any KafkaTopic with either name or topicName %s; error %w", endpoint.Ref.Name, err) } - for i := range objects.Items { - topic := strimziv1.KafkaTopic{} - if err := runtime.DefaultUnstructuredConverter.FromUnstructured(objects.Items[i].Object, &topic); err != nil { - return nil, err - } - if topic.Status.TopicName == endpoint.Ref.Name { - return &topic, nil - } + if len(topics.Items) == 0 { + return nil, fmt.Errorf("couldn't find any KafkaTopic with either name or topicName %s", endpoint.Ref.Name) } - return nil, fmt.Errorf("couldn't find any KafkaTopic with either name or topicName %s", endpoint.Ref.Name) + return &topics.Items[0], nil } // Order --. diff --git a/pkg/util/bindings/strimzi_test.go b/pkg/util/bindings/strimzi_test.go index 26efccf2e..09d54cfc9 100644 --- a/pkg/util/bindings/strimzi_test.go +++ b/pkg/util/bindings/strimzi_test.go @@ -20,23 +20,20 @@ package bindings import ( "context" "errors" - "net/http" - "net/http/httptest" "testing" camelv1 "github.com/apache/camel-k/v2/pkg/apis/camel/v1" strimziv1 "github.com/apache/camel-k/v2/pkg/apis/duck/strimzi/v1" + "github.com/apache/camel-k/v2/pkg/client/strimzi/clientset/internalclientset/fake" "github.com/apache/camel-k/v2/pkg/internal" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" v1 "k8s.io/api/core/v1" k8serrors "k8s.io/apimachinery/pkg/api/errors" - "k8s.io/apimachinery/pkg/api/meta" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/runtime/schema" - "k8s.io/client-go/rest" - ctrl "sigs.k8s.io/controller-runtime/pkg/client" + "k8s.io/apimachinery/pkg/runtime" + k8stesting "k8s.io/client-go/testing" ) func TestStrimziDirect(t *testing.T) { @@ -105,15 +102,13 @@ func TestStrimziLookup(t *testing.T) { }, } - client, err := internal.NewFakeClient() - require.NoError(t, err) - require.NoError(t, client.Create(ctx, &cluster)) - require.NoError(t, client.Create(ctx, &topic)) - provider := StrimziBindingProvider{} + client := fake.NewSimpleClientset(&cluster, &topic) + provider := StrimziBindingProvider{ + Client: client, + } bindingContext := BindingContext{ Ctx: ctx, - Client: client, Namespace: "test", Profile: camelv1.TraitProfileKubernetes, } @@ -170,15 +165,13 @@ func TestStrimziLookupByTopicName(t *testing.T) { }, } - client, err := internal.NewFakeClient() - require.NoError(t, err) - require.NoError(t, client.Create(ctx, &cluster)) - require.NoError(t, client.Create(ctx, &topic)) - provider := StrimziBindingProvider{} + client := fake.NewSimpleClientset(&cluster, &topic) + provider := StrimziBindingProvider{ + Client: client, + } bindingContext := BindingContext{ Ctx: ctx, - Client: client, Namespace: "test", Profile: camelv1.TraitProfileKubernetes, } @@ -222,14 +215,13 @@ func TestStrimziKafkaCR(t *testing.T) { }, } - client, err := internal.NewFakeClient() - require.NoError(t, err) - require.NoError(t, client.Create(ctx, &cluster)) - provider := StrimziBindingProvider{} + client := fake.NewSimpleClientset(&cluster) + provider := StrimziBindingProvider{ + Client: client, + } bindingContext := BindingContext{ Ctx: ctx, - Client: client, Namespace: "test", Profile: camelv1.TraitProfileKubernetes, } @@ -276,14 +268,13 @@ func TestStrimziPassThrough(t *testing.T) { }, } - client, err := internal.NewFakeClient() - require.NoError(t, err) - require.NoError(t, client.Create(ctx, &cluster)) - provider := StrimziBindingProvider{} + client := fake.NewSimpleClientset(&cluster) + provider := StrimziBindingProvider{ + Client: client, + } bindingContext := BindingContext{ Ctx: ctx, - Client: client, Namespace: "test", Profile: camelv1.TraitProfileKubernetes, } @@ -303,120 +294,91 @@ func TestStrimziPassThrough(t *testing.T) { assert.Nil(t, binding) } -func TestStrimziLookupTopicNamespace(t *testing.T) { +func TestStrimziNamespace(t *testing.T) { for _, tc := range []struct { name string namespace string - topicName string - wantName string + wantURI string }{ - {name: "context namespace", topicName: "shared", wantName: "local"}, - {name: "explicit namespace", namespace: "other", topicName: "shared", wantName: "remote"}, - {name: "resource name", topicName: "local", wantName: "local"}, - {name: "missing topic", topicName: "missing"}, + {name: "context namespace", wantURI: "kafka:topic?brokers=local%3A9092"}, + {name: "explicit namespace", namespace: "other", wantURI: "kafka:topic?brokers=remote%3A9092"}, } { t.Run(tc.name, func(t *testing.T) { - ctx := context.Background() - client, err := internal.NewFakeClient() - require.NoError(t, err) - for _, topic := range []strimziv1.KafkaTopic{ - {ObjectMeta: metav1.ObjectMeta{Name: "aaa-unrelated", Namespace: "test"}}, - {ObjectMeta: metav1.ObjectMeta{Name: "local", Namespace: "test"}, Status: strimziv1.KafkaTopicStatus{TopicName: "shared"}}, - {ObjectMeta: metav1.ObjectMeta{Name: "remote", Namespace: "other"}, Status: strimziv1.KafkaTopicStatus{TopicName: "shared"}}, - } { - require.NoError(t, client.Create(ctx, &topic)) - } - topic, err := (StrimziBindingProvider{}).lookupTopic(BindingContext{ - Ctx: ctx, Client: client, Namespace: "test", - }, camelv1.Endpoint{Ref: &v1.ObjectReference{Name: tc.topicName, Namespace: tc.namespace}}) - if tc.wantName == "" { - require.Error(t, err) - assert.Nil(t, topic) - return - } + client := fake.NewSimpleClientset( + &strimziv1.KafkaTopic{Name: "topic", Namespace: "test", Labels: map[string]string{strimziv1.StrimziKafkaClusterLabel: "local-cluster"}}, + &strimziv1.KafkaTopic{Name: "topic", Namespace: "other", Labels: map[string]string{strimziv1.StrimziKafkaClusterLabel: "remote-cluster"}}, + &strimziv1.Kafka{Name: "local-cluster", Namespace: "test", Status: strimziv1.KafkaStatus{Listeners: []strimziv1.KafkaStatusListener{{Name: "plain", BootstrapServers: "local:9092"}}}}, + &strimziv1.Kafka{Name: "remote-cluster", Namespace: "other", Status: strimziv1.KafkaStatus{Listeners: []strimziv1.KafkaStatusListener{{Name: "plain", BootstrapServers: "remote:9092"}}}}, + ) + binding, err := (StrimziBindingProvider{Client: client}).Translate(BindingContext{ + Ctx: context.Background(), Namespace: "test", + }, EndpointContext{}, camelv1.Endpoint{Ref: &v1.ObjectReference{ + APIVersion: "kafka.strimzi.io/v1", Kind: "KafkaTopic", Name: "topic", Namespace: tc.namespace, + }}) require.NoError(t, err) - require.NotNil(t, topic) - assert.Equal(t, tc.wantName, topic.Name) + require.NotNil(t, binding) + assert.Equal(t, tc.wantURI, binding.URI) }) } } func TestStrimziMissingCluster(t *testing.T) { - client, err := internal.NewFakeClient() - require.NoError(t, err) - servers, err := (StrimziBindingProvider{}).getBootstrapServers(BindingContext{ - Ctx: context.Background(), Client: client, - }, "missing", "test") + provider := StrimziBindingProvider{Client: fake.NewSimpleClientset()} + servers, err := provider.getBootstrapServers(BindingContext{Ctx: context.Background()}, "missing", "test") require.Error(t, err) assert.True(t, k8serrors.IsNotFound(err)) assert.Empty(t, servers) } -func TestStrimziBypassesCache(t *testing.T) { +func TestStrimziListenerValidation(t *testing.T) { for _, tc := range []struct { - name string - kind string - refName string + name string + listeners []strimziv1.KafkaStatusListener + wantError string }{ - {name: "cluster", kind: "Kafka", refName: "cluster"}, - {name: "topic", kind: "KafkaTopic", refName: "topic"}, - {name: "topic status name", kind: "KafkaTopic", refName: "shared"}, + {name: "no listeners", wantError: `cluster "cluster" has no listeners of name "plain"`}, + {name: "only TLS", listeners: []strimziv1.KafkaStatusListener{{Name: "tls", BootstrapServers: "tls:9093"}}, wantError: `cluster "cluster" has no listeners of name "plain"`}, + {name: "empty bootstrap servers", listeners: []strimziv1.KafkaStatusListener{{Name: "plain"}}, wantError: `cluster "cluster" has no bootstrap servers in "plain" listener`}, } { t.Run(tc.name, func(t *testing.T) { - server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - w.Header().Set("Content-Type", "application/json") - var response string - switch r.URL.Path { - case "/apis/kafka.strimzi.io/v1/namespaces/other/kafkas/cluster": - response = `{"apiVersion":"kafka.strimzi.io/v1","kind":"Kafka","metadata":{"name":"cluster","namespace":"other"},"status":{"listeners":[{"name":"plain","bootstrapServers":"live:9092"}]}}` - case "/apis/kafka.strimzi.io/v1/namespaces/other/kafkatopics/topic": - response = `{"apiVersion":"kafka.strimzi.io/v1","kind":"KafkaTopic","metadata":{"name":"topic","namespace":"other","labels":{"strimzi.io/cluster":"cluster"}}}` - case "/apis/kafka.strimzi.io/v1/namespaces/other/kafkatopics/shared": - w.WriteHeader(http.StatusNotFound) - response = `{"kind":"Status","apiVersion":"v1","status":"Failure","reason":"NotFound","code":404}` - case "/apis/kafka.strimzi.io/v1/namespaces/other/kafkatopics": - response = `{"apiVersion":"kafka.strimzi.io/v1","kind":"KafkaTopicList","items":[{"metadata":{"name":"topic","namespace":"other","labels":{"strimzi.io/cluster":"cluster"}},"status":{"topicName":"shared"}}]}` - default: - http.NotFound(w, r) - return - } - _, err := w.Write([]byte(response)) - assert.NoError(t, err) - })) - defer server.Close() - mapper := meta.NewDefaultRESTMapper([]schema.GroupVersion{strimziv1.SchemeGroupVersion}) - mapper.Add(strimziv1.SchemeGroupVersion.WithKind("Kafka"), meta.RESTScopeNamespace) - mapper.Add(strimziv1.SchemeGroupVersion.WithKind("KafkaTopic"), meta.RESTScopeNamespace) - client, err := internal.NewFakeClient() - require.NoError(t, err) - liveClient, err := ctrl.New(&rest.Config{Host: server.URL}, ctrl.Options{ - Scheme: client.GetScheme(), Mapper: mapper, - Cache: &ctrl.CacheOptions{Reader: strimziRejectingCache{}}, - }) - require.NoError(t, err) - client.(*internal.FakeClient).Client = liveClient - endpoint := camelv1.Endpoint{Ref: &v1.ObjectReference{ - APIVersion: "kafka.strimzi.io/v1", Kind: tc.kind, Name: tc.refName, Namespace: "other", - }} - if tc.kind == "Kafka" { - endpoint.Properties = asEndpointProperties(map[string]string{"topic": tc.refName}) - } - binding, err := (StrimziBindingProvider{}).Translate(BindingContext{ - Ctx: context.Background(), Client: client, Namespace: "test", - }, EndpointContext{}, endpoint) - require.NoError(t, err) - require.NotNil(t, binding) - assert.Equal(t, "kafka:"+tc.refName+"?brokers=live%3A9092", binding.URI) + provider := StrimziBindingProvider{Client: fake.NewSimpleClientset(&strimziv1.Kafka{ + Name: "cluster", Namespace: "test", + Status: strimziv1.KafkaStatus{Listeners: tc.listeners}, + })} + servers, err := provider.getBootstrapServers(BindingContext{Ctx: context.Background()}, "cluster", "test") + require.EqualError(t, err, tc.wantError) + assert.Empty(t, servers) }) } } -type strimziRejectingCache struct{} - -func (strimziRejectingCache) Get(context.Context, ctrl.ObjectKey, ctrl.Object, ...ctrl.GetOption) error { - return errors.New("Strimzi reads must not use the cache") +func TestStrimziLookupErrors(t *testing.T) { + for _, tc := range []struct { + name string + verb string + }{ + {name: "get failure", verb: "get"}, + {name: "list failure", verb: "list"}, + } { + t.Run(tc.name, func(t *testing.T) { + client := fake.NewSimpleClientset() + wantErr := errors.New("API request failed") + client.PrependReactor(tc.verb, "kafkatopics", func(k8stesting.Action) (bool, runtime.Object, error) { + return true, nil, wantErr + }) + topic, err := (StrimziBindingProvider{Client: client}).lookupTopic(BindingContext{ + Ctx: context.Background(), Namespace: "test", + }, camelv1.Endpoint{Ref: &v1.ObjectReference{Name: "missing"}}) + require.ErrorIs(t, err, wantErr) + assert.Nil(t, topic) + }) + } } -func (strimziRejectingCache) List(context.Context, ctrl.ObjectList, ...ctrl.ListOption) error { - return errors.New("Strimzi reads must not use the cache") +func TestStrimziMissingTopic(t *testing.T) { + provider := StrimziBindingProvider{Client: fake.NewSimpleClientset()} + topic, err := provider.lookupTopic(BindingContext{Ctx: context.Background(), Namespace: "test"}, + camelv1.Endpoint{Ref: &v1.ObjectReference{Name: "missing"}}) + require.EqualError(t, err, "couldn't find any KafkaTopic with either name or topicName missing") + assert.Nil(t, topic) }
