diff --git a/charts/fission-all/templates/deployment.yaml b/charts/fission-all/templates/deployment.yaml index 0e9c635a..7fc1b96e 100644 --- a/charts/fission-all/templates/deployment.yaml +++ b/charts/fission-all/templates/deployment.yaml @@ -85,6 +85,7 @@ rules: resources: - deployments - deployments/scale + - replicasets verbs: - '*' - apiGroups: diff --git a/pkg/crd/key.go b/pkg/crd/key.go index 425fe86e..e99e5b36 100644 --- a/pkg/crd/key.go +++ b/pkg/crd/key.go @@ -31,3 +31,10 @@ import ( func CacheKey(metadata *metav1.ObjectMeta) string { return fmt.Sprintf("%v_%v", metadata.UID, metadata.ResourceVersion) } + +// CacheKeyForUID create a key that uniquely identifies the +// of the object. Since resourceVersion changes on every update and +// UIDs are unique, we don't use resource version here +func CacheKeyUID(metadata *metav1.ObjectMeta) string { + return fmt.Sprintf("%v", metadata.UID) +} diff --git a/pkg/executor/executor.go b/pkg/executor/executor.go index 91f5f2e0..80340fb9 100644 --- a/pkg/executor/executor.go +++ b/pkg/executor/executor.go @@ -45,6 +45,7 @@ import ( "github.com/fission/fission/pkg/executor/util" fetcherConfig "github.com/fission/fission/pkg/fetcher/config" genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" + "github.com/fission/fission/pkg/utils" ) type ( @@ -281,32 +282,50 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames pkgInformer := informerFactory.Core().V1().Packages() envInformer := informerFactory.Core().V1().Environments() + gpmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypePoolmgr, time.Minute*30) + if err != nil { + return err + } + gpmPodInformer := gpmInformerFactory.Core().V1().Pods() + gpmRsInformer := gpmInformerFactory.Apps().V1().ReplicaSets() gpm, err := poolmgr.MakeGenericPoolManager( logger, fissionClient, kubernetesClient, metricsClient, functionNamespace, fetcherConfig, executorInstanceID, - funcInformer, pkgInformer, - ) + funcInformer, pkgInformer, envInformer, + gpmPodInformer, gpmRsInformer) if err != nil { return errors.Wrap(err, "pool manager creation faied") } + ndmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypeNewdeploy, time.Minute*30) + if err != nil { + return err + } + ndmDeplInformer := ndmInformerFactory.Apps().V1().Deployments() + ndmSvcInformer := ndmInformerFactory.Core().V1().Services() ndm, err := newdeploy.MakeNewDeploy( logger, fissionClient, kubernetesClient, functionNamespace, fetcherConfig, executorInstanceID, funcInformer, envInformer, - ) + ndmDeplInformer, ndmSvcInformer) if err != nil { return errors.Wrap(err, "new deploy manager creation faied") } + cnmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypeContainer, time.Minute*30) + if err != nil { + return err + } + cnmDeplInformer := cnmInformerFactory.Apps().V1().Deployments() + cnmSvcInformer := cnmInformerFactory.Core().V1().Services() ctx := context.Background() cnm, err := container.MakeContainer( - ctx, - logger, + ctx, logger, fissionClient, kubernetesClient, - functionNamespace, executorInstanceID, funcInformer) + functionNamespace, executorInstanceID, funcInformer, + cnmDeplInformer, cnmSvcInformer) if err != nil { return errors.Wrap(err, "container manager creation faied") } @@ -339,13 +358,23 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames cms := cms.MakeConfigSecretController(ctx, logger, fissionClient, kubernetesClient, executorTypes, configmapInformer, secretInformer) - api, err := MakeExecutor(ctx, logger, cms, fissionClient, executorTypes, []k8sCache.SharedIndexInformer{ - funcInformer.Informer(), pkgInformer.Informer(), envInformer.Informer(), configmapInformer.Informer(), secretInformer.Informer(), - }) + api, err := MakeExecutor(ctx, logger, cms, fissionClient, executorTypes, + []k8sCache.SharedIndexInformer{ + funcInformer.Informer(), + pkgInformer.Informer(), + envInformer.Informer(), + configmapInformer.Informer(), + secretInformer.Informer(), + gpmPodInformer.Informer(), + gpmRsInformer.Informer(), + ndmDeplInformer.Informer(), + ndmSvcInformer.Informer(), + cnmDeplInformer.Informer(), + cnmSvcInformer.Informer(), + }) if err != nil { return err } - go reaper.CleanupRoleBindings(logger, kubernetesClient, fissionClient, functionNamespace, envBuilderNamespace, time.Minute*30) go api.Serve(port, openTracingEnabled) go serveMetric(logger) diff --git a/pkg/executor/executortype/container/containermgr.go b/pkg/executor/executortype/container/containermgr.go index f249be58..d28b4d70 100644 --- a/pkg/executor/executortype/container/containermgr.go +++ b/pkg/executor/executortype/container/containermgr.go @@ -29,13 +29,16 @@ import ( multierror "github.com/hashicorp/go-multierror" "github.com/pkg/errors" "go.uber.org/zap" - appsv1 "k8s.io/api/apps/v1" apiv1 "k8s.io/api/core/v1" k8sErrs "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" k8sTypes "k8s.io/apimachinery/pkg/types" + appsinformers "k8s.io/client-go/informers/apps/v1" + coreinformers "k8s.io/client-go/informers/core/v1" "k8s.io/client-go/kubernetes" + appslisters "k8s.io/client-go/listers/apps/v1" + corelisters "k8s.io/client-go/listers/core/v1" k8sCache "k8s.io/client-go/tools/cache" fv1 "github.com/fission/fission/pkg/apis/core/v1" @@ -69,10 +72,13 @@ type ( throttler *throttler.Throttler - serviceInformer k8sCache.SharedIndexInformer - deploymentInformer k8sCache.SharedIndexInformer - defaultIdlePodReapTime time.Duration + + deplLister appslisters.DeploymentLister + svcLister corelisters.ServiceLister + + deplListerSynced k8sCache.InformerSynced + svcListerSynced k8sCache.InformerSynced } ) @@ -84,7 +90,10 @@ func MakeContainer( kubernetesClient *kubernetes.Clientset, namespace string, instanceID string, - funcInformer finformerv1.FunctionInformer) (executortype.ExecutorType, error) { + funcInformer finformerv1.FunctionInformer, + deplInformer appsinformers.DeploymentInformer, + svcInformer coreinformers.ServiceInformer, +) (executortype.ExecutorType, error) { enableIstio := false if len(os.Getenv("ENABLE_ISTIO")) > 0 { istio, err := strconv.ParseBool(os.Getenv("ENABLE_ISTIO")) @@ -110,21 +119,21 @@ func MakeContainer( // Time is set slightly higher than NewDeploy as cold starts are longer for CaaF defaultIdlePodReapTime: 1 * time.Minute, } + caaf.deplLister = deplInformer.Lister() + caaf.deplListerSynced = deplInformer.Informer().HasSynced + + caaf.svcLister = svcInformer.Lister() + caaf.svcListerSynced = svcInformer.Informer().HasSynced funcInformer.Informer().AddEventHandler(caaf.FuncInformerHandler(ctx)) - - informerFactory, err := utils.GetInformerFactoryByExecutor(caaf.kubernetesClient, fv1.ExecutorTypeContainer) - if err != nil { - return nil, err - } - caaf.serviceInformer = informerFactory.Core().V1().Services().Informer() - caaf.deploymentInformer = informerFactory.Apps().V1().Deployments().Informer() - return caaf, nil } // 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 { + caaf.logger.Fatal("failed to wait for caches to sync") + } go caaf.idleObjectReaper() } @@ -174,40 +183,6 @@ func (caaf *Container) TapService(ctx context.Context, svcHost string) error { return nil } -func (caaf *Container) getServiceInfo(ctx context.Context, obj apiv1.ObjectReference) (*apiv1.Service, error) { - item, exists, err := utils.GetCachedItem(obj, caaf.serviceInformer) - - if err != nil || !exists { - caaf.logger.Debug( - "Falling back to getting service info from k8s API -- this may cause performance issues for your function.", - zap.Bool("exists", exists), - zap.Error(err), - ) - service, err := caaf.kubernetesClient.CoreV1().Services(obj.Namespace).Get(ctx, obj.Name, metav1.GetOptions{}) - return service, err - } - - service := item.(*apiv1.Service) - return service, nil -} - -func (caaf *Container) getDeploymentInfo(ctx context.Context, obj apiv1.ObjectReference) (*appsv1.Deployment, error) { - item, exists, err := utils.GetCachedItem(obj, caaf.deploymentInformer) - - if err != nil || !exists { - caaf.logger.Debug( - "Falling back to getting deployment info from k8s API -- this may cause performance issues for your function.", - zap.Bool("exists", exists), - zap.Error(err), - ) - deployment, err := caaf.kubernetesClient.AppsV1().Deployments(obj.Namespace).Get(ctx, obj.Name, metav1.GetOptions{}) - return deployment, err - } - - deployment := item.(*appsv1.Deployment) - return deployment, nil -} - // IsValid does a get on the service address to ensure it's a valid service, then // scale deployment to 1 replica if there are no available replicas for function. // Return true if no error occurs, return false otherwise. @@ -222,16 +197,15 @@ func (caaf *Container) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) bool } for _, obj := range fsvc.KubernetesObjects { if strings.ToLower(obj.Kind) == "service" { - _, err := caaf.getServiceInfo(ctx, obj) + _, err := caaf.svcLister.Services(obj.Namespace).Get(obj.Name) if err != nil { if !k8sErrs.IsNotFound(err) { caaf.logger.Error("error validating function service", zap.String("function", fsvc.Function.Name), zap.Error(err)) } return false } - } else if strings.ToLower(obj.Kind) == "deployment" { - currentDeploy, err := caaf.getDeploymentInfo(ctx, obj) + currentDeploy, err := caaf.deplLister.Deployments(obj.Namespace).Get(obj.Name) if err != nil { if !k8sErrs.IsNotFound(err) { caaf.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 7fab0888..7b20b803 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr.go @@ -28,14 +28,17 @@ import ( multierror "github.com/hashicorp/go-multierror" "github.com/pkg/errors" "go.uber.org/zap" - appsv1 "k8s.io/api/apps/v1" autoscalingv1 "k8s.io/api/autoscaling/v1" apiv1 "k8s.io/api/core/v1" k8sErrs "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" k8sTypes "k8s.io/apimachinery/pkg/types" + appsinformers "k8s.io/client-go/informers/apps/v1" + coreinformers "k8s.io/client-go/informers/core/v1" "k8s.io/client-go/kubernetes" + appslisters "k8s.io/client-go/listers/apps/v1" + corelisters "k8s.io/client-go/listers/core/v1" k8sCache "k8s.io/client-go/tools/cache" fv1 "github.com/fission/fission/pkg/apis/core/v1" @@ -70,10 +73,13 @@ type ( throttler *throttler.Throttler - serviceInformer k8sCache.SharedIndexInformer - deploymentInformer k8sCache.SharedIndexInformer - defaultIdlePodReapTime time.Duration + + deplLister appslisters.DeploymentLister + svcLister corelisters.ServiceLister + + deplListerSynced k8sCache.InformerSynced + svcListerSynced k8sCache.InformerSynced } ) @@ -87,6 +93,8 @@ func MakeNewDeploy( instanceID string, funcInformer finformerv1.FunctionInformer, envInformer finformerv1.EnvironmentInformer, + deplInformer appsinformers.DeploymentInformer, + svcInformer coreinformers.ServiceInformer, ) (executortype.ExecutorType, error) { enableIstio := false if len(os.Getenv("ENABLE_ISTIO")) > 0 { @@ -115,22 +123,23 @@ func MakeNewDeploy( defaultIdlePodReapTime: 2 * time.Minute, } + nd.deplLister = deplInformer.Lister() + nd.deplListerSynced = deplInformer.Informer().HasSynced + + nd.svcLister = svcInformer.Lister() + nd.svcListerSynced = svcInformer.Informer().HasSynced + funcInformer.Informer().AddEventHandler(nd.FunctionEventHandlers()) envInformer.Informer().AddEventHandler(nd.EnvEventHandlers()) - informerFactory, err := utils.GetInformerFactoryByExecutor(nd.kubernetesClient, fv1.ExecutorTypePoolmgr) - if err != nil { - return nil, err - } - nd.serviceInformer = informerFactory.Core().V1().Services().Informer() - nd.deploymentInformer = informerFactory.Apps().V1().Deployments().Informer() return nd, nil } // Run start the function and environment controller along with an object reaper. func (deploy *NewDeploy) Run(ctx context.Context) { - go deploy.serviceInformer.Run(ctx.Done()) - go deploy.deploymentInformer.Run(ctx.Done()) + if ok := k8sCache.WaitForCacheSync(ctx.Done(), deploy.deplListerSynced, deploy.svcListerSynced); !ok { + deploy.logger.Fatal("failed to wait for caches to sync") + } go deploy.idleObjectReaper() } @@ -177,40 +186,6 @@ func (deploy *NewDeploy) TapService(ctx context.Context, svcHost string) error { return nil } -func (deploy *NewDeploy) getServiceInfo(ctx context.Context, obj apiv1.ObjectReference) (*apiv1.Service, error) { - item, exists, err := utils.GetCachedItem(obj, deploy.serviceInformer) - - if err != nil || !exists { - deploy.logger.Debug( - "Falling back to getting service info from k8s API -- this may cause performance issues for your function.", - zap.Bool("exists", exists), - zap.Error(err), - ) - service, err := deploy.kubernetesClient.CoreV1().Services(obj.Namespace).Get(ctx, obj.Name, metav1.GetOptions{}) - return service, err - } - - service := item.(*apiv1.Service) - return service, nil -} - -func (deploy *NewDeploy) getDeploymentInfo(ctx context.Context, obj apiv1.ObjectReference) (*appsv1.Deployment, error) { - item, exists, err := utils.GetCachedItem(obj, deploy.deploymentInformer) - - if err != nil || !exists { - deploy.logger.Debug( - "Falling back to getting deployment info from k8s API -- this may cause performance issues for your function.", - zap.Bool("exists", exists), - zap.Error(err), - ) - deployment, err := deploy.kubernetesClient.AppsV1().Deployments(obj.Namespace).Get(ctx, obj.Name, metav1.GetOptions{}) - return deployment, err - } - - deployment := item.(*appsv1.Deployment) - return deployment, nil -} - // IsValid does a get on the service address to ensure it's a valid service, then // scale deployment to 1 replica if there are no available replicas for function. // Return true if no error occurs, return false otherwise. @@ -225,7 +200,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.getServiceInfo(ctx, obj) + _, err := deploy.svcLister.Services(obj.Namespace).Get(obj.Name) if err != nil { if !k8sErrs.IsNotFound(err) { deploy.logger.Error("error validating function service", zap.String("function", fsvc.Function.Name), zap.Error(err)) @@ -234,7 +209,7 @@ func (deploy *NewDeploy) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) boo } } else if strings.ToLower(obj.Kind) == "deployment" { - currentDeploy, err := deploy.getDeploymentInfo(ctx, obj) + currentDeploy, err := deploy.deplLister.Deployments(obj.Namespace).Get(obj.Name) if err != nil { if !k8sErrs.IsNotFound(err) { deploy.logger.Error("error validating function deployment", zap.String("function", fsvc.Function.Name), zap.Error(err)) diff --git a/pkg/executor/executortype/poolmgr/common.go b/pkg/executor/executortype/poolmgr/common.go index 5cf07bea..bfb27d5e 100644 --- a/pkg/executor/executortype/poolmgr/common.go +++ b/pkg/executor/executortype/poolmgr/common.go @@ -27,3 +27,13 @@ func getEnvPoolSize(env *fv1.Environment) int32 { } return poolsize } + +func getSpecializedPodLabels(env *fv1.Environment) map[string]string { + specialPodLabels := make(map[string]string) + specialPodLabels[fv1.EXECUTOR_TYPE] = string(fv1.ExecutorTypePoolmgr) + specialPodLabels[fv1.ENVIRONMENT_NAME] = env.ObjectMeta.Name + specialPodLabels[fv1.ENVIRONMENT_NAMESPACE] = env.ObjectMeta.Namespace + specialPodLabels[fv1.ENVIRONMENT_UID] = string(env.ObjectMeta.UID) + specialPodLabels["managed"] = "false" + return specialPodLabels +} diff --git a/pkg/executor/executortype/poolmgr/gp.go b/pkg/executor/executortype/poolmgr/gp.go index 82500963..59d750e1 100644 --- a/pkg/executor/executortype/poolmgr/gp.go +++ b/pkg/executor/executortype/poolmgr/gp.go @@ -90,7 +90,7 @@ func MakeGenericPool( fsCache *fscache.FunctionServiceCache, fetcherConfig *fetcherConfig.Config, instanceID string, - enableIstio bool) (*GenericPool, error) { + enableIstio bool) *GenericPool { gpLogger := logger.Named("generic_pool") @@ -129,24 +129,25 @@ func MakeGenericPool( gp.runtimeImagePullPolicy = utils.GetImagePullPolicy(os.Getenv("RUNTIME_IMAGE_PULL_POLICY")) + return gp +} + +func (gp *GenericPool) setup(ctx context.Context) error { // create fetcher SA in this ns, if not already created - err = fetcherConfig.SetupServiceAccount(gp.kubernetesClient, gp.namespace, nil) + err := gp.fetcherConfig.SetupServiceAccount(gp.kubernetesClient, gp.namespace, nil) if err != nil { - return nil, errors.Wrapf(err, "error creating fetcher service account in namespace %q", gp.namespace) + return errors.Wrapf(err, "error creating fetcher service account in namespace %q", gp.namespace) } - // Labels for generic deployment/RS/pods. - //gp.labelsForPool = gp.getDeployLabels() - // create the pool - err = gp.createPoolDeployment(context.Background(), env) + err = gp.createPoolDeployment(ctx, gp.env) if err != nil { - return nil, err + return err } go gp.startReadyPodController() go gp.updateCPUUtilizationSvc() - return gp, nil + return nil } func (gp *GenericPool) getEnvironmentPoolLabels(env *fv1.Environment) map[string]string { @@ -575,7 +576,7 @@ func (gp *GenericPool) getPercent(cpuUsage resource.Quantity, percentage float64 } // destroys the pool -- the deployment, replicaset and pods -func (gp *GenericPool) destroy() error { +func (gp *GenericPool) destroy(ctx context.Context) error { close(gp.stopReadyPodControllerCh) deletePropagation := metav1.DeletePropagationBackground @@ -584,7 +585,7 @@ func (gp *GenericPool) destroy() error { } err := gp.kubernetesClient.AppsV1(). - Deployments(gp.namespace).Delete(context.TODO(), gp.deployment.ObjectMeta.Name, delOpt) + Deployments(gp.namespace).Delete(ctx, gp.deployment.ObjectMeta.Name, delOpt) if err != nil { gp.logger.Error("error destroying deployment", zap.Error(err), diff --git a/pkg/executor/executortype/poolmgr/gp_deployment.go b/pkg/executor/executortype/poolmgr/gp_deployment.go index 95e91cbd..89a452ed 100644 --- a/pkg/executor/executortype/poolmgr/gp_deployment.go +++ b/pkg/executor/executortype/poolmgr/gp_deployment.go @@ -203,3 +203,41 @@ func (gp *GenericPool) createPoolDeployment(ctx context.Context, env *fv1.Enviro return nil } + +func (gp *GenericPool) updatePoolDeployment(ctx context.Context, env *fv1.Environment) error { + logger := gp.logger.With(zap.String("env", env.Name), zap.String("namespace", env.Namespace)) + if gp.env.ObjectMeta.ResourceVersion == env.ObjectMeta.ResourceVersion { + logger.Debug("env resource version matching with pool env") + return nil + } + newDeployment := gp.deployment.DeepCopy() + spec, err := gp.genDeploymentSpec(env) + if err != nil { + logger.Error("error generating deployment spec", zap.Error(err)) + return err + } + newDeployment.Spec = *spec + deployMeta := gp.genDeploymentMeta(env) + deployMeta.Name = gp.deployment.Name + newDeployment.ObjectMeta = deployMeta + + poolsize := getEnvPoolSize(env) + switch env.Spec.AllowedFunctionsPerContainer { + case fv1.AllowedFunctionsPerContainerInfinite: + poolsize = 1 + } + newDeployment.Spec.Replicas = &poolsize + + depl, err := gp.kubernetesClient.AppsV1().Deployments(gp.namespace).Update(ctx, newDeployment, metav1.UpdateOptions{}) + if err != nil { + logger.Error("error updating deployment in kubernetes", zap.Error(err), zap.String("deployment", depl.Name)) + return err + } + // possible concurrency issue here as + // gp.env and gp.deployment referenced at few places + // we can move update pool to gpm.service if required + gp.env = env + gp.deployment = depl + logger.Info("Updated deployment for pool", zap.String("deployment", depl.Name)) + return nil +} diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index ea0c20a0..2879d602 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -35,7 +35,10 @@ import ( "k8s.io/apimachinery/pkg/runtime" k8sTypes "k8s.io/apimachinery/pkg/types" "k8s.io/apimachinery/pkg/watch" + appsinformers "k8s.io/client-go/informers/apps/v1" + coreinformers "k8s.io/client-go/informers/core/v1" "k8s.io/client-go/kubernetes" + corelisters "k8s.io/client-go/listers/core/v1" k8sCache "k8s.io/client-go/tools/cache" metricsclient "k8s.io/metrics/pkg/client/clientset/versioned" @@ -56,7 +59,7 @@ type requestType int const ( GET_POOL requestType = iota - CLEANUP_POOLS + CLEANUP_POOL ) type ( @@ -77,14 +80,20 @@ type ( enableIstio bool fetcherConfig *fetcherConfig.Config - podInformer k8sCache.SharedIndexInformer + // podLister can list/get pods from the shared informer's store + podLister corelisters.PodLister + + // podListerSynced returns true if the pod store has been synced at least once. + podListerSynced k8sCache.InformerSynced defaultIdlePodReapTime time.Duration + + poolPodC *PoolPodController } request struct { requestType + ctx context.Context env *fv1.Environment - envList []fv1.Environment responseChannel chan *response } response struct { @@ -104,6 +113,9 @@ func MakeGenericPoolManager( instanceID string, funcInformer finformerv1.FunctionInformer, pkgInformer finformerv1.PackageInformer, + envInformer finformerv1.EnvironmentInformer, + podInformer coreinformers.PodInformer, + rsInformer appsinformers.ReplicaSetInformer, ) (executortype.ExecutorType, error) { gpmLogger := logger.Named("generic_pool_manager") @@ -117,6 +129,9 @@ func MakeGenericPoolManager( enableIstio = istio } + poolPodC := NewPoolPodController(gpmLogger, kubernetesClient, functionNamespace, + enableIstio, funcInformer, pkgInformer, envInformer, rsInformer) + gpm := &GenericPoolManager{ logger: gpmLogger, pools: make(map[string]*GenericPool), @@ -131,29 +146,24 @@ func MakeGenericPoolManager( defaultIdlePodReapTime: 2 * time.Minute, fetcherConfig: fetcherConfig, enableIstio: enableIstio, + poolPodC: poolPodC, } + gpm.podLister = podInformer.Lister() + gpm.podListerSynced = podInformer.Informer().HasSynced - go gpm.service() - - 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 { - return nil, err - } - gpm.podInformer = kubeInformerFactory.Core().V1().Pods().Informer() return gpm, nil } func (gpm *GenericPoolManager) Run(ctx context.Context) { - // eagerPoolCreator must run after CleanupOldExecutorObjects. - // Otherwise, the poolmanager may wrongly delete the deployment. - go gpm.eagerPoolCreator() - go gpm.podInformer.Run(ctx.Done()) + if ok := k8sCache.WaitForCacheSync(ctx.Done(), gpm.podListerSynced); !ok { + gpm.logger.Fatal("failed to wait for caches to sync") + } + go gpm.service() + gpm.poolPodC.InjectGpm(gpm) go gpm.WebsocketStartEventChecker(gpm.kubernetesClient) go gpm.NoActiveConnectionEventChecker(gpm.kubernetesClient) go gpm.idleObjectReaper() + go gpm.poolPodC.Run(ctx.Done()) } func (gpm *GenericPoolManager) GetTypeName(ctx context.Context) fv1.ExecutorType { @@ -168,7 +178,7 @@ func (gpm *GenericPoolManager) GetFuncSvc(ctx context.Context, fn *fv1.Function) return nil, err } - pool, created, err := gpm.getPool(env) + pool, created, err := gpm.getPool(ctx, env) if err != nil { return nil, err } @@ -207,30 +217,12 @@ func (gpm *GenericPoolManager) TapService(ctx context.Context, svcHost string) e return nil } -func (gpm *GenericPoolManager) getPodInfo(ctx context.Context, obj apiv1.ObjectReference) (*apiv1.Pod, error) { - store := gpm.podInformer.GetStore() - - item, exists, err := store.Get(obj) - if err != nil || !exists { - item, exists, err = store.GetByKey(fmt.Sprintf("%s/%s", obj.Namespace, obj.Name)) - } - - if err != nil || !exists { - gpm.logger.Debug("Falling back to getting pod info from k8s API -- this may cause performance issues for your function.") - pod, err := gpm.kubernetesClient.CoreV1().Pods(obj.Namespace).Get(ctx, obj.Name, metav1.GetOptions{}) - return pod, err - } - - pod := item.(*apiv1.Pod) - return pod, nil -} - // IsValid checks if pod is not deleted and that it has the address passed as the argument. Also checks that all the // containers in it are reporting a ready status for the healthCheck. func (gpm *GenericPoolManager) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) bool { for _, obj := range fsvc.KubernetesObjects { if strings.ToLower(obj.Kind) == "pod" { - pod, err := gpm.getPodInfo(ctx, obj) + pod, err := gpm.podLister.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]. @@ -258,7 +250,7 @@ func (gpm *GenericPoolManager) RefreshFuncPods(ctx context.Context, logger *zap. return err } - gp, created, err := gpm.getPool(env) + gp, created, err := gpm.getPool(ctx, env) if err != nil { return err } @@ -314,7 +306,7 @@ func (gpm *GenericPoolManager) AdoptExistingResources(ctx context.Context) { wg.Add(1) go func() { defer wg.Done() - _, created, err := gpm.getPool(&env) + _, created, err := gpm.getPool(ctx, &env) if err != nil { gpm.logger.Error("adopt pool failed", zap.Error(err)) } @@ -461,7 +453,7 @@ func (gpm *GenericPoolManager) service() { // just because they are missing in the cache, we end up creating another duplicate pool. var err error created := false - pool, ok := gpm.pools[crd.CacheKey(&req.env.ObjectMeta)] + pool, ok := gpm.pools[crd.CacheKeyUID(&req.env.ObjectMeta)] if !ok { // To support backward compatibility, if envs are created in default ns, we go ahead // and create pools in fission-function ns as earlier. @@ -469,43 +461,47 @@ func (gpm *GenericPoolManager) service() { if req.env.ObjectMeta.Namespace != metav1.NamespaceDefault { ns = req.env.ObjectMeta.Namespace } - - pool, err = MakeGenericPool(gpm.logger, - gpm.fissionClient, gpm.kubernetesClient, gpm.metricsClient, req.env, ns, - gpm.namespace, gpm.fsCache, gpm.fetcherConfig, gpm.instanceID, gpm.enableIstio) + pool = MakeGenericPool(gpm.logger, gpm.fissionClient, gpm.kubernetesClient, + gpm.metricsClient, req.env, ns, gpm.namespace, gpm.fsCache, + gpm.fetcherConfig, gpm.instanceID, gpm.enableIstio) + err = pool.setup(req.ctx) if err != nil { req.responseChannel <- &response{error: err} continue } - gpm.pools[crd.CacheKey(&req.env.ObjectMeta)] = pool + gpm.pools[crd.CacheKeyUID(&req.env.ObjectMeta)] = pool created = true } req.responseChannel <- &response{pool: pool, created: created} - case CLEANUP_POOLS: - latestEnvPoolsize := make(map[string]int) - for _, env := range req.envList { - latestEnvPoolsize[crd.CacheKey(&env.ObjectMeta)] = int(getEnvPoolSize(&env)) + case CLEANUP_POOL: + env := *req.env + gpm.logger.Info("destroying pool", + zap.String("environment", env.ObjectMeta.Name), + zap.String("namespace", env.ObjectMeta.Namespace)) + + key := crd.CacheKeyUID(&req.env.ObjectMeta) + pool, ok := gpm.pools[key] + if !ok { + gpm.logger.Error("Could not find pool", zap.String("environment", env.ObjectMeta.Name), zap.String("namespace", env.ObjectMeta.Namespace)) + return } - for key, pool := range gpm.pools { - poolsize, ok := latestEnvPoolsize[key] - if !ok || poolsize == 0 { - // Env no longer exists or pool size changed to zero - - gpm.logger.Info("destroying generic pool", zap.Any("environment", pool.env.ObjectMeta)) - delete(gpm.pools, key) - - // and delete the pool asynchronously. - go pool.destroy() //nolint errcheck - } + delete(gpm.pools, key) + err := pool.destroy(req.ctx) + if err != nil { + gpm.logger.Error("failed to destroy pool", + zap.String("environment", env.ObjectMeta.Name), + zap.String("namespace", env.ObjectMeta.Namespace), + zap.Error(err)) } // no response, caller doesn't wait } } } -func (gpm *GenericPoolManager) getPool(env *fv1.Environment) (*GenericPool, bool, error) { +func (gpm *GenericPoolManager) getPool(ctx context.Context, env *fv1.Environment) (*GenericPool, bool, error) { c := make(chan *response) gpm.requestChannel <- &request{ + ctx: ctx, requestType: GET_POOL, env: env, responseChannel: c, @@ -514,10 +510,11 @@ func (gpm *GenericPoolManager) getPool(env *fv1.Environment) (*GenericPool, bool return resp.pool, resp.created, resp.error } -func (gpm *GenericPoolManager) cleanupPools(envs []fv1.Environment) { +func (gpm *GenericPoolManager) cleanupPool(ctx context.Context, env *fv1.Environment) { gpm.requestChannel <- &request{ - requestType: CLEANUP_POOLS, - envList: envs, + ctx: ctx, + requestType: CLEANUP_POOL, + env: env, } } @@ -551,53 +548,6 @@ func (gpm *GenericPoolManager) getFunctionEnv(ctx context.Context, fn *fv1.Funct return env, nil } -func (gpm *GenericPoolManager) eagerPoolCreator() { - pollSleep := 2 * time.Second - for { - // get list of envs from controller - envs, err := gpm.fissionClient.CoreV1().Environments(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{}) - if err != nil { - if utils.IsNetworkError(err) { - gpm.logger.Error("encountered network error, retrying", zap.Error(err)) - } else { - gpm.logger.Error("failed to get environment list", zap.Error(err)) - } - time.Sleep(5 * time.Second) - continue - } - - // Create pools for all envs. TODO: we should make this a bit less eager, only - // creating pools for envs that are actually used by functions. Also we might want - // to keep these eagerly created pools smaller than the ones created when there are - // actual function calls. - - wg := &sync.WaitGroup{} - - for i := range envs.Items { - env := envs.Items[i] - // Create pool only if poolsize greater than zero - if getEnvPoolSize(&env) > 0 { - wg.Add(1) - go func() { - defer wg.Done() - _, created, err := gpm.getPool(&env) - if err != nil { - gpm.logger.Error("eager-create pool failed", zap.Error(err)) - } - if created { - gpm.logger.Info("created pool for the environment", zap.String("env", env.ObjectMeta.Name), zap.String("namespace", gpm.namespace)) - } - }() - } - } - - // Clean up pools whose env was deleted - gpm.cleanupPools(envs.Items) - wg.Wait() - time.Sleep(pollSleep) - } -} - // idleObjectReaper reaps objects after certain idle time func (gpm *GenericPoolManager) idleObjectReaper() { ctx := context.Background() diff --git a/pkg/executor/executortype/poolmgr/poolpodcontroller.go b/pkg/executor/executortype/poolmgr/poolpodcontroller.go new file mode 100644 index 00000000..70248b46 --- /dev/null +++ b/pkg/executor/executortype/poolmgr/poolpodcontroller.go @@ -0,0 +1,414 @@ +/* +Copyright 2021 The Fission 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 poolmgr + +import ( + "context" + "strings" + "time" + + "go.uber.org/zap" + apps "k8s.io/api/apps/v1" + v1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "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" + "k8s.io/client-go/kubernetes" + k8sCache "k8s.io/client-go/tools/cache" + "k8s.io/client-go/util/workqueue" + + 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" + flisterv1 "github.com/fission/fission/pkg/generated/listers/core/v1" +) + +type ( + PoolPodController struct { + logger *zap.Logger + kubernetesClient *kubernetes.Clientset + namespace string + enableIstio bool + + envLister flisterv1.EnvironmentLister + envListerSynced k8sCache.InformerSynced + + envCreateUpdateQueue workqueue.RateLimitingInterface + envDeleteQueue workqueue.RateLimitingInterface + + spCleanupPodQueue workqueue.RateLimitingInterface + + gpm *GenericPoolManager + } +) + +func NewPoolPodController(logger *zap.Logger, + kubernetesClient *kubernetes.Clientset, + namespace string, + enableIstio bool, + funcInformer finformerv1.FunctionInformer, + pkgInformer finformerv1.PackageInformer, + envInformer finformerv1.EnvironmentInformer, + rsInformer appsinformers.ReplicaSetInformer) *PoolPodController { + p := &PoolPodController{ + logger: logger, + kubernetesClient: kubernetesClient, + namespace: namespace, + enableIstio: enableIstio, + + envCreateUpdateQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "EnvAddUpdateQueue"), + envDeleteQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "EnvDeleteQueue"), + spCleanupPodQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "SpecializedPodCleanupQueue"), + } + funcInformer.Informer().AddEventHandler(FunctionEventHandlers(p.logger, p.kubernetesClient, p.namespace, p.enableIstio)) + pkgInformer.Informer().AddEventHandler(PackageEventHandlers(p.logger, p.kubernetesClient, p.namespace)) + envInformer.Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ + AddFunc: p.enqueueEnvAdd, + UpdateFunc: p.enqueueEnvUpdate, + DeleteFunc: p.enqueueEnvDelete, + }) + rsInformer.Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ + AddFunc: p.handleRSAdd, + UpdateFunc: p.handleRSUpdate, + DeleteFunc: p.handleRSDelete, + }) + + p.envLister = envInformer.Lister() + p.envListerSynced = envInformer.Informer().HasSynced + p.logger.Info("pool pod controller handlers registered") + return p +} + +func (p *PoolPodController) InjectGpm(gpm *GenericPoolManager) { + p.gpm = gpm +} + +func IsPodActive(p *v1.Pod) bool { + return v1.PodSucceeded != p.Status.Phase && + v1.PodFailed != p.Status.Phase && + p.DeletionTimestamp == nil +} + +func (p *PoolPodController) processRS(rs *apps.ReplicaSet) { + if *(rs.Spec.Replicas) != 0 { + return + } + logger := p.logger.With(zap.String("rs", rs.Name), zap.String("namespace", rs.Namespace)) + logger.Debug("replica set has zero replica count") + // List all specialized pods and schedule for cleanup + rsLabelMap, err := metav1.LabelSelectorAsMap(rs.Spec.Selector) + if err != nil { + p.logger.Error("Failed to parse label selector", zap.Error(err)) + return + } + rsLabelMap["managed"] = "false" + specializedPods, err := p.gpm.podLister.Pods(rs.Namespace).List(labels.SelectorFromSet(rsLabelMap)) + if err != nil { + logger.Error("Failed to list specialized pods", zap.Error(err)) + } + if len(specializedPods) == 0 { + return + } + logger.Info("specialized pods identified for cleanup with RS", zap.Int("numPods", len(specializedPods))) + for _, pod := range specializedPods { + if !IsPodActive(pod) { + continue + } + key, err := k8sCache.MetaNamespaceKeyFunc(pod) + if err != nil { + logger.Error("Failed to get key for pod", zap.Error(err)) + continue + } + p.spCleanupPodQueue.Add(key) + } +} + +func (p *PoolPodController) handleRSAdd(obj interface{}) { + rs, ok := obj.(*apps.ReplicaSet) + if !ok { + p.logger.Error("unexpected type when adding rs to pool pod controller", zap.Any("obj", obj)) + return + } + p.processRS(rs) +} + +func (p *PoolPodController) handleRSUpdate(oldObj interface{}, newObj interface{}) { + rs, ok := newObj.(*apps.ReplicaSet) + if !ok { + p.logger.Error("unexpected type when updating rs to pool pod controller", zap.Any("obj", newObj)) + return + } + p.processRS(rs) +} + +func (p *PoolPodController) handleRSDelete(obj interface{}) { + rs, ok := obj.(*apps.ReplicaSet) + if !ok { + tombstone, ok := obj.(k8sCache.DeletedFinalStateUnknown) + if !ok { + p.logger.Error("couldnt get object from tombstone", zap.Any("obj", obj)) + return + } + rs, ok = tombstone.Obj.(*apps.ReplicaSet) + if !ok { + p.logger.Error("tombstone contained object that is not a replicaset", zap.Any("obj", obj)) + return + } + } + p.processRS(rs) +} + +func (p *PoolPodController) enqueueEnvAdd(obj interface{}) { + key, err := k8sCache.MetaNamespaceKeyFunc(obj) + if err != nil { + p.logger.Error("error retrieving key from object in poolPodController", zap.Any("obj", obj)) + return + } + p.logger.Debug("enqueue env add", zap.String("key", key)) + p.envCreateUpdateQueue.Add(key) +} + +func (p *PoolPodController) enqueueEnvUpdate(oldObj, newObj interface{}) { + key, err := k8sCache.MetaNamespaceKeyFunc(newObj) + if err != nil { + p.logger.Error("error retrieving key from object in poolPodController", zap.Any("obj", key)) + return + } + p.logger.Debug("enqueue env update", zap.String("key", key)) + p.envCreateUpdateQueue.Add(key) +} + +func (p *PoolPodController) enqueueEnvDelete(obj interface{}) { + env, ok := obj.(*fv1.Environment) + if !ok { + p.logger.Error("unexpected type when deleting env to pool pod controller", zap.Any("obj", obj)) + return + } + p.logger.Debug("enqueue env delete", zap.Any("env", env)) + p.envDeleteQueue.Add(env) +} + +func (p *PoolPodController) Run(stopCh <-chan struct{}) { + defer utilruntime.HandleCrash() + defer p.envCreateUpdateQueue.ShutDown() + defer p.envDeleteQueue.ShutDown() + defer p.spCleanupPodQueue.ShutDown() + + // Wait for the caches to be synced before starting workers + p.logger.Info("Waiting for informer caches to sync") + if ok := k8sCache.WaitForCacheSync(stopCh, p.envListerSynced); !ok { + p.logger.Fatal("failed to wait for caches to sync") + } + for i := 0; i < 4; i++ { + go wait.Until(p.workerRun("envCreateUpdate", p.envCreateUpdateQueueProcessFunc), time.Second, stopCh) + } + go wait.Until(p.workerRun("envDeleteQueue", p.envDeleteQueueProcessFunc), time.Second, stopCh) + go wait.Until(p.workerRun("spCleanupPodQueue", p.spCleanupPodQueueProcessFunc), time.Second, stopCh) + p.logger.Info("Started workers for poolPodController") + <-stopCh + p.logger.Info("Shutting down workers for poolPodController") +} + +func (p *PoolPodController) workerRun(name string, processFunc func() bool) func() { + return func() { + p.logger.Debug("Starting worker with func", zap.String("name", name)) + for { + if quit := processFunc(); quit { + p.logger.Info("Shutting down worker", zap.String("name", name)) + return + } + } + } +} + +func (p *PoolPodController) envCreateUpdateQueueProcessFunc() bool { + maxRetries := 3 + handleEnv := func(ctx context.Context, env *fv1.Environment) error { + log := p.logger.With(zap.String("env", env.ObjectMeta.Name), zap.String("namespace", env.ObjectMeta.Namespace)) + log.Debug("env reconsile request processing") + pool, created, err := p.gpm.getPool(ctx, env) + if err != nil { + log.Error("error getting pool", zap.Error(err)) + return err + } + if created { + log.Info("created pool for the environment") + return nil + } + poolsize := getEnvPoolSize(env) + if poolsize == 0 { + log.Info("pool size is zero") + p.gpm.cleanupPool(ctx, env) + return nil + } + err = pool.updatePoolDeployment(ctx, env) + if err != nil { + log.Error("error updating pool", zap.Error(err)) + return err + } + // If any specialized pods are running, those would be + // deleted by replicaSet controller. + return nil + } + + obj, quit := p.envCreateUpdateQueue.Get() + if quit { + return true + } + key := obj.(string) + defer p.envCreateUpdateQueue.Done(key) + + namespace, name, err := k8sCache.SplitMetaNamespaceKey(key) + if err != nil { + p.logger.Error("error splitting key", zap.Error(err)) + p.envCreateUpdateQueue.Forget(key) + return false + } + env, err := p.envLister.Environments(namespace).Get(name) + if apierrors.IsNotFound(err) { + p.logger.Info("env not found", zap.String("key", key)) + p.envCreateUpdateQueue.Forget(key) + return false + } + + if err != nil { + if p.envCreateUpdateQueue.NumRequeues(key) < maxRetries { + p.envCreateUpdateQueue.AddRateLimited(key) + p.logger.Error("error getting env, retrying", zap.Error(err)) + } else { + p.envCreateUpdateQueue.Forget(key) + p.logger.Error("error getting env, retrying, max retries reached", zap.Error(err)) + } + return false + } + + ctx := context.Background() + err = handleEnv(ctx, env) + if err != nil { + if p.envCreateUpdateQueue.NumRequeues(key) < maxRetries { + p.envCreateUpdateQueue.AddRateLimited(key) + p.logger.Error("error handling env from envInformer, retrying", zap.String("key", key), zap.Error(err)) + } else { + p.envCreateUpdateQueue.Forget(key) + p.logger.Error("error handling env from envInformer, max retries reached", zap.String("key", key), zap.Error(err)) + } + return false + } + p.envCreateUpdateQueue.Forget(key) + return false +} + +func (p *PoolPodController) envDeleteQueueProcessFunc() bool { + obj, quit := p.envDeleteQueue.Get() + if quit { + return true + } + defer p.envDeleteQueue.Done(obj) + env, ok := obj.(*fv1.Environment) + if !ok { + p.logger.Error("unexpected type when deleting env to pool pod controller", zap.Any("obj", obj)) + p.envDeleteQueue.Forget(obj) + return false + } + ctx := context.Background() + p.logger.Debug("env delete request processing") + p.gpm.cleanupPool(ctx, env) + specializePodLables := getSpecializedPodLabels(env) + specializedPods, err := p.gpm.podLister.Pods(p.gpm.namespace).List(labels.SelectorFromSet(specializePodLables)) + if err != nil { + p.logger.Error("failed to list specialized pods", zap.Error(err)) + p.envDeleteQueue.Forget(obj) + return false + } + if len(specializedPods) == 0 { + p.envDeleteQueue.Forget(obj) + return false + } + p.logger.Info("specialized pods identified for cleanup after env delete", zap.String("env", env.ObjectMeta.Name), zap.String("namespace", env.ObjectMeta.Namespace), zap.Int("count", len(specializedPods))) + for _, pod := range specializedPods { + if !IsPodActive(pod) { + continue + } + key, err := k8sCache.MetaNamespaceKeyFunc(pod) + if err != nil { + p.logger.Error("Failed to get key for pod", zap.Error(err)) + continue + } + p.spCleanupPodQueue.Add(key) + } + p.envDeleteQueue.Forget(obj) + return false +} + +func (p *PoolPodController) spCleanupPodQueueProcessFunc() bool { + maxRetries := 3 + obj, quit := p.spCleanupPodQueue.Get() + if quit { + return true + } + key := obj.(string) + defer p.spCleanupPodQueue.Done(key) + namespace, name, err := k8sCache.SplitMetaNamespaceKey(key) + if err != nil { + p.logger.Error("error splitting key", zap.Error(err)) + p.spCleanupPodQueue.Forget(key) + return false + } + pod, err := p.gpm.podLister.Pods(namespace).Get(name) + if apierrors.IsNotFound(err) { + p.logger.Info("pod not found", zap.String("key", key)) + p.spCleanupPodQueue.Forget(key) + return false + } + if !IsPodActive(pod) { + p.logger.Info("pod not active", zap.String("key", key)) + p.spCleanupPodQueue.Forget(key) + return false + } + if err != nil { + if p.spCleanupPodQueue.NumRequeues(key) < maxRetries { + p.spCleanupPodQueue.AddRateLimited(key) + p.logger.Error("error getting pod, retrying", zap.Error(err)) + } else { + p.spCleanupPodQueue.Forget(key) + p.logger.Error("error getting pod, max retries reached", zap.Error(err)) + } + return false + } + podName := strings.SplitAfter(pod.GetName(), ".") + if fsvc, ok := p.gpm.fsCache.PodToFsvc.Load(strings.TrimSuffix(podName[0], ".")); ok { + fsvc, ok := fsvc.(*fscache.FuncSvc) + if ok { + p.gpm.fsCache.DeleteFunctionSvc(fsvc) + p.gpm.fsCache.DeleteEntry(fsvc) + } else { + p.logger.Error("could not covert item from PodToFsvc", zap.String("key", key)) + } + } + err = p.kubernetesClient.CoreV1().Pods(p.namespace).Delete(context.TODO(), pod.Name, metav1.DeleteOptions{}) + if err != nil { + p.logger.Error("failed to delete pod", zap.Error(err), zap.String("pod", pod.ObjectMeta.Name), zap.String("pod_namespace", pod.ObjectMeta.Namespace)) + return false + } + p.logger.Info("cleaned specialized pod as environment update/deleted", + zap.String("pod", pod.ObjectMeta.Name), zap.String("pod_namespace", pod.ObjectMeta.Namespace), + zap.String("address", pod.Status.PodIP)) + p.spCleanupPodQueue.Forget(key) + return false +} diff --git a/pkg/utils/informer.go b/pkg/utils/informer.go index cceeed88..662b72ac 100644 --- a/pkg/utils/informer.go +++ b/pkg/utils/informer.go @@ -1,28 +1,26 @@ package utils import ( - "fmt" + "time" - apiv1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/selection" k8sInformers "k8s.io/client-go/informers" "k8s.io/client-go/kubernetes" - k8sCache "k8s.io/client-go/tools/cache" metricsapi "k8s.io/metrics/pkg/apis/metrics" v1 "github.com/fission/fission/pkg/apis/core/v1" ) -func GetInformerFactoryByExecutor(client *kubernetes.Clientset, executorType v1.ExecutorType) (k8sInformers.SharedInformerFactory, error) { +func GetInformerFactoryByExecutor(client *kubernetes.Clientset, executorType v1.ExecutorType, defaultResync time.Duration) (k8sInformers.SharedInformerFactory, error) { executorLabel, err := labels.NewRequirement(v1.EXECUTOR_TYPE, selection.DoubleEquals, []string{string(executorType)}) if err != nil { return nil, err } labelSelector := labels.NewSelector() labelSelector.Add(*executorLabel) - informerFactory := k8sInformers.NewSharedInformerFactoryWithOptions(client, 0, + informerFactory := k8sInformers.NewSharedInformerFactoryWithOptions(client, defaultResync, k8sInformers.WithTweakListOptions(func(options *metav1.ListOptions) { options.LabelSelector = labelSelector.String() })) @@ -47,14 +45,3 @@ func SupportedMetricsAPIVersionAvailable(discoveredAPIGroups *metav1.APIGroupLis } return false } - -func GetCachedItem(obj apiv1.ObjectReference, informer k8sCache.SharedIndexInformer) (item interface{}, exists bool, err error) { - store := informer.GetStore() - - item, exists, err = store.Get(obj) - if err != nil || !exists { - item, exists, err = store.GetByKey(fmt.Sprintf("%s/%s", obj.Namespace, obj.Name)) - } - - return item, exists, err -}