diff --git a/pkg/buildermgr/common.go b/pkg/buildermgr/common.go index e4f98f48..46abfc93 100644 --- a/pkg/buildermgr/common.go +++ b/pkg/buildermgr/common.go @@ -47,7 +47,7 @@ import ( func buildPackage(ctx context.Context, logger *zap.Logger, fissionClient *crd.FissionClient, envBuilderNamespace string, storageSvcUrl string, pkg *fv1.Package) (uploadResp *fetcher.ArchiveUploadResponse, buildLogs string, err error) { - env, err := fissionClient.CoreV1().Environments(pkg.Spec.Environment.Namespace).Get(context.TODO(), pkg.Spec.Environment.Name, metav1.GetOptions{}) + env, err := fissionClient.CoreV1().Environments(pkg.Spec.Environment.Namespace).Get(ctx, pkg.Spec.Environment.Name, metav1.GetOptions{}) if err != nil { e := "error getting environment CRD info" logger.Error(e, zap.Error(err)) diff --git a/pkg/buildermgr/pkgwatcher.go b/pkg/buildermgr/pkgwatcher.go index 13ef25b7..3d87e451 100644 --- a/pkg/buildermgr/pkgwatcher.go +++ b/pkg/buildermgr/pkgwatcher.go @@ -73,7 +73,7 @@ func makePackageWatcher(logger *zap.Logger, fissionClient *crd.FissionClient, k8 // 5. Update package resource in package ref of functions that share the same package // 6. Update package status to succeed state // *. Update package status to failed state,if any one of steps above failed/time out -func (pkgw *packageWatcher) build(srcpkg *fv1.Package) { +func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) { // Ignore duplicate build requests key := fmt.Sprintf("%v-%v", srcpkg.ObjectMeta.Name, srcpkg.ObjectMeta.ResourceVersion) _, err := pkgw.buildCache.Set(key, srcpkg) @@ -95,7 +95,7 @@ func (pkgw *packageWatcher) build(srcpkg *fv1.Package) { return } - env, err := pkgw.fissionClient.CoreV1().Environments(pkg.Spec.Environment.Namespace).Get(context.TODO(), pkg.Spec.Environment.Name, metav1.GetOptions{}) + env, err := pkgw.fissionClient.CoreV1().Environments(pkg.Spec.Environment.Namespace).Get(ctx, pkg.Spec.Environment.Name, metav1.GetOptions{}) if k8serrors.IsNotFound(err) { e := "environment does not exist" pkgw.logger.Error(e, zap.String("environment", pkg.Spec.Environment.Name)) @@ -167,7 +167,7 @@ func (pkgw *packageWatcher) build(srcpkg *fv1.Package) { // Add the package getter rolebinding to builder sa // we continue here if role binding was not setup successfully. this is because without this, the fetcher wont be able to fetch the source pkg into the container and // the build will fail eventually - err := utils.SetupRoleBinding(pkgw.logger, pkgw.k8sClient, fv1.PackageGetterRB, pkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionBuilderSA, builderNs) + err := utils.SetupRoleBinding(ctx, pkgw.logger, pkgw.k8sClient, fv1.PackageGetterRB, pkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionBuilderSA, builderNs) if err != nil { pkgw.logger.Error("error setting up role binding for package", zap.Error(err), @@ -181,7 +181,6 @@ func (pkgw *packageWatcher) build(srcpkg *fv1.Package) { zap.String("package", fmt.Sprintf("%s.%s", pkg.ObjectMeta.Name, pkg.ObjectMeta.Namespace))) } - ctx := context.Background() uploadResp, buildLogs, err := buildPackage(ctx, pkgw.logger, pkgw.fissionClient, builderNs, pkgw.storageSvcUrl, pkg) if err != nil { pkgw.logger.Error("error building package", zap.Error(err), zap.String("package_name", pkg.ObjectMeta.Name)) @@ -200,7 +199,7 @@ func (pkgw *packageWatcher) build(srcpkg *fv1.Package) { pkgw.logger.Info("starting package info update", zap.String("package_name", pkg.ObjectMeta.Name)) fnList, err := pkgw.fissionClient.CoreV1(). - Functions(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{}) + Functions(metav1.NamespaceAll).List(ctx, metav1.ListOptions{}) if err != nil { e := "error getting function list" pkgw.logger.Error(e, zap.Error(err)) @@ -224,7 +223,7 @@ func (pkgw *packageWatcher) build(srcpkg *fv1.Package) { fn.Spec.Package.PackageRef.ResourceVersion != pkg.ObjectMeta.ResourceVersion { fn.Spec.Package.PackageRef.ResourceVersion = pkg.ObjectMeta.ResourceVersion // update CRD - _, err = pkgw.fissionClient.CoreV1().Functions(fn.ObjectMeta.Namespace).Update(context.TODO(), &fn, metav1.UpdateOptions{}) + _, err = pkgw.fissionClient.CoreV1().Functions(fn.ObjectMeta.Namespace).Update(ctx, &fn, metav1.UpdateOptions{}) if err != nil { e := "error updating function package resource version" pkgw.logger.Error(e, zap.Error(err)) @@ -296,7 +295,8 @@ func (pkgw *packageWatcher) packageInformerHandler() k8sCache.ResourceEventHandl } // Only build pending state packages. if pkg.Status.BuildStatus == fv1.BuildStatusPending { - go pkgw.build(pkg) + ctx := context.Background() + go pkgw.build(ctx, pkg) } } return k8sCache.ResourceEventHandlerFuncs{ diff --git a/pkg/executor/executortype/container/common.go b/pkg/executor/executortype/container/common.go index 70087355..695d5a79 100644 --- a/pkg/executor/executortype/container/common.go +++ b/pkg/executor/executortype/container/common.go @@ -65,10 +65,10 @@ func (cn *Container) getResources(fn *fv1.Function) apiv1.ResourceRequirements { } // cleanupContainer cleans all kubernetes objects related to function -func (cn *Container) cleanupContainer(ns string, name string) error { +func (cn *Container) cleanupContainer(ctx context.Context, ns string, name string) error { result := &multierror.Error{} - err := cn.deleteSvc(ns, name) + err := cn.deleteSvc(ctx, ns, name) if err != nil && !k8s_err.IsNotFound(err) { cn.logger.Error("error deleting service for Container function", zap.Error(err), @@ -77,7 +77,7 @@ func (cn *Container) cleanupContainer(ns string, name string) error { result = multierror.Append(result, err) } - err = cn.deleteHpa(ns, name) + err = cn.deleteHpa(ctx, ns, name) if err != nil && !k8s_err.IsNotFound(err) { cn.logger.Error("error deleting HPA for Container function", zap.Error(err), @@ -86,7 +86,7 @@ func (cn *Container) cleanupContainer(ns string, name string) error { result = multierror.Append(result, err) } - err = cn.deleteDeployment(ns, name) + err = cn.deleteDeployment(ctx, ns, name) if err != nil && !k8s_err.IsNotFound(err) { cn.logger.Error("error deleting deployment for Container function", zap.Error(err), diff --git a/pkg/executor/executortype/container/containermgr.go b/pkg/executor/executortype/container/containermgr.go index cecfb17f..f249be58 100644 --- a/pkg/executor/executortype/container/containermgr.go +++ b/pkg/executor/executortype/container/containermgr.go @@ -174,7 +174,7 @@ func (caaf *Container) TapService(ctx context.Context, svcHost string) error { return nil } -func (caaf *Container) getServiceInfo(obj apiv1.ObjectReference) (*apiv1.Service, error) { +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 { @@ -183,7 +183,7 @@ func (caaf *Container) getServiceInfo(obj apiv1.ObjectReference) (*apiv1.Service zap.Bool("exists", exists), zap.Error(err), ) - service, err := caaf.kubernetesClient.CoreV1().Services(obj.Namespace).Get(context.TODO(), obj.Name, metav1.GetOptions{}) + service, err := caaf.kubernetesClient.CoreV1().Services(obj.Namespace).Get(ctx, obj.Name, metav1.GetOptions{}) return service, err } @@ -191,7 +191,7 @@ func (caaf *Container) getServiceInfo(obj apiv1.ObjectReference) (*apiv1.Service return service, nil } -func (caaf *Container) getDeploymentInfo(obj apiv1.ObjectReference) (*appsv1.Deployment, error) { +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 { @@ -200,7 +200,7 @@ func (caaf *Container) getDeploymentInfo(obj apiv1.ObjectReference) (*appsv1.Dep zap.Bool("exists", exists), zap.Error(err), ) - deployment, err := caaf.kubernetesClient.AppsV1().Deployments(obj.Namespace).Get(context.TODO(), obj.Name, metav1.GetOptions{}) + deployment, err := caaf.kubernetesClient.AppsV1().Deployments(obj.Namespace).Get(ctx, obj.Name, metav1.GetOptions{}) return deployment, err } @@ -222,7 +222,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.getServiceInfo(obj) + _, err := caaf.getServiceInfo(ctx, obj) if err != nil { if !k8sErrs.IsNotFound(err) { caaf.logger.Error("error validating function service", zap.String("function", fsvc.Function.Name), zap.Error(err)) @@ -231,7 +231,7 @@ func (caaf *Container) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) bool } } else if strings.ToLower(obj.Kind) == "deployment" { - currentDeploy, err := caaf.getDeploymentInfo(obj) + currentDeploy, err := caaf.getDeploymentInfo(ctx, obj) if err != nil { if !k8sErrs.IsNotFound(err) { caaf.logger.Error("error validating function deployment", zap.String("function", fsvc.Function.Name), zap.Error(err)) @@ -317,17 +317,17 @@ func (caaf *Container) CleanupOldExecutorObjects(ctx context.Context) { LabelSelector: labels.Set(map[string]string{fv1.EXECUTOR_TYPE: string(fv1.ExecutorTypeContainer)}).AsSelector().String(), } - err := reaper.CleanupHpa(caaf.logger, caaf.kubernetesClient, caaf.instanceID, listOpts) + err := reaper.CleanupHpa(ctx, caaf.logger, caaf.kubernetesClient, caaf.instanceID, listOpts) if err != nil { errs = multierror.Append(errs, err) } - err = reaper.CleanupDeployments(caaf.logger, caaf.kubernetesClient, caaf.instanceID, listOpts) + err = reaper.CleanupDeployments(ctx, caaf.logger, caaf.kubernetesClient, caaf.instanceID, listOpts) if err != nil { errs = multierror.Append(errs, err) } - err = reaper.CleanupServices(caaf.logger, caaf.kubernetesClient, caaf.instanceID, listOpts) + err = reaper.CleanupServices(ctx, caaf.logger, caaf.kubernetesClient, caaf.instanceID, listOpts) if err != nil { errs = multierror.Append(errs, err) } @@ -366,11 +366,11 @@ func (caaf *Container) createFunction(ctx context.Context, fn *fv1.Function) (*f return fsvc, err } -func (caaf *Container) deleteFunction(fn *fv1.Function) error { +func (caaf *Container) deleteFunction(ctx context.Context, fn *fv1.Function) error { if fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType != fv1.ExecutorTypeContainer { return nil } - err := caaf.fnDelete(fn) + err := caaf.fnDelete(ctx, fn) if err != nil { err = errors.Wrapf(err, "error deleting kubernetes objects of function %v", fn.ObjectMeta) } @@ -379,7 +379,7 @@ func (caaf *Container) deleteFunction(fn *fv1.Function) error { func (caaf *Container) fnCreate(ctx context.Context, fn *fv1.Function) (*fscache.FuncSvc, error) { cleanupFunc := func(ns string, name string) { - err := caaf.cleanupContainer(ns, name) + err := caaf.cleanupContainer(ctx, ns, name) if err != nil { caaf.logger.Error("received error while cleaning function resources", zap.String("namespace", ns), zap.String("name", name)) @@ -401,7 +401,7 @@ func (caaf *Container) fnCreate(ctx context.Context, fn *fv1.Function) (*fscache // Since Container waits for pods of deployment to be ready, // change the order of kubeObject creation (create service first, // then deployment) to take advantage of waiting time. - svc, err := caaf.createOrGetSvc(fn, deployLabels, deployAnnotations, objName, ns) + svc, err := caaf.createOrGetSvc(ctx, fn, deployLabels, deployAnnotations, objName, ns) if err != nil { caaf.logger.Error("error creating service", zap.Error(err), zap.String("service", objName)) go cleanupFunc(ns, objName) @@ -416,7 +416,7 @@ func (caaf *Container) fnCreate(ctx context.Context, fn *fv1.Function) (*fscache return nil, errors.Wrapf(err, "error creating deployment %v", objName) } - hpa, err := caaf.createOrGetHpa(objName, &fn.Spec.InvokeStrategy.ExecutionStrategy, depl, deployLabels, deployAnnotations) + hpa, err := caaf.createOrGetHpa(ctx, objName, &fn.Spec.InvokeStrategy.ExecutionStrategy, depl, deployLabels, deployAnnotations) if err != nil { caaf.logger.Error("error creating HPA", zap.Error(err), zap.String("hpa", objName)) go cleanupFunc(ns, objName) @@ -488,7 +488,7 @@ func (caaf *Container) updateFunction(ctx context.Context, oldFn *fv1.Function, caaf.logger.Info("function does not use new deployment executor anymore, deleting resources", zap.Any("function", newFn)) // IMP - pass the oldFn, as the new/modified function is not in cache - return caaf.deleteFunction(oldFn) + return caaf.deleteFunction(ctx, oldFn) } // Executor type changed to Container from something else @@ -519,7 +519,7 @@ func (caaf *Container) updateFunction(ctx context.Context, oldFn *fv1.Function, return err } - hpa, err := caaf.getHpa(ns, fsvc.Name) + hpa, err := caaf.getHpa(ctx, ns, fsvc.Name) if err != nil { caaf.updateStatus(oldFn, err, "error getting HPA while updating function") return err @@ -545,7 +545,7 @@ func (caaf *Container) updateFunction(ctx context.Context, oldFn *fv1.Function, } if hpaChanged { - err := caaf.updateHpa(hpa) + err := caaf.updateHpa(ctx, hpa) if err != nil { caaf.updateStatus(oldFn, err, "error updating HPA while updating function") return err @@ -623,7 +623,7 @@ func (caaf *Container) updateFuncDeployment(ctx context.Context, fn *fv1.Functio return err } - err = caaf.updateDeployment(newDeployment, ns) + err = caaf.updateDeployment(ctx, newDeployment, ns) if err != nil { caaf.updateStatus(fn, err, "failed to update deployment while updating function") return err @@ -632,7 +632,7 @@ func (caaf *Container) updateFuncDeployment(ctx context.Context, fn *fv1.Functio return nil } -func (caaf *Container) fnDelete(fn *fv1.Function) error { +func (caaf *Container) fnDelete(ctx context.Context, fn *fv1.Function) error { multierr := &multierror.Error{} // GetByFunction uses resource version as part of cache key, however, @@ -661,7 +661,7 @@ func (caaf *Container) fnDelete(fn *fv1.Function) error { ns = fn.ObjectMeta.Namespace } - err = caaf.cleanupContainer(ns, objName) + err = caaf.cleanupContainer(ctx, ns, objName) multierr = multierror.Append(multierr, err) return multierr.ErrorOrNil() @@ -716,6 +716,7 @@ func (caaf *Container) updateStatus(fn *fv1.Function, err error, message string) // idleObjectReaper reaps objects after certain idle time func (caaf *Container) idleObjectReaper() { + ctx := context.Background() pollSleep := 5 * time.Second for { @@ -734,7 +735,7 @@ func (caaf *Container) idleObjectReaper() { continue } - fn, err := caaf.fissionClient.CoreV1().Functions(fsvc.Function.Namespace).Get(context.TODO(), fsvc.Function.Name, metav1.GetOptions{}) + fn, err := caaf.fissionClient.CoreV1().Functions(fsvc.Function.Namespace).Get(ctx, fsvc.Function.Name, metav1.GetOptions{}) if err != nil { // CaaF manager handles the function delete event and clean cache/kubeobjs itself, // so we ignore the not found error for functions with CaaF executor type here. @@ -762,7 +763,7 @@ func (caaf *Container) idleObjectReaper() { } currentDeploy, err := caaf.kubernetesClient.AppsV1(). - Deployments(deployObj.Namespace).Get(context.TODO(), deployObj.Name, metav1.GetOptions{}) + Deployments(deployObj.Namespace).Get(ctx, deployObj.Name, metav1.GetOptions{}) if err != nil { caaf.logger.Error("error getting function deployment", zap.Error(err), zap.String("function", fsvc.Function.Name)) return @@ -775,7 +776,7 @@ func (caaf *Container) idleObjectReaper() { return } - err = caaf.scaleDeployment(deployObj.Namespace, deployObj.Name, minScale) + err = caaf.scaleDeployment(ctx, deployObj.Namespace, deployObj.Name, minScale) if err != nil { caaf.logger.Error("error scaling down function deployment", zap.Error(err), zap.String("function", fsvc.Function.Name)) } diff --git a/pkg/executor/executortype/container/deployment.go b/pkg/executor/executortype/container/deployment.go index ea6852d6..c53c7782 100644 --- a/pkg/executor/executortype/container/deployment.go +++ b/pkg/executor/executortype/container/deployment.go @@ -75,14 +75,14 @@ func (cn *Container) createOrGetDeployment(ctx context.Context, fn *fv1.Function } if *existingDepl.Spec.Replicas < minScale { - err = cn.scaleDeployment(existingDepl.Namespace, existingDepl.Name, minScale) + err = cn.scaleDeployment(ctx, existingDepl.Namespace, existingDepl.Name, minScale) if err != nil { cn.logger.Error("error scaling up function deployment", zap.Error(err), zap.String("function", fn.ObjectMeta.Name)) return nil, err } } if existingDepl.Status.AvailableReplicas < minScale { - existingDepl, err = cn.waitForDeploy(existingDepl, minScale, specializationTimeout) + existingDepl, err = cn.waitForDeploy(ctx, existingDepl, minScale, specializationTimeout) } return existingDepl, err @@ -102,28 +102,28 @@ func (cn *Container) createOrGetDeployment(ctx context.Context, fn *fv1.Function } } if minScale > 0 { - depl, err = cn.waitForDeploy(depl, minScale, specializationTimeout) + depl, err = cn.waitForDeploy(ctx, depl, minScale, specializationTimeout) } return depl, err } return nil, err } -func (cn *Container) updateDeployment(deployment *appsv1.Deployment, ns string) error { - _, err := cn.kubernetesClient.AppsV1().Deployments(ns).Update(context.TODO(), deployment, metav1.UpdateOptions{}) +func (cn *Container) updateDeployment(ctx context.Context, deployment *appsv1.Deployment, ns string) error { + _, err := cn.kubernetesClient.AppsV1().Deployments(ns).Update(ctx, deployment, metav1.UpdateOptions{}) return err } -func (cn *Container) deleteDeployment(ns string, name string) error { +func (cn *Container) deleteDeployment(ctx context.Context, ns string, name string) error { // DeletePropagationBackground deletes the object immediately and dependent are deleted later // DeletePropagationForeground not advisable; it marks for deleteion and API can still serve those objects deletePropagation := metav1.DeletePropagationBackground - return cn.kubernetesClient.AppsV1().Deployments(ns).Delete(context.TODO(), name, metav1.DeleteOptions{ + return cn.kubernetesClient.AppsV1().Deployments(ns).Delete(ctx, name, metav1.DeleteOptions{ PropagationPolicy: &deletePropagation, }) } -func (cn *Container) waitForDeploy(depl *appsv1.Deployment, replicas int32, specializationTimeout int) (latestDepl *appsv1.Deployment, err error) { +func (cn *Container) waitForDeploy(ctx context.Context, depl *appsv1.Deployment, replicas int32, specializationTimeout int) (latestDepl *appsv1.Deployment, err error) { oldStatus := depl.Status // if no specializationTimeout is set, use default value @@ -132,7 +132,7 @@ func (cn *Container) waitForDeploy(depl *appsv1.Deployment, replicas int32, spec } for i := 0; i < specializationTimeout; i++ { - latestDepl, err := cn.kubernetesClient.AppsV1().Deployments(depl.ObjectMeta.Namespace).Get(context.TODO(), depl.Name, metav1.GetOptions{}) + latestDepl, err := cn.kubernetesClient.AppsV1().Deployments(depl.ObjectMeta.Namespace).Get(ctx, depl.Name, metav1.GetOptions{}) if err != nil { return nil, err } @@ -196,7 +196,7 @@ func (cn *Container) getDeploymentSpec(ctx context.Context, fn *fv1.Function, ta resources := cn.getResources(fn) // Other executor types rely on Environments to add configmaps and secrets - envFromSources, err := util.ConvertConfigSecrets(fn, cn.kubernetesClient) + envFromSources, err := util.ConvertConfigSecrets(ctx, fn, cn.kubernetesClient) if err != nil { return nil, err } @@ -277,12 +277,12 @@ func (cn *Container) getDeploymentSpec(ctx context.Context, fn *fv1.Function, ta return deployment, nil } -func (caaf *Container) scaleDeployment(deplNS string, deplName string, replicas int32) error { +func (caaf *Container) scaleDeployment(ctx context.Context, deplNS string, deplName string, replicas int32) error { caaf.logger.Info("scaling deployment", zap.String("deployment", deplName), zap.String("namespace", deplNS), zap.Int32("replicas", replicas)) - _, err := caaf.kubernetesClient.AppsV1().Deployments(deplNS).UpdateScale(context.TODO(), deplName, &autoscalingv1.Scale{ + _, err := caaf.kubernetesClient.AppsV1().Deployments(deplNS).UpdateScale(ctx, deplName, &autoscalingv1.Scale{ ObjectMeta: metav1.ObjectMeta{ Name: deplName, Namespace: deplNS, diff --git a/pkg/executor/executortype/container/funchandlers.go b/pkg/executor/executortype/container/funchandlers.go index d3ac420b..bae022e9 100644 --- a/pkg/executor/executortype/container/funchandlers.go +++ b/pkg/executor/executortype/container/funchandlers.go @@ -55,7 +55,7 @@ func (caaf *Container) FuncInformerHandler(ctx context.Context) k8sCache.Resourc go func() { log := caaf.logger.With(zap.String("function_name", fn.ObjectMeta.Name), zap.String("function_namespace", fn.ObjectMeta.Namespace)) log.Debug("start function delete handler") - err := caaf.deleteFunction(fn) + err := caaf.deleteFunction(ctx, fn) if err != nil { log.Error("error deleting function", zap.Error(err)) } diff --git a/pkg/executor/executortype/container/hpa.go b/pkg/executor/executortype/container/hpa.go index 005a6944..ed486cad 100644 --- a/pkg/executor/executortype/container/hpa.go +++ b/pkg/executor/executortype/container/hpa.go @@ -34,7 +34,7 @@ const ( DeploymentVersion = "apps/v1" ) -func (cn *Container) createOrGetHpa(hpaName string, execStrategy *fv1.ExecutionStrategy, +func (cn *Container) createOrGetHpa(ctx context.Context, hpaName string, execStrategy *fv1.ExecutionStrategy, depl *appsv1.Deployment, deployLabels map[string]string, deployAnnotations map[string]string) (*asv1.HorizontalPodAutoscaler, error) { if depl == nil { @@ -69,14 +69,14 @@ func (cn *Container) createOrGetHpa(hpaName string, execStrategy *fv1.ExecutionS }, } - existingHpa, err := cn.getHpa(depl.ObjectMeta.Namespace, hpaName) + existingHpa, err := cn.getHpa(ctx, depl.ObjectMeta.Namespace, hpaName) if err == nil { // to adopt orphan service if existingHpa.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] != cn.instanceID { existingHpa.Annotations = hpa.Annotations existingHpa.Labels = hpa.Labels existingHpa.Spec = hpa.Spec - existingHpa, err = cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Update(context.TODO(), existingHpa, metav1.UpdateOptions{}) + existingHpa, err = cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Update(ctx, existingHpa, metav1.UpdateOptions{}) if err != nil { cn.logger.Warn("error adopting HPA", zap.Error(err), zap.String("HPA", hpaName), zap.String("ns", depl.ObjectMeta.Namespace)) @@ -85,10 +85,10 @@ func (cn *Container) createOrGetHpa(hpaName string, execStrategy *fv1.ExecutionS } return existingHpa, err } else if k8s_err.IsNotFound(err) { - cHpa, err := cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Create(context.TODO(), hpa, metav1.CreateOptions{}) + cHpa, err := cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Create(ctx, hpa, metav1.CreateOptions{}) if err != nil { if k8s_err.IsAlreadyExists(err) { - cHpa, err = cn.getHpa(depl.ObjectMeta.Namespace, hpaName) + cHpa, err = cn.getHpa(ctx, depl.ObjectMeta.Namespace, hpaName) } if err != nil { return nil, err @@ -99,15 +99,15 @@ func (cn *Container) createOrGetHpa(hpaName string, execStrategy *fv1.ExecutionS return nil, err } -func (cn *Container) getHpa(ns, name string) (*asv1.HorizontalPodAutoscaler, error) { - return cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Get(context.TODO(), name, metav1.GetOptions{}) +func (cn *Container) getHpa(ctx context.Context, ns, name string) (*asv1.HorizontalPodAutoscaler, error) { + return cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Get(ctx, name, metav1.GetOptions{}) } -func (cn *Container) updateHpa(hpa *asv1.HorizontalPodAutoscaler) error { - _, err := cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Update(context.TODO(), hpa, metav1.UpdateOptions{}) +func (cn *Container) updateHpa(ctx context.Context, hpa *asv1.HorizontalPodAutoscaler) error { + _, err := cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Update(ctx, hpa, metav1.UpdateOptions{}) return err } -func (cn *Container) deleteHpa(ns string, name string) error { - return cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Delete(context.TODO(), name, metav1.DeleteOptions{}) +func (cn *Container) deleteHpa(ctx context.Context, ns string, name string) error { + return cn.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Delete(ctx, name, metav1.DeleteOptions{}) } diff --git a/pkg/executor/executortype/container/svc.go b/pkg/executor/executortype/container/svc.go index 0c5e5648..ed576239 100644 --- a/pkg/executor/executortype/container/svc.go +++ b/pkg/executor/executortype/container/svc.go @@ -42,7 +42,7 @@ func (cn *Container) getSvPort(fn *fv1.Function) (port int32, err error) { return fn.Spec.PodSpec.Containers[0].Ports[0].ContainerPort, nil } -func (cn *Container) createOrGetSvc(fn *fv1.Function, deployLabels map[string]string, deployAnnotations map[string]string, svcName string, svcNamespace string) (*apiv1.Service, error) { +func (cn *Container) createOrGetSvc(ctx context.Context, fn *fv1.Function, deployLabels map[string]string, deployAnnotations map[string]string, svcName string, svcNamespace string) (*apiv1.Service, error) { targetPort, err := cn.getSvPort(fn) if err != nil { return nil, err @@ -66,7 +66,7 @@ func (cn *Container) createOrGetSvc(fn *fv1.Function, deployLabels map[string]st }, } - existingSvc, err := cn.kubernetesClient.CoreV1().Services(svcNamespace).Get(context.TODO(), svcName, metav1.GetOptions{}) + existingSvc, err := cn.kubernetesClient.CoreV1().Services(svcNamespace).Get(ctx, svcName, metav1.GetOptions{}) if err == nil { // to adopt orphan service if existingSvc.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] != cn.instanceID { @@ -75,7 +75,7 @@ func (cn *Container) createOrGetSvc(fn *fv1.Function, deployLabels map[string]st existingSvc.Spec.Ports = service.Spec.Ports existingSvc.Spec.Selector = service.Spec.Selector existingSvc.Spec.Type = service.Spec.Type - existingSvc, err = cn.kubernetesClient.CoreV1().Services(svcNamespace).Update(context.TODO(), existingSvc, metav1.UpdateOptions{}) + existingSvc, err = cn.kubernetesClient.CoreV1().Services(svcNamespace).Update(ctx, existingSvc, metav1.UpdateOptions{}) if err != nil { cn.logger.Warn("error adopting service", zap.Error(err), zap.String("service", svcName), zap.String("ns", svcNamespace)) @@ -84,10 +84,10 @@ func (cn *Container) createOrGetSvc(fn *fv1.Function, deployLabels map[string]st } return existingSvc, err } else if k8s_err.IsNotFound(err) { - svc, err := cn.kubernetesClient.CoreV1().Services(svcNamespace).Create(context.TODO(), service, metav1.CreateOptions{}) + svc, err := cn.kubernetesClient.CoreV1().Services(svcNamespace).Create(ctx, service, metav1.CreateOptions{}) if err != nil { if k8s_err.IsAlreadyExists(err) { - svc, err = cn.kubernetesClient.CoreV1().Services(svcNamespace).Get(context.TODO(), svcName, metav1.GetOptions{}) + svc, err = cn.kubernetesClient.CoreV1().Services(svcNamespace).Get(ctx, svcName, metav1.GetOptions{}) } if err != nil { return nil, err @@ -98,6 +98,6 @@ func (cn *Container) createOrGetSvc(fn *fv1.Function, deployLabels map[string]st return nil, err } -func (cn *Container) deleteSvc(ns string, name string) error { - return cn.kubernetesClient.CoreV1().Services(ns).Delete(context.TODO(), name, metav1.DeleteOptions{}) +func (cn *Container) deleteSvc(ctx context.Context, ns string, name string) error { + return cn.kubernetesClient.CoreV1().Services(ns).Delete(ctx, name, metav1.DeleteOptions{}) } diff --git a/pkg/executor/executortype/newdeploy/envhandlers.go b/pkg/executor/executortype/newdeploy/envhandlers.go index e091b234..7040783d 100644 --- a/pkg/executor/executortype/newdeploy/envhandlers.go +++ b/pkg/executor/executortype/newdeploy/envhandlers.go @@ -32,17 +32,18 @@ func (deploy *NewDeploy) EnvEventHandlers() k8sCache.ResourceEventHandlerFuncs { UpdateFunc: func(oldObj interface{}, newObj interface{}) { newEnv := newObj.(*fv1.Environment) oldEnv := oldObj.(*fv1.Environment) + ctx := context.Background() // Currently only an image update in environment calls for function's deployment recreation. In future there might be more attributes which would want to do it if oldEnv.Spec.Runtime.Image != newEnv.Spec.Runtime.Image { deploy.logger.Debug("Updating all function of the environment that changed, old env:", zap.Any("environment", oldEnv)) - funcs := deploy.getEnvFunctions(&newEnv.ObjectMeta) + funcs := deploy.getEnvFunctions(ctx, &newEnv.ObjectMeta) for _, f := range funcs { - function, err := deploy.fissionClient.CoreV1().Functions(f.ObjectMeta.Namespace).Get(context.TODO(), f.ObjectMeta.Name, metav1.GetOptions{}) + function, err := deploy.fissionClient.CoreV1().Functions(f.ObjectMeta.Namespace).Get(ctx, f.ObjectMeta.Name, metav1.GetOptions{}) if err != nil { deploy.logger.Error("Error getting function", zap.Error(err), zap.Any("function", function)) continue } - err = deploy.updateFuncDeployment(function, newEnv) + err = deploy.updateFuncDeployment(ctx, function, newEnv) if err != nil { deploy.logger.Error("Error updating function", zap.Error(err), zap.Any("function", function)) continue diff --git a/pkg/executor/executortype/newdeploy/funchandlers.go b/pkg/executor/executortype/newdeploy/funchandlers.go index 0c551330..b764bb7f 100644 --- a/pkg/executor/executortype/newdeploy/funchandlers.go +++ b/pkg/executor/executortype/newdeploy/funchandlers.go @@ -16,6 +16,8 @@ limitations under the License. package newdeploy import ( + "context" + "go.uber.org/zap" k8sCache "k8s.io/client-go/tools/cache" @@ -29,9 +31,10 @@ func (deploy *NewDeploy) FunctionEventHandlers() k8sCache.ResourceEventHandlerFu // and worker pattern to process items instead of moving process to another goroutine. // example: https://github.com/kubernetes/kubernetes/blob/master/pkg/controller/job/job_controller.go go func() { + ctx := context.Background() fn := obj.(*fv1.Function) deploy.logger.Debug("create deployment for function", zap.Any("fn", fn.ObjectMeta), zap.Any("fnspec", fn.Spec)) - _, err := deploy.createFunction(fn) + _, err := deploy.createFunction(ctx, fn) if err != nil { deploy.logger.Error("error eager creating function", zap.Error(err), @@ -43,7 +46,8 @@ func (deploy *NewDeploy) FunctionEventHandlers() k8sCache.ResourceEventHandlerFu DeleteFunc: func(obj interface{}) { fn := obj.(*fv1.Function) go func() { - err := deploy.deleteFunction(fn) + ctx := context.Background() + err := deploy.deleteFunction(ctx, fn) if err != nil { deploy.logger.Error("error deleting function", zap.Error(err), @@ -55,7 +59,8 @@ func (deploy *NewDeploy) FunctionEventHandlers() k8sCache.ResourceEventHandlerFu oldFn := oldObj.(*fv1.Function) newFn := newObj.(*fv1.Function) go func() { - err := deploy.updateFunction(oldFn, newFn) + ctx := context.Background() + err := deploy.updateFunction(ctx, oldFn, newFn) if err != nil { deploy.logger.Error("error updating function", zap.Error(err), diff --git a/pkg/executor/executortype/newdeploy/newdeploy.go b/pkg/executor/executortype/newdeploy/newdeploy.go index e4574987..8a6474e0 100644 --- a/pkg/executor/executortype/newdeploy/newdeploy.go +++ b/pkg/executor/executortype/newdeploy/newdeploy.go @@ -45,7 +45,7 @@ const ( DeploymentVersion = "apps/v1" ) -func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Environment, +func (deploy *NewDeploy) createOrGetDeployment(ctx context.Context, fn *fv1.Function, env *fv1.Environment, deployName string, deployLabels map[string]string, deployAnnotations map[string]string, deployNamespace string) (*appsv1.Deployment, error) { specializationTimeout := fn.Spec.InvokeStrategy.ExecutionStrategy.SpecializationTimeout @@ -58,12 +58,12 @@ func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Enviro minScale = 1 } - deployment, err := deploy.getDeploymentSpec(fn, env, &minScale, deployName, deployNamespace, deployLabels, deployAnnotations) + deployment, err := deploy.getDeploymentSpec(ctx, fn, env, &minScale, deployName, deployNamespace, deployLabels, deployAnnotations) if err != nil { return nil, err } - existingDepl, err := deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Get(context.TODO(), deployName, metav1.GetOptions{}) + existingDepl, err := deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Get(ctx, deployName, metav1.GetOptions{}) if err == nil { // Try to adopt orphan deployment created by the old executor. if existingDepl.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] != deploy.instanceID { @@ -75,7 +75,7 @@ func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Enviro // Update with the latest deployment spec. Kubernetes will trigger // rolling update if spec is different from the one in the cluster. - existingDepl, err = deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Update(context.TODO(), existingDepl, metav1.UpdateOptions{}) + existingDepl, err = deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Update(ctx, existingDepl, metav1.UpdateOptions{}) if err != nil { deploy.logger.Warn("error adopting deploy", zap.Error(err), zap.String("deploy", deployName), zap.String("ns", deployNamespace)) @@ -86,27 +86,27 @@ func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Enviro } if *existingDepl.Spec.Replicas < minScale { - err = deploy.scaleDeployment(existingDepl.Namespace, existingDepl.Name, minScale) + err = deploy.scaleDeployment(ctx, existingDepl.Namespace, existingDepl.Name, minScale) if err != nil { deploy.logger.Error("error scaling up function deployment", zap.Error(err), zap.String("function", fn.ObjectMeta.Name)) return nil, err } } if existingDepl.Status.AvailableReplicas < minScale { - existingDepl, err = deploy.waitForDeploy(existingDepl, minScale, specializationTimeout) + existingDepl, err = deploy.waitForDeploy(ctx, existingDepl, minScale, specializationTimeout) } return existingDepl, err } else if k8s_err.IsNotFound(err) { - err := deploy.setupRBACObjs(deployNamespace, fn) + err := deploy.setupRBACObjs(ctx, deployNamespace, fn) if err != nil { return nil, err } - depl, err := deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Create(context.TODO(), deployment, metav1.CreateOptions{}) + depl, err := deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Create(ctx, deployment, metav1.CreateOptions{}) if err != nil { if k8s_err.IsAlreadyExists(err) { - depl, err = deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Get(context.TODO(), deployName, metav1.GetOptions{}) + depl, err = deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Get(ctx, deployName, metav1.GetOptions{}) } if err != nil { deploy.logger.Error("error while creating function deployment", @@ -118,14 +118,14 @@ func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Enviro } } if minScale > 0 { - depl, err = deploy.waitForDeploy(depl, minScale, specializationTimeout) + depl, err = deploy.waitForDeploy(ctx, depl, minScale, specializationTimeout) } return depl, err } return nil, err } -func (deploy *NewDeploy) setupRBACObjs(deployNamespace string, fn *fv1.Function) error { +func (deploy *NewDeploy) setupRBACObjs(ctx context.Context, deployNamespace string, fn *fv1.Function) error { // create fetcher SA in this ns, if not already created err := deploy.fetcherConfig.SetupServiceAccount(deploy.kubernetesClient, deployNamespace, fn.ObjectMeta) if err != nil { @@ -139,7 +139,7 @@ func (deploy *NewDeploy) setupRBACObjs(deployNamespace string, fn *fv1.Function) } // create a cluster role binding for the fetcher SA, if not already created, granting access to do a get on packages in any ns - err = utils.SetupRoleBinding(deploy.logger, deploy.kubernetesClient, fv1.PackageGetterRB, fn.Spec.Package.PackageRef.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, deployNamespace) + err = utils.SetupRoleBinding(ctx, deploy.logger, deploy.kubernetesClient, fv1.PackageGetterRB, fn.Spec.Package.PackageRef.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, deployNamespace) if err != nil { deploy.logger.Error("error creating role binding for function", zap.Error(err), @@ -150,7 +150,7 @@ func (deploy *NewDeploy) setupRBACObjs(deployNamespace string, fn *fv1.Function) } // create rolebinding in function namespace for fetcherSA.envNamespace to be able to get secrets and configmaps - err = utils.SetupRoleBinding(deploy.logger, deploy.kubernetesClient, fv1.SecretConfigMapGetterRB, fn.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, deployNamespace) + err = utils.SetupRoleBinding(ctx, deploy.logger, deploy.kubernetesClient, fv1.SecretConfigMapGetterRB, fn.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, deployNamespace) if err != nil { deploy.logger.Error("error creating role binding for function", zap.Error(err), @@ -166,21 +166,21 @@ func (deploy *NewDeploy) setupRBACObjs(deployNamespace string, fn *fv1.Function) return nil } -func (deploy *NewDeploy) updateDeployment(deployment *appsv1.Deployment, ns string) error { - _, err := deploy.kubernetesClient.AppsV1().Deployments(ns).Update(context.TODO(), deployment, metav1.UpdateOptions{}) +func (deploy *NewDeploy) updateDeployment(ctx context.Context, deployment *appsv1.Deployment, ns string) error { + _, err := deploy.kubernetesClient.AppsV1().Deployments(ns).Update(ctx, deployment, metav1.UpdateOptions{}) return err } -func (deploy *NewDeploy) deleteDeployment(ns string, name string) error { +func (deploy *NewDeploy) deleteDeployment(ctx context.Context, ns string, name string) error { // DeletePropagationBackground deletes the object immediately and dependent are deleted later // DeletePropagationForeground not advisable; it marks for deleteion and API can still serve those objects deletePropagation := metav1.DeletePropagationBackground - return deploy.kubernetesClient.AppsV1().Deployments(ns).Delete(context.TODO(), name, metav1.DeleteOptions{ + return deploy.kubernetesClient.AppsV1().Deployments(ns).Delete(ctx, name, metav1.DeleteOptions{ PropagationPolicy: &deletePropagation, }) } -func (deploy *NewDeploy) getDeploymentSpec(fn *fv1.Function, env *fv1.Environment, targetReplicas *int32, +func (deploy *NewDeploy) getDeploymentSpec(ctx context.Context, fn *fv1.Function, env *fv1.Environment, targetReplicas *int32, deployName string, deployNamespace string, deployLabels map[string]string, deployAnnotations map[string]string) (*appsv1.Deployment, error) { replicas := int32(fn.Spec.InvokeStrategy.ExecutionStrategy.MinScale) @@ -233,7 +233,7 @@ func (deploy *NewDeploy) getDeploymentSpec(fn *fv1.Function, env *fv1.Environmen // rollback, set RevisionHistoryLimit to 0 to disable this feature. revisionHistoryLimit := int32(0) - rvCount, err := referencedResourcesRVSum(deploy.kubernetesClient, fn.ObjectMeta.Namespace, fn.Spec.Secrets, fn.Spec.ConfigMaps) + rvCount, err := referencedResourcesRVSum(ctx, deploy.kubernetesClient, fn.ObjectMeta.Namespace, fn.Spec.Secrets, fn.Spec.ConfigMaps) if err != nil { return nil, err } @@ -366,7 +366,7 @@ func (deploy *NewDeploy) getResources(env *fv1.Environment, fn *fv1.Function) ap return resources } -func (deploy *NewDeploy) createOrGetHpa(hpaName string, execStrategy *fv1.ExecutionStrategy, +func (deploy *NewDeploy) createOrGetHpa(ctx context.Context, hpaName string, execStrategy *fv1.ExecutionStrategy, depl *appsv1.Deployment, deployLabels map[string]string, deployAnnotations map[string]string) (*asv1.HorizontalPodAutoscaler, error) { if depl == nil { @@ -401,14 +401,14 @@ func (deploy *NewDeploy) createOrGetHpa(hpaName string, execStrategy *fv1.Execut }, } - existingHpa, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Get(context.TODO(), hpaName, metav1.GetOptions{}) + existingHpa, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Get(ctx, hpaName, metav1.GetOptions{}) if err == nil { // to adopt orphan service if existingHpa.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] != deploy.instanceID { existingHpa.Annotations = hpa.Annotations existingHpa.Labels = hpa.Labels existingHpa.Spec = hpa.Spec - existingHpa, err = deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Update(context.TODO(), existingHpa, metav1.UpdateOptions{}) + existingHpa, err = deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Update(ctx, existingHpa, metav1.UpdateOptions{}) if err != nil { deploy.logger.Warn("error adopting HPA", zap.Error(err), zap.String("HPA", hpaName), zap.String("ns", depl.ObjectMeta.Namespace)) @@ -417,10 +417,10 @@ func (deploy *NewDeploy) createOrGetHpa(hpaName string, execStrategy *fv1.Execut } return existingHpa, err } else if k8s_err.IsNotFound(err) { - cHpa, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Create(context.TODO(), hpa, metav1.CreateOptions{}) + cHpa, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Create(ctx, hpa, metav1.CreateOptions{}) if err != nil { if k8s_err.IsAlreadyExists(err) { - cHpa, err = deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Get(context.TODO(), hpaName, metav1.GetOptions{}) + cHpa, err = deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Get(ctx, hpaName, metav1.GetOptions{}) } if err != nil { return nil, err @@ -431,20 +431,20 @@ func (deploy *NewDeploy) createOrGetHpa(hpaName string, execStrategy *fv1.Execut return nil, err } -func (deploy *NewDeploy) getHpa(ns, name string) (*asv1.HorizontalPodAutoscaler, error) { - return deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Get(context.TODO(), name, metav1.GetOptions{}) +func (deploy *NewDeploy) getHpa(ctx context.Context, ns, name string) (*asv1.HorizontalPodAutoscaler, error) { + return deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Get(ctx, name, metav1.GetOptions{}) } -func (deploy *NewDeploy) updateHpa(hpa *asv1.HorizontalPodAutoscaler) error { - _, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Update(context.TODO(), hpa, metav1.UpdateOptions{}) +func (deploy *NewDeploy) updateHpa(ctx context.Context, hpa *asv1.HorizontalPodAutoscaler) error { + _, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Update(ctx, hpa, metav1.UpdateOptions{}) return err } -func (deploy *NewDeploy) deleteHpa(ns string, name string) error { - return deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Delete(context.TODO(), name, metav1.DeleteOptions{}) +func (deploy *NewDeploy) deleteHpa(ctx context.Context, ns string, name string) error { + return deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Delete(ctx, name, metav1.DeleteOptions{}) } -func (deploy *NewDeploy) createOrGetSvc(deployLabels map[string]string, deployAnnotations map[string]string, svcName string, svcNamespace string) (*apiv1.Service, error) { +func (deploy *NewDeploy) createOrGetSvc(ctx context.Context, deployLabels map[string]string, deployAnnotations map[string]string, svcName string, svcNamespace string) (*apiv1.Service, error) { service := &apiv1.Service{ ObjectMeta: metav1.ObjectMeta{ Name: svcName, @@ -465,7 +465,7 @@ func (deploy *NewDeploy) createOrGetSvc(deployLabels map[string]string, deployAn }, } - existingSvc, err := deploy.kubernetesClient.CoreV1().Services(svcNamespace).Get(context.TODO(), svcName, metav1.GetOptions{}) + existingSvc, err := deploy.kubernetesClient.CoreV1().Services(svcNamespace).Get(ctx, svcName, metav1.GetOptions{}) if err == nil { // to adopt orphan service if existingSvc.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] != deploy.instanceID { @@ -474,7 +474,7 @@ func (deploy *NewDeploy) createOrGetSvc(deployLabels map[string]string, deployAn existingSvc.Spec.Ports = service.Spec.Ports existingSvc.Spec.Selector = service.Spec.Selector existingSvc.Spec.Type = service.Spec.Type - existingSvc, err = deploy.kubernetesClient.CoreV1().Services(svcNamespace).Update(context.TODO(), existingSvc, metav1.UpdateOptions{}) + existingSvc, err = deploy.kubernetesClient.CoreV1().Services(svcNamespace).Update(ctx, existingSvc, metav1.UpdateOptions{}) if err != nil { deploy.logger.Warn("error adopting service", zap.Error(err), zap.String("service", svcName), zap.String("ns", svcNamespace)) @@ -483,10 +483,10 @@ func (deploy *NewDeploy) createOrGetSvc(deployLabels map[string]string, deployAn } return existingSvc, err } else if k8s_err.IsNotFound(err) { - svc, err := deploy.kubernetesClient.CoreV1().Services(svcNamespace).Create(context.TODO(), service, metav1.CreateOptions{}) + svc, err := deploy.kubernetesClient.CoreV1().Services(svcNamespace).Create(ctx, service, metav1.CreateOptions{}) if err != nil { if k8s_err.IsAlreadyExists(err) { - svc, err = deploy.kubernetesClient.CoreV1().Services(svcNamespace).Get(context.TODO(), svcName, metav1.GetOptions{}) + svc, err = deploy.kubernetesClient.CoreV1().Services(svcNamespace).Get(ctx, svcName, metav1.GetOptions{}) } if err != nil { return nil, err @@ -497,11 +497,11 @@ func (deploy *NewDeploy) createOrGetSvc(deployLabels map[string]string, deployAn return nil, err } -func (deploy *NewDeploy) deleteSvc(ns string, name string) error { - return deploy.kubernetesClient.CoreV1().Services(ns).Delete(context.TODO(), name, metav1.DeleteOptions{}) +func (deploy *NewDeploy) deleteSvc(ctx context.Context, ns string, name string) error { + return deploy.kubernetesClient.CoreV1().Services(ns).Delete(ctx, name, metav1.DeleteOptions{}) } -func (deploy *NewDeploy) waitForDeploy(depl *appsv1.Deployment, replicas int32, specializationTimeout int) (latestDepl *appsv1.Deployment, err error) { +func (deploy *NewDeploy) waitForDeploy(ctx context.Context, depl *appsv1.Deployment, replicas int32, specializationTimeout int) (latestDepl *appsv1.Deployment, err error) { oldStatus := depl.Status // if no specializationTimeout is set, use default value @@ -510,7 +510,7 @@ func (deploy *NewDeploy) waitForDeploy(depl *appsv1.Deployment, replicas int32, } for i := 0; i < specializationTimeout; i++ { - latestDepl, err = deploy.kubernetesClient.AppsV1().Deployments(depl.ObjectMeta.Namespace).Get(context.TODO(), depl.Name, metav1.GetOptions{}) + latestDepl, err = deploy.kubernetesClient.AppsV1().Deployments(depl.ObjectMeta.Namespace).Get(ctx, depl.Name, metav1.GetOptions{}) if err != nil { return nil, err } @@ -533,10 +533,10 @@ func (deploy *NewDeploy) waitForDeploy(depl *appsv1.Deployment, replicas int32, } // cleanupNewdeploy cleans all kubernetes objects related to function -func (deploy *NewDeploy) cleanupNewdeploy(ns string, name string) error { +func (deploy *NewDeploy) cleanupNewdeploy(ctx context.Context, ns string, name string) error { result := &multierror.Error{} - err := deploy.deleteSvc(ns, name) + err := deploy.deleteSvc(ctx, ns, name) if err != nil && !k8s_err.IsNotFound(err) { deploy.logger.Error("error deleting service for newdeploy function", zap.Error(err), @@ -545,7 +545,7 @@ func (deploy *NewDeploy) cleanupNewdeploy(ns string, name string) error { result = multierror.Append(result, err) } - err = deploy.deleteHpa(ns, name) + err = deploy.deleteHpa(ctx, ns, name) if err != nil && !k8s_err.IsNotFound(err) { deploy.logger.Error("error deleting HPA for newdeploy function", zap.Error(err), @@ -554,7 +554,7 @@ func (deploy *NewDeploy) cleanupNewdeploy(ns string, name string) error { result = multierror.Append(result, err) } - err = deploy.deleteDeployment(ns, name) + err = deploy.deleteDeployment(ctx, ns, name) if err != nil && !k8s_err.IsNotFound(err) { deploy.logger.Error("error deleting deployment for newdeploy function", zap.Error(err), @@ -574,11 +574,11 @@ func (deploy *NewDeploy) cleanupNewdeploy(ns string, name string) error { // identical way to get a value that can reflect resources changed without affecting by the time. // To achieve this goal, the sum of the resource version of all referenced resources is a good fit for our // scenario since the sum of the resource version is always the same as long as no resources changed. -func referencedResourcesRVSum(client *kubernetes.Clientset, namespace string, secrets []fv1.SecretReference, cfgmaps []fv1.ConfigMapReference) (int, error) { +func referencedResourcesRVSum(ctx context.Context, client *kubernetes.Clientset, namespace string, secrets []fv1.SecretReference, cfgmaps []fv1.ConfigMapReference) (int, error) { rvCount := 0 if len(secrets) > 0 { - list, err := client.CoreV1().Secrets(namespace).List(context.TODO(), metav1.ListOptions{}) + list, err := client.CoreV1().Secrets(namespace).List(ctx, metav1.ListOptions{}) if err != nil { return 0, err } @@ -598,7 +598,7 @@ func referencedResourcesRVSum(client *kubernetes.Clientset, namespace string, se } if len(cfgmaps) > 0 { - list, err := client.CoreV1().ConfigMaps(namespace).List(context.TODO(), metav1.ListOptions{}) + list, err := client.CoreV1().ConfigMaps(namespace).List(ctx, metav1.ListOptions{}) if err != nil { return 0, err } diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr.go b/pkg/executor/executortype/newdeploy/newdeploymgr.go index 132cb124..7fab0888 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr.go @@ -144,7 +144,7 @@ func (deploy *NewDeploy) GetFuncSvc(ctx context.Context, fn *fv1.Function) (*fsc // TODO: client-go doesn't support to pass in context. // Once it supports context, we should change the signature of method. // https://github.com/kubernetes/kubernetes/issues/46503 - return deploy.createFunction(fn) + return deploy.createFunction(ctx, fn) } // GetFuncSvcFromCache returns a function service from cache; error otherwise. @@ -194,7 +194,7 @@ func (deploy *NewDeploy) getServiceInfo(ctx context.Context, obj apiv1.ObjectRef return service, nil } -func (deploy *NewDeploy) getDeploymentInfo(obj apiv1.ObjectReference) (*appsv1.Deployment, error) { +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 { @@ -203,7 +203,7 @@ func (deploy *NewDeploy) getDeploymentInfo(obj apiv1.ObjectReference) (*appsv1.D zap.Bool("exists", exists), zap.Error(err), ) - deployment, err := deploy.kubernetesClient.AppsV1().Deployments(obj.Namespace).Get(context.TODO(), obj.Name, metav1.GetOptions{}) + deployment, err := deploy.kubernetesClient.AppsV1().Deployments(obj.Namespace).Get(ctx, obj.Name, metav1.GetOptions{}) return deployment, err } @@ -234,7 +234,7 @@ func (deploy *NewDeploy) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) boo } } else if strings.ToLower(obj.Kind) == "deployment" { - currentDeploy, err := deploy.getDeploymentInfo(obj) + currentDeploy, err := deploy.getDeploymentInfo(ctx, obj) if err != nil { if !k8sErrs.IsNotFound(err) { deploy.logger.Error("error validating function deployment", zap.String("function", fsvc.Function.Name), zap.Error(err)) @@ -273,7 +273,7 @@ func (deploy *NewDeploy) RefreshFuncPods(ctx context.Context, logger *zap.Logger // Ideally there should be only one deployment but for now we rely on label/selector to ensure that condition for _, deployment := range dep.Items { - rvCount, err := referencedResourcesRVSum(deploy.kubernetesClient, deployment.Namespace, f.Spec.Secrets, f.Spec.ConfigMaps) + rvCount, err := referencedResourcesRVSum(ctx, deploy.kubernetesClient, deployment.Namespace, f.Spec.Secrets, f.Spec.ConfigMaps) if err != nil { return err } @@ -308,7 +308,7 @@ func (deploy *NewDeploy) AdoptExistingResources(ctx context.Context) { go func() { defer wg.Done() - _, err = deploy.fnCreate(fn) + _, err = deploy.fnCreate(ctx, fn) if err != nil { deploy.logger.Warn("failed to adopt resources for function", zap.Error(err)) return @@ -330,17 +330,17 @@ func (deploy *NewDeploy) CleanupOldExecutorObjects(ctx context.Context) { LabelSelector: labels.Set(map[string]string{fv1.EXECUTOR_TYPE: string(fv1.ExecutorTypeNewdeploy)}).AsSelector().String(), } - err := reaper.CleanupHpa(deploy.logger, deploy.kubernetesClient, deploy.instanceID, listOpts) + err := reaper.CleanupHpa(ctx, 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) + err = reaper.CleanupDeployments(ctx, 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) + err = reaper.CleanupServices(ctx, deploy.logger, deploy.kubernetesClient, deploy.instanceID, listOpts) if err != nil { errs = multierror.Append(errs, err) } @@ -351,8 +351,8 @@ func (deploy *NewDeploy) CleanupOldExecutorObjects(ctx context.Context) { } } -func (deploy *NewDeploy) getEnvFunctions(m *metav1.ObjectMeta) []fv1.Function { - funcList, err := deploy.fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{}) +func (deploy *NewDeploy) getEnvFunctions(ctx context.Context, m *metav1.ObjectMeta) []fv1.Function { + funcList, err := deploy.fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(ctx, metav1.ListOptions{}) if err != nil { deploy.logger.Error("Error getting functions for env", zap.Error(err), zap.Any("environment", m)) } @@ -365,14 +365,14 @@ func (deploy *NewDeploy) getEnvFunctions(m *metav1.ObjectMeta) []fv1.Function { return relatedFunctions } -func (deploy *NewDeploy) createFunction(fn *fv1.Function) (*fscache.FuncSvc, error) { +func (deploy *NewDeploy) createFunction(ctx context.Context, fn *fv1.Function) (*fscache.FuncSvc, error) { if fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType != fv1.ExecutorTypeNewdeploy { return nil, nil } fsvcObj, err := deploy.throttler.RunOnce(string(fn.ObjectMeta.UID), func(ableToCreate bool) (interface{}, error) { if ableToCreate { - return deploy.fnCreate(fn) + return deploy.fnCreate(ctx, fn) } return deploy.fsCache.GetByFunctionUID(fn.ObjectMeta.UID) }) @@ -394,20 +394,21 @@ func (deploy *NewDeploy) createFunction(fn *fv1.Function) (*fscache.FuncSvc, err return fsvc, err } -func (deploy *NewDeploy) deleteFunction(fn *fv1.Function) error { +func (deploy *NewDeploy) deleteFunction(ctx context.Context, fn *fv1.Function) error { if fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType != fv1.ExecutorTypeNewdeploy { return nil } - err := deploy.fnDelete(fn) + err := deploy.fnDelete(ctx, fn) if err != nil { err = errors.Wrapf(err, "error deleting kubernetes objects of function %v", fn.ObjectMeta) } return err } -func (deploy *NewDeploy) fnCreate(fn *fv1.Function) (*fscache.FuncSvc, error) { +func (deploy *NewDeploy) fnCreate(ctx context.Context, fn *fv1.Function) (*fscache.FuncSvc, error) { cleanupFunc := func(ns string, name string) { - err := deploy.cleanupNewdeploy(ns, name) + ctx := context.Background() + err := deploy.cleanupNewdeploy(ctx, ns, name) if err != nil { deploy.logger.Error("received error while cleaning function resources", zap.String("namespace", ns), zap.String("name", name)) @@ -415,7 +416,7 @@ func (deploy *NewDeploy) fnCreate(fn *fv1.Function) (*fscache.FuncSvc, error) { } env, err := deploy.fissionClient.CoreV1(). Environments(fn.Spec.Environment.Namespace). - Get(context.TODO(), fn.Spec.Environment.Name, metav1.GetOptions{}) + Get(ctx, fn.Spec.Environment.Name, metav1.GetOptions{}) if err != nil { return nil, err } @@ -436,7 +437,7 @@ func (deploy *NewDeploy) fnCreate(fn *fv1.Function) (*fscache.FuncSvc, error) { // Since newdeploy waits for pods of deployment to be ready, // change the order of kubeObject creation (create service first, // then deployment) to take advantage of waiting time. - svc, err := deploy.createOrGetSvc(deployLabels, deployAnnotations, objName, ns) + svc, err := deploy.createOrGetSvc(ctx, deployLabels, deployAnnotations, objName, ns) if err != nil { deploy.logger.Error("error creating service", zap.Error(err), zap.String("service", objName)) go cleanupFunc(ns, objName) @@ -444,14 +445,14 @@ func (deploy *NewDeploy) fnCreate(fn *fv1.Function) (*fscache.FuncSvc, error) { } svcAddress := fmt.Sprintf("%v.%v", svc.Name, svc.Namespace) - depl, err := deploy.createOrGetDeployment(fn, env, objName, deployLabels, deployAnnotations, ns) + depl, err := deploy.createOrGetDeployment(ctx, fn, env, objName, deployLabels, deployAnnotations, ns) if err != nil { deploy.logger.Error("error creating deployment", zap.Error(err), zap.String("deployment", objName)) go cleanupFunc(ns, objName) return nil, errors.Wrapf(err, "error creating deployment %v", objName) } - hpa, err := deploy.createOrGetHpa(objName, &fn.Spec.InvokeStrategy.ExecutionStrategy, depl, deployLabels, deployAnnotations) + hpa, err := deploy.createOrGetHpa(ctx, objName, &fn.Spec.InvokeStrategy.ExecutionStrategy, depl, deployLabels, deployAnnotations) if err != nil { deploy.logger.Error("error creating HPA", zap.Error(err), zap.String("hpa", objName)) go cleanupFunc(ns, objName) @@ -506,7 +507,7 @@ func (deploy *NewDeploy) fnCreate(fn *fv1.Function) (*fscache.FuncSvc, error) { return fsvc, nil } -func (deploy *NewDeploy) updateFunction(oldFn *fv1.Function, newFn *fv1.Function) error { +func (deploy *NewDeploy) updateFunction(ctx context.Context, oldFn *fv1.Function, newFn *fv1.Function) error { if oldFn.ObjectMeta.ResourceVersion == newFn.ObjectMeta.ResourceVersion { return nil @@ -524,7 +525,7 @@ func (deploy *NewDeploy) updateFunction(oldFn *fv1.Function, newFn *fv1.Function deploy.logger.Info("function does not use new deployment executor anymore, deleting resources", zap.Any("function", newFn)) // IMP - pass the oldFn, as the new/modified function is not in cache - return deploy.deleteFunction(oldFn) + return deploy.deleteFunction(ctx, oldFn) } // Executor type changed to New Deployment from something else @@ -533,7 +534,7 @@ func (deploy *NewDeploy) updateFunction(oldFn *fv1.Function, newFn *fv1.Function deploy.logger.Info("function type changed to new deployment, creating resources", zap.Any("old_function", oldFn.ObjectMeta), zap.Any("new_function", newFn.ObjectMeta)) - _, err := deploy.createFunction(newFn) + _, err := deploy.createFunction(ctx, newFn) if err != nil { deploy.updateStatus(oldFn, err, "error changing the function's type to newdeploy") } @@ -557,7 +558,7 @@ func (deploy *NewDeploy) updateFunction(oldFn *fv1.Function, newFn *fv1.Function return err } - hpa, err := deploy.getHpa(ns, fsvc.Name) + hpa, err := deploy.getHpa(ctx, ns, fsvc.Name) if err != nil { deploy.updateStatus(oldFn, err, "error getting HPA while updating function") return err @@ -583,7 +584,7 @@ func (deploy *NewDeploy) updateFunction(oldFn *fv1.Function, newFn *fv1.Function } if hpaChanged { - err := deploy.updateHpa(hpa) + err := deploy.updateHpa(ctx, hpa) if err != nil { deploy.updateStatus(oldFn, err, "error updating HPA while updating function") return err @@ -621,18 +622,18 @@ func (deploy *NewDeploy) updateFunction(oldFn *fv1.Function, newFn *fv1.Function if deployChanged { env, err := deploy.fissionClient.CoreV1().Environments(newFn.Spec.Environment.Namespace). - Get(context.TODO(), newFn.Spec.Environment.Name, metav1.GetOptions{}) + Get(ctx, newFn.Spec.Environment.Name, metav1.GetOptions{}) if err != nil { deploy.updateStatus(oldFn, err, "failed to get environment while updating function") return err } - return deploy.updateFuncDeployment(newFn, env) + return deploy.updateFuncDeployment(ctx, newFn, env) } return nil } -func (deploy *NewDeploy) updateFuncDeployment(fn *fv1.Function, env *fv1.Environment) error { +func (deploy *NewDeploy) updateFuncDeployment(ctx context.Context, fn *fv1.Function, env *fv1.Environment) error { fsvc, err := deploy.fsCache.GetByFunctionUID(fn.ObjectMeta.UID) if err != nil { err = errors.Wrapf(err, "error updating function due to unable to find function service cache: %v", fn) @@ -651,7 +652,7 @@ func (deploy *NewDeploy) updateFuncDeployment(fn *fv1.Function, env *fv1.Environ ns = fn.ObjectMeta.Namespace } - existingDepl, err := deploy.kubernetesClient.AppsV1().Deployments(ns).Get(context.TODO(), fnObjName, metav1.GetOptions{}) + existingDepl, err := deploy.kubernetesClient.AppsV1().Deployments(ns).Get(ctx, fnObjName, metav1.GetOptions{}) if err != nil { return err } @@ -659,7 +660,7 @@ func (deploy *NewDeploy) updateFuncDeployment(fn *fv1.Function, env *fv1.Environ // the resource version inside function packageRef is changed, // so the content of fetchRequest in deployment cmd is different. // Therefore, the deployment update will trigger a rolling update. - newDeployment, err := deploy.getDeploymentSpec(fn, env, + newDeployment, err := deploy.getDeploymentSpec(ctx, fn, env, existingDepl.Spec.Replicas, // use current replicas instead of minscale in the ExecutionStrategy. fnObjName, ns, deployLabels, deploy.getDeployAnnotations(fn.ObjectMeta, env.ObjectMeta)) if err != nil { @@ -667,7 +668,7 @@ func (deploy *NewDeploy) updateFuncDeployment(fn *fv1.Function, env *fv1.Environ return err } - err = deploy.updateDeployment(newDeployment, ns) + err = deploy.updateDeployment(ctx, newDeployment, ns) if err != nil { deploy.updateStatus(fn, err, "failed to update deployment while updating function") return err @@ -676,7 +677,7 @@ func (deploy *NewDeploy) updateFuncDeployment(fn *fv1.Function, env *fv1.Environ return nil } -func (deploy *NewDeploy) fnDelete(fn *fv1.Function) error { +func (deploy *NewDeploy) fnDelete(ctx context.Context, fn *fv1.Function) error { multierr := &multierror.Error{} // GetByFunction uses resource version as part of cache key, however, @@ -705,7 +706,7 @@ func (deploy *NewDeploy) fnDelete(fn *fv1.Function) error { ns = fn.ObjectMeta.Namespace } - err = deploy.cleanupNewdeploy(ns, objName) + err = deploy.cleanupNewdeploy(ctx, ns, objName) multierr = multierror.Append(multierr, err) return multierr.ErrorOrNil() @@ -767,12 +768,12 @@ func (deploy *NewDeploy) updateStatus(fn *fv1.Function, err error, message strin // idleObjectReaper reaps objects after certain idle time func (deploy *NewDeploy) idleObjectReaper() { - + ctx := context.Background() pollSleep := 5 * time.Second for { time.Sleep(pollSleep) - envs, err := deploy.fissionClient.CoreV1().Environments(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{}) + envs, err := deploy.fissionClient.CoreV1().Environments(metav1.NamespaceAll).List(ctx, metav1.ListOptions{}) if err != nil { deploy.logger.Fatal("failed to get environment list", zap.Error(err)) } @@ -803,7 +804,7 @@ func (deploy *NewDeploy) idleObjectReaper() { zap.String("function", fsvc.Name)) } - fn, err := deploy.fissionClient.CoreV1().Functions(fsvc.Function.Namespace).Get(context.TODO(), fsvc.Function.Name, metav1.GetOptions{}) + fn, err := deploy.fissionClient.CoreV1().Functions(fsvc.Function.Namespace).Get(ctx, fsvc.Function.Name, metav1.GetOptions{}) if err != nil { // Newdeploy manager handles the function delete event and clean cache/kubeobjs itself, // so we ignore the not found error for functions with newdeploy executor type here. @@ -834,7 +835,7 @@ func (deploy *NewDeploy) idleObjectReaper() { } currentDeploy, err := deploy.kubernetesClient.AppsV1(). - Deployments(deployObj.Namespace).Get(context.TODO(), deployObj.Name, metav1.GetOptions{}) + Deployments(deployObj.Namespace).Get(ctx, deployObj.Name, metav1.GetOptions{}) if err != nil { deploy.logger.Error("error getting function deployment", zap.Error(err), zap.String("function", fsvc.Function.Name)) return @@ -847,7 +848,7 @@ func (deploy *NewDeploy) idleObjectReaper() { return } - err = deploy.scaleDeployment(deployObj.Namespace, deployObj.Name, minScale) + err = deploy.scaleDeployment(ctx, deployObj.Namespace, deployObj.Name, minScale) if err != nil { deploy.logger.Error("error scaling down function deployment", zap.Error(err), zap.String("function", fsvc.Function.Name)) } @@ -867,12 +868,12 @@ func getDeploymentObj(kubeobjs []apiv1.ObjectReference) *apiv1.ObjectReference { return nil } -func (deploy *NewDeploy) scaleDeployment(deplNS string, deplName string, replicas int32) error { +func (deploy *NewDeploy) scaleDeployment(ctx context.Context, deplNS string, deplName string, replicas int32) error { deploy.logger.Info("scaling deployment", zap.String("deployment", deplName), zap.String("namespace", deplNS), zap.Int32("replicas", replicas)) - _, err := deploy.kubernetesClient.AppsV1().Deployments(deplNS).UpdateScale(context.TODO(), deplName, &autoscalingv1.Scale{ + _, err := deploy.kubernetesClient.AppsV1().Deployments(deplNS).UpdateScale(ctx, deplName, &autoscalingv1.Scale{ ObjectMeta: metav1.ObjectMeta{ Name: deplName, Namespace: deplNS, diff --git a/pkg/executor/executortype/poolmgr/funchandlers.go b/pkg/executor/executortype/poolmgr/funchandlers.go index 12e6f5e2..442d1d97 100644 --- a/pkg/executor/executortype/poolmgr/funchandlers.go +++ b/pkg/executor/executortype/poolmgr/funchandlers.go @@ -44,6 +44,7 @@ func getIstioServiceLabels(fnName string) map[string]string { func FunctionEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clientset, fissionfnNamespace string, istioEnabled bool) k8sCache.ResourceEventHandlerFuncs { return k8sCache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { + ctx := context.Background() fn := obj.(*fv1.Function) // Since istio only allows accessing pod through k8s service, @@ -69,7 +70,7 @@ func FunctionEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clie // setup rolebinding is tried, if it fails, we don't return. we just log an error and move on, because : // 1. not all functions have secrets and/or configmaps, so things will work without this rolebinding in that case. // 2. on the contrary, when the route is tried, the env fetcher logs will show a 403 forbidden message and same will be relayed to executor. - err := utils.SetupRoleBinding(logger, kubernetesClient, fv1.SecretConfigMapGetterRB, fn.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs) + err := utils.SetupRoleBinding(ctx, logger, kubernetesClient, fv1.SecretConfigMapGetterRB, fn.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs) if err != nil { logger.Error("error creating rolebinding", zap.Error(err), zap.String("role_binding", fv1.SecretConfigMapGetterRB)) } else { @@ -120,7 +121,7 @@ func FunctionEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clie } // create function istio service if it does not exist - _, err = kubernetesClient.CoreV1().Services(envNs).Create(context.TODO(), &svc, metav1.CreateOptions{}) + _, err = kubernetesClient.CoreV1().Services(envNs).Create(ctx, &svc, metav1.CreateOptions{}) if err != nil && !kerrors.IsAlreadyExists(err) { logger.Error("error creating istio service for function", zap.Error(err), @@ -132,6 +133,7 @@ func FunctionEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clie }, DeleteFunc: func(obj interface{}) { + ctx := context.Background() fn := obj.(*fv1.Function) fnExecutorType := fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType @@ -147,7 +149,7 @@ func FunctionEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clie if istioEnabled { svcName := utils.GetFunctionIstioServiceName(fn.ObjectMeta.Name, fn.ObjectMeta.Namespace) // delete function istio service - err := kubernetesClient.CoreV1().Services(envNs).Delete(context.TODO(), svcName, metav1.DeleteOptions{}) + err := kubernetesClient.CoreV1().Services(envNs).Delete(ctx, svcName, metav1.DeleteOptions{}) if err != nil && !kerrors.IsNotFound(err) { logger.Error("error deleting istio service for function", zap.Error(err), @@ -181,8 +183,8 @@ func FunctionEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clie if newFunc.Spec.Environment.Namespace != metav1.NamespaceDefault { envNs = newFunc.Spec.Environment.Namespace } - - err := utils.SetupRoleBinding(logger, kubernetesClient, fv1.SecretConfigMapGetterRB, + ctx := context.Background() + err := utils.SetupRoleBinding(ctx, logger, kubernetesClient, fv1.SecretConfigMapGetterRB, newFunc.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs) diff --git a/pkg/executor/executortype/poolmgr/gp.go b/pkg/executor/executortype/poolmgr/gp.go index 5cf03e42..82500963 100644 --- a/pkg/executor/executortype/poolmgr/gp.go +++ b/pkg/executor/executortype/poolmgr/gp.go @@ -143,7 +143,6 @@ func MakeGenericPool( if err != nil { return nil, err } - gpLogger.Info("deployment created", zap.Any("environment", env.ObjectMeta)) go gp.startReadyPodController() go gp.updateCPUUtilizationSvc() @@ -223,7 +222,7 @@ func (gp *GenericPool) updateCPUUtilizationSvc() { // choosePod picks a ready pod from the pool and relabels it, waiting if necessary. // returns the key and pod API object. -func (gp *GenericPool) choosePod(newLabels map[string]string) (string, *apiv1.Pod, error) { +func (gp *GenericPool) choosePod(ctx context.Context, newLabels map[string]string) (string, *apiv1.Pod, error) { startTime := time.Now() expoDelay := 100 * time.Millisecond for { @@ -277,7 +276,7 @@ func (gp *GenericPool) choosePod(newLabels map[string]string) (string, *apiv1.Po patch := fmt.Sprintf(`{"metadata":{"annotations":%v, "labels":%v}}`, string(annotationPatch), string(labelPatch)) gp.logger.Info("relabel pod", zap.String("pod", patch)) - newPod, err := gp.kubernetesClient.CoreV1().Pods(chosenPod.Namespace).Patch(context.TODO(), chosenPod.Name, k8sTypes.StrategicMergePatchType, []byte(patch), metav1.PatchOptions{}) + newPod, err := gp.kubernetesClient.CoreV1().Pods(chosenPod.Namespace).Patch(ctx, chosenPod.Name, k8sTypes.StrategicMergePatchType, []byte(patch), metav1.PatchOptions{}) if err != nil { gp.logger.Error("failed to relabel pod", zap.Error(err), zap.String("pod", chosenPod.Name), zap.Duration("delay", expoDelay)) gp.readyPodQueue.Done(key) @@ -397,7 +396,7 @@ func (gp *GenericPool) specializePod(ctx context.Context, pod *apiv1.Pod, fn *fv return nil } -func (gp *GenericPool) createSvc(name string, labels map[string]string) (*apiv1.Service, error) { +func (gp *GenericPool) createSvc(ctx context.Context, name string, labels map[string]string) (*apiv1.Service, error) { service := apiv1.Service{ ObjectMeta: metav1.ObjectMeta{ Name: name, @@ -415,7 +414,7 @@ func (gp *GenericPool) createSvc(name string, labels map[string]string) (*apiv1. Selector: labels, }, } - svc, err := gp.kubernetesClient.CoreV1().Services(gp.namespace).Create(context.TODO(), &service, metav1.CreateOptions{}) + svc, err := gp.kubernetesClient.CoreV1().Services(gp.namespace).Create(ctx, &service, metav1.CreateOptions{}) return svc, err } @@ -466,7 +465,7 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac } } - key, pod, err := gp.choosePod(funcLabels) + key, pod, err := gp.choosePod(ctx, funcLabels) if err != nil { return nil, err } @@ -485,7 +484,7 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac svcName = fmt.Sprintf("%s-%v", svcName, fn.ObjectMeta.UID) } - svc, err := gp.createSvc(svcName, funcLabels) + svc, err := gp.createSvc(ctx, svcName, funcLabels) if err != nil { gp.scheduleDeletePod(pod.ObjectMeta.Name) return nil, err diff --git a/pkg/executor/executortype/poolmgr/gp_deployment.go b/pkg/executor/executortype/poolmgr/gp_deployment.go index 87881294..95e91cbd 100644 --- a/pkg/executor/executortype/poolmgr/gp_deployment.go +++ b/pkg/executor/executortype/poolmgr/gp_deployment.go @@ -199,5 +199,7 @@ func (gp *GenericPool) createPoolDeployment(ctx context.Context, env *fv1.Enviro } gp.deployment = depl + gp.logger.Info("deployment created", zap.String("deployment", depl.Name), zap.String("ns", depl.Namespace), zap.Any("environment", env.ObjectMeta)) + return nil } diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index 46a4785f..ea0c20a0 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -163,7 +163,7 @@ func (gpm *GenericPoolManager) GetTypeName(ctx context.Context) fv1.ExecutorType func (gpm *GenericPoolManager) GetFuncSvc(ctx context.Context, fn *fv1.Function) (*fscache.FuncSvc, error) { // from Func -> get Env gpm.logger.Debug("getting environment for function", zap.String("function", fn.ObjectMeta.Name)) - env, err := gpm.getFunctionEnv(fn) + env, err := gpm.getFunctionEnv(ctx, fn) if err != nil { return nil, err } @@ -207,7 +207,7 @@ func (gpm *GenericPoolManager) TapService(ctx context.Context, svcHost string) e return nil } -func (gpm *GenericPoolManager) getPodInfo(obj apiv1.ObjectReference) (*apiv1.Pod, error) { +func (gpm *GenericPoolManager) getPodInfo(ctx context.Context, obj apiv1.ObjectReference) (*apiv1.Pod, error) { store := gpm.podInformer.GetStore() item, exists, err := store.Get(obj) @@ -217,7 +217,7 @@ func (gpm *GenericPoolManager) getPodInfo(obj apiv1.ObjectReference) (*apiv1.Pod 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(context.TODO(), obj.Name, metav1.GetOptions{}) + pod, err := gpm.kubernetesClient.CoreV1().Pods(obj.Namespace).Get(ctx, obj.Name, metav1.GetOptions{}) return pod, err } @@ -230,7 +230,7 @@ func (gpm *GenericPoolManager) getPodInfo(obj apiv1.ObjectReference) (*apiv1.Pod 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(obj) + pod, err := gpm.getPodInfo(ctx, obj) 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]. @@ -437,12 +437,12 @@ func (gpm *GenericPoolManager) CleanupOldExecutorObjects(ctx context.Context) { LabelSelector: labels.Set(map[string]string{fv1.EXECUTOR_TYPE: string(fv1.ExecutorTypePoolmgr)}).AsSelector().String(), } - err := reaper.CleanupDeployments(gpm.logger, gpm.kubernetesClient, gpm.instanceID, listOpts) + err := reaper.CleanupDeployments(ctx, 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) + err = reaper.CleanupPods(ctx, gpm.logger, gpm.kubernetesClient, gpm.instanceID, listOpts) if err != nil { errs = multierror.Append(errs, err) } @@ -521,7 +521,7 @@ func (gpm *GenericPoolManager) cleanupPools(envs []fv1.Environment) { } } -func (gpm *GenericPoolManager) getFunctionEnv(fn *fv1.Function) (*fv1.Environment, error) { +func (gpm *GenericPoolManager) getFunctionEnv(ctx context.Context, fn *fv1.Function) (*fv1.Environment, error) { var env *fv1.Environment // Cached ? @@ -533,7 +533,7 @@ func (gpm *GenericPoolManager) getFunctionEnv(fn *fv1.Function) (*fv1.Environmen } // Get env from controller - env, err = gpm.fissionClient.CoreV1().Environments(fn.Spec.Environment.Namespace).Get(context.TODO(), fn.Spec.Environment.Name, metav1.GetOptions{}) + env, err = gpm.fissionClient.CoreV1().Environments(fn.Spec.Environment.Namespace).Get(ctx, fn.Spec.Environment.Name, metav1.GetOptions{}) if err != nil { return nil, err } @@ -600,13 +600,13 @@ func (gpm *GenericPoolManager) eagerPoolCreator() { // idleObjectReaper reaps objects after certain idle time func (gpm *GenericPoolManager) idleObjectReaper() { - + ctx := context.Background() pollSleep := 5 * time.Second for { time.Sleep(pollSleep) - envs, err := gpm.fissionClient.CoreV1().Environments(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{}) + envs, err := gpm.fissionClient.CoreV1().Environments(metav1.NamespaceAll).List(ctx, metav1.ListOptions{}) if err != nil { gpm.logger.Error("failed to get environment list", zap.Error(err)) continue @@ -617,7 +617,7 @@ func (gpm *GenericPoolManager) idleObjectReaper() { envList[env.ObjectMeta.UID] = struct{}{} } - fns, err := gpm.fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{}) + fns, err := gpm.fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(ctx, metav1.ListOptions{}) if err != nil { gpm.logger.Error("failed to get environment list", zap.Error(err)) continue @@ -685,7 +685,7 @@ func (gpm *GenericPoolManager) idleObjectReaper() { zap.String("executor", string(fsvc.Executor)), zap.String("pod", fsvc.Name), ) - reaper.CleanupKubeObject(gpm.logger, gpm.kubernetesClient, &fsvc.KubernetesObjects[i]) + reaper.CleanupKubeObject(ctx, gpm.logger, gpm.kubernetesClient, &fsvc.KubernetesObjects[i]) time.Sleep(50 * time.Millisecond) gpm.fsCache.ReapTime(fsvc.Function.Name, fsvc.Address, time.Since(startTime).Seconds()) } @@ -777,7 +777,8 @@ func (gpm *GenericPoolManager) NoActiveConnectionEventChecker(kubeClient *kubern zap.String("executor", string(fsvc.Executor)), zap.String("pod", fsvc.Name), ) - reaper.CleanupKubeObject(gpm.logger, gpm.kubernetesClient, &fsvc.KubernetesObjects[i]) + ctx := context.Background() + reaper.CleanupKubeObject(ctx, gpm.logger, gpm.kubernetesClient, &fsvc.KubernetesObjects[i]) time.Sleep(50 * time.Millisecond) } } diff --git a/pkg/executor/executortype/poolmgr/packagehandlers.go b/pkg/executor/executortype/poolmgr/packagehandlers.go index 30d7b445..bfa29ba3 100644 --- a/pkg/executor/executortype/poolmgr/packagehandlers.go +++ b/pkg/executor/executortype/poolmgr/packagehandlers.go @@ -17,6 +17,8 @@ limitations under the License. package poolmgr import ( + "context" + "go.uber.org/zap" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/client-go/kubernetes" @@ -42,10 +44,10 @@ func PackageEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clien if pkg.Spec.Environment.Namespace != metav1.NamespaceDefault { envNs = pkg.Spec.Environment.Namespace } - + ctx := context.Background() // here, we return if we hit an error during rolebinding setup. this is because this rolebinding is mandatory for // every function's package to be loaded into its env. without that, there's no point to move forward. - err := utils.SetupRoleBinding(logger, kubernetesClient, fv1.PackageGetterRB, pkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs) + err := utils.SetupRoleBinding(ctx, logger, kubernetesClient, fv1.PackageGetterRB, pkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs) if err != nil { logger.Error("error creating rolebinding for package", zap.Error(err), @@ -79,7 +81,8 @@ func PackageEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clien envNs = newPkg.Spec.Environment.Namespace } - err := utils.SetupRoleBinding(logger, kubernetesClient, fv1.PackageGetterRB, + ctx := context.Background() + err := utils.SetupRoleBinding(ctx, logger, kubernetesClient, fv1.PackageGetterRB, newPkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs) if err != nil { diff --git a/pkg/executor/reaper/reaper.go b/pkg/executor/reaper/reaper.go index b0c93906..335871c1 100644 --- a/pkg/executor/reaper/reaper.go +++ b/pkg/executor/reaper/reaper.go @@ -37,28 +37,28 @@ var ( ) // CleanupKubeObject deletes given kubernetes object -func CleanupKubeObject(logger *zap.Logger, kubeClient *kubernetes.Clientset, kubeobj *apiv1.ObjectReference) { +func CleanupKubeObject(ctx context.Context, logger *zap.Logger, kubeClient *kubernetes.Clientset, kubeobj *apiv1.ObjectReference) { switch strings.ToLower(kubeobj.Kind) { case "pod": - err := kubeClient.CoreV1().Pods(kubeobj.Namespace).Delete(context.TODO(), kubeobj.Name, meta_v1.DeleteOptions{}) + err := kubeClient.CoreV1().Pods(kubeobj.Namespace).Delete(ctx, kubeobj.Name, meta_v1.DeleteOptions{}) if err != nil { logger.Error("error cleaning up pod", zap.Error(err), zap.String("pod", kubeobj.Name)) } case "service": - err := kubeClient.CoreV1().Services(kubeobj.Namespace).Delete(context.TODO(), kubeobj.Name, meta_v1.DeleteOptions{}) + err := kubeClient.CoreV1().Services(kubeobj.Namespace).Delete(ctx, kubeobj.Name, meta_v1.DeleteOptions{}) if err != nil { logger.Error("error cleaning up service", zap.Error(err), zap.String("service", kubeobj.Name)) } case "deployment": - err := kubeClient.AppsV1().Deployments(kubeobj.Namespace).Delete(context.TODO(), kubeobj.Name, delOpt) + err := kubeClient.AppsV1().Deployments(kubeobj.Namespace).Delete(ctx, kubeobj.Name, delOpt) if err != nil { logger.Error("error cleaning up deployment", zap.Error(err), zap.String("deployment", kubeobj.Name)) } case "horizontalpodautoscaler": - err := kubeClient.AutoscalingV1().HorizontalPodAutoscalers(kubeobj.Namespace).Delete(context.TODO(), kubeobj.Name, meta_v1.DeleteOptions{}) + err := kubeClient.AutoscalingV1().HorizontalPodAutoscalers(kubeobj.Namespace).Delete(ctx, kubeobj.Name, meta_v1.DeleteOptions{}) if err != nil { logger.Error("error cleaning up horizontalpodautoscaler", zap.Error(err), zap.String("horizontalpodautoscaler", kubeobj.Name)) } @@ -70,8 +70,8 @@ func CleanupKubeObject(logger *zap.Logger, kubeClient *kubernetes.Clientset, kub } // CleanupDeployments deletes deployment(s) for a given instanceID -func CleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error { - deploymentList, err := client.AppsV1().Deployments(meta_v1.NamespaceAll).List(context.TODO(), listOps) +func CleanupDeployments(ctx context.Context, logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error { + deploymentList, err := client.AppsV1().Deployments(meta_v1.NamespaceAll).List(ctx, listOps) if err != nil { return err } @@ -83,7 +83,7 @@ func CleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instan } if ok && id != instanceID { logger.Info("cleaning up deployment", zap.String("deployment", dep.ObjectMeta.Name)) - err := client.AppsV1().Deployments(dep.ObjectMeta.Namespace).Delete(context.TODO(), dep.ObjectMeta.Name, delOpt) + err := client.AppsV1().Deployments(dep.ObjectMeta.Namespace).Delete(ctx, dep.ObjectMeta.Name, delOpt) if err != nil { logger.Error("error cleaning up deployment", zap.Error(err), @@ -97,8 +97,8 @@ func CleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instan } // CleanupPods deletes pod(s) for a given instanceID -func CleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error { - podList, err := client.CoreV1().Pods(meta_v1.NamespaceAll).List(context.TODO(), listOps) +func CleanupPods(ctx context.Context, logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error { + podList, err := client.CoreV1().Pods(meta_v1.NamespaceAll).List(ctx, listOps) if err != nil { return err } @@ -110,7 +110,7 @@ func CleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceID st } if ok && id != instanceID { logger.Info("cleaning up pod", zap.String("pod", pod.ObjectMeta.Name)) - err := client.CoreV1().Pods(pod.ObjectMeta.Namespace).Delete(context.TODO(), pod.ObjectMeta.Name, meta_v1.DeleteOptions{}) + err := client.CoreV1().Pods(pod.ObjectMeta.Namespace).Delete(ctx, pod.ObjectMeta.Name, meta_v1.DeleteOptions{}) if err != nil { logger.Error("error cleaning up pod", zap.Error(err), @@ -124,8 +124,8 @@ func CleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceID st } // CleanupServices deletes service(s) for a given instanceID -func CleanupServices(logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error { - svcList, err := client.CoreV1().Services(meta_v1.NamespaceAll).List(context.TODO(), listOps) +func CleanupServices(ctx context.Context, logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error { + svcList, err := client.CoreV1().Services(meta_v1.NamespaceAll).List(ctx, listOps) if err != nil { return err } @@ -137,7 +137,7 @@ func CleanupServices(logger *zap.Logger, client *kubernetes.Clientset, instanceI } if ok && id != instanceID { logger.Info("cleaning up service", zap.String("service", svc.ObjectMeta.Name)) - err := client.CoreV1().Services(svc.ObjectMeta.Namespace).Delete(context.TODO(), svc.ObjectMeta.Name, meta_v1.DeleteOptions{}) + err := client.CoreV1().Services(svc.ObjectMeta.Namespace).Delete(ctx, svc.ObjectMeta.Name, meta_v1.DeleteOptions{}) if err != nil { logger.Error("error cleaning up service", zap.Error(err), @@ -151,8 +151,8 @@ func CleanupServices(logger *zap.Logger, client *kubernetes.Clientset, instanceI } // CleanupHpa deletes horizontal pod autoscaler(s) for a given instanceID -func CleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error { - hpaList, err := client.AutoscalingV1().HorizontalPodAutoscalers(meta_v1.NamespaceAll).List(context.TODO(), listOps) +func CleanupHpa(ctx context.Context, logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error { + hpaList, err := client.AutoscalingV1().HorizontalPodAutoscalers(meta_v1.NamespaceAll).List(ctx, listOps) if err != nil { return err } @@ -165,7 +165,7 @@ func CleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceID str } if ok && id != instanceID { logger.Info("cleaning up HPA", zap.String("hpa", hpa.ObjectMeta.Name)) - err := client.AutoscalingV1().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Delete(context.TODO(), hpa.ObjectMeta.Name, meta_v1.DeleteOptions{}) + err := client.AutoscalingV1().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Delete(ctx, hpa.ObjectMeta.Name, meta_v1.DeleteOptions{}) if err != nil { logger.Error("error cleaning up HPA", zap.Error(err), @@ -181,13 +181,14 @@ func CleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceID str // CleanupRoleBindings periodically lists rolebindings across all namespaces and removes Service Accounts from them or // deletes the rolebindings completely if there are no Service Accounts in a rolebinding object. func CleanupRoleBindings(logger *zap.Logger, client *kubernetes.Clientset, fissionClient *crd.FissionClient, functionNs, envBuilderNs string, cleanupRoleBindingInterval time.Duration) { + ctx := context.Background() for { // some sleep before the next reaper iteration time.Sleep(cleanupRoleBindingInterval) logger.Debug("starting cleanupRoleBindings cycle") // get all rolebindings ( just to be efficient, one call to kubernetes ) - rbList, err := client.RbacV1().RoleBindings(meta_v1.NamespaceAll).List(context.TODO(), meta_v1.ListOptions{}) + rbList, err := client.RbacV1().RoleBindings(meta_v1.NamespaceAll).List(ctx, meta_v1.ListOptions{}) if err != nil { // something wrong, but next iteration hopefully succeeds logger.Error("error listing role bindings in all namespaces", zap.Error(err)) @@ -208,7 +209,7 @@ func CleanupRoleBindings(logger *zap.Logger, client *kubernetes.Clientset, fissi // in order to find out if there are any functions that need this role-binding in role-binding namespace, // we can list the functions once per role-binding. - funcList, err := fissionClient.CoreV1().Functions(roleBinding.Namespace).List(context.TODO(), meta_v1.ListOptions{}) + funcList, err := fissionClient.CoreV1().Functions(roleBinding.Namespace).List(ctx, meta_v1.ListOptions{}) if err != nil { logger.Error("error fetching function list in namespace", zap.Error(err), zap.String("namespace", roleBinding.Namespace)) continue @@ -259,7 +260,7 @@ func CleanupRoleBindings(logger *zap.Logger, client *kubernetes.Clientset, fissi // else if its a secret-configmap-rb, we have only one SA which is fission-fetcher if roleBinding.Name == fv1.PackageGetterRB { // check if there is an env obj in saNs - envList, err := fissionClient.CoreV1().Environments(saNs).List(context.TODO(), meta_v1.ListOptions{}) + envList, err := fissionClient.CoreV1().Environments(saNs).List(ctx, meta_v1.ListOptions{}) if err != nil { logger.Error("error fetching environment list in service account namespace", zap.Error(err), zap.String("namespace", saNs)) continue @@ -305,7 +306,7 @@ func CleanupRoleBindings(logger *zap.Logger, client *kubernetes.Clientset, fissi zap.String("role_binding_namespace", roleBinding.Namespace)) // call this once in the end for each role-binding - err = utils.RemoveSAFromRoleBindingWithRetries(logger, client, roleBinding.Name, roleBinding.Namespace, saToRemove) + err = utils.RemoveSAFromRoleBindingWithRetries(ctx, logger, client, roleBinding.Name, roleBinding.Namespace, saToRemove) if err != nil { // if there's an error, we just log it and proceed with the next role-binding, hoping that this role-binding // will be processed in next iteration. diff --git a/pkg/executor/util/util.go b/pkg/executor/util/util.go index b40ff27b..58df2c13 100644 --- a/pkg/executor/util/util.go +++ b/pkg/executor/util/util.go @@ -57,7 +57,7 @@ func WaitTimeout(wg *sync.WaitGroup, timeout time.Duration) { } // ConvertConfigSecrets returns envFromSource which can be passed directly into the pod spec -func ConvertConfigSecrets(fn *fv1.Function, kc *kubernetes.Clientset) ([]apiv1.EnvFromSource, error) { +func ConvertConfigSecrets(ctx context.Context, fn *fv1.Function, kc *kubernetes.Clientset) ([]apiv1.EnvFromSource, error) { cmList := fn.Spec.ConfigMaps secList := fn.Spec.Secrets @@ -65,9 +65,9 @@ func ConvertConfigSecrets(fn *fv1.Function, kc *kubernetes.Clientset) ([]apiv1.E secEnvSources := make([]*apiv1.SecretEnvSource, 0) for _, cm := range cmList { if cm.Namespace != fn.Namespace { - return nil, errors.New("Function should not reference config map of different namespace") + return nil, errors.New("function should not reference config map of different namespace") } - _, err := kc.CoreV1().ConfigMaps(cm.Namespace).Get(context.TODO(), cm.Name, metav1.GetOptions{}) + _, err := kc.CoreV1().ConfigMaps(cm.Namespace).Get(ctx, cm.Name, metav1.GetOptions{}) if err != nil { return nil, err } @@ -81,9 +81,9 @@ func ConvertConfigSecrets(fn *fv1.Function, kc *kubernetes.Clientset) ([]apiv1.E for _, sec := range secList { if sec.Namespace != fn.Namespace { - return nil, errors.New("Function should not reference secret of different namespace") + return nil, errors.New("function should not reference secret of different namespace") } - _, err := kc.CoreV1().Secrets(sec.Namespace).Get(context.TODO(), sec.Name, metav1.GetOptions{}) + _, err := kc.CoreV1().Secrets(sec.Namespace).Get(ctx, sec.Name, metav1.GetOptions{}) if err != nil { return nil, err } diff --git a/pkg/utils/rbacutils.go b/pkg/utils/rbacutils.go index af8e9375..cdaea7e7 100644 --- a/pkg/utils/rbacutils.go +++ b/pkg/utils/rbacutils.go @@ -104,7 +104,7 @@ type PatchSpec struct { } // AddSaToRoleBindingWithRetries adds a service account to a rolebinding object. IT retries on already exists and conflict errors. -func AddSaToRoleBindingWithRetries(logger *zap.Logger, k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs, sa, saNamespace, role, roleKind string) (err error) { +func AddSaToRoleBindingWithRetries(ctx context.Context, logger *zap.Logger, k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs, sa, saNamespace, role, roleKind string) (err error) { patch := PatchSpec{} patch.Op = "add" patch.Path = "/subjects/-" @@ -121,7 +121,7 @@ func AddSaToRoleBindingWithRetries(logger *zap.Logger, k8sClient *kubernetes.Cli } for i := 0; i < maxRetries; i++ { - _, err = k8sClient.RbacV1().RoleBindings(roleBindingNs).Patch(context.TODO(), roleBinding, types.JSONPatchType, patchJson, metav1.PatchOptions{}) + _, err = k8sClient.RbacV1().RoleBindings(roleBindingNs).Patch(ctx, roleBinding, types.JSONPatchType, patchJson, metav1.PatchOptions{}) if err == nil { logger.Debug("patched rolebinding", zap.String("role_binding", roleBinding), @@ -137,7 +137,7 @@ func AddSaToRoleBindingWithRetries(logger *zap.Logger, k8sClient *kubernetes.Cli // someone may have deleted the object between us checking if the object is present and deciding to patch // so just create the object again rbObj := makeRoleBindingObj(roleBinding, roleBindingNs, role, roleKind, sa, saNamespace) - _, err = k8sClient.RbacV1().RoleBindings(roleBindingNs).Create(context.TODO(), rbObj, metav1.CreateOptions{}) + _, err = k8sClient.RbacV1().RoleBindings(roleBindingNs).Create(ctx, rbObj, metav1.CreateOptions{}) if err == nil { logger.Debug("created rolebinding", zap.String("role_binding", roleBinding), @@ -176,10 +176,10 @@ func AddSaToRoleBindingWithRetries(logger *zap.Logger, k8sClient *kubernetes.Cli // RemoveSAFromRoleBindingWithRetries removes an SA from the rolebinding passed as parameter. If this is the only SA in // the rolebinding, then it deletes the rolebinding object. -func RemoveSAFromRoleBindingWithRetries(logger *zap.Logger, k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs string, saToRemove map[string]bool) (err error) { +func RemoveSAFromRoleBindingWithRetries(ctx context.Context, logger *zap.Logger, k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs string, saToRemove map[string]bool) (err error) { for i := 0; i < maxRetries; i++ { rbObj, err := k8sClient.RbacV1().RoleBindings(roleBindingNs).Get( - context.TODO(), + ctx, roleBinding, metav1.GetOptions{}) if err != nil { // silently ignoring the error. there's no need for us to remove sa anymore. @@ -206,13 +206,13 @@ func RemoveSAFromRoleBindingWithRetries(logger *zap.Logger, k8sClient *kubernete }) } if len(newSubjects) == 0 { - return DeleteRoleBinding(k8sClient, roleBinding, roleBindingNs) + return DeleteRoleBinding(ctx, k8sClient, roleBinding, roleBindingNs) } rbObj.Subjects = newSubjects // cant use patch for deletes, the results become in-deterministic, so using update. - _, err = k8sClient.RbacV1().RoleBindings(rbObj.Namespace).Update(context.TODO(), rbObj, metav1.UpdateOptions{}) + _, err = k8sClient.RbacV1().RoleBindings(rbObj.Namespace).Update(ctx, rbObj, metav1.UpdateOptions{}) switch { case err == nil: logger.Debug("removed service accounts from rolebinding", @@ -236,10 +236,10 @@ func RemoveSAFromRoleBindingWithRetries(logger *zap.Logger, k8sClient *kubernete // SetupRoleBinding adds a role to a service account if the rolebinding object is already present in the namespace. // if not, it creates a rolebinding object granting the role to the SA in the namespace. -func SetupRoleBinding(logger *zap.Logger, k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs, role, roleKind, sa, saNamespace string) error { +func SetupRoleBinding(ctx context.Context, logger *zap.Logger, k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs, role, roleKind, sa, saNamespace string) error { // get the role binding object rbObj, err := k8sClient.RbacV1().RoleBindings(roleBindingNs).Get( - context.TODO(), + ctx, roleBinding, metav1.GetOptions{}) if err == nil { @@ -249,7 +249,7 @@ func SetupRoleBinding(logger *zap.Logger, k8sClient *kubernetes.Clientset, roleB zap.String("service_account_namespace", saNamespace), zap.String("role_binding", roleBinding), zap.String("role_binding_namespace", roleBindingNs)) - return AddSaToRoleBindingWithRetries(logger, k8sClient, roleBinding, roleBindingNs, sa, saNamespace, role, roleKind) + return AddSaToRoleBindingWithRetries(ctx, logger, k8sClient, roleBinding, roleBindingNs, sa, saNamespace, role, roleKind) } logger.Debug("service account already present in rolebinding so nothing to add", zap.String("service_account_name", sa), @@ -266,14 +266,14 @@ func SetupRoleBinding(logger *zap.Logger, k8sClient *kubernetes.Clientset, roleB zap.String("role_binding", roleBinding), zap.String("role_binding_namespace", roleBindingNs)) rbObj = makeRoleBindingObj(roleBinding, roleBindingNs, role, roleKind, sa, saNamespace) - _, err = k8sClient.RbacV1().RoleBindings(roleBindingNs).Create(context.TODO(), rbObj, metav1.CreateOptions{}) + _, err = k8sClient.RbacV1().RoleBindings(roleBindingNs).Create(ctx, rbObj, metav1.CreateOptions{}) if k8serrors.IsAlreadyExists(err) { logger.Debug("rolebinding already exists in namespace - adding service account to rolebinding", zap.String("service_account_name", sa), zap.String("service_account_namespace", saNamespace), zap.String("role_binding", roleBinding), zap.String("role_binding_namespace", roleBindingNs)) - err = AddSaToRoleBindingWithRetries(logger, k8sClient, roleBinding, roleBindingNs, sa, saNamespace, role, roleKind) + err = AddSaToRoleBindingWithRetries(ctx, logger, k8sClient, roleBinding, roleBindingNs, sa, saNamespace, role, roleKind) } } @@ -282,10 +282,10 @@ func SetupRoleBinding(logger *zap.Logger, k8sClient *kubernetes.Clientset, roleB // DeleteRoleBinding deletes a rolebinding object. if k8s throws an error that the rolebinding is not there, it just // returns silently. -func DeleteRoleBinding(k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs string) error { +func DeleteRoleBinding(ctx context.Context, k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs string) error { // if deleteRoleBinding is invoked by 2 fission services at the same time for the same rolebinding, // the first call will succeed while the 2nd will fail with isNotFound. but we dont want to error out then. - err := k8sClient.RbacV1().RoleBindings(roleBindingNs).Delete(context.TODO(), roleBinding, metav1.DeleteOptions{}) + err := k8sClient.RbacV1().RoleBindings(roleBindingNs).Delete(ctx, roleBinding, metav1.DeleteOptions{}) if err == nil || k8serrors.IsNotFound(err) { return nil }