test/apimachinery/watchlist: prove dynamic client's List method not streaming

This commit is contained in:
Lukasz Szaszkiewicz
2025-06-11 09:45:38 +02:00
parent 3038f3530d
commit edadfee47d

View File

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