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
This commit is contained in:
Shubham Bansal
2022-11-04 11:57:37 +05:30
committed by GitHub
parent 6af53807aa
commit 261bf24974
2 changed files with 117 additions and 190 deletions
+2 -2
View File
@@ -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()
+115 -188
View File
@@ -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{