diff --git a/pkg/descheduler/descheduler.go b/pkg/descheduler/descheduler.go index 277365b44..c8c1c3f95 100644 --- a/pkg/descheduler/descheduler.go +++ b/pkg/descheduler/descheduler.go @@ -20,9 +20,7 @@ import ( "context" "fmt" "math" - "net/http" "strconv" - "sync" "time" promapi "github.com/prometheus/client_golang/api" @@ -31,21 +29,15 @@ import ( "go.opentelemetry.io/otel/trace" v1 "k8s.io/api/core/v1" - apierrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/meta" "k8s.io/apimachinery/pkg/labels" - utilruntime "k8s.io/apimachinery/pkg/util/runtime" utilversion "k8s.io/apimachinery/pkg/util/version" "k8s.io/apimachinery/pkg/util/wait" "k8s.io/client-go/discovery" "k8s.io/client-go/informers" clientset "k8s.io/client-go/kubernetes" fakeclientset "k8s.io/client-go/kubernetes/fake" - corev1listers "k8s.io/client-go/listers/core/v1" - "k8s.io/client-go/rest" - "k8s.io/client-go/tools/cache" "k8s.io/client-go/tools/events" - "k8s.io/client-go/util/workqueue" componentbaseconfig "k8s.io/component-base/config" "k8s.io/klog/v2" @@ -67,9 +59,7 @@ import ( ) const ( - prometheusAuthTokenSecretKey = "prometheusAuthToken" - workQueueKey = "key" - indexerNodeSelectorGlobal = "indexer_node_selector_global" + indexerNodeSelectorGlobal = "indexer_node_selector_global" ) type eprunner func(ctx context.Context, nodes []*v1.Node) *frameworktypes.Status @@ -79,29 +69,6 @@ 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 @@ -117,58 +84,6 @@ type descheduler struct { profileRunners []profileRunner } -type ( - createPrometheusClientFunc func(url, token string) (promapi.Client, *http.Transport, error) - inClusterConfigFunc func() (*rest.Config, error) -) - -func newInClusterPromClientController(prometheusClient promapi.Client, prometheusConfig *api.Prometheus) *inClusterPromClientController { - return &inClusterPromClientController{ - promClient: prometheusClient, - prometheusConfig: prometheusConfig, - createPrometheusClient: client.CreatePrometheusClient, - inClusterConfig: rest.InClusterConfig, - } -} - -func newSecretBasedPromClientController(prometheusClient promapi.Client, prometheusConfig *api.Prometheus, namespacedSharedInformerFactory informers.SharedInformerFactory) (*secretBasedPromClientController, error) { - if prometheusConfig == nil || prometheusConfig.AuthToken == nil || prometheusConfig.AuthToken.SecretReference == nil { - return nil, fmt.Errorf("prometheus metrics source configuration is missing authentication token secret") - } - authTokenSecret := prometheusConfig.AuthToken.SecretReference - if authTokenSecret.Name == "" || authTokenSecret.Namespace == "" { - return nil, fmt.Errorf("prometheus metrics source configuration is missing authentication token secret") - } - - if namespacedSharedInformerFactory == nil { - return nil, fmt.Errorf("namespacedSharedInformerFactory not configured") - } - - ctrl := &secretBasedPromClientController{ - promClient: prometheusClient, - queue: workqueue.NewRateLimitingQueueWithConfig(workqueue.DefaultControllerRateLimiter(), workqueue.RateLimitingQueueConfig{Name: "descheduler"}), - prometheusConfig: prometheusConfig, - createPrometheusClient: client.CreatePrometheusClient, - } - - namespacedSharedInformerFactory.Core().V1().Secrets().Informer().AddEventHandler(ctrl.eventHandler()) - ctrl.namespacedSecretsLister = namespacedSharedInformerFactory.Core().V1().Secrets().Lister().Secrets(authTokenSecret.Namespace) - - return ctrl, nil -} - -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 -} - func nodeSelectorFromPolicy(deschedulerPolicy *api.DeschedulerPolicy) (labels.Selector, error) { nodeSelector := labels.Everything() if deschedulerPolicy.NodeSelector != nil { @@ -308,135 +223,6 @@ func newDescheduler(ctx context.Context, rs *options.DeschedulerServer, deschedu return desch, nil } -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.prometheusConfig.URL, cfg.BearerToken) - if err != nil { - d.clearConnection() - return fmt.Errorf("unable to create a prometheus client: %v", err) - } - d.promClient = prometheusClient - if d.previousPrometheusClientTransport != nil { - d.previousPrometheusClientTransport.CloseIdleConnections() - } - d.previousPrometheusClientTransport = transport - d.currentPrometheusAuthToken = cfg.BearerToken - } - return nil - } - if err == rest.ErrNotInCluster { - return nil - } - return fmt.Errorf("unexpected error when reading in cluster config: %v", err) -} - -func clearPromClientConnection(currentPrometheusAuthToken *string, previousPrometheusClientTransport **http.Transport, promClient *promapi.Client) { - *currentPrometheusAuthToken = "" - if *previousPrometheusClientTransport != nil { - (*previousPrometheusClientTransport).CloseIdleConnections() - } - *previousPrometheusClientTransport = nil - *promClient = nil -} - -func (d *inClusterPromClientController) clearConnection() { - clearPromClientConnection(&d.currentPrometheusAuthToken, &d.previousPrometheusClientTransport, &d.promClient) -} - -func (d *secretBasedPromClientController) clearConnection() { - clearPromClientConnection(&d.currentPrometheusAuthToken, &d.previousPrometheusClientTransport, &d.promClient) -} - -func (d *secretBasedPromClientController) runAuthenticationSecretReconciler(ctx context.Context) { - defer utilruntime.HandleCrash() - defer d.queue.ShutDown() - - klog.Infof("Starting authentication secret reconciler") - defer klog.Infof("Shutting down authentication secret reconciler") - - go wait.UntilWithContext(ctx, d.runAuthenticationSecretReconcilerWorker, time.Second) - - <-ctx.Done() -} - -func (d *secretBasedPromClientController) runAuthenticationSecretReconcilerWorker(ctx context.Context) { - for d.processNextWorkItem(ctx) { - } -} - -func (d *secretBasedPromClientController) processNextWorkItem(ctx context.Context) bool { - dsKey, quit := d.queue.Get() - if quit { - return false - } - defer d.queue.Done(dsKey) - - err := d.sync() - if err == nil { - d.queue.Forget(dsKey) - return true - } - - utilruntime.HandleError(fmt.Errorf("%v failed with : %v", dsKey, err)) - d.queue.AddRateLimited(dsKey) - - return true -} - -func (d *secretBasedPromClientController) sync() error { - d.mu.Lock() - defer d.mu.Unlock() - - prometheusConfig := d.prometheusConfig - ns := prometheusConfig.AuthToken.SecretReference.Namespace - name := prometheusConfig.AuthToken.SecretReference.Name - secretObj, err := d.namespacedSecretsLister.Get(name) - if err != nil { - // clear the token if the secret is not found - if apierrors.IsNotFound(err) { - d.clearConnection() - } - return fmt.Errorf("unable to get %v/%v secret", ns, name) - } - authToken := string(secretObj.Data[prometheusAuthTokenSecretKey]) - if authToken == "" { - d.clearConnection() - return fmt.Errorf("prometheus authentication token secret missing %q data or empty", prometheusAuthTokenSecretKey) - } - if d.currentPrometheusAuthToken == authToken { - return nil - } - - klog.V(2).Infof("authentication secret token updated, recreating prometheus client") - prometheusClient, transport, err := d.createPrometheusClient(prometheusConfig.URL, authToken) - if err != nil { - d.clearConnection() - return fmt.Errorf("unable to create a prometheus client: %v", err) - } - d.promClient = prometheusClient - if d.previousPrometheusClientTransport != nil { - d.previousPrometheusClientTransport.CloseIdleConnections() - } - d.previousPrometheusClientTransport = transport - d.currentPrometheusAuthToken = authToken - return nil -} - -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) }, - DeleteFunc: func(obj interface{}) { d.queue.Add(workQueueKey) }, - } -} - func (d *descheduler) runDeschedulerLoop(ctx context.Context) error { var span trace.Span ctx, span = tracing.Tracer().Start(ctx, "runDeschedulerLoop") diff --git a/pkg/descheduler/descheduler_test.go b/pkg/descheduler/descheduler_test.go index 7a7dc5889..b968e93c5 100644 --- a/pkg/descheduler/descheduler_test.go +++ b/pkg/descheduler/descheduler_test.go @@ -8,7 +8,6 @@ import ( "net/http" "net/url" "strings" - "sync" "testing" "time" @@ -1889,779 +1888,6 @@ func newPrometheusConfig() *api.Prometheus { } } -func TestPromClientControllerSync_InvalidConfig(t *testing.T) { - testCases := []struct { - name string - objects []runtime.Object - prometheusConfig *api.Prometheus - expectedErr error - }{ - { - name: "empty prometheus config", - prometheusConfig: nil, - expectedErr: fmt.Errorf("prometheus metrics source configuration is missing authentication token secret"), - }, - { - name: "missing prometheus config", - prometheusConfig: &api.Prometheus{ - URL: prometheusURL, - }, - expectedErr: fmt.Errorf("prometheus metrics source configuration is missing authentication token secret"), - }, - { - name: "missing auth token config", - prometheusConfig: &api.Prometheus{ - URL: prometheusURL, - AuthToken: nil, - }, - expectedErr: fmt.Errorf("prometheus metrics source configuration is missing authentication token secret"), - }, - { - name: "missing secret reference", - prometheusConfig: &api.Prometheus{ - URL: prometheusURL, - AuthToken: &api.AuthToken{ - SecretReference: nil, - }, - }, - expectedErr: fmt.Errorf("prometheus metrics source configuration is missing authentication token secret"), - }, - { - name: "missing secret reference name", - prometheusConfig: &api.Prometheus{ - URL: prometheusURL, - AuthToken: &api.AuthToken{ - SecretReference: &api.SecretReference{ - Namespace: "kube-system", - }, - }, - }, - expectedErr: fmt.Errorf("prometheus metrics source configuration is missing authentication token secret"), - }, - { - name: "missing secret reference namespace", - prometheusConfig: &api.Prometheus{ - URL: prometheusURL, - AuthToken: &api.AuthToken{ - SecretReference: &api.SecretReference{ - Name: "prom-token", - }, - }, - }, - expectedErr: fmt.Errorf("prometheus metrics source configuration is missing authentication token secret"), - }, - } - - for _, tc := range testCases { - t.Run(tc.name, func(t *testing.T) { - ctx, cancel := context.WithCancel(context.TODO()) - defer cancel() - _, err := setupPromClientControllerTest(ctx, t, tc.objects, tc.prometheusConfig, false) - - // Verify error expectations - if tc.expectedErr != nil { - if err == nil { - t.Errorf("Expected error %q but got none", tc.expectedErr) - } else if err.Error() != tc.expectedErr.Error() { - t.Errorf("Expected error %q but got %q", tc.expectedErr, err.Error()) - } - } else { - t.Errorf("Expected an error, got none") - } - }) - } -} - -func TestPromClientControllerSync_InvalidSecret(t *testing.T) { - testCases := []struct { - name string - objects []runtime.Object - prometheusConfig *api.Prometheus - expectedErr error - }{ - { - 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 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"), - }, - } - - for _, tc := range testCases { - t.Run(tc.name, func(t *testing.T) { - ctx, cancel := context.WithCancel(context.TODO()) - defer cancel() - setup, err := setupPromClientControllerTest(ctx, t, tc.objects, tc.prometheusConfig, true) - if err != nil { - t.Fatal(err) - } - - // Call sync - err = setup.ctrl.sync() - - // Verify error expectations - if tc.expectedErr != nil { - if err == nil { - t.Errorf("Expected error %q but got none", tc.expectedErr) - } else if err.Error() != tc.expectedErr.Error() { - t.Errorf("Expected error %q but got %q", tc.expectedErr, err.Error()) - } - } else { - t.Errorf("Expected an error, got none") - } - }) - } -} - -func TestPromClientControllerSync_ClientCreation(t *testing.T) { - testCases := []struct { - name string - objects []runtime.Object - currentAuthToken string - createPrometheusClientFunc func(url, token string) (promapi.Client, *http.Transport, error) - expectedErr error - expectClientCreated bool - expectCurrentTokenCleared bool - expectPreviousTransportCleared bool - }{ - { - name: "secret not found", - currentAuthToken: "old-token", - expectedErr: fmt.Errorf("unable to get kube-system/prom-token secret"), - createPrometheusClientFunc: func(url, token string) (promapi.Client, *http.Transport, error) { - t.Fatalf("unexpected create client invocation") - return nil, nil, fmt.Errorf("unexpected create client invocation") - }, - expectCurrentTokenCleared: true, - expectPreviousTransportCleared: true, - }, - { - name: "token unchanged - no client creation", - objects: []runtime.Object{newPrometheusAuthSecret(withToken("same-token"))}, - currentAuthToken: "same-token", - createPrometheusClientFunc: func(url, token string) (promapi.Client, *http.Transport, error) { - t.Fatalf("unexpected create client invocation") - return nil, nil, fmt.Errorf("unexpected create client invocation") - }, - expectClientCreated: false, - }, - { - name: "token changed - client created successfully", - objects: []runtime.Object{newPrometheusAuthSecret(withToken("new-token"))}, - currentAuthToken: "old-token", - createPrometheusClientFunc: func(url, token string) (promapi.Client, *http.Transport, error) { - return &mockPrometheusClient{name: "new-client"}, &http.Transport{}, nil - }, - expectClientCreated: true, - }, - { - name: "token changed - client creation fails", - objects: []runtime.Object{newPrometheusAuthSecret(withToken("new-token"))}, - currentAuthToken: "old-token", - createPrometheusClientFunc: func(url, token string) (promapi.Client, *http.Transport, error) { - return nil, nil, fmt.Errorf("failed to create client") - }, - expectedErr: fmt.Errorf("unable to create a prometheus client: failed to create client"), - expectClientCreated: false, - }, - } - - for _, tc := range testCases { - t.Run(tc.name, func(t *testing.T) { - for _, setupMode := range []struct { - name string - 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) *secretBasedPromClientController { - setup, err := setupPromClientControllerTest(ctx, t, objects, newPrometheusConfig(), true) - if err != nil { - t.Fatal(err) - } - return setup.ctrl - }, - }, - { - name: "running with full descheduler", - setupFn: func(ctx context.Context, t *testing.T, objects []runtime.Object) *secretBasedPromClientController { - deschedulerPolicy := &api.DeschedulerPolicy{ - MetricsProviders: []api.MetricsProvider{ - { - Source: api.PrometheusMetrics, - Prometheus: newPrometheusConfig(), - }, - }, - } - _, descheduler, _, _ := initDescheduler(t, ctx, initFeatureGates(), deschedulerPolicy, nil, false, objects...) - return descheduler.secretBasedPromClientCtrl - }, - }, - } { - t.Run(setupMode.name, func(t *testing.T) { - ctx, cancel := context.WithCancel(context.TODO()) - defer cancel() - - ctrl := setupMode.setupFn(ctx, t, tc.objects) - - // Set additional test-specific fields - ctrl.currentPrometheusAuthToken = tc.currentAuthToken - if tc.currentAuthToken != "" { - ctrl.previousPrometheusClientTransport = &http.Transport{} - } - - // Mock createPrometheusClient - clientCreated := false - if tc.createPrometheusClientFunc != nil { - ctrl.createPrometheusClient = func(url, token string) (promapi.Client, *http.Transport, error) { - client, transport, err := tc.createPrometheusClientFunc(url, token) - if err == nil { - clientCreated = true - } - return client, transport, err - } - } - - // Call sync - err := ctrl.sync() - - // Verify error expectations - if tc.expectedErr != nil { - if err == nil { - t.Errorf("Expected error %q but got none", tc.expectedErr) - } else if err.Error() != tc.expectedErr.Error() { - t.Errorf("Expected error %q but got %q", tc.expectedErr, err.Error()) - } - } else { - if err != nil { - t.Errorf("Expected no error but got: %v", err) - } - } - - // Verify client creation expectations - if tc.expectClientCreated && !clientCreated { - t.Errorf("Expected prometheus client to be created but it wasn't") - } - if !tc.expectClientCreated && clientCreated { - t.Errorf("Expected prometheus client not to be created but it was") - } - - // Verify token cleared expectations - if tc.expectCurrentTokenCleared && ctrl.currentPrometheusAuthToken != "" { - t.Errorf("Expected current auth token to be cleared but it wasn't") - } - - // Verify previous transport cleared expectations - if tc.expectPreviousTransportCleared && ctrl.previousPrometheusClientTransport != nil { - t.Errorf("Expected previous transport to be cleared but it wasn't") - } - - // Verify promClient cleared when secret not found - if tc.expectPreviousTransportCleared && ctrl.prometheusClient() != nil { - t.Errorf("Expected promClient to be cleared but it wasn't") - } - - // Verify token updated when client created - if tc.expectClientCreated && len(tc.objects) > 0 { - if secret, ok := tc.objects[0].(*v1.Secret); ok && secret.Data != nil { - expectedToken := string(secret.Data[prometheusAuthTokenSecretKey]) - if ctrl.currentPrometheusAuthToken != expectedToken { - t.Errorf("Expected current auth token to be %q but got %q", expectedToken, ctrl.currentPrometheusAuthToken) - } - } - } - }) - } - }) - } -} - -func TestPromClientControllerSync_EventHandler(t *testing.T) { - testCases := []struct { - name string - operation func(ctx context.Context, fakeClient *fakeclientset.Clientset) error - processItem bool - expectedPromClientSet bool - expectedCreatedClientsCount int - expectedCurrentToken string - expectedPreviousTransportCleared bool - expectDifferentClients bool - expectCreatePrometheusClientError bool - }{ - // Check initial conditions - { - name: "no secret initially", - operation: func(ctx context.Context, fakeClient *fakeclientset.Clientset) error { return nil }, - processItem: false, - expectedPromClientSet: false, - expectedCreatedClientsCount: 0, - expectedCurrentToken: "", - }, - // Change conditions - { - name: "add secret", - operation: func(ctx context.Context, fakeClient *fakeclientset.Clientset) error { - secret := newPrometheusAuthSecret(withToken("token-1")) - _, err := fakeClient.CoreV1().Secrets(secret.Namespace).Create(ctx, secret, metav1.CreateOptions{}) - return err - }, - processItem: true, - expectedPromClientSet: true, - expectedCreatedClientsCount: 1, - expectedCurrentToken: "token-1", - }, - { - name: "update secret", - operation: func(ctx context.Context, fakeClient *fakeclientset.Clientset) error { - secret := newPrometheusAuthSecret(withToken("token-2")) - _, err := fakeClient.CoreV1().Secrets(secret.Namespace).Update(ctx, secret, metav1.UpdateOptions{}) - return err - }, - processItem: true, - expectedPromClientSet: true, - expectedCreatedClientsCount: 2, - expectedCurrentToken: "token-2", - expectDifferentClients: true, - }, - { - name: "update secret with invalid data", - operation: func(ctx context.Context, fakeClient *fakeclientset.Clientset) error { - secret := newPrometheusAuthSecret(withToken("token-3")) - secret.Data[prometheusAuthTokenSecretKey] = []byte{} - _, err := fakeClient.CoreV1().Secrets(secret.Namespace).Update(ctx, secret, metav1.UpdateOptions{}) - return err - }, - processItem: true, - expectedPromClientSet: false, - expectedCreatedClientsCount: 2, - expectedCurrentToken: "", - expectDifferentClients: true, - }, - { - name: "update secret with valid data", - operation: func(ctx context.Context, fakeClient *fakeclientset.Clientset) error { - secret := newPrometheusAuthSecret(withToken("token-4")) - _, err := fakeClient.CoreV1().Secrets(secret.Namespace).Update(ctx, secret, metav1.UpdateOptions{}) - return err - }, - processItem: true, - expectedPromClientSet: true, - expectedCreatedClientsCount: 3, - expectedCurrentToken: "token-4", - expectDifferentClients: true, - }, - { - name: "update secret with valid data but createPrometheusClient failing", - operation: func(ctx context.Context, fakeClient *fakeclientset.Clientset) error { - secret := newPrometheusAuthSecret(withToken("token-5")) - _, err := fakeClient.CoreV1().Secrets(secret.Namespace).Update(ctx, secret, metav1.UpdateOptions{}) - return err - }, - processItem: true, - expectedPromClientSet: false, - expectedCreatedClientsCount: 3, - expectedCurrentToken: "", - expectDifferentClients: true, - expectCreatePrometheusClientError: true, - }, - { - name: "delete secret", - operation: func(ctx context.Context, fakeClient *fakeclientset.Clientset) error { - secret := newPrometheusAuthSecret(withToken("token-5")) - return fakeClient.CoreV1().Secrets(secret.Namespace).Delete(ctx, secret.Name, metav1.DeleteOptions{}) - }, - processItem: true, - expectedPromClientSet: false, - expectedCreatedClientsCount: 3, - expectedCurrentToken: "", - expectedPreviousTransportCleared: true, - }, - } - - for _, setupMode := range []struct { - name string - 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 *secretBasedPromClientController, fakeClient *fakeclientset.Clientset) { - setup, err := setupPromClientControllerTest(ctx, t, nil, newPrometheusConfig(), true) - if err != nil { - t.Fatal(err) - } - - // Start the reconciler to process queue items - go setup.ctrl.runAuthenticationSecretReconciler(ctx) - - return setup.ctrl, setup.fakeClient - }, - }, - { - name: "running with full descheduler", - init: func(t *testing.T, ctx context.Context) (ctrl *secretBasedPromClientController, fakeClient *fakeclientset.Clientset) { - deschedulerPolicy := &api.DeschedulerPolicy{ - MetricsProviders: []api.MetricsProvider{ - { - Source: api.PrometheusMetrics, - Prometheus: newPrometheusConfig(), - }, - }, - } - - _, descheduler, _, client := initDescheduler(t, ctx, initFeatureGates(), deschedulerPolicy, nil, false) - // The reconciler is already started by initDescheduler via bootstrapDescheduler - - return descheduler.secretBasedPromClientCtrl, client - }, - }, - } { - t.Run(setupMode.name, func(t *testing.T) { - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() - - ctrl, fakeClient := setupMode.init(t, ctx) - - // Track created clients to verify different instances - var createdClients []promapi.Client - var createdClientsMu sync.Mutex - - for _, tc := range testCases { - t.Run(tc.name, func(t *testing.T) { - ctrl.mu.Lock() - ctrl.createPrometheusClient = func(url, token string) (promapi.Client, *http.Transport, error) { - if tc.expectCreatePrometheusClientError { - return nil, &http.Transport{}, fmt.Errorf("error creating a prometheus client") - } - client := &mockPrometheusClient{name: "client-" + token} - createdClientsMu.Lock() - createdClients = append(createdClients, client) - createdClientsMu.Unlock() - return client, &http.Transport{}, nil - } - ctrl.mu.Unlock() - - if err := tc.operation(ctx, fakeClient); err != nil { - t.Fatalf("Failed to execute operation: %v", err) - } - - 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 (with mutex protection) - ctrl.mu.RLock() - promClient := ctrl.promClient - currentToken := ctrl.currentPrometheusAuthToken - previousTransport := ctrl.previousPrometheusClientTransport - ctrl.mu.RUnlock() - - t.Logf("promClient: %v\n", promClient) - if tc.expectedPromClientSet { - if promClient == nil { - return false, nil - } - } else { - if promClient != nil { - return false, nil - } - } - - createdClientsMu.Lock() - createdClientsLen := len(createdClients) - createdClientsMu.Unlock() - t.Logf("createdClientsLen: %v\n", createdClientsLen) - if createdClientsLen != tc.expectedCreatedClientsCount { - return false, nil - } - t.Logf("currentToken: %v\n", currentToken) - if currentToken != tc.expectedCurrentToken { - return false, nil - } - t.Logf("previousTransport: %v\n", previousTransport) - if tc.expectedPreviousTransportCleared { - if previousTransport != nil { - return false, nil - } - } - - return true, nil - }) - if err != nil { - t.Fatalf("Timed out waiting for expected conditions: %v", err) - } - - // Log all expected conditions that were met - t.Logf("All expected conditions met: promClientSet=%v, createdClientsCount=%d, currentToken=%q, previousTransportCleared=%v", - tc.expectedPromClientSet, tc.expectedCreatedClientsCount, tc.expectedCurrentToken, tc.expectedPreviousTransportCleared) - } - - // Validate post-conditions - if tc.expectedPromClientSet { - if ctrl.prometheusClient() == nil { - t.Error("Expected prometheus client to be set, but it was nil") - } - } else { - promClient := ctrl.prometheusClient() - if promClient != nil { - t.Errorf("Expected prometheus client to be nil, but got: %v", promClient) - } - } - - createdClientsMu.Lock() - createdClientsLen := len(createdClients) - createdClientsMu.Unlock() - if createdClientsLen != tc.expectedCreatedClientsCount { - t.Errorf("Expected %d clients created, but got %d", tc.expectedCreatedClientsCount, len(createdClients)) - } - - if ctrl.currentPrometheusAuthToken != tc.expectedCurrentToken { - t.Errorf("Expected current token to be %q, got %q", tc.expectedCurrentToken, ctrl.currentPrometheusAuthToken) - } - - if tc.expectedPreviousTransportCleared { - if ctrl.previousPrometheusClientTransport != nil { - t.Error("Expected previous transport to be cleared, but it was set") - } - } - - if tc.expectDifferentClients && len(createdClients) >= 2 { - createdClientsMu.Lock() - defer createdClientsMu.Unlock() - if createdClients[0] == createdClients[1] { - t.Error("Expected different client instances") - } - } - }) - } - }) - } -} - -func TestReconcileInClusterSAToken(t *testing.T) { - testCases := []struct { - name string - currentAuthToken string - inClusterConfigFunc func() (*rest.Config, error) - createPrometheusClientFunc func(url, token string) (promapi.Client, *http.Transport, error) - expectedErr error - expectClientCreated bool - expectCurrentToken string - expectPreviousTransportCleared bool - expectPromClientCleared bool - }{ - { - name: "token unchanged - no client creation", - currentAuthToken: "same-token", - inClusterConfigFunc: func() (*rest.Config, error) { - return &rest.Config{BearerToken: "same-token"}, nil - }, - expectClientCreated: false, - expectCurrentToken: "same-token", - }, - { - name: "token changed - client created successfully", - currentAuthToken: "old-token", - inClusterConfigFunc: func() (*rest.Config, error) { - return &rest.Config{BearerToken: "new-token"}, nil - }, - createPrometheusClientFunc: func(url, token string) (promapi.Client, *http.Transport, error) { - if token != "new-token" { - t.Errorf("Expected token to be %q, got %q", "new-token", token) - } - return &mockPrometheusClient{name: "new-client"}, &http.Transport{}, nil - }, - expectClientCreated: true, - expectCurrentToken: "new-token", - }, - { - name: "token changed - client creation fails", - currentAuthToken: "old-token", - inClusterConfigFunc: func() (*rest.Config, error) { - return &rest.Config{BearerToken: "new-token"}, nil - }, - createPrometheusClientFunc: func(url, token string) (promapi.Client, *http.Transport, error) { - return nil, nil, fmt.Errorf("failed to create client") - }, - expectedErr: fmt.Errorf("unable to create a prometheus client: failed to create client"), - expectClientCreated: false, - expectCurrentToken: "", - expectPreviousTransportCleared: false, - expectPromClientCleared: false, - }, - { - name: "not in cluster - no error", - currentAuthToken: "current-token", - inClusterConfigFunc: func() (*rest.Config, error) { - return nil, rest.ErrNotInCluster - }, - expectClientCreated: false, - expectCurrentToken: "current-token", - }, - { - name: "unexpected error", - currentAuthToken: "current-token", - inClusterConfigFunc: func() (*rest.Config, error) { - return nil, fmt.Errorf("unexpected error") - }, - expectedErr: fmt.Errorf("unexpected error when reading in cluster config: unexpected error"), - expectClientCreated: false, - expectCurrentToken: "current-token", - }, - { - name: "first token - client created successfully", - currentAuthToken: "", - inClusterConfigFunc: func() (*rest.Config, error) { - return &rest.Config{BearerToken: "first-token"}, nil - }, - createPrometheusClientFunc: func(url, token string) (promapi.Client, *http.Transport, error) { - return &mockPrometheusClient{name: "first-client"}, &http.Transport{}, nil - }, - expectClientCreated: true, - expectCurrentToken: "first-token", - }, - { - name: "token changed with previous transport - clears previous transport", - currentAuthToken: "old-token", - inClusterConfigFunc: func() (*rest.Config, error) { - return &rest.Config{BearerToken: "new-token"}, nil - }, - createPrometheusClientFunc: func(url, token string) (promapi.Client, *http.Transport, error) { - return &mockPrometheusClient{name: "new-client"}, &http.Transport{}, nil - }, - expectClientCreated: true, - expectCurrentToken: "new-token", - expectPreviousTransportCleared: true, - }, - } - - for _, tc := range testCases { - t.Run(tc.name, func(t *testing.T) { - for _, setupMode := range []struct { - name string - init func(t *testing.T) *inClusterPromClientController - }{ - { - name: "running with prom reconciler directly", - init: func(t *testing.T) *inClusterPromClientController { - return &inClusterPromClientController{ - currentPrometheusAuthToken: tc.currentAuthToken, - prometheusConfig: &api.Prometheus{ - URL: prometheusURL, - }, - inClusterConfig: tc.inClusterConfigFunc, - } - }, - }, - { - name: "running with full descheduler", - init: func(t *testing.T) *inClusterPromClientController { - ctx := context.Background() - prometheusConfig := &api.Prometheus{ - URL: prometheusURL, - } - deschedulerPolicy := &api.DeschedulerPolicy{ - MetricsProviders: []api.MetricsProvider{ - { - Source: api.PrometheusMetrics, - Prometheus: prometheusConfig, - }, - }, - } - _, descheduler, _, _ := initDescheduler(t, ctx, initFeatureGates(), deschedulerPolicy, nil, false) - - // 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 - }, - }, - } { - t.Run(setupMode.name, func(t *testing.T) { - ctrl := setupMode.init(t) - - // Set previous transport and client if test expects them to be cleared - if tc.expectPreviousTransportCleared { - ctrl.previousPrometheusClientTransport = &http.Transport{} - } - if tc.expectPromClientCleared { - ctrl.promClient = &mockPrometheusClient{name: "old-client"} - } - - // Mock createPrometheusClient - clientCreated := false - if tc.createPrometheusClientFunc != nil { - ctrl.createPrometheusClient = func(url, token string) (promapi.Client, *http.Transport, error) { - client, transport, err := tc.createPrometheusClientFunc(url, token) - if err == nil { - clientCreated = true - } - return client, transport, err - } - } - - // Call reconcileInClusterSAToken - err := ctrl.reconcileInClusterSAToken() - - // Verify error expectations - if tc.expectedErr != nil { - if err == nil { - t.Errorf("Expected error %q but got none", tc.expectedErr) - } else if err.Error() != tc.expectedErr.Error() { - t.Errorf("Expected error %q but got %q", tc.expectedErr, err.Error()) - } - } else { - if err != nil { - t.Errorf("Expected no error but got: %v", err) - } - } - - // Verify client creation expectations - if tc.expectClientCreated && !clientCreated { - t.Errorf("Expected prometheus client to be created but it wasn't") - } - if !tc.expectClientCreated && clientCreated { - t.Errorf("Expected prometheus client not to be created but it was") - } - - // Verify token expectations - if ctrl.currentPrometheusAuthToken != tc.expectCurrentToken { - t.Errorf("Expected current token to be %q but got %q", tc.expectCurrentToken, ctrl.currentPrometheusAuthToken) - } - - // Verify previous transport cleared when expected - if tc.expectPreviousTransportCleared { - if tc.expectClientCreated { - // Success case: new transport should be set - if ctrl.previousPrometheusClientTransport == nil { - t.Error("Expected previous transport to be set to new transport, but it was nil") - } - } else if tc.expectedErr != nil { - // Failure case: transport should be nil - if ctrl.previousPrometheusClientTransport != nil { - t.Error("Expected previous transport to be cleared on error, but it was set") - } - } - } - - // Verify promClient cleared when expected - if tc.expectPromClientCleared { - if ctrl.promClient != nil { - t.Error("Expected promClient to be cleared, but it was set") - } - } - }) - } - }) - } -} // TestPluginInformerRegistration tests that plugin-specific informers are registered during newDescheduler func TestPluginInformerRegistration(t *testing.T) { diff --git a/pkg/descheduler/prom_client_controller.go b/pkg/descheduler/prom_client_controller.go new file mode 100644 index 000000000..5cd33a2b5 --- /dev/null +++ b/pkg/descheduler/prom_client_controller.go @@ -0,0 +1,249 @@ +/* +Copyright 2026 The Kubernetes Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package descheduler + +import ( + "context" + "fmt" + "net/http" + "sync" + "time" + + promapi "github.com/prometheus/client_golang/api" + + apierrors "k8s.io/apimachinery/pkg/api/errors" + utilruntime "k8s.io/apimachinery/pkg/util/runtime" + "k8s.io/apimachinery/pkg/util/wait" + "k8s.io/client-go/informers" + corev1listers "k8s.io/client-go/listers/core/v1" + "k8s.io/client-go/rest" + "k8s.io/client-go/tools/cache" + "k8s.io/client-go/util/workqueue" + "k8s.io/klog/v2" + + "sigs.k8s.io/descheduler/pkg/api" + "sigs.k8s.io/descheduler/pkg/descheduler/client" +) + +const ( + prometheusAuthTokenSecretKey = "prometheusAuthToken" + workQueueKey = "key" +) + +// 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 ( + createPrometheusClientFunc func(url, token string) (promapi.Client, *http.Transport, error) + inClusterConfigFunc func() (*rest.Config, error) +) + +func newInClusterPromClientController(prometheusClient promapi.Client, prometheusConfig *api.Prometheus) *inClusterPromClientController { + return &inClusterPromClientController{ + promClient: prometheusClient, + prometheusConfig: prometheusConfig, + createPrometheusClient: client.CreatePrometheusClient, + inClusterConfig: rest.InClusterConfig, + } +} + +func newSecretBasedPromClientController(prometheusClient promapi.Client, prometheusConfig *api.Prometheus, namespacedSharedInformerFactory informers.SharedInformerFactory) (*secretBasedPromClientController, error) { + if prometheusConfig == nil || prometheusConfig.AuthToken == nil || prometheusConfig.AuthToken.SecretReference == nil { + return nil, fmt.Errorf("prometheus metrics source configuration is missing authentication token secret") + } + authTokenSecret := prometheusConfig.AuthToken.SecretReference + if authTokenSecret.Name == "" || authTokenSecret.Namespace == "" { + return nil, fmt.Errorf("prometheus metrics source configuration is missing authentication token secret") + } + + if namespacedSharedInformerFactory == nil { + return nil, fmt.Errorf("namespacedSharedInformerFactory not configured") + } + + ctrl := &secretBasedPromClientController{ + promClient: prometheusClient, + queue: workqueue.NewRateLimitingQueueWithConfig(workqueue.DefaultControllerRateLimiter(), workqueue.RateLimitingQueueConfig{Name: "descheduler"}), + prometheusConfig: prometheusConfig, + createPrometheusClient: client.CreatePrometheusClient, + } + + namespacedSharedInformerFactory.Core().V1().Secrets().Informer().AddEventHandler(ctrl.eventHandler()) + ctrl.namespacedSecretsLister = namespacedSharedInformerFactory.Core().V1().Secrets().Lister().Secrets(authTokenSecret.Namespace) + + return ctrl, nil +} + +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 +} + +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.prometheusConfig.URL, cfg.BearerToken) + if err != nil { + d.clearConnection() + return fmt.Errorf("unable to create a prometheus client: %v", err) + } + d.promClient = prometheusClient + if d.previousPrometheusClientTransport != nil { + d.previousPrometheusClientTransport.CloseIdleConnections() + } + d.previousPrometheusClientTransport = transport + d.currentPrometheusAuthToken = cfg.BearerToken + } + return nil + } + if err == rest.ErrNotInCluster { + return nil + } + return fmt.Errorf("unexpected error when reading in cluster config: %v", err) +} + +func clearPromClientConnection(currentPrometheusAuthToken *string, previousPrometheusClientTransport **http.Transport, promClient *promapi.Client) { + *currentPrometheusAuthToken = "" + if *previousPrometheusClientTransport != nil { + (*previousPrometheusClientTransport).CloseIdleConnections() + } + *previousPrometheusClientTransport = nil + *promClient = nil +} + +func (d *inClusterPromClientController) clearConnection() { + clearPromClientConnection(&d.currentPrometheusAuthToken, &d.previousPrometheusClientTransport, &d.promClient) +} + +func (d *secretBasedPromClientController) clearConnection() { + clearPromClientConnection(&d.currentPrometheusAuthToken, &d.previousPrometheusClientTransport, &d.promClient) +} + +func (d *secretBasedPromClientController) runAuthenticationSecretReconciler(ctx context.Context) { + defer utilruntime.HandleCrash() + defer d.queue.ShutDown() + + klog.Infof("Starting authentication secret reconciler") + defer klog.Infof("Shutting down authentication secret reconciler") + + go wait.UntilWithContext(ctx, d.runAuthenticationSecretReconcilerWorker, time.Second) + + <-ctx.Done() +} + +func (d *secretBasedPromClientController) runAuthenticationSecretReconcilerWorker(ctx context.Context) { + for d.processNextWorkItem(ctx) { + } +} + +func (d *secretBasedPromClientController) processNextWorkItem(ctx context.Context) bool { + dsKey, quit := d.queue.Get() + if quit { + return false + } + defer d.queue.Done(dsKey) + + err := d.sync() + if err == nil { + d.queue.Forget(dsKey) + return true + } + + utilruntime.HandleError(fmt.Errorf("%v failed with : %v", dsKey, err)) + d.queue.AddRateLimited(dsKey) + + return true +} + +func (d *secretBasedPromClientController) sync() error { + d.mu.Lock() + defer d.mu.Unlock() + + prometheusConfig := d.prometheusConfig + ns := prometheusConfig.AuthToken.SecretReference.Namespace + name := prometheusConfig.AuthToken.SecretReference.Name + secretObj, err := d.namespacedSecretsLister.Get(name) + if err != nil { + // clear the token if the secret is not found + if apierrors.IsNotFound(err) { + d.clearConnection() + } + return fmt.Errorf("unable to get %v/%v secret", ns, name) + } + authToken := string(secretObj.Data[prometheusAuthTokenSecretKey]) + if authToken == "" { + d.clearConnection() + return fmt.Errorf("prometheus authentication token secret missing %q data or empty", prometheusAuthTokenSecretKey) + } + if d.currentPrometheusAuthToken == authToken { + return nil + } + + klog.V(2).Infof("authentication secret token updated, recreating prometheus client") + prometheusClient, transport, err := d.createPrometheusClient(prometheusConfig.URL, authToken) + if err != nil { + d.clearConnection() + return fmt.Errorf("unable to create a prometheus client: %v", err) + } + d.promClient = prometheusClient + if d.previousPrometheusClientTransport != nil { + d.previousPrometheusClientTransport.CloseIdleConnections() + } + d.previousPrometheusClientTransport = transport + d.currentPrometheusAuthToken = authToken + return nil +} + +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) }, + DeleteFunc: func(obj interface{}) { d.queue.Add(workQueueKey) }, + } +} diff --git a/pkg/descheduler/prom_client_controller_test.go b/pkg/descheduler/prom_client_controller_test.go new file mode 100644 index 000000000..903cc798b --- /dev/null +++ b/pkg/descheduler/prom_client_controller_test.go @@ -0,0 +1,811 @@ +/* +Copyright 2026 The Kubernetes Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package descheduler + +import ( + "context" + "fmt" + "net/http" + "sync" + "testing" + "time" + + promapi "github.com/prometheus/client_golang/api" + + v1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/util/wait" + fakeclientset "k8s.io/client-go/kubernetes/fake" + "k8s.io/client-go/rest" + + "sigs.k8s.io/descheduler/pkg/api" +) + +func TestPromClientControllerSync_InvalidConfig(t *testing.T) { + testCases := []struct { + name string + objects []runtime.Object + prometheusConfig *api.Prometheus + expectedErr error + }{ + { + name: "empty prometheus config", + prometheusConfig: nil, + expectedErr: fmt.Errorf("prometheus metrics source configuration is missing authentication token secret"), + }, + { + name: "missing prometheus config", + prometheusConfig: &api.Prometheus{ + URL: prometheusURL, + }, + expectedErr: fmt.Errorf("prometheus metrics source configuration is missing authentication token secret"), + }, + { + name: "missing auth token config", + prometheusConfig: &api.Prometheus{ + URL: prometheusURL, + AuthToken: nil, + }, + expectedErr: fmt.Errorf("prometheus metrics source configuration is missing authentication token secret"), + }, + { + name: "missing secret reference", + prometheusConfig: &api.Prometheus{ + URL: prometheusURL, + AuthToken: &api.AuthToken{ + SecretReference: nil, + }, + }, + expectedErr: fmt.Errorf("prometheus metrics source configuration is missing authentication token secret"), + }, + { + name: "missing secret reference name", + prometheusConfig: &api.Prometheus{ + URL: prometheusURL, + AuthToken: &api.AuthToken{ + SecretReference: &api.SecretReference{ + Namespace: "kube-system", + }, + }, + }, + expectedErr: fmt.Errorf("prometheus metrics source configuration is missing authentication token secret"), + }, + { + name: "missing secret reference namespace", + prometheusConfig: &api.Prometheus{ + URL: prometheusURL, + AuthToken: &api.AuthToken{ + SecretReference: &api.SecretReference{ + Name: "prom-token", + }, + }, + }, + expectedErr: fmt.Errorf("prometheus metrics source configuration is missing authentication token secret"), + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + ctx, cancel := context.WithCancel(context.TODO()) + defer cancel() + _, err := setupPromClientControllerTest(ctx, t, tc.objects, tc.prometheusConfig, false) + + // Verify error expectations + if tc.expectedErr != nil { + if err == nil { + t.Errorf("Expected error %q but got none", tc.expectedErr) + } else if err.Error() != tc.expectedErr.Error() { + t.Errorf("Expected error %q but got %q", tc.expectedErr, err.Error()) + } + } else { + t.Errorf("Expected an error, got none") + } + }) + } +} + +func TestPromClientControllerSync_InvalidSecret(t *testing.T) { + testCases := []struct { + name string + objects []runtime.Object + prometheusConfig *api.Prometheus + expectedErr error + }{ + { + 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 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"), + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + ctx, cancel := context.WithCancel(context.TODO()) + defer cancel() + setup, err := setupPromClientControllerTest(ctx, t, tc.objects, tc.prometheusConfig, true) + if err != nil { + t.Fatal(err) + } + + // Call sync + err = setup.ctrl.sync() + + // Verify error expectations + if tc.expectedErr != nil { + if err == nil { + t.Errorf("Expected error %q but got none", tc.expectedErr) + } else if err.Error() != tc.expectedErr.Error() { + t.Errorf("Expected error %q but got %q", tc.expectedErr, err.Error()) + } + } else { + t.Errorf("Expected an error, got none") + } + }) + } +} + +func TestPromClientControllerSync_ClientCreation(t *testing.T) { + testCases := []struct { + name string + objects []runtime.Object + currentAuthToken string + createPrometheusClientFunc func(url, token string) (promapi.Client, *http.Transport, error) + expectedErr error + expectClientCreated bool + expectCurrentTokenCleared bool + expectPreviousTransportCleared bool + }{ + { + name: "secret not found", + currentAuthToken: "old-token", + expectedErr: fmt.Errorf("unable to get kube-system/prom-token secret"), + createPrometheusClientFunc: func(url, token string) (promapi.Client, *http.Transport, error) { + t.Fatalf("unexpected create client invocation") + return nil, nil, fmt.Errorf("unexpected create client invocation") + }, + expectCurrentTokenCleared: true, + expectPreviousTransportCleared: true, + }, + { + name: "token unchanged - no client creation", + objects: []runtime.Object{newPrometheusAuthSecret(withToken("same-token"))}, + currentAuthToken: "same-token", + createPrometheusClientFunc: func(url, token string) (promapi.Client, *http.Transport, error) { + t.Fatalf("unexpected create client invocation") + return nil, nil, fmt.Errorf("unexpected create client invocation") + }, + expectClientCreated: false, + }, + { + name: "token changed - client created successfully", + objects: []runtime.Object{newPrometheusAuthSecret(withToken("new-token"))}, + currentAuthToken: "old-token", + createPrometheusClientFunc: func(url, token string) (promapi.Client, *http.Transport, error) { + return &mockPrometheusClient{name: "new-client"}, &http.Transport{}, nil + }, + expectClientCreated: true, + }, + { + name: "token changed - client creation fails", + objects: []runtime.Object{newPrometheusAuthSecret(withToken("new-token"))}, + currentAuthToken: "old-token", + createPrometheusClientFunc: func(url, token string) (promapi.Client, *http.Transport, error) { + return nil, nil, fmt.Errorf("failed to create client") + }, + expectedErr: fmt.Errorf("unable to create a prometheus client: failed to create client"), + expectClientCreated: false, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + for _, setupMode := range []struct { + name string + 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) *secretBasedPromClientController { + setup, err := setupPromClientControllerTest(ctx, t, objects, newPrometheusConfig(), true) + if err != nil { + t.Fatal(err) + } + return setup.ctrl + }, + }, + { + name: "running with full descheduler", + setupFn: func(ctx context.Context, t *testing.T, objects []runtime.Object) *secretBasedPromClientController { + deschedulerPolicy := &api.DeschedulerPolicy{ + MetricsProviders: []api.MetricsProvider{ + { + Source: api.PrometheusMetrics, + Prometheus: newPrometheusConfig(), + }, + }, + } + _, descheduler, _, _ := initDescheduler(t, ctx, initFeatureGates(), deschedulerPolicy, nil, false, objects...) + return descheduler.secretBasedPromClientCtrl + }, + }, + } { + t.Run(setupMode.name, func(t *testing.T) { + ctx, cancel := context.WithCancel(context.TODO()) + defer cancel() + + ctrl := setupMode.setupFn(ctx, t, tc.objects) + + // Set additional test-specific fields + ctrl.currentPrometheusAuthToken = tc.currentAuthToken + if tc.currentAuthToken != "" { + ctrl.previousPrometheusClientTransport = &http.Transport{} + } + + // Mock createPrometheusClient + clientCreated := false + if tc.createPrometheusClientFunc != nil { + ctrl.createPrometheusClient = func(url, token string) (promapi.Client, *http.Transport, error) { + client, transport, err := tc.createPrometheusClientFunc(url, token) + if err == nil { + clientCreated = true + } + return client, transport, err + } + } + + // Call sync + err := ctrl.sync() + + // Verify error expectations + if tc.expectedErr != nil { + if err == nil { + t.Errorf("Expected error %q but got none", tc.expectedErr) + } else if err.Error() != tc.expectedErr.Error() { + t.Errorf("Expected error %q but got %q", tc.expectedErr, err.Error()) + } + } else { + if err != nil { + t.Errorf("Expected no error but got: %v", err) + } + } + + // Verify client creation expectations + if tc.expectClientCreated && !clientCreated { + t.Errorf("Expected prometheus client to be created but it wasn't") + } + if !tc.expectClientCreated && clientCreated { + t.Errorf("Expected prometheus client not to be created but it was") + } + + // Verify token cleared expectations + if tc.expectCurrentTokenCleared && ctrl.currentPrometheusAuthToken != "" { + t.Errorf("Expected current auth token to be cleared but it wasn't") + } + + // Verify previous transport cleared expectations + if tc.expectPreviousTransportCleared && ctrl.previousPrometheusClientTransport != nil { + t.Errorf("Expected previous transport to be cleared but it wasn't") + } + + // Verify promClient cleared when secret not found + if tc.expectPreviousTransportCleared && ctrl.prometheusClient() != nil { + t.Errorf("Expected promClient to be cleared but it wasn't") + } + + // Verify token updated when client created + if tc.expectClientCreated && len(tc.objects) > 0 { + if secret, ok := tc.objects[0].(*v1.Secret); ok && secret.Data != nil { + expectedToken := string(secret.Data[prometheusAuthTokenSecretKey]) + if ctrl.currentPrometheusAuthToken != expectedToken { + t.Errorf("Expected current auth token to be %q but got %q", expectedToken, ctrl.currentPrometheusAuthToken) + } + } + } + }) + } + }) + } +} + +func TestPromClientControllerSync_EventHandler(t *testing.T) { + testCases := []struct { + name string + operation func(ctx context.Context, fakeClient *fakeclientset.Clientset) error + processItem bool + expectedPromClientSet bool + expectedCreatedClientsCount int + expectedCurrentToken string + expectedPreviousTransportCleared bool + expectDifferentClients bool + expectCreatePrometheusClientError bool + }{ + // Check initial conditions + { + name: "no secret initially", + operation: func(ctx context.Context, fakeClient *fakeclientset.Clientset) error { return nil }, + processItem: false, + expectedPromClientSet: false, + expectedCreatedClientsCount: 0, + expectedCurrentToken: "", + }, + // Change conditions + { + name: "add secret", + operation: func(ctx context.Context, fakeClient *fakeclientset.Clientset) error { + secret := newPrometheusAuthSecret(withToken("token-1")) + _, err := fakeClient.CoreV1().Secrets(secret.Namespace).Create(ctx, secret, metav1.CreateOptions{}) + return err + }, + processItem: true, + expectedPromClientSet: true, + expectedCreatedClientsCount: 1, + expectedCurrentToken: "token-1", + }, + { + name: "update secret", + operation: func(ctx context.Context, fakeClient *fakeclientset.Clientset) error { + secret := newPrometheusAuthSecret(withToken("token-2")) + _, err := fakeClient.CoreV1().Secrets(secret.Namespace).Update(ctx, secret, metav1.UpdateOptions{}) + return err + }, + processItem: true, + expectedPromClientSet: true, + expectedCreatedClientsCount: 2, + expectedCurrentToken: "token-2", + expectDifferentClients: true, + }, + { + name: "update secret with invalid data", + operation: func(ctx context.Context, fakeClient *fakeclientset.Clientset) error { + secret := newPrometheusAuthSecret(withToken("token-3")) + secret.Data[prometheusAuthTokenSecretKey] = []byte{} + _, err := fakeClient.CoreV1().Secrets(secret.Namespace).Update(ctx, secret, metav1.UpdateOptions{}) + return err + }, + processItem: true, + expectedPromClientSet: false, + expectedCreatedClientsCount: 2, + expectedCurrentToken: "", + expectDifferentClients: true, + }, + { + name: "update secret with valid data", + operation: func(ctx context.Context, fakeClient *fakeclientset.Clientset) error { + secret := newPrometheusAuthSecret(withToken("token-4")) + _, err := fakeClient.CoreV1().Secrets(secret.Namespace).Update(ctx, secret, metav1.UpdateOptions{}) + return err + }, + processItem: true, + expectedPromClientSet: true, + expectedCreatedClientsCount: 3, + expectedCurrentToken: "token-4", + expectDifferentClients: true, + }, + { + name: "update secret with valid data but createPrometheusClient failing", + operation: func(ctx context.Context, fakeClient *fakeclientset.Clientset) error { + secret := newPrometheusAuthSecret(withToken("token-5")) + _, err := fakeClient.CoreV1().Secrets(secret.Namespace).Update(ctx, secret, metav1.UpdateOptions{}) + return err + }, + processItem: true, + expectedPromClientSet: false, + expectedCreatedClientsCount: 3, + expectedCurrentToken: "", + expectDifferentClients: true, + expectCreatePrometheusClientError: true, + }, + { + name: "delete secret", + operation: func(ctx context.Context, fakeClient *fakeclientset.Clientset) error { + secret := newPrometheusAuthSecret(withToken("token-5")) + return fakeClient.CoreV1().Secrets(secret.Namespace).Delete(ctx, secret.Name, metav1.DeleteOptions{}) + }, + processItem: true, + expectedPromClientSet: false, + expectedCreatedClientsCount: 3, + expectedCurrentToken: "", + expectedPreviousTransportCleared: true, + }, + } + + for _, setupMode := range []struct { + name string + 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 *secretBasedPromClientController, fakeClient *fakeclientset.Clientset) { + setup, err := setupPromClientControllerTest(ctx, t, nil, newPrometheusConfig(), true) + if err != nil { + t.Fatal(err) + } + + // Start the reconciler to process queue items + go setup.ctrl.runAuthenticationSecretReconciler(ctx) + + return setup.ctrl, setup.fakeClient + }, + }, + { + name: "running with full descheduler", + init: func(t *testing.T, ctx context.Context) (ctrl *secretBasedPromClientController, fakeClient *fakeclientset.Clientset) { + deschedulerPolicy := &api.DeschedulerPolicy{ + MetricsProviders: []api.MetricsProvider{ + { + Source: api.PrometheusMetrics, + Prometheus: newPrometheusConfig(), + }, + }, + } + + _, descheduler, _, client := initDescheduler(t, ctx, initFeatureGates(), deschedulerPolicy, nil, false) + // The reconciler is already started by initDescheduler via bootstrapDescheduler + + return descheduler.secretBasedPromClientCtrl, client + }, + }, + } { + t.Run(setupMode.name, func(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + ctrl, fakeClient := setupMode.init(t, ctx) + + // Track created clients to verify different instances + var createdClients []promapi.Client + var createdClientsMu sync.Mutex + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + ctrl.mu.Lock() + ctrl.createPrometheusClient = func(url, token string) (promapi.Client, *http.Transport, error) { + if tc.expectCreatePrometheusClientError { + return nil, &http.Transport{}, fmt.Errorf("error creating a prometheus client") + } + client := &mockPrometheusClient{name: "client-" + token} + createdClientsMu.Lock() + createdClients = append(createdClients, client) + createdClientsMu.Unlock() + return client, &http.Transport{}, nil + } + ctrl.mu.Unlock() + + if err := tc.operation(ctx, fakeClient); err != nil { + t.Fatalf("Failed to execute operation: %v", err) + } + + 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 (with mutex protection) + ctrl.mu.RLock() + promClient := ctrl.promClient + currentToken := ctrl.currentPrometheusAuthToken + previousTransport := ctrl.previousPrometheusClientTransport + ctrl.mu.RUnlock() + + t.Logf("promClient: %v\n", promClient) + if tc.expectedPromClientSet { + if promClient == nil { + return false, nil + } + } else { + if promClient != nil { + return false, nil + } + } + + createdClientsMu.Lock() + createdClientsLen := len(createdClients) + createdClientsMu.Unlock() + t.Logf("createdClientsLen: %v\n", createdClientsLen) + if createdClientsLen != tc.expectedCreatedClientsCount { + return false, nil + } + t.Logf("currentToken: %v\n", currentToken) + if currentToken != tc.expectedCurrentToken { + return false, nil + } + t.Logf("previousTransport: %v\n", previousTransport) + if tc.expectedPreviousTransportCleared { + if previousTransport != nil { + return false, nil + } + } + + return true, nil + }) + if err != nil { + t.Fatalf("Timed out waiting for expected conditions: %v", err) + } + + // Log all expected conditions that were met + t.Logf("All expected conditions met: promClientSet=%v, createdClientsCount=%d, currentToken=%q, previousTransportCleared=%v", + tc.expectedPromClientSet, tc.expectedCreatedClientsCount, tc.expectedCurrentToken, tc.expectedPreviousTransportCleared) + } + + // Validate post-conditions + if tc.expectedPromClientSet { + if ctrl.prometheusClient() == nil { + t.Error("Expected prometheus client to be set, but it was nil") + } + } else { + promClient := ctrl.prometheusClient() + if promClient != nil { + t.Errorf("Expected prometheus client to be nil, but got: %v", promClient) + } + } + + createdClientsMu.Lock() + createdClientsLen := len(createdClients) + createdClientsMu.Unlock() + if createdClientsLen != tc.expectedCreatedClientsCount { + t.Errorf("Expected %d clients created, but got %d", tc.expectedCreatedClientsCount, len(createdClients)) + } + + if ctrl.currentPrometheusAuthToken != tc.expectedCurrentToken { + t.Errorf("Expected current token to be %q, got %q", tc.expectedCurrentToken, ctrl.currentPrometheusAuthToken) + } + + if tc.expectedPreviousTransportCleared { + if ctrl.previousPrometheusClientTransport != nil { + t.Error("Expected previous transport to be cleared, but it was set") + } + } + + if tc.expectDifferentClients && len(createdClients) >= 2 { + createdClientsMu.Lock() + defer createdClientsMu.Unlock() + if createdClients[0] == createdClients[1] { + t.Error("Expected different client instances") + } + } + }) + } + }) + } +} + +func TestReconcileInClusterSAToken(t *testing.T) { + testCases := []struct { + name string + currentAuthToken string + inClusterConfigFunc func() (*rest.Config, error) + createPrometheusClientFunc func(url, token string) (promapi.Client, *http.Transport, error) + expectedErr error + expectClientCreated bool + expectCurrentToken string + expectPreviousTransportCleared bool + expectPromClientCleared bool + }{ + { + name: "token unchanged - no client creation", + currentAuthToken: "same-token", + inClusterConfigFunc: func() (*rest.Config, error) { + return &rest.Config{BearerToken: "same-token"}, nil + }, + expectClientCreated: false, + expectCurrentToken: "same-token", + }, + { + name: "token changed - client created successfully", + currentAuthToken: "old-token", + inClusterConfigFunc: func() (*rest.Config, error) { + return &rest.Config{BearerToken: "new-token"}, nil + }, + createPrometheusClientFunc: func(url, token string) (promapi.Client, *http.Transport, error) { + if token != "new-token" { + t.Errorf("Expected token to be %q, got %q", "new-token", token) + } + return &mockPrometheusClient{name: "new-client"}, &http.Transport{}, nil + }, + expectClientCreated: true, + expectCurrentToken: "new-token", + }, + { + name: "token changed - client creation fails", + currentAuthToken: "old-token", + inClusterConfigFunc: func() (*rest.Config, error) { + return &rest.Config{BearerToken: "new-token"}, nil + }, + createPrometheusClientFunc: func(url, token string) (promapi.Client, *http.Transport, error) { + return nil, nil, fmt.Errorf("failed to create client") + }, + expectedErr: fmt.Errorf("unable to create a prometheus client: failed to create client"), + expectClientCreated: false, + expectCurrentToken: "", + expectPreviousTransportCleared: false, + expectPromClientCleared: false, + }, + { + name: "not in cluster - no error", + currentAuthToken: "current-token", + inClusterConfigFunc: func() (*rest.Config, error) { + return nil, rest.ErrNotInCluster + }, + expectClientCreated: false, + expectCurrentToken: "current-token", + }, + { + name: "unexpected error", + currentAuthToken: "current-token", + inClusterConfigFunc: func() (*rest.Config, error) { + return nil, fmt.Errorf("unexpected error") + }, + expectedErr: fmt.Errorf("unexpected error when reading in cluster config: unexpected error"), + expectClientCreated: false, + expectCurrentToken: "current-token", + }, + { + name: "first token - client created successfully", + currentAuthToken: "", + inClusterConfigFunc: func() (*rest.Config, error) { + return &rest.Config{BearerToken: "first-token"}, nil + }, + createPrometheusClientFunc: func(url, token string) (promapi.Client, *http.Transport, error) { + return &mockPrometheusClient{name: "first-client"}, &http.Transport{}, nil + }, + expectClientCreated: true, + expectCurrentToken: "first-token", + }, + { + name: "token changed with previous transport - clears previous transport", + currentAuthToken: "old-token", + inClusterConfigFunc: func() (*rest.Config, error) { + return &rest.Config{BearerToken: "new-token"}, nil + }, + createPrometheusClientFunc: func(url, token string) (promapi.Client, *http.Transport, error) { + return &mockPrometheusClient{name: "new-client"}, &http.Transport{}, nil + }, + expectClientCreated: true, + expectCurrentToken: "new-token", + expectPreviousTransportCleared: true, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + for _, setupMode := range []struct { + name string + init func(t *testing.T) *inClusterPromClientController + }{ + { + name: "running with prom reconciler directly", + init: func(t *testing.T) *inClusterPromClientController { + return &inClusterPromClientController{ + currentPrometheusAuthToken: tc.currentAuthToken, + prometheusConfig: &api.Prometheus{ + URL: prometheusURL, + }, + inClusterConfig: tc.inClusterConfigFunc, + } + }, + }, + { + name: "running with full descheduler", + init: func(t *testing.T) *inClusterPromClientController { + ctx := context.Background() + prometheusConfig := &api.Prometheus{ + URL: prometheusURL, + } + deschedulerPolicy := &api.DeschedulerPolicy{ + MetricsProviders: []api.MetricsProvider{ + { + Source: api.PrometheusMetrics, + Prometheus: prometheusConfig, + }, + }, + } + _, descheduler, _, _ := initDescheduler(t, ctx, initFeatureGates(), deschedulerPolicy, nil, false) + + // 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 + }, + }, + } { + t.Run(setupMode.name, func(t *testing.T) { + ctrl := setupMode.init(t) + + // Set previous transport and client if test expects them to be cleared + if tc.expectPreviousTransportCleared { + ctrl.previousPrometheusClientTransport = &http.Transport{} + } + if tc.expectPromClientCleared { + ctrl.promClient = &mockPrometheusClient{name: "old-client"} + } + + // Mock createPrometheusClient + clientCreated := false + if tc.createPrometheusClientFunc != nil { + ctrl.createPrometheusClient = func(url, token string) (promapi.Client, *http.Transport, error) { + client, transport, err := tc.createPrometheusClientFunc(url, token) + if err == nil { + clientCreated = true + } + return client, transport, err + } + } + + // Call reconcileInClusterSAToken + err := ctrl.reconcileInClusterSAToken() + + // Verify error expectations + if tc.expectedErr != nil { + if err == nil { + t.Errorf("Expected error %q but got none", tc.expectedErr) + } else if err.Error() != tc.expectedErr.Error() { + t.Errorf("Expected error %q but got %q", tc.expectedErr, err.Error()) + } + } else { + if err != nil { + t.Errorf("Expected no error but got: %v", err) + } + } + + // Verify client creation expectations + if tc.expectClientCreated && !clientCreated { + t.Errorf("Expected prometheus client to be created but it wasn't") + } + if !tc.expectClientCreated && clientCreated { + t.Errorf("Expected prometheus client not to be created but it was") + } + + // Verify token expectations + if ctrl.currentPrometheusAuthToken != tc.expectCurrentToken { + t.Errorf("Expected current token to be %q but got %q", tc.expectCurrentToken, ctrl.currentPrometheusAuthToken) + } + + // Verify previous transport cleared when expected + if tc.expectPreviousTransportCleared { + if tc.expectClientCreated { + // Success case: new transport should be set + if ctrl.previousPrometheusClientTransport == nil { + t.Error("Expected previous transport to be set to new transport, but it was nil") + } + } else if tc.expectedErr != nil { + // Failure case: transport should be nil + if ctrl.previousPrometheusClientTransport != nil { + t.Error("Expected previous transport to be cleared on error, but it was set") + } + } + } + + // Verify promClient cleared when expected + if tc.expectPromClientCleared { + if ctrl.promClient != nil { + t.Error("Expected promClient to be cleared, but it was set") + } + } + }) + } + }) + } +}