From 964df4ce952fcfe75035aeec5429b919d188c90d Mon Sep 17 00:00:00 2001 From: Jan Chaloupka Date: Wed, 4 Feb 2026 14:43:33 +0100 Subject: [PATCH] refactor(promClientController): split it into two prom client controllers --- pkg/descheduler/descheduler.go | 206 +++++++++++++++++----------- pkg/descheduler/descheduler_test.go | 183 +++++++++++++++--------- 2 files changed, 242 insertions(+), 147 deletions(-) diff --git a/pkg/descheduler/descheduler.go b/pkg/descheduler/descheduler.go index 6ebe35f46..d6fa187ca 100644 --- a/pkg/descheduler/descheduler.go +++ b/pkg/descheduler/descheduler.go @@ -22,6 +22,7 @@ import ( "math" "net/http" "strconv" + "sync" "time" promapi "github.com/prometheus/client_golang/api" @@ -78,17 +79,41 @@ type profileRunner struct { descheduleEPs, balanceEPs eprunner } +// inClusterPromClientController manages prometheus client using in-cluster SA token +type inClusterPromClientController struct { + mu sync.RWMutex + promClient promapi.Client + previousPrometheusClientTransport *http.Transport + currentPrometheusAuthToken string + prometheusConfig *api.Prometheus + createPrometheusClient createPrometheusClientFunc + inClusterConfig inClusterConfigFunc +} + +// secretBasedPromClientController manages prometheus client using Kubernetes secret +type secretBasedPromClientController struct { + mu sync.RWMutex + promClient promapi.Client + previousPrometheusClientTransport *http.Transport + queue workqueue.RateLimitingInterface + currentPrometheusAuthToken string + namespacedSecretsLister corev1listers.SecretNamespaceLister + prometheusConfig *api.Prometheus + createPrometheusClient createPrometheusClientFunc +} + type descheduler struct { - rs *options.DeschedulerServer - client clientset.Interface - kubeClientSandbox *kubeClientSandbox - getPodsAssignedToNode podutil.GetPodsAssignedToNodeFunc - sharedInformerFactory informers.SharedInformerFactory - deschedulerPolicy *api.DeschedulerPolicy - eventRecorder events.EventRecorder - podEvictor *evictions.PodEvictor - metricsCollector *metricscollector.MetricsCollector - promClientCtrl *promClientController + rs *options.DeschedulerServer + client clientset.Interface + kubeClientSandbox *kubeClientSandbox + getPodsAssignedToNode podutil.GetPodsAssignedToNodeFunc + sharedInformerFactory informers.SharedInformerFactory + deschedulerPolicy *api.DeschedulerPolicy + eventRecorder events.EventRecorder + podEvictor *evictions.PodEvictor + metricsCollector *metricscollector.MetricsCollector + inClusterPromClientCtrl *inClusterPromClientController + secretBasedPromClientCtrl *secretBasedPromClientController } type ( @@ -96,28 +121,35 @@ type ( inClusterConfigFunc func() (*rest.Config, error) ) -type promClientController struct { - promClient promapi.Client - previousPrometheusClientTransport *http.Transport - queue workqueue.RateLimitingInterface - currentPrometheusAuthToken string - namespacedSecretsLister corev1listers.SecretNamespaceLister - metricsProviders map[api.MetricsSource]*api.MetricsProvider - createPrometheusClient createPrometheusClientFunc - inClusterConfig inClusterConfigFunc -} - -func newPromClientController(prometheusClient promapi.Client, metricsProviders map[api.MetricsSource]*api.MetricsProvider) *promClientController { - return &promClientController{ +func newInClusterPromClientController(prometheusClient promapi.Client, prometheusConfig *api.Prometheus) *inClusterPromClientController { + return &inClusterPromClientController{ promClient: prometheusClient, - queue: workqueue.NewRateLimitingQueueWithConfig(workqueue.DefaultControllerRateLimiter(), workqueue.RateLimitingQueueConfig{Name: "descheduler"}), - metricsProviders: metricsProviders, + prometheusConfig: prometheusConfig, createPrometheusClient: client.CreatePrometheusClient, inClusterConfig: rest.InClusterConfig, } } -func (d *promClientController) prometheusClient() promapi.Client { +func newSecretBasedPromClientController(prometheusClient promapi.Client, prometheusConfig *api.Prometheus) *secretBasedPromClientController { + ctrl := &secretBasedPromClientController{ + promClient: prometheusClient, + queue: workqueue.NewRateLimitingQueueWithConfig(workqueue.DefaultControllerRateLimiter(), workqueue.RateLimitingQueueConfig{Name: "descheduler"}), + prometheusConfig: prometheusConfig, + createPrometheusClient: client.CreatePrometheusClient, + } + + return ctrl +} + +func (d *inClusterPromClientController) prometheusClient() promapi.Client { + d.mu.RLock() + defer d.mu.RUnlock() + return d.promClient +} + +func (d *secretBasedPromClientController) prometheusClient() promapi.Client { + d.mu.RLock() + defer d.mu.RUnlock() return d.promClient } @@ -140,31 +172,47 @@ func addNodeSelectorIndexer(sharedInformerFactory informers.SharedInformerFactor func metricsProviderListToMap(providersList []api.MetricsProvider) map[api.MetricsSource]*api.MetricsProvider { providersMap := make(map[api.MetricsSource]*api.MetricsProvider) for _, provider := range providersList { + provider := provider providersMap[provider.Source] = &provider } return providersMap } // setupPrometheusProvider sets up the prometheus provider on the descheduler if configured -func setupPrometheusProvider(d *descheduler, namespacedSharedInformerFactory informers.SharedInformerFactory) error { - if d.promClientCtrl == nil { +func setupPrometheusProvider(ctrl *secretBasedPromClientController, namespacedSharedInformerFactory informers.SharedInformerFactory) error { + if ctrl == nil { return nil } - prometheusProvider := d.promClientCtrl.metricsProviders[api.PrometheusMetrics] - if prometheusProvider != nil && prometheusProvider.Prometheus != nil && prometheusProvider.Prometheus.AuthToken != nil { - authTokenSecret := prometheusProvider.Prometheus.AuthToken.SecretReference + prometheusConfig := ctrl.prometheusConfig + if prometheusConfig == nil { + return fmt.Errorf("prometheus configuration is missing") + } + if prometheusConfig.AuthToken != nil { + authTokenSecret := prometheusConfig.AuthToken.SecretReference if authTokenSecret == nil || authTokenSecret.Namespace == "" { return fmt.Errorf("prometheus metrics source configuration is missing authentication token secret") } if namespacedSharedInformerFactory == nil { return fmt.Errorf("namespacedSharedInformerFactory not configured") } - namespacedSharedInformerFactory.Core().V1().Secrets().Informer().AddEventHandler(d.promClientCtrl.eventHandler()) - d.promClientCtrl.namespacedSecretsLister = namespacedSharedInformerFactory.Core().V1().Secrets().Lister().Secrets(authTokenSecret.Namespace) + namespacedSharedInformerFactory.Core().V1().Secrets().Informer().AddEventHandler(ctrl.eventHandler()) + ctrl.namespacedSecretsLister = namespacedSharedInformerFactory.Core().V1().Secrets().Lister().Secrets(authTokenSecret.Namespace) } return nil } +func getPrometheusConfig(providersList []api.MetricsProvider) *api.Prometheus { + prometheusProvider := metricsProviderListToMap(providersList)[api.PrometheusMetrics] + if prometheusProvider == nil { + return nil + } + return prometheusProvider.Prometheus +} + +func configureSecretPromClientReconciler(prometheusConfig *api.Prometheus) bool { + return prometheusConfig.URL != "" && prometheusConfig.AuthToken != nil +} + func newDescheduler(ctx context.Context, rs *options.DeschedulerServer, deschedulerPolicy *api.DeschedulerPolicy, evictionPolicyGroupVersion string, eventRecorder events.EventRecorder, client clientset.Interface, sharedInformerFactory informers.SharedInformerFactory, kubeClientSandbox *kubeClientSandbox) (*descheduler, error) { podInformer := sharedInformerFactory.Core().V1().Pods().Informer() // Temporarily register the PVC because it is used by the DefaultEvictor plugin during @@ -205,7 +253,17 @@ func newDescheduler(ctx context.Context, rs *options.DeschedulerServer, deschedu deschedulerPolicy: deschedulerPolicy, eventRecorder: eventRecorder, podEvictor: podEvictor, - promClientCtrl: newPromClientController(rs.PrometheusClient, metricsProviderListToMap(deschedulerPolicy.MetricsProviders)), + } + + prometheusConfig := getPrometheusConfig(deschedulerPolicy.MetricsProviders) + if prometheusConfig != nil && prometheusConfig.URL != "" { + if configureSecretPromClientReconciler(prometheusConfig) { + // Secret-based mode + desch.secretBasedPromClientCtrl = newSecretBasedPromClientController(rs.PrometheusClient, prometheusConfig) + } else { + // In-cluster mode + desch.inClusterPromClientCtrl = newInClusterPromClientController(rs.PrometheusClient, prometheusConfig) + } } nodeSelector, err := nodeSelectorFromPolicy(deschedulerPolicy) @@ -224,13 +282,16 @@ func newDescheduler(ctx context.Context, rs *options.DeschedulerServer, deschedu return desch, nil } -func (d *promClientController) reconcileInClusterSAToken() error { +func (d *inClusterPromClientController) reconcileInClusterSAToken() error { + d.mu.Lock() + defer d.mu.Unlock() + // Read the sa token and assume it has the sufficient permissions to authenticate cfg, err := d.inClusterConfig() if err == nil { if d.currentPrometheusAuthToken != cfg.BearerToken { klog.V(2).Infof("Creating Prometheus client (with SA token)") - prometheusClient, transport, err := d.createPrometheusClient(d.metricsProviders[api.PrometheusMetrics].Prometheus.URL, cfg.BearerToken) + prometheusClient, transport, err := d.createPrometheusClient(d.prometheusConfig.URL, cfg.BearerToken) if err != nil { return fmt.Errorf("unable to create a prometheus client: %v", err) } @@ -249,7 +310,7 @@ func (d *promClientController) reconcileInClusterSAToken() error { return fmt.Errorf("unexpected error when reading in cluster config: %v", err) } -func (d *promClientController) runAuthenticationSecretReconciler(ctx context.Context) { +func (d *secretBasedPromClientController) runAuthenticationSecretReconciler(ctx context.Context) { defer utilruntime.HandleCrash() defer d.queue.ShutDown() @@ -261,12 +322,12 @@ func (d *promClientController) runAuthenticationSecretReconciler(ctx context.Con <-ctx.Done() } -func (d *promClientController) runAuthenticationSecretReconcilerWorker(ctx context.Context) { +func (d *secretBasedPromClientController) runAuthenticationSecretReconcilerWorker(ctx context.Context) { for d.processNextWorkItem(ctx) { } } -func (d *promClientController) processNextWorkItem(ctx context.Context) bool { +func (d *secretBasedPromClientController) processNextWorkItem(ctx context.Context) bool { dsKey, quit := d.queue.Get() if quit { return false @@ -285,8 +346,11 @@ func (d *promClientController) processNextWorkItem(ctx context.Context) bool { return true } -func (d *promClientController) sync() error { - prometheusConfig := d.metricsProviders[api.PrometheusMetrics].Prometheus +func (d *secretBasedPromClientController) sync() error { + d.mu.Lock() + defer d.mu.Unlock() + + prometheusConfig := d.prometheusConfig if prometheusConfig == nil || prometheusConfig.AuthToken == nil || prometheusConfig.AuthToken.SecretReference == nil { return fmt.Errorf("prometheus metrics source configuration is missing authentication token secret") } @@ -327,7 +391,7 @@ func (d *promClientController) sync() error { return nil } -func (d *promClientController) eventHandler() cache.ResourceEventHandler { +func (d *secretBasedPromClientController) eventHandler() cache.ResourceEventHandler { return cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { d.queue.Add(workQueueKey) }, UpdateFunc: func(old, new interface{}) { d.queue.Add(workQueueKey) }, @@ -397,8 +461,14 @@ func (d *descheduler) runProfiles(ctx context.Context) { var profileRunners []profileRunner for idx, profile := range d.deschedulerPolicy.Profiles { var promClient promapi.Client - if d.promClientCtrl != nil { - promClient = d.promClientCtrl.prometheusClient() + if d.inClusterPromClientCtrl != nil && d.secretBasedPromClientCtrl != nil { + klog.Error(fmt.Errorf("At most one of inClusterPromClientCtrl and secretBasedPromClientCtrl can be set, got both")) + continue + } + if d.inClusterPromClientCtrl != nil { + promClient = d.inClusterPromClientCtrl.prometheusClient() + } else if d.secretBasedPromClientCtrl != nil { + promClient = d.secretBasedPromClientCtrl.prometheusClient() } currProfile, err := frameworkprofile.NewProfile( ctx, @@ -539,14 +609,6 @@ func validateVersionCompatibility(discovery discovery.DiscoveryInterface, desche return nil } -type tokenReconciliation int - -const ( - noReconciliation tokenReconciliation = iota - inClusterReconciliation - secretReconciliation -) - type runFncType func(context.Context) error func bootstrapDescheduler( @@ -554,7 +616,6 @@ func bootstrapDescheduler( rs *options.DeschedulerServer, deschedulerPolicy *api.DeschedulerPolicy, evictionPolicyGroupVersion string, - metricProviderTokenReconciliation tokenReconciliation, sharedInformerFactory, namespacedSharedInformerFactory informers.SharedInformerFactory, eventRecorder events.EventRecorder, ) (*descheduler, runFncType, error) { @@ -565,7 +626,7 @@ func bootstrapDescheduler( } // Setup Prometheus provider - if err := setupPrometheusProvider(descheduler, namespacedSharedInformerFactory); err != nil { + if err := setupPrometheusProvider(descheduler.secretBasedPromClientCtrl, namespacedSharedInformerFactory); err != nil { return nil, nil, fmt.Errorf("failed to setup Prometheus provider: %v", err) } @@ -588,7 +649,7 @@ func bootstrapDescheduler( } // Setup Prometheus provider (with the real shared informer factory as the secret is only read) - if err := setupPrometheusProvider(descheduler, namespacedSharedInformerFactory); err != nil { + if err := setupPrometheusProvider(descheduler.secretBasedPromClientCtrl, namespacedSharedInformerFactory); err != nil { return nil, nil, fmt.Errorf("failed to setup Prometheus provider for the dry run descheduler: %v", err) } } @@ -604,12 +665,12 @@ func bootstrapDescheduler( descheduler.kubeClientSandbox.fakeSharedInformerFactory().WaitForCacheSync(ctx.Done()) } sharedInformerFactory.Start(ctx.Done()) - if metricProviderTokenReconciliation == secretReconciliation { + if descheduler.secretBasedPromClientCtrl != nil { namespacedSharedInformerFactory.Start(ctx.Done()) } sharedInformerFactory.WaitForCacheSync(ctx.Done()) - if metricProviderTokenReconciliation == secretReconciliation { + if descheduler.secretBasedPromClientCtrl != nil { namespacedSharedInformerFactory.WaitForCacheSync(ctx.Done()) } @@ -629,8 +690,8 @@ func bootstrapDescheduler( } } - if metricProviderTokenReconciliation == secretReconciliation { - go descheduler.promClientCtrl.runAuthenticationSecretReconciler(ctx) + if descheduler.secretBasedPromClientCtrl != nil { + go descheduler.secretBasedPromClientCtrl.runAuthenticationSecretReconciler(ctx) } return nil @@ -641,9 +702,9 @@ func bootstrapDescheduler( } runFnc := func(ctx context.Context) error { - if metricProviderTokenReconciliation == inClusterReconciliation { + if descheduler.inClusterPromClientCtrl != nil { // Read the sa token and assume it has the sufficient permissions to authenticate - if err := descheduler.promClientCtrl.reconcileInClusterSAToken(); err != nil { + if err := descheduler.inClusterPromClientCtrl.reconcileInClusterSAToken(); err != nil { return fmt.Errorf("unable to reconcile an in cluster SA token: %v", err) } } @@ -659,19 +720,6 @@ func bootstrapDescheduler( return descheduler, runFnc, nil } -func prometheusProviderToTokenReconciliation(prometheusProvider *api.MetricsProvider) tokenReconciliation { - if prometheusProvider != nil && prometheusProvider.Prometheus != nil && prometheusProvider.Prometheus.URL != "" { - if prometheusProvider.Prometheus.AuthToken != nil { - // Will get reconciled - return secretReconciliation - } else { - // Use the sa token and assume it has the sufficient permissions to authenticate - return inClusterReconciliation - } - } - return noReconciliation -} - func RunDeschedulerStrategies(ctx context.Context, rs *options.DeschedulerServer, deschedulerPolicy *api.DeschedulerPolicy, evictionPolicyGroupVersion string) error { ctx, cancel := context.WithCancel(ctx) defer cancel() @@ -692,14 +740,12 @@ func RunDeschedulerStrategies(ctx context.Context, rs *options.DeschedulerServer defer eventBroadcaster.Shutdown() var namespacedSharedInformerFactory informers.SharedInformerFactory - - prometheusProvider := metricsProviderListToMap(deschedulerPolicy.MetricsProviders)[api.PrometheusMetrics] - metricProviderTokenReconciliation := prometheusProviderToTokenReconciliation(prometheusProvider) - if metricProviderTokenReconciliation == secretReconciliation { - namespacedSharedInformerFactory = informers.NewSharedInformerFactoryWithOptions(rs.Client, 0, informers.WithTransform(trimManagedFields), informers.WithNamespace(prometheusProvider.Prometheus.AuthToken.SecretReference.Namespace)) + prometheusConfig := getPrometheusConfig(deschedulerPolicy.MetricsProviders) + if prometheusConfig != nil && configureSecretPromClientReconciler(prometheusConfig) { + namespacedSharedInformerFactory = informers.NewSharedInformerFactoryWithOptions(rs.Client, 0, informers.WithTransform(trimManagedFields), informers.WithNamespace(prometheusConfig.AuthToken.SecretReference.Namespace)) } - _, runLoop, err := bootstrapDescheduler(ctx, rs, deschedulerPolicy, evictionPolicyGroupVersion, metricProviderTokenReconciliation, sharedInformerFactory, namespacedSharedInformerFactory, eventRecorder) + _, runLoop, err := bootstrapDescheduler(ctx, rs, deschedulerPolicy, evictionPolicyGroupVersion, sharedInformerFactory, namespacedSharedInformerFactory, eventRecorder) if err != nil { span.AddEvent("Failed to bootstrap a descheduler", trace.WithAttributes(attribute.String("err", err.Error()))) return err diff --git a/pkg/descheduler/descheduler_test.go b/pkg/descheduler/descheduler_test.go index 372eaf2bd..492351e05 100644 --- a/pkg/descheduler/descheduler_test.go +++ b/pkg/descheduler/descheduler_test.go @@ -26,6 +26,7 @@ import ( fakediscovery "k8s.io/client-go/discovery/fake" "k8s.io/client-go/informers" fakeclientset "k8s.io/client-go/kubernetes/fake" + corev1listers "k8s.io/client-go/listers/core/v1" "k8s.io/client-go/rest" core "k8s.io/client-go/testing" "k8s.io/client-go/tools/cache" @@ -220,15 +221,13 @@ func initDescheduler(t *testing.T, ctx context.Context, featureGates featuregate eventBroadcaster, eventRecorder := utils.GetRecorderAndBroadcaster(ctx, client) var namespacedSharedInformerFactory informers.SharedInformerFactory - - prometheusProvider := metricsProviderListToMap(internalDeschedulerPolicy.MetricsProviders)[api.PrometheusMetrics] - metricProviderTokenReconciliation := prometheusProviderToTokenReconciliation(prometheusProvider) - if metricProviderTokenReconciliation == secretReconciliation { - namespacedSharedInformerFactory = informers.NewSharedInformerFactoryWithOptions(rs.Client, 0, informers.WithTransform(trimManagedFields), informers.WithNamespace(prometheusProvider.Prometheus.AuthToken.SecretReference.Namespace)) + prometheusConfig := getPrometheusConfig(internalDeschedulerPolicy.MetricsProviders) + if prometheusConfig != nil && configureSecretPromClientReconciler(prometheusConfig) { + namespacedSharedInformerFactory = informers.NewSharedInformerFactoryWithOptions(rs.Client, 0, informers.WithTransform(trimManagedFields), informers.WithNamespace(prometheusConfig.AuthToken.SecretReference.Namespace)) } // Always create descheduler with real client/factory first to register all informers - descheduler, runFnc, err := bootstrapDescheduler(ctx, rs, internalDeschedulerPolicy, "v1", metricProviderTokenReconciliation, sharedInformerFactory, namespacedSharedInformerFactory, eventRecorder) + descheduler, runFnc, err := bootstrapDescheduler(ctx, rs, internalDeschedulerPolicy, "v1", sharedInformerFactory, namespacedSharedInformerFactory, eventRecorder) if err != nil { eventBroadcaster.Shutdown() t.Fatalf("Failed to bootstrap a descheduler: %v", err) @@ -237,6 +236,10 @@ func initDescheduler(t *testing.T, ctx context.Context, featureGates featuregate if dryRun { if err := wait.PollUntilContextTimeout(ctx, 100*time.Millisecond, 5*time.Second, true, func(ctx context.Context) (bool, error) { for _, obj := range objects { + // Only check for nodes - secrets are handled by namespacedSharedInformerFactory + if _, ok := obj.(*v1.Node); !ok { + continue + } exists, err := descheduler.kubeClientSandbox.hasRuntimeObjectInIndexer(obj) if err != nil { return false, err @@ -261,6 +264,27 @@ func initDescheduler(t *testing.T, ctx context.Context, featureGates featuregate return rs, descheduler, runFnc, client } +func waitForSecretsPropagation(ctx context.Context, secretsLister corev1listers.SecretNamespaceLister, objects []runtime.Object) error { + return wait.PollUntilContextTimeout(ctx, 100*time.Millisecond, 5*time.Second, true, func(ctx context.Context) (bool, error) { + for _, obj := range objects { + secret, ok := obj.(*v1.Secret) + if !ok { + continue + } + _, err := secretsLister.Get(secret.Name) + if err != nil { + if apierrors.IsNotFound(err) { + klog.Infof("Secret %s/%s has not propagated to the indexer", secret.Namespace, secret.Name) + return false, nil + } + return false, err + } + klog.Infof("Secret %s/%s has propagated to the indexer", secret.Namespace, secret.Name) + } + return true, nil + }) +} + func TestTaintsUpdated(t *testing.T) { initPluginRegistry() @@ -1552,8 +1576,9 @@ func TestPluginPrometheusClientAccess_Secret(t *testing.T) { t.Logf("Cycle %d: %s", i+1, cycle.name) // Set the descheduler's Prometheus client - t.Logf("Setting descheduler.promClientCtrl.promClient from %v to %v", descheduler.promClientCtrl.promClient, cycle.client) - descheduler.promClientCtrl.createPrometheusClient = func(url, token string) (promapi.Client, *http.Transport, error) { + descheduler.secretBasedPromClientCtrl.mu.Lock() + t.Logf("Setting descheduler.secretBasedPromClientCtrl.promClient from %v to %v", descheduler.secretBasedPromClientCtrl.promClient, cycle.client) + descheduler.secretBasedPromClientCtrl.createPrometheusClient = func(url, token string) (promapi.Client, *http.Transport, error) { if token != cycle.token { t.Fatalf("Expected token to be %q, got %q", cycle.token, token) } @@ -1562,6 +1587,7 @@ func TestPluginPrometheusClientAccess_Secret(t *testing.T) { } return cycle.client, &http.Transport{}, nil } + descheduler.secretBasedPromClientCtrl.mu.Unlock() if err := cycle.operation(); err != nil { t.Fatalf("operation failed: %v", err) @@ -1569,7 +1595,7 @@ func TestPluginPrometheusClientAccess_Secret(t *testing.T) { if !cycle.skipWaiting { err := wait.PollUntilContextTimeout(ctx, 50*time.Millisecond, 200*time.Millisecond, true, func(ctx context.Context) (bool, error) { - currentPromClient := descheduler.promClientCtrl.prometheusClient() + currentPromClient := descheduler.secretBasedPromClientCtrl.prometheusClient() if currentPromClient != cycle.client { t.Logf("Waiting for prometheus client to be set to %v, got %v instead, waiting", cycle.client, currentPromClient) return false, nil @@ -1590,7 +1616,7 @@ func TestPluginPrometheusClientAccess_Secret(t *testing.T) { t.Fatalf("Unexpected error during running a descheduling cycle: %v", err) } - t.Logf("After cycle %d: prometheusClientFromReactor=%v, descheduler.promClientCtrl.promClient=%v", i+1, prometheusClientFromReactor, descheduler.promClientCtrl.promClient) + t.Logf("After cycle %d: prometheusClientFromReactor=%v, descheduler.secretBasedPromClientCtrl.promClient=%v", i+1, prometheusClientFromReactor, descheduler.secretBasedPromClientCtrl.prometheusClient()) if !newInvoked { t.Fatalf("Expected plugin new to be invoked during cycle %d", i+1) @@ -1600,7 +1626,7 @@ func TestPluginPrometheusClientAccess_Secret(t *testing.T) { t.Fatalf("Expected deschedule reactor to be invoked during cycle %d", i+1) } - verifyAllPrometheusClientsEqual(t, cycle.client, prometheusClientFromReactor, prometheusClientFromPluginNewHandle, descheduler.promClientCtrl.promClient) + verifyAllPrometheusClientsEqual(t, cycle.client, prometheusClientFromReactor, prometheusClientFromPluginNewHandle, descheduler.secretBasedPromClientCtrl.prometheusClient()) } }) } @@ -1722,11 +1748,12 @@ func TestPluginPrometheusClientAccess_InCluster(t *testing.T) { t.Logf("Cycle %d: %s", i+1, cycle.name) // Set the descheduler's Prometheus client - t.Logf("Setting descheduler.promClientCtrl.promClient from %v to %v", descheduler.promClientCtrl.promClient, cycle.client) - descheduler.promClientCtrl.inClusterConfig = func() (*rest.Config, error) { + descheduler.inClusterPromClientCtrl.mu.Lock() + t.Logf("Setting descheduler.inClusterPromClientCtrl.promClient from %v to %v", descheduler.inClusterPromClientCtrl.promClient, cycle.client) + descheduler.inClusterPromClientCtrl.inClusterConfig = func() (*rest.Config, error) { return &rest.Config{BearerToken: cycle.token}, nil } - descheduler.promClientCtrl.createPrometheusClient = func(url, token string) (promapi.Client, *http.Transport, error) { + descheduler.inClusterPromClientCtrl.createPrometheusClient = func(url, token string) (promapi.Client, *http.Transport, error) { if token != cycle.token { t.Errorf("Expected token to be %q, got %q", cycle.token, token) } @@ -1735,6 +1762,7 @@ func TestPluginPrometheusClientAccess_InCluster(t *testing.T) { } return cycle.client, &http.Transport{}, nil } + descheduler.inClusterPromClientCtrl.mu.Unlock() newInvoked = false reactorInvoked = false @@ -1745,7 +1773,7 @@ func TestPluginPrometheusClientAccess_InCluster(t *testing.T) { t.Fatalf("Unexpected error during running a descheduling cycle: %v", err) } - t.Logf("After cycle %d: prometheusClientFromReactor=%v, descheduler.promClientCtrl.promClient=%v", i+1, prometheusClientFromReactor, descheduler.promClientCtrl.promClient) + t.Logf("After cycle %d: prometheusClientFromReactor=%v, descheduler.inClusterPromClientCtrl.promClient=%v", i+1, prometheusClientFromReactor, descheduler.inClusterPromClientCtrl.prometheusClient()) if !newInvoked { t.Fatalf("Expected plugin new to be invoked during cycle %d", i+1) @@ -1755,7 +1783,7 @@ func TestPluginPrometheusClientAccess_InCluster(t *testing.T) { t.Fatalf("Expected deschedule reactor to be invoked during cycle %d", i+1) } - verifyAllPrometheusClientsEqual(t, cycle.client, prometheusClientFromReactor, prometheusClientFromPluginNewHandle, descheduler.promClientCtrl.promClient) + verifyAllPrometheusClientsEqual(t, cycle.client, prometheusClientFromReactor, prometheusClientFromPluginNewHandle, descheduler.inClusterPromClientCtrl.prometheusClient()) } }) } @@ -1769,6 +1797,10 @@ func withToken(token string) func(*v1.Secret) { func newPrometheusAuthSecret(apply func(*v1.Secret)) *v1.Secret { secret := &v1.Secret{ + TypeMeta: metav1.TypeMeta{ + APIVersion: "v1", + Kind: "Secret", + }, ObjectMeta: metav1.ObjectMeta{ Name: "prom-token", Namespace: "kube-system", @@ -1787,11 +1819,11 @@ type promClientControllerTestSetup struct { fakeClient *fakeclientset.Clientset namespacedInformerFactory informers.SharedInformerFactory metricsProviders map[api.MetricsSource]*api.MetricsProvider - ctrl *promClientController + ctrl *secretBasedPromClientController namespace string } -func setupPromClientControllerTest(ctx context.Context, objects []runtime.Object, prometheusConfig *api.Prometheus) *promClientControllerTestSetup { +func setupPromClientControllerTest(ctx context.Context, t *testing.T, objects []runtime.Object, prometheusConfig *api.Prometheus, setNamespacedSharedInformerFactory bool) *promClientControllerTestSetup { fakeClient := fakeclientset.NewSimpleClientset(objects...) namespace := "default" @@ -1807,15 +1839,26 @@ func setupPromClientControllerTest(ctx context.Context, objects []runtime.Object Prometheus: prometheusConfig, }}) - ctrl := newPromClientController(nil, metricsProviders) - if prometheusConfig != nil && prometheusConfig.AuthToken != nil && prometheusConfig.AuthToken.SecretReference != nil { - ctrl.namespacedSecretsLister = namespacedInformerFactory.Core().V1().Secrets().Lister().Secrets(namespace) - namespacedInformerFactory.Core().V1().Secrets().Informer().AddEventHandler(ctrl.eventHandler()) + ctrl := newSecretBasedPromClientController(nil, prometheusConfig) + if setNamespacedSharedInformerFactory { + if err := setupPrometheusProvider(ctrl, namespacedInformerFactory); err != nil { + t.Fatal(err) + } } namespacedInformerFactory.Start(ctx.Done()) namespacedInformerFactory.WaitForCacheSync(ctx.Done()) + // Wait for secrets to propagate to the informer + if len(objects) > 0 && setNamespacedSharedInformerFactory { + ctx, cancel := context.WithTimeout(ctx, 5*time.Second) + defer cancel() + secretsLister := namespacedInformerFactory.Core().V1().Secrets().Lister().Secrets(namespace) + if err := waitForSecretsPropagation(ctx, secretsLister, objects); err != nil { + t.Fatalf("secrets did not propagate to the indexer: %v", err) + } + } + return &promClientControllerTestSetup{ fakeClient: fakeClient, namespacedInformerFactory: namespacedInformerFactory, @@ -1839,10 +1882,11 @@ func newPrometheusConfig() *api.Prometheus { func TestPromClientControllerSync_InvalidConfig(t *testing.T) { testCases := []struct { - name string - objects []runtime.Object - prometheusConfig *api.Prometheus - expectedErr error + name string + objects []runtime.Object + prometheusConfig *api.Prometheus + expectedErr error + setupPromClientControllerTest bool }{ { name: "empty prometheus config", @@ -1875,18 +1919,20 @@ func TestPromClientControllerSync_InvalidConfig(t *testing.T) { expectedErr: fmt.Errorf("prometheus metrics source configuration is missing authentication token secret"), }, { - name: "secret exists but empty token", - objects: []runtime.Object{newPrometheusAuthSecret(withToken(""))}, - prometheusConfig: newPrometheusConfig(), - expectedErr: fmt.Errorf("prometheus authentication token secret missing \"prometheusAuthToken\" data or empty"), + name: "secret exists but empty token", + objects: []runtime.Object{newPrometheusAuthSecret(withToken(""))}, + prometheusConfig: newPrometheusConfig(), + setupPromClientControllerTest: true, + expectedErr: fmt.Errorf("prometheus authentication token secret missing \"prometheusAuthToken\" data or empty"), }, { name: "secret exists but missing token key", objects: []runtime.Object{newPrometheusAuthSecret(func(s *v1.Secret) { s.Data = map[string][]byte{} })}, - prometheusConfig: newPrometheusConfig(), - expectedErr: fmt.Errorf("prometheus authentication token secret missing \"prometheusAuthToken\" data or empty"), + prometheusConfig: newPrometheusConfig(), + setupPromClientControllerTest: true, + expectedErr: fmt.Errorf("prometheus authentication token secret missing \"prometheusAuthToken\" data or empty"), }, } @@ -1894,7 +1940,7 @@ func TestPromClientControllerSync_InvalidConfig(t *testing.T) { t.Run(tc.name, func(t *testing.T) { ctx, cancel := context.WithCancel(context.TODO()) defer cancel() - setup := setupPromClientControllerTest(ctx, tc.objects, tc.prometheusConfig) + setup := setupPromClientControllerTest(ctx, t, tc.objects, tc.prometheusConfig, tc.setupPromClientControllerTest) // Call sync err := setup.ctrl.sync() @@ -1970,18 +2016,18 @@ func TestPromClientControllerSync_ClientCreation(t *testing.T) { t.Run(tc.name, func(t *testing.T) { for _, setupMode := range []struct { name string - setupFn func(context.Context, *testing.T, []runtime.Object) *promClientController + setupFn func(context.Context, *testing.T, []runtime.Object) *secretBasedPromClientController }{ { name: "running with prom reconciler directly", - setupFn: func(ctx context.Context, t *testing.T, objects []runtime.Object) *promClientController { - setup := setupPromClientControllerTest(ctx, objects, newPrometheusConfig()) + setupFn: func(ctx context.Context, t *testing.T, objects []runtime.Object) *secretBasedPromClientController { + setup := setupPromClientControllerTest(ctx, t, objects, newPrometheusConfig(), true) return setup.ctrl }, }, { name: "running with full descheduler", - setupFn: func(ctx context.Context, t *testing.T, objects []runtime.Object) *promClientController { + setupFn: func(ctx context.Context, t *testing.T, objects []runtime.Object) *secretBasedPromClientController { deschedulerPolicy := &api.DeschedulerPolicy{ MetricsProviders: []api.MetricsProvider{ { @@ -1991,7 +2037,7 @@ func TestPromClientControllerSync_ClientCreation(t *testing.T) { }, } _, descheduler, _, _ := initDescheduler(t, ctx, initFeatureGates(), deschedulerPolicy, nil, false, objects...) - return descheduler.promClientCtrl + return descheduler.secretBasedPromClientCtrl }, }, } { @@ -2054,7 +2100,7 @@ func TestPromClientControllerSync_ClientCreation(t *testing.T) { } // Verify promClient cleared when secret not found - if tc.expectPreviousTransportCleared && ctrl.promClient != nil { + if tc.expectPreviousTransportCleared && ctrl.prometheusClient() != nil { t.Errorf("Expected promClient to be cleared but it wasn't") } @@ -2135,12 +2181,12 @@ func TestPromClientControllerSync_EventHandler(t *testing.T) { for _, setupMode := range []struct { name string - init func(t *testing.T, ctx context.Context) (ctrl *promClientController, fakeClient *fakeclientset.Clientset) + init func(t *testing.T, ctx context.Context) (ctrl *secretBasedPromClientController, fakeClient *fakeclientset.Clientset) }{ { name: "running with prom reconciler directly", - init: func(t *testing.T, ctx context.Context) (ctrl *promClientController, fakeClient *fakeclientset.Clientset) { - setup := setupPromClientControllerTest(ctx, nil, newPrometheusConfig()) + init: func(t *testing.T, ctx context.Context) (ctrl *secretBasedPromClientController, fakeClient *fakeclientset.Clientset) { + setup := setupPromClientControllerTest(ctx, t, nil, newPrometheusConfig(), true) // Start the reconciler to process queue items go setup.ctrl.runAuthenticationSecretReconciler(ctx) @@ -2150,7 +2196,7 @@ func TestPromClientControllerSync_EventHandler(t *testing.T) { }, { name: "running with full descheduler", - init: func(t *testing.T, ctx context.Context) (ctrl *promClientController, fakeClient *fakeclientset.Clientset) { + init: func(t *testing.T, ctx context.Context) (ctrl *secretBasedPromClientController, fakeClient *fakeclientset.Clientset) { deschedulerPolicy := &api.DeschedulerPolicy{ MetricsProviders: []api.MetricsProvider{ { @@ -2163,7 +2209,7 @@ func TestPromClientControllerSync_EventHandler(t *testing.T) { _, descheduler, _, client := initDescheduler(t, ctx, initFeatureGates(), deschedulerPolicy, nil, false) // The reconciler is already started by initDescheduler via bootstrapDescheduler - return descheduler.promClientCtrl, client + return descheduler.secretBasedPromClientCtrl, client }, }, } { @@ -2193,13 +2239,19 @@ func TestPromClientControllerSync_EventHandler(t *testing.T) { if tc.processItem { // Wait for event to be processed by the reconciler err := wait.PollUntilContextTimeout(ctx, 10*time.Millisecond, 2*time.Second, true, func(ctx context.Context) (bool, error) { - // Check if all expected conditions are met + // Check if all expected conditions are met (with mutex protection) + ctrl.mu.RLock() + promClient := ctrl.promClient + currentToken := ctrl.currentPrometheusAuthToken + previousTransport := ctrl.previousPrometheusClientTransport + ctrl.mu.RUnlock() + if tc.expectedPromClientSet { - if ctrl.promClient == nil { + if promClient == nil { return false, nil } } else { - if ctrl.promClient != nil { + if promClient != nil { return false, nil } } @@ -2211,12 +2263,12 @@ func TestPromClientControllerSync_EventHandler(t *testing.T) { return false, nil } - if ctrl.currentPrometheusAuthToken != tc.expectedCurrentToken { + if currentToken != tc.expectedCurrentToken { return false, nil } if tc.expectedPreviousTransportCleared { - if ctrl.previousPrometheusClientTransport != nil { + if previousTransport != nil { return false, nil } } @@ -2234,12 +2286,13 @@ func TestPromClientControllerSync_EventHandler(t *testing.T) { // Validate post-conditions if tc.expectedPromClientSet { - if ctrl.promClient == nil { + if ctrl.prometheusClient() == nil { t.Error("Expected prometheus client to be set, but it was nil") } } else { - if ctrl.promClient != nil { - t.Errorf("Expected prometheus client to be nil, but got: %v", ctrl.promClient) + promClient := ctrl.prometheusClient() + if promClient != nil { + t.Errorf("Expected prometheus client to be nil, but got: %v", promClient) } } @@ -2374,20 +2427,15 @@ func TestReconcileInClusterSAToken(t *testing.T) { t.Run(tc.name, func(t *testing.T) { for _, setupMode := range []struct { name string - init func(t *testing.T) *promClientController + init func(t *testing.T) *inClusterPromClientController }{ { name: "running with prom reconciler directly", - init: func(t *testing.T) *promClientController { - return &promClientController{ + init: func(t *testing.T) *inClusterPromClientController { + return &inClusterPromClientController{ currentPrometheusAuthToken: tc.currentAuthToken, - metricsProviders: map[api.MetricsSource]*api.MetricsProvider{ - api.PrometheusMetrics: { - Source: api.PrometheusMetrics, - Prometheus: &api.Prometheus{ - URL: prometheusURL, - }, - }, + prometheusConfig: &api.Prometheus{ + URL: prometheusURL, }, inClusterConfig: tc.inClusterConfigFunc, } @@ -2395,7 +2443,7 @@ func TestReconcileInClusterSAToken(t *testing.T) { }, { name: "running with full descheduler", - init: func(t *testing.T) *promClientController { + init: func(t *testing.T) *inClusterPromClientController { ctx := context.Background() prometheusConfig := &api.Prometheus{ URL: prometheusURL, @@ -2410,10 +2458,11 @@ func TestReconcileInClusterSAToken(t *testing.T) { } _, descheduler, _, _ := initDescheduler(t, ctx, initFeatureGates(), deschedulerPolicy, nil, false) - // Override the fields needed for this test - descheduler.promClientCtrl.currentPrometheusAuthToken = tc.currentAuthToken - descheduler.promClientCtrl.inClusterConfig = tc.inClusterConfigFunc - return descheduler.promClientCtrl + // Override the fields needed for this test (no need for a mutex + // since reconcileInClusterSAToken gets to run later) + descheduler.inClusterPromClientCtrl.currentPrometheusAuthToken = tc.currentAuthToken + descheduler.inClusterPromClientCtrl.inClusterConfig = tc.inClusterConfigFunc + return descheduler.inClusterPromClientCtrl }, }, } {