From 3ae1742953a6ad35261882c8f7cb678505e6a6f6 Mon Sep 17 00:00:00 2001 From: Sanket Sudake Date: Sun, 11 Dec 2022 20:41:12 +0530 Subject: [PATCH] Executor user informer factory in executors in place of informers (#2666) * Use informerfactory across executor * Run function informer for poolpodcontroller if istio enabled * Use same namespace for secret as keda mqtriggers Signed-off-by: Sanket Sudake --- pkg/crd/client.go | 7 +++- pkg/executor/executor.go | 33 +++++---------- .../executortype/container/containermgr.go | 8 ++-- .../executortype/newdeploy/newdeploymgr.go | 13 +++--- .../newdeploy/newdeploymgr_test.go | 28 +++---------- .../executortype/poolmgr/funchandlers.go | 14 ++----- pkg/executor/executortype/poolmgr/gpm.go | 7 ++-- .../executortype/poolmgr/poolpodcontroller.go | 21 +++++----- .../poolmgr/poolpodcontroller_test.go | 40 ++++--------------- pkg/mqtrigger/scalermanager.go | 4 +- 10 files changed, 60 insertions(+), 115 deletions(-) diff --git a/pkg/crd/client.go b/pkg/crd/client.go index e77e33b3..080b9558 100644 --- a/pkg/crd/client.go +++ b/pkg/crd/client.go @@ -33,6 +33,7 @@ import ( metricsclient "k8s.io/metrics/pkg/client/clientset/versioned" "github.com/fission/fission/pkg/generated/clientset/versioned" + "github.com/fission/fission/pkg/utils" ) // GetKubernetesClient gets a kubernetes client using the kubeconfig file at the @@ -91,9 +92,13 @@ func MakeFissionClient() (versioned.Interface, kubernetes.Interface, apiextensio // WaitForCRDs does a timeout to check if CRDs have been installed func WaitForCRDs(ctx context.Context, logger *zap.Logger, fissionClient versioned.Interface) error { logger.Info("Waiting for CRDs to be installed") + defaultNs := utils.DefaultNSResolver().DefaultNamespace + if defaultNs == "" { + defaultNs = metav1.NamespaceDefault + } start := time.Now() for { - fi := fissionClient.CoreV1().Functions(metav1.NamespaceDefault) + fi := fissionClient.CoreV1().Functions(defaultNs) _, err := fi.List(ctx, metav1.ListOptions{}) if err != nil { time.Sleep(100 * time.Millisecond) diff --git a/pkg/executor/executor.go b/pkg/executor/executor.go index 9531e211..733c1447 100644 --- a/pkg/executor/executor.go +++ b/pkg/executor/executor.go @@ -42,7 +42,6 @@ import ( fetcherConfig "github.com/fission/fission/pkg/fetcher/config" "github.com/fission/fission/pkg/generated/clientset/versioned" genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" - finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1" "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/metrics" otelUtils "github.com/fission/fission/pkg/utils/otel" @@ -276,13 +275,9 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error { logger.Info("Starting executor", zap.String("instanceID", executorInstanceID)) - funcInformer := make(map[string]finformerv1.FunctionInformer, 0) - envInformer := make(map[string]finformerv1.EnvironmentInformer, 0) - + finformerFactory := make(map[string]genInformer.SharedInformerFactory, 0) for _, ns := range utils.DefaultNSResolver().FissionResourceNS { - factory := genInformer.NewFilteredSharedInformerFactory(fissionClient, time.Minute*30, ns, nil) - funcInformer[ns] = factory.Core().V1().Functions() - envInformer[ns] = factory.Core().V1().Environments() + finformerFactory[ns] = genInformer.NewFilteredSharedInformerFactory(fissionClient, time.Minute*30, ns, nil) } executorLabel, err := utils.GetInformerLabelByExecutor(fv1.ExecutorTypePoolmgr) @@ -294,7 +289,7 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error { logger, fissionClient, kubernetesClient, metricsClient, fetcherConfig, executorInstanceID, - funcInformer, envInformer, + finformerFactory, gpmInformerFactory, podSpecPatch) if err != nil { return errors.Wrap(err, "pool manager creation failed") @@ -309,7 +304,7 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error { logger, fissionClient, kubernetesClient, fetcherConfig, executorInstanceID, - funcInformer, envInformer, + finformerFactory, ndmInformerFactory, podSpecPatch) if err != nil { return errors.Wrap(err, "new deploy manager creation failed") @@ -323,7 +318,7 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error { cnm, err := container.MakeContainer( ctx, logger, fissionClient, kubernetesClient, - executorInstanceID, funcInformer, + executorInstanceID, finformerFactory, cnmInformerFactory) if err != nil { return errors.Wrap(err, "container manager creation failed") @@ -356,29 +351,23 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error { cms := cms.MakeConfigSecretController(ctx, logger, fissionClient, kubernetesClient, executorTypes, configMapInformer, secretInformer) fissionInformers := make([]k8sCache.SharedIndexInformer, 0) - for _, informer := range funcInformer { - fissionInformers = append(fissionInformers, informer.Informer()) - } - for _, informer := range envInformer { - fissionInformers = append(fissionInformers, informer.Informer()) - } for _, informer := range configMapInformer { fissionInformers = append(fissionInformers, informer) } for _, informer := range secretInformer { fissionInformers = append(fissionInformers, informer) } + for _, factory := range finformerFactory { + factory.Start(ctx.Done()) + } for _, informerFactory := range gpmInformerFactory { - fissionInformers = append(fissionInformers, informerFactory.Core().V1().Pods().Informer()) - fissionInformers = append(fissionInformers, informerFactory.Apps().V1().ReplicaSets().Informer()) + informerFactory.Start(ctx.Done()) } for _, informerFactory := range ndmInformerFactory { - fissionInformers = append(fissionInformers, informerFactory.Apps().V1().Deployments().Informer()) - fissionInformers = append(fissionInformers, informerFactory.Core().V1().Services().Informer()) + informerFactory.Start(ctx.Done()) } for _, informerFactory := range cnmInformerFactory { - fissionInformers = append(fissionInformers, informerFactory.Apps().V1().Deployments().Informer()) - fissionInformers = append(fissionInformers, informerFactory.Core().V1().Services().Informer()) + informerFactory.Start(ctx.Done()) } api, err := MakeExecutor(ctx, logger, cms, fissionClient, executorTypes, diff --git a/pkg/executor/executortype/container/containermgr.go b/pkg/executor/executortype/container/containermgr.go index 9013862f..1b931da3 100644 --- a/pkg/executor/executortype/container/containermgr.go +++ b/pkg/executor/executortype/container/containermgr.go @@ -49,7 +49,7 @@ import ( executorUtils "github.com/fission/fission/pkg/executor/util" hpautils "github.com/fission/fission/pkg/executor/util/hpa" "github.com/fission/fission/pkg/generated/clientset/versioned" - finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1" + genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" "github.com/fission/fission/pkg/throttler" "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/maps" @@ -98,7 +98,7 @@ func MakeContainer( fissionClient versioned.Interface, kubernetesClient kubernetes.Interface, instanceID string, - funcInformer map[string]finformerv1.FunctionInformer, + finformerFactory map[string]genInformer.SharedInformerFactory, cnmInformerFactory map[string]k8sInformers.SharedInformerFactory, ) (executortype.ExecutorType, error) { enableIstio := false @@ -139,8 +139,8 @@ func MakeContainer( caaf.svcLister[ns] = informerFactory.Core().V1().Services().Lister() caaf.svcListerSynced[ns] = informerFactory.Core().V1().Services().Informer().HasSynced } - for _, informer := range funcInformer { - informer.Informer().AddEventHandler(caaf.FuncInformerHandler(ctx)) + for _, factory := range finformerFactory { + factory.Core().V1().Functions().Informer().AddEventHandler(caaf.FuncInformerHandler(ctx)) } return caaf, nil } diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr.go b/pkg/executor/executortype/newdeploy/newdeploymgr.go index 2ab199e1..efc841bd 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr.go @@ -51,7 +51,7 @@ import ( hpautils "github.com/fission/fission/pkg/executor/util/hpa" fetcherConfig "github.com/fission/fission/pkg/fetcher/config" "github.com/fission/fission/pkg/generated/clientset/versioned" - finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1" + genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" "github.com/fission/fission/pkg/throttler" "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/maps" @@ -103,8 +103,7 @@ func MakeNewDeploy( kubernetesClient kubernetes.Interface, fetcherConfig *fetcherConfig.Config, instanceID string, - funcInformer map[string]finformerv1.FunctionInformer, - envInformer map[string]finformerv1.EnvironmentInformer, + finformerFactory map[string]genInformer.SharedInformerFactory, ndmInformerFactory map[string]k8sInformers.SharedInformerFactory, podSpecPatch *apiv1.PodSpec, ) (executortype.ExecutorType, error) { @@ -148,11 +147,11 @@ func MakeNewDeploy( nd.svcLister[ns] = informerFactory.Core().V1().Services().Lister() nd.svcListerSynced[ns] = informerFactory.Core().V1().Services().Informer().HasSynced } - for _, fnInformer := range funcInformer { - fnInformer.Informer().AddEventHandler(nd.FunctionEventHandlers(ctx)) + for _, factory := range finformerFactory { + factory.Core().V1().Functions().Informer().AddEventHandler(nd.FunctionEventHandlers(ctx)) } - for _, envInformer := range envInformer { - envInformer.Informer().AddEventHandler(nd.EnvEventHandlers(ctx)) + for _, factory := range finformerFactory { + factory.Core().V1().Environments().Informer().AddEventHandler(nd.EnvEventHandlers(ctx)) } return nd, nil } diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr_test.go b/pkg/executor/executortype/newdeploy/newdeploymgr_test.go index 4ebf65de..3d2eea05 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr_test.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr_test.go @@ -21,7 +21,6 @@ import ( fetcherConfig "github.com/fission/fission/pkg/fetcher/config" fClient "github.com/fission/fission/pkg/generated/clientset/versioned/fake" genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" - finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1" "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/loggerfactory" ) @@ -35,25 +34,13 @@ const ( configmapName string = "newdeploy-test-configmap" ) -func runInformers(ctx context.Context, informers []k8sCache.SharedIndexInformer) { - // Run all informers - for _, informer := range informers { - go informer.Run(ctx.Done()) - } -} - func TestRefreshFuncPods(t *testing.T) { os.Setenv("DEBUG_ENV", "true") logger := loggerfactory.GetLogger() kubernetesClient := fake.NewSimpleClientset() fissionClient := fClient.NewSimpleClientset() - informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) - funcInformer := map[string]finformerv1.FunctionInformer{ - metav1.NamespaceAll: informerFactory.Core().V1().Functions(), - } - envInformer := map[string]finformerv1.EnvironmentInformer{ - metav1.NamespaceAll: informerFactory.Core().V1().Environments(), - } + factory := make(map[string]genInformer.SharedInformerFactory, 0) + factory[metav1.NamespaceAll] = genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) executorLabel, err := utils.GetInformerLabelByExecutor(fv1.ExecutorTypeNewdeploy) if err != nil { @@ -70,7 +57,7 @@ func TestRefreshFuncPods(t *testing.T) { } executor, err := MakeNewDeploy(ctx, logger, fissionClient, kubernetesClient, fetcherConfig, "test", - funcInformer, envInformer, ndmInformerFactory, nil) + factory, ndmInformerFactory, nil) if err != nil { t.Fatalf("new deploy manager creation failed: %s", err) } @@ -87,16 +74,13 @@ func TestRefreshFuncPods(t *testing.T) { go ndm.Run(ctx) t.Log("New deploy manager started") - informer := []k8sCache.SharedIndexInformer{ - envInformer[metav1.NamespaceAll].Informer(), - funcInformer[metav1.NamespaceAll].Informer(), + for _, f := range factory { + f.Start(ctx.Done()) } for _, informerFactory := range ndmInformerFactory { - informer = append(informer, informerFactory.Apps().V1().Deployments().Informer()) - informer = append(informer, informerFactory.Core().V1().Services().Informer()) + informerFactory.Start(ctx.Done()) } - runInformers(ctx, informer) t.Log("Informers required for new deploy manager started") waitSynced := make([]k8sCache.InformerSynced, 0) diff --git a/pkg/executor/executortype/poolmgr/funchandlers.go b/pkg/executor/executortype/poolmgr/funchandlers.go index 53c314f6..d4015524 100644 --- a/pkg/executor/executortype/poolmgr/funchandlers.go +++ b/pkg/executor/executortype/poolmgr/funchandlers.go @@ -59,12 +59,6 @@ func FunctionEventHandlers(ctx context.Context, logger *zap.Logger, kubernetesCl return } - // create or update role-binding - envNs := fissionfnNamespace - if fn.Spec.Environment.Namespace != metav1.NamespaceDefault { - envNs = fn.Spec.Environment.Namespace - } - if istioEnabled { // create a same name service for function // since istio only allows the traffic to service @@ -74,6 +68,7 @@ func FunctionEventHandlers(ctx context.Context, logger *zap.Logger, kubernetesCl } svcName := utils.GetFunctionIstioServiceName(fn.ObjectMeta.Name, fn.ObjectMeta.Namespace) + envNs := utils.DefaultNSResolver().GetFunctionNS(fn.Spec.Environment.Namespace) // service for accepting user traffic svc := apiv1.Service{ @@ -124,12 +119,9 @@ func FunctionEventHandlers(ctx context.Context, logger *zap.Logger, kubernetesCl return } - envNs := fissionfnNamespace - if fn.Spec.Environment.Namespace != metav1.NamespaceDefault { - envNs = fn.Spec.Environment.Namespace - } - if istioEnabled { + envNs := utils.DefaultNSResolver().GetFunctionNS(fn.Spec.Environment.Namespace) + svcName := utils.GetFunctionIstioServiceName(fn.ObjectMeta.Name, fn.ObjectMeta.Namespace) // delete function istio service err := kubernetesClient.CoreV1().Services(envNs).Delete(ctx, svcName, metav1.DeleteOptions{}) diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index 92079e0d..33ef7de1 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -50,7 +50,7 @@ import ( executorUtils "github.com/fission/fission/pkg/executor/util" fetcherConfig "github.com/fission/fission/pkg/fetcher/config" "github.com/fission/fission/pkg/generated/clientset/versioned" - finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1" + genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" "github.com/fission/fission/pkg/utils" otelUtils "github.com/fission/fission/pkg/utils/otel" ) @@ -117,8 +117,7 @@ func MakeGenericPoolManager(ctx context.Context, metricsClient metricsclient.Interface, fetcherConfig *fetcherConfig.Config, instanceID string, - funcInformer map[string]finformerv1.FunctionInformer, - envInformer map[string]finformerv1.EnvironmentInformer, + finformerFactory map[string]genInformer.SharedInformerFactory, gpmInformerFactory map[string]k8sInformers.SharedInformerFactory, podSpecPatch *apiv1.PodSpec, ) (executortype.ExecutorType, error) { @@ -135,7 +134,7 @@ func MakeGenericPoolManager(ctx context.Context, } poolPodC := NewPoolPodController(ctx, gpmLogger, kubernetesClient, - enableIstio, funcInformer, envInformer, gpmInformerFactory) + enableIstio, finformerFactory, gpmInformerFactory) gpm := &GenericPoolManager{ logger: gpmLogger, diff --git a/pkg/executor/executortype/poolmgr/poolpodcontroller.go b/pkg/executor/executortype/poolmgr/poolpodcontroller.go index 064a6b58..79243a46 100644 --- a/pkg/executor/executortype/poolmgr/poolpodcontroller.go +++ b/pkg/executor/executortype/poolmgr/poolpodcontroller.go @@ -37,7 +37,7 @@ import ( fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/executor/fscache" - finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1" + genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" flisterv1 "github.com/fission/fission/pkg/generated/listers/core/v1" "github.com/fission/fission/pkg/utils" ) @@ -70,8 +70,7 @@ type ( func NewPoolPodController(ctx context.Context, logger *zap.Logger, kubernetesClient kubernetes.Interface, enableIstio bool, - funcInformer map[string]finformerv1.FunctionInformer, - envInformer map[string]finformerv1.EnvironmentInformer, + finformerFactory map[string]genInformer.SharedInformerFactory, gpmInformerFactory map[string]k8sInformers.SharedInformerFactory) *PoolPodController { logger = logger.Named("pool_pod_controller") p := &PoolPodController{ @@ -87,17 +86,19 @@ func NewPoolPodController(ctx context.Context, logger *zap.Logger, envDeleteQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "EnvDeleteQueue"), spCleanupPodQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "SpecializedPodCleanupQueue"), } - for _, informer := range funcInformer { - informer.Informer().AddEventHandler(FunctionEventHandlers(ctx, p.logger, p.kubernetesClient, p.nsResolver.ResolveNamespace(p.nsResolver.FunctionNamespace), p.enableIstio)) + if p.enableIstio { + for _, factory := range finformerFactory { + factory.Core().V1().Functions().Informer().AddEventHandler(FunctionEventHandlers(ctx, p.logger, p.kubernetesClient, p.nsResolver.ResolveNamespace(p.nsResolver.FunctionNamespace), p.enableIstio)) + } } - for ns, informer := range envInformer { - informer.Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ + for ns, informer := range finformerFactory { + informer.Core().V1().Environments().Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ AddFunc: p.enqueueEnvAdd, UpdateFunc: p.enqueueEnvUpdate, DeleteFunc: p.enqueueEnvDelete, }) - p.envLister[ns] = informer.Lister() - p.envListerSynced[ns] = informer.Informer().HasSynced + p.envLister[ns] = informer.Core().V1().Environments().Lister() + p.envListerSynced[ns] = informer.Core().V1().Environments().Informer().HasSynced } for ns, informerFactory := range gpmInformerFactory { informerFactory.Apps().V1().ReplicaSets().Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ @@ -105,8 +106,8 @@ func NewPoolPodController(ctx context.Context, logger *zap.Logger, UpdateFunc: p.handleRSUpdate, DeleteFunc: p.handleRSDelete, }) - p.podLister[ns] = informerFactory.Core().V1().Pods().Lister() p.podListerSynced[ns] = informerFactory.Core().V1().Pods().Informer().HasSynced + p.podLister[ns] = informerFactory.Core().V1().Pods().Lister() } p.logger.Info("pool pod controller handlers registered") diff --git a/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go b/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go index ad9dd33c..2b5f126d 100644 --- a/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go +++ b/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go @@ -32,34 +32,18 @@ import ( fetcherConfig "github.com/fission/fission/pkg/fetcher/config" fClient "github.com/fission/fission/pkg/generated/clientset/versioned/fake" genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" - finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1" "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/loggerfactory" ) -func runInformers(ctx context.Context, informers []k8sCache.SharedIndexInformer) { - // Run all informers - for _, informer := range informers { - go informer.Run(ctx.Done()) - } -} - func TestPoolPodControllerPodCleanup(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() logger := loggerfactory.GetLogger() kubernetesClient := fake.NewSimpleClientset() fissionClient := fClient.NewSimpleClientset() - informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) - funcInformer := map[string]finformerv1.FunctionInformer{ - metav1.NamespaceAll: informerFactory.Core().V1().Functions(), - } - pkgInformer := map[string]finformerv1.PackageInformer{ - metav1.NamespaceAll: informerFactory.Core().V1().Packages(), - } - envInformer := map[string]finformerv1.EnvironmentInformer{ - metav1.NamespaceAll: informerFactory.Core().V1().Environments(), - } + factory := make(map[string]genInformer.SharedInformerFactory, 0) + factory[metav1.NamespaceDefault] = genInformer.NewFilteredSharedInformerFactory(fissionClient, time.Minute*30, metav1.NamespaceDefault, nil) executorLabel, err := utils.GetInformerLabelByExecutor(fv1.ExecutorTypePoolmgr) if err != nil { @@ -68,9 +52,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) { gpmInformerFactory := utils.GetInformerFactoryByExecutor(kubernetesClient, executorLabel, time.Minute*30) ppc := NewPoolPodController(ctx, logger, kubernetesClient, false, - funcInformer, - envInformer, - gpmInformerFactory) + factory, gpmInformerFactory) executorInstanceID := strings.ToLower(uniuri.NewLen(8)) metricsClient := metricsclient.NewSimpleClientset() @@ -82,8 +64,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) { logger, fissionClient, kubernetesClient, metricsClient, fetcherConfig, executorInstanceID, - funcInformer, envInformer, - gpmInformerFactory, nil) + factory, gpmInformerFactory, nil) if err != nil { t.Fatalf("Error creating generic pool manager: %v", err) } @@ -92,23 +73,18 @@ func TestPoolPodControllerPodCleanup(t *testing.T) { go ppc.Run(ctx, ctx.Done()) - informers := []k8sCache.SharedIndexInformer{ - funcInformer[metav1.NamespaceAll].Informer(), - pkgInformer[metav1.NamespaceAll].Informer(), - envInformer[metav1.NamespaceAll].Informer(), + for _, f := range factory { + f.Start(ctx.Done()) } for _, informerFactory := range gpmInformerFactory { - informers = append(informers, informerFactory.Core().V1().Pods().Informer()) - informers = append(informers, informerFactory.Apps().V1().ReplicaSets().Informer()) + informerFactory.Start(ctx.Done()) } - runInformers(ctx, informers) - pod := &corev1.Pod{ ObjectMeta: metav1.ObjectMeta{ Name: "test-pod", - Namespace: "test-different-namespace", + Namespace: metav1.NamespaceDefault, }, Status: corev1.PodStatus{ Phase: corev1.PodRunning, diff --git a/pkg/mqtrigger/scalermanager.go b/pkg/mqtrigger/scalermanager.go index 5d6b5d80..c83062e8 100644 --- a/pkg/mqtrigger/scalermanager.go +++ b/pkg/mqtrigger/scalermanager.go @@ -217,7 +217,7 @@ func getEnvVarlist(ctx context.Context, mqt *fv1.MessageQueueTrigger, routerURL // Add Auth Fields secretName := mqt.Spec.Secret if len(secretName) > 0 { - secret, err := kubeClient.CoreV1().Secrets(apiv1.NamespaceDefault).Get(ctx, secretName, metav1.GetOptions{}) + secret, err := kubeClient.CoreV1().Secrets(mqt.Namespace).Get(ctx, secretName, metav1.GetOptions{}) if err != nil { return nil, err } @@ -304,7 +304,7 @@ func checkAndUpdateTriggerFields(mqt, newMqt *fv1.MessageQueueTrigger) bool { } func getAuthTriggerSpec(ctx context.Context, mqt *fv1.MessageQueueTrigger, authenticationRef string, kubeClient kubernetes.Interface) (*unstructured.Unstructured, error) { - secret, err := kubeClient.CoreV1().Secrets(apiv1.NamespaceDefault).Get(ctx, mqt.Spec.Secret, metav1.GetOptions{}) + secret, err := kubeClient.CoreV1().Secrets(mqt.Namespace).Get(ctx, mqt.Spec.Secret, metav1.GetOptions{}) if err != nil { return nil, err }