mirror of
https://github.com/kubernetes-sigs/descheduler.git
synced 2026-08-23 22:46:35 +00:00
refactor: move prometheus client controller related code under a seperate file
This commit is contained in:
@@ -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")
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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) },
|
||||
}
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user