diff --git a/pkg/executor/cms/cmscontroller.go b/pkg/executor/cms/cmscontroller.go index 8f84866d..876b476e 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" - k8sCache "k8s.io/client-go/tools/cache" fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/crd" @@ -34,9 +34,6 @@ type ( ConfigSecretController struct { logger *zap.Logger - configmapInformer *k8sCache.SharedIndexInformer - secretInformer *k8sCache.SharedIndexInformer - fissionClient *crd.FissionClient } ) @@ -44,17 +41,15 @@ type ( // MakeConfigSecretController makes a controller for configmaps and secrets which changes related functions func MakeConfigSecretController(ctx context.Context, logger *zap.Logger, fissionClient *crd.FissionClient, kubernetesClient *kubernetes.Clientset, types map[fv1.ExecutorType]executortype.ExecutorType, - configmapInformer *k8sCache.SharedIndexInformer, - secretInformer *k8sCache.SharedIndexInformer) *ConfigSecretController { + configmapInformer informerv1.ConfigMapInformer, + secretInformer informerv1.SecretInformer) *ConfigSecretController { logger.Debug("Creating ConfigMap & Secret Controller") cmsController := &ConfigSecretController{ - logger: logger, - configmapInformer: configmapInformer, - secretInformer: secretInformer, - fissionClient: fissionClient, + logger: logger, + fissionClient: fissionClient, } - (*configmapInformer).AddEventHandler(ConfigMapEventHandlers(ctx, logger, fissionClient, kubernetesClient, types)) - (*secretInformer).AddEventHandler(SecretEventHandlers(ctx, logger, fissionClient, kubernetesClient, types)) + configmapInformer.Informer().AddEventHandler(ConfigMapEventHandlers(ctx, logger, fissionClient, kubernetesClient, types)) + secretInformer.Informer().AddEventHandler(SecretEventHandlers(ctx, logger, fissionClient, kubernetesClient, types)) return cmsController } diff --git a/pkg/executor/executor.go b/pkg/executor/executor.go index 5c2b97a7..91f5f2e0 100644 --- a/pkg/executor/executor.go +++ b/pkg/executor/executor.go @@ -277,15 +277,15 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames logger.Info("Starting executor", zap.String("instanceID", executorInstanceID)) informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) - funcInformer := informerFactory.Core().V1().Functions().Informer() - pkgInformer := informerFactory.Core().V1().Packages().Informer() - envInformer := informerFactory.Core().V1().Environments().Informer() + funcInformer := informerFactory.Core().V1().Functions() + pkgInformer := informerFactory.Core().V1().Packages() + envInformer := informerFactory.Core().V1().Environments() gpm, err := poolmgr.MakeGenericPoolManager( logger, fissionClient, kubernetesClient, metricsClient, functionNamespace, fetcherConfig, executorInstanceID, - &funcInformer, &pkgInformer, + funcInformer, pkgInformer, ) if err != nil { return errors.Wrap(err, "pool manager creation faied") @@ -295,7 +295,7 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames logger, fissionClient, kubernetesClient, functionNamespace, fetcherConfig, executorInstanceID, - &funcInformer, &envInformer, + funcInformer, envInformer, ) if err != nil { return errors.Wrap(err, "new deploy manager creation faied") @@ -306,7 +306,7 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames ctx, logger, fissionClient, kubernetesClient, - functionNamespace, executorInstanceID, &funcInformer) + functionNamespace, executorInstanceID, funcInformer) if err != nil { return errors.Wrap(err, "container manager creation faied") } @@ -334,13 +334,13 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames util.WaitTimeout(wg, 30*time.Second) k8sInformerFactory := k8sInformers.NewSharedInformerFactory(kubernetesClient, time.Minute*30) - configmapInformer := k8sInformerFactory.Core().V1().ConfigMaps().Informer() - secretInformer := k8sInformerFactory.Core().V1().Secrets().Informer() + configmapInformer := k8sInformerFactory.Core().V1().ConfigMaps() + secretInformer := k8sInformerFactory.Core().V1().Secrets() - cms := cms.MakeConfigSecretController(ctx, logger, fissionClient, kubernetesClient, executorTypes, &configmapInformer, &secretInformer) + cms := cms.MakeConfigSecretController(ctx, logger, fissionClient, kubernetesClient, executorTypes, configmapInformer, secretInformer) api, err := MakeExecutor(ctx, logger, cms, fissionClient, executorTypes, []k8sCache.SharedIndexInformer{ - funcInformer, pkgInformer, envInformer, configmapInformer, secretInformer, + funcInformer.Informer(), pkgInformer.Informer(), envInformer.Informer(), configmapInformer.Informer(), secretInformer.Informer(), }) if err != nil { return err diff --git a/pkg/executor/executortype/container/containermgr.go b/pkg/executor/executortype/container/containermgr.go index 35c544e5..cecfb17f 100644 --- a/pkg/executor/executortype/container/containermgr.go +++ b/pkg/executor/executortype/container/containermgr.go @@ -43,6 +43,7 @@ import ( "github.com/fission/fission/pkg/executor/executortype" "github.com/fission/fission/pkg/executor/fscache" "github.com/fission/fission/pkg/executor/reaper" + finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1" "github.com/fission/fission/pkg/throttler" "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/maps" @@ -68,7 +69,6 @@ type ( throttler *throttler.Throttler - funcInformer *k8sCache.SharedIndexInformer serviceInformer k8sCache.SharedIndexInformer deploymentInformer k8sCache.SharedIndexInformer @@ -84,7 +84,7 @@ func MakeContainer( kubernetesClient *kubernetes.Clientset, namespace string, instanceID string, - funcInformer *k8sCache.SharedIndexInformer) (executortype.ExecutorType, error) { + funcInformer finformerv1.FunctionInformer) (executortype.ExecutorType, error) { enableIstio := false if len(os.Getenv("ENABLE_ISTIO")) > 0 { istio, err := strconv.ParseBool(os.Getenv("ENABLE_ISTIO")) @@ -105,15 +105,13 @@ func MakeContainer( fsCache: fscache.MakeFunctionServiceCache(logger), throttler: throttler.MakeThrottler(1 * time.Minute), - funcInformer: funcInformer, - runtimeImagePullPolicy: utils.GetImagePullPolicy(os.Getenv("RUNTIME_IMAGE_PULL_POLICY")), useIstio: enableIstio, // Time is set slightly higher than NewDeploy as cold starts are longer for CaaF defaultIdlePodReapTime: 1 * time.Minute, } - (*caaf.funcInformer).AddEventHandler(caaf.FuncInformerHandler(ctx)) + funcInformer.Informer().AddEventHandler(caaf.FuncInformerHandler(ctx)) informerFactory, err := utils.GetInformerFactoryByExecutor(caaf.kubernetesClient, fv1.ExecutorTypeContainer) if err != nil { diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr.go b/pkg/executor/executortype/newdeploy/newdeploymgr.go index 99ba566c..132cb124 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr.go @@ -44,6 +44,7 @@ import ( "github.com/fission/fission/pkg/executor/fscache" "github.com/fission/fission/pkg/executor/reaper" fetcherConfig "github.com/fission/fission/pkg/fetcher/config" + finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1" "github.com/fission/fission/pkg/throttler" "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/maps" @@ -69,9 +70,6 @@ type ( throttler *throttler.Throttler - funcInformer *k8sCache.SharedIndexInformer - envInformer *k8sCache.SharedIndexInformer - serviceInformer k8sCache.SharedIndexInformer deploymentInformer k8sCache.SharedIndexInformer @@ -87,8 +85,8 @@ func MakeNewDeploy( namespace string, fetcherConfig *fetcherConfig.Config, instanceID string, - funcInformer *k8sCache.SharedIndexInformer, - envInformer *k8sCache.SharedIndexInformer, + funcInformer finformerv1.FunctionInformer, + envInformer finformerv1.EnvironmentInformer, ) (executortype.ExecutorType, error) { enableIstio := false if len(os.Getenv("ENABLE_ISTIO")) > 0 { @@ -115,12 +113,10 @@ func MakeNewDeploy( useIstio: enableIstio, defaultIdlePodReapTime: 2 * time.Minute, - funcInformer: funcInformer, - envInformer: envInformer, } - (*nd.funcInformer).AddEventHandler(nd.FunctionEventHandlers()) - (*nd.envInformer).AddEventHandler(nd.EnvEventHandlers()) + funcInformer.Informer().AddEventHandler(nd.FunctionEventHandlers()) + envInformer.Informer().AddEventHandler(nd.EnvEventHandlers()) informerFactory, err := utils.GetInformerFactoryByExecutor(nd.kubernetesClient, fv1.ExecutorTypePoolmgr) if err != nil { diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index 7bf8a887..46a4785f 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -46,6 +46,7 @@ import ( "github.com/fission/fission/pkg/executor/fscache" "github.com/fission/fission/pkg/executor/reaper" fetcherConfig "github.com/fission/fission/pkg/fetcher/config" + finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1" "github.com/fission/fission/pkg/utils" ) @@ -76,9 +77,6 @@ type ( enableIstio bool fetcherConfig *fetcherConfig.Config - funcInformer *k8sCache.SharedIndexInformer - pkgInformer *k8sCache.SharedIndexInformer - podInformer k8sCache.SharedIndexInformer defaultIdlePodReapTime time.Duration @@ -104,8 +102,8 @@ func MakeGenericPoolManager( functionNamespace string, fetcherConfig *fetcherConfig.Config, instanceID string, - funcInformer *k8sCache.SharedIndexInformer, - pkgInformer *k8sCache.SharedIndexInformer, + funcInformer finformerv1.FunctionInformer, + pkgInformer finformerv1.PackageInformer, ) (executortype.ExecutorType, error) { gpmLogger := logger.Named("generic_pool_manager") @@ -133,14 +131,12 @@ func MakeGenericPoolManager( defaultIdlePodReapTime: 2 * time.Minute, fetcherConfig: fetcherConfig, enableIstio: enableIstio, - funcInformer: funcInformer, - pkgInformer: pkgInformer, } go gpm.service() - (*gpm.funcInformer).AddEventHandler(FunctionEventHandlers(gpm.logger, gpm.kubernetesClient, gpm.namespace, gpm.enableIstio)) - (*gpm.pkgInformer).AddEventHandler(PackageEventHandlers(gpm.logger, gpm.kubernetesClient, gpm.namespace)) + funcInformer.Informer().AddEventHandler(FunctionEventHandlers(gpm.logger, gpm.kubernetesClient, gpm.namespace, gpm.enableIstio)) + pkgInformer.Informer().AddEventHandler(PackageEventHandlers(gpm.logger, gpm.kubernetesClient, gpm.namespace)) kubeInformerFactory, err := utils.GetInformerFactoryByExecutor(gpm.kubernetesClient, fv1.ExecutorTypePoolmgr) if err != nil {