From edadfee47d3346874d4ba37b4e7aa495966e3e17 Mon Sep 17 00:00:00 2001 From: Lukasz Szaszkiewicz Date: Wed, 11 Jun 2025 09:45:38 +0200 Subject: [PATCH] test/apimachinery/watchlist: prove dynamic client's List method not streaming --- test/e2e/apimachinery/watchlist.go | 163 ++++++++++++++++------------- 1 file changed, 88 insertions(+), 75 deletions(-) diff --git a/test/e2e/apimachinery/watchlist.go b/test/e2e/apimachinery/watchlist.go index 7db72c4fc20..720e21aae72 100644 --- a/test/e2e/apimachinery/watchlist.go +++ b/test/e2e/apimachinery/watchlist.go @@ -30,6 +30,8 @@ import ( "github.com/onsi/gomega" v1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/api/meta" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime" @@ -44,6 +46,7 @@ import ( "k8s.io/client-go/rest" "k8s.io/client-go/tools/cache" "k8s.io/client-go/util/consistencydetector" + "k8s.io/client-go/util/watchlist" "k8s.io/component-base/featuregate" featuregatetesting "k8s.io/component-base/featuregate/testing" "k8s.io/kubernetes/test/e2e/framework" @@ -92,7 +95,7 @@ var _ = SIGDescribe("API Streaming (aka. WatchList)", framework.WithFeatureGate( secret, err = f.ClientSet.CoreV1().Secrets(f.Namespace.Name).Update(ctx, secret, metav1.UpdateOptions{}) framework.ExpectNoError(err) - expectedSecrets[0] = *secret + expectedSecrets[0] = secret verifyStore(ctx, expectedSecrets, secretInformer.GetStore()) }) ginkgo.It("should be requested by client-go's List method when WatchListClient is enabled", func(ctx context.Context) { @@ -110,13 +113,13 @@ var _ = SIGDescribe("API Streaming (aka. WatchList)", framework.WithFeatureGate( ginkgo.By("Verifying if the secret list was properly streamed") streamedSecrets := secretList.Items - gomega.Expect(cmp.Equal(expectedSecrets, streamedSecrets)).To(gomega.BeTrueBecause("data received via watchlist must match the added data")) + gomega.Expect(cmp.Equal(expectedSecrets, toSecretPointerSlice(streamedSecrets))).To(gomega.BeTrueBecause("data received via watchlist must match the added data")) ginkgo.By("Verifying if expected requests were sent to the server") expectedRequestsMadeByKubeClient := getExpectedRequestsMadeByClientFor(secretList.ResourceVersion) gomega.Expect(rt.actualRequests).To(gomega.Equal(expectedRequestsMadeByKubeClient)) }) - ginkgo.It("should be requested by dynamic client's List method when WatchListClient is enabled", func(ctx context.Context) { + ginkgo.It("should NOT be requested by dynamic client's List method when WatchListClient is enabled", func(ctx context.Context) { featuregatetesting.SetFeatureGateDuringTest(ginkgo.GinkgoTB(), utilfeature.DefaultFeatureGate, featuregate.Feature(clientfeatures.WatchListClient), true) ginkgo.By(fmt.Sprintf("Adding 5 secrets to %s namespace", f.Namespace.Name)) @@ -126,17 +129,17 @@ var _ = SIGDescribe("API Streaming (aka. WatchList)", framework.WithFeatureGate( wrappedDynamicClient, err := dynamic.NewForConfig(clientConfig) framework.ExpectNoError(err) - ginkgo.By("Streaming secrets from the server") + ginkgo.By("Getting secrets from the server") secretList, err := wrappedDynamicClient.Resource(v1.SchemeGroupVersion.WithResource("secrets")).Namespace(f.Namespace.Name).List(ctx, metav1.ListOptions{LabelSelector: "watchlist=true"}) framework.ExpectNoError(err) ginkgo.By("Verifying if the secret list was properly streamed") - streamedSecrets := secretList.Items - gomega.Expect(cmp.Equal(expectedSecrets, streamedSecrets)).To(gomega.BeTrueBecause("data received via watchlist must match the added data")) + actualSecrets := secretList.Items + gomega.Expect(cmp.Equal(expectedSecrets, toSecretPointerSlice(actualSecrets))).To(gomega.BeTrueBecause("data received via list must match the added data")) gomega.Expect(secretList.GetObjectKind().GroupVersionKind()).To(gomega.Equal(v1.SchemeGroupVersion.WithKind("SecretList"))) ginkgo.By("Verifying if expected requests were sent to the server") - expectedRequestsMadeByDynamicClient := getExpectedRequestsMadeByClientFor(secretList.GetResourceVersion()) + expectedRequestsMadeByDynamicClient := []string{expectedListRequestMadeByClient} gomega.Expect(rt.actualRequests).To(gomega.Equal(expectedRequestsMadeByDynamicClient)) }) ginkgo.It("should NOT be requested by metadata client's List method when WatchListClient is enabled", func(ctx context.Context) { @@ -170,35 +173,25 @@ var _ = SIGDescribe("API Streaming (aka. WatchList)", framework.WithFeatureGate( // Validates unsupported Accept headers in WatchList. // Sets AcceptContentType to "application/json;as=Table", which the API doesn't support, returning a 406 error. - // After the 406, the client falls back to a regular list request. ginkgo.It("doesn't support receiving resources as Tables", func(ctx context.Context) { featuregatetesting.SetFeatureGateDuringTest(ginkgo.GinkgoTB(), utilfeature.DefaultFeatureGate, featuregate.Feature(clientfeatures.WatchListClient), true) - ginkgo.By(fmt.Sprintf("Adding 5 secrets to %s namespace", f.Namespace.Name)) - _ = addWellKnownUnstructuredSecrets(ctx, f) - - rt, clientConfig := clientConfigWithRoundTripper(f) - modifiedClientConfig := dynamic.ConfigFor(clientConfig) + modifiedClientConfig := dynamic.ConfigFor(f.ClientConfig()) modifiedClientConfig.AcceptContentTypes = strings.Join([]string{ fmt.Sprintf("application/json;as=Table;v=%s;g=%s", metav1.SchemeGroupVersion.Version, metav1.GroupName), }, ",") modifiedClientConfig.GroupVersion = &v1.SchemeGroupVersion restClient, err := rest.RESTClientFor(modifiedClientConfig) framework.ExpectNoError(err) - wrappedDynamicClient := dynamic.New(restClient) + dynamicClient := dynamic.New(restClient) - // note that the client in case of an error (406) will fall back - // to a standard list request thus the overall call passes - ginkgo.By("Streaming secrets as Table from the server") - secretTable, err := wrappedDynamicClient.Resource(v1.SchemeGroupVersion.WithResource("secrets")).Namespace(f.Namespace.Name).List(ctx, metav1.ListOptions{LabelSelector: "watchlist=true"}) + opts, hasPreparedOptions, err := watchlist.PrepareWatchListOptionsFromListOptions(metav1.ListOptions{}) framework.ExpectNoError(err) - gomega.Expect(secretTable.GetObjectKind().GroupVersionKind()).To(gomega.Equal(metav1.SchemeGroupVersion.WithKind("Table"))) - - ginkgo.By("Verifying if expected response was sent by the server") - gomega.Expect(rt.actualResponseStatuses[0]).To(gomega.Equal("406 Not Acceptable")) - expectedRequestsMadeByDynamicClient := getExpectedRequestsMadeByClientWhenFallbackToListFor(secretTable.GetResourceVersion()) - gomega.Expect(rt.actualRequests).To(gomega.Equal(expectedRequestsMadeByDynamicClient)) + gomega.Expect(hasPreparedOptions).To(gomega.BeTrueBecause("it should be possible to prepare watchlist opts from an empty ListOptions")) + _, err = dynamicClient.Resource(v1.SchemeGroupVersion.WithResource("secrets")).Namespace("default").Watch(ctx, opts) + gomega.Expect(err).To(gomega.HaveOccurred()) + gomega.Expect(err.(apierrors.APIStatus)).To(gomega.HaveField("Status().Code", gomega.Equal(int32(406)))) }) // Sets AcceptContentType to both "application/json;as=Table" and "application/json". @@ -206,31 +199,46 @@ var _ = SIGDescribe("API Streaming (aka. WatchList)", framework.WithFeatureGate( ginkgo.It("falls backs to supported content type when when receiving resources as Tables was requested", func(ctx context.Context) { featuregatetesting.SetFeatureGateDuringTest(ginkgo.GinkgoTB(), utilfeature.DefaultFeatureGate, featuregate.Feature(clientfeatures.WatchListClient), true) - ginkgo.By(fmt.Sprintf("Adding 5 secrets to %s namespace", f.Namespace.Name)) - expectedSecrets := addWellKnownUnstructuredSecrets(ctx, f) - - rt, clientConfig := clientConfigWithRoundTripper(f) - modifiedClientConfig := dynamic.ConfigFor(clientConfig) + modifiedClientConfig := f.ClientConfig() modifiedClientConfig.AcceptContentTypes = strings.Join([]string{ fmt.Sprintf("application/json;as=Table;v=%s;g=%s", metav1.SchemeGroupVersion.Version, metav1.GroupName), "application/json", }, ",") modifiedClientConfig.GroupVersion = &v1.SchemeGroupVersion - restClient, err := rest.RESTClientFor(modifiedClientConfig) - framework.ExpectNoError(err) - wrappedDynamicClient := dynamic.New(restClient) - - ginkgo.By("Streaming secrets from the server") - secretList, err := wrappedDynamicClient.Resource(v1.SchemeGroupVersion.WithResource("secrets")).Namespace(f.Namespace.Name).List(ctx, metav1.ListOptions{LabelSelector: "watchlist=true"}) + dynamicClient, err := dynamic.NewForConfig(modifiedClientConfig) framework.ExpectNoError(err) - ginkgo.By("Verifying if the secret list was properly streamed") - streamedSecrets := secretList.Items - gomega.Expect(cmp.Equal(expectedSecrets, streamedSecrets)).To(gomega.BeTrueBecause("data received via watchlist must match the added data")) + stopCh := make(chan struct{}) + defer close(stopCh) - ginkgo.By("Verifying if expected requests were sent to the server") - expectedRequestsMadeByDynamicClient := getExpectedRequestsMadeByClientFor(secretList.GetResourceVersion()) - gomega.Expect(rt.actualRequests).To(gomega.Equal(expectedRequestsMadeByDynamicClient)) + secretInformer := cache.NewSharedIndexInformer( + &cache.ListWatch{ + ListFunc: func(options metav1.ListOptions) (runtime.Object, error) { + return nil, fmt.Errorf("unexpected list call") + }, + WatchFunc: func(options metav1.ListOptions) (watch.Interface, error) { + options.LabelSelector = "watchlist=true" + return dynamicClient.Resource(v1.SchemeGroupVersion.WithResource("secrets")).Namespace(f.Namespace.Name).Watch(context.TODO(), options) + }, + }, + &unstructured.Unstructured{}, + time.Duration(0), + nil, + ) + + expectedSecrets := addWellKnownUnstructuredSecrets(ctx, f) + + ginkgo.By("Starting the secret informer") + go secretInformer.Run(stopCh) + + ginkgo.By("Waiting until the secret informer is fully synchronised") + err = wait.PollUntilContextTimeout(ctx, 100*time.Millisecond, 30*time.Second, false, func(context.Context) (done bool, err error) { + return secretInformer.HasSynced(), nil + }) + framework.ExpectNoError(err, "Failed waiting for the secret informer in %s namespace to be synced", f.Namespace.Namespace) + + ginkgo.By("Verifying if the secret informer was properly synchronised") + verifyStore[unstructured.Unstructured](ctx, expectedSecrets, secretInformer.GetStore()) }) }) @@ -262,17 +270,19 @@ func clientConfigWithRoundTripper(f *framework.Framework) (*roundTripper, *rest. return rt, clientConfig } -func verifyStore(ctx context.Context, expectedSecrets []v1.Secret, store cache.Store) { +func verifyStore[T any](ctx context.Context, expectedSecrets []*T, store cache.Store) { err := wait.PollUntilContextTimeout(ctx, 100*time.Millisecond, 30*time.Second, true, func(ctx context.Context) (done bool, err error) { ginkgo.By("Comparing secrets retrieved directly from the server with the ones that have been streamed to the secret informer") - rawStreamedSecrets := store.List() - streamedSecrets := make([]v1.Secret, 0, len(rawStreamedSecrets)) - for _, rawSecret := range rawStreamedSecrets { - streamedSecrets = append(streamedSecrets, *rawSecret.(*v1.Secret)) - } - sort.Sort(byName(expectedSecrets)) - sort.Sort(byName(streamedSecrets)) - return cmp.Equal(expectedSecrets, streamedSecrets), nil + + expectedSecretsAsMetaObject, err := toMetaObjectSlice(expectedSecrets) + framework.ExpectNoError(err) + actualSecretsAsMetaObject, err := toMetaObjectSlice(store.List()) + framework.ExpectNoError(err) + + sort.Sort(byName(expectedSecretsAsMetaObject)) + sort.Sort(byName(actualSecretsAsMetaObject)) + + return cmp.Equal(expectedSecretsAsMetaObject, actualSecretsAsMetaObject), nil }) framework.ExpectNoError(err) } @@ -313,52 +323,35 @@ func getExpectedRequestsMadeByClientFor(rv string) []string { return expectedRequestMadeByClient } -func getExpectedRequestsMadeByClientWhenFallbackToListFor(rv string) []string { - expectedRequestMadeByClient := []string{ - expectedStreamingRequestMadeByClient, - // corresponds to a list request made by the client - func() string { - params := url.Values{} - params.Add("labelSelector", "watchlist=true") - return params.Encode() - }(), - } - if consistencydetector.IsDataConsistencyDetectionForListEnabled() { - // corresponds to a standard list request made by the consistency detector build in into the client - expectedRequestMadeByClient = append(expectedRequestMadeByClient, getExpectedListRequestMadeByConsistencyDetectorFor(rv)) - } - return expectedRequestMadeByClient -} - -func addWellKnownSecrets(ctx context.Context, f *framework.Framework) []v1.Secret { +func addWellKnownSecrets(ctx context.Context, f *framework.Framework) []*v1.Secret { ginkgo.By(fmt.Sprintf("Adding 5 secrets to %s namespace", f.Namespace.Name)) - var secrets []v1.Secret + var secrets []*v1.Secret for i := 1; i <= 5; i++ { secret, err := f.ClientSet.CoreV1().Secrets(f.Namespace.Name).Create(ctx, newSecret(fmt.Sprintf("secret-%d", i)), metav1.CreateOptions{}) framework.ExpectNoError(err) - secrets = append(secrets, *secret) + secrets = append(secrets, secret) } return secrets } // addWellKnownUnstructuredSecrets exists because secrets from addWellKnownSecrets // don't have type info and cannot be converted. -func addWellKnownUnstructuredSecrets(ctx context.Context, f *framework.Framework) []unstructured.Unstructured { - var secrets []unstructured.Unstructured +func addWellKnownUnstructuredSecrets(ctx context.Context, f *framework.Framework) []*unstructured.Unstructured { + var secrets []*unstructured.Unstructured for i := 1; i <= 5; i++ { unstructuredSecret, err := runtime.DefaultUnstructuredConverter.ToUnstructured(newSecret(fmt.Sprintf("secret-%d", i))) framework.ExpectNoError(err) secret, err := f.DynamicClient.Resource(v1.SchemeGroupVersion.WithResource("secrets")).Namespace(f.Namespace.Name).Create(ctx, &unstructured.Unstructured{Object: unstructuredSecret}, metav1.CreateOptions{}) framework.ExpectNoError(err) - secrets = append(secrets, *secret) + secrets = append(secrets, secret) } return secrets } -type byName []v1.Secret +type byName []metav1.Object func (a byName) Len() int { return len(a) } -func (a byName) Less(i, j int) bool { return a[i].Name < a[j].Name } +func (a byName) Less(i, j int) bool { return a[i].GetName() < a[j].GetName() } func (a byName) Swap(i, j int) { a[i], a[j] = a[j], a[i] } func newSecret(name string) *v1.Secret { @@ -369,3 +362,23 @@ func newSecret(name string) *v1.Secret { }, } } + +func toMetaObjectSlice[T any](s []T) ([]metav1.Object, error) { + result := make([]metav1.Object, len(s)) + for i, v := range s { + m, err := meta.Accessor(v) + if err != nil { + return nil, err + } + result[i] = m + } + return result, nil +} + +func toSecretPointerSlice[T any](items []T) []*T { + result := make([]*T, 0, len(items)) + for i := range items { + result = append(result, &items[i]) + } + return result +}