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)
 }

Reply via email to