From 3fa0f4bde333ec6f65a27dd3da9747ce1c62f2e9 Mon Sep 17 00:00:00 2001 From: Sanket Sudake Date: Mon, 26 Sep 2022 16:05:45 +0530 Subject: [PATCH] Ensuring passing context across fission (#2555) Signed-off-by: Sanket Sudake --- pkg/buildermgr/buildermgr.go | 2 +- pkg/buildermgr/common.go | 4 +- pkg/buildermgr/envwatcher.go | 67 ++++++++-------- pkg/buildermgr/pkgwatcher.go | 34 ++++---- pkg/canaryconfigmgr/canaryConfigMgr.go | 52 ++++++------- pkg/controller/api_test.go | 2 +- pkg/controller/config.go | 10 +-- pkg/executor/client/client.go | 2 +- pkg/executor/executor.go | 4 +- pkg/executor/executor_test.go | 44 +++++------ .../executortype/newdeploy/envhandlers.go | 3 +- .../executortype/newdeploy/funchandlers.go | 5 +- .../executortype/newdeploy/newdeploy.go | 2 +- .../executortype/newdeploy/newdeploymgr.go | 14 ++-- .../newdeploy/newdeploymgr_test.go | 2 +- .../executortype/poolmgr/funchandlers.go | 5 +- pkg/executor/executortype/poolmgr/gp.go | 48 ++++++------ pkg/executor/executortype/poolmgr/gpm.go | 23 +++--- .../executortype/poolmgr/packagehandlers.go | 4 +- .../executortype/poolmgr/poolpodcontroller.go | 27 +++---- .../poolmgr/poolpodcontroller_test.go | 11 ++- .../fscache/functionServiceCache_test.go | 4 +- pkg/executor/util/util_test.go | 6 +- pkg/fetcher/config/config.go | 5 +- pkg/fission-cli/cmd/check/check.go | 2 +- pkg/fission-cli/cmd/plugin/list.go | 2 +- .../cmd/support/resources/fissionversion.go | 2 +- pkg/fission-cli/cmd/version/version.go | 2 +- pkg/fission-cli/util/util.go | 4 +- pkg/healthcheck/healthcheck.go | 38 ++++----- pkg/kubewatcher/main.go | 2 +- pkg/kubewatcher/watchSync.go | 8 +- pkg/mqtrigger/scalermanager.go | 78 +++++++++---------- pkg/mqtrigger/scalermanager_test.go | 13 ++-- pkg/plugin/plugin.go | 14 ++-- pkg/plugin/plugin_test.go | 6 +- pkg/poolcache/poolcache_test.go | 3 +- pkg/router/functionHandler.go | 6 +- pkg/router/httpTriggers.go | 6 +- pkg/router/ingress.go | 22 +++--- pkg/tracker/tracker_test.go | 4 +- pkg/utils/otel/provider_test.go | 3 +- pkg/utils/rbacutils.go | 6 +- 43 files changed, 301 insertions(+), 300 deletions(-) diff --git a/pkg/buildermgr/buildermgr.go b/pkg/buildermgr/buildermgr.go index c5642b29..b94cbbd6 100644 --- a/pkg/buildermgr/buildermgr.go +++ b/pkg/buildermgr/buildermgr.go @@ -64,7 +64,7 @@ func Start(ctx context.Context, logger *zap.Logger, storageSvcUrl string, envBui } envWatcher := makeEnvironmentWatcher(bmLogger, fissionClient, kubernetesClient, fetcherConfig, envBuilderNamespace, podSpecPatch) - go envWatcher.watchEnvironments() + go envWatcher.watchEnvironments(ctx) k8sInformerFactory := k8sInformers.NewSharedInformerFactory(kubernetesClient, time.Minute*30) informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) diff --git a/pkg/buildermgr/common.go b/pkg/buildermgr/common.go index 0f7d6a97..c3bbbf5d 100644 --- a/pkg/buildermgr/common.go +++ b/pkg/buildermgr/common.go @@ -121,7 +121,7 @@ func buildPackage(ctx context.Context, logger *zap.Logger, fissionClient version return uploadResp, buildResp.BuildLogs, nil } -func updatePackage(logger *zap.Logger, fissionClient versioned.Interface, +func updatePackage(ctx context.Context, logger *zap.Logger, fissionClient versioned.Interface, pkg *fv1.Package, status fv1.BuildStatus, buildLogs string, uploadResp *fetcher.ArchiveUploadResponse) (*fv1.Package, error) { @@ -140,7 +140,7 @@ func updatePackage(logger *zap.Logger, fissionClient versioned.Interface, } // update package spec - pkg, err := fissionClient.CoreV1().Packages(pkg.ObjectMeta.Namespace).Update(context.TODO(), pkg, metav1.UpdateOptions{}) + pkg, err := fissionClient.CoreV1().Packages(pkg.ObjectMeta.Namespace).Update(ctx, pkg, metav1.UpdateOptions{}) if err != nil { e := "error updating package" logger.Error(e, zap.Error(err)) diff --git a/pkg/buildermgr/envwatcher.go b/pkg/buildermgr/envwatcher.go index d34e9999..7925316a 100644 --- a/pkg/buildermgr/envwatcher.go +++ b/pkg/buildermgr/envwatcher.go @@ -67,6 +67,7 @@ type ( envwRequest struct { requestType + ctx context.Context env *fv1.Environment envList []fv1.Environment respChan chan envwResponse @@ -148,10 +149,10 @@ func (envw *environmentWatcher) getLabels(envName string, envNamespace string, e } } -func (envw *environmentWatcher) watchEnvironments() { +func (envw *environmentWatcher) watchEnvironments(ctx context.Context) { rv := "" for { - wi, err := envw.fissionClient.CoreV1().Environments(metav1.NamespaceAll).Watch(context.TODO(), + wi, err := envw.fissionClient.CoreV1().Environments(metav1.NamespaceAll).Watch(ctx, metav1.ListOptions{ ResourceVersion: rv, }) @@ -178,15 +179,15 @@ func (envw *environmentWatcher) watchEnvironments() { } env := ev.Object.(*fv1.Environment) rv = env.ObjectMeta.ResourceVersion - envw.sync() + envw.sync(ctx) } } } -func (envw *environmentWatcher) sync() { +func (envw *environmentWatcher) sync(ctx context.Context) { maxRetries := 10 for i := 0; i < maxRetries; i++ { - envList, err := envw.fissionClient.CoreV1().Environments(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{}) + envList, err := envw.fissionClient.CoreV1().Environments(metav1.NamespaceAll).List(ctx, metav1.ListOptions{}) if err != nil { if utils.IsNetworkError(err) { envw.logger.Error("error syncing environment CRD resources due to network error, retrying later", zap.Error(err)) @@ -204,14 +205,14 @@ func (envw *environmentWatcher) sync() { len(env.Spec.Builder.Image) == 0 { // ignore env without builder image continue } - _, err := envw.getEnvBuilder(&env) + _, err := envw.getEnvBuilder(ctx, &env) if err != nil { envw.logger.Error("error creating builder", zap.Error(err), zap.String("builder_target", env.ObjectMeta.Name)) } } // Remove environment builders no longer needed - envw.cleanupEnvBuilders(envList.Items) + envw.cleanupEnvBuilders(ctx, envList.Items) break } } @@ -231,7 +232,7 @@ func (envw *environmentWatcher) service() { key := envw.getCacheKey(req.env.ObjectMeta.Name, ns, req.env.ObjectMeta.ResourceVersion) builderInfo, ok := envw.cache[key] if !ok { - builderInfo, err := envw.createBuilder(req.env, ns) + builderInfo, err := envw.createBuilder(req.ctx, req.env, ns) if err != nil { req.respChan <- envwResponse{err: err} continue @@ -260,7 +261,7 @@ func (envw *environmentWatcher) service() { // cache and CRD. We need to iterate over the services & // deployments to remove both normal and orphan builders. - svcList, err := envw.getBuilderServiceList(envw.getLabelForDeploymentOwner(), metav1.NamespaceAll) + svcList, err := envw.getBuilderServiceList(req.ctx, envw.getLabelForDeploymentOwner(), metav1.NamespaceAll) if err != nil { envw.logger.Error("error getting the builder service list", zap.Error(err)) } @@ -270,7 +271,7 @@ func (envw *environmentWatcher) service() { envResourceVersion := svc.ObjectMeta.Labels[LABEL_ENV_RESOURCEVERSION] key := envw.getCacheKey(envName, envNamespace, envResourceVersion) if _, ok := latestEnvList[key]; !ok { - err := envw.deleteBuilderServiceByName(svc.ObjectMeta.Name, svc.ObjectMeta.Namespace) + err := envw.deleteBuilderServiceByName(req.ctx, svc.ObjectMeta.Name, svc.ObjectMeta.Namespace) if err != nil { envw.logger.Error("error removing builder service", zap.Error(err), zap.String("service_name", svc.ObjectMeta.Name), @@ -280,7 +281,7 @@ func (envw *environmentWatcher) service() { delete(envw.cache, key) } - deployList, err := envw.getBuilderDeploymentList(envw.getLabelForDeploymentOwner(), metav1.NamespaceAll) + deployList, err := envw.getBuilderDeploymentList(req.ctx, envw.getLabelForDeploymentOwner(), metav1.NamespaceAll) if err != nil { envw.logger.Error("error getting the builder deployment list", zap.Error(err)) } @@ -290,7 +291,7 @@ func (envw *environmentWatcher) service() { envResourceVersion := deploy.ObjectMeta.Labels[LABEL_ENV_RESOURCEVERSION] key := envw.getCacheKey(envName, envNamespace, envResourceVersion) if _, ok := latestEnvList[key]; !ok { - err := envw.deleteBuilderDeploymentByName(deploy.ObjectMeta.Name, deploy.ObjectMeta.Namespace) + err := envw.deleteBuilderDeploymentByName(req.ctx, deploy.ObjectMeta.Name, deploy.ObjectMeta.Namespace) if err != nil { envw.logger.Error("error removing builder deployment", zap.Error(err), zap.String("deployment_name", deploy.ObjectMeta.Name), @@ -303,10 +304,11 @@ func (envw *environmentWatcher) service() { } } -func (envw *environmentWatcher) getEnvBuilder(env *fv1.Environment) (*builderInfo, error) { +func (envw *environmentWatcher) getEnvBuilder(ctx context.Context, env *fv1.Environment) (*builderInfo, error) { respChan := make(chan envwResponse) envw.requestChan <- envwRequest{ requestType: GET_BUILDER, + ctx: ctx, env: env, respChan: respChan, } @@ -314,26 +316,27 @@ func (envw *environmentWatcher) getEnvBuilder(env *fv1.Environment) (*builderInf return resp.builderInfo, resp.err } -func (envw *environmentWatcher) cleanupEnvBuilders(envs []fv1.Environment) { +func (envw *environmentWatcher) cleanupEnvBuilders(ctx context.Context, envs []fv1.Environment) { envw.requestChan <- envwRequest{ requestType: CLEANUP_BUILDERS, + ctx: ctx, envList: envs, } } -func (envw *environmentWatcher) createBuilder(env *fv1.Environment, ns string) (*builderInfo, error) { +func (envw *environmentWatcher) createBuilder(ctx context.Context, env *fv1.Environment, ns string) (*builderInfo, error) { var svc *apiv1.Service var deploy *appsv1.Deployment sel := envw.getLabels(env.ObjectMeta.Name, ns, env.ObjectMeta.ResourceVersion) - svcList, err := envw.getBuilderServiceList(sel, ns) + svcList, err := envw.getBuilderServiceList(ctx, sel, ns) if err != nil { return nil, err } // there should be only one service in svcList if len(svcList) == 0 { - svc, err = envw.createBuilderService(env, ns) + svc, err = envw.createBuilderService(ctx, env, ns) if err != nil { return nil, errors.Wrap(err, "error creating builder service") } @@ -343,19 +346,19 @@ func (envw *environmentWatcher) createBuilder(env *fv1.Environment, ns string) ( return nil, fmt.Errorf("found more than one builder service for environment %q", env.ObjectMeta.Name) } - deployList, err := envw.getBuilderDeploymentList(sel, ns) + deployList, err := envw.getBuilderDeploymentList(ctx, sel, ns) if err != nil { return nil, err } // there should be only one deploy in deployList if len(deployList) == 0 { // create builder SA in this ns, if not already created - _, err := utils.SetupSA(envw.kubernetesClient, fv1.FissionBuilderSA, ns) + _, err := utils.SetupSA(ctx, envw.kubernetesClient, fv1.FissionBuilderSA, ns) if err != nil { return nil, errors.Wrapf(err, "error creating %q in ns: %s", fv1.FissionBuilderSA, ns) } - deploy, err = envw.createBuilderDeployment(env, ns) + deploy, err = envw.createBuilderDeployment(ctx, env, ns) if err != nil { return nil, errors.Wrap(err, "error creating builder deployment") } @@ -372,29 +375,29 @@ func (envw *environmentWatcher) createBuilder(env *fv1.Environment, ns string) ( }, nil } -func (envw *environmentWatcher) deleteBuilderServiceByName(name, namespace string) error { +func (envw *environmentWatcher) deleteBuilderServiceByName(ctx context.Context, name, namespace string) error { err := envw.kubernetesClient.CoreV1(). Services(namespace). - Delete(context.TODO(), name, delOpt) + Delete(ctx, name, delOpt) if err != nil { return errors.Wrapf(err, "error deleting builder service %s.%s", name, namespace) } return nil } -func (envw *environmentWatcher) deleteBuilderDeploymentByName(name, namespace string) error { +func (envw *environmentWatcher) deleteBuilderDeploymentByName(ctx context.Context, name, namespace string) error { err := envw.kubernetesClient.AppsV1(). Deployments(namespace). - Delete(context.TODO(), name, delOpt) + Delete(ctx, name, delOpt) if err != nil { return errors.Wrapf(err, "error deleting builder deployment %s.%s", name, namespace) } return nil } -func (envw *environmentWatcher) getBuilderServiceList(sel map[string]string, ns string) ([]apiv1.Service, error) { +func (envw *environmentWatcher) getBuilderServiceList(ctx context.Context, sel map[string]string, ns string) ([]apiv1.Service, error) { svcList, err := envw.kubernetesClient.CoreV1().Services(ns).List( - context.TODO(), + ctx, metav1.ListOptions{ LabelSelector: labels.Set(sel).AsSelector().String(), }) @@ -404,7 +407,7 @@ func (envw *environmentWatcher) getBuilderServiceList(sel map[string]string, ns return svcList.Items, nil } -func (envw *environmentWatcher) createBuilderService(env *fv1.Environment, ns string) (*apiv1.Service, error) { +func (envw *environmentWatcher) createBuilderService(ctx context.Context, env *fv1.Environment, ns string) (*apiv1.Service, error) { name := fmt.Sprintf("%v-%v", env.ObjectMeta.Name, env.ObjectMeta.ResourceVersion) sel := envw.getLabels(env.ObjectMeta.Name, ns, env.ObjectMeta.ResourceVersion) service := apiv1.Service{ @@ -439,16 +442,16 @@ func (envw *environmentWatcher) createBuilderService(env *fv1.Environment, ns st }, } envw.logger.Info("creating builder service", zap.String("service_name", name)) - _, err := envw.kubernetesClient.CoreV1().Services(ns).Create(context.TODO(), &service, metav1.CreateOptions{}) + _, err := envw.kubernetesClient.CoreV1().Services(ns).Create(ctx, &service, metav1.CreateOptions{}) if err != nil { return nil, err } return &service, nil } -func (envw *environmentWatcher) getBuilderDeploymentList(sel map[string]string, ns string) ([]appsv1.Deployment, error) { +func (envw *environmentWatcher) getBuilderDeploymentList(ctx context.Context, sel map[string]string, ns string) ([]appsv1.Deployment, error) { deployList, err := envw.kubernetesClient.AppsV1().Deployments(ns).List( - context.TODO(), + ctx, metav1.ListOptions{ LabelSelector: labels.Set(sel).AsSelector().String(), }) @@ -458,7 +461,7 @@ func (envw *environmentWatcher) getBuilderDeploymentList(sel map[string]string, return deployList.Items, nil } -func (envw *environmentWatcher) createBuilderDeployment(env *fv1.Environment, ns string) (*appsv1.Deployment, error) { +func (envw *environmentWatcher) createBuilderDeployment(ctx context.Context, env *fv1.Environment, ns string) (*appsv1.Deployment, error) { name := fmt.Sprintf("%v-%v", env.ObjectMeta.Name, env.ObjectMeta.ResourceVersion) sel := envw.getLabels(env.ObjectMeta.Name, ns, env.ObjectMeta.ResourceVersion) var replicas int32 = 1 @@ -546,7 +549,7 @@ func (envw *environmentWatcher) createBuilderDeployment(env *fv1.Environment, ns deployment.Spec.Template.Spec = *newPodSpec } - _, err = envw.kubernetesClient.AppsV1().Deployments(ns).Create(context.TODO(), deployment, metav1.CreateOptions{}) + _, err = envw.kubernetesClient.AppsV1().Deployments(ns).Create(ctx, deployment, metav1.CreateOptions{}) if err != nil { return nil, err } diff --git a/pkg/buildermgr/pkgwatcher.go b/pkg/buildermgr/pkgwatcher.go index c98b5d8e..c38c7891 100644 --- a/pkg/buildermgr/pkgwatcher.go +++ b/pkg/buildermgr/pkgwatcher.go @@ -90,7 +90,7 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) { pkgw.logger.Info("starting build for package", zap.String("package_name", srcpkg.ObjectMeta.Name), zap.String("resource_version", srcpkg.ObjectMeta.ResourceVersion)) - pkg, err := updatePackage(pkgw.logger, pkgw.fissionClient, srcpkg, fv1.BuildStatusRunning, "", nil) + pkg, err := updatePackage(ctx, pkgw.logger, pkgw.fissionClient, srcpkg, fv1.BuildStatusRunning, "", nil) if err != nil { pkgw.logger.Error("error setting package pending state", zap.Error(err)) return @@ -100,7 +100,7 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) { if k8serrors.IsNotFound(err) { e := "environment does not exist" pkgw.logger.Error(e, zap.String("environment", pkg.Spec.Environment.Name)) - _, er := updatePackage(pkgw.logger, pkgw.fissionClient, pkg, + _, er := updatePackage(ctx, pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, fmt.Sprintf("%s: %q", e, pkg.Spec.Environment.Name), nil) if er != nil { pkgw.logger.Error( @@ -185,7 +185,7 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) { 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)) - _, er := updatePackage(pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil) + _, er := updatePackage(ctx, pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil) if er != nil { pkgw.logger.Error( "error updating package", @@ -205,7 +205,7 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) { e := "error getting function list" pkgw.logger.Error(e, zap.Error(err)) buildLogs += fmt.Sprintf("%s: %v\n", e, err) - _, er := updatePackage(pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil) + _, er := updatePackage(ctx, pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil) if er != nil { pkgw.logger.Error( "error updating package", @@ -229,7 +229,7 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) { e := "error updating function package resource version" pkgw.logger.Error(e, zap.Error(err)) buildLogs += fmt.Sprintf("%s: %v\n", e, err) - _, er := updatePackage(pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil) + _, er := updatePackage(ctx, pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil) if er != nil { pkgw.logger.Error( "error updating package", @@ -243,11 +243,11 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) { } } - _, err = updatePackage(pkgw.logger, pkgw.fissionClient, pkg, + _, err = updatePackage(ctx, pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusSucceeded, buildLogs, uploadResp) if err != nil { pkgw.logger.Error("error updating package info", zap.Error(err), zap.String("package_name", pkg.ObjectMeta.Name)) - _, er := updatePackage(pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil) + _, er := updatePackage(ctx, pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil) if er != nil { pkgw.logger.Error( "error updating package", @@ -265,7 +265,7 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) { time.Sleep(healthCheckBackOff.GetNext()) } // build timeout - _, err = updatePackage(pkgw.logger, pkgw.fissionClient, pkg, + _, err = updatePackage(ctx, pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, "Build timeout due to environment builder not ready", nil) if err != nil { pkgw.logger.Error( @@ -280,12 +280,11 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) { zap.String("package", fmt.Sprintf("%s.%s", pkg.ObjectMeta.Name, pkg.ObjectMeta.Namespace))) } -func (pkgw *packageWatcher) packageInformerHandler() k8sCache.ResourceEventHandlerFuncs { - processPkg := func(pkg *fv1.Package) { +func (pkgw *packageWatcher) packageInformerHandler(ctx context.Context) k8sCache.ResourceEventHandlerFuncs { + processPkg := func(ctx context.Context, pkg *fv1.Package) { var err error - if len(pkg.Status.BuildStatus) == 0 { - _, err = setInitialBuildStatus(pkgw.fissionClient, pkg) + _, err = setInitialBuildStatus(ctx, pkgw.fissionClient, pkg) if err != nil { pkgw.logger.Error("error filling package status", zap.Error(err)) } @@ -296,14 +295,13 @@ func (pkgw *packageWatcher) packageInformerHandler() k8sCache.ResourceEventHandl } // Only build pending state packages. if pkg.Status.BuildStatus == fv1.BuildStatusPending { - ctx := context.Background() go pkgw.build(ctx, pkg) } } return k8sCache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { pkg := obj.(*fv1.Package) - processPkg(pkg) + processPkg(ctx, pkg) }, UpdateFunc: func(oldObj, newObj interface{}) { oldPkg := oldObj.(*fv1.Package) @@ -318,7 +316,7 @@ func (pkgw *packageWatcher) packageInformerHandler() k8sCache.ResourceEventHandl pkg.Status.BuildStatus != fv1.BuildStatusPending { return } - processPkg(pkg) + processPkg(ctx, pkg) }, } } @@ -326,14 +324,14 @@ func (pkgw *packageWatcher) packageInformerHandler() k8sCache.ResourceEventHandl func (pkgw *packageWatcher) Run(ctx context.Context) { go metrics.ServeMetrics(ctx, pkgw.logger) go (*pkgw.podInformer).Run(ctx.Done()) - (*pkgw.pkgInformer).AddEventHandler(pkgw.packageInformerHandler()) + (*pkgw.pkgInformer).AddEventHandler(pkgw.packageInformerHandler(ctx)) (*pkgw.pkgInformer).Run(ctx.Done()) } // setInitialBuildStatus sets initial build status to a package if it is empty. // This normally occurs when the user applies package YAML files that have no status field // through kubectl. -func setInitialBuildStatus(fissionClient versioned.Interface, pkg *fv1.Package) (*fv1.Package, error) { +func setInitialBuildStatus(ctx context.Context, fissionClient versioned.Interface, pkg *fv1.Package) (*fv1.Package, error) { pkg.Status = fv1.PackageStatus{ LastUpdateTimestamp: metav1.Time{Time: time.Now().UTC()}, } @@ -351,5 +349,5 @@ func setInitialBuildStatus(fissionClient versioned.Interface, pkg *fv1.Package) } // TODO: use UpdateStatus to update status - return fissionClient.CoreV1().Packages(pkg.Namespace).Update(context.TODO(), pkg, metav1.UpdateOptions{}) + return fissionClient.CoreV1().Packages(pkg.Namespace).Update(ctx, pkg, metav1.UpdateOptions{}) } diff --git a/pkg/canaryconfigmgr/canaryConfigMgr.go b/pkg/canaryconfigmgr/canaryConfigMgr.go index 3c5289b5..2111fd3a 100644 --- a/pkg/canaryconfigmgr/canaryConfigMgr.go +++ b/pkg/canaryconfigmgr/canaryConfigMgr.go @@ -49,7 +49,7 @@ type canaryConfigMgr struct { canaryCfgCancelFuncMap *canaryConfigCancelFuncMap } -func MakeCanaryConfigMgr(logger *zap.Logger, fissionClient versioned.Interface, kubeClient kubernetes.Interface, prometheusSvc string) (*canaryConfigMgr, error) { +func MakeCanaryConfigMgr(ctx context.Context, logger *zap.Logger, fissionClient versioned.Interface, kubeClient kubernetes.Interface, prometheusSvc string) (*canaryConfigMgr, error) { if prometheusSvc == "" { logger.Info("try to retrieve prometheus server information from environment variables") @@ -96,16 +96,16 @@ func MakeCanaryConfigMgr(logger *zap.Logger, fissionClient versioned.Interface, informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) informer := informerFactory.Core().V1().CanaryConfigs().Informer() configMgr.canaryConfigInformer = &informer - configMgr.CanaryConfigEventHandlers() + configMgr.CanaryConfigEventHandlers(ctx) return configMgr, nil } -func (canaryCfgMgr *canaryConfigMgr) CanaryConfigEventHandlers() { +func (canaryCfgMgr *canaryConfigMgr) CanaryConfigEventHandlers(ctx context.Context) { (*canaryCfgMgr.canaryConfigInformer).AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { canaryConfig := obj.(*fv1.CanaryConfig) if canaryConfig.Status.Status == fv1.CanaryConfigStatusPending { - go canaryCfgMgr.addCanaryConfig(canaryConfig) + go canaryCfgMgr.addCanaryConfig(ctx, canaryConfig) } }, DeleteFunc: func(obj interface{}) { @@ -121,9 +121,9 @@ func (canaryCfgMgr *canaryConfigMgr) CanaryConfigEventHandlers() { zap.String("name", newConfig.ObjectMeta.Name), zap.String("namespace", newConfig.ObjectMeta.Namespace), zap.String("version", newConfig.ObjectMeta.ResourceVersion)) - go canaryCfgMgr.updateCanaryConfig(oldConfig, newConfig) + go canaryCfgMgr.updateCanaryConfig(ctx, oldConfig, newConfig) } - go canaryCfgMgr.reSyncCanaryConfigs() + go canaryCfgMgr.reSyncCanaryConfigs(ctx) }, }) @@ -134,7 +134,7 @@ func (canaryCfgMgr *canaryConfigMgr) Run(ctx context.Context) { canaryCfgMgr.logger.Info("started canary configmgr controller") } -func (canaryCfgMgr *canaryConfigMgr) addCanaryConfig(canaryConfig *fv1.CanaryConfig) { +func (canaryCfgMgr *canaryConfigMgr) addCanaryConfig(ctx context.Context, canaryConfig *fv1.CanaryConfig) { canaryCfgMgr.logger.Debug("addCanaryConfig called", zap.String("canary_config", canaryConfig.ObjectMeta.Name)) // for each canary config, create a ticker with increment interval @@ -152,7 +152,7 @@ func (canaryCfgMgr *canaryConfigMgr) addCanaryConfig(canaryConfig *fv1.CanaryCon // create a context cancel func for each canary config. this will be used to cancel the processing of this canary // config in the event that it's deleted - ctx, cancel := context.WithCancel(context.Background()) + ctx, cancel := context.WithCancel(ctx) cacheValue := &CanaryProcessingInfo{ CancelFunc: &cancel, @@ -319,7 +319,7 @@ func (canaryCfgMgr *canaryConfigMgr) RollForwardOrBack(ctx context.Context, cana zap.String("namespace", canaryConfig.ObjectMeta.Namespace), zap.String("version", canaryConfig.ObjectMeta.ResourceVersion)) ticker.Stop() - err := canaryCfgMgr.rollback(canaryConfig, triggerObj) + err := canaryCfgMgr.rollback(ctx, canaryConfig, triggerObj) if err != nil { canaryCfgMgr.logger.Error("error rolling back canary config", zap.Error(err), @@ -332,7 +332,7 @@ func (canaryCfgMgr *canaryConfigMgr) RollForwardOrBack(ctx context.Context, cana } } - doneProcessingCanaryConfig, err := canaryCfgMgr.rollForward(canaryConfig, triggerObj) + doneProcessingCanaryConfig, err := canaryCfgMgr.rollForward(ctx, canaryConfig, triggerObj) if err != nil { // just log the error and hope that next iteration will succeed canaryCfgMgr.logger.Error("error incrementing weights for trigger", @@ -348,7 +348,7 @@ func (canaryCfgMgr *canaryConfigMgr) RollForwardOrBack(ctx context.Context, cana ticker.Stop() // update the status of canary config as done processing, we don't care if we aren't able to update because // resync takes care of the update - err = canaryCfgMgr.updateCanaryConfigStatusWithRetries(canaryConfig.ObjectMeta.Name, canaryConfig.ObjectMeta.Namespace, + err = canaryCfgMgr.updateCanaryConfigStatusWithRetries(ctx, canaryConfig.ObjectMeta.Name, canaryConfig.ObjectMeta.Namespace, fv1.CanaryConfigStatusSucceeded) if err != nil { // can't do much after max retries other than logging it. @@ -368,9 +368,9 @@ func (canaryCfgMgr *canaryConfigMgr) RollForwardOrBack(ctx context.Context, cana } } -func (canaryCfgMgr *canaryConfigMgr) updateHttpTriggerWithRetries(triggerName, triggerNamespace string, fnWeights map[string]int) (err error) { +func (canaryCfgMgr *canaryConfigMgr) updateHttpTriggerWithRetries(ctx context.Context, triggerName, triggerNamespace string, fnWeights map[string]int) (err error) { for i := 0; i < maxRetries; i++ { - triggerObj, err := canaryCfgMgr.fissionClient.CoreV1().HTTPTriggers(triggerNamespace).Get(context.TODO(), triggerName, metav1.GetOptions{}) + triggerObj, err := canaryCfgMgr.fissionClient.CoreV1().HTTPTriggers(triggerNamespace).Get(ctx, triggerName, metav1.GetOptions{}) if err != nil { e := "error getting http trigger object" canaryCfgMgr.logger.Error(e, zap.Error(err), zap.String("trigger_name", triggerName), zap.String("trigger_namespace", triggerNamespace)) @@ -379,7 +379,7 @@ func (canaryCfgMgr *canaryConfigMgr) updateHttpTriggerWithRetries(triggerName, t triggerObj.Spec.FunctionReference.FunctionWeights = fnWeights - _, err = canaryCfgMgr.fissionClient.CoreV1().HTTPTriggers(triggerNamespace).Update(context.TODO(), triggerObj, metav1.UpdateOptions{}) + _, err = canaryCfgMgr.fissionClient.CoreV1().HTTPTriggers(triggerNamespace).Update(ctx, triggerObj, metav1.UpdateOptions{}) switch { case err == nil: canaryCfgMgr.logger.Debug("updated http trigger", zap.String("trigger_name", triggerName), zap.String("trigger_namespace", triggerNamespace)) @@ -403,9 +403,9 @@ func (canaryCfgMgr *canaryConfigMgr) updateHttpTriggerWithRetries(triggerName, t return err } -func (canaryCfgMgr *canaryConfigMgr) updateCanaryConfigStatusWithRetries(cfgName, cfgNamespace string, status string) (err error) { +func (canaryCfgMgr *canaryConfigMgr) updateCanaryConfigStatusWithRetries(ctx context.Context, cfgName, cfgNamespace string, status string) (err error) { for i := 0; i < maxRetries; i++ { - canaryCfgObj, err := canaryCfgMgr.fissionClient.CoreV1().CanaryConfigs(cfgNamespace).Get(context.TODO(), cfgName, metav1.GetOptions{}) + canaryCfgObj, err := canaryCfgMgr.fissionClient.CoreV1().CanaryConfigs(cfgNamespace).Get(ctx, cfgName, metav1.GetOptions{}) if err != nil { e := "error getting http canary config object" canaryCfgMgr.logger.Error(e, @@ -423,7 +423,7 @@ func (canaryCfgMgr *canaryConfigMgr) updateCanaryConfigStatusWithRetries(cfgName canaryCfgObj.Status.Status = status - _, err = canaryCfgMgr.fissionClient.CoreV1().CanaryConfigs(cfgNamespace).Update(context.TODO(), canaryCfgObj, metav1.UpdateOptions{}) + _, err = canaryCfgMgr.fissionClient.CoreV1().CanaryConfigs(cfgNamespace).Update(ctx, canaryCfgObj, metav1.UpdateOptions{}) switch { case err == nil: canaryCfgMgr.logger.Info("updated canary config", @@ -449,23 +449,23 @@ func (canaryCfgMgr *canaryConfigMgr) updateCanaryConfigStatusWithRetries(cfgName return err } -func (canaryCfgMgr *canaryConfigMgr) rollback(canaryConfig *fv1.CanaryConfig, trigger *fv1.HTTPTrigger) error { +func (canaryCfgMgr *canaryConfigMgr) rollback(ctx context.Context, canaryConfig *fv1.CanaryConfig, trigger *fv1.HTTPTrigger) error { functionWeights := trigger.Spec.FunctionReference.FunctionWeights functionWeights[canaryConfig.Spec.NewFunction] = 0 functionWeights[canaryConfig.Spec.OldFunction] = 100 - err := canaryCfgMgr.updateHttpTriggerWithRetries(trigger.ObjectMeta.Name, trigger.ObjectMeta.Namespace, functionWeights) + err := canaryCfgMgr.updateHttpTriggerWithRetries(ctx, trigger.ObjectMeta.Name, trigger.ObjectMeta.Namespace, functionWeights) if err != nil { return err } - err = canaryCfgMgr.updateCanaryConfigStatusWithRetries(canaryConfig.ObjectMeta.Name, canaryConfig.ObjectMeta.Namespace, + err = canaryCfgMgr.updateCanaryConfigStatusWithRetries(ctx, canaryConfig.ObjectMeta.Name, canaryConfig.ObjectMeta.Namespace, fv1.CanaryConfigStatusFailed) return err } -func (canaryCfgMgr *canaryConfigMgr) rollForward(canaryConfig *fv1.CanaryConfig, trigger *fv1.HTTPTrigger) (bool, error) { +func (canaryCfgMgr *canaryConfigMgr) rollForward(ctx context.Context, canaryConfig *fv1.CanaryConfig, trigger *fv1.HTTPTrigger) (bool, error) { doneProcessingCanaryConfig := false functionWeights := trigger.Spec.FunctionReference.FunctionWeights @@ -487,11 +487,11 @@ func (canaryCfgMgr *canaryConfigMgr) rollForward(canaryConfig *fv1.CanaryConfig, zap.String("namespace", canaryConfig.ObjectMeta.Namespace), zap.Any("function_weights", functionWeights)) - err := canaryCfgMgr.updateHttpTriggerWithRetries(trigger.ObjectMeta.Name, trigger.ObjectMeta.Namespace, functionWeights) + err := canaryCfgMgr.updateHttpTriggerWithRetries(ctx, trigger.ObjectMeta.Name, trigger.ObjectMeta.Namespace, functionWeights) return doneProcessingCanaryConfig, err } -func (canaryCfgMgr *canaryConfigMgr) reSyncCanaryConfigs() { +func (canaryCfgMgr *canaryConfigMgr) reSyncCanaryConfigs(ctx context.Context) { for _, obj := range (*canaryCfgMgr.canaryConfigInformer).GetStore().List() { canaryConfig := obj.(*fv1.CanaryConfig) _, err := canaryCfgMgr.canaryCfgCancelFuncMap.lookup(&canaryConfig.ObjectMeta) @@ -502,7 +502,7 @@ func (canaryCfgMgr *canaryConfigMgr) reSyncCanaryConfigs() { zap.String("version", canaryConfig.ObjectMeta.ResourceVersion)) // new canaryConfig detected, add it to our cache and start processing it - go canaryCfgMgr.addCanaryConfig(canaryConfig) + go canaryCfgMgr.addCanaryConfig(ctx, canaryConfig) } } } @@ -527,7 +527,7 @@ func (canaryCfgMgr *canaryConfigMgr) deleteCanaryConfig(canaryConfig *fv1.Canary (*canaryProcessingInfo.CancelFunc)() } -func (canaryCfgMgr *canaryConfigMgr) updateCanaryConfig(oldCanaryConfig *fv1.CanaryConfig, newCanaryConfig *fv1.CanaryConfig) { +func (canaryCfgMgr *canaryConfigMgr) updateCanaryConfig(ctx context.Context, oldCanaryConfig *fv1.CanaryConfig, newCanaryConfig *fv1.CanaryConfig) { // before removing the object from cache, we need to get it's cancel func and cancel it canaryCfgMgr.deleteCanaryConfig(oldCanaryConfig) @@ -540,7 +540,7 @@ func (canaryCfgMgr *canaryConfigMgr) updateCanaryConfig(oldCanaryConfig *fv1.Can zap.String("version", oldCanaryConfig.ObjectMeta.ResourceVersion)) return } - canaryCfgMgr.addCanaryConfig(newCanaryConfig) + canaryCfgMgr.addCanaryConfig(ctx, newCanaryConfig) } func getEnvValue(envVar string) string { diff --git a/pkg/controller/api_test.go b/pkg/controller/api_test.go index e9627184..2aeefe20 100644 --- a/pkg/controller/api_test.go +++ b/pkg/controller/api_test.go @@ -519,7 +519,7 @@ func TestMain(m *testing.M) { exitVal := m.Run() logger.Info("Deleting test namespace", zap.String("namespace", testNS)) gracePeriod := int64(0) - err = kubeClient.CoreV1().Namespaces().Delete(context.TODO(), testNS, metav1.DeleteOptions{GracePeriodSeconds: &gracePeriod}) + err = kubeClient.CoreV1().Namespaces().Delete(ctx, testNS, metav1.DeleteOptions{GracePeriodSeconds: &gracePeriod}) if err != nil { logger.Error("error deleting test namespace", zap.String("namespace", testNS), zap.Error(err)) } diff --git a/pkg/controller/config.go b/pkg/controller/config.go index 04756911..02ea8bff 100644 --- a/pkg/controller/config.go +++ b/pkg/controller/config.go @@ -28,15 +28,15 @@ import ( "github.com/fission/fission/pkg/generated/clientset/versioned" ) -func ConfigCanaryFeature(context context.Context, logger *zap.Logger, fissionClient versioned.Interface, kubeClient kubernetes.Interface, featureConfig *config.FeatureConfig, featureStatus map[string]string) error { +func ConfigCanaryFeature(ctx context.Context, logger *zap.Logger, fissionClient versioned.Interface, kubeClient kubernetes.Interface, featureConfig *config.FeatureConfig, featureStatus map[string]string) error { // start the appropriate controller if featureConfig.CanaryConfig.IsEnabled { - canaryCfgMgr, err := canaryconfigmgr.MakeCanaryConfigMgr(logger, fissionClient, kubeClient, featureConfig.CanaryConfig.PrometheusSvc) + canaryCfgMgr, err := canaryconfigmgr.MakeCanaryConfigMgr(ctx, logger, fissionClient, kubeClient, featureConfig.CanaryConfig.PrometheusSvc) if err != nil { featureStatus[config.CanaryFeature] = err.Error() return errors.Wrap(err, "failed to start canary config manager") } - canaryCfgMgr.Run(context) + canaryCfgMgr.Run(ctx) logger.Info("started canary config manager") } @@ -44,7 +44,7 @@ func ConfigCanaryFeature(context context.Context, logger *zap.Logger, fissionCli } // ConfigureFeatures gets the feature config and configures the features that are enabled -func ConfigureFeatures(context context.Context, logger *zap.Logger, unitTestMode bool, fissionClient versioned.Interface, kubeClient kubernetes.Interface) (map[string]string, error) { +func ConfigureFeatures(ctx context.Context, logger *zap.Logger, unitTestMode bool, fissionClient versioned.Interface, kubeClient kubernetes.Interface) (map[string]string, error) { // set feature enabled to false if unitTestMode if unitTestMode { return nil, nil @@ -61,6 +61,6 @@ func ConfigureFeatures(context context.Context, logger *zap.Logger, unitTestMode // configure respective features // in the future when new optional features are added, we need to add corresponding feature handlers and invoke them here - err = ConfigCanaryFeature(context, logger, fissionClient, kubeClient, featureConfig, featureStatus) + err = ConfigCanaryFeature(ctx, logger, fissionClient, kubeClient, featureConfig, featureStatus) return featureStatus, err } diff --git a/pkg/executor/client/client.go b/pkg/executor/client/client.go index 5a96d3d1..99d66f47 100644 --- a/pkg/executor/client/client.go +++ b/pkg/executor/client/client.go @@ -153,7 +153,7 @@ func (c *Client) service() { svcReqs = append(svcReqs, req) } c.logger.Debug("tapped services in batch", zap.Int("service_count", len(urls))) - err := c._tapService(context.TODO(), svcReqs) + err := c._tapService(context.Background(), svcReqs) if err != nil { c.logger.Error("error tapping function service address", zap.Error(err)) } diff --git a/pkg/executor/executor.go b/pkg/executor/executor.go index d17bd0df..59f02300 100644 --- a/pkg/executor/executor.go +++ b/pkg/executor/executor.go @@ -295,7 +295,7 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace st } gpmPodInformer := gpmInformerFactory.Core().V1().Pods() gpmRsInformer := gpmInformerFactory.Apps().V1().ReplicaSets() - gpm, err := poolmgr.MakeGenericPoolManager( + gpm, err := poolmgr.MakeGenericPoolManager(ctx, logger, fissionClient, kubernetesClient, metricsClient, functionNamespace, fetcherConfig, executorInstanceID, @@ -311,7 +311,7 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace st } ndmDeplInformer := ndmInformerFactory.Apps().V1().Deployments() ndmSvcInformer := ndmInformerFactory.Core().V1().Services() - ndm, err := newdeploy.MakeNewDeploy( + ndm, err := newdeploy.MakeNewDeploy(ctx, logger, fissionClient, kubernetesClient, functionNamespace, fetcherConfig, executorInstanceID, diff --git a/pkg/executor/executor_test.go b/pkg/executor/executor_test.go index cd74296e..7a3a43ac 100644 --- a/pkg/executor/executor_test.go +++ b/pkg/executor/executor_test.go @@ -51,8 +51,8 @@ func panicIf(err error) { } // return the number of pods in the given namespace matching the given labels -func countPods(kubeClient kubernetes.Interface, ns string, labelz map[string]string) int { - pods, err := kubeClient.CoreV1().Pods(ns).List(context.TODO(), metav1.ListOptions{ +func countPods(ctx context.Context, kubeClient kubernetes.Interface, ns string, labelz map[string]string) int { + pods, err := kubeClient.CoreV1().Pods(ns).List(ctx, metav1.ListOptions{ LabelSelector: labels.Set(labelz).AsSelector().String(), }) if err != nil { @@ -61,8 +61,8 @@ func countPods(kubeClient kubernetes.Interface, ns string, labelz map[string]str return len(pods.Items) } -func createTestNamespace(kubeClient kubernetes.Interface, ns string) { - _, err := kubeClient.CoreV1().Namespaces().Create(context.TODO(), &apiv1.Namespace{ +func createTestNamespace(ctx context.Context, kubeClient kubernetes.Interface, ns string) { + _, err := kubeClient.CoreV1().Namespaces().Create(ctx, &apiv1.Namespace{ ObjectMeta: metav1.ObjectMeta{ Name: ns, }, @@ -74,8 +74,8 @@ func createTestNamespace(kubeClient kubernetes.Interface, ns string) { } // create a nodeport service -func createSvc(kubeClient kubernetes.Interface, ns string, name string, targetPort int, nodePort int32, labels map[string]string) *apiv1.Service { - svc, err := kubeClient.CoreV1().Services(ns).Create(context.TODO(), &apiv1.Service{ +func createSvc(ctx context.Context, kubeClient kubernetes.Interface, ns string, name string, targetPort int, nodePort int32, labels map[string]string) *apiv1.Service { + svc, err := kubeClient.CoreV1().Services(ns).Create(ctx, &apiv1.Service{ ObjectMeta: metav1.ObjectMeta{ Name: name, }, @@ -120,18 +120,19 @@ func TestExecutor(t *testing.T) { log.Panicf("failed to connect: %v", err) } + ctx := context.Background() // create the test's namespaces - createTestNamespace(kubeClient, fissionNs) + createTestNamespace(ctx, kubeClient, fissionNs) defer func() { - err := kubeClient.CoreV1().Namespaces().Delete(context.TODO(), fissionNs, metav1.DeleteOptions{}) + err := kubeClient.CoreV1().Namespaces().Delete(ctx, fissionNs, metav1.DeleteOptions{}) if err != nil { log.Fatalf("failed to delete namespace: %v", err) } }() - createTestNamespace(kubeClient, functionNs) + createTestNamespace(ctx, kubeClient, functionNs) defer func() { - err := kubeClient.CoreV1().Namespaces().Delete(context.TODO(), functionNs, metav1.DeleteOptions{}) + err := kubeClient.CoreV1().Namespaces().Delete(ctx, functionNs, metav1.DeleteOptions{}) if err != nil { log.Fatalf("failed to delete namespace: %v", err) } @@ -143,18 +144,18 @@ func TestExecutor(t *testing.T) { panicIf(err) // make sure CRD types exist on cluster - err = crd.EnsureFissionCRDs(context.TODO(), logger, apiExtClient) + err = crd.EnsureFissionCRDs(ctx, logger, apiExtClient) if err != nil { log.Panicf("failed to ensure crds: %v", err) } - err = crd.WaitForCRDs(context.TODO(), fissionClient) + err = crd.WaitForCRDs(ctx, fissionClient) if err != nil { log.Panicf("failed to wait crds: %v", err) } // create an env on the cluster - env, err := fissionClient.CoreV1().Environments(fissionNs).Create(context.TODO(), &fv1.Environment{ + env, err := fissionClient.CoreV1().Environments(fissionNs).Create(ctx, &fv1.Environment{ ObjectMeta: metav1.ObjectMeta{ Name: "nodejs", Namespace: fissionNs, @@ -173,7 +174,6 @@ func TestExecutor(t *testing.T) { // create poolmgr port := 9999 - ctx := context.Background() err = StartExecutor(ctx, logger, functionNs, "fission-builder", port) if err != nil { log.Panicf("failed to start poolmgr: %v", err) @@ -208,7 +208,7 @@ func TestExecutor(t *testing.T) { Deployment: deployment, }, } - p, err = fissionClient.CoreV1().Packages(fissionNs).Create(context.TODO(), p, metav1.CreateOptions{}) + p, err = fissionClient.CoreV1().Packages(fissionNs).Create(ctx, p, metav1.CreateOptions{}) if err != nil { log.Panicf("failed to create package: %v", err) } @@ -230,7 +230,7 @@ func TestExecutor(t *testing.T) { }, }, } - _, err = fissionClient.CoreV1().Functions(fissionNs).Create(context.TODO(), f, metav1.CreateOptions{}) + _, err = fissionClient.CoreV1().Functions(fissionNs).Create(ctx, f, metav1.CreateOptions{}) if err != nil { log.Panicf("failed to create function: %v", err) } @@ -238,18 +238,18 @@ func TestExecutor(t *testing.T) { // create a service to call fetcher and the env container labels := map[string]string{"functionName": f.ObjectMeta.Name} var fetcherPort int32 = 30001 - fetcherSvc := createSvc(kubeClient, functionNs, fmt.Sprintf("%v-%v", f.ObjectMeta.Name, "fetcher"), 8000, fetcherPort, labels) + fetcherSvc := createSvc(ctx, kubeClient, functionNs, fmt.Sprintf("%v-%v", f.ObjectMeta.Name, "fetcher"), 8000, fetcherPort, labels) defer func() { - err := kubeClient.CoreV1().Services(functionNs).Delete(context.TODO(), fetcherSvc.ObjectMeta.Name, metav1.DeleteOptions{}) + err := kubeClient.CoreV1().Services(functionNs).Delete(ctx, fetcherSvc.ObjectMeta.Name, metav1.DeleteOptions{}) if err != nil { log.Fatalf("failed to delete service: %v", err) } }() var funcSvcPort int32 = 30002 - functionSvc := createSvc(kubeClient, functionNs, f.ObjectMeta.Name, 8888, funcSvcPort, labels) + functionSvc := createSvc(ctx, kubeClient, functionNs, f.ObjectMeta.Name, 8888, funcSvcPort, labels) defer func() { - err := kubeClient.CoreV1().Services(functionNs).Delete(context.TODO(), functionSvc.ObjectMeta.Name, metav1.DeleteOptions{}) + err := kubeClient.CoreV1().Services(functionNs).Delete(ctx, functionSvc.ObjectMeta.Name, metav1.DeleteOptions{}) if err != nil { log.Fatalf("failed to delete service: %v", err) } @@ -257,14 +257,14 @@ func TestExecutor(t *testing.T) { // the main test: get a service for a given function t1 := time.Now() - svc, err := poolmgrClient.GetServiceForFunction(context.TODO(), f) + svc, err := poolmgrClient.GetServiceForFunction(ctx, f) if err != nil { log.Panicf("failed to get func svc: %v", err) } log.Printf("svc for function created at: %v (in %v)", svc, time.Since(t1)) // ensure that a pod with the label functionName=f.ObjectMeta.Name exists - podCount := countPods(kubeClient, functionNs, map[string]string{"functionName": f.ObjectMeta.Name}) + podCount := countPods(ctx, kubeClient, functionNs, map[string]string{"functionName": f.ObjectMeta.Name}) if podCount != 1 { log.Panicf("expected 1 function pod, found %v", podCount) } diff --git a/pkg/executor/executortype/newdeploy/envhandlers.go b/pkg/executor/executortype/newdeploy/envhandlers.go index ddd92420..00cb65be 100644 --- a/pkg/executor/executortype/newdeploy/envhandlers.go +++ b/pkg/executor/executortype/newdeploy/envhandlers.go @@ -25,14 +25,13 @@ import ( fv1 "github.com/fission/fission/pkg/apis/core/v1" ) -func (deploy *NewDeploy) EnvEventHandlers() k8sCache.ResourceEventHandlerFuncs { +func (deploy *NewDeploy) EnvEventHandlers(ctx context.Context) k8sCache.ResourceEventHandlerFuncs { return k8sCache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) {}, DeleteFunc: func(obj interface{}) {}, 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)) diff --git a/pkg/executor/executortype/newdeploy/funchandlers.go b/pkg/executor/executortype/newdeploy/funchandlers.go index bffaec01..f70d461f 100644 --- a/pkg/executor/executortype/newdeploy/funchandlers.go +++ b/pkg/executor/executortype/newdeploy/funchandlers.go @@ -24,14 +24,13 @@ import ( fv1 "github.com/fission/fission/pkg/apis/core/v1" ) -func (deploy *NewDeploy) FunctionEventHandlers() k8sCache.ResourceEventHandlerFuncs { +func (deploy *NewDeploy) FunctionEventHandlers(ctx context.Context) k8sCache.ResourceEventHandlerFuncs { return k8sCache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { // TODO: A workaround to process items in parallel. We should use workqueue ("k8s.io/client-go/util/workqueue") // 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(ctx, fn) @@ -46,7 +45,6 @@ func (deploy *NewDeploy) FunctionEventHandlers() k8sCache.ResourceEventHandlerFu DeleteFunc: func(obj interface{}) { fn := obj.(*fv1.Function) go func() { - ctx := context.Background() err := deploy.deleteFunction(ctx, fn) if err != nil { deploy.logger.Error("error deleting function", @@ -59,7 +57,6 @@ func (deploy *NewDeploy) FunctionEventHandlers() k8sCache.ResourceEventHandlerFu oldFn := oldObj.(*fv1.Function) newFn := newObj.(*fv1.Function) go func() { - ctx := context.Background() err := deploy.updateFunction(ctx, oldFn, newFn) if err != nil { deploy.logger.Error("error updating function", diff --git a/pkg/executor/executortype/newdeploy/newdeploy.go b/pkg/executor/executortype/newdeploy/newdeploy.go index d2c42ce4..069de80c 100644 --- a/pkg/executor/executortype/newdeploy/newdeploy.go +++ b/pkg/executor/executortype/newdeploy/newdeploy.go @@ -120,7 +120,7 @@ func (deploy *NewDeploy) createOrGetDeployment(ctx context.Context, fn *fv1.Func 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) + err := deploy.fetcherConfig.SetupServiceAccount(ctx, deploy.kubernetesClient, deployNamespace, fn.ObjectMeta) if err != nil { deploy.logger.Error("error creating fission fetcher service account for function", zap.Error(err), diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr.go b/pkg/executor/executortype/newdeploy/newdeploymgr.go index 4abdc3c8..3d3ba360 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr.go @@ -96,6 +96,7 @@ type ( // MakeNewDeploy initializes and returns an instance of NewDeploy. func MakeNewDeploy( + ctx context.Context, logger *zap.Logger, fissionClient versioned.Interface, kubernetesClient kubernetes.Interface, @@ -145,8 +146,8 @@ func MakeNewDeploy( nd.svcLister = svcInformer.Lister() nd.svcListerSynced = svcInformer.Informer().HasSynced - funcInformer.Informer().AddEventHandler(nd.FunctionEventHandlers()) - envInformer.Informer().AddEventHandler(nd.EnvEventHandlers()) + funcInformer.Informer().AddEventHandler(nd.FunctionEventHandlers(ctx)) + envInformer.Informer().AddEventHandler(nd.EnvEventHandlers(ctx)) return nd, nil } @@ -402,8 +403,7 @@ func (deploy *NewDeploy) deleteFunction(ctx context.Context, fn *fv1.Function) e } func (deploy *NewDeploy) fnCreate(ctx context.Context, fn *fv1.Function) (*fscache.FuncSvc, error) { - cleanupFunc := func(ns string, name string) { - ctx := context.Background() + cleanupFunc := func(ctx context.Context, ns string, name string) { err := deploy.cleanupNewdeploy(ctx, ns, name) if err != nil { deploy.logger.Error("received error while cleaning function resources", @@ -436,7 +436,7 @@ func (deploy *NewDeploy) fnCreate(ctx context.Context, fn *fv1.Function) (*fscac 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) + go cleanupFunc(context.Background(), ns, objName) return nil, errors.Wrapf(err, "error creating service %v", objName) } svcAddress := fmt.Sprintf("%v.%v", svc.Name, svc.Namespace) @@ -444,14 +444,14 @@ func (deploy *NewDeploy) fnCreate(ctx context.Context, fn *fv1.Function) (*fscac 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) + go cleanupFunc(context.Background(), ns, objName) return nil, errors.Wrapf(err, "error creating deployment %v", objName) } hpa, err := deploy.hpaops.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) + go cleanupFunc(context.Background(), ns, objName) return nil, errors.Wrapf(err, "error creating the HPA %v", objName) } diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr_test.go b/pkg/executor/executortype/newdeploy/newdeploymgr_test.go index 06673067..13e0bf3d 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr_test.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr_test.go @@ -76,7 +76,7 @@ func TestRefreshFuncPods(t *testing.T) { t.Fatalf("Error creating fetcher config: %s", err) } - executor, err := MakeNewDeploy(logger, fissionClient, kubernetesClient, functionNamespace, fetcherConfig, "test", + executor, err := MakeNewDeploy(ctx, logger, fissionClient, kubernetesClient, functionNamespace, fetcherConfig, "test", funcInformer, envInformer, deployInformer, svcInformer, podSpecPatch) if err != nil { t.Fatalf("new deploy manager creation failed: %s", err) diff --git a/pkg/executor/executortype/poolmgr/funchandlers.go b/pkg/executor/executortype/poolmgr/funchandlers.go index f276af86..1eb2e2b4 100644 --- a/pkg/executor/executortype/poolmgr/funchandlers.go +++ b/pkg/executor/executortype/poolmgr/funchandlers.go @@ -41,10 +41,9 @@ func getIstioServiceLabels(fnName string) map[string]string { // Based on function create/update/delete event, we create role binding // for the secret/configmap access which is used by fetcher component. // If istio is enabled, we create a service for the function. -func FunctionEventHandlers(logger *zap.Logger, kubernetesClient kubernetes.Interface, fissionfnNamespace string, istioEnabled bool) k8sCache.ResourceEventHandlerFuncs { +func FunctionEventHandlers(ctx context.Context, logger *zap.Logger, kubernetesClient kubernetes.Interface, 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, @@ -133,7 +132,6 @@ func FunctionEventHandlers(logger *zap.Logger, kubernetesClient kubernetes.Inter }, DeleteFunc: func(obj interface{}) { - ctx := context.Background() fn := obj.(*fv1.Function) fnExecutorType := fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType @@ -183,7 +181,6 @@ func FunctionEventHandlers(logger *zap.Logger, kubernetesClient kubernetes.Inter if newFunc.Spec.Environment.Namespace != metav1.NamespaceDefault { envNs = newFunc.Spec.Environment.Namespace } - ctx := context.Background() err := utils.SetupRoleBinding(ctx, logger, kubernetesClient, fv1.SecretConfigMapGetterRB, newFunc.ObjectMeta.Namespace, utils.GetSecretConfigMapGetterCR(), fv1.ClusterRole, fv1.FissionFetcherSA, envNs) diff --git a/pkg/executor/executortype/poolmgr/gp.go b/pkg/executor/executortype/poolmgr/gp.go index 8de1f820..2e63f51a 100644 --- a/pkg/executor/executortype/poolmgr/gp.go +++ b/pkg/executor/executortype/poolmgr/gp.go @@ -142,7 +142,7 @@ func MakeGenericPool( func (gp *GenericPool) setup(ctx context.Context) error { // create fetcher SA in this ns, if not already created - err := gp.fetcherConfig.SetupServiceAccount(gp.kubernetesClient, gp.namespace, nil) + err := gp.fetcherConfig.SetupServiceAccount(ctx, gp.kubernetesClient, gp.namespace, nil) if err != nil { return errors.Wrapf(err, "error creating fetcher service account in namespace %q", gp.namespace) } @@ -156,7 +156,7 @@ func (gp *GenericPool) setup(ctx context.Context) error { if err != nil { return err } - go gp.updateCPUUtilizationSvc() + go gp.updateCPUUtilizationSvc(ctx) return nil } @@ -185,7 +185,7 @@ func (gp *GenericPool) checkMetricsApi() bool { return utils.SupportedMetricsAPIVersionAvailable(apiGroups) } -func (gp *GenericPool) updateCPUUtilizationSvc() { +func (gp *GenericPool) updateCPUUtilizationSvc(ctx context.Context) { var metricsApiAvailabe bool checkDuration := 30 @@ -194,8 +194,8 @@ func (gp *GenericPool) updateCPUUtilizationSvc() { gp.logger.Warn("Metrics API not available") } - serviceFunc := func() { - podMetricsList, err := gp.metricsClient.MetricsV1beta1().PodMetricses(gp.namespace).List(context.TODO(), metav1.ListOptions{ + serviceFunc := func(ctx context.Context) { + podMetricsList, err := gp.metricsClient.MetricsV1beta1().PodMetricses(gp.namespace).List(ctx, metav1.ListOptions{ LabelSelector: "managed=false", }) if err != nil { @@ -220,7 +220,7 @@ func (gp *GenericPool) updateCPUUtilizationSvc() { for { if metricsApiAvailabe { - serviceFunc() + serviceFunc(ctx) } else { if gp.checkMetricsApi() { metricsApiAvailabe = true @@ -342,22 +342,20 @@ func (gp *GenericPool) labelsForFunction(metadata *metav1.ObjectMeta) map[string return label } -func (gp *GenericPool) scheduleDeletePod(name string) { - go func() { - // The sleep allows debugging or collecting logs from the pod before it's - // cleaned up. (We need a better solutions for both those things; log - // aggregation and storage will help.) - gp.logger.Error("error in pod - scheduling cleanup", zap.String("pod", name)) - err := gp.kubernetesClient.CoreV1().Pods(gp.namespace).Delete(context.TODO(), name, metav1.DeleteOptions{}) - if err != nil { - gp.logger.Error( - "error deleting pod", - zap.String("name", name), - zap.String("namespace", gp.namespace), - zap.Error(err), - ) - } - }() +func (gp *GenericPool) scheduleDeletePod(ctx context.Context, name string) { + // The sleep allows debugging or collecting logs from the pod before it's + // cleaned up. (We need a better solutions for both those things; log + // aggregation and storage will help.) + gp.logger.Error("error in pod - scheduling cleanup", zap.String("pod", name)) + err := gp.kubernetesClient.CoreV1().Pods(gp.namespace).Delete(ctx, name, metav1.DeleteOptions{}) + if err != nil { + gp.logger.Error( + "error deleting pod", + zap.String("name", name), + zap.String("namespace", gp.namespace), + zap.Error(err), + ) + } } // IsIPv6 validates if the podIP follows to IPv6 protocol @@ -502,7 +500,7 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac gp.readyPodQueue.Done(key) err = gp.specializePod(ctx, pod, fn) if err != nil { - gp.scheduleDeletePod(pod.ObjectMeta.Name) + go gp.scheduleDeletePod(context.Background(), pod.ObjectMeta.Name) return nil, err } logger.Info("specialized pod", zap.String("pod", pod.ObjectMeta.Name), zap.String("podNamespace", pod.ObjectMeta.Namespace), zap.String("podIP", pod.Status.PodIP)) @@ -516,11 +514,11 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac svc, err := gp.createSvc(ctx, svcName, funcLabels) if err != nil { - gp.scheduleDeletePod(pod.ObjectMeta.Name) + go gp.scheduleDeletePod(context.Background(), pod.ObjectMeta.Name) return nil, err } if svc.ObjectMeta.Name != svcName { - gp.scheduleDeletePod(pod.ObjectMeta.Name) + go gp.scheduleDeletePod(context.Background(), pod.ObjectMeta.Name) return nil, errors.Errorf("sanity check failed for svc %v", svc.ObjectMeta.Name) } diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index db637baa..ec3b1238 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -111,7 +111,7 @@ type ( } ) -func MakeGenericPoolManager( +func MakeGenericPoolManager(ctx context.Context, logger *zap.Logger, fissionClient versioned.Interface, kubernetesClient kubernetes.Interface, @@ -138,7 +138,7 @@ func MakeGenericPoolManager( enableIstio = istio } - poolPodC := NewPoolPodController(gpmLogger, kubernetesClient, functionNamespace, + poolPodC := NewPoolPodController(ctx, gpmLogger, kubernetesClient, functionNamespace, enableIstio, funcInformer, pkgInformer, envInformer, rsInformer, podInformer) gpm := &GenericPoolManager{ @@ -171,10 +171,10 @@ func (gpm *GenericPoolManager) Run(ctx context.Context) { } go gpm.service() gpm.poolPodC.InjectGpm(gpm) - go gpm.WebsocketStartEventChecker(gpm.kubernetesClient) - go gpm.NoActiveConnectionEventChecker(gpm.kubernetesClient) + go gpm.WebsocketStartEventChecker(ctx, gpm.kubernetesClient) + go gpm.NoActiveConnectionEventChecker(ctx, gpm.kubernetesClient) go gpm.idleObjectReaper(ctx) - go gpm.poolPodC.Run(ctx.Done()) + go gpm.poolPodC.Run(ctx, ctx.Done()) } func (gpm *GenericPoolManager) GetTypeName(ctx context.Context) fv1.ExecutorType { @@ -664,17 +664,17 @@ func (gpm *GenericPoolManager) doIdleObjectReaper(ctx context.Context) { } // WebsocketStartEventChecker checks if the pod has emitted a websocket connection start event -func (gpm *GenericPoolManager) WebsocketStartEventChecker(kubeClient kubernetes.Interface) { +func (gpm *GenericPoolManager) WebsocketStartEventChecker(ctx context.Context, kubeClient kubernetes.Interface) { informer := k8sCache.NewSharedInformer( &k8sCache.ListWatch{ ListFunc: func(options metav1.ListOptions) (runtime.Object, error) { options.FieldSelector = "involvedObject.kind=Pod,type=Normal,reason=WsConnectionStarted" - return kubeClient.CoreV1().Events(apiv1.NamespaceAll).List(context.TODO(), options) + return kubeClient.CoreV1().Events(apiv1.NamespaceAll).List(ctx, options) }, WatchFunc: func(options metav1.ListOptions) (watch.Interface, error) { options.FieldSelector = "involvedObject.kind=Pod,type=Normal,reason=WsConnectionStarted" - return kubeClient.CoreV1().Events(apiv1.NamespaceAll).Watch(context.TODO(), options) + return kubeClient.CoreV1().Events(apiv1.NamespaceAll).Watch(ctx, options) }, }, &apiv1.Event{}, @@ -705,17 +705,17 @@ func (gpm *GenericPoolManager) WebsocketStartEventChecker(kubeClient kubernetes. } // NoActiveConnectionEventChecker checks if the pod has emitted an inactive event -func (gpm *GenericPoolManager) NoActiveConnectionEventChecker(kubeClient kubernetes.Interface) { +func (gpm *GenericPoolManager) NoActiveConnectionEventChecker(ctx context.Context, kubeClient kubernetes.Interface) { informer := k8sCache.NewSharedInformer( &k8sCache.ListWatch{ ListFunc: func(options metav1.ListOptions) (runtime.Object, error) { options.FieldSelector = "involvedObject.kind=Pod,type=Normal,reason=NoActiveConnections" - return kubeClient.CoreV1().Events(apiv1.NamespaceAll).List(context.TODO(), options) + return kubeClient.CoreV1().Events(apiv1.NamespaceAll).List(ctx, options) }, WatchFunc: func(options metav1.ListOptions) (watch.Interface, error) { options.FieldSelector = "involvedObject.kind=Pod,type=Normal,reason=NoActiveConnections" - return kubeClient.CoreV1().Events(apiv1.NamespaceAll).Watch(context.TODO(), options) + return kubeClient.CoreV1().Events(apiv1.NamespaceAll).Watch(ctx, options) }, }, &apiv1.Event{}, @@ -737,7 +737,6 @@ func (gpm *GenericPoolManager) NoActiveConnectionEventChecker(kubeClient kuberne gpm.logger.Error("could not convert value from PodToFsvc") return } - ctx := context.Background() gpm.fsCache.DeleteFunctionSvc(ctx, fsvc) for i := range fsvc.KubernetesObjects { gpm.logger.Info("release idle function resources due to inactivity", diff --git a/pkg/executor/executortype/poolmgr/packagehandlers.go b/pkg/executor/executortype/poolmgr/packagehandlers.go index 3639c888..c09fce5e 100644 --- a/pkg/executor/executortype/poolmgr/packagehandlers.go +++ b/pkg/executor/executortype/poolmgr/packagehandlers.go @@ -31,7 +31,7 @@ import ( // PackageEventHandlers provides handlers for package events. // Based on package create/update event, we create role binding // for the package which is used by fetcher component. -func PackageEventHandlers(logger *zap.Logger, kubernetesClient kubernetes.Interface, fissionfnNamespace string) k8sCache.ResourceEventHandlerFuncs { +func PackageEventHandlers(ctx context.Context, logger *zap.Logger, kubernetesClient kubernetes.Interface, fissionfnNamespace string) k8sCache.ResourceEventHandlerFuncs { return k8sCache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { pkg := obj.(*fv1.Package) @@ -44,7 +44,6 @@ func PackageEventHandlers(logger *zap.Logger, kubernetesClient kubernetes.Interf 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(ctx, logger, kubernetesClient, fv1.PackageGetterRB, pkg.ObjectMeta.Namespace, utils.GetPackageGetterCR(), fv1.ClusterRole, fv1.FissionFetcherSA, envNs) @@ -81,7 +80,6 @@ func PackageEventHandlers(logger *zap.Logger, kubernetesClient kubernetes.Interf envNs = newPkg.Spec.Environment.Namespace } - ctx := context.Background() err := utils.SetupRoleBinding(ctx, logger, kubernetesClient, fv1.PackageGetterRB, newPkg.ObjectMeta.Namespace, utils.GetPackageGetterCR(), fv1.ClusterRole, fv1.FissionFetcherSA, envNs) diff --git a/pkg/executor/executortype/poolmgr/poolpodcontroller.go b/pkg/executor/executortype/poolmgr/poolpodcontroller.go index 845e237f..5c3abd60 100644 --- a/pkg/executor/executortype/poolmgr/poolpodcontroller.go +++ b/pkg/executor/executortype/poolmgr/poolpodcontroller.go @@ -66,7 +66,7 @@ type ( } ) -func NewPoolPodController(logger *zap.Logger, +func NewPoolPodController(ctx context.Context, logger *zap.Logger, kubernetesClient kubernetes.Interface, namespace string, enableIstio bool, @@ -86,8 +86,8 @@ func NewPoolPodController(logger *zap.Logger, envDeleteQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "EnvDeleteQueue"), spCleanupPodQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "SpecializedPodCleanupQueue"), } - funcInformer.Informer().AddEventHandler(FunctionEventHandlers(p.logger, p.kubernetesClient, p.namespace, p.enableIstio)) - pkgInformer.Informer().AddEventHandler(PackageEventHandlers(p.logger, p.kubernetesClient, p.namespace)) + funcInformer.Informer().AddEventHandler(FunctionEventHandlers(ctx, p.logger, p.kubernetesClient, p.namespace, p.enableIstio)) + pkgInformer.Informer().AddEventHandler(PackageEventHandlers(ctx, p.logger, p.kubernetesClient, p.namespace)) envInformer.Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ AddFunc: p.enqueueEnvAdd, UpdateFunc: p.enqueueEnvUpdate, @@ -216,7 +216,7 @@ func (p *PoolPodController) enqueueEnvDelete(obj interface{}) { p.envDeleteQueue.Add(env) } -func (p *PoolPodController) Run(stopCh <-chan struct{}) { +func (p *PoolPodController) Run(ctx context.Context, stopCh <-chan struct{}) { defer utilruntime.HandleCrash() defer p.envCreateUpdateQueue.ShutDown() defer p.envDeleteQueue.ShutDown() @@ -228,20 +228,20 @@ func (p *PoolPodController) Run(stopCh <-chan struct{}) { p.logger.Fatal("failed to wait for caches to sync") } for i := 0; i < 4; i++ { - go wait.Until(p.workerRun("envCreateUpdate", p.envCreateUpdateQueueProcessFunc), time.Second, stopCh) + go wait.Until(p.workerRun(ctx, "envCreateUpdate", p.envCreateUpdateQueueProcessFunc), time.Second, stopCh) } - go wait.Until(p.workerRun("envDeleteQueue", p.envDeleteQueueProcessFunc), time.Second, stopCh) - go wait.Until(p.workerRun("spCleanupPodQueue", p.spCleanupPodQueueProcessFunc), time.Second, stopCh) + go wait.Until(p.workerRun(ctx, "envDeleteQueue", p.envDeleteQueueProcessFunc), time.Second, stopCh) + go wait.Until(p.workerRun(ctx, "spCleanupPodQueue", p.spCleanupPodQueueProcessFunc), time.Second, stopCh) p.logger.Info("Started workers for poolPodController") <-stopCh p.logger.Info("Shutting down workers for poolPodController") } -func (p *PoolPodController) workerRun(name string, processFunc func() bool) func() { +func (p *PoolPodController) workerRun(ctx context.Context, name string, processFunc func(ctx context.Context) bool) func() { return func() { p.logger.Debug("Starting worker with func", zap.String("name", name)) for { - if quit := processFunc(); quit { + if quit := processFunc(ctx); quit { p.logger.Info("Shutting down worker", zap.String("name", name)) return } @@ -249,7 +249,7 @@ func (p *PoolPodController) workerRun(name string, processFunc func() bool) func } } -func (p *PoolPodController) envCreateUpdateQueueProcessFunc() bool { +func (p *PoolPodController) envCreateUpdateQueueProcessFunc(ctx context.Context) bool { maxRetries := 3 handleEnv := func(ctx context.Context, env *fv1.Environment) error { log := p.logger.With(zap.String("env", env.ObjectMeta.Name), zap.String("namespace", env.ObjectMeta.Namespace)) @@ -310,7 +310,6 @@ func (p *PoolPodController) envCreateUpdateQueueProcessFunc() bool { return false } - ctx := context.Background() err = handleEnv(ctx, env) if err != nil { if p.envCreateUpdateQueue.NumRequeues(key) < maxRetries { @@ -326,7 +325,7 @@ func (p *PoolPodController) envCreateUpdateQueueProcessFunc() bool { return false } -func (p *PoolPodController) envDeleteQueueProcessFunc() bool { +func (p *PoolPodController) envDeleteQueueProcessFunc(ctx context.Context) bool { obj, quit := p.envDeleteQueue.Get() if quit { return true @@ -338,7 +337,6 @@ func (p *PoolPodController) envDeleteQueueProcessFunc() bool { p.envDeleteQueue.Forget(obj) return false } - ctx := context.Background() p.logger.Debug("env delete request processing") p.gpm.cleanupPool(ctx, env) specializePodLables := getSpecializedPodLabels(env) @@ -368,7 +366,7 @@ func (p *PoolPodController) envDeleteQueueProcessFunc() bool { return false } -func (p *PoolPodController) spCleanupPodQueueProcessFunc() bool { +func (p *PoolPodController) spCleanupPodQueueProcessFunc(ctx context.Context) bool { maxRetries := 3 obj, quit := p.spCleanupPodQueue.Get() if quit { @@ -403,7 +401,6 @@ func (p *PoolPodController) spCleanupPodQueueProcessFunc() bool { } return false } - ctx := context.Background() podName := strings.SplitAfter(pod.GetName(), ".") if fsvc, ok := p.gpm.fsCache.PodToFsvc.Load(strings.TrimSuffix(podName[0], ".")); ok { fsvc, ok := fsvc.(*fscache.FuncSvc) diff --git a/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go b/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go index 55b27684..cf2d99ca 100644 --- a/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go +++ b/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go @@ -44,6 +44,8 @@ func runInformers(ctx context.Context, informers []k8sCache.SharedIndexInformer) } func TestPoolPodControllerPodCleanup(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() logger := loggerfactory.GetLogger() kubernetesClient := fake.NewSimpleClientset() fissionClient := fClient.NewSimpleClientset() @@ -60,7 +62,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) { gpmRsInformer := gpmInformerFactory.Apps().V1().ReplicaSets() fnNamespace := "fission-function" - ppc := NewPoolPodController(logger, kubernetesClient, fnNamespace, false, + ppc := NewPoolPodController(ctx, logger, kubernetesClient, fnNamespace, false, funcInformer, pkgInformer, envInformer, @@ -73,7 +75,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) { if err != nil { t.Fatalf("Error creating fetcher config: %v", err) } - executor, err := MakeGenericPoolManager( + executor, err := MakeGenericPoolManager(ctx, logger, fissionClient, kubernetesClient, metricsClient, fnNamespace, fetcherConfig, executorInstanceID, @@ -85,10 +87,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) { gpm := executor.(*GenericPoolManager) ppc.InjectGpm(gpm) - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() - - go ppc.Run(ctx.Done()) + go ppc.Run(ctx, ctx.Done()) podInformer := gpmPodInformer.Informer() diff --git a/pkg/executor/fscache/functionServiceCache_test.go b/pkg/executor/fscache/functionServiceCache_test.go index 48dfabdd..4fb705dd 100644 --- a/pkg/executor/fscache/functionServiceCache_test.go +++ b/pkg/executor/fscache/functionServiceCache_test.go @@ -184,7 +184,9 @@ func TestFunctionServiceNewCache(t *testing.T) { }, } - ctx := context.Background() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + fsc.AddFunc(ctx, *fsvc) _, active, err := fsc.GetFuncSvc(ctx, fsvc.Function, 5) if err != nil { diff --git a/pkg/executor/util/util_test.go b/pkg/executor/util/util_test.go index 78c7fcb5..cdc25f73 100644 --- a/pkg/executor/util/util_test.go +++ b/pkg/executor/util/util_test.go @@ -60,7 +60,9 @@ securityContext: Data: configMapData, } - configmap, err := kubeClient.CoreV1().ConfigMaps("fission").Create(context.Background(), &testConfigMap, metav1.CreateOptions{}) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + configmap, err := kubeClient.CoreV1().ConfigMaps("fission").Create(ctx, &testConfigMap, metav1.CreateOptions{}) if err != nil { t.Errorf("Error creating configmap %v", err) } @@ -106,7 +108,7 @@ securityContext: } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - got, err := GetSpecFromConfigMap(context.Background(), kubeClient, tt.cm, tt.cmns) + got, err := GetSpecFromConfigMap(ctx, kubeClient, tt.cm, tt.cmns) if (err != nil) != tt.wantErr { t.Errorf("GetSpecFromConfigMap() error = %v, wantErr %v", err, tt.wantErr) return diff --git a/pkg/fetcher/config/config.go b/pkg/fetcher/config/config.go index 0e59d077..06179797 100644 --- a/pkg/fetcher/config/config.go +++ b/pkg/fetcher/config/config.go @@ -1,6 +1,7 @@ package container import ( + "context" "encoding/json" "fmt" "log" @@ -90,8 +91,8 @@ func MakeFetcherConfig(sharedMountPath string) (*Config, error) { }, nil } -func (cfg *Config) SetupServiceAccount(kubernetesClient kubernetes.Interface, namespace string, context interface{}) error { - _, err := utils.SetupSA(kubernetesClient, fv1.FissionFetcherSA, namespace) +func (cfg *Config) SetupServiceAccount(ctx context.Context, kubernetesClient kubernetes.Interface, namespace string, context interface{}) error { + _, err := utils.SetupSA(ctx, kubernetesClient, fv1.FissionFetcherSA, namespace) if err != nil { log.Printf("Error : %v creating %s in ns : %s for: %#v", err, fv1.FissionFetcherSA, namespace, context) return err diff --git a/pkg/fission-cli/cmd/check/check.go b/pkg/fission-cli/cmd/check/check.go index cfda16e8..c74df766 100644 --- a/pkg/fission-cli/cmd/check/check.go +++ b/pkg/fission-cli/cmd/check/check.go @@ -45,6 +45,6 @@ func (opts *CheckSubCommand) do(input cli.Input) error { FissionClient: opts.Client(), }) - healthcheck.RunChecks(hc) + healthcheck.RunChecks(input.Context(), hc) return nil } diff --git a/pkg/fission-cli/cmd/plugin/list.go b/pkg/fission-cli/cmd/plugin/list.go index 02ae5216..3546d625 100644 --- a/pkg/fission-cli/cmd/plugin/list.go +++ b/pkg/fission-cli/cmd/plugin/list.go @@ -37,7 +37,7 @@ func List(input cli.Input) error { func (opts *ListSubCommand) do(input cli.Input) error { w := tabwriter.NewWriter(os.Stdout, 0, 0, 1, ' ', 0) fmt.Fprintln(w, "NAME\tVERSION\tPATH") - for _, p := range plugin.FindAll() { + for _, p := range plugin.FindAll(input.Context()) { fmt.Fprintf(w, "%v\t%v\t%v\n", p.Name, p.Version, p.Path) } w.Flush() diff --git a/pkg/fission-cli/cmd/support/resources/fissionversion.go b/pkg/fission-cli/cmd/support/resources/fissionversion.go index 24c4a844..5e162d04 100644 --- a/pkg/fission-cli/cmd/support/resources/fissionversion.go +++ b/pkg/fission-cli/cmd/support/resources/fissionversion.go @@ -34,7 +34,7 @@ func NewFissionVersion(client client.Interface) Resource { } func (res FissionVersion) Dump(ctx context.Context, dumpDir string) { - ver := util.GetVersion(res.client) + ver := util.GetVersion(ctx, res.client) file := filepath.Clean(fmt.Sprintf("%v/%v", dumpDir, "fission-version.txt")) writeToFile(file, ver) } diff --git a/pkg/fission-cli/cmd/version/version.go b/pkg/fission-cli/cmd/version/version.go index 39d472da..cb1ef9a3 100644 --- a/pkg/fission-cli/cmd/version/version.go +++ b/pkg/fission-cli/cmd/version/version.go @@ -36,7 +36,7 @@ func Version(input cli.Input) error { } func (opts *VersionSubCommand) do(input cli.Input) error { - ver := util.GetVersion(opts.Client()) + ver := util.GetVersion(input.Context(), opts.Client()) bs, err := yaml.Marshal(ver) if err != nil { return errors.Wrap(err, "error formatting versions") diff --git a/pkg/fission-cli/util/util.go b/pkg/fission-cli/util/util.go index 686afaae..0ad8e0c4 100644 --- a/pkg/fission-cli/util/util.go +++ b/pkg/fission-cli/util/util.go @@ -171,7 +171,7 @@ func CheckFunctionExistence(client client.Interface, functions []string, fnNames return nil } -func GetVersion(client client.Interface) info.Versions { +func GetVersion(ctx context.Context, client client.Interface) info.Versions { // Fetch client versions versions := info.Versions{ Client: map[string]info.BuildMeta{ @@ -179,7 +179,7 @@ func GetVersion(client client.Interface) info.Versions { }, } - for _, pmd := range plugin.FindAll() { + for _, pmd := range plugin.FindAll(ctx) { versions.Client[pmd.Name] = info.BuildMeta{ Version: pmd.Version, } diff --git a/pkg/healthcheck/healthcheck.go b/pkg/healthcheck/healthcheck.go index dd336c10..72f2a7f0 100644 --- a/pkg/healthcheck/healthcheck.go +++ b/pkg/healthcheck/healthcheck.go @@ -50,7 +50,7 @@ type Category struct { type Checker struct { successMsg string - check func() error + check func(ctx context.Context) error } type Options struct { @@ -102,8 +102,8 @@ func (hc *HealthChecker) CheckKubeVersion() (err error) { return nil } -func (hc *HealthChecker) CheckServiceStatus(namespace string, name string) (err error) { - depl, err := hc.kubeAPI.AppsV1().Deployments(namespace).Get(context.TODO(), name, metav1.GetOptions{}) +func (hc *HealthChecker) CheckServiceStatus(ctx context.Context, namespace string, name string) (err error) { + depl, err := hc.kubeAPI.AppsV1().Deployments(namespace).Get(ctx, name, metav1.GetOptions{}) if err != nil { return fmt.Errorf("failed to get %s deployment status", name) } @@ -112,7 +112,7 @@ func (hc *HealthChecker) CheckServiceStatus(namespace string, name string) (err return fmt.Errorf("%s deployment is not running", name) } - _, err = hc.kubeAPI.CoreV1().Services(namespace).Get(context.TODO(), name, metav1.GetOptions{}) + _, err = hc.kubeAPI.CoreV1().Services(namespace).Get(ctx, name, metav1.GetOptions{}) if err != nil { return fmt.Errorf("failed to get %s service status", name) } @@ -120,8 +120,8 @@ func (hc *HealthChecker) CheckServiceStatus(namespace string, name string) (err return nil } -func (hc *HealthChecker) CheckFissionVersion() error { - ver := util.GetVersion(hc.FissionClient) +func (hc *HealthChecker) CheckFissionVersion(ctx context.Context) error { + ver := util.GetVersion(ctx, hc.FissionClient) clientVersion := ver.Client["fission/core"].Version serverVersion := ver.Server["fission/core"].Version @@ -148,7 +148,7 @@ func (hc *HealthChecker) allCategories() []*Category { []Checker{ { successMsg: "kubernetes version is compatible", - check: func() (err error) { + check: func(ctx context.Context) (err error) { return hc.CheckKubeVersion() }, }, @@ -160,26 +160,26 @@ func (hc *HealthChecker) allCategories() []*Category { []Checker{ { successMsg: "controller is running fine", - check: func() error { - return hc.CheckServiceStatus(hc.fissionNamespace, "controller") + check: func(ctx context.Context) error { + return hc.CheckServiceStatus(ctx, hc.fissionNamespace, "controller") }, }, { successMsg: "executor is running fine", - check: func() error { - return hc.CheckServiceStatus(hc.fissionNamespace, "executor") + check: func(ctx context.Context) error { + return hc.CheckServiceStatus(ctx, hc.fissionNamespace, "executor") }, }, { successMsg: "router is running fine", - check: func() error { - return hc.CheckServiceStatus(hc.fissionNamespace, "router") + check: func(ctx context.Context) error { + return hc.CheckServiceStatus(ctx, hc.fissionNamespace, "router") }, }, { successMsg: "storagesvc is running fine", - check: func() error { - return hc.CheckServiceStatus(hc.fissionNamespace, "storagesvc") + check: func(ctx context.Context) error { + return hc.CheckServiceStatus(ctx, hc.fissionNamespace, "storagesvc") }, }, }, @@ -190,8 +190,8 @@ func (hc *HealthChecker) allCategories() []*Category { []Checker{ { successMsg: "fission is up-to-date", - check: func() error { - return hc.CheckFissionVersion() + check: func(ctx context.Context) error { + return hc.CheckFissionVersion(ctx) }, }, }, @@ -224,13 +224,13 @@ func NewHealthChecker(categoryIDs []CategoryID, options *Options) *HealthChecker return hc } -func RunChecks(hc *HealthChecker) { +func RunChecks(ctx context.Context, hc *HealthChecker) { for _, c := range hc.categories { if c.enabled { fmt.Println(c.ID) fmt.Println(strings.Repeat("-", 20)) for _, checker := range c.checkers { - err := checker.check() + err := checker.check(ctx) if err != nil { fmt.Printf("%s %s\n", failStatus, err) } else { diff --git a/pkg/kubewatcher/main.go b/pkg/kubewatcher/main.go index 41540293..f0249fad 100644 --- a/pkg/kubewatcher/main.go +++ b/pkg/kubewatcher/main.go @@ -39,7 +39,7 @@ func Start(ctx context.Context, logger *zap.Logger, routerUrl string) error { poster := publisher.MakeWebhookPublisher(logger, routerUrl) kubeWatch := MakeKubeWatcher(ctx, logger, kubeClient, poster) - MakeWatchSync(logger, fissionClient, kubeWatch) + MakeWatchSync(ctx, logger, fissionClient, kubeWatch) return nil } diff --git a/pkg/kubewatcher/watchSync.go b/pkg/kubewatcher/watchSync.go index 161ce219..35dfd74a 100644 --- a/pkg/kubewatcher/watchSync.go +++ b/pkg/kubewatcher/watchSync.go @@ -34,20 +34,20 @@ type ( } ) -func MakeWatchSync(logger *zap.Logger, client versioned.Interface, kubeWatcher *KubeWatcher) *WatchSync { +func MakeWatchSync(ctx context.Context, logger *zap.Logger, client versioned.Interface, kubeWatcher *KubeWatcher) *WatchSync { ws := &WatchSync{ logger: logger.Named("watch_sync"), client: client, kubeWatcher: kubeWatcher, } - go ws.syncSvc() + go ws.syncSvc(ctx) return ws } -func (ws *WatchSync) syncSvc() { +func (ws *WatchSync) syncSvc(ctx context.Context) { // TODO watch instead of polling for { - watches, err := ws.client.CoreV1().KubernetesWatchTriggers(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{}) + watches, err := ws.client.CoreV1().KubernetesWatchTriggers(metav1.NamespaceAll).List(ctx, metav1.ListOptions{}) if err != nil { ws.logger.Fatal("failed to get Kubernetes watch trigger list", zap.Error(err)) } diff --git a/pkg/mqtrigger/scalermanager.go b/pkg/mqtrigger/scalermanager.go index 6a3bfbe8..80068c6b 100644 --- a/pkg/mqtrigger/scalermanager.go +++ b/pkg/mqtrigger/scalermanager.go @@ -66,7 +66,7 @@ func getAuthTriggerClient(namespace string) (dynamic.ResourceInterface, error) { return dynamicClient.Resource(authTriggerGVR).Namespace(namespace), nil } -func mqTriggerEventHandlers(logger *zap.Logger, kubeClient kubernetes.Interface, routerURL string) k8sCache.ResourceEventHandlerFuncs { +func mqTriggerEventHandlers(ctx context.Context, logger *zap.Logger, kubeClient kubernetes.Interface, routerURL string) k8sCache.ResourceEventHandlerFuncs { return k8sCache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { go func() { @@ -79,17 +79,17 @@ func mqTriggerEventHandlers(logger *zap.Logger, kubeClient kubernetes.Interface, authenticationRef := "" if len(mqt.Spec.Secret) > 0 { authenticationRef = fmt.Sprintf("%s-auth-trigger", mqt.ObjectMeta.Name) - err := createAuthTrigger(mqt, authenticationRef, kubeClient) + err := createAuthTrigger(ctx, mqt, authenticationRef, kubeClient) if err != nil { logger.Error("Failed to create Authentication Trigger", zap.Error(err)) return } } - if err := createDeployment(mqt, routerURL, kubeClient); err != nil { + if err := createDeployment(ctx, mqt, routerURL, kubeClient); err != nil { logger.Error("Failed to create Deployment", zap.Error(err)) if len(authenticationRef) > 0 { - err = deleteAuthTrigger(authenticationRef, mqt.ObjectMeta.Namespace) + err = deleteAuthTrigger(ctx, authenticationRef, mqt.ObjectMeta.Namespace) if err != nil { logger.Error("Failed to delete Authentication Trigger", zap.Error(err)) } @@ -97,14 +97,14 @@ func mqTriggerEventHandlers(logger *zap.Logger, kubeClient kubernetes.Interface, return } - if err := createScaledObject(mqt, authenticationRef); err != nil { + if err := createScaledObject(ctx, mqt, authenticationRef); err != nil { logger.Error("Failed to create ScaledObject", zap.Error(err)) if len(authenticationRef) > 0 { - if err = deleteAuthTrigger(authenticationRef, mqt.ObjectMeta.Namespace); err != nil { + if err = deleteAuthTrigger(ctx, authenticationRef, mqt.ObjectMeta.Namespace); err != nil { logger.Error("Failed to delete Authentication Trigger", zap.Error(err)) } } - if err = deleteDeployment(mqt.ObjectMeta.Name, mqt.ObjectMeta.Namespace, kubeClient); err != nil { + if err = deleteDeployment(ctx, mqt.ObjectMeta.Name, mqt.ObjectMeta.Namespace, kubeClient); err != nil { logger.Error("Failed to delete Deployment", zap.Error(err)) } } @@ -126,18 +126,18 @@ func mqTriggerEventHandlers(logger *zap.Logger, kubeClient kubernetes.Interface, authenticationRef := "" if len(newMqt.Spec.Secret) > 0 && newMqt.Spec.Secret != mqt.Spec.Secret { authenticationRef = fmt.Sprintf("%s-auth-trigger", mqt.ObjectMeta.Name) - if err := updateAuthTrigger(mqt, authenticationRef, kubeClient); err != nil { + if err := updateAuthTrigger(ctx, mqt, authenticationRef, kubeClient); err != nil { logger.Error("Failed to update Authentication Trigger", zap.Error(err)) return } } - if err := updateDeployment(mqt, routerURL, kubeClient); err != nil { + if err := updateDeployment(ctx, mqt, routerURL, kubeClient); err != nil { logger.Error("Failed to Update Deployment", zap.Error(err)) return } - if err := updateScaledObject(mqt, authenticationRef); err != nil { + if err := updateScaledObject(ctx, mqt, authenticationRef); err != nil { logger.Error("Failed to Update ScaledObject", zap.Error(err)) return } @@ -160,7 +160,7 @@ func StartScalerManager(ctx context.Context, logger *zap.Logger, routerURL strin } informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) mqTriggerInformer := informerFactory.Core().V1().MessageQueueTriggers().Informer() - mqTriggerInformer.AddEventHandler(mqTriggerEventHandlers(logger, kubeClient, routerURL)) + mqTriggerInformer.AddEventHandler(mqTriggerEventHandlers(ctx, logger, kubeClient, routerURL)) mqTriggerInformer.Run(ctx.Done()) return nil } @@ -171,7 +171,7 @@ func toEnvVar(str string) string { return strings.ToUpper(envVar) } -func getEnvVarlist(mqt *fv1.MessageQueueTrigger, routerURL string, kubeClient kubernetes.Interface) ([]apiv1.EnvVar, error) { +func getEnvVarlist(ctx context.Context, mqt *fv1.MessageQueueTrigger, routerURL string, kubeClient kubernetes.Interface) ([]apiv1.EnvVar, error) { url := routerURL + "/" + strings.TrimPrefix(utils.UrlForFunction(mqt.Spec.FunctionReference.Name, mqt.ObjectMeta.Namespace), "/") envVars := []apiv1.EnvVar{ { @@ -214,7 +214,7 @@ func getEnvVarlist(mqt *fv1.MessageQueueTrigger, routerURL string, kubeClient ku // Add Auth Fields secretName := mqt.Spec.Secret if len(secretName) > 0 { - secret, err := kubeClient.CoreV1().Secrets(apiv1.NamespaceDefault).Get(context.TODO(), secretName, metav1.GetOptions{}) + secret, err := kubeClient.CoreV1().Secrets(apiv1.NamespaceDefault).Get(ctx, secretName, metav1.GetOptions{}) if err != nil { return nil, err } @@ -300,8 +300,8 @@ func checkAndUpdateTriggerFields(mqt, newMqt *fv1.MessageQueueTrigger) bool { return updated } -func getAuthTriggerSpec(mqt *fv1.MessageQueueTrigger, authenticationRef string, kubeClient kubernetes.Interface) (*unstructured.Unstructured, error) { - secret, err := kubeClient.CoreV1().Secrets(apiv1.NamespaceDefault).Get(context.TODO(), mqt.Spec.Secret, metav1.GetOptions{}) +func getAuthTriggerSpec(ctx context.Context, mqt *fv1.MessageQueueTrigger, authenticationRef string, kubeClient kubernetes.Interface) (*unstructured.Unstructured, error) { + secret, err := kubeClient.CoreV1().Secrets(apiv1.NamespaceDefault).Get(ctx, mqt.Spec.Secret, metav1.GetOptions{}) if err != nil { return nil, err } @@ -338,8 +338,8 @@ func getAuthTriggerSpec(mqt *fv1.MessageQueueTrigger, authenticationRef string, return authTriggerObj, nil } -func createAuthTrigger(mqt *fv1.MessageQueueTrigger, authenticationRef string, kubeClient kubernetes.Interface) error { - authTriggerObj, err := getAuthTriggerSpec(mqt, authenticationRef, kubeClient) +func createAuthTrigger(ctx context.Context, mqt *fv1.MessageQueueTrigger, authenticationRef string, kubeClient kubernetes.Interface) error { + authTriggerObj, err := getAuthTriggerSpec(ctx, mqt, authenticationRef, kubeClient) if err != nil { return err } @@ -347,50 +347,50 @@ func createAuthTrigger(mqt *fv1.MessageQueueTrigger, authenticationRef string, k if err != nil { return err } - _, err = authTriggerClient.Create(context.Background(), authTriggerObj, metav1.CreateOptions{}) + _, err = authTriggerClient.Create(ctx, authTriggerObj, metav1.CreateOptions{}) if err != nil { return err } return nil } -func updateAuthTrigger(mqt *fv1.MessageQueueTrigger, authenticationRef string, kubeClient kubernetes.Interface) error { +func updateAuthTrigger(ctx context.Context, mqt *fv1.MessageQueueTrigger, authenticationRef string, kubeClient kubernetes.Interface) error { authTriggerClient, err := getAuthTriggerClient(mqt.ObjectMeta.Namespace) if err != nil { return err } - oldAuthTriggerObj, err := authTriggerClient.Get(context.Background(), authenticationRef, metav1.GetOptions{}) + oldAuthTriggerObj, err := authTriggerClient.Get(ctx, authenticationRef, metav1.GetOptions{}) if err != nil { return err } resourceVersion := oldAuthTriggerObj.GetResourceVersion() - authTriggerObj, err := getAuthTriggerSpec(mqt, authenticationRef, kubeClient) + authTriggerObj, err := getAuthTriggerSpec(ctx, mqt, authenticationRef, kubeClient) if err != nil { return err } authTriggerObj.SetResourceVersion(resourceVersion) - _, err = authTriggerClient.Update(context.Background(), authTriggerObj, metav1.UpdateOptions{}) + _, err = authTriggerClient.Update(ctx, authTriggerObj, metav1.UpdateOptions{}) if err != nil { return err } return nil } -func deleteAuthTrigger(name, namespace string) error { +func deleteAuthTrigger(ctx context.Context, name, namespace string) error { authTriggerClient, err := getAuthTriggerClient(namespace) if err != nil { return err } - err = authTriggerClient.Delete(context.Background(), name, metav1.DeleteOptions{}) + err = authTriggerClient.Delete(ctx, name, metav1.DeleteOptions{}) if err != nil { return err } return nil } -func getDeploymentSpec(mqt *fv1.MessageQueueTrigger, routerURL string, kubeClient kubernetes.Interface) (*appsv1.Deployment, error) { - envVars, err := getEnvVarlist(mqt, routerURL, kubeClient) +func getDeploymentSpec(ctx context.Context, mqt *fv1.MessageQueueTrigger, routerURL string, kubeClient kubernetes.Interface) (*appsv1.Deployment, error) { + envVars, err := getEnvVarlist(ctx, mqt, routerURL, kubeClient) if err != nil { return nil, err } @@ -448,33 +448,33 @@ func getDeploymentSpec(mqt *fv1.MessageQueueTrigger, routerURL string, kubeClien }, nil } -func createDeployment(mqt *fv1.MessageQueueTrigger, routerURL string, kubeClient kubernetes.Interface) error { - deployment, err := getDeploymentSpec(mqt, routerURL, kubeClient) +func createDeployment(ctx context.Context, mqt *fv1.MessageQueueTrigger, routerURL string, kubeClient kubernetes.Interface) error { + deployment, err := getDeploymentSpec(ctx, mqt, routerURL, kubeClient) if err != nil { return err } - _, err = kubeClient.AppsV1().Deployments(mqt.ObjectMeta.Namespace).Create(context.TODO(), deployment, metav1.CreateOptions{}) + _, err = kubeClient.AppsV1().Deployments(mqt.ObjectMeta.Namespace).Create(ctx, deployment, metav1.CreateOptions{}) if err != nil { return err } return nil } -func updateDeployment(mqt *fv1.MessageQueueTrigger, routerURL string, kubeClient kubernetes.Interface) error { - deployment, err := getDeploymentSpec(mqt, routerURL, kubeClient) +func updateDeployment(ctx context.Context, mqt *fv1.MessageQueueTrigger, routerURL string, kubeClient kubernetes.Interface) error { + deployment, err := getDeploymentSpec(ctx, mqt, routerURL, kubeClient) if err != nil { return err } - _, err = kubeClient.AppsV1().Deployments(mqt.ObjectMeta.Namespace).Update(context.TODO(), deployment, metav1.UpdateOptions{}) + _, err = kubeClient.AppsV1().Deployments(mqt.ObjectMeta.Namespace).Update(ctx, deployment, metav1.UpdateOptions{}) if err != nil { return err } return nil } -func deleteDeployment(name string, namespace string, kubeClient kubernetes.Interface) error { +func deleteDeployment(ctx context.Context, name string, namespace string, kubeClient kubernetes.Interface) error { deletePolicy := metav1.DeletePropagationForeground - if err := kubeClient.AppsV1().Deployments(namespace).Delete(context.TODO(), name, metav1.DeleteOptions{ + if err := kubeClient.AppsV1().Deployments(namespace).Delete(ctx, name, metav1.DeleteOptions{ PropagationPolicy: &deletePolicy, }); err != nil { return err @@ -522,25 +522,25 @@ func getScaledObject(mqt *fv1.MessageQueueTrigger, authenticationRef string) *un } } -func createScaledObject(mqt *fv1.MessageQueueTrigger, authenticationRef string) error { +func createScaledObject(ctx context.Context, mqt *fv1.MessageQueueTrigger, authenticationRef string) error { scaledObject := getScaledObject(mqt, authenticationRef) kedaClient, err := getScaledObjectClient(mqt.ObjectMeta.Namespace) if err != nil { return err } - _, err = kedaClient.Create(context.Background(), scaledObject, metav1.CreateOptions{}) + _, err = kedaClient.Create(ctx, scaledObject, metav1.CreateOptions{}) if err != nil { return err } return nil } -func updateScaledObject(mqt *fv1.MessageQueueTrigger, authenticationRef string) error { +func updateScaledObject(ctx context.Context, mqt *fv1.MessageQueueTrigger, authenticationRef string) error { kedaClient, err := getScaledObjectClient(mqt.ObjectMeta.Namespace) if err != nil { return err } - oldScaledObject, err := kedaClient.Get(context.Background(), mqt.ObjectMeta.Name, metav1.GetOptions{}) + oldScaledObject, err := kedaClient.Get(ctx, mqt.ObjectMeta.Name, metav1.GetOptions{}) if err != nil { return err } @@ -549,7 +549,7 @@ func updateScaledObject(mqt *fv1.MessageQueueTrigger, authenticationRef string) scaledObject := getScaledObject(mqt, authenticationRef) scaledObject.SetResourceVersion(resourceVersion) - _, err = kedaClient.Update(context.Background(), scaledObject, metav1.UpdateOptions{}) + _, err = kedaClient.Update(ctx, scaledObject, metav1.UpdateOptions{}) if err != nil { return err } diff --git a/pkg/mqtrigger/scalermanager_test.go b/pkg/mqtrigger/scalermanager_test.go index 9ec89b36..2cee7453 100644 --- a/pkg/mqtrigger/scalermanager_test.go +++ b/pkg/mqtrigger/scalermanager_test.go @@ -40,6 +40,8 @@ func Test_toEnvVar(t *testing.T) { } func Test_getEnvVarlist(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() // Kafka Test with Valid Secret pollingInterval := int32(30) cooldownPeriod := int32(300) @@ -97,7 +99,7 @@ func Test_getEnvVarlist(t *testing.T) { } kubeClient := fake.NewSimpleClientset() - _, err := kubeClient.CoreV1().Secrets(namespace).Create(context.Background(), secret, metav1.CreateOptions{}) + _, err := kubeClient.CoreV1().Secrets(namespace).Create(ctx, secret, metav1.CreateOptions{}) if err != nil { assert.Equal(t, nil, err) } @@ -217,7 +219,7 @@ func Test_getEnvVarlist(t *testing.T) { } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - got, err := getEnvVarlist(tt.args.mqt, tt.args.routerURL, tt.args.kubeClient) + got, err := getEnvVarlist(ctx, tt.args.mqt, tt.args.routerURL, tt.args.kubeClient) sort.Slice(got, func(i, j int) bool { return got[i].Name < got[j].Name }) @@ -371,7 +373,8 @@ func Test_checkAndUpdateTriggerFields(t *testing.T) { } func Test_getAuthTriggerSpec(t *testing.T) { - + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() // Valid - with Secret pollingInterval := int32(30) cooldownPeriod := int32(300) @@ -428,7 +431,7 @@ func Test_getAuthTriggerSpec(t *testing.T) { } kubeClient := fake.NewSimpleClientset() - _, err := kubeClient.CoreV1().Secrets(namespace).Create(context.Background(), secret, metav1.CreateOptions{}) + _, err := kubeClient.CoreV1().Secrets(namespace).Create(ctx, secret, metav1.CreateOptions{}) if err != nil { assert.Equal(t, nil, err) } @@ -537,7 +540,7 @@ func Test_getAuthTriggerSpec(t *testing.T) { } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - got, err := getAuthTriggerSpec(tt.args.mqt, tt.args.authenticationRef, tt.args.kubeClient) + got, err := getAuthTriggerSpec(ctx, tt.args.mqt, tt.args.authenticationRef, tt.args.kubeClient) if (err != nil) != tt.wantErr { t.Errorf("getAuthTriggerSpec() error = %v, wantErr %v", err, tt.wantErr) return diff --git a/pkg/plugin/plugin.go b/pkg/plugin/plugin.go index c746ddb5..4b6e832c 100644 --- a/pkg/plugin/plugin.go +++ b/pkg/plugin/plugin.go @@ -69,13 +69,13 @@ func (md *Metadata) HasAlias(needle string) bool { // Find searches the machine for the given plugin, returning the metadata of the plugin. // The only metadata that is guaranteed to be non-empty is the path and Name. All other fields are considered optional. // If found it returns the plugin, otherwise it returns ErrPluginNotFound if the plugin was not found. -func Find(pluginName string) (*Metadata, error) { +func Find(ctx context.Context, pluginName string) (*Metadata, error) { // Search PATH for plugin as command-name // To check if plugin is actually there still. pluginPath, err := findPluginOnPath(pluginName) if err != nil { // Fallback: Search for alias in each command - mds := FindAll() + mds := FindAll(ctx) for _, md := range mds { if md.HasAlias(pluginName) { return md, nil @@ -84,7 +84,7 @@ func Find(pluginName string) (*Metadata, error) { return nil, ErrPluginNotFound } - md, err := fetchPluginMetadata(pluginPath) + md, err := fetchPluginMetadata(ctx, pluginPath) if err != nil { return nil, err } @@ -102,7 +102,7 @@ func Exec(md *Metadata, args []string) error { } // FindAll searches the machine for all plugins currently present. -func FindAll() map[string]*Metadata { +func FindAll(ctx context.Context) map[string]*Metadata { plugins := map[string]*Metadata{} dirs := strings.Split(os.Getenv("PATH"), ":") @@ -116,7 +116,7 @@ func FindAll() map[string]*Metadata { continue } fp := path.Join(dir, f.Name()) - md, err := fetchPluginMetadata(fp) + md, err := fetchPluginMetadata(ctx, fp) if err != nil { continue } @@ -142,7 +142,7 @@ func findPluginOnPath(pluginName string) (path string, err error) { } // fetchPluginMetadata attempts to fetch the plugin metadata given the plugin path. -func fetchPluginMetadata(pluginPath string) (*Metadata, error) { +func fetchPluginMetadata(ctx context.Context, pluginPath string) (*Metadata, error) { d, err := os.Stat(pluginPath) if err != nil { return nil, ErrPluginNotFound @@ -153,7 +153,7 @@ func fetchPluginMetadata(pluginPath string) (*Metadata, error) { // Fetch the metadata from the plugin itself. buf := bytes.NewBuffer(nil) - ctx, cancel := context.WithTimeout(context.Background(), cmdTimeout) + ctx, cancel := context.WithTimeout(ctx, cmdTimeout) defer cancel() cmd := exec.CommandContext(ctx, pluginPath, cmdMetadataArgs) // Note: issue can occur with signal propagation diff --git a/pkg/plugin/plugin_test.go b/pkg/plugin/plugin_test.go index 727e5a69..3218ad08 100644 --- a/pkg/plugin/plugin_test.go +++ b/pkg/plugin/plugin_test.go @@ -17,6 +17,7 @@ limitations under the License. package plugin import ( + "context" "encoding/json" "fmt" "os" @@ -57,7 +58,10 @@ func TestFind(t *testing.T) { } Prefix = "" - found, err := Find(md.Name) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + found, err := Find(ctx, md.Name) os.RemoveAll(testDir) assert.NoError(t, err) assert.NotEmpty(t, found) diff --git a/pkg/poolcache/poolcache_test.go b/pkg/poolcache/poolcache_test.go index 4079160b..643f2bc4 100644 --- a/pkg/poolcache/poolcache_test.go +++ b/pkg/poolcache/poolcache_test.go @@ -17,7 +17,8 @@ func checkErr(err error) { } func TestPoolCache(t *testing.T) { - ctx := context.Background() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() logger := loggerfactory.GetLogger() c := NewPoolCache(logger) diff --git a/pkg/router/functionHandler.go b/pkg/router/functionHandler.go index 4f6a7553..47cf3a04 100644 --- a/pkg/router/functionHandler.go +++ b/pkg/router/functionHandler.go @@ -258,7 +258,7 @@ func (roundTripper *RetryingRoundTripper) RoundTrip(req *http.Request) (*http.Re "service-entry": roundTripper.serviceURL.String()})...) if roundTripper.funcHandler.function.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType == fv1.ExecutorTypePoolmgr { defer func(ctx context.Context, fn *fv1.Function, serviceURL *url.URL) { - go roundTripper.funcHandler.unTapService(fn, serviceURL) //nolint errcheck + go roundTripper.funcHandler.unTapService(context.Background(), fn, serviceURL) //nolint errcheck }(ctx, roundTripper.funcHandler.function, roundTripper.serviceURL) } @@ -584,9 +584,9 @@ func (roundTripper RetryingRoundTripper) addForwardedHostHeader(req *http.Reques } // unTapservice marks the serviceURL in executor's cache as inactive, so that it can be reused -func (fh functionHandler) unTapService(fn *fv1.Function, serviceUrl *url.URL) error { +func (fh functionHandler) unTapService(ctx context.Context, fn *fv1.Function, serviceUrl *url.URL) error { fh.logger.Debug("UnTapService Called") - ctx, cancel := context.WithTimeout(context.Background(), fh.unTapServiceTimeout) + ctx, cancel := context.WithTimeout(ctx, fh.unTapServiceTimeout) defer cancel() err := fh.executor.UnTapService(ctx, fn.ObjectMeta, fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType, serviceUrl) if err != nil { diff --git a/pkg/router/httpTriggers.go b/pkg/router/httpTriggers.go index 47c4a547..09690ef5 100644 --- a/pkg/router/httpTriggers.go +++ b/pkg/router/httpTriggers.go @@ -283,13 +283,13 @@ func (ts *HTTPTriggerSet) addTriggerHandlers() { ts.triggerInformer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { trigger := obj.(*fv1.HTTPTrigger) - go createIngress(ts.logger, trigger, ts.kubeClient) + go createIngress(context.Background(), ts.logger, trigger, ts.kubeClient) ts.syncTriggers() }, DeleteFunc: func(obj interface{}) { ts.syncTriggers() trigger := obj.(*fv1.HTTPTrigger) - go deleteIngress(ts.logger, trigger, ts.kubeClient) + go deleteIngress(context.Background(), ts.logger, trigger, ts.kubeClient) }, UpdateFunc: func(oldObj interface{}, newObj interface{}) { oldTrigger := oldObj.(*fv1.HTTPTrigger) @@ -299,7 +299,7 @@ func (ts *HTTPTriggerSet) addTriggerHandlers() { return } - go updateIngress(ts.logger, oldTrigger, newTrigger, ts.kubeClient) + go updateIngress(context.Background(), ts.logger, oldTrigger, newTrigger, ts.kubeClient) ts.syncTriggers() }, }) diff --git a/pkg/router/ingress.go b/pkg/router/ingress.go index cdfd3ebb..d4bb654d 100644 --- a/pkg/router/ingress.go +++ b/pkg/router/ingress.go @@ -39,11 +39,11 @@ func init() { } } -func createIngress(logger *zap.Logger, trigger *fv1.HTTPTrigger, kubeClient kubernetes.Interface) { +func createIngress(ctx context.Context, logger *zap.Logger, trigger *fv1.HTTPTrigger, kubeClient kubernetes.Interface) { if !trigger.Spec.CreateIngress { return } - _, err := kubeClient.NetworkingV1().Ingresses(podNamespace).Create(context.TODO(), util.GetIngressSpec(podNamespace, trigger), v1.CreateOptions{}) + _, err := kubeClient.NetworkingV1().Ingresses(podNamespace).Create(ctx, util.GetIngressSpec(podNamespace, trigger), v1.CreateOptions{}) if err != nil && !k8serrors.IsAlreadyExists(err) { logger.Error("failed to create ingress", zap.Error(err)) return @@ -51,18 +51,18 @@ func createIngress(logger *zap.Logger, trigger *fv1.HTTPTrigger, kubeClient kube logger.Debug("created ingress successfully for trigger", zap.String("trigger", trigger.ObjectMeta.Name)) } -func deleteIngress(logger *zap.Logger, trigger *fv1.HTTPTrigger, kubeClient kubernetes.Interface) { +func deleteIngress(ctx context.Context, logger *zap.Logger, trigger *fv1.HTTPTrigger, kubeClient kubernetes.Interface) { if !trigger.Spec.CreateIngress { return } - ingress, err := kubeClient.NetworkingV1().Ingresses(podNamespace).Get(context.TODO(), trigger.ObjectMeta.Name, v1.GetOptions{}) + ingress, err := kubeClient.NetworkingV1().Ingresses(podNamespace).Get(ctx, trigger.ObjectMeta.Name, v1.GetOptions{}) if err != nil && !k8serrors.IsNotFound(err) { logger.Error("failed to get ingress when deleting trigger", zap.Error(err), zap.String("trigger", trigger.ObjectMeta.Name)) return } - err = kubeClient.NetworkingV1().Ingresses(podNamespace).Delete(context.TODO(), ingress.Name, v1.DeleteOptions{}) + err = kubeClient.NetworkingV1().Ingresses(podNamespace).Delete(ctx, ingress.Name, v1.DeleteOptions{}) if err != nil && !k8serrors.IsNotFound(err) { logger.Error("failed to delete ingress for trigger", zap.Error(err), @@ -71,25 +71,25 @@ func deleteIngress(logger *zap.Logger, trigger *fv1.HTTPTrigger, kubeClient kube } } -func updateIngress(logger *zap.Logger, oldT *fv1.HTTPTrigger, newT *fv1.HTTPTrigger, kubeClient kubernetes.Interface) { +func updateIngress(ctx context.Context, logger *zap.Logger, oldT *fv1.HTTPTrigger, newT *fv1.HTTPTrigger, kubeClient kubernetes.Interface) { if !oldT.Spec.CreateIngress && !newT.Spec.CreateIngress { return } if !oldT.Spec.CreateIngress && newT.Spec.CreateIngress { - createIngress(logger, newT, kubeClient) + createIngress(ctx, logger, newT, kubeClient) return } if !newT.Spec.CreateIngress && oldT.Spec.CreateIngress { - deleteIngress(logger, oldT, kubeClient) + deleteIngress(ctx, logger, oldT, kubeClient) return } - oldIngress, err := kubeClient.NetworkingV1().Ingresses(podNamespace).Get(context.TODO(), oldT.ObjectMeta.Name, v1.GetOptions{}) + oldIngress, err := kubeClient.NetworkingV1().Ingresses(podNamespace).Get(ctx, oldT.ObjectMeta.Name, v1.GetOptions{}) if err != nil { if k8serrors.IsNotFound(err) { - createIngress(logger, newT, kubeClient) + createIngress(ctx, logger, newT, kubeClient) } logger.Error("failed to get ingress when updating trigger", zap.Error(err), @@ -123,7 +123,7 @@ func updateIngress(logger *zap.Logger, oldT *fv1.HTTPTrigger, newT *fv1.HTTPTrig } if changes { - _, err = kubeClient.NetworkingV1().Ingresses(podNamespace).Update(context.TODO(), oldIngress, v1.UpdateOptions{}) + _, err = kubeClient.NetworkingV1().Ingresses(podNamespace).Update(ctx, oldIngress, v1.UpdateOptions{}) if err != nil { logger.Error("failed to update ingress for trigger", zap.Error(err), zap.String("trigger", oldT.ObjectMeta.Name)) return diff --git a/pkg/tracker/tracker_test.go b/pkg/tracker/tracker_test.go index 0a04fefd..fcf5b055 100644 --- a/pkg/tracker/tracker_test.go +++ b/pkg/tracker/tracker_test.go @@ -14,6 +14,8 @@ import ( ) func TestTracker(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() t.Run("NewTracker", func(test *testing.T) { for _, test := range []struct { name string @@ -101,7 +103,7 @@ func TestTracker(t *testing.T) { defer server.Close() tr.gaAPIURL = server.URL - t := tr.SendEvent(context.Background(), *test.request) + t := tr.SendEvent(ctx, *test.request) if test.status == http.StatusOK { assert.Nil(testing, t, test.expected) } else { diff --git a/pkg/utils/otel/provider_test.go b/pkg/utils/otel/provider_test.go index 8f217b7e..943d8b86 100644 --- a/pkg/utils/otel/provider_test.go +++ b/pkg/utils/otel/provider_test.go @@ -62,7 +62,8 @@ func TestGetTraceExporter(t *testing.T) { t.Errorf("Expected OTEL_EXPORTER_OTLP_INSECURE to be set, got %s", OtelInsecureEnvVar) } logger := loggerfactory.GetLogger() - ctx := context.Background() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() tests := []struct { oltpEndpoint string oltpInsecure string diff --git a/pkg/utils/rbacutils.go b/pkg/utils/rbacutils.go index 6ea9e969..905b4c2b 100644 --- a/pkg/utils/rbacutils.go +++ b/pkg/utils/rbacutils.go @@ -50,15 +50,15 @@ func MakeSAObj(sa, ns string) *apiv1.ServiceAccount { } // SetupSA checks if a service account is present in the namespace, if not creates it. -func SetupSA(k8sClient kubernetes.Interface, sa, ns string) (*apiv1.ServiceAccount, error) { - saObj, err := k8sClient.CoreV1().ServiceAccounts(ns).Get(context.TODO(), sa, metav1.GetOptions{}) +func SetupSA(ctx context.Context, k8sClient kubernetes.Interface, sa, ns string) (*apiv1.ServiceAccount, error) { + saObj, err := k8sClient.CoreV1().ServiceAccounts(ns).Get(ctx, sa, metav1.GetOptions{}) if err == nil { return saObj, nil } if k8serrors.IsNotFound(err) { saObj = MakeSAObj(sa, ns) - saObj, err = k8sClient.CoreV1().ServiceAccounts(ns).Create(context.TODO(), saObj, metav1.CreateOptions{}) + saObj, err = k8sClient.CoreV1().ServiceAccounts(ns).Create(ctx, saObj, metav1.CreateOptions{}) } return saObj, err