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