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 91e473ba705ea8718ee84aae20c871ed4d1f3d04
Author: harbinresearcher <[email protected]>
AuthorDate: Wed Sep 30 21:42:08 2026 +0800

    Fix #6846: Use context client for Strimzi bindings
    
    Signed-off-by: harbinresearcher <[email protected]>
---
 pkg/util/bindings/strimzi.go      |  65 +++++++--------
 pkg/util/bindings/strimzi_test.go | 165 ++++++++++++++++++++++++++++++++++----
 2 files changed, 181 insertions(+), 49 deletions(-)

diff --git a/pkg/util/bindings/strimzi.go b/pkg/util/bindings/strimzi.go
index bf51296fd..3ddeb491b 100644
--- a/pkg/util/bindings/strimzi.go
+++ b/pkg/util/bindings/strimzi.go
@@ -23,11 +23,13 @@ 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/client/strimzi/clientset/internalclientset"
+       "github.com/apache/camel-k/v2/pkg/util/kubernetes"
        "github.com/apache/camel-k/v2/pkg/util/uri"
        k8serrors "k8s.io/apimachinery/pkg/api/errors"
-       v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
+       "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
+       "k8s.io/apimachinery/pkg/runtime"
        "k8s.io/apimachinery/pkg/runtime/schema"
+       "sigs.k8s.io/controller-runtime/pkg/client"
 )
 
 func init() {
@@ -41,9 +43,7 @@ type camelKafka struct {
 }
 
 // StrimziBindingProvider allows to connect to a Kafka topic via Binding.
-type StrimziBindingProvider struct {
-       Client internalclientset.Interface
-}
+type StrimziBindingProvider struct{}
 
 func (s StrimziBindingProvider) ID() string {
        return "strimzi"
@@ -98,16 +98,7 @@ 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
@@ -151,15 +142,6 @@ 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
@@ -182,10 +164,16 @@ func (s StrimziBindingProvider) 
lookupBootstrapServers(ctx BindingContext, endpo
 }
 
 func (s StrimziBindingProvider) getBootstrapServers(ctx BindingContext, 
clusterName, namespace string) (string, error) {
-       cluster, err := s.Client.KafkaV1().Kafkas(namespace).Get(ctx.Ctx, 
clusterName, v1.GetOptions{})
+       // 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)
        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 {
@@ -206,27 +194,40 @@ func (s StrimziBindingProvider) lookupTopic(ctx 
BindingContext, endpoint camelv1
                namespace = ctx.Namespace
        }
        // first check by KafkaTopic name
-       topic, err := s.Client.KafkaV1().KafkaTopics(namespace).Get(ctx.Ctx, 
endpoint.Ref.Name, v1.GetOptions{})
+       object, err := kubernetes.GetUnstructured(ctx.Ctx, ctx.Client,
+               
strimziv1.SchemeGroupVersion.WithKind(strimziv1.StrimziKindTopic), 
endpoint.Ref.Name, namespace)
        if err != nil && !k8serrors.IsNotFound(err) {
                return nil, err
        }
        if err == nil {
-               return topic, nil
+               topic := strimziv1.KafkaTopic{}
+               if err := 
runtime.DefaultUnstructuredConverter.FromUnstructured(object.Object, &topic); 
err != nil {
+                       return nil, err
+               }
+
+               return &topic, nil
        }
 
        // if not found, then, look at the .status.topicName (it may be 
autogenerated)
-       topics, err := s.Client.KafkaV1().KafkaTopics(namespace).List(ctx.Ctx, 
v1.ListOptions{
-               FieldSelector: "status.topicName=" + endpoint.Ref.Name,
-       })
+       // 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))
 
        if err != nil {
                return nil, fmt.Errorf("couldn't find any KafkaTopic with 
either name or topicName %s; error %w", endpoint.Ref.Name, err)
        }
-       if len(topics.Items) == 0 {
-               return nil, fmt.Errorf("couldn't find any KafkaTopic with 
either name or topicName %s", endpoint.Ref.Name)
+       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
+               }
        }
 
-       return &topics.Items[0], nil
+       return nil, fmt.Errorf("couldn't find any KafkaTopic with either name 
or topicName %s", endpoint.Ref.Name)
 }
 
 // Order --.
diff --git a/pkg/util/bindings/strimzi_test.go 
b/pkg/util/bindings/strimzi_test.go
index 8a05021c6..26efccf2e 100644
--- a/pkg/util/bindings/strimzi_test.go
+++ b/pkg/util/bindings/strimzi_test.go
@@ -19,17 +19,24 @@ 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"
 )
 
 func TestStrimziDirect(t *testing.T) {
@@ -98,13 +105,15 @@ func TestStrimziLookup(t *testing.T) {
                },
        }
 
-       client := fake.NewSimpleClientset(&cluster, &topic)
-       provider := StrimziBindingProvider{
-               Client: client,
-       }
+       client, err := internal.NewFakeClient()
+       require.NoError(t, err)
+       require.NoError(t, client.Create(ctx, &cluster))
+       require.NoError(t, client.Create(ctx, &topic))
+       provider := StrimziBindingProvider{}
 
        bindingContext := BindingContext{
                Ctx:       ctx,
+               Client:    client,
                Namespace: "test",
                Profile:   camelv1.TraitProfileKubernetes,
        }
@@ -161,13 +170,15 @@ func TestStrimziLookupByTopicName(t *testing.T) {
                },
        }
 
-       client := fake.NewSimpleClientset(&cluster, &topic)
-       provider := StrimziBindingProvider{
-               Client: client,
-       }
+       client, err := internal.NewFakeClient()
+       require.NoError(t, err)
+       require.NoError(t, client.Create(ctx, &cluster))
+       require.NoError(t, client.Create(ctx, &topic))
+       provider := StrimziBindingProvider{}
 
        bindingContext := BindingContext{
                Ctx:       ctx,
+               Client:    client,
                Namespace: "test",
                Profile:   camelv1.TraitProfileKubernetes,
        }
@@ -211,13 +222,14 @@ func TestStrimziKafkaCR(t *testing.T) {
                },
        }
 
-       client := fake.NewSimpleClientset(&cluster)
-       provider := StrimziBindingProvider{
-               Client: client,
-       }
+       client, err := internal.NewFakeClient()
+       require.NoError(t, err)
+       require.NoError(t, client.Create(ctx, &cluster))
+       provider := StrimziBindingProvider{}
 
        bindingContext := BindingContext{
                Ctx:       ctx,
+               Client:    client,
                Namespace: "test",
                Profile:   camelv1.TraitProfileKubernetes,
        }
@@ -264,13 +276,14 @@ func TestStrimziPassThrough(t *testing.T) {
                },
        }
 
-       client := fake.NewSimpleClientset(&cluster)
-       provider := StrimziBindingProvider{
-               Client: client,
-       }
+       client, err := internal.NewFakeClient()
+       require.NoError(t, err)
+       require.NoError(t, client.Create(ctx, &cluster))
+       provider := StrimziBindingProvider{}
 
        bindingContext := BindingContext{
                Ctx:       ctx,
+               Client:    client,
                Namespace: "test",
                Profile:   camelv1.TraitProfileKubernetes,
        }
@@ -289,3 +302,121 @@ func TestStrimziPassThrough(t *testing.T) {
        require.NoError(t, err)
        assert.Nil(t, binding)
 }
+
+func TestStrimziLookupTopicNamespace(t *testing.T) {
+       for _, tc := range []struct {
+               name      string
+               namespace string
+               topicName string
+               wantName  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"},
+       } {
+               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
+                       }
+                       require.NoError(t, err)
+                       require.NotNil(t, topic)
+                       assert.Equal(t, tc.wantName, topic.Name)
+               })
+       }
+}
+
+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")
+       require.Error(t, err)
+       assert.True(t, k8serrors.IsNotFound(err))
+       assert.Empty(t, servers)
+}
+
+func TestStrimziBypassesCache(t *testing.T) {
+       for _, tc := range []struct {
+               name    string
+               kind    string
+               refName string
+       }{
+               {name: "cluster", kind: "Kafka", refName: "cluster"},
+               {name: "topic", kind: "KafkaTopic", refName: "topic"},
+               {name: "topic status name", kind: "KafkaTopic", refName: 
"shared"},
+       } {
+               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)
+               })
+       }
+}
+
+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 (strimziRejectingCache) List(context.Context, ctrl.ObjectList, 
...ctrl.ListOption) error {
+       return errors.New("Strimzi reads must not use the cache")
+}

Reply via email to