From 691feaa84f81deafa41a1ea9a908393af893209c Mon Sep 17 00:00:00 2001 From: Shubham Bansal <62992590+shubham-bansal96@users.noreply.github.com> Date: Sun, 4 Dec 2022 21:01:31 +0530 Subject: [PATCH] K8s informer to work with specific namespaces for executor (#2651) Consider specific namespaces mentioned by the user in building informers in the executor - Confimaps - Secrets - Deployments - Services - Pods - Replicasets --- pkg/executor/cms/cmscontroller.go | 10 ++-- pkg/executor/executor.go | 58 ++++++++----------- .../executortype/container/containermgr.go | 44 ++++++++------ .../executortype/newdeploy/newdeploymgr.go | 46 +++++++++------ .../newdeploy/newdeploymgr_test.go | 35 +++++++---- pkg/executor/executortype/poolmgr/gpm.go | 28 +++++---- .../executortype/poolmgr/poolpodcontroller.go | 39 +++++++------ .../poolmgr/poolpodcontroller_test.go | 29 +++++----- pkg/utils/informer.go | 26 ++++++--- 9 files changed, 180 insertions(+), 135 deletions(-) diff --git a/pkg/executor/cms/cmscontroller.go b/pkg/executor/cms/cmscontroller.go index f5ce9631..06691eca 100644 --- a/pkg/executor/cms/cmscontroller.go +++ b/pkg/executor/cms/cmscontroller.go @@ -21,8 +21,8 @@ import ( "github.com/pkg/errors" "go.uber.org/zap" - informerv1 "k8s.io/client-go/informers/core/v1" "k8s.io/client-go/kubernetes" + "k8s.io/client-go/tools/cache" fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/executor/executortype" @@ -41,18 +41,18 @@ type ( // MakeConfigSecretController makes a controller for configmaps and secrets which changes related functions func MakeConfigSecretController(ctx context.Context, logger *zap.Logger, fissionClient versioned.Interface, kubernetesClient kubernetes.Interface, types map[fv1.ExecutorType]executortype.ExecutorType, - configmapInformer map[string]informerv1.ConfigMapInformer, - secretInformer map[string]informerv1.SecretInformer) *ConfigSecretController { + configmapInformer, + secretInformer map[string]cache.SharedIndexInformer) *ConfigSecretController { logger.Debug("Creating ConfigMap & Secret Controller") cmsController := &ConfigSecretController{ logger: logger, fissionClient: fissionClient, } for _, informer := range configmapInformer { - informer.Informer().AddEventHandler(ConfigMapEventHandlers(ctx, logger, fissionClient, kubernetesClient, types)) + informer.AddEventHandler(ConfigMapEventHandlers(ctx, logger, fissionClient, kubernetesClient, types)) } for _, informer := range secretInformer { - informer.Informer().AddEventHandler(SecretEventHandlers(ctx, logger, fissionClient, kubernetesClient, types)) + informer.AddEventHandler(SecretEventHandlers(ctx, logger, fissionClient, kubernetesClient, types)) } return cmsController diff --git a/pkg/executor/executor.go b/pkg/executor/executor.go index f677bb0e..4b174ebb 100644 --- a/pkg/executor/executor.go +++ b/pkg/executor/executor.go @@ -29,8 +29,6 @@ import ( "github.com/pkg/errors" "go.uber.org/zap" apiv1 "k8s.io/api/core/v1" - k8sInformers "k8s.io/client-go/informers" - k8sInformersv1 "k8s.io/client-go/informers/core/v1" k8sCache "k8s.io/client-go/tools/cache" fv1 "github.com/fission/fission/pkg/apis/core/v1" @@ -297,49 +295,46 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error { pkgInformer[ns] = factory.Core().V1().Packages() } - gpmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypePoolmgr, time.Minute*30) + executorLabel, err := utils.GetInformerLabelByExecutor(fv1.ExecutorTypePoolmgr) if err != nil { return err } - gpmPodInformer := gpmInformerFactory.Core().V1().Pods() - gpmRsInformer := gpmInformerFactory.Apps().V1().ReplicaSets() + gpmInformerFactory := utils.GetInformerFactoryByExecutor(kubernetesClient, executorLabel, time.Minute*30) gpm, err := poolmgr.MakeGenericPoolManager(ctx, logger, fissionClient, kubernetesClient, metricsClient, fetcherConfig, executorInstanceID, funcInformer, pkgInformer, envInformer, - gpmPodInformer, gpmRsInformer, podSpecPatch) + gpmInformerFactory, podSpecPatch) if err != nil { return errors.Wrap(err, "pool manager creation failed") } - ndmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypeNewdeploy, time.Minute*30) + executorLabel, err = utils.GetInformerLabelByExecutor(fv1.ExecutorTypeNewdeploy) if err != nil { return err } - ndmDeplInformer := ndmInformerFactory.Apps().V1().Deployments() - ndmSvcInformer := ndmInformerFactory.Core().V1().Services() + ndmInformerFactory := utils.GetInformerFactoryByExecutor(kubernetesClient, executorLabel, time.Minute*30) ndm, err := newdeploy.MakeNewDeploy(ctx, logger, fissionClient, kubernetesClient, fetcherConfig, executorInstanceID, funcInformer, envInformer, - ndmDeplInformer, ndmSvcInformer, podSpecPatch) + ndmInformerFactory, podSpecPatch) if err != nil { return errors.Wrap(err, "new deploy manager creation failed") } - cnmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypeContainer, time.Minute*30) + executorLabel, err = utils.GetInformerLabelByExecutor(fv1.ExecutorTypeContainer) if err != nil { return err } - cnmDeplInformer := cnmInformerFactory.Apps().V1().Deployments() - cnmSvcInformer := cnmInformerFactory.Core().V1().Services() + cnmInformerFactory := utils.GetInformerFactoryByExecutor(kubernetesClient, executorLabel, time.Minute*30) cnm, err := container.MakeContainer( ctx, logger, fissionClient, kubernetesClient, executorInstanceID, funcInformer, - cnmDeplInformer, cnmSvcInformer) + cnmInformerFactory) if err != nil { return errors.Wrap(err, "container manager creation failed") } @@ -366,15 +361,8 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error { // TODO: use context to control the waiting time once kubernetes client supports it. util.WaitTimeout(wg, 30*time.Second) - configMapInformer := make(map[string]k8sInformersv1.ConfigMapInformer, 0) - secretInformer := make(map[string]k8sInformersv1.SecretInformer, 0) - - for _, ns := range utils.DefaultNSResolver().FissionResourceNS { - factory := k8sInformers.NewFilteredSharedInformerFactory(kubernetesClient, time.Minute*30, ns, nil) - configMapInformer[ns] = factory.Core().V1().ConfigMaps() - secretInformer[ns] = factory.Core().V1().Secrets() - } - + configMapInformer := utils.GetK8sInformersForNamespaces(kubernetesClient, time.Minute*30, fv1.ConfigMaps) + secretInformer := utils.GetK8sInformersForNamespaces(kubernetesClient, time.Minute*30, fv1.Secrets) cms := cms.MakeConfigSecretController(ctx, logger, fissionClient, kubernetesClient, executorTypes, configMapInformer, secretInformer) fissionInformers := make([]k8sCache.SharedIndexInformer, 0) @@ -388,20 +376,24 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error { fissionInformers = append(fissionInformers, informer.Informer()) } for _, informer := range configMapInformer { - fissionInformers = append(fissionInformers, informer.Informer()) + fissionInformers = append(fissionInformers, informer) } for _, informer := range secretInformer { - fissionInformers = append(fissionInformers, informer.Informer()) + fissionInformers = append(fissionInformers, informer) + } + for _, informerFactory := range gpmInformerFactory { + fissionInformers = append(fissionInformers, informerFactory.Core().V1().Pods().Informer()) + fissionInformers = append(fissionInformers, informerFactory.Apps().V1().ReplicaSets().Informer()) + } + for _, informerFactory := range ndmInformerFactory { + fissionInformers = append(fissionInformers, informerFactory.Apps().V1().Deployments().Informer()) + fissionInformers = append(fissionInformers, informerFactory.Core().V1().Services().Informer()) + } + for _, informerFactory := range cnmInformerFactory { + fissionInformers = append(fissionInformers, informerFactory.Apps().V1().Deployments().Informer()) + fissionInformers = append(fissionInformers, informerFactory.Core().V1().Services().Informer()) } - fissionInformers = append(fissionInformers, - gpmPodInformer.Informer(), - gpmRsInformer.Informer(), - ndmDeplInformer.Informer(), - ndmSvcInformer.Informer(), - cnmDeplInformer.Informer(), - cnmSvcInformer.Informer(), - ) api, err := MakeExecutor(ctx, logger, cms, fissionClient, executorTypes, fissionInformers..., ) diff --git a/pkg/executor/executortype/container/containermgr.go b/pkg/executor/executortype/container/containermgr.go index dd815231..9013862f 100644 --- a/pkg/executor/executortype/container/containermgr.go +++ b/pkg/executor/executortype/container/containermgr.go @@ -35,8 +35,7 @@ import ( "k8s.io/apimachinery/pkg/labels" k8sTypes "k8s.io/apimachinery/pkg/types" "k8s.io/apimachinery/pkg/util/wait" - appsinformers "k8s.io/client-go/informers/apps/v1" - coreinformers "k8s.io/client-go/informers/core/v1" + k8sInformers "k8s.io/client-go/informers" "k8s.io/client-go/kubernetes" appslisters "k8s.io/client-go/listers/apps/v1" corelisters "k8s.io/client-go/listers/core/v1" @@ -81,11 +80,11 @@ type ( defaultIdlePodReapTime time.Duration - deplLister appslisters.DeploymentLister - svcLister corelisters.ServiceLister + deplLister map[string]appslisters.DeploymentLister + svcLister map[string]corelisters.ServiceLister - deplListerSynced k8sCache.InformerSynced - svcListerSynced k8sCache.InformerSynced + deplListerSynced map[string]k8sCache.InformerSynced + svcListerSynced map[string]k8sCache.InformerSynced hpaops *hpautils.HpaOperations objectReaperIntervalSecond time.Duration @@ -100,8 +99,7 @@ func MakeContainer( kubernetesClient kubernetes.Interface, instanceID string, funcInformer map[string]finformerv1.FunctionInformer, - deplInformer appsinformers.DeploymentInformer, - svcInformer coreinformers.ServiceInformer, + cnmInformerFactory map[string]k8sInformers.SharedInformerFactory, ) (executortype.ExecutorType, error) { enableIstio := false if len(os.Getenv("ENABLE_ISTIO")) > 0 { @@ -129,14 +127,18 @@ func MakeContainer( defaultIdlePodReapTime: 1 * time.Minute, objectReaperIntervalSecond: time.Duration(executorUtils.GetObjectReaperInterval(logger, fv1.ExecutorTypeContainer, 5)) * time.Second, hpaops: hpautils.NewHpaOperations(logger, kubernetesClient, instanceID), + deplLister: make(map[string]appslisters.DeploymentLister), + deplListerSynced: make(map[string]k8sCache.InformerSynced), + svcLister: make(map[string]corelisters.ServiceLister), + svcListerSynced: make(map[string]k8sCache.InformerSynced), } - caaf.deplLister = deplInformer.Lister() - caaf.deplListerSynced = deplInformer.Informer().HasSynced - - caaf.svcLister = svcInformer.Lister() - caaf.svcListerSynced = svcInformer.Informer().HasSynced - + for ns, informerFactory := range cnmInformerFactory { + caaf.deplLister[ns] = informerFactory.Apps().V1().Deployments().Lister() + caaf.deplListerSynced[ns] = informerFactory.Apps().V1().Deployments().Informer().HasSynced + 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)) } @@ -145,7 +147,15 @@ func MakeContainer( // Run start the function along with an object reaper. func (caaf *Container) Run(ctx context.Context) { - if ok := k8sCache.WaitForCacheSync(ctx.Done(), caaf.deplListerSynced, caaf.svcListerSynced); !ok { + waitSynced := make([]k8sCache.InformerSynced, 0) + for _, deplListerSynced := range caaf.deplListerSynced { + waitSynced = append(waitSynced, deplListerSynced) + } + for _, svcListerSynced := range caaf.svcListerSynced { + waitSynced = append(waitSynced, svcListerSynced) + } + + if ok := k8sCache.WaitForCacheSync(ctx.Done(), waitSynced...); !ok { caaf.logger.Fatal("failed to wait for caches to sync") } go caaf.idleObjectReaper(ctx) @@ -214,7 +224,7 @@ func (caaf *Container) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) bool } for _, obj := range fsvc.KubernetesObjects { if strings.ToLower(obj.Kind) == "service" { - _, err := caaf.svcLister.Services(obj.Namespace).Get(obj.Name) + _, err := caaf.svcLister[obj.Namespace].Services(obj.Namespace).Get(obj.Name) if err != nil { if !k8sErrs.IsNotFound(err) { logger.Error("error validating function service", zap.String("function", fsvc.Function.Name), zap.Error(err)) @@ -222,7 +232,7 @@ func (caaf *Container) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) bool return false } } else if strings.ToLower(obj.Kind) == "deployment" { - currentDeploy, err := caaf.deplLister.Deployments(obj.Namespace).Get(obj.Name) + currentDeploy, err := caaf.deplLister[obj.Namespace].Deployments(obj.Namespace).Get(obj.Name) if err != nil { if !k8sErrs.IsNotFound(err) { logger.Error("error validating function deployment", zap.String("function", fsvc.Function.Name), zap.Error(err)) diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr.go b/pkg/executor/executortype/newdeploy/newdeploymgr.go index 3d706840..2ab199e1 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr.go @@ -36,8 +36,7 @@ import ( "k8s.io/apimachinery/pkg/labels" k8sTypes "k8s.io/apimachinery/pkg/types" "k8s.io/apimachinery/pkg/util/wait" - appsinformers "k8s.io/client-go/informers/apps/v1" - coreinformers "k8s.io/client-go/informers/core/v1" + k8sInformers "k8s.io/client-go/informers" "k8s.io/client-go/kubernetes" appslisters "k8s.io/client-go/listers/apps/v1" corelisters "k8s.io/client-go/listers/core/v1" @@ -83,11 +82,11 @@ type ( defaultIdlePodReapTime time.Duration - deplLister appslisters.DeploymentLister - svcLister corelisters.ServiceLister + deplLister map[string]appslisters.DeploymentLister + svcLister map[string]corelisters.ServiceLister - deplListerSynced k8sCache.InformerSynced - svcListerSynced k8sCache.InformerSynced + deplListerSynced map[string]k8sCache.InformerSynced + svcListerSynced map[string]k8sCache.InformerSynced hpaops *hpautils.HpaOperations @@ -106,8 +105,7 @@ func MakeNewDeploy( instanceID string, funcInformer map[string]finformerv1.FunctionInformer, envInformer map[string]finformerv1.EnvironmentInformer, - deplInformer appsinformers.DeploymentInformer, - svcInformer coreinformers.ServiceInformer, + ndmInformerFactory map[string]k8sInformers.SharedInformerFactory, podSpecPatch *apiv1.PodSpec, ) (executortype.ExecutorType, error) { enableIstio := false @@ -137,15 +135,19 @@ func MakeNewDeploy( objectReaperIntervalSecond: time.Duration(executorUtils.GetObjectReaperInterval(logger, fv1.ExecutorTypeNewdeploy, 5)) * time.Second, hpaops: hpautils.NewHpaOperations(logger, kubernetesClient, instanceID), - podSpecPatch: podSpecPatch, + podSpecPatch: podSpecPatch, + deplLister: make(map[string]appslisters.DeploymentLister), + deplListerSynced: make(map[string]k8sCache.InformerSynced), + svcLister: make(map[string]corelisters.ServiceLister), + svcListerSynced: make(map[string]k8sCache.InformerSynced), } - nd.deplLister = deplInformer.Lister() - nd.deplListerSynced = deplInformer.Informer().HasSynced - - nd.svcLister = svcInformer.Lister() - nd.svcListerSynced = svcInformer.Informer().HasSynced - + for ns, informerFactory := range ndmInformerFactory { + nd.deplLister[ns] = informerFactory.Apps().V1().Deployments().Lister() + nd.deplListerSynced[ns] = informerFactory.Apps().V1().Deployments().Informer().HasSynced + 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)) } @@ -157,7 +159,15 @@ func MakeNewDeploy( // Run start the function and environment controller along with an object reaper. func (deploy *NewDeploy) Run(ctx context.Context) { - if ok := k8sCache.WaitForCacheSync(ctx.Done(), deploy.deplListerSynced, deploy.svcListerSynced); !ok { + waitSynced := make([]k8sCache.InformerSynced, 0) + for _, deplListerSynced := range deploy.deplListerSynced { + waitSynced = append(waitSynced, deplListerSynced) + } + for _, svcListerSynced := range deploy.svcListerSynced { + waitSynced = append(waitSynced, svcListerSynced) + } + + if ok := k8sCache.WaitForCacheSync(ctx.Done(), waitSynced...); !ok { deploy.logger.Fatal("failed to wait for caches to sync") } go deploy.idleObjectReaper(ctx) @@ -222,7 +232,7 @@ func (deploy *NewDeploy) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) boo } for _, obj := range fsvc.KubernetesObjects { if strings.ToLower(obj.Kind) == "service" { - _, err := deploy.svcLister.Services(obj.Namespace).Get(obj.Name) + _, err := deploy.svcLister[obj.Namespace].Services(obj.Namespace).Get(obj.Name) if err != nil { if !k8sErrs.IsNotFound(err) { logger.Error("error validating function service", zap.String("function", fsvc.Function.Name), zap.Error(err)) @@ -231,7 +241,7 @@ func (deploy *NewDeploy) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) boo } } else if strings.ToLower(obj.Kind) == "deployment" { - currentDeploy, err := deploy.deplLister.Deployments(obj.Namespace).Get(obj.Name) + currentDeploy, err := deploy.deplLister[obj.Namespace].Deployments(obj.Namespace).Get(obj.Name) if err != nil { if !k8sErrs.IsNotFound(err) { logger.Error("error validating function deployment", zap.String("function", fsvc.Function.Name), zap.Error(err)) diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr_test.go b/pkg/executor/executortype/newdeploy/newdeploymgr_test.go index 4bef4476..39ddefed 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr_test.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr_test.go @@ -55,13 +55,12 @@ func TestRefreshFuncPods(t *testing.T) { envInformer := map[string]finformerv1.EnvironmentInformer{ metav1.NamespaceAll: informerFactory.Core().V1().Environments(), } - newDeployInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypeNewdeploy, time.Minute*30) - if err != nil { - t.Fatalf("Error creating informer factory: %s", err) - } - deployInformer := newDeployInformerFactory.Apps().V1().Deployments() - svcInformer := newDeployInformerFactory.Core().V1().Services() + executorLabel, err := utils.GetInformerLabelByExecutor(fv1.ExecutorTypeNewdeploy) + if err != nil { + t.Fatalf("Error creating labels for informer: %s", err) + } + ndmInformerFactory := utils.GetInformerFactoryByExecutor(kubernetesClient, executorLabel, time.Minute*30) ctx, cancel := context.WithCancel(context.Background()) defer cancel() @@ -82,7 +81,7 @@ func TestRefreshFuncPods(t *testing.T) { } executor, err := MakeNewDeploy(ctx, logger, fissionClient, kubernetesClient, fetcherConfig, "test", - funcInformer, envInformer, deployInformer, svcInformer, podSpecPatch) + funcInformer, envInformer, ndmInformerFactory, podSpecPatch) if err != nil { t.Fatalf("new deploy manager creation failed: %s", err) } @@ -99,15 +98,27 @@ func TestRefreshFuncPods(t *testing.T) { go ndm.Run(ctx) t.Log("New deploy manager started") - runInformers(ctx, []k8sCache.SharedIndexInformer{ + informer := []k8sCache.SharedIndexInformer{ envInformer[metav1.NamespaceAll].Informer(), funcInformer[metav1.NamespaceAll].Informer(), - deployInformer.Informer(), - svcInformer.Informer(), - }) + } + for _, informerFactory := range ndmInformerFactory { + informer = append(informer, informerFactory.Apps().V1().Deployments().Informer()) + informer = append(informer, informerFactory.Core().V1().Services().Informer()) + } + + runInformers(ctx, informer) t.Log("Informers required for new deploy manager started") - if ok := k8sCache.WaitForCacheSync(ctx.Done(), ndm.deplListerSynced, ndm.svcListerSynced); !ok { + waitSynced := make([]k8sCache.InformerSynced, 0) + for _, deplListerSynced := range ndm.deplListerSynced { + waitSynced = append(waitSynced, deplListerSynced) + } + for _, svcListerSynced := range ndm.svcListerSynced { + waitSynced = append(waitSynced, svcListerSynced) + } + + if ok := k8sCache.WaitForCacheSync(ctx.Done(), waitSynced...); !ok { t.Fatal("Timed out waiting for caches to sync") } diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index 6d06896b..3cc48157 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -37,8 +37,7 @@ import ( k8sTypes "k8s.io/apimachinery/pkg/types" "k8s.io/apimachinery/pkg/util/wait" "k8s.io/apimachinery/pkg/watch" - appsinformers "k8s.io/client-go/informers/apps/v1" - coreinformers "k8s.io/client-go/informers/core/v1" + k8sInformers "k8s.io/client-go/informers" "k8s.io/client-go/kubernetes" corelisters "k8s.io/client-go/listers/core/v1" k8sCache "k8s.io/client-go/tools/cache" @@ -88,10 +87,10 @@ type ( fetcherConfig *fetcherConfig.Config // podLister can list/get pods from the shared informer's store - podLister corelisters.PodLister + podLister map[string]corelisters.PodLister // podListerSynced returns true if the pod store has been synced at least once. - podListerSynced k8sCache.InformerSynced + podListerSynced map[string]k8sCache.InformerSynced defaultIdlePodReapTime time.Duration @@ -123,8 +122,7 @@ func MakeGenericPoolManager(ctx context.Context, funcInformer map[string]finformerv1.FunctionInformer, pkgInformer map[string]finformerv1.PackageInformer, envInformer map[string]finformerv1.EnvironmentInformer, - podInformer coreinformers.PodInformer, - rsInformer appsinformers.ReplicaSetInformer, + gpmInformerFactory map[string]k8sInformers.SharedInformerFactory, podSpecPatch *apiv1.PodSpec, ) (executortype.ExecutorType, error) { @@ -140,7 +138,7 @@ func MakeGenericPoolManager(ctx context.Context, } poolPodC := NewPoolPodController(ctx, gpmLogger, kubernetesClient, - enableIstio, funcInformer, pkgInformer, envInformer, rsInformer, podInformer) + enableIstio, funcInformer, pkgInformer, envInformer, gpmInformerFactory) gpm := &GenericPoolManager{ logger: gpmLogger, @@ -159,9 +157,13 @@ func MakeGenericPoolManager(ctx context.Context, poolPodC: poolPodC, podSpecPatch: podSpecPatch, objectReaperIntervalSecond: time.Duration(executorUtils.GetObjectReaperInterval(logger, fv1.ExecutorTypePoolmgr, 5)) * time.Second, + podLister: make(map[string]corelisters.PodLister), + podListerSynced: make(map[string]k8sCache.InformerSynced), + } + for ns, informerFactory := range gpmInformerFactory { + gpm.podLister[ns] = informerFactory.Core().V1().Pods().Lister() + gpm.podListerSynced[ns] = informerFactory.Core().V1().Pods().Informer().HasSynced } - gpm.podLister = podInformer.Lister() - gpm.podListerSynced = podInformer.Informer().HasSynced gpm.logger.Debug("inside MakeGenericPoolManager") @@ -169,8 +171,10 @@ func MakeGenericPoolManager(ctx context.Context, } func (gpm *GenericPoolManager) Run(ctx context.Context) { - if ok := k8sCache.WaitForCacheSync(ctx.Done(), gpm.podListerSynced); !ok { - gpm.logger.Fatal("failed to wait for caches to sync") + for _, podListerSynced := range gpm.podListerSynced { + if ok := k8sCache.WaitForCacheSync(ctx.Done(), podListerSynced); !ok { + gpm.logger.Fatal("failed to wait for caches to sync") + } } go gpm.service() gpm.poolPodC.InjectGpm(gpm) @@ -246,7 +250,7 @@ func (gpm *GenericPoolManager) IsValid(ctx context.Context, fsvc *fscache.FuncSv otelUtils.SpanTrackEvent(ctx, "IsValid", fscache.GetAttributesForFuncSvc(fsvc)...) for _, obj := range fsvc.KubernetesObjects { if strings.ToLower(obj.Kind) == "pod" { - pod, err := gpm.podLister.Pods(obj.Namespace).Get(obj.Name) + pod, err := gpm.podLister[obj.Namespace].Pods(obj.Namespace).Get(obj.Name) if err == nil && utils.IsReadyPod(pod) { // Normally, the address format is http://[pod-ip]:[port], however, if the // Istio is enabled the address format changes to http://[svc-name]:[port]. diff --git a/pkg/executor/executortype/poolmgr/poolpodcontroller.go b/pkg/executor/executortype/poolmgr/poolpodcontroller.go index f0fa4dbe..e4282f81 100644 --- a/pkg/executor/executortype/poolmgr/poolpodcontroller.go +++ b/pkg/executor/executortype/poolmgr/poolpodcontroller.go @@ -29,8 +29,7 @@ import ( "k8s.io/apimachinery/pkg/labels" utilruntime "k8s.io/apimachinery/pkg/util/runtime" "k8s.io/apimachinery/pkg/util/wait" - appsinformers "k8s.io/client-go/informers/apps/v1" - coreinformers "k8s.io/client-go/informers/core/v1" + k8sInformers "k8s.io/client-go/informers" "k8s.io/client-go/kubernetes" corelisters "k8s.io/client-go/listers/core/v1" k8sCache "k8s.io/client-go/tools/cache" @@ -54,10 +53,10 @@ type ( envListerSynced map[string]k8sCache.InformerSynced // podLister can list/get pods from the shared informer's store - podLister corelisters.PodLister + podLister map[string]corelisters.PodLister // podListerSynced returns true if the pod store has been synced at least once. - podListerSynced k8sCache.InformerSynced + podListerSynced map[string]k8sCache.InformerSynced envCreateUpdateQueue workqueue.RateLimitingInterface envDeleteQueue workqueue.RateLimitingInterface @@ -74,8 +73,7 @@ func NewPoolPodController(ctx context.Context, logger *zap.Logger, funcInformer map[string]finformerv1.FunctionInformer, pkgInformer map[string]finformerv1.PackageInformer, envInformer map[string]finformerv1.EnvironmentInformer, - rsInformer appsinformers.ReplicaSetInformer, - podInformer coreinformers.PodInformer) *PoolPodController { + gpmInformerFactory map[string]k8sInformers.SharedInformerFactory) *PoolPodController { logger = logger.Named("pool_pod_controller") p := &PoolPodController{ logger: logger, @@ -84,6 +82,8 @@ func NewPoolPodController(ctx context.Context, logger *zap.Logger, enableIstio: enableIstio, envLister: make(map[string]flisterv1.EnvironmentLister, 0), envListerSynced: make(map[string]k8sCache.InformerSynced, 0), + podLister: make(map[string]corelisters.PodLister), + podListerSynced: make(map[string]k8sCache.InformerSynced), envCreateUpdateQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "EnvAddUpdateQueue"), envDeleteQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "EnvDeleteQueue"), spCleanupPodQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "SpecializedPodCleanupQueue"), @@ -103,14 +103,16 @@ func NewPoolPodController(ctx context.Context, logger *zap.Logger, p.envLister[ns] = informer.Lister() p.envListerSynced[ns] = informer.Informer().HasSynced } - rsInformer.Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ - AddFunc: p.handleRSAdd, - UpdateFunc: p.handleRSUpdate, - DeleteFunc: p.handleRSDelete, - }) + for ns, informerFactory := range gpmInformerFactory { + informerFactory.Apps().V1().ReplicaSets().Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ + AddFunc: p.handleRSAdd, + 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 = podInformer.Lister() - p.podListerSynced = podInformer.Informer().HasSynced p.logger.Info("pool pod controller handlers registered") return p } @@ -138,7 +140,7 @@ func (p *PoolPodController) processRS(rs *apps.ReplicaSet) { return } rsLabelMap["managed"] = "false" - specializedPods, err := p.podLister.Pods(rs.Namespace).List(labels.SelectorFromSet(rsLabelMap)) + specializedPods, err := p.podLister[rs.Namespace].Pods(rs.Namespace).List(labels.SelectorFromSet(rsLabelMap)) if err != nil { logger.Error("Failed to list specialized pods", zap.Error(err)) } @@ -233,7 +235,9 @@ func (p *PoolPodController) Run(ctx context.Context, stopCh <-chan struct{}) { p.logger.Info("Waiting for informer caches to sync") waitSynced := make([]k8sCache.InformerSynced, 0) - waitSynced = append(waitSynced, p.podListerSynced) + for _, synced := range p.podListerSynced { + waitSynced = append(waitSynced, synced) + } for _, synced := range p.envListerSynced { waitSynced = append(waitSynced, synced) } @@ -373,7 +377,8 @@ func (p *PoolPodController) envDeleteQueueProcessFunc(ctx context.Context) bool p.logger.Debug("env delete request processing") p.gpm.cleanupPool(ctx, env) specializePodLables := getSpecializedPodLabels(env) - specializedPods, err := p.podLister.Pods(p.nsResolver.ResolveNamespace(p.nsResolver.FunctionNamespace)).List(labels.SelectorFromSet(specializePodLables)) + ns := p.nsResolver.ResolveNamespace(p.nsResolver.FunctionNamespace) + specializedPods, err := p.podLister[ns].Pods(ns).List(labels.SelectorFromSet(specializePodLables)) if err != nil { p.logger.Error("failed to list specialized pods", zap.Error(err)) p.envDeleteQueue.Forget(obj) @@ -413,7 +418,7 @@ func (p *PoolPodController) spCleanupPodQueueProcessFunc(ctx context.Context) bo p.spCleanupPodQueue.Forget(key) return false } - pod, err := p.podLister.Pods(namespace).Get(name) + pod, err := p.podLister[namespace].Pods(namespace).Get(name) if apierrors.IsNotFound(err) { p.logger.Info("pod not found", zap.String("key", key)) p.spCleanupPodQueue.Forget(key) diff --git a/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go b/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go index 4c45a657..22a27135 100644 --- a/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go +++ b/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go @@ -61,19 +61,17 @@ func TestPoolPodControllerPodCleanup(t *testing.T) { metav1.NamespaceAll: informerFactory.Core().V1().Environments(), } - gpmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypePoolmgr, time.Minute*30) + executorLabel, err := utils.GetInformerLabelByExecutor(fv1.ExecutorTypePoolmgr) if err != nil { - t.Fatalf("Error creating informer factory: %v", err) + t.Fatalf("Error creating labels for informer: %v", err) } - gpmPodInformer := gpmInformerFactory.Core().V1().Pods() - gpmRsInformer := gpmInformerFactory.Apps().V1().ReplicaSets() + gpmInformerFactory := utils.GetInformerFactoryByExecutor(kubernetesClient, executorLabel, time.Minute*30) ppc := NewPoolPodController(ctx, logger, kubernetesClient, false, funcInformer, pkgInformer, envInformer, - gpmRsInformer, - gpmPodInformer) + gpmInformerFactory) executorInstanceID := strings.ToLower(uniuri.NewLen(8)) metricsClient := metricsclient.NewSimpleClientset() @@ -86,7 +84,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) { fissionClient, kubernetesClient, metricsClient, fetcherConfig, executorInstanceID, funcInformer, pkgInformer, envInformer, - gpmPodInformer, gpmRsInformer, nil) + gpmInformerFactory, nil) if err != nil { t.Fatalf("Error creating generic pool manager: %v", err) } @@ -95,15 +93,18 @@ func TestPoolPodControllerPodCleanup(t *testing.T) { go ppc.Run(ctx, ctx.Done()) - podInformer := gpmPodInformer.Informer() - - runInformers(ctx, []k8sCache.SharedIndexInformer{ + informers := []k8sCache.SharedIndexInformer{ funcInformer[metav1.NamespaceAll].Informer(), pkgInformer[metav1.NamespaceAll].Informer(), envInformer[metav1.NamespaceAll].Informer(), - podInformer, - gpmRsInformer.Informer(), - }) + } + + for _, informerFactory := range gpmInformerFactory { + informers = append(informers, informerFactory.Core().V1().Pods().Informer()) + informers = append(informers, informerFactory.Apps().V1().ReplicaSets().Informer()) + } + + runInformers(ctx, informers) pod := &corev1.Pod{ ObjectMeta: metav1.ObjectMeta{ @@ -124,7 +125,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) { found := false for found == false && time.Since(start) < time.Second*5 { t.Log("Waiting for pod to be added to pool") - pod, err := ppc.podLister.Pods(pod.Namespace).Get(pod.Name) + pod, err := ppc.podLister[pod.Namespace].Pods(pod.Namespace).Get(pod.Name) if err == nil { found = true t.Logf("Found pod %#v", pod.ObjectMeta) diff --git a/pkg/utils/informer.go b/pkg/utils/informer.go index 70d0e79b..1f8097ea 100644 --- a/pkg/utils/informer.go +++ b/pkg/utils/informer.go @@ -69,6 +69,22 @@ func GetK8sInformersForNamespaces(client kubernetes.Interface, defaultSync time. return informers } +func GetInformerFactoryByExecutor(client kubernetes.Interface, labels labels.Selector, defaultResync time.Duration) map[string]k8sInformers.SharedInformerFactory { + informerFactory := make(map[string]k8sInformers.SharedInformerFactory) + + namespaces := DefaultNSResolver() + for _, ns := range namespaces.FissionNSWithOptions(WithBuilderNs(), WithFunctionNs(), WithDefaultNs()) { + factory := k8sInformers.NewSharedInformerFactoryWithOptions(client, defaultResync, + k8sInformers.WithTweakListOptions(func(options *metav1.ListOptions) { + options.LabelSelector = labels.String() + }), + k8sInformers.WithNamespace(ns)) + + informerFactory[ns] = factory + } + return informerFactory +} + func GetInformerFactoryByReadyPod(client kubernetes.Interface, namespace string, labelSelector *metav1.LabelSelector) (k8sInformers.SharedInformerFactory, error) { informerFactory := k8sInformers.NewSharedInformerFactoryWithOptions(client, 0, k8sInformers.WithNamespace(namespace), @@ -79,20 +95,16 @@ func GetInformerFactoryByReadyPod(client kubernetes.Interface, namespace string, return informerFactory, nil } -func GetInformerFactoryByExecutor(client kubernetes.Interface, executorType fv1.ExecutorType, defaultResync time.Duration) (k8sInformers.SharedInformerFactory, error) { +func GetInformerLabelByExecutor(executorType fv1.ExecutorType) (labels.Selector, error) { executorLabel, err := labels.NewRequirement(fv1.EXECUTOR_TYPE, selection.DoubleEquals, []string{string(executorType)}) if err != nil { return nil, err } labelSelector := labels.NewSelector() labelSelector.Add(*executorLabel) - informerFactory := k8sInformers.NewSharedInformerFactoryWithOptions(client, defaultResync, - k8sInformers.WithTweakListOptions(func(options *metav1.ListOptions) { - options.LabelSelector = labelSelector.String() - })) - return informerFactory, nil -} + return labelSelector, nil +} func SupportedMetricsAPIVersionAvailable(discoveredAPIGroups *metav1.APIGroupList) bool { var supportedMetricsAPIVersions = []string{ "v1beta1",