From 261bf2497411b47ca87016ae0a3f67a2335a3d5e Mon Sep 17 00:00:00 2001 From: Shubham Bansal <62992590+shubham-bansal96@users.noreply.github.com> Date: Fri, 4 Nov 2022 11:57:37 +0530 Subject: [PATCH] Use informer for environment handling in buildermanager with multiple namespaces (#2603) * changes to add informer for environment * remove unnecessary code * code refactor * code review changes --- pkg/buildermgr/buildermgr.go | 4 +- pkg/buildermgr/envwatcher.go | 303 +++++++++++++---------------------- 2 files changed, 117 insertions(+), 190 deletions(-) diff --git a/pkg/buildermgr/buildermgr.go b/pkg/buildermgr/buildermgr.go index 5d3c3f6e..7244657f 100644 --- a/pkg/buildermgr/buildermgr.go +++ b/pkg/buildermgr/buildermgr.go @@ -62,8 +62,8 @@ func Start(ctx context.Context, logger *zap.Logger, storageSvcUrl string, envBui } } - envWatcher := makeEnvironmentWatcher(bmLogger, fissionClient, kubernetesClient, fetcherConfig, envBuilderNamespace, podSpecPatch) - go envWatcher.watchEnvironments(ctx) + envWatcher := makeEnvironmentWatcher(ctx, bmLogger, fissionClient, kubernetesClient, fetcherConfig, envBuilderNamespace, podSpecPatch) + envWatcher.Run(ctx) k8sInformerFactory := k8sInformers.NewSharedInformerFactory(kubernetesClient, time.Minute*30) podInformer := k8sInformerFactory.Core().V1().Pods().Informer() diff --git a/pkg/buildermgr/envwatcher.go b/pkg/buildermgr/envwatcher.go index b4046b0b..597d9945 100644 --- a/pkg/buildermgr/envwatcher.go +++ b/pkg/buildermgr/envwatcher.go @@ -30,22 +30,18 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/util/intstr" - "k8s.io/apimachinery/pkg/watch" "k8s.io/client-go/kubernetes" + k8sCache "k8s.io/client-go/tools/cache" fv1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/executor/util" fetcherConfig "github.com/fission/fission/pkg/fetcher/config" "github.com/fission/fission/pkg/generated/clientset/versioned" "github.com/fission/fission/pkg/utils" ) -type requestType int - const ( - GET_BUILDER requestType = iota - CLEANUP_BUILDERS - LABEL_ENV_NAME = "envName" LABEL_ENV_NAMESPACE = "envNamespace" LABEL_ENV_RESOURCEVERSION = "envResourceVersion" @@ -65,23 +61,9 @@ type ( service *apiv1.Service } - envwRequest struct { - requestType - ctx context.Context - env *fv1.Environment - envList []fv1.Environment - respChan chan envwResponse - } - - envwResponse struct { - builderInfo *builderInfo - err error - } - environmentWatcher struct { logger *zap.Logger cache map[string]*builderInfo - requestChan chan envwRequest builderNamespace string fissionClient versioned.Interface kubernetesClient kubernetes.Interface @@ -89,10 +71,12 @@ type ( builderImagePullPolicy apiv1.PullPolicy useIstio bool podSpecPatch *apiv1.PodSpec + envWatchInformer map[string]k8sCache.SharedIndexInformer } ) func makeEnvironmentWatcher( + ctx context.Context, logger *zap.Logger, fissionClient versioned.Interface, kubernetesClient kubernetes.Interface, @@ -115,7 +99,6 @@ func makeEnvironmentWatcher( envWatcher := &environmentWatcher{ logger: logger.Named("environment_watcher"), cache: make(map[string]*builderInfo), - requestChan: make(chan envwRequest), builderNamespace: builderNamespace, fissionClient: fissionClient, kubernetesClient: kubernetesClient, @@ -123,17 +106,13 @@ func makeEnvironmentWatcher( useIstio: useIstio, fetcherConfig: fetcherConfig, podSpecPatch: podSpecPatch, + envWatchInformer: utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.EnvironmentResource), } - go envWatcher.service() - + envWatcher.EnvWatchEventHandlers(ctx) return envWatcher } -func (envw *environmentWatcher) getCacheKey(envName string, envNamespace string, envResourceVersion string) string { - return fmt.Sprintf("%v-%v-%v", envName, envNamespace, envResourceVersion) -} - func (env *environmentWatcher) getLabelForDeploymentOwner() map[string]string { return map[string]string{ LABEL_DEPLOYMENT_OWNER: BUILDER_MGR, @@ -149,178 +128,123 @@ func (envw *environmentWatcher) getLabels(envName string, envNamespace string, e } } -func (envw *environmentWatcher) watchEnvironments(ctx context.Context) { - rv := "" - for { - wi, err := envw.fissionClient.CoreV1().Environments(metav1.NamespaceAll).Watch(ctx, - metav1.ListOptions{ - ResourceVersion: rv, - }) - if err != nil { - if utils.IsNetworkError(err) { - envw.logger.Error("encountered network error, retrying later", zap.Error(err)) - time.Sleep(5 * time.Second) - continue - } - envw.logger.Fatal("error watching environment list", zap.Error(err)) - } - - for { - ev, more := <-wi.ResultChan() - if !more { - // restart watch from last rv - break - } - if ev.Type == watch.Error { - // restart watch from the start - rv = "" - time.Sleep(time.Second) - break - } - env := ev.Object.(*fv1.Environment) - rv = env.ObjectMeta.ResourceVersion - envw.sync(ctx) - } +func (envw *environmentWatcher) Run(ctx context.Context) { + for _, informer := range envw.envWatchInformer { + go informer.Run(ctx.Done()) } } -func (envw *environmentWatcher) sync(ctx context.Context) { - maxRetries := 10 - for i := 0; i < maxRetries; i++ { - 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)) - time.Sleep(50 * time.Duration(2*i) * time.Millisecond) - continue - } - envw.logger.Fatal("error syncing environment CRD resources", zap.Error(err)) - } +func (envw *environmentWatcher) EnvWatchEventHandlers(ctx context.Context) { + for _, informer := range envw.envWatchInformer { + informer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ + AddFunc: func(obj interface{}) { + envObj := obj.(*fv1.Environment) + envw.AddUpdateBuilder(ctx, envObj) + }, + UpdateFunc: func(oldObj interface{}, newObj interface{}) { + oldEnvObj := oldObj.(*fv1.Environment) + newEnvObj := newObj.(*fv1.Environment) + if oldEnvObj.ObjectMeta.ResourceVersion != newEnvObj.ObjectMeta.ResourceVersion { + envw.AddUpdateBuilder(ctx, newEnvObj) + } + }, + DeleteFunc: func(obj interface{}) { + envObj := obj.(*fv1.Environment) + envw.DeleteBuilder(ctx, envObj) + }, + }) + } +} - // Create environment builders for all environments - for i := range envList.Items { - env := envList.Items[i] - - if env.Spec.Version == 1 || // builder is not supported with v1 interface - len(env.Spec.Builder.Image) == 0 { // ignore env without builder image - continue - } - _, err := envw.getEnvBuilder(ctx, &env) +func (envw *environmentWatcher) AddUpdateBuilder(ctx context.Context, env *fv1.Environment) { + //builder is not supported with v1 interface and ignore env without builder image + if env.Spec.Version != 1 && len(env.Spec.Builder.Image) != 0 { + if _, ok := envw.cache[crd.CacheKeyUID(&env.ObjectMeta)]; !ok { + builderInfo, err := envw.createBuilder(ctx, env, envw.getNamespace(env)) if err != nil { - envw.logger.Error("error creating builder", zap.Error(err), zap.String("builder_target", env.ObjectMeta.Name)) + envw.logger.Error("error creating builder service", zap.Error(err)) + return } - } - - // Remove environment builders no longer needed - envw.cleanupEnvBuilders(ctx, envList.Items) - break - } -} - -func (envw *environmentWatcher) service() { - for { - req := <-envw.requestChan - switch req.requestType { - case GET_BUILDER: - // In order to support backward compatibility, for all environments with builder image created in default env, - // the pods will be created in fission-builder namespace - ns := envw.builderNamespace - if req.env.ObjectMeta.Namespace != metav1.NamespaceDefault { - ns = req.env.ObjectMeta.Namespace - } - - key := envw.getCacheKey(req.env.ObjectMeta.Name, ns, req.env.ObjectMeta.ResourceVersion) - builderInfo, ok := envw.cache[key] - if !ok { - builderInfo, err := envw.createBuilder(req.ctx, req.env, ns) - if err != nil { - req.respChan <- envwResponse{err: err} - continue - } - envw.cache[key] = builderInfo - } - req.respChan <- envwResponse{builderInfo: builderInfo} - - case CLEANUP_BUILDERS: - latestEnvList := make(map[string]*fv1.Environment) - for i := range req.envList { - env := req.envList[i] - // In order to support backward compatibility, for all builder images created in default - // env, the pods are created in fission-builder namespace - ns := envw.builderNamespace - if env.ObjectMeta.Namespace != metav1.NamespaceDefault { - ns = env.ObjectMeta.Namespace - } - key := envw.getCacheKey(env.ObjectMeta.Name, ns, env.ObjectMeta.ResourceVersion) - latestEnvList[key] = &env - } - - // If an environment is deleted when builder manager down, - // the builder belongs to the environment will be out-of- - // control (an orphan builder) since there is no record in - // cache and CRD. We need to iterate over the services & - // deployments to remove both normal and orphan builders. - - svcList, err := envw.getBuilderServiceList(req.ctx, envw.getLabelForDeploymentOwner(), metav1.NamespaceAll) + envw.cache[crd.CacheKeyUID(&env.ObjectMeta)] = builderInfo + } else { + envw.DeleteBuilder(ctx, env) + // once older builder deleted then add new builder service + builderInfo, err := envw.createBuilder(ctx, env, envw.getNamespace(env)) if err != nil { - envw.logger.Error("error getting the builder service list", zap.Error(err)) - } - for _, svc := range svcList { - envName := svc.ObjectMeta.Labels[LABEL_ENV_NAME] - envNamespace := svc.ObjectMeta.Labels[LABEL_ENV_NAMESPACE] - envResourceVersion := svc.ObjectMeta.Labels[LABEL_ENV_RESOURCEVERSION] - key := envw.getCacheKey(envName, envNamespace, envResourceVersion) - if _, ok := latestEnvList[key]; !ok { - 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), - zap.String("service_namespace", svc.ObjectMeta.Namespace)) - } - } - delete(envw.cache, key) - } - - 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)) - } - for _, deploy := range deployList { - envName := deploy.ObjectMeta.Labels[LABEL_ENV_NAME] - envNamespace := deploy.ObjectMeta.Labels[LABEL_ENV_NAMESPACE] - envResourceVersion := deploy.ObjectMeta.Labels[LABEL_ENV_RESOURCEVERSION] - key := envw.getCacheKey(envName, envNamespace, envResourceVersion) - if _, ok := latestEnvList[key]; !ok { - 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), - zap.String("deployment_namespace", deploy.ObjectMeta.Namespace)) - } - } - delete(envw.cache, key) + envw.logger.Error("error updating builder service", zap.Error(err)) + return } + envw.cache[crd.CacheKeyUID(&env.ObjectMeta)] = builderInfo } } } -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, +func (envw *environmentWatcher) DeleteBuilder(ctx context.Context, env *fv1.Environment) { + if _, ok := envw.cache[crd.CacheKeyUID(&env.ObjectMeta)]; ok { + envw.DeleteBuilderService(ctx, env) + envw.DeleteBuilderDeployment(ctx, env) + delete(envw.cache, crd.CacheKeyUID(&env.ObjectMeta)) + envw.logger.Info("builder service deleted", zap.String("env_name", env.ObjectMeta.Name), zap.String("namespace", envw.getNamespace(env))) + } else { + envw.logger.Debug("builder service not found", zap.String("env_name", env.ObjectMeta.Name), zap.String("namespace", envw.getNamespace(env))) } - resp := <-respChan - return resp.builderInfo, resp.err } -func (envw *environmentWatcher) cleanupEnvBuilders(ctx context.Context, envs []fv1.Environment) { - envw.requestChan <- envwRequest{ - requestType: CLEANUP_BUILDERS, - ctx: ctx, - envList: envs, +func (envw *environmentWatcher) getNamespace(env *fv1.Environment) string { + // In order to support backward compatibility, for all environments with builder image created in default env, + // the pods will be created in fission-builder namespace + ns := envw.builderNamespace + if env.ObjectMeta.Namespace != metav1.NamespaceDefault { + ns = env.ObjectMeta.Namespace + } + return ns +} + +func (envw *environmentWatcher) DeleteBuilderService(ctx context.Context, env *fv1.Environment) { + ns := envw.getNamespace(env) + svcList, err := envw.getBuilderServiceList(ctx, envw.getLabelForDeploymentOwner(), ns) + if err != nil { + envw.logger.Error("error getting the builder service list", zap.Error(err)) + } + for _, svc := range svcList { + envName := svc.ObjectMeta.Labels[LABEL_ENV_NAME] + if _, ok := envw.cache[crd.CacheKeyUID(&env.ObjectMeta)]; ok { + err := envw.deleteBuilderServiceByName(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), + zap.String("service_namespace", svc.ObjectMeta.Namespace), + zap.String("env_name", envName)) + } + break + } else { + envw.logger.Error("builder service not found", + zap.String("service_name", svc.ObjectMeta.Name), + zap.String("service_namespace", svc.ObjectMeta.Namespace)) + } + } +} + +func (envw *environmentWatcher) DeleteBuilderDeployment(ctx context.Context, env *fv1.Environment) { + ns := envw.getNamespace(env) + deployList, err := envw.getBuilderDeploymentList(ctx, envw.getLabelForDeploymentOwner(), ns) + if err != nil { + envw.logger.Error("error getting the builder deployment list", zap.Error(err)) + } + for _, deploy := range deployList { + if _, ok := envw.cache[crd.CacheKeyUID(&env.ObjectMeta)]; ok { + err := envw.deleteBuilderDeploymentByName(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), + zap.String("deployment_namespace", deploy.ObjectMeta.Namespace)) + } + break + } else { + envw.logger.Error("builder deployment not found", zap.Error(err), + zap.String("deployment_name", deploy.ObjectMeta.Name), + zap.String("deployment_namespace", deploy.ObjectMeta.Namespace)) + } } } @@ -338,12 +262,14 @@ func (envw *environmentWatcher) createBuilder(ctx context.Context, env *fv1.Envi if len(svcList) == 0 { svc, err = envw.createBuilderService(ctx, env, ns) if err != nil { - return nil, errors.Wrap(err, "error creating builder service") + return nil, errors.Wrap(err, + fmt.Sprintf("error creating builder service for environment in namespace %s %s", env.ObjectMeta.Name, ns)) + } } else if len(svcList) == 1 { svc = &svcList[0] } else { - return nil, fmt.Errorf("found more than one builder service for environment %q", env.ObjectMeta.Name) + return nil, fmt.Errorf("found more than one builder service for environment in namespace %s %s", env.ObjectMeta.Name, ns) } deployList, err := envw.getBuilderDeploymentList(ctx, sel, ns) @@ -360,12 +286,13 @@ func (envw *environmentWatcher) createBuilder(ctx context.Context, env *fv1.Envi deploy, err = envw.createBuilderDeployment(ctx, env, ns) if err != nil { - return nil, errors.Wrap(err, "error creating builder deployment") + return nil, errors.Wrap(err, fmt.Sprintf("error creating builder deployment for environment in namespace %s %s", env.ObjectMeta.Name, ns)) + } } else if len(deployList) == 1 { deploy = &deployList[0] } else { - return nil, fmt.Errorf("found more than one builder deployment for environment %q", env.ObjectMeta.Name) + return nil, fmt.Errorf("found more than one builder deployment for environment in namespace %s %s", env.ObjectMeta.Name, ns) } return &builderInfo{