diff --git a/pkg/executor/executor.go b/pkg/executor/executor.go index c3aedda3..86e52e2a 100644 --- a/pkg/executor/executor.go +++ b/pkg/executor/executor.go @@ -244,18 +244,22 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames executorTypes[gpm.GetTypeName()] = gpm executorTypes[ndm.GetTypeName()] = ndm - if ok, _ := strconv.ParseBool(os.Getenv("ADOPT_EXISTING_RESOURCES")); ok { - wg := &sync.WaitGroup{} - for _, et := range executorTypes { - wg.Add(1) - go func(et executortype.ExecutorType) { - defer wg.Done() + adoptExistingResources, _ := strconv.ParseBool(os.Getenv("ADOPT_EXISTING_RESOURCES")) + + wg := &sync.WaitGroup{} + for _, et := range executorTypes { + wg.Add(1) + go func(et executortype.ExecutorType) { + defer wg.Done() + if adoptExistingResources { et.AdoptExistingResources() - }(et) - } - // set hard timeout for resource adoption - util.WaitTimeout(wg, 30*time.Second) + } + et.CleanupOldExecutorObjects() + }(et) } + // set hard timeout for resource adoption + // TODO: use context to control the waiting time once kubernetes client supports it. + util.WaitTimeout(wg, 30*time.Second) cms := cms.MakeConfigSecretController(logger, fissionClient, kubernetesClient, executorTypes) @@ -264,7 +268,6 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames return err } - go reaper.CleanupOldExecutorObjects(logger, kubernetesClient, executorInstanceID) go reaper.CleanupRoleBindings(logger, kubernetesClient, fissionClient, functionNamespace, envBuilderNamespace, time.Minute*30) go api.Serve(port) go serveMetric(logger) diff --git a/pkg/executor/executortype/executortype.go b/pkg/executor/executortype/executortype.go index ad126738..864d2aae 100644 --- a/pkg/executor/executortype/executortype.go +++ b/pkg/executor/executortype/executortype.go @@ -53,4 +53,7 @@ type ExecutorType interface { // AdoptOrphanResources adopts existing resources created by the deleted executor. AdoptExistingResources() + + // CleanupOldExecutorObjects cleans up resources created by old executor instances + CleanupOldExecutorObjects() } diff --git a/pkg/executor/executortype/newdeploy/newdeploy.go b/pkg/executor/executortype/newdeploy/newdeploy.go index bea5f8cd..9aa19d79 100644 --- a/pkg/executor/executortype/newdeploy/newdeploy.go +++ b/pkg/executor/executortype/newdeploy/newdeploy.go @@ -101,23 +101,24 @@ func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Enviro depl, err := deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Create(deployment) if err != nil { - deploy.logger.Error("error while creating function deployment", - zap.Error(err), - zap.String("function", fn.Metadata.Name), - zap.String("deployment_name", deployName), - zap.String("deployment_namespace", deployNamespace)) - return nil, err + if k8s_err.IsAlreadyExists(err) { + depl, err = deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Get(deployName, metav1.GetOptions{}) + } + if err != nil { + deploy.logger.Error("error while creating function deployment", + zap.Error(err), + zap.String("function", fn.Metadata.Name), + zap.String("deployment_name", deployName), + zap.String("deployment_namespace", deployNamespace)) + return nil, err + } } - if minScale > 0 { depl, err = deploy.waitForDeploy(depl, minScale, specializationTimeout) } - return depl, err } - return nil, err - } func (deploy *NewDeploy) setupRBACObjs(deployNamespace string, fn *fv1.Function) error { @@ -403,13 +404,16 @@ func (deploy *NewDeploy) createOrGetHpa(hpaName string, execStrategy *fv1.Execut } else if k8s_err.IsNotFound(err) { cHpa, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Create(hpa) if err != nil { - return nil, err + if k8s_err.IsAlreadyExists(err) { + cHpa, err = deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Get(hpaName, metav1.GetOptions{}) + } + if err != nil { + return nil, err + } } return cHpa, nil } - return nil, err - } func (deploy *NewDeploy) getHpa(ns, name string) (*asv1.HorizontalPodAutoscaler, error) { @@ -465,7 +469,12 @@ func (deploy *NewDeploy) createOrGetSvc(deployLabels map[string]string, deployAn } else if k8s_err.IsNotFound(err) { svc, err := deploy.kubernetesClient.CoreV1().Services(svcNamespace).Create(service) if err != nil { - return nil, err + if k8s_err.IsAlreadyExists(err) { + svc, err = deploy.kubernetesClient.CoreV1().Services(svcNamespace).Get(svcName, metav1.GetOptions{}) + } + if err != nil { + return nil, err + } } return svc, nil } diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr.go b/pkg/executor/executortype/newdeploy/newdeploymgr.go index 1962b79f..4d5d44d6 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr.go @@ -43,6 +43,7 @@ import ( "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/executor/executortype" "github.com/fission/fission/pkg/executor/fscache" + "github.com/fission/fission/pkg/executor/reaper" fetcherConfig "github.com/fission/fission/pkg/fetcher/config" "github.com/fission/fission/pkg/throttler" "github.com/fission/fission/pkg/types" @@ -272,6 +273,35 @@ func (deploy *NewDeploy) AdoptExistingResources() { wg.Wait() } +func (deploy *NewDeploy) CleanupOldExecutorObjects() { + deploy.logger.Info("Newdeploy starts to clean orphaned resources", zap.String("instanceID", deploy.instanceID)) + + errs := &multierror.Error{} + listOpts := metav1.ListOptions{ + LabelSelector: labels.Set(map[string]string{types.EXECUTOR_TYPE: string(fv1.ExecutorTypeNewdeploy)}).AsSelector().String(), + } + + err := reaper.CleanupHpa(deploy.logger, deploy.kubernetesClient, deploy.instanceID, listOpts) + if err != nil { + errs = multierror.Append(errs, err) + } + + err = reaper.CleanupDeployments(deploy.logger, deploy.kubernetesClient, deploy.instanceID, listOpts) + if err != nil { + errs = multierror.Append(errs, err) + } + + err = reaper.CleanupServices(deploy.logger, deploy.kubernetesClient, deploy.instanceID, listOpts) + if err != nil { + errs = multierror.Append(errs, err) + } + + if errs.ErrorOrNil() != nil { + // TODO retry reaper; logged and ignored for now + deploy.logger.Error("Failed to cleanup old executor objects", zap.Error(err)) + } +} + func (deploy *NewDeploy) initFuncController() (k8sCache.Store, k8sCache.Controller) { resyncPeriod := 30 * time.Second listWatch := k8sCache.NewListWatchFromClient(deploy.crdClient, "functions", metav1.NamespaceAll, fields.Everything()) diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index b435450e..d5056a42 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -26,6 +26,7 @@ import ( "sync" "time" + "github.com/hashicorp/go-multierror" "go.uber.org/zap" apiv1 "k8s.io/api/core/v1" k8serrors "k8s.io/apimachinery/pkg/api/errors" @@ -376,6 +377,30 @@ func (gpm *GenericPoolManager) AdoptExistingResources() { wg.Wait() } +func (gpm *GenericPoolManager) CleanupOldExecutorObjects() { + gpm.logger.Info("Poolmanager starts to clean orphaned resources", zap.String("instanceID", gpm.instanceId)) + + errs := &multierror.Error{} + listOpts := metav1.ListOptions{ + LabelSelector: labels.Set(map[string]string{types.EXECUTOR_TYPE: string(fv1.ExecutorTypePoolmgr)}).AsSelector().String(), + } + + err := reaper.CleanupDeployments(gpm.logger, gpm.kubernetesClient, gpm.instanceId, listOpts) + if err != nil { + errs = multierror.Append(errs, err) + } + + err = reaper.CleanupPods(gpm.logger, gpm.kubernetesClient, gpm.instanceId, listOpts) + if err != nil { + errs = multierror.Append(errs, err) + } + + if errs.ErrorOrNil() != nil { + // TODO retry reaper; logged and ignored for now + gpm.logger.Error("Failed to cleanup old executor objects", zap.Error(err)) + } +} + func (gpm *GenericPoolManager) service() { for { req := <-gpm.requestChannel diff --git a/pkg/executor/reaper/reaper.go b/pkg/executor/reaper/reaper.go index fe237710..40e90b3c 100644 --- a/pkg/executor/reaper/reaper.go +++ b/pkg/executor/reaper/reaper.go @@ -20,7 +20,6 @@ import ( "strings" "time" - "github.com/hashicorp/go-multierror" "go.uber.org/zap" apiv1 "k8s.io/api/core/v1" meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -37,43 +36,6 @@ var ( delOpt = meta_v1.DeleteOptions{PropagationPolicy: &deletePropagation} ) -// CleanupOldExecutorObjects cleans up resources created by old executor instances -func CleanupOldExecutorObjects(logger *zap.Logger, kubernetesClient *kubernetes.Clientset, instanceId string) { - errs := &multierror.Error{} - - err := cleanupHpa(logger, kubernetesClient, instanceId) - if err != nil { - errs = multierror.Append(errs, err) - } - - err = cleanupDeployments(logger, kubernetesClient, instanceId) - if err != nil { - errs = multierror.Append(errs, err) - } - - // Pods might be running user functions, so we give them - // a few minutes before terminating them. This time is the - // maximum function runtime, plus the time a router might - // still route to an old instance, i.e. router cache expiry - // time. - time.Sleep(6 * time.Minute) - - err = cleanupServices(logger, kubernetesClient, instanceId) - if err != nil { - errs = multierror.Append(errs, err) - } - - err = cleanupPods(logger, kubernetesClient, instanceId) - if err != nil { - errs = multierror.Append(errs, err) - } - - if errs.ErrorOrNil() != nil { - // TODO retry reaper; logged and ignored for now - logger.Error("Failed to cleanup old executor objects", zap.Error(err)) - } -} - // CleanupKubeObject deletes given kubernetes object func CleanupKubeObject(logger *zap.Logger, kubeClient *kubernetes.Clientset, kubeobj *apiv1.ObjectReference) { switch strings.ToLower(kubeobj.Kind) { @@ -107,8 +69,8 @@ func CleanupKubeObject(logger *zap.Logger, kubeClient *kubernetes.Clientset, kub } } -func cleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instanceId string) error { - deploymentList, err := client.AppsV1().Deployments(meta_v1.NamespaceAll).List(meta_v1.ListOptions{}) +func CleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instanceId string, listOps meta_v1.ListOptions) error { + deploymentList, err := client.AppsV1().Deployments(meta_v1.NamespaceAll).List(listOps) if err != nil { return err } @@ -119,7 +81,7 @@ func cleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instan id, ok = dep.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL] } if ok && id != instanceId { - logger.Debug("cleaning up deployment", zap.String("deployment", dep.ObjectMeta.Name)) + logger.Info("cleaning up deployment", zap.String("deployment", dep.ObjectMeta.Name)) err := client.AppsV1().Deployments(dep.ObjectMeta.Namespace).Delete(dep.ObjectMeta.Name, &delOpt) if err != nil { logger.Error("error cleaning up deployment", @@ -133,8 +95,8 @@ func cleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instan return nil } -func cleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceId string) error { - podList, err := client.CoreV1().Pods(meta_v1.NamespaceAll).List(meta_v1.ListOptions{}) +func CleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceId string, listOps meta_v1.ListOptions) error { + podList, err := client.CoreV1().Pods(meta_v1.NamespaceAll).List(listOps) if err != nil { return err } @@ -145,7 +107,7 @@ func cleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceId st id, ok = pod.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL] } if ok && id != instanceId { - logger.Debug("cleaning up pod", zap.String("pod", pod.ObjectMeta.Name)) + logger.Info("cleaning up pod", zap.String("pod", pod.ObjectMeta.Name)) err := client.CoreV1().Pods(pod.ObjectMeta.Namespace).Delete(pod.ObjectMeta.Name, nil) if err != nil { logger.Error("error cleaning up pod", @@ -159,8 +121,8 @@ func cleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceId st return nil } -func cleanupServices(logger *zap.Logger, client *kubernetes.Clientset, instanceId string) error { - svcList, err := client.CoreV1().Services(meta_v1.NamespaceAll).List(meta_v1.ListOptions{}) +func CleanupServices(logger *zap.Logger, client *kubernetes.Clientset, instanceId string, listOps meta_v1.ListOptions) error { + svcList, err := client.CoreV1().Services(meta_v1.NamespaceAll).List(listOps) if err != nil { return err } @@ -171,7 +133,7 @@ func cleanupServices(logger *zap.Logger, client *kubernetes.Clientset, instanceI id, ok = svc.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL] } if ok && id != instanceId { - logger.Debug("cleaning up service", zap.String("service", svc.ObjectMeta.Name)) + logger.Info("cleaning up service", zap.String("service", svc.ObjectMeta.Name)) err := client.CoreV1().Services(svc.ObjectMeta.Namespace).Delete(svc.ObjectMeta.Name, nil) if err != nil { logger.Error("error cleaning up service", @@ -185,8 +147,8 @@ func cleanupServices(logger *zap.Logger, client *kubernetes.Clientset, instanceI return nil } -func cleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceId string) error { - hpaList, err := client.AutoscalingV1().HorizontalPodAutoscalers(meta_v1.NamespaceAll).List(meta_v1.ListOptions{}) +func CleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceId string, listOps meta_v1.ListOptions) error { + hpaList, err := client.AutoscalingV1().HorizontalPodAutoscalers(meta_v1.NamespaceAll).List(listOps) if err != nil { return err } @@ -198,7 +160,7 @@ func cleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceId str id, ok = hpa.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL] } if ok && id != instanceId { - logger.Debug("cleaning up HPA", zap.String("hpa", hpa.ObjectMeta.Name)) + logger.Info("cleaning up HPA", zap.String("hpa", hpa.ObjectMeta.Name)) err := client.AutoscalingV1().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Delete(hpa.ObjectMeta.Name, nil) if err != nil { logger.Error("error cleaning up HPA",