From f86fde81e63e98ace46a8de151c4ce789d0cb011 Mon Sep 17 00:00:00 2001 From: Sanket Sudake Date: Mon, 5 Jul 2021 20:28:26 +0530 Subject: [PATCH] Shared informers (#2092) * Use shared informers across executor for fission CRs Signed-off-by: Sanket Sudake * Add informer for configmap and secrets Signed-off-by: Sanket Sudake * Use shared informers in router Signed-off-by: Sanket Sudake * Add shared informer for canary config Signed-off-by: Sanket Sudake * Add sharedinformer for ready pod check Signed-off-by: Sanket Sudake * Refer informer instead of store in resolver Signed-off-by: Sanket Sudake * Run informers before executors Signed-off-by: Sanket Sudake * Add missing canary config handler calls Signed-off-by: Sanket Sudake * Enable websocket test Signed-off-by: Sanket Sudake * Run dump collection always Signed-off-by: Sanket Sudake --- .github/workflows/push_pr.yaml | 1 + .github/workflows/upgrade_test.yaml | 1 + pkg/canaryconfigmgr/canaryConfigMgr.go | 74 +++---- pkg/executor/cms/cmhandler.go | 73 ++++++ pkg/executor/cms/cmscontroller.go | 129 ++--------- pkg/executor/cms/secrethandler.go | 72 ++++++ pkg/executor/executor.go | 42 +++- .../executortype/newdeploy/envhandlers.go | 54 +++++ .../executortype/newdeploy/funchandlers.go | 68 ++++++ .../executortype/newdeploy/newdeploymgr.go | 123 ++--------- .../executortype/poolmgr/funchandlers.go | 198 +++++++++++++++++ .../executortype/poolmgr/funcwatcher.go | 208 ------------------ pkg/executor/executortype/poolmgr/gp.go | 5 +- pkg/executor/executortype/poolmgr/gpm.go | 28 +-- .../executortype/poolmgr/packagehandlers.go | 99 +++++++++ .../executortype/poolmgr/packagewatcher.go | 111 ---------- .../poolmgr/readyPodController.go | 24 +- pkg/router/functionReferenceResolver.go | 15 +- pkg/router/httpTriggers.go | 153 ++++++------- pkg/router/router.go | 14 +- pkg/utils/informer.go | 17 +- test/kind_CI.sh | 2 +- 22 files changed, 793 insertions(+), 718 deletions(-) create mode 100644 pkg/executor/cms/cmhandler.go create mode 100644 pkg/executor/cms/secrethandler.go create mode 100644 pkg/executor/executortype/newdeploy/envhandlers.go create mode 100644 pkg/executor/executortype/newdeploy/funchandlers.go create mode 100644 pkg/executor/executortype/poolmgr/funchandlers.go delete mode 100644 pkg/executor/executortype/poolmgr/funcwatcher.go create mode 100644 pkg/executor/executortype/poolmgr/packagehandlers.go delete mode 100644 pkg/executor/executortype/poolmgr/packagewatcher.go diff --git a/.github/workflows/push_pr.yaml b/.github/workflows/push_pr.yaml index bd38dcde..f021fb5c 100644 --- a/.github/workflows/push_pr.yaml +++ b/.github/workflows/push_pr.yaml @@ -108,6 +108,7 @@ jobs: run: ./test/kind_CI.sh - name: Collect Fission Dump + if: ${{ always() }} run: | command -v fission && fission support dump diff --git a/.github/workflows/upgrade_test.yaml b/.github/workflows/upgrade_test.yaml index 88877272..b24f2e0f 100644 --- a/.github/workflows/upgrade_test.yaml +++ b/.github/workflows/upgrade_test.yaml @@ -84,6 +84,7 @@ jobs: source ./test/upgrade_test/fission_objects.sh test_fission_objects - name: Collect Fission Dump + if: ${{ always() }} run: | command -v fission && fission support dump diff --git a/pkg/canaryconfigmgr/canaryConfigMgr.go b/pkg/canaryconfigmgr/canaryConfigMgr.go index 68c46522..a0c5e8db 100644 --- a/pkg/canaryconfigmgr/canaryConfigMgr.go +++ b/pkg/canaryconfigmgr/canaryConfigMgr.go @@ -28,13 +28,13 @@ import ( "go.uber.org/zap" k8serrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/fields" "k8s.io/client-go/kubernetes" "k8s.io/client-go/rest" k8sCache "k8s.io/client-go/tools/cache" fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/crd" + genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" ) const ( @@ -45,8 +45,7 @@ type canaryConfigMgr struct { logger *zap.Logger fissionClient *crd.FissionClient kubeClient *kubernetes.Clientset - canaryConfigStore k8sCache.Store - canaryConfigController k8sCache.Controller + canaryConfigInformer *k8sCache.SharedIndexInformer promClient *PrometheusApiClient crdClient rest.Interface canaryCfgCancelFuncMap *canaryConfigCancelFuncMap @@ -97,49 +96,44 @@ func MakeCanaryConfigMgr(logger *zap.Logger, fissionClient *crd.FissionClient, k canaryCfgCancelFuncMap: makecanaryConfigCancelFuncMap(), } - store, controller := configMgr.initCanaryConfigController() - configMgr.canaryConfigStore = store - configMgr.canaryConfigController = controller - + informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Second*30) + informer := informerFactory.Core().V1().CanaryConfigs().Informer() + configMgr.canaryConfigInformer = &informer + configMgr.CanaryConfigEventHandlers() return configMgr, nil } -func (canaryCfgMgr *canaryConfigMgr) initCanaryConfigController() (k8sCache.Store, k8sCache.Controller) { - resyncPeriod := 30 * time.Second - listWatch := k8sCache.NewListWatchFromClient(canaryCfgMgr.crdClient, "canaryconfigs", metav1.NamespaceAll, fields.Everything()) - store, controller := k8sCache.NewInformer(listWatch, &fv1.CanaryConfig{}, resyncPeriod, - k8sCache.ResourceEventHandlerFuncs{ - AddFunc: func(obj interface{}) { - canaryConfig := obj.(*fv1.CanaryConfig) - if canaryConfig.Status.Status == fv1.CanaryConfigStatusPending { - go canaryCfgMgr.addCanaryConfig(canaryConfig) - } - }, - DeleteFunc: func(obj interface{}) { - canaryConfig := obj.(*fv1.CanaryConfig) - go canaryCfgMgr.deleteCanaryConfig(canaryConfig) - }, - UpdateFunc: func(oldObj interface{}, newObj interface{}) { - oldConfig := oldObj.(*fv1.CanaryConfig) - newConfig := newObj.(*fv1.CanaryConfig) - if oldConfig.ObjectMeta.ResourceVersion != newConfig.ObjectMeta.ResourceVersion && - newConfig.Status.Status == fv1.CanaryConfigStatusPending { - canaryCfgMgr.logger.Info("update canary config invoked", - 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.reSyncCanaryConfigs() +func (canaryCfgMgr *canaryConfigMgr) CanaryConfigEventHandlers() { + (*canaryCfgMgr.canaryConfigInformer).AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ + AddFunc: func(obj interface{}) { + canaryConfig := obj.(*fv1.CanaryConfig) + if canaryConfig.Status.Status == fv1.CanaryConfigStatusPending { + go canaryCfgMgr.addCanaryConfig(canaryConfig) + } + }, + DeleteFunc: func(obj interface{}) { + canaryConfig := obj.(*fv1.CanaryConfig) + go canaryCfgMgr.deleteCanaryConfig(canaryConfig) + }, + UpdateFunc: func(oldObj interface{}, newObj interface{}) { + oldConfig := oldObj.(*fv1.CanaryConfig) + newConfig := newObj.(*fv1.CanaryConfig) + if oldConfig.ObjectMeta.ResourceVersion != newConfig.ObjectMeta.ResourceVersion && + newConfig.Status.Status == fv1.CanaryConfigStatusPending { + canaryCfgMgr.logger.Info("update canary config invoked", + 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.reSyncCanaryConfigs() - }, - }) - - return store, controller + }, + }) } func (canaryCfgMgr *canaryConfigMgr) Run(ctx context.Context) { - go canaryCfgMgr.canaryConfigController.Run(ctx.Done()) + go (*canaryCfgMgr.canaryConfigInformer).Run(ctx.Done()) canaryCfgMgr.logger.Info("started canary configmgr controller") } @@ -501,7 +495,7 @@ func (canaryCfgMgr *canaryConfigMgr) rollForward(canaryConfig *fv1.CanaryConfig, } func (canaryCfgMgr *canaryConfigMgr) reSyncCanaryConfigs() { - for _, obj := range canaryCfgMgr.canaryConfigStore.List() { + for _, obj := range (*canaryCfgMgr.canaryConfigInformer).GetStore().List() { canaryConfig := obj.(*fv1.CanaryConfig) _, err := canaryCfgMgr.canaryCfgCancelFuncMap.lookup(&canaryConfig.ObjectMeta) if err != nil && canaryConfig.Status.Status == fv1.CanaryConfigStatusPending { diff --git a/pkg/executor/cms/cmhandler.go b/pkg/executor/cms/cmhandler.go new file mode 100644 index 00000000..b28ac372 --- /dev/null +++ b/pkg/executor/cms/cmhandler.go @@ -0,0 +1,73 @@ +/* +Copyright 2021 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ +package cms + +import ( + "context" + + "go.uber.org/zap" + apiv1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "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/executortype" +) + +func getConfigmapRelatedFuncs(logger *zap.Logger, m *metav1.ObjectMeta, fissionClient *crd.FissionClient) ([]fv1.Function, error) { + funcList, err := fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{}) + if err != nil { + return nil, err + } + // In future a cache that populates at start and is updated on changes might be better solution + relatedFunctions := make([]fv1.Function, 0) + for _, f := range funcList.Items { + for _, cm := range f.Spec.ConfigMaps { + if (cm.Name == m.Name) && (cm.Namespace == m.Namespace) { + relatedFunctions = append(relatedFunctions, f) + break + } + } + } + return relatedFunctions, nil +} + +func ConfigMapEventHandlers(logger *zap.Logger, fissionClient *crd.FissionClient, + kubernetesClient *kubernetes.Clientset, types map[fv1.ExecutorType]executortype.ExecutorType) k8sCache.ResourceEventHandlerFuncs { + + return k8sCache.ResourceEventHandlerFuncs{ + AddFunc: func(obj interface{}) {}, + DeleteFunc: func(obj interface{}) {}, + UpdateFunc: func(oldObj interface{}, newObj interface{}) { + oldCm := oldObj.(*apiv1.ConfigMap) + newCm := newObj.(*apiv1.ConfigMap) + if oldCm.ObjectMeta.ResourceVersion != newCm.ObjectMeta.ResourceVersion { + if newCm.ObjectMeta.Namespace != "kube-system" { + logger.Debug("Configmap changed", + zap.String("configmap_name", newCm.ObjectMeta.Name), + zap.String("configmap_namespace", newCm.ObjectMeta.Namespace)) + } + funcs, err := getConfigmapRelatedFuncs(logger, &newCm.ObjectMeta, fissionClient) + if err != nil { + logger.Error("Failed to get functions related to configmap", zap.String("configmap_name", newCm.ObjectMeta.Name), zap.String("configmap_namespace", newCm.ObjectMeta.Namespace)) + } + refreshPods(logger, funcs, types) + } + }, + } +} diff --git a/pkg/executor/cms/cmscontroller.go b/pkg/executor/cms/cmscontroller.go index 9e76a23a..c5fdb4b0 100644 --- a/pkg/executor/cms/cmscontroller.go +++ b/pkg/executor/cms/cmscontroller.go @@ -17,16 +17,10 @@ limitations under the License. package cms import ( - "context" - "time" - "github.com/pkg/errors" "go.uber.org/zap" - apiv1 "k8s.io/api/core/v1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/fields" "k8s.io/client-go/kubernetes" - "k8s.io/client-go/tools/cache" + k8sCache "k8s.io/client-go/tools/cache" fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/crd" @@ -38,126 +32,31 @@ type ( ConfigSecretController struct { logger *zap.Logger - configmapController cache.Controller - secretController cache.Controller + configmapInformer *k8sCache.SharedIndexInformer + secretInformer *k8sCache.SharedIndexInformer fissionClient *crd.FissionClient } ) -//MakeConfigSecretController makes a controller for configmaps and secrets which changes related functions +// MakeConfigSecretController makes a controller for configmaps and secrets which changes related functions func MakeConfigSecretController(logger *zap.Logger, fissionClient *crd.FissionClient, - kubernetesClient *kubernetes.Clientset, types map[fv1.ExecutorType]executortype.ExecutorType) *ConfigSecretController { + kubernetesClient *kubernetes.Clientset, types map[fv1.ExecutorType]executortype.ExecutorType, + configmapInformer *k8sCache.SharedIndexInformer, + secretInformer *k8sCache.SharedIndexInformer) *ConfigSecretController { logger.Debug("Creating ConfigMap & Secret Controller") - _, cmcontroller := initConfigmapController(logger, fissionClient, kubernetesClient, types) - _, scontroller := initSecretController(logger, fissionClient, kubernetesClient, types) cmsController := &ConfigSecretController{ - logger: logger, - configmapController: cmcontroller, - secretController: scontroller, - fissionClient: fissionClient, + logger: logger, + configmapInformer: configmapInformer, + secretInformer: secretInformer, + fissionClient: fissionClient, } + (*configmapInformer).AddEventHandler(ConfigMapEventHandlers(logger, fissionClient, kubernetesClient, types)) + (*secretInformer).AddEventHandler(SecretEventHandlers(logger, fissionClient, kubernetesClient, types)) + return cmsController } -//Run runs the controllers for configmaps and secrets -func (csController *ConfigSecretController) Run(ctx context.Context) { - go csController.configmapController.Run(ctx.Done()) - go csController.secretController.Run(ctx.Done()) -} - -func initConfigmapController(logger *zap.Logger, fissionClient *crd.FissionClient, - kubernetesClient *kubernetes.Clientset, types map[fv1.ExecutorType]executortype.ExecutorType) (cache.Store, cache.Controller) { - resyncPeriod := 30 * time.Second - listWatch := cache.NewListWatchFromClient(kubernetesClient.CoreV1().RESTClient(), "configmaps", metav1.NamespaceAll, fields.Everything()) - store, controller := cache.NewInformer(listWatch, &apiv1.ConfigMap{}, resyncPeriod, cache.ResourceEventHandlerFuncs{ - AddFunc: func(obj interface{}) {}, - DeleteFunc: func(obj interface{}) {}, - UpdateFunc: func(oldObj interface{}, newObj interface{}) { - oldCm := oldObj.(*apiv1.ConfigMap) - newCm := newObj.(*apiv1.ConfigMap) - if oldCm.ObjectMeta.ResourceVersion != newCm.ObjectMeta.ResourceVersion { - if newCm.ObjectMeta.Namespace != "kube-system" { - logger.Debug("Configmap changed", - zap.String("configmap_name", newCm.ObjectMeta.Name), - zap.String("configmap_namespace", newCm.ObjectMeta.Namespace)) - } - funcs, err := getConfigmapRelatedFuncs(logger, &newCm.ObjectMeta, fissionClient) - if err != nil { - logger.Error("Failed to get functions related to configmap", zap.String("configmap_name", newCm.ObjectMeta.Name), zap.String("configmap_namespace", newCm.ObjectMeta.Namespace)) - } - refreshPods(logger, funcs, types) - } - }, - }) - return store, controller -} - -func getConfigmapRelatedFuncs(logger *zap.Logger, m *metav1.ObjectMeta, fissionClient *crd.FissionClient) ([]fv1.Function, error) { - funcList, err := fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{}) - if err != nil { - return nil, err - } - // In future a cache that populates at start and is updated on changes might be better solution - relatedFunctions := make([]fv1.Function, 0) - for _, f := range funcList.Items { - for _, cm := range f.Spec.ConfigMaps { - if (cm.Name == m.Name) && (cm.Namespace == m.Namespace) { - relatedFunctions = append(relatedFunctions, f) - break - } - } - } - return relatedFunctions, nil -} - -func initSecretController(logger *zap.Logger, fissionClient *crd.FissionClient, - kubernetesClient *kubernetes.Clientset, types map[fv1.ExecutorType]executortype.ExecutorType) (cache.Store, cache.Controller) { - resyncPeriod := 30 * time.Second - listWatch := cache.NewListWatchFromClient(kubernetesClient.CoreV1().RESTClient(), "secrets", metav1.NamespaceAll, fields.Everything()) - store, controller := cache.NewInformer(listWatch, &apiv1.Secret{}, resyncPeriod, cache.ResourceEventHandlerFuncs{ - AddFunc: func(obj interface{}) {}, - DeleteFunc: func(obj interface{}) {}, - UpdateFunc: func(oldObj interface{}, newObj interface{}) { - oldS := oldObj.(*apiv1.Secret) - newS := newObj.(*apiv1.Secret) - if oldS.ObjectMeta.ResourceVersion != newS.ObjectMeta.ResourceVersion { - if newS.ObjectMeta.Namespace != "kube-system" { - logger.Debug("Secret changed", - zap.String("configmap_name", newS.ObjectMeta.Name), - zap.String("configmap_namespace", newS.ObjectMeta.Namespace)) - - } - funcs, err := getSecretRelatedFuncs(logger, &newS.ObjectMeta, fissionClient) - if err != nil { - logger.Error("Failed to get functions related to secret", zap.String("secret_name", newS.ObjectMeta.Name), zap.String("secret_namespace", newS.ObjectMeta.Namespace)) - } - refreshPods(logger, funcs, types) - } - }, - }) - return store, controller - -} - -func getSecretRelatedFuncs(logger *zap.Logger, m *metav1.ObjectMeta, fissionClient *crd.FissionClient) ([]fv1.Function, error) { - funcList, err := fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{}) - if err != nil { - return nil, err - } - // In future a cache that populates at start and is updated on changes might be better solution - relatedFunctions := make([]fv1.Function, 0) - for _, f := range funcList.Items { - for _, secret := range f.Spec.Secrets { - if (secret.Name == m.Name) && (secret.Namespace == m.Namespace) { - relatedFunctions = append(relatedFunctions, f) - break - } - } - } - return relatedFunctions, nil -} - func refreshPods(logger *zap.Logger, funcs []fv1.Function, types map[fv1.ExecutorType]executortype.ExecutorType) { for _, f := range funcs { var err error diff --git a/pkg/executor/cms/secrethandler.go b/pkg/executor/cms/secrethandler.go new file mode 100644 index 00000000..36f8b41b --- /dev/null +++ b/pkg/executor/cms/secrethandler.go @@ -0,0 +1,72 @@ +/* +Copyright 2021 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ +package cms + +import ( + "context" + + "go.uber.org/zap" + apiv1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "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/executortype" +) + +func getSecretRelatedFuncs(logger *zap.Logger, m *metav1.ObjectMeta, fissionClient *crd.FissionClient) ([]fv1.Function, error) { + funcList, err := fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{}) + if err != nil { + return nil, err + } + // In future a cache that populates at start and is updated on changes might be better solution + relatedFunctions := make([]fv1.Function, 0) + for _, f := range funcList.Items { + for _, secret := range f.Spec.Secrets { + if (secret.Name == m.Name) && (secret.Namespace == m.Namespace) { + relatedFunctions = append(relatedFunctions, f) + break + } + } + } + return relatedFunctions, nil +} + +func SecretEventHandlers(logger *zap.Logger, fissionClient *crd.FissionClient, + kubernetesClient *kubernetes.Clientset, types map[fv1.ExecutorType]executortype.ExecutorType) k8sCache.ResourceEventHandlerFuncs { + return k8sCache.ResourceEventHandlerFuncs{ + AddFunc: func(obj interface{}) {}, + DeleteFunc: func(obj interface{}) {}, + UpdateFunc: func(oldObj interface{}, newObj interface{}) { + oldS := oldObj.(*apiv1.Secret) + newS := newObj.(*apiv1.Secret) + if oldS.ObjectMeta.ResourceVersion != newS.ObjectMeta.ResourceVersion { + if newS.ObjectMeta.Namespace != "kube-system" { + logger.Debug("Secret changed", + zap.String("configmap_name", newS.ObjectMeta.Name), + zap.String("configmap_namespace", newS.ObjectMeta.Namespace)) + } + funcs, err := getSecretRelatedFuncs(logger, &newS.ObjectMeta, fissionClient) + if err != nil { + logger.Error("Failed to get functions related to secret", zap.String("secret_name", newS.ObjectMeta.Name), zap.String("secret_namespace", newS.ObjectMeta.Namespace)) + } + refreshPods(logger, funcs, types) + } + }, + } +} diff --git a/pkg/executor/executor.go b/pkg/executor/executor.go index 87c78ac9..e25b2ae0 100644 --- a/pkg/executor/executor.go +++ b/pkg/executor/executor.go @@ -30,6 +30,8 @@ import ( "github.com/pkg/errors" "github.com/prometheus/client_golang/prometheus/promhttp" "go.uber.org/zap" + k8sInformers "k8s.io/client-go/informers" + k8sCache "k8s.io/client-go/tools/cache" fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/crd" @@ -41,6 +43,7 @@ import ( "github.com/fission/fission/pkg/executor/reaper" "github.com/fission/fission/pkg/executor/util" fetcherConfig "github.com/fission/fission/pkg/fetcher/config" + genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" ) type ( @@ -69,8 +72,9 @@ type ( ) // MakeExecutor returns an Executor for given ExecutorType(s). -func MakeExecutor(logger *zap.Logger, cms *cms.ConfigSecretController, - fissionClient *crd.FissionClient, types map[fv1.ExecutorType]executortype.ExecutorType) (*Executor, error) { +func MakeExecutor(ctx context.Context, logger *zap.Logger, cms *cms.ConfigSecretController, + fissionClient *crd.FissionClient, types map[fv1.ExecutorType]executortype.ExecutorType, + informers []k8sCache.SharedIndexInformer) (*Executor, error) { executor := &Executor{ logger: logger.Named("executor"), cms: cms, @@ -80,12 +84,18 @@ func MakeExecutor(logger *zap.Logger, cms *cms.ConfigSecretController, requestChan: make(chan *createFuncServiceRequest), fsCreateWg: make(map[string]*sync.WaitGroup), } + + // Run all informers + for _, informer := range informers { + go informer.Run(ctx.Done()) + } + for _, et := range types { go func(et executortype.ExecutorType) { - et.Run(context.Background()) + et.Run(ctx) }(et) } - go cms.Run(context.Background()) + go executor.serveCreateFuncServices() return executor, nil @@ -257,10 +267,17 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames logger.Info("Starting executor", zap.String("instanceID", executorInstanceID)) + informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Second*30) + funcInformer := informerFactory.Core().V1().Functions().Informer() + pkgInformer := informerFactory.Core().V1().Packages().Informer() + envInformer := informerFactory.Core().V1().Environments().Informer() + gpm, err := poolmgr.MakeGenericPoolManager( logger, fissionClient, kubernetesClient, metricsClient, - functionNamespace, fetcherConfig, executorInstanceID) + functionNamespace, fetcherConfig, executorInstanceID, + &funcInformer, &pkgInformer, + ) if err != nil { return errors.Wrap(err, "pool manager creation faied") } @@ -268,7 +285,9 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames ndm, err := newdeploy.MakeNewDeploy( logger, fissionClient, kubernetesClient, fissionClient.CoreV1().RESTClient(), - functionNamespace, fetcherConfig, executorInstanceID) + functionNamespace, fetcherConfig, executorInstanceID, + &funcInformer, &envInformer, + ) if err != nil { return errors.Wrap(err, "new deploy manager creation faied") } @@ -294,9 +313,16 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames // TODO: use context to control the waiting time once kubernetes client supports it. util.WaitTimeout(wg, 30*time.Second) - cms := cms.MakeConfigSecretController(logger, fissionClient, kubernetesClient, executorTypes) + k8sInformerFactory := k8sInformers.NewSharedInformerFactory(kubernetesClient, time.Second*30) + configmapInformer := k8sInformerFactory.Core().V1().ConfigMaps().Informer() + secretInformer := k8sInformerFactory.Core().V1().Secrets().Informer() - api, err := MakeExecutor(logger, cms, fissionClient, executorTypes) + cms := cms.MakeConfigSecretController(logger, fissionClient, kubernetesClient, executorTypes, &configmapInformer, &secretInformer) + + ctx := context.Background() + api, err := MakeExecutor(ctx, logger, cms, fissionClient, executorTypes, []k8sCache.SharedIndexInformer{ + funcInformer, pkgInformer, envInformer, configmapInformer, secretInformer, + }) if err != nil { return err } diff --git a/pkg/executor/executortype/newdeploy/envhandlers.go b/pkg/executor/executortype/newdeploy/envhandlers.go new file mode 100644 index 00000000..e091b234 --- /dev/null +++ b/pkg/executor/executortype/newdeploy/envhandlers.go @@ -0,0 +1,54 @@ +/* +Copyright 2021 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ +package newdeploy + +import ( + "context" + + "go.uber.org/zap" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + k8sCache "k8s.io/client-go/tools/cache" + + fv1 "github.com/fission/fission/pkg/apis/core/v1" +) + +func (deploy *NewDeploy) EnvEventHandlers() 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) + // Currently only an image update in environment calls for function's deployment recreation. In future there might be more attributes which would want to do it + if oldEnv.Spec.Runtime.Image != newEnv.Spec.Runtime.Image { + deploy.logger.Debug("Updating all function of the environment that changed, old env:", zap.Any("environment", oldEnv)) + funcs := deploy.getEnvFunctions(&newEnv.ObjectMeta) + for _, f := range funcs { + function, err := deploy.fissionClient.CoreV1().Functions(f.ObjectMeta.Namespace).Get(context.TODO(), f.ObjectMeta.Name, metav1.GetOptions{}) + if err != nil { + deploy.logger.Error("Error getting function", zap.Error(err), zap.Any("function", function)) + continue + } + err = deploy.updateFuncDeployment(function, newEnv) + if err != nil { + deploy.logger.Error("Error updating function", zap.Error(err), zap.Any("function", function)) + continue + } + } + } + }, + } +} diff --git a/pkg/executor/executortype/newdeploy/funchandlers.go b/pkg/executor/executortype/newdeploy/funchandlers.go new file mode 100644 index 00000000..0c551330 --- /dev/null +++ b/pkg/executor/executortype/newdeploy/funchandlers.go @@ -0,0 +1,68 @@ +/* +Copyright 2021 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ +package newdeploy + +import ( + "go.uber.org/zap" + k8sCache "k8s.io/client-go/tools/cache" + + fv1 "github.com/fission/fission/pkg/apis/core/v1" +) + +func (deploy *NewDeploy) FunctionEventHandlers() 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() { + fn := obj.(*fv1.Function) + deploy.logger.Debug("create deployment for function", zap.Any("fn", fn.ObjectMeta), zap.Any("fnspec", fn.Spec)) + _, err := deploy.createFunction(fn) + if err != nil { + deploy.logger.Error("error eager creating function", + zap.Error(err), + zap.Any("function", fn)) + } + deploy.logger.Debug("end create deployment for function", zap.Any("fn", fn.ObjectMeta), zap.Any("fnspec", fn.Spec)) + }() + }, + DeleteFunc: func(obj interface{}) { + fn := obj.(*fv1.Function) + go func() { + err := deploy.deleteFunction(fn) + if err != nil { + deploy.logger.Error("error deleting function", + zap.Error(err), + zap.Any("function", fn)) + } + }() + }, + UpdateFunc: func(oldObj interface{}, newObj interface{}) { + oldFn := oldObj.(*fv1.Function) + newFn := newObj.(*fv1.Function) + go func() { + err := deploy.updateFunction(oldFn, newFn) + if err != nil { + deploy.logger.Error("error updating function", + zap.Error(err), + zap.Any("old_function", oldFn), + zap.Any("new_function", newFn)) + } + }() + }, + } +} diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr.go b/pkg/executor/executortype/newdeploy/newdeploymgr.go index 14d21ebc..d4acb99c 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr.go @@ -33,7 +33,6 @@ import ( apiv1 "k8s.io/api/core/v1" k8sErrs "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/fields" "k8s.io/apimachinery/pkg/labels" k8sTypes "k8s.io/apimachinery/pkg/types" "k8s.io/client-go/kubernetes" @@ -69,12 +68,11 @@ type ( fsCache *fscache.FunctionServiceCache // cache funcSvc's by function, address and pod name - throttler *throttler.Throttler - funcStore k8sCache.Store - funcController k8sCache.Controller + throttler *throttler.Throttler + + funcInformer *k8sCache.SharedIndexInformer + envInformer *k8sCache.SharedIndexInformer - envStore k8sCache.Store - envController k8sCache.Controller serviceInformer k8sCache.SharedIndexInformer deploymentInformer k8sCache.SharedIndexInformer @@ -91,6 +89,8 @@ func MakeNewDeploy( namespace string, fetcherConfig *fetcherConfig.Config, instanceID string, + funcInformer *k8sCache.SharedIndexInformer, + envInformer *k8sCache.SharedIndexInformer, ) (executortype.ExecutorType, error) { enableIstio := false if len(os.Getenv("ENABLE_ISTIO")) > 0 { @@ -118,19 +118,16 @@ func MakeNewDeploy( useIstio: enableIstio, defaultIdlePodReapTime: 2 * time.Minute, + funcInformer: funcInformer, + envInformer: envInformer, } if nd.crdClient != nil { - fnStore, fnController := nd.initFuncController() - nd.funcStore = fnStore - nd.funcController = fnController - - envStore, envController := nd.initEnvController() - nd.envStore = envStore - nd.envController = envController + (*nd.funcInformer).AddEventHandler(nd.FunctionEventHandlers()) + (*nd.envInformer).AddEventHandler(nd.EnvEventHandlers()) } - informerFactory, err := utils.GetInformerFacoryByExecutor(nd.kubernetesClient, fv1.ExecutorTypePoolmgr) + informerFactory, err := utils.GetInformerFactoryByExecutor(nd.kubernetesClient, fv1.ExecutorTypePoolmgr) if err != nil { return nil, err } @@ -141,8 +138,6 @@ func MakeNewDeploy( // Run start the function and environment controller along with an object reaper. func (deploy *NewDeploy) Run(ctx context.Context) { - go deploy.funcController.Run(ctx.Done()) - go deploy.envController.Run(ctx.Done()) go deploy.serviceInformer.Run(ctx.Done()) go deploy.deploymentInformer.Run(ctx.Done()) go deploy.idleObjectReaper() @@ -191,19 +186,8 @@ func (deploy *NewDeploy) TapService(svcHost string) error { return nil } -func getCachedItem(obj apiv1.ObjectReference, informer k8sCache.SharedIndexInformer) (item interface{}, exists bool, err error) { - store := informer.GetStore() - - item, exists, err = store.Get(obj) - if err != nil || !exists { - item, exists, err = store.GetByKey(fmt.Sprintf("%s/%s", obj.Namespace, obj.Name)) - } - - return item, exists, err -} - func (deploy *NewDeploy) getServiceInfo(obj apiv1.ObjectReference) (*apiv1.Service, error) { - item, exists, err := getCachedItem(obj, deploy.serviceInformer) + item, exists, err := utils.GetCachedItem(obj, deploy.serviceInformer) if err != nil || !exists { deploy.logger.Debug( @@ -220,7 +204,7 @@ func (deploy *NewDeploy) getServiceInfo(obj apiv1.ObjectReference) (*apiv1.Servi } func (deploy *NewDeploy) getDeploymentInfo(obj apiv1.ObjectReference) (*appsv1.Deployment, error) { - item, exists, err := getCachedItem(obj, deploy.deploymentInformer) + item, exists, err := utils.GetCachedItem(obj, deploy.deploymentInformer) if err != nil || !exists { deploy.logger.Debug( @@ -378,85 +362,6 @@ func (deploy *NewDeploy) CleanupOldExecutorObjects() { } } -func (deploy *NewDeploy) initFuncController() (k8sCache.Store, k8sCache.Controller) { - resyncPeriod := 30 * time.Second - listWatch := k8sCache.NewListWatchFromClient(deploy.crdClient, "functions", metav1.NamespaceAll, fields.Everything()) - store, controller := k8sCache.NewInformer(listWatch, &fv1.Function{}, resyncPeriod, 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() { - fn := obj.(*fv1.Function) - deploy.logger.Debug("create deployment for function", zap.Any("fn", fn.ObjectMeta), zap.Any("fnspec", fn.Spec)) - _, err := deploy.createFunction(fn) - if err != nil { - deploy.logger.Error("error eager creating function", - zap.Error(err), - zap.Any("function", fn)) - } - deploy.logger.Debug("end create deployment for function", zap.Any("fn", fn.ObjectMeta), zap.Any("fnspec", fn.Spec)) - }() - }, - DeleteFunc: func(obj interface{}) { - fn := obj.(*fv1.Function) - go func() { - err := deploy.deleteFunction(fn) - if err != nil { - deploy.logger.Error("error deleting function", - zap.Error(err), - zap.Any("function", fn)) - } - }() - }, - UpdateFunc: func(oldObj interface{}, newObj interface{}) { - oldFn := oldObj.(*fv1.Function) - newFn := newObj.(*fv1.Function) - go func() { - err := deploy.updateFunction(oldFn, newFn) - if err != nil { - deploy.logger.Error("error updating function", - zap.Error(err), - zap.Any("old_function", oldFn), - zap.Any("new_function", newFn)) - } - }() - }, - }) - return store, controller -} - -func (deploy *NewDeploy) initEnvController() (k8sCache.Store, k8sCache.Controller) { - resyncPeriod := 30 * time.Second - listWatch := k8sCache.NewListWatchFromClient(deploy.crdClient, "environments", metav1.NamespaceAll, fields.Everything()) - store, controller := k8sCache.NewInformer(listWatch, &fv1.Environment{}, resyncPeriod, k8sCache.ResourceEventHandlerFuncs{ - AddFunc: func(obj interface{}) {}, - DeleteFunc: func(obj interface{}) {}, - UpdateFunc: func(oldObj interface{}, newObj interface{}) { - newEnv := newObj.(*fv1.Environment) - oldEnv := oldObj.(*fv1.Environment) - // Currently only an image update in environment calls for function's deployment recreation. In future there might be more attributes which would want to do it - if oldEnv.Spec.Runtime.Image != newEnv.Spec.Runtime.Image { - deploy.logger.Debug("Updating all function of the environment that changed, old env:", zap.Any("environment", oldEnv)) - funcs := deploy.getEnvFunctions(&newEnv.ObjectMeta) - for _, f := range funcs { - function, err := deploy.fissionClient.CoreV1().Functions(f.ObjectMeta.Namespace).Get(context.TODO(), f.ObjectMeta.Name, metav1.GetOptions{}) - if err != nil { - deploy.logger.Error("Error getting function", zap.Error(err), zap.Any("function", function)) - continue - } - err = deploy.updateFuncDeployment(function, newEnv) - if err != nil { - deploy.logger.Error("Error updating function", zap.Error(err), zap.Any("function", function)) - continue - } - } - } - }, - }) - return store, controller -} - func (deploy *NewDeploy) getEnvFunctions(m *metav1.ObjectMeta) []fv1.Function { funcList, err := deploy.fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{}) if err != nil { @@ -876,7 +781,7 @@ func (deploy *NewDeploy) idleObjectReaper() { // For function with the environment that no longer exists, executor // scales down the deployment as usual and prints log to notify user. if _, ok := envList[fsvc.Environment.ObjectMeta.UID]; !ok { - deploy.logger.Error("function environment no longer exists", + deploy.logger.Warn("function environment no longer exists", zap.String("environment", fsvc.Environment.ObjectMeta.Name), zap.String("function", fsvc.Name)) } diff --git a/pkg/executor/executortype/poolmgr/funchandlers.go b/pkg/executor/executortype/poolmgr/funchandlers.go new file mode 100644 index 00000000..e16dea16 --- /dev/null +++ b/pkg/executor/executortype/poolmgr/funchandlers.go @@ -0,0 +1,198 @@ +/* +Copyright 2018 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package poolmgr + +import ( + "context" + + "go.uber.org/zap" + apiv1 "k8s.io/api/core/v1" + kerrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/util/intstr" + "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/utils" +) + +func getIstioServiceLabels(fnName string) map[string]string { + return map[string]string{ + "functionName": fnName, + } +} + +func (gpm *GenericPoolManager) FunctionEventHandlers(kubernetesClient *kubernetes.Clientset, fissionfnNamespace string, istioEnabled bool) k8sCache.ResourceEventHandlerFuncs { + return k8sCache.ResourceEventHandlerFuncs{ + AddFunc: func(obj interface{}) { + fn := obj.(*fv1.Function) + + // Since istio only allows accessing pod through k8s service, + // for the functions with executor type "poolmgr" we need to + // create a service for sending requests to pod in pool. + // Functions with executor type "Newdeploy" is specialized at + // pod starts. In this case, just ignore such functions. + fnExecutorType := fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType + + // In some cases, user may not enter the executorType explicitly, for example in his spec.yaml. + // we assume it to be of type poolmgr + if fnExecutorType != "" && fnExecutorType != fv1.ExecutorTypePoolmgr { + return + } + + // create or update role-binding + envNs := fissionfnNamespace + if fn.Spec.Environment.Namespace != metav1.NamespaceDefault { + envNs = fn.Spec.Environment.Namespace + } + + // TODO : Just bring to your attention during review : + // setup rolebinding is tried, if it fails, we don't return. we just log an error and move on, because : + // 1. not all functions have secrets and/or configmaps, so things will work without this rolebinding in that case. + // 2. on the contrary, when the route is tried, the env fetcher logs will show a 403 forbidden message and same will be relayed to executor. + err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.SecretConfigMapGetterRB, fn.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs) + if err != nil { + gpm.logger.Error("error creating rolebinding", zap.Error(err), zap.String("role_binding", fv1.SecretConfigMapGetterRB)) + } else { + gpm.logger.Debug("successfully set up rolebinding for fetcher service account for function", + zap.String("service_account", fv1.FissionFetcherSA), + zap.String("service_account_namepsace", envNs), + zap.String("function_name", fn.ObjectMeta.Name), + zap.String("function_namespace", fn.ObjectMeta.Namespace)) + } + + if istioEnabled { + // create a same name service for function + // since istio only allows the traffic to service + sel := map[string]string{ + "functionName": fn.ObjectMeta.Name, + "functionUid": string(fn.ObjectMeta.UID), + } + + svcName := utils.GetFunctionIstioServiceName(fn.ObjectMeta.Name, fn.ObjectMeta.Namespace) + + // service for accepting user traffic + svc := apiv1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Namespace: envNs, + Name: svcName, + Labels: getIstioServiceLabels(fn.ObjectMeta.Name), + }, + Spec: apiv1.ServiceSpec{ + Type: apiv1.ServiceTypeClusterIP, + Ports: []apiv1.ServicePort{ + // Service port name should begin with a recognized prefix, or the traffic will be + // treated as TCP traffic. (https://istio.io/docs/setup/kubernetes/additional-setup/requirements/) + { + Name: "http-fetcher", + Protocol: apiv1.ProtocolTCP, + Port: 8000, + TargetPort: intstr.FromInt(8000), + }, + { + Name: "http-env", + Protocol: apiv1.ProtocolTCP, + Port: 8888, + TargetPort: intstr.FromInt(8888), + }, + }, + Selector: sel, + }, + } + + // create function istio service if it does not exist + _, err = kubernetesClient.CoreV1().Services(envNs).Create(context.TODO(), &svc, metav1.CreateOptions{}) + if err != nil && !kerrors.IsAlreadyExists(err) { + gpm.logger.Error("error creating istio service for function", + zap.Error(err), + zap.String("service_name", svcName), + zap.String("function_name", fn.ObjectMeta.Name), + zap.Any("selectors", sel)) + } + } + }, + + DeleteFunc: func(obj interface{}) { + fn := obj.(*fv1.Function) + + fnExecutorType := fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType + if fnExecutorType != "" && fnExecutorType != fv1.ExecutorTypePoolmgr { + return + } + + envNs := fissionfnNamespace + if fn.Spec.Environment.Namespace != metav1.NamespaceDefault { + envNs = fn.Spec.Environment.Namespace + } + + if istioEnabled { + svcName := utils.GetFunctionIstioServiceName(fn.ObjectMeta.Name, fn.ObjectMeta.Namespace) + // delete function istio service + err := kubernetesClient.CoreV1().Services(envNs).Delete(context.TODO(), svcName, metav1.DeleteOptions{}) + if err != nil && !kerrors.IsNotFound(err) { + gpm.logger.Error("error deleting istio service for function", + zap.Error(err), + zap.String("service_name", svcName), + zap.String("function_name", fn.ObjectMeta.Name)) + + } + } + }, + + UpdateFunc: func(oldObj, newObj interface{}) { + oldFunc := oldObj.(*fv1.Function) + newFunc := newObj.(*fv1.Function) + + if oldFunc.ObjectMeta.ResourceVersion == newFunc.ObjectMeta.ResourceVersion { + return + } + + envChanged := (oldFunc.Spec.Environment.Namespace != newFunc.Spec.Environment.Namespace) + + executorTypeChangedToPM := (oldFunc.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType != fv1.ExecutorTypePoolmgr && + newFunc.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType == fv1.ExecutorTypePoolmgr) + + // if a func's env reference gets updated and the newly referenced env is in a different ns, + // we need to create a rolebinding in func's ns so that the fetcher-sa in env ns has access + // to fetch secrets and config maps from the func's ns. + // similarly if executorType changed to Pool Manager, we now need a rolebinding in the func ns for fetcher sa + // present in env ns because for newdeploy, the fetcher sa is in function namespace + if envChanged || executorTypeChangedToPM { + envNs := fissionfnNamespace + if newFunc.Spec.Environment.Namespace != metav1.NamespaceDefault { + envNs = newFunc.Spec.Environment.Namespace + } + + err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.SecretConfigMapGetterRB, + newFunc.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, + fv1.FissionFetcherSA, envNs) + + if err != nil { + gpm.logger.Error("error creating rolebinding", zap.Error(err), zap.String("role_binding", fv1.SecretConfigMapGetterRB)) + } else { + gpm.logger.Debug("successfully set up rolebinding for fetcher service account for function", + zap.String("service_account", fv1.FissionFetcherSA), + zap.String("service_account_namepsace", envNs), + zap.String("function_name", newFunc.ObjectMeta.Name), + zap.String("function_namespace", newFunc.ObjectMeta.Namespace)) + } + } + }, + } + +} diff --git a/pkg/executor/executortype/poolmgr/funcwatcher.go b/pkg/executor/executortype/poolmgr/funcwatcher.go deleted file mode 100644 index 335ff386..00000000 --- a/pkg/executor/executortype/poolmgr/funcwatcher.go +++ /dev/null @@ -1,208 +0,0 @@ -/* -Copyright 2018 The Fission Authors. - -Licensed under the Apache License, Version 2.0 (the "License"); -you may not use this file except in compliance with the License. -You may obtain a copy of the License at - - http://www.apache.org/licenses/LICENSE-2.0 - -Unless required by applicable law or agreed to in writing, software -distributed under the License is distributed on an "AS IS" BASIS, -WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -See the License for the specific language governing permissions and -limitations under the License. -*/ - -package poolmgr - -import ( - "context" - "time" - - "go.uber.org/zap" - apiv1 "k8s.io/api/core/v1" - kerrors "k8s.io/apimachinery/pkg/api/errors" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/fields" - "k8s.io/apimachinery/pkg/util/intstr" - "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/utils" -) - -func getIstioServiceLabels(fnName string) map[string]string { - return map[string]string{ - "functionName": fnName, - } -} - -func (gpm *GenericPoolManager) makeFuncController(fissionClient *crd.FissionClient, - kubernetesClient *kubernetes.Clientset, fissionfnNamespace string, istioEnabled bool) (k8sCache.Store, k8sCache.Controller) { - - resyncPeriod := 30 * time.Second - lw := k8sCache.NewListWatchFromClient(fissionClient.CoreV1().RESTClient(), "functions", metav1.NamespaceAll, fields.Everything()) - - funcStore, controller := k8sCache.NewInformer(lw, &fv1.Function{}, resyncPeriod, - k8sCache.ResourceEventHandlerFuncs{ - AddFunc: func(obj interface{}) { - fn := obj.(*fv1.Function) - - // Since istio only allows accessing pod through k8s service, - // for the functions with executor type "poolmgr" we need to - // create a service for sending requests to pod in pool. - // Functions with executor type "Newdeploy" is specialized at - // pod starts. In this case, just ignore such functions. - fnExecutorType := fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType - - // In some cases, user may not enter the executorType explicitly, for example in his spec.yaml. - // we assume it to be of type poolmgr - if fnExecutorType != "" && fnExecutorType != fv1.ExecutorTypePoolmgr { - return - } - - // create or update role-binding - envNs := fissionfnNamespace - if fn.Spec.Environment.Namespace != metav1.NamespaceDefault { - envNs = fn.Spec.Environment.Namespace - } - - // TODO : Just bring to your attention during review : - // setup rolebinding is tried, if it fails, we don't return. we just log an error and move on, because : - // 1. not all functions have secrets and/or configmaps, so things will work without this rolebinding in that case. - // 2. on the contrary, when the route is tried, the env fetcher logs will show a 403 forbidden message and same will be relayed to executor. - err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.SecretConfigMapGetterRB, fn.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs) - if err != nil { - gpm.logger.Error("error creating rolebinding", zap.Error(err), zap.String("role_binding", fv1.SecretConfigMapGetterRB)) - } else { - gpm.logger.Debug("successfully set up rolebinding for fetcher service account for function", - zap.String("service_account", fv1.FissionFetcherSA), - zap.String("service_account_namepsace", envNs), - zap.String("function_name", fn.ObjectMeta.Name), - zap.String("function_namespace", fn.ObjectMeta.Namespace)) - } - - if istioEnabled { - // create a same name service for function - // since istio only allows the traffic to service - sel := map[string]string{ - "functionName": fn.ObjectMeta.Name, - "functionUid": string(fn.ObjectMeta.UID), - } - - svcName := utils.GetFunctionIstioServiceName(fn.ObjectMeta.Name, fn.ObjectMeta.Namespace) - - // service for accepting user traffic - svc := apiv1.Service{ - ObjectMeta: metav1.ObjectMeta{ - Namespace: envNs, - Name: svcName, - Labels: getIstioServiceLabels(fn.ObjectMeta.Name), - }, - Spec: apiv1.ServiceSpec{ - Type: apiv1.ServiceTypeClusterIP, - Ports: []apiv1.ServicePort{ - // Service port name should begin with a recognized prefix, or the traffic will be - // treated as TCP traffic. (https://istio.io/docs/setup/kubernetes/additional-setup/requirements/) - { - Name: "http-fetcher", - Protocol: apiv1.ProtocolTCP, - Port: 8000, - TargetPort: intstr.FromInt(8000), - }, - { - Name: "http-env", - Protocol: apiv1.ProtocolTCP, - Port: 8888, - TargetPort: intstr.FromInt(8888), - }, - }, - Selector: sel, - }, - } - - // create function istio service if it does not exist - _, err = kubernetesClient.CoreV1().Services(envNs).Create(context.TODO(), &svc, metav1.CreateOptions{}) - if err != nil && !kerrors.IsAlreadyExists(err) { - gpm.logger.Error("error creating istio service for function", - zap.Error(err), - zap.String("service_name", svcName), - zap.String("function_name", fn.ObjectMeta.Name), - zap.Any("selectors", sel)) - } - } - }, - - DeleteFunc: func(obj interface{}) { - fn := obj.(*fv1.Function) - - fnExecutorType := fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType - if fnExecutorType != "" && fnExecutorType != fv1.ExecutorTypePoolmgr { - return - } - - envNs := fissionfnNamespace - if fn.Spec.Environment.Namespace != metav1.NamespaceDefault { - envNs = fn.Spec.Environment.Namespace - } - - if istioEnabled { - svcName := utils.GetFunctionIstioServiceName(fn.ObjectMeta.Name, fn.ObjectMeta.Namespace) - // delete function istio service - err := kubernetesClient.CoreV1().Services(envNs).Delete(context.TODO(), svcName, metav1.DeleteOptions{}) - if err != nil && !kerrors.IsNotFound(err) { - gpm.logger.Error("error deleting istio service for function", - zap.Error(err), - zap.String("service_name", svcName), - zap.String("function_name", fn.ObjectMeta.Name)) - - } - } - }, - - UpdateFunc: func(oldObj, newObj interface{}) { - oldFunc := oldObj.(*fv1.Function) - newFunc := newObj.(*fv1.Function) - - if oldFunc.ObjectMeta.ResourceVersion == newFunc.ObjectMeta.ResourceVersion { - return - } - - envChanged := (oldFunc.Spec.Environment.Namespace != newFunc.Spec.Environment.Namespace) - - executorTypeChangedToPM := (oldFunc.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType != fv1.ExecutorTypePoolmgr && - newFunc.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType == fv1.ExecutorTypePoolmgr) - - // if a func's env reference gets updated and the newly referenced env is in a different ns, - // we need to create a rolebinding in func's ns so that the fetcher-sa in env ns has access - // to fetch secrets and config maps from the func's ns. - // similarly if executorType changed to Pool Manager, we now need a rolebinding in the func ns for fetcher sa - // present in env ns because for newdeploy, the fetcher sa is in function namespace - if envChanged || executorTypeChangedToPM { - envNs := fissionfnNamespace - if newFunc.Spec.Environment.Namespace != metav1.NamespaceDefault { - envNs = newFunc.Spec.Environment.Namespace - } - - err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.SecretConfigMapGetterRB, - newFunc.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, - fv1.FissionFetcherSA, envNs) - - if err != nil { - gpm.logger.Error("error creating rolebinding", zap.Error(err), zap.String("role_binding", fv1.SecretConfigMapGetterRB)) - } else { - gpm.logger.Debug("successfully set up rolebinding for fetcher service account for function", - zap.String("service_account", fv1.FissionFetcherSA), - zap.String("service_account_namepsace", envNs), - zap.String("function_name", newFunc.ObjectMeta.Name), - zap.String("function_namespace", newFunc.ObjectMeta.Namespace)) - } - } - }, - }) - - return funcStore, controller -} diff --git a/pkg/executor/executortype/poolmgr/gp.go b/pkg/executor/executortype/poolmgr/gp.go index 1c506ed2..d65dd923 100644 --- a/pkg/executor/executortype/poolmgr/gp.go +++ b/pkg/executor/executortype/poolmgr/gp.go @@ -71,8 +71,7 @@ type ( fissionClient *crd.FissionClient fetcherConfig *fetcherConfig.Config stopReadyPodControllerCh chan struct{} - readyPodController cache.Controller - readyPodIndexer cache.Indexer + readyPodInformer cache.SharedIndexInformer readyPodQueue workqueue.DelayingInterface poolInstanceID string // small random string to uniquify pod names instanceID string // poolmgr instance id @@ -249,7 +248,7 @@ func (gp *GenericPool) choosePod(newLabels map[string]string) (string, *apiv1.Po key = item.(string) gp.logger.Debug("got key from the queue", zap.String("key", key)) - obj, exists, err := gp.readyPodIndexer.GetByKey(key) + obj, exists, err := gp.readyPodInformer.GetIndexer().GetByKey(key) if err != nil { gp.logger.Error("fetching object from store failed", zap.String("key", key), zap.Error(err)) return "", nil, err diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index 9f98ffd8..fe00b4e0 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -76,11 +76,10 @@ type ( enableIstio bool fetcherConfig *fetcherConfig.Config - funcStore k8sCache.Store - funcController k8sCache.Controller - pkgStore k8sCache.Store - pkgController k8sCache.Controller - podInformer k8sCache.SharedIndexInformer + funcInformer *k8sCache.SharedIndexInformer + pkgInformer *k8sCache.SharedIndexInformer + + podInformer k8sCache.SharedIndexInformer defaultIdlePodReapTime time.Duration } @@ -103,7 +102,10 @@ func MakeGenericPoolManager( metricsClient *metricsclient.Clientset, functionNamespace string, fetcherConfig *fetcherConfig.Config, - instanceID string) (executortype.ExecutorType, error) { + instanceID string, + funcInformer *k8sCache.SharedIndexInformer, + pkgInformer *k8sCache.SharedIndexInformer, +) (executortype.ExecutorType, error) { gpmLogger := logger.Named("generic_pool_manager") @@ -120,6 +122,8 @@ func MakeGenericPoolManager( requestChannel: make(chan *request), defaultIdlePodReapTime: 2 * time.Minute, fetcherConfig: fetcherConfig, + funcInformer: funcInformer, + pkgInformer: pkgInformer, } go gpm.service() @@ -132,16 +136,14 @@ func MakeGenericPoolManager( gpm.enableIstio = istio } - gpm.funcStore, gpm.funcController = gpm.makeFuncController( - gpm.fissionClient, gpm.kubernetesClient, gpm.namespace, gpm.enableIstio) + (*gpm.funcInformer).AddEventHandler(gpm.FunctionEventHandlers(gpm.kubernetesClient, gpm.namespace, gpm.enableIstio)) + (*gpm.pkgInformer).AddEventHandler(gpm.PackageEventHandlers(gpm.kubernetesClient, gpm.namespace)) - gpm.pkgStore, gpm.pkgController = gpm.makePkgController(gpm.fissionClient, gpm.kubernetesClient, gpm.namespace) - - informerFactory, err := utils.GetInformerFacoryByExecutor(gpm.kubernetesClient, fv1.ExecutorTypePoolmgr) + kubeInformerFactory, err := utils.GetInformerFactoryByExecutor(gpm.kubernetesClient, fv1.ExecutorTypePoolmgr) if err != nil { return nil, err } - gpm.podInformer = informerFactory.Core().V1().Pods().Informer() + gpm.podInformer = kubeInformerFactory.Core().V1().Pods().Informer() return gpm, nil } @@ -149,8 +151,6 @@ func (gpm *GenericPoolManager) Run(ctx context.Context) { // eagerPoolCreator must run after CleanupOldExecutorObjects. // Otherwise, the poolmanager may wrongly delete the deployment. go gpm.eagerPoolCreator() - go gpm.funcController.Run(ctx.Done()) - go gpm.pkgController.Run(ctx.Done()) go gpm.podInformer.Run(ctx.Done()) go gpm.idleObjectReaper() } diff --git a/pkg/executor/executortype/poolmgr/packagehandlers.go b/pkg/executor/executortype/poolmgr/packagehandlers.go new file mode 100644 index 00000000..0c6c025a --- /dev/null +++ b/pkg/executor/executortype/poolmgr/packagehandlers.go @@ -0,0 +1,99 @@ +/* +Copyright 2018 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package poolmgr + +import ( + "go.uber.org/zap" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "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/utils" +) + +func (gpm *GenericPoolManager) PackageEventHandlers(kubernetesClient *kubernetes.Clientset, fissionfnNamespace string) k8sCache.ResourceEventHandlerFuncs { + return k8sCache.ResourceEventHandlerFuncs{ + AddFunc: func(obj interface{}) { + pkg := obj.(*fv1.Package) + gpm.logger.Debug("list watch for package reported a new package addition", + zap.String("package_name", pkg.ObjectMeta.Name), + zap.String("package_namespace", pkg.ObjectMeta.Namespace)) + + // create or update role-binding for fetcher sa in env ns to be able to get the pkg contents from pkg namespace + envNs := fissionfnNamespace + if pkg.Spec.Environment.Namespace != metav1.NamespaceDefault { + envNs = pkg.Spec.Environment.Namespace + } + + // 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(gpm.logger, kubernetesClient, fv1.PackageGetterRB, pkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs) + if err != nil { + gpm.logger.Error("error creating rolebinding for package", + zap.Error(err), + zap.String("role_binding", fv1.PackageGetterRB), + zap.String("package_name", pkg.ObjectMeta.Name), + zap.String("package_namespace", pkg.ObjectMeta.Namespace)) + return + } + + gpm.logger.Debug("successfully set up rolebinding for fetcher service account", + zap.String("service_account", fv1.FissionFetcherSA), + zap.String("service_account_namespace", envNs), + zap.String("package_name", pkg.ObjectMeta.Name), + zap.String("package_namespace", pkg.ObjectMeta.Namespace)) + }, + + UpdateFunc: func(oldObj, newObj interface{}) { + oldPkg := oldObj.(*fv1.Package) + newPkg := newObj.(*fv1.Package) + + if oldPkg.ObjectMeta.ResourceVersion == newPkg.ObjectMeta.ResourceVersion { + return + } + + // if a pkg's env reference gets updated and the newly referenced env is in a different ns, + // we need to update the role-binding in pkg ns to grant permissions to the fetcher-sa in env ns + // to do a get on pkg + if oldPkg.Spec.Environment.Namespace != newPkg.Spec.Environment.Namespace { + envNs := fissionfnNamespace + if newPkg.Spec.Environment.Namespace != metav1.NamespaceDefault { + envNs = newPkg.Spec.Environment.Namespace + } + + err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.PackageGetterRB, + newPkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, + fv1.FissionFetcherSA, envNs) + if err != nil { + gpm.logger.Error("error updating rolebinding for package", + zap.Error(err), + zap.String("role_binding", fv1.PackageGetterRB), + zap.String("package_name", newPkg.ObjectMeta.Name), + zap.String("package_namespace", newPkg.ObjectMeta.Namespace)) + return + } + + gpm.logger.Debug("successfully updated rolebinding for fetcher service account", + zap.String("service_account", fv1.FissionFetcherSA), + zap.String("service_account_namespace", envNs), + zap.String("package_name", newPkg.ObjectMeta.Name), + zap.String("package_namespace", newPkg.ObjectMeta.Namespace)) + } + }, + } +} diff --git a/pkg/executor/executortype/poolmgr/packagewatcher.go b/pkg/executor/executortype/poolmgr/packagewatcher.go deleted file mode 100644 index 296304be..00000000 --- a/pkg/executor/executortype/poolmgr/packagewatcher.go +++ /dev/null @@ -1,111 +0,0 @@ -/* -Copyright 2018 The Fission Authors. - -Licensed under the Apache License, Version 2.0 (the "License"); -you may not use this file except in compliance with the License. -You may obtain a copy of the License at - - http://www.apache.org/licenses/LICENSE-2.0 - -Unless required by applicable law or agreed to in writing, software -distributed under the License is distributed on an "AS IS" BASIS, -WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -See the License for the specific language governing permissions and -limitations under the License. -*/ - -package poolmgr - -import ( - "time" - - "go.uber.org/zap" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/fields" - "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/utils" -) - -// TODO : It may make sense to make each of add, update, delete funcs run as separate go routines. -func (gpm *GenericPoolManager) makePkgController(fissionClient *crd.FissionClient, - kubernetesClient *kubernetes.Clientset, fissionfnNamespace string) (k8sCache.Store, k8sCache.Controller) { - - resyncPeriod := 30 * time.Second - lw := k8sCache.NewListWatchFromClient(fissionClient.CoreV1().RESTClient(), "packages", metav1.NamespaceAll, fields.Everything()) - pkgStore, controller := k8sCache.NewInformer(lw, &fv1.Package{}, resyncPeriod, - k8sCache.ResourceEventHandlerFuncs{ - AddFunc: func(obj interface{}) { - pkg := obj.(*fv1.Package) - gpm.logger.Debug("list watch for package reported a new package addition", - zap.String("package_name", pkg.ObjectMeta.Name), - zap.String("package_namespace", pkg.ObjectMeta.Namespace)) - - // create or update role-binding for fetcher sa in env ns to be able to get the pkg contents from pkg namespace - envNs := fissionfnNamespace - if pkg.Spec.Environment.Namespace != metav1.NamespaceDefault { - envNs = pkg.Spec.Environment.Namespace - } - - // 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(gpm.logger, kubernetesClient, fv1.PackageGetterRB, pkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs) - if err != nil { - gpm.logger.Error("error creating rolebinding for package", - zap.Error(err), - zap.String("role_binding", fv1.PackageGetterRB), - zap.String("package_name", pkg.ObjectMeta.Name), - zap.String("package_namespace", pkg.ObjectMeta.Namespace)) - return - } - - gpm.logger.Debug("successfully set up rolebinding for fetcher service account", - zap.String("service_account", fv1.FissionFetcherSA), - zap.String("service_account_namespace", envNs), - zap.String("package_name", pkg.ObjectMeta.Name), - zap.String("package_namespace", pkg.ObjectMeta.Namespace)) - }, - - UpdateFunc: func(oldObj, newObj interface{}) { - oldPkg := oldObj.(*fv1.Package) - newPkg := newObj.(*fv1.Package) - - if oldPkg.ObjectMeta.ResourceVersion == newPkg.ObjectMeta.ResourceVersion { - return - } - - // if a pkg's env reference gets updated and the newly referenced env is in a different ns, - // we need to update the role-binding in pkg ns to grant permissions to the fetcher-sa in env ns - // to do a get on pkg - if oldPkg.Spec.Environment.Namespace != newPkg.Spec.Environment.Namespace { - envNs := fissionfnNamespace - if newPkg.Spec.Environment.Namespace != metav1.NamespaceDefault { - envNs = newPkg.Spec.Environment.Namespace - } - - err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.PackageGetterRB, - newPkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, - fv1.FissionFetcherSA, envNs) - if err != nil { - gpm.logger.Error("error updating rolebinding for package", - zap.Error(err), - zap.String("role_binding", fv1.PackageGetterRB), - zap.String("package_name", newPkg.ObjectMeta.Name), - zap.String("package_namespace", newPkg.ObjectMeta.Namespace)) - return - } - - gpm.logger.Debug("successfully updated rolebinding for fetcher service account", - zap.String("service_account", fv1.FissionFetcherSA), - zap.String("service_account_namespace", envNs), - zap.String("package_name", newPkg.ObjectMeta.Name), - zap.String("package_namespace", newPkg.ObjectMeta.Namespace)) - } - }, - }) - - return pkgStore, controller -} diff --git a/pkg/executor/executortype/poolmgr/readyPodController.go b/pkg/executor/executortype/poolmgr/readyPodController.go index d1ef8d14..00ceca52 100644 --- a/pkg/executor/executortype/poolmgr/readyPodController.go +++ b/pkg/executor/executortype/poolmgr/readyPodController.go @@ -4,27 +4,31 @@ import ( "time" "go.uber.org/zap" - apiv1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" + informers "k8s.io/client-go/informers/core/v1" "k8s.io/client-go/tools/cache" "k8s.io/client-go/util/workqueue" ) +func (gp *GenericPool) newPodInformer() cache.SharedIndexInformer { + optionsModifier := func(options *metav1.ListOptions) { + options.LabelSelector = labels.Set( + gp.deployment.Spec.Selector.MatchLabels).AsSelector().String() + options.FieldSelector = "status.phase=Running" + } + return informers.NewFilteredPodInformer(gp.kubernetesClient, gp.namespace, 0, nil, optionsModifier) +} + func (gp *GenericPool) startReadyPodController() { // create the pod watcher to filter by labels // Filtering pod by phase=Running. In some cases the pod can be in // different state than Running, for example Kubernetes sets a // pod to Termination while k8s waits for the grace period of // the pod, even if all the containers are in Ready state. - optionsModifier := func(options *metav1.ListOptions) { - options.LabelSelector = labels.Set( - gp.deployment.Spec.Selector.MatchLabels).AsSelector().String() - options.FieldSelector = "status.phase=Running" - } - readyPodWatcher := cache.NewFilteredListWatchFromClient(gp.kubernetesClient.CoreV1().RESTClient(), "pods", gp.namespace, optionsModifier) gp.readyPodQueue = workqueue.NewDelayingQueue() - gp.readyPodIndexer, gp.readyPodController = cache.NewIndexerInformer(readyPodWatcher, &apiv1.Pod{}, 0, cache.ResourceEventHandlerFuncs{ + gp.readyPodInformer = gp.newPodInformer() + gp.readyPodInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { key, err := cache.MetaNamespaceKeyFunc(obj) if err == nil { @@ -39,7 +43,7 @@ func (gp *GenericPool) startReadyPodController() { gp.logger.Debug("delete func called", zap.String("key", key)) } }, - }, cache.Indexers{}) - go gp.readyPodController.Run(gp.stopReadyPodControllerCh) + }) + go gp.readyPodInformer.Run(gp.stopReadyPodControllerCh) gp.logger.Info("readyPod controller started", zap.String("env", gp.env.ObjectMeta.Name), zap.String("envID", string(gp.env.ObjectMeta.UID))) } diff --git a/pkg/router/functionReferenceResolver.go b/pkg/router/functionReferenceResolver.go index d13f228c..6e691548 100644 --- a/pkg/router/functionReferenceResolver.go +++ b/pkg/router/functionReferenceResolver.go @@ -33,8 +33,9 @@ type ( // reference into a resolveResult functionReferenceResolver struct { // FunctionReference -> function metadata - refCache *cache.Cache - store k8sCache.Store + refCache *cache.Cache + funcInformer *k8sCache.SharedIndexInformer + // store k8sCache.Store } resolveResultType int @@ -68,10 +69,10 @@ const ( resolveResultMultipleFunctions ) -func makeFunctionReferenceResolver(store k8sCache.Store) *functionReferenceResolver { +func makeFunctionReferenceResolver(funcInformer *k8sCache.SharedIndexInformer) *functionReferenceResolver { frr := &functionReferenceResolver{ - refCache: cache.MakeCache(time.Minute, 0), - store: store, + refCache: cache.MakeCache(time.Minute, 0), + funcInformer: funcInformer, } return frr } @@ -120,7 +121,7 @@ func (frr *functionReferenceResolver) resolve(trigger fv1.HTTPTrigger) (*resolve // resolveByName simply looks up function by name in a namespace. func (frr *functionReferenceResolver) resolveByName(namespace, name string) (*resolveResult, error) { // get function from cache - obj, isExist, err := frr.store.Get(&fv1.Function{ + obj, isExist, err := (*frr.funcInformer).GetStore().Get(&fv1.Function{ ObjectMeta: metav1.ObjectMeta{ Namespace: namespace, Name: name, @@ -154,7 +155,7 @@ func (frr *functionReferenceResolver) resolveByFunctionWeights(namespace string, for functionName, functionWeight := range fr.FunctionWeights { // get function from cache - obj, isExist, err := frr.store.Get(&fv1.Function{ + obj, isExist, err := (*frr.funcInformer).GetStore().Get(&fv1.Function{ ObjectMeta: metav1.ObjectMeta{ Namespace: namespace, Name: functionName, diff --git a/pkg/router/httpTriggers.go b/pkg/router/httpTriggers.go index 17f3846f..777cdbfd 100644 --- a/pkg/router/httpTriggers.go +++ b/pkg/router/httpTriggers.go @@ -23,8 +23,6 @@ import ( "github.com/gorilla/mux" "go.uber.org/zap" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/fields" "k8s.io/apimachinery/pkg/types" "k8s.io/client-go/kubernetes" "k8s.io/client-go/rest" @@ -33,6 +31,7 @@ import ( fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/crd" executorClient "github.com/fission/fission/pkg/executor/client" + genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" "github.com/fission/fission/pkg/throttler" "github.com/fission/fission/pkg/utils" ) @@ -49,11 +48,9 @@ type HTTPTriggerSet struct { resolver *functionReferenceResolver crdClient rest.Interface triggers []fv1.HTTPTrigger - triggerStore k8sCache.Store - triggerController k8sCache.Controller + triggerInformer k8sCache.SharedIndexInformer functions []fv1.Function - funcStore k8sCache.Store - funcController k8sCache.Controller + funcInformer k8sCache.SharedIndexInformer updateRouterRequestChannel chan struct{} tsRoundTripperParams *tsRoundTripperParams isDebugEnv bool @@ -62,7 +59,7 @@ type HTTPTriggerSet struct { } func makeHTTPTriggerSet(logger *zap.Logger, fmap *functionServiceMap, fissionClient *crd.FissionClient, - kubeClient *kubernetes.Clientset, executor *executorClient.Client, crdClient rest.Interface, params *tsRoundTripperParams, isDebugEnv bool, unTapServiceTimeout time.Duration, actionThrottler *throttler.Throttler) (*HTTPTriggerSet, k8sCache.Store, k8sCache.Store) { + kubeClient *kubernetes.Clientset, executor *executorClient.Client, crdClient rest.Interface, params *tsRoundTripperParams, isDebugEnv bool, unTapServiceTimeout time.Duration, actionThrottler *throttler.Throttler) *HTTPTriggerSet { httpTriggerSet := &HTTPTriggerSet{ logger: logger.Named("http_trigger_set"), @@ -78,21 +75,19 @@ func makeHTTPTriggerSet(logger *zap.Logger, fmap *functionServiceMap, fissionCli svcAddrUpdateThrottler: actionThrottler, unTapServiceTimeout: unTapServiceTimeout, } - var tStore, fnStore k8sCache.Store - var tController, fnController k8sCache.Controller - if httpTriggerSet.crdClient != nil { - tStore, tController = httpTriggerSet.initTriggerController() - httpTriggerSet.triggerStore = tStore - httpTriggerSet.triggerController = tController - fnStore, fnController = httpTriggerSet.initFunctionController() - httpTriggerSet.funcStore = fnStore - httpTriggerSet.funcController = fnController + informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Second*30) + httpTriggerSet.triggerInformer = informerFactory.Core().V1().HTTPTriggers().Informer() + httpTriggerSet.funcInformer = informerFactory.Core().V1().Functions().Informer() + + httpTriggerSet.addTriggerHandlers() + httpTriggerSet.addFunctionHandlers() } - return httpTriggerSet, tStore, fnStore + return httpTriggerSet } -func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mr *mutableRouter, resolver *functionReferenceResolver) { +func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mr *mutableRouter) { + resolver := makeFunctionReferenceResolver(&ts.funcInformer) ts.resolver = resolver ts.mutableRouter = mr @@ -104,8 +99,8 @@ func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mr *mutableRouter } go ts.updateRouter() go ts.syncTriggers() - go ts.runWatcher(ctx, ts.funcController) - go ts.runWatcher(ctx, ts.triggerController) + go ts.runInformer(ctx, ts.funcInformer) + go ts.runInformer(ctx, ts.triggerInformer) } func defaultHomeHandler(w http.ResponseWriter, r *http.Request) { @@ -234,78 +229,70 @@ func (ts *HTTPTriggerSet) updateTriggerStatusFailed(ht *fv1.HTTPTrigger, err err // TODO } -func (ts *HTTPTriggerSet) initTriggerController() (k8sCache.Store, k8sCache.Controller) { - resyncPeriod := 30 * time.Second - listWatch := k8sCache.NewListWatchFromClient(ts.crdClient, "httptriggers", metav1.NamespaceAll, fields.Everything()) - store, controller := k8sCache.NewInformer(listWatch, &fv1.HTTPTrigger{}, resyncPeriod, - k8sCache.ResourceEventHandlerFuncs{ - AddFunc: func(obj interface{}) { - trigger := obj.(*fv1.HTTPTrigger) - go createIngress(ts.logger, trigger, ts.kubeClient) - ts.syncTriggers() - }, - DeleteFunc: func(obj interface{}) { - ts.syncTriggers() - trigger := obj.(*fv1.HTTPTrigger) - go deleteIngress(ts.logger, trigger, ts.kubeClient) - }, - UpdateFunc: func(oldObj interface{}, newObj interface{}) { - oldTrigger := oldObj.(*fv1.HTTPTrigger) - newTrigger := newObj.(*fv1.HTTPTrigger) +func (ts *HTTPTriggerSet) addTriggerHandlers() { + ts.triggerInformer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ + AddFunc: func(obj interface{}) { + trigger := obj.(*fv1.HTTPTrigger) + go createIngress(ts.logger, trigger, ts.kubeClient) + ts.syncTriggers() + }, + DeleteFunc: func(obj interface{}) { + ts.syncTriggers() + trigger := obj.(*fv1.HTTPTrigger) + go deleteIngress(ts.logger, trigger, ts.kubeClient) + }, + UpdateFunc: func(oldObj interface{}, newObj interface{}) { + oldTrigger := oldObj.(*fv1.HTTPTrigger) + newTrigger := newObj.(*fv1.HTTPTrigger) - if oldTrigger.ObjectMeta.ResourceVersion == newTrigger.ObjectMeta.ResourceVersion { - return - } + if oldTrigger.ObjectMeta.ResourceVersion == newTrigger.ObjectMeta.ResourceVersion { + return + } - go updateIngress(ts.logger, oldTrigger, newTrigger, ts.kubeClient) - ts.syncTriggers() - }, - }) - return store, controller + go updateIngress(ts.logger, oldTrigger, newTrigger, ts.kubeClient) + ts.syncTriggers() + }, + }) } -func (ts *HTTPTriggerSet) initFunctionController() (k8sCache.Store, k8sCache.Controller) { - resyncPeriod := 30 * time.Second - listWatch := k8sCache.NewListWatchFromClient(ts.crdClient, "functions", metav1.NamespaceAll, fields.Everything()) - store, controller := k8sCache.NewInformer(listWatch, &fv1.Function{}, resyncPeriod, - k8sCache.ResourceEventHandlerFuncs{ - AddFunc: func(obj interface{}) { - ts.syncTriggers() - }, - DeleteFunc: func(obj interface{}) { - ts.syncTriggers() - }, - UpdateFunc: func(oldObj interface{}, newObj interface{}) { - oldFn := oldObj.(*fv1.Function) - fn := newObj.(*fv1.Function) +func (ts *HTTPTriggerSet) addFunctionHandlers() { + ts.funcInformer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ + AddFunc: func(obj interface{}) { + ts.syncTriggers() + }, + DeleteFunc: func(obj interface{}) { + ts.syncTriggers() + }, + UpdateFunc: func(oldObj interface{}, newObj interface{}) { + oldFn := oldObj.(*fv1.Function) + fn := newObj.(*fv1.Function) - if oldFn.ObjectMeta.ResourceVersion == fn.ObjectMeta.ResourceVersion { - return - } + if oldFn.ObjectMeta.ResourceVersion == fn.ObjectMeta.ResourceVersion { + return + } - // update resolver function reference cache - for key, rr := range ts.resolver.copy() { - if key.namespace == fn.ObjectMeta.Namespace && - rr.functionMap[fn.ObjectMeta.Name] != nil && - rr.functionMap[fn.ObjectMeta.Name].ObjectMeta.ResourceVersion != fn.ObjectMeta.ResourceVersion { - // invalidate resolver cache - ts.logger.Debug("invalidating resolver cache") - err := ts.resolver.delete(key.namespace, key.triggerName, key.triggerResourceVersion) - if err != nil { - ts.logger.Error("error deleting functionReferenceResolver cache", zap.Error(err)) - } - break + // update resolver function reference cache + for key, rr := range ts.resolver.copy() { + if key.namespace == fn.ObjectMeta.Namespace && + rr.functionMap[fn.ObjectMeta.Name] != nil && + rr.functionMap[fn.ObjectMeta.Name].ObjectMeta.ResourceVersion != fn.ObjectMeta.ResourceVersion { + // invalidate resolver cache + ts.logger.Debug("invalidating resolver cache") + err := ts.resolver.delete(key.namespace, key.triggerName, key.triggerResourceVersion) + if err != nil { + ts.logger.Error("error deleting functionReferenceResolver cache", zap.Error(err)) } + break } - ts.syncTriggers() - }, - }) - return store, controller + } + ts.syncTriggers() + }, + }) } -func (ts *HTTPTriggerSet) runWatcher(ctx context.Context, controller k8sCache.Controller) { +func (ts *HTTPTriggerSet) runInformer(ctx context.Context, informer k8sCache.SharedIndexInformer) { go func() { - controller.Run(ctx.Done()) + informer.Run(ctx.Done()) }() } @@ -316,7 +303,7 @@ func (ts *HTTPTriggerSet) syncTriggers() { func (ts *HTTPTriggerSet) updateRouter() { for range ts.updateRouterRequestChannel { // get triggers - latestTriggers := ts.triggerStore.List() + latestTriggers := ts.triggerInformer.GetStore().List() triggers := make([]fv1.HTTPTrigger, len(latestTriggers)) for _, t := range latestTriggers { triggers = append(triggers, *t.(*fv1.HTTPTrigger)) @@ -324,7 +311,7 @@ func (ts *HTTPTriggerSet) updateRouter() { ts.triggers = triggers // get functions - latestFunctions := ts.funcStore.List() + latestFunctions := ts.funcInformer.GetStore().List() functionTimeout := make(map[types.UID]int, len(latestFunctions)) functions := make([]fv1.Function, len(latestFunctions)) for _, f := range latestFunctions { diff --git a/pkg/router/router.go b/pkg/router/router.go index 0cb97981..3a9b117e 100644 --- a/pkg/router/router.go +++ b/pkg/router/router.go @@ -63,7 +63,7 @@ import ( // request url ---[trigger]---> Function(name, deployment) ----[deployment]----> Function(name, uid) ----[pool mgr]---> k8s service url -func router(ctx context.Context, logger *zap.Logger, httpTriggerSet *HTTPTriggerSet, resolver *functionReferenceResolver) *mutableRouter { +func router(ctx context.Context, logger *zap.Logger, httpTriggerSet *HTTPTriggerSet) *mutableRouter { var mr *mutableRouter // see issue https://github.com/fission/fission/issues/1317 @@ -74,13 +74,13 @@ func router(ctx context.Context, logger *zap.Logger, httpTriggerSet *HTTPTrigger mr = newMutableRouter(logger, mux.NewRouter()) } - httpTriggerSet.subscribeRouter(ctx, mr, resolver) + httpTriggerSet.subscribeRouter(ctx, mr) return mr } func serve(ctx context.Context, logger *zap.Logger, port int, tracingSamplingRate float64, - httpTriggerSet *HTTPTriggerSet, resolver *functionReferenceResolver, displayAccessLog bool) { - mr := router(ctx, logger, httpTriggerSet, resolver) + httpTriggerSet *HTTPTriggerSet, displayAccessLog bool) { + mr := router(ctx, logger, httpTriggerSet) url := fmt.Sprintf(":%v", port) err := http.ListenAndServe(url, &ochttp.Handler{ @@ -236,7 +236,7 @@ func Start(logger *zap.Logger, port int, executorURL string) { zap.Bool("default", displayAccessLog)) } - triggers, _, fnStore := makeHTTPTriggerSet(logger.Named("triggerset"), fmap, fissionClient, kubeClient, executor, fissionClient.CoreV1().RESTClient(), &tsRoundTripperParams{ + triggers := makeHTTPTriggerSet(logger.Named("triggerset"), fmap, fissionClient, kubeClient, executor, fissionClient.CoreV1().RESTClient(), &tsRoundTripperParams{ timeout: timeout, timeoutExponent: timeoutExponent, disableKeepAlive: disableKeepAlive, @@ -245,12 +245,10 @@ func Start(logger *zap.Logger, port int, executorURL string) { svcAddrRetryCount: svcAddrRetryCount, }, isDebugEnv, unTapServiceTimeout, throttler.MakeThrottler(svcAddrUpdateTimeout)) - resolver := makeFunctionReferenceResolver(fnStore) - go serveMetric(logger) logger.Info("starting router", zap.Int("port", port)) ctx, cancel := context.WithCancel(context.Background()) defer cancel() - serve(ctx, logger, port, tracingSamplingRate, triggers, resolver, displayAccessLog) + serve(ctx, logger, port, tracingSamplingRate, triggers, displayAccessLog) } diff --git a/pkg/utils/informer.go b/pkg/utils/informer.go index 8444fa47..cceeed88 100644 --- a/pkg/utils/informer.go +++ b/pkg/utils/informer.go @@ -1,17 +1,21 @@ package utils import ( + "fmt" + + apiv1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/selection" k8sInformers "k8s.io/client-go/informers" "k8s.io/client-go/kubernetes" + k8sCache "k8s.io/client-go/tools/cache" metricsapi "k8s.io/metrics/pkg/apis/metrics" v1 "github.com/fission/fission/pkg/apis/core/v1" ) -func GetInformerFacoryByExecutor(client *kubernetes.Clientset, executorType v1.ExecutorType) (k8sInformers.SharedInformerFactory, error) { +func GetInformerFactoryByExecutor(client *kubernetes.Clientset, executorType v1.ExecutorType) (k8sInformers.SharedInformerFactory, error) { executorLabel, err := labels.NewRequirement(v1.EXECUTOR_TYPE, selection.DoubleEquals, []string{string(executorType)}) if err != nil { return nil, err @@ -43,3 +47,14 @@ func SupportedMetricsAPIVersionAvailable(discoveredAPIGroups *metav1.APIGroupLis } return false } + +func GetCachedItem(obj apiv1.ObjectReference, informer k8sCache.SharedIndexInformer) (item interface{}, exists bool, err error) { + store := informer.GetStore() + + item, exists, err = store.Get(obj) + if err != nil || !exists { + item, exists, err = store.GetByKey(fmt.Sprintf("%s/%s", obj.Namespace, obj.Name)) + } + + return item, exists, err +} diff --git a/test/kind_CI.sh b/test/kind_CI.sh index a7237931..21856305 100755 --- a/test/kind_CI.sh +++ b/test/kind_CI.sh @@ -84,7 +84,7 @@ main() { $ROOT/test/tests/mqtrigger/nats/test_mqtrigger.sh \ $ROOT/test/tests/mqtrigger/nats/test_mqtrigger_error.sh \ $ROOT/test/tests/test_huge_response/test_huge_response.sh \ - $ROOT/test/tests/test_kubectl/test_kubectl.sh + $ROOT/test/tests/test_kubectl/test_kubectl.sh \ $ROOT/test/tests/websocket/test_ws.sh export JOBS=3