From 827baea974d9c54e9fb3508d80824490f4c2e5cf Mon Sep 17 00:00:00 2001 From: Sanket Sudake Date: Wed, 19 Oct 2022 15:48:47 +0530 Subject: [PATCH] Allow namespace configuration for different CRD resources in Fission (#2539) * Allow multiple namespaces for builder manager * Enable multiple namespaces for executor informers * Added missing context * helm chart support for multiple namespaces * Directly consume map type from GetInformerForNamespaces fn * Optimize function resolver by choosing namespace-specific informer * helm chart support for multiple namespaces * consider default namespace and move duplicate code to helm template * Improve documentation for fission namespace values Signed-off-by: Sanket Sudake Co-authored-by: shubham bansal --- charts/fission-all/templates/_helpers.tpl | 9 + .../templates/buildermgr/deployment.yaml | 1 + .../templates/controller/deployment.yaml | 1 + .../templates/executor/deployment.yaml | 1 + .../templates/kubewatcher/deployment.yaml | 1 + .../mqt-fission-kafka/deployment.yaml | 1 + .../templates/mqt-keda/deployment.yaml | 1 + .../templates/router/deployment.yaml | 1 + .../templates/storagesvc/deployment.yaml | 1 + .../templates/timer/deployment.yaml | 1 + charts/fission-all/values.yaml | 30 +++- pkg/apis/core/v1/const.go | 11 ++ pkg/buildermgr/buildermgr.go | 6 +- pkg/buildermgr/pkgwatcher.go | 18 +- pkg/canaryconfigmgr/canaryConfigMgr.go | 87 +++++----- pkg/executor/executor.go | 60 +++++-- .../executortype/container/containermgr.go | 6 +- .../executortype/newdeploy/newdeploymgr.go | 13 +- .../newdeploy/newdeploymgr_test.go | 14 +- pkg/executor/executortype/poolmgr/gpm.go | 12 +- .../executortype/poolmgr/poolpodcontroller.go | 77 ++++++--- .../poolmgr/poolpodcontroller_test.go | 19 ++- pkg/mqtrigger/mqtmanager.go | 14 +- pkg/mqtrigger/scalermanager.go | 13 +- pkg/router/functionReferenceResolver.go | 38 ++++- pkg/router/httpTriggers.go | 157 +++++++++--------- pkg/utils/informer.go | 56 +++++++ 27 files changed, 430 insertions(+), 219 deletions(-) diff --git a/charts/fission-all/templates/_helpers.tpl b/charts/fission-all/templates/_helpers.tpl index f7167e3f..ed32d600 100644 --- a/charts/fission-all/templates/_helpers.tpl +++ b/charts/fission-all/templates/_helpers.tpl @@ -71,3 +71,12 @@ This template generates the image name for the deployment depending on the value - name: OTEL_PROPAGATORS value: "{{ .Values.openTelemetry.propagators }}" {{- end }} + +{{- define "fission-resource-namespace.envs" }} +- name: FISSION_RESOURCE_NAMESPACES +{{- if not .Values.singleDefaultNamespace }} + value: "{{ .Values.defaultNamespace }},{{ join "," .Values.additionalFissionNamespaces }}" +{{- else }} + value: {{ .Values.defaultNamespace }} +{{- end }} +{{- end }} \ No newline at end of file diff --git a/charts/fission-all/templates/buildermgr/deployment.yaml b/charts/fission-all/templates/buildermgr/deployment.yaml index 1d67708a..f6b8f761 100644 --- a/charts/fission-all/templates/buildermgr/deployment.yaml +++ b/charts/fission-all/templates/buildermgr/deployment.yaml @@ -55,6 +55,7 @@ spec: value: {{ .Values.pprof.enabled | quote }} - name: HELM_RELEASE_NAME value: {{ .Release.Name | quote }} + {{- include "fission-resource-namespace.envs" . | indent 8 }} {{- include "opentelemtry.envs" . | indent 8 }} ports: - containerPort: 8080 diff --git a/charts/fission-all/templates/controller/deployment.yaml b/charts/fission-all/templates/controller/deployment.yaml index d5ee524d..b3fca557 100644 --- a/charts/fission-all/templates/controller/deployment.yaml +++ b/charts/fission-all/templates/controller/deployment.yaml @@ -38,6 +38,7 @@ spec: value: {{ .Values.debugEnv | quote }} - name: PPROF_ENABLED value: {{ .Values.pprof.enabled | quote }} + {{- include "fission-resource-namespace.envs" . | indent 8 }} - name: POD_NAMESPACE valueFrom: fieldRef: diff --git a/charts/fission-all/templates/executor/deployment.yaml b/charts/fission-all/templates/executor/deployment.yaml index b3f6e937..ed2f24c0 100644 --- a/charts/fission-all/templates/executor/deployment.yaml +++ b/charts/fission-all/templates/executor/deployment.yaml @@ -71,6 +71,7 @@ spec: - name: CONTAINER_OBJECT_REAPER_INTERVAL value: {{ .Values.executor.container.objectReaperInterval | quote }} {{- end}} + {{- include "fission-resource-namespace.envs" . | indent 8 }} - name: HELM_RELEASE_NAME value: {{ .Release.Name | quote }} {{- include "opentelemtry.envs" . | indent 8 }} diff --git a/charts/fission-all/templates/kubewatcher/deployment.yaml b/charts/fission-all/templates/kubewatcher/deployment.yaml index 2dc47352..ce72860e 100644 --- a/charts/fission-all/templates/kubewatcher/deployment.yaml +++ b/charts/fission-all/templates/kubewatcher/deployment.yaml @@ -29,6 +29,7 @@ spec: value: {{ .Values.debugEnv | quote }} - name: PPROF_ENABLED value: {{ .Values.pprof.enabled | quote }} + {{- include "fission-resource-namespace.envs" . | indent 8 }} {{- include "opentelemtry.envs" . | indent 8 }} resources: {{- toYaml .Values.kubewatcher.resources | nindent 10 }} diff --git a/charts/fission-all/templates/mqt-fission-kafka/deployment.yaml b/charts/fission-all/templates/mqt-fission-kafka/deployment.yaml index a8367c0f..89ad630a 100644 --- a/charts/fission-all/templates/mqt-fission-kafka/deployment.yaml +++ b/charts/fission-all/templates/mqt-fission-kafka/deployment.yaml @@ -47,6 +47,7 @@ spec: value: {{ .Values.debugEnv | quote }} - name: PPROF_ENABLED value: {{ .Values.pprof.enabled | quote }} + {{- include "fission-resource-namespace.envs" . | indent 8 }} {{- include "opentelemtry.envs" . | indent 8 }} # TLS authentication is TLS with authentication (2 way) # More info: https://docs.confluent.io/current/kafka/authentication_ssl.html#ssl-overview diff --git a/charts/fission-all/templates/mqt-keda/deployment.yaml b/charts/fission-all/templates/mqt-keda/deployment.yaml index 6a4aed5a..517b6217 100644 --- a/charts/fission-all/templates/mqt-keda/deployment.yaml +++ b/charts/fission-all/templates/mqt-keda/deployment.yaml @@ -46,6 +46,7 @@ spec: value: "{{ .Values.mqt_keda.connector_images.gcp_pubsub.image }}:{{ .Values.mqt_keda.connector_images.gcp_pubsub.tag }}" - name: REDIS_IMAGE value: "{{ .Values.mqt_keda.connector_images.redis.image }}:{{ .Values.mqt_keda.connector_images.redis.tag }}" + {{- include "fission-resource-namespace.envs" . | indent 8 }} {{- include "opentelemtry.envs" . | indent 8 }} resources: {{- toYaml .Values.mqt_keda.resources | nindent 10 }} diff --git a/charts/fission-all/templates/router/deployment.yaml b/charts/fission-all/templates/router/deployment.yaml index ff6bbffe..77b853d3 100644 --- a/charts/fission-all/templates/router/deployment.yaml +++ b/charts/fission-all/templates/router/deployment.yaml @@ -83,6 +83,7 @@ spec: value: {{ .Values.pprof.enabled | quote }} - name: DISPLAY_ACCESS_LOG value: {{ .Values.router.displayAccessLog | default false | quote }} + {{- include "fission-resource-namespace.envs" . | indent 8 }} {{- include "opentelemtry.envs" . | indent 8 }} resources: {{- toYaml .Values.router.resources | nindent 10 }} diff --git a/charts/fission-all/templates/storagesvc/deployment.yaml b/charts/fission-all/templates/storagesvc/deployment.yaml index 03208beb..66fcdb8a 100644 --- a/charts/fission-all/templates/storagesvc/deployment.yaml +++ b/charts/fission-all/templates/storagesvc/deployment.yaml @@ -60,6 +60,7 @@ spec: - name: STORAGE_S3_REGION value: {{ .Values.persistence.s3.region }} {{- end }} + {{- include "fission-resource-namespace.envs" . | indent 8 }} {{- include "opentelemtry.envs" . | indent 8 }} resources: {{- toYaml .Values.storagesvc.resources | nindent 10 }} diff --git a/charts/fission-all/templates/timer/deployment.yaml b/charts/fission-all/templates/timer/deployment.yaml index c413f4c5..63316598 100644 --- a/charts/fission-all/templates/timer/deployment.yaml +++ b/charts/fission-all/templates/timer/deployment.yaml @@ -29,6 +29,7 @@ spec: value: {{ .Values.debugEnv | quote }} - name: PPROF_ENABLED value: {{ .Values.pprof.enabled | quote }} + {{- include "fission-resource-namespace.envs" . | indent 8 }} {{- include "opentelemtry.envs" . | indent 8 }} resources: {{- toYaml .Values.timer.resources | nindent 10 }} diff --git a/charts/fission-all/values.yaml b/charts/fission-all/values.yaml index 45721c57..4103f952 100644 --- a/charts/fission-all/values.yaml +++ b/charts/fission-all/values.yaml @@ -71,10 +71,26 @@ functionNamespace: fission-function ## builderNamespace: fission-builder -## defaultNamespace represents the default namespace in Kubernetes. +## defaultNamespace represents the namespace in which Fission custom resources will be created by the Fission user. +## This is different from the release namespace. +## Please consider setting `singleDefaultNamespace` and `additionalFissionNamespaces` if you want +## more than one namespace to be used for Fission custom resources. +## Fission will watch the defaultNamespace only if `singleDefaultNamespace` is true. ## defaultNamespace: default +## If true, fission will only watch for fission custom resources created in the `defaultNamespace` above. +## +singleDefaultNamespace: true + +## Fission will watch the following namespaces along with the `defaultNamespace` for fission custom resources. +## Only works if `singleDefaultNamespace` is false. +## additionalFissionNamespaces: +## - namespace1 +## - namespace2 +## - namespace3 +additionalFissionNamespaces: [] + ## createNamespace decides to create namespaces by the chart. ## If set to true, functionNamespace and builderNamespace namespaces mentioned above will be created by the chart. ## Set to false if you want to create the namespaces manually. @@ -251,7 +267,7 @@ router: maxRetries: 10 ## Extend the container specs for the core fission pods. - ## Can be used to add things like affinty/tolerations/nodeSelectors/etc. + ## Can be used to add things like affinity/tolerations/nodeSelectors/etc. ## For example: ## extraCoreComponentPodConfig: ## affinity: @@ -479,8 +495,8 @@ serviceMonitor: ##namespace in which you want to deploy servicemonitor ## namespace: "" - ## Map of additional lables to add to the ServiceMonitor resources - # to allow selecting sepcific ServiceMonitors + ## Map of additional labels to add to the ServiceMonitor resources + # to allow selecting specific ServiceMonitors # in case of multiple prometheus deployments additionalServiceMonitorLabels: {} # release: "monitoring" @@ -493,8 +509,8 @@ podMonitor: ##namespace in which you want to deploy podmonitor ## namespace: "" - ## Map of additional lables to add to the PodMonitor resources - # to allow selecting sepcific PodMonitor + ## Map of additional labels to add to the PodMonitor resources + # to allow selecting specific PodMonitor # in case of multiple prometheus deployments additionalPodMonitorLabels: {} # release: "monitoring" @@ -542,7 +558,7 @@ persistence: size: 8Gi ## Extend the container specs for the core fission pods. -## Can be used to add things like affinty/tolerations/nodeSelectors/etc. +## Can be used to add things like affinity/tolerations/nodeSelectors/etc. ## For example: ## extraCoreComponentPodConfig: ## affinity: diff --git a/pkg/apis/core/v1/const.go b/pkg/apis/core/v1/const.go index f4d19220..1a394f94 100644 --- a/pkg/apis/core/v1/const.go +++ b/pkg/apis/core/v1/const.go @@ -153,3 +153,14 @@ const ( ClusterRole = "ClusterRole" ) + +const ( + CanaryConfigResource = "canaryconfigs" + EnvironmentResource = "environments" + FunctionResource = "functions" + HttpTriggerResource = "httptriggers" + KubernetesWatchResource = "kuberneteswatchtriggers" + MessageQueueResource = "messagequeuetriggers" + PackagesResource = "packages" + TimeTriggerResource = "timetriggers" +) diff --git a/pkg/buildermgr/buildermgr.go b/pkg/buildermgr/buildermgr.go index ecccafda..5d3c3f6e 100644 --- a/pkg/buildermgr/buildermgr.go +++ b/pkg/buildermgr/buildermgr.go @@ -29,7 +29,6 @@ import ( "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/executor/util" fetcherConfig "github.com/fission/fission/pkg/fetcher/config" - genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" "github.com/fission/fission/pkg/utils" ) @@ -67,11 +66,10 @@ func Start(ctx context.Context, logger *zap.Logger, storageSvcUrl string, envBui go envWatcher.watchEnvironments(ctx) k8sInformerFactory := k8sInformers.NewSharedInformerFactory(kubernetesClient, time.Minute*30) - informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) podInformer := k8sInformerFactory.Core().V1().Pods().Informer() - pkgInformer := informerFactory.Core().V1().Packages().Informer() pkgWatcher := makePackageWatcher(bmLogger, fissionClient, - kubernetesClient, envBuilderNamespace, storageSvcUrl, &podInformer, &pkgInformer) + kubernetesClient, envBuilderNamespace, storageSvcUrl, podInformer, + utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.PackagesResource)) pkgWatcher.Run(ctx) return nil } diff --git a/pkg/buildermgr/pkgwatcher.go b/pkg/buildermgr/pkgwatcher.go index c38c7891..16b8e8cf 100644 --- a/pkg/buildermgr/pkgwatcher.go +++ b/pkg/buildermgr/pkgwatcher.go @@ -40,8 +40,8 @@ type ( logger *zap.Logger fissionClient versioned.Interface k8sClient kubernetes.Interface - podInformer *k8sCache.SharedIndexInformer - pkgInformer *k8sCache.SharedIndexInformer + podInformer k8sCache.SharedIndexInformer + pkgInformer map[string]k8sCache.SharedIndexInformer builderNamespace string storageSvcUrl string buildCache *cache.Cache @@ -49,8 +49,8 @@ type ( ) func makePackageWatcher(logger *zap.Logger, fissionClient versioned.Interface, k8sClientSet kubernetes.Interface, - builderNamespace string, storageSvcUrl string, podInformer *k8sCache.SharedIndexInformer, - pkgInformer *k8sCache.SharedIndexInformer) *packageWatcher { + builderNamespace string, storageSvcUrl string, podInformer k8sCache.SharedIndexInformer, + pkgInformer map[string]k8sCache.SharedIndexInformer) *packageWatcher { pkgw := &packageWatcher{ logger: logger.Named("package_watcher"), fissionClient: fissionClient, @@ -122,7 +122,7 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) { for healthCheckBackOff.NextExists() { // Informer store is not able to use label to find the pod, // iterate all available environment builders. - items := (*pkgw.podInformer).GetStore().List() + items := pkgw.podInformer.GetStore().List() if err != nil { pkgw.logger.Error("error retrieving pod information for environment", zap.Error(err), zap.String("environment", env.ObjectMeta.Name)) return @@ -323,9 +323,11 @@ func (pkgw *packageWatcher) packageInformerHandler(ctx context.Context) k8sCache func (pkgw *packageWatcher) Run(ctx context.Context) { go metrics.ServeMetrics(ctx, pkgw.logger) - go (*pkgw.podInformer).Run(ctx.Done()) - (*pkgw.pkgInformer).AddEventHandler(pkgw.packageInformerHandler(ctx)) - (*pkgw.pkgInformer).Run(ctx.Done()) + go pkgw.podInformer.Run(ctx.Done()) + for _, pkgInformer := range pkgw.pkgInformer { + pkgInformer.AddEventHandler(pkgw.packageInformerHandler(ctx)) + pkgInformer.Run(ctx.Done()) + } } // setInitialBuildStatus sets initial build status to a package if it is empty. diff --git a/pkg/canaryconfigmgr/canaryConfigMgr.go b/pkg/canaryconfigmgr/canaryConfigMgr.go index 2111fd3a..d13d9546 100644 --- a/pkg/canaryconfigmgr/canaryConfigMgr.go +++ b/pkg/canaryconfigmgr/canaryConfigMgr.go @@ -33,7 +33,7 @@ import ( fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/generated/clientset/versioned" - genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" + "github.com/fission/fission/pkg/utils" ) const ( @@ -44,7 +44,7 @@ type canaryConfigMgr struct { logger *zap.Logger fissionClient versioned.Interface kubeClient kubernetes.Interface - canaryConfigInformer *k8sCache.SharedIndexInformer + canaryConfigInformer map[string]k8sCache.SharedIndexInformer promClient *PrometheusApiClient canaryCfgCancelFuncMap *canaryConfigCancelFuncMap } @@ -92,45 +92,46 @@ func MakeCanaryConfigMgr(ctx context.Context, logger *zap.Logger, fissionClient promClient: promClient, canaryCfgCancelFuncMap: makecanaryConfigCancelFuncMap(), } - - informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) - informer := informerFactory.Core().V1().CanaryConfigs().Informer() - configMgr.canaryConfigInformer = &informer + configMgr.canaryConfigInformer = utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.CanaryConfigResource) configMgr.CanaryConfigEventHandlers(ctx) return configMgr, nil } func (canaryCfgMgr *canaryConfigMgr) CanaryConfigEventHandlers(ctx context.Context) { - (*canaryCfgMgr.canaryConfigInformer).AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ - AddFunc: func(obj interface{}) { - canaryConfig := obj.(*fv1.CanaryConfig) - if canaryConfig.Status.Status == fv1.CanaryConfigStatusPending { - go canaryCfgMgr.addCanaryConfig(ctx, 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(ctx, oldConfig, newConfig) - } - go canaryCfgMgr.reSyncCanaryConfigs(ctx) + for _, informer := range canaryCfgMgr.canaryConfigInformer { + informer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ + AddFunc: func(obj interface{}) { + canaryConfig := obj.(*fv1.CanaryConfig) + if canaryConfig.Status.Status == fv1.CanaryConfigStatusPending { + go canaryCfgMgr.addCanaryConfig(ctx, 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(ctx, oldConfig, newConfig) + } + go canaryCfgMgr.reSyncCanaryConfigs(ctx) - }, - }) + }, + }) + } } func (canaryCfgMgr *canaryConfigMgr) Run(ctx context.Context) { - go (*canaryCfgMgr.canaryConfigInformer).Run(ctx.Done()) + for _, informer := range canaryCfgMgr.canaryConfigInformer { + go informer.Run(ctx.Done()) + } canaryCfgMgr.logger.Info("started canary configmgr controller") } @@ -492,17 +493,19 @@ func (canaryCfgMgr *canaryConfigMgr) rollForward(ctx context.Context, canaryConf } func (canaryCfgMgr *canaryConfigMgr) reSyncCanaryConfigs(ctx context.Context) { - 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 { - canaryCfgMgr.logger.Debug("adding canary config from resync loop", - zap.String("name", canaryConfig.ObjectMeta.Name), - zap.String("namespace", canaryConfig.ObjectMeta.Namespace), - zap.String("version", canaryConfig.ObjectMeta.ResourceVersion)) + for _, informer := range canaryCfgMgr.canaryConfigInformer { + for _, obj := range informer.GetStore().List() { + canaryConfig := obj.(*fv1.CanaryConfig) + _, err := canaryCfgMgr.canaryCfgCancelFuncMap.lookup(&canaryConfig.ObjectMeta) + if err != nil && canaryConfig.Status.Status == fv1.CanaryConfigStatusPending { + canaryCfgMgr.logger.Debug("adding canary config from resync loop", + zap.String("name", canaryConfig.ObjectMeta.Name), + zap.String("namespace", canaryConfig.ObjectMeta.Namespace), + zap.String("version", canaryConfig.ObjectMeta.ResourceVersion)) - // new canaryConfig detected, add it to our cache and start processing it - go canaryCfgMgr.addCanaryConfig(ctx, canaryConfig) + // new canaryConfig detected, add it to our cache and start processing it + go canaryCfgMgr.addCanaryConfig(ctx, canaryConfig) + } } } } diff --git a/pkg/executor/executor.go b/pkg/executor/executor.go index f3ec9741..7eb77bbb 100644 --- a/pkg/executor/executor.go +++ b/pkg/executor/executor.go @@ -45,6 +45,7 @@ import ( fetcherConfig "github.com/fission/fission/pkg/fetcher/config" "github.com/fission/fission/pkg/generated/clientset/versioned" genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" + finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1" "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/metrics" otelUtils "github.com/fission/fission/pkg/utils/otel" @@ -79,7 +80,7 @@ type ( // MakeExecutor returns an Executor for given ExecutorType(s). func MakeExecutor(ctx context.Context, logger *zap.Logger, cms *cms.ConfigSecretController, fissionClient versioned.Interface, types map[fv1.ExecutorType]executortype.ExecutorType, - informers []k8sCache.SharedIndexInformer) (*Executor, error) { + informers ...k8sCache.SharedIndexInformer) (*Executor, error) { executor := &Executor{ logger: logger.Named("executor"), cms: cms, @@ -284,10 +285,24 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace st logger.Info("Starting executor", zap.String("instanceID", executorInstanceID)) - informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) - funcInformer := informerFactory.Core().V1().Functions() - pkgInformer := informerFactory.Core().V1().Packages() - envInformer := informerFactory.Core().V1().Environments() + funcInformer := make(map[string]finformerv1.FunctionInformer, 0) + envInformer := make(map[string]finformerv1.EnvironmentInformer, 0) + pkgInformer := make(map[string]finformerv1.PackageInformer, 0) + + for _, ns := range utils.GetNamespaces() { + factory := genInformer.NewFilteredSharedInformerFactory(fissionClient, time.Minute*30, ns, nil) + funcInformer[ns] = factory.Core().V1().Functions() + } + + for _, ns := range utils.GetNamespaces() { + factory := genInformer.NewFilteredSharedInformerFactory(fissionClient, time.Minute*30, ns, nil) + envInformer[ns] = factory.Core().V1().Environments() + } + + for _, ns := range utils.GetNamespaces() { + factory := genInformer.NewFilteredSharedInformerFactory(fissionClient, time.Minute*30, ns, nil) + pkgInformer[ns] = factory.Core().V1().Packages() + } gpmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypePoolmgr, time.Minute*30) if err != nil { @@ -364,20 +379,29 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace st cms := cms.MakeConfigSecretController(ctx, logger, fissionClient, kubernetesClient, executorTypes, configmapInformer, secretInformer) + fissionInformers := make([]k8sCache.SharedIndexInformer, 0) + for _, informer := range funcInformer { + fissionInformers = append(fissionInformers, informer.Informer()) + } + for _, informer := range envInformer { + fissionInformers = append(fissionInformers, informer.Informer()) + } + for _, informer := range pkgInformer { + fissionInformers = append(fissionInformers, informer.Informer()) + } + fissionInformers = append(fissionInformers, + configmapInformer.Informer(), + secretInformer.Informer(), + gpmPodInformer.Informer(), + gpmRsInformer.Informer(), + ndmDeplInformer.Informer(), + ndmSvcInformer.Informer(), + cnmDeplInformer.Informer(), + cnmSvcInformer.Informer(), + ) api, err := MakeExecutor(ctx, logger, cms, fissionClient, executorTypes, - []k8sCache.SharedIndexInformer{ - funcInformer.Informer(), - pkgInformer.Informer(), - envInformer.Informer(), - configmapInformer.Informer(), - secretInformer.Informer(), - gpmPodInformer.Informer(), - gpmRsInformer.Informer(), - ndmDeplInformer.Informer(), - ndmSvcInformer.Informer(), - cnmDeplInformer.Informer(), - cnmSvcInformer.Informer(), - }) + fissionInformers..., + ) if err != nil { return err } diff --git a/pkg/executor/executortype/container/containermgr.go b/pkg/executor/executortype/container/containermgr.go index 51509b61..341e3591 100644 --- a/pkg/executor/executortype/container/containermgr.go +++ b/pkg/executor/executortype/container/containermgr.go @@ -98,7 +98,7 @@ func MakeContainer( kubernetesClient kubernetes.Interface, namespace string, instanceID string, - funcInformer finformerv1.FunctionInformer, + funcInformer map[string]finformerv1.FunctionInformer, deplInformer appsinformers.DeploymentInformer, svcInformer coreinformers.ServiceInformer, ) (executortype.ExecutorType, error) { @@ -135,7 +135,9 @@ func MakeContainer( caaf.svcLister = svcInformer.Lister() caaf.svcListerSynced = svcInformer.Informer().HasSynced - funcInformer.Informer().AddEventHandler(caaf.FuncInformerHandler(ctx)) + for _, informer := range funcInformer { + informer.Informer().AddEventHandler(caaf.FuncInformerHandler(ctx)) + } return caaf, nil } diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr.go b/pkg/executor/executortype/newdeploy/newdeploymgr.go index 3d3ba360..5d498b22 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr.go @@ -103,8 +103,8 @@ func MakeNewDeploy( namespace string, fetcherConfig *fetcherConfig.Config, instanceID string, - funcInformer finformerv1.FunctionInformer, - envInformer finformerv1.EnvironmentInformer, + funcInformer map[string]finformerv1.FunctionInformer, + envInformer map[string]finformerv1.EnvironmentInformer, deplInformer appsinformers.DeploymentInformer, svcInformer coreinformers.ServiceInformer, podSpecPatch *apiv1.PodSpec, @@ -146,9 +146,12 @@ func MakeNewDeploy( nd.svcLister = svcInformer.Lister() nd.svcListerSynced = svcInformer.Informer().HasSynced - funcInformer.Informer().AddEventHandler(nd.FunctionEventHandlers(ctx)) - envInformer.Informer().AddEventHandler(nd.EnvEventHandlers(ctx)) - + for _, fnInformer := range funcInformer { + fnInformer.Informer().AddEventHandler(nd.FunctionEventHandlers(ctx)) + } + for _, envInformer := range envInformer { + envInformer.Informer().AddEventHandler(nd.EnvEventHandlers(ctx)) + } return nd, nil } diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr_test.go b/pkg/executor/executortype/newdeploy/newdeploymgr_test.go index 13e0bf3d..2252e1eb 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr_test.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr_test.go @@ -22,6 +22,7 @@ import ( fetcherConfig "github.com/fission/fission/pkg/fetcher/config" fClient "github.com/fission/fission/pkg/generated/clientset/versioned/fake" genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" + finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1" "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/loggerfactory" ) @@ -47,9 +48,12 @@ func TestRefreshFuncPods(t *testing.T) { kubernetesClient := fake.NewSimpleClientset() fissionClient := fClient.NewSimpleClientset() informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) - funcInformer := informerFactory.Core().V1().Functions() - envInformer := informerFactory.Core().V1().Environments() - + funcInformer := map[string]finformerv1.FunctionInformer{ + metav1.NamespaceAll: informerFactory.Core().V1().Functions(), + } + envInformer := map[string]finformerv1.EnvironmentInformer{ + metav1.NamespaceAll: informerFactory.Core().V1().Environments(), + } newDeployInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypeNewdeploy, time.Minute*30) if err != nil { t.Fatalf("Error creating informer factory: %s", err) @@ -88,8 +92,8 @@ func TestRefreshFuncPods(t *testing.T) { t.Log("New deploy manager started") runInformers(ctx, []k8sCache.SharedIndexInformer{ - envInformer.Informer(), - funcInformer.Informer(), + envInformer[metav1.NamespaceAll].Informer(), + funcInformer[metav1.NamespaceAll].Informer(), deployInformer.Informer(), svcInformer.Informer(), }) diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index ec3b1238..0f13a231 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -119,9 +119,9 @@ func MakeGenericPoolManager(ctx context.Context, functionNamespace string, fetcherConfig *fetcherConfig.Config, instanceID string, - funcInformer finformerv1.FunctionInformer, - pkgInformer finformerv1.PackageInformer, - envInformer finformerv1.EnvironmentInformer, + funcInformer map[string]finformerv1.FunctionInformer, + pkgInformer map[string]finformerv1.PackageInformer, + envInformer map[string]finformerv1.EnvironmentInformer, podInformer coreinformers.PodInformer, rsInformer appsinformers.ReplicaSetInformer, podSpecPatch *apiv1.PodSpec, @@ -554,7 +554,11 @@ func (gpm *GenericPoolManager) getFunctionEnv(ctx context.Context, fn *fv1.Funct } // Get env from controller - env, err = gpm.poolPodC.envLister.Environments(fn.Spec.Environment.Namespace).Get(fn.Spec.Environment.Name) + envLister, err := gpm.poolPodC.getEnvLister(fn.Spec.Environment.Namespace) + if err != nil { + return nil, err + } + env, err = envLister.Environments(fn.Spec.Environment.Namespace).Get(fn.Spec.Environment.Name) if err != nil { return nil, err } diff --git a/pkg/executor/executortype/poolmgr/poolpodcontroller.go b/pkg/executor/executortype/poolmgr/poolpodcontroller.go index 5c3abd60..9e6fd29a 100644 --- a/pkg/executor/executortype/poolmgr/poolpodcontroller.go +++ b/pkg/executor/executortype/poolmgr/poolpodcontroller.go @@ -17,6 +17,7 @@ package poolmgr import ( "context" + "fmt" "strings" "time" @@ -48,8 +49,8 @@ type ( namespace string enableIstio bool - envLister flisterv1.EnvironmentLister - envListerSynced k8sCache.InformerSynced + envLister map[string]flisterv1.EnvironmentLister + envListerSynced map[string]k8sCache.InformerSynced // podLister can list/get pods from the shared informer's store podLister corelisters.PodLister @@ -70,37 +71,44 @@ func NewPoolPodController(ctx context.Context, logger *zap.Logger, kubernetesClient kubernetes.Interface, namespace string, enableIstio bool, - funcInformer finformerv1.FunctionInformer, - pkgInformer finformerv1.PackageInformer, - envInformer finformerv1.EnvironmentInformer, + funcInformer map[string]finformerv1.FunctionInformer, + pkgInformer map[string]finformerv1.PackageInformer, + envInformer map[string]finformerv1.EnvironmentInformer, rsInformer appsinformers.ReplicaSetInformer, podInformer coreinformers.PodInformer) *PoolPodController { logger = logger.Named("pool_pod_controller") p := &PoolPodController{ - logger: logger, - kubernetesClient: kubernetesClient, - namespace: namespace, - enableIstio: enableIstio, - + logger: logger, + kubernetesClient: kubernetesClient, + namespace: namespace, + enableIstio: enableIstio, + envLister: make(map[string]flisterv1.EnvironmentLister, 0), + envListerSynced: make(map[string]k8sCache.InformerSynced, 0), envCreateUpdateQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "EnvAddUpdateQueue"), envDeleteQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "EnvDeleteQueue"), spCleanupPodQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "SpecializedPodCleanupQueue"), } - funcInformer.Informer().AddEventHandler(FunctionEventHandlers(ctx, p.logger, p.kubernetesClient, p.namespace, p.enableIstio)) - pkgInformer.Informer().AddEventHandler(PackageEventHandlers(ctx, p.logger, p.kubernetesClient, p.namespace)) - envInformer.Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ - AddFunc: p.enqueueEnvAdd, - UpdateFunc: p.enqueueEnvUpdate, - DeleteFunc: p.enqueueEnvDelete, - }) + for _, informer := range funcInformer { + informer.Informer().AddEventHandler(FunctionEventHandlers(ctx, p.logger, p.kubernetesClient, p.namespace, p.enableIstio)) + } + for _, informer := range pkgInformer { + informer.Informer().AddEventHandler(PackageEventHandlers(ctx, p.logger, p.kubernetesClient, p.namespace)) + } + for ns, informer := range envInformer { + informer.Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ + AddFunc: p.enqueueEnvAdd, + UpdateFunc: p.enqueueEnvUpdate, + DeleteFunc: p.enqueueEnvDelete, + }) + p.envLister[ns] = informer.Lister() + p.envListerSynced[ns] = informer.Informer().HasSynced + } rsInformer.Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ AddFunc: p.handleRSAdd, UpdateFunc: p.handleRSUpdate, DeleteFunc: p.handleRSDelete, }) - p.envLister = envInformer.Lister() - p.envListerSynced = envInformer.Informer().HasSynced p.podLister = podInformer.Lister() p.podListerSynced = podInformer.Informer().HasSynced p.logger.Info("pool pod controller handlers registered") @@ -221,10 +229,15 @@ func (p *PoolPodController) Run(ctx context.Context, stopCh <-chan struct{}) { defer p.envCreateUpdateQueue.ShutDown() defer p.envDeleteQueue.ShutDown() defer p.spCleanupPodQueue.ShutDown() - // Wait for the caches to be synced before starting workers p.logger.Info("Waiting for informer caches to sync") - if ok := k8sCache.WaitForCacheSync(stopCh, p.envListerSynced, p.podListerSynced); !ok { + + waitSynced := make([]k8sCache.InformerSynced, 0) + waitSynced = append(waitSynced, p.podListerSynced) + for _, synced := range p.envListerSynced { + waitSynced = append(waitSynced, synced) + } + if ok := k8sCache.WaitForCacheSync(stopCh, waitSynced...); !ok { p.logger.Fatal("failed to wait for caches to sync") } for i := 0; i < 4; i++ { @@ -249,6 +262,20 @@ func (p *PoolPodController) workerRun(ctx context.Context, name string, processF } } +func (p *PoolPodController) getEnvLister(namespace string) (flisterv1.EnvironmentLister, error) { + lister, ok := p.envLister[metav1.NamespaceAll] + if ok { + return lister, nil + } + for ns, lister := range p.envLister { + if ns == namespace { + return lister, nil + } + } + p.logger.Error("no environment lister found for namespace", zap.String("namespace", namespace)) + return nil, fmt.Errorf("no environment lister found for namespace %s", namespace) +} + func (p *PoolPodController) envCreateUpdateQueueProcessFunc(ctx context.Context) bool { maxRetries := 3 handleEnv := func(ctx context.Context, env *fv1.Environment) error { @@ -292,7 +319,13 @@ func (p *PoolPodController) envCreateUpdateQueueProcessFunc(ctx context.Context) p.envCreateUpdateQueue.Forget(key) return false } - env, err := p.envLister.Environments(namespace).Get(name) + envLister, err := p.getEnvLister(namespace) + if err != nil { + p.logger.Error("error getting environment lister", zap.Error(err)) + p.envCreateUpdateQueue.Forget(key) + return false + } + env, err := envLister.Environments(namespace).Get(name) if apierrors.IsNotFound(err) { p.logger.Info("env not found", zap.String("key", key)) p.envCreateUpdateQueue.Forget(key) diff --git a/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go b/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go index cf2d99ca..b7f5870d 100644 --- a/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go +++ b/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go @@ -32,6 +32,7 @@ import ( fetcherConfig "github.com/fission/fission/pkg/fetcher/config" fClient "github.com/fission/fission/pkg/generated/clientset/versioned/fake" genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" + finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1" "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/loggerfactory" ) @@ -50,9 +51,15 @@ func TestPoolPodControllerPodCleanup(t *testing.T) { kubernetesClient := fake.NewSimpleClientset() fissionClient := fClient.NewSimpleClientset() informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) - funcInformer := informerFactory.Core().V1().Functions() - pkgInformer := informerFactory.Core().V1().Packages() - envInformer := informerFactory.Core().V1().Environments() + funcInformer := map[string]finformerv1.FunctionInformer{ + metav1.NamespaceAll: informerFactory.Core().V1().Functions(), + } + pkgInformer := map[string]finformerv1.PackageInformer{ + metav1.NamespaceAll: informerFactory.Core().V1().Packages(), + } + envInformer := map[string]finformerv1.EnvironmentInformer{ + metav1.NamespaceAll: informerFactory.Core().V1().Environments(), + } gpmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypePoolmgr, time.Minute*30) if err != nil { @@ -92,9 +99,9 @@ func TestPoolPodControllerPodCleanup(t *testing.T) { podInformer := gpmPodInformer.Informer() runInformers(ctx, []k8sCache.SharedIndexInformer{ - funcInformer.Informer(), - pkgInformer.Informer(), - envInformer.Informer(), + funcInformer[metav1.NamespaceAll].Informer(), + pkgInformer[metav1.NamespaceAll].Informer(), + envInformer[metav1.NamespaceAll].Informer(), podInformer, gpmRsInformer.Informer(), }) diff --git a/pkg/mqtrigger/mqtmanager.go b/pkg/mqtrigger/mqtmanager.go index b416f63b..f3cc6138 100644 --- a/pkg/mqtrigger/mqtmanager.go +++ b/pkg/mqtrigger/mqtmanager.go @@ -26,8 +26,8 @@ import ( fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/generated/clientset/versioned" - genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" "github.com/fission/fission/pkg/mqtrigger/messageQueue" + "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/metrics" ) @@ -80,12 +80,12 @@ func MakeMessageQueueTriggerManager(logger *zap.Logger, func (mqt *MessageQueueTriggerManager) Run(ctx context.Context) { go mqt.service() - informerFactory := genInformer.NewSharedInformerFactory(mqt.fissionClient, time.Minute*30) - mqTriggerInformer := informerFactory.Core().V1().MessageQueueTriggers().Informer() - mqTriggerInformer.AddEventHandler(mqt.mqtInformerHandlers()) - go mqTriggerInformer.Run(ctx.Done()) - if ok := k8sCache.WaitForCacheSync(ctx.Done(), mqTriggerInformer.HasSynced); !ok { - mqt.logger.Fatal("failed to wait for caches to sync") + for _, informer := range utils.GetInformersForNamespaces(mqt.fissionClient, time.Minute*30, fv1.MessageQueueResource) { + informer.AddEventHandler(mqt.mqtInformerHandlers()) + go informer.Run(ctx.Done()) + if ok := k8sCache.WaitForCacheSync(ctx.Done(), informer.HasSynced); !ok { + mqt.logger.Fatal("failed to wait for caches to sync") + } } go metrics.ServeMetrics(ctx, mqt.logger) } diff --git a/pkg/mqtrigger/scalermanager.go b/pkg/mqtrigger/scalermanager.go index e40e7dbd..5d6b5d80 100644 --- a/pkg/mqtrigger/scalermanager.go +++ b/pkg/mqtrigger/scalermanager.go @@ -24,7 +24,6 @@ import ( fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/executor/util" - genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" "github.com/fission/fission/pkg/utils" ) @@ -158,10 +157,14 @@ func StartScalerManager(ctx context.Context, logger *zap.Logger, routerURL strin if err != nil { return errors.Wrap(err, "error waiting for CRDs") } - informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) - mqTriggerInformer := informerFactory.Core().V1().MessageQueueTriggers().Informer() - mqTriggerInformer.AddEventHandler(mqTriggerEventHandlers(ctx, logger, kubeClient, routerURL)) - mqTriggerInformer.Run(ctx.Done()) + + for _, informer := range utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.MessageQueueResource) { + informer.AddEventHandler(mqTriggerEventHandlers(ctx, logger, kubeClient, routerURL)) + go informer.Run(ctx.Done()) + if ok := k8sCache.WaitForCacheSync(ctx.Done(), informer.HasSynced); !ok { + logger.Fatal("failed to wait for caches to sync") + } + } return nil } diff --git a/pkg/router/functionReferenceResolver.go b/pkg/router/functionReferenceResolver.go index 6e691548..09c71f9b 100644 --- a/pkg/router/functionReferenceResolver.go +++ b/pkg/router/functionReferenceResolver.go @@ -21,6 +21,7 @@ import ( "time" "github.com/pkg/errors" + "go.uber.org/zap" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" k8sCache "k8s.io/client-go/tools/cache" @@ -34,7 +35,8 @@ type ( functionReferenceResolver struct { // FunctionReference -> function metadata refCache *cache.Cache - funcInformer *k8sCache.SharedIndexInformer + funcInformer map[string]k8sCache.SharedIndexInformer + logger *zap.Logger // store k8sCache.Store } @@ -69,10 +71,11 @@ const ( resolveResultMultipleFunctions ) -func makeFunctionReferenceResolver(funcInformer *k8sCache.SharedIndexInformer) *functionReferenceResolver { +func makeFunctionReferenceResolver(logger *zap.Logger, funcInformer map[string]k8sCache.SharedIndexInformer) *functionReferenceResolver { frr := &functionReferenceResolver{ refCache: cache.MakeCache(time.Minute, 0), funcInformer: funcInformer, + logger: logger.Named("function_ref_resolver"), } return frr } @@ -118,10 +121,24 @@ func (frr *functionReferenceResolver) resolve(trigger fv1.HTTPTrigger) (*resolve return rr, nil } +func (frr *functionReferenceResolver) getInformerByNamespace(namespace string) (k8sCache.SharedIndexInformer, error) { + if informer, ok := frr.funcInformer[metav1.NamespaceAll]; ok { + return informer, nil + } + if informer, ok := frr.funcInformer[namespace]; ok { + return informer, nil + } + return nil, fmt.Errorf("informer for namespace %s not found", namespace) +} + // 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.funcInformer).GetStore().Get(&fv1.Function{ + informer, err := frr.getInformerByNamespace(namespace) + if err != nil { + return nil, err + } + obj, isExist, err := informer.GetStore().Get(&fv1.Function{ ObjectMeta: metav1.ObjectMeta{ Namespace: namespace, Name: name, @@ -131,10 +148,11 @@ func (frr *functionReferenceResolver) resolveByName(namespace, name string) (*re return nil, err } if !isExist { - return nil, errors.Errorf("function %v does not exist", name) + frr.logger.Error("function does not exists", zap.String("name", name), zap.String("namespace", namespace)) + return nil, errors.Errorf("function %s/%s does not exist", namespace, name) } - f := obj.(*fv1.Function) + functionMap := map[string]*fv1.Function{ f.ObjectMeta.Name: f, } @@ -155,7 +173,11 @@ func (frr *functionReferenceResolver) resolveByFunctionWeights(namespace string, for functionName, functionWeight := range fr.FunctionWeights { // get function from cache - obj, isExist, err := (*frr.funcInformer).GetStore().Get(&fv1.Function{ + informer, err := frr.getInformerByNamespace(namespace) + if err != nil { + return nil, err + } + obj, isExist, err := informer.GetStore().Get(&fv1.Function{ ObjectMeta: metav1.ObjectMeta{ Namespace: namespace, Name: functionName, @@ -165,9 +187,9 @@ func (frr *functionReferenceResolver) resolveByFunctionWeights(namespace string, return nil, err } if !isExist { - return nil, fmt.Errorf("function %v does not exist", functionName) + frr.logger.Error("function does not exists", zap.String("name", functionName), zap.String("namespace", namespace)) + return nil, fmt.Errorf("function %s/%s does not exist", namespace, functionName) } - f := obj.(*fv1.Function) functionMap[f.ObjectMeta.Name] = f sumPrefix = sumPrefix + functionWeight diff --git a/pkg/router/httpTriggers.go b/pkg/router/httpTriggers.go index 09690ef5..86f1bf48 100644 --- a/pkg/router/httpTriggers.go +++ b/pkg/router/httpTriggers.go @@ -32,7 +32,6 @@ import ( executorClient "github.com/fission/fission/pkg/executor/client" config "github.com/fission/fission/pkg/featureconfig" "github.com/fission/fission/pkg/generated/clientset/versioned" - genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" "github.com/fission/fission/pkg/throttler" "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/metrics" @@ -50,9 +49,9 @@ type HTTPTriggerSet struct { executor *executorClient.Client resolver *functionReferenceResolver triggers []fv1.HTTPTrigger - triggerInformer k8sCache.SharedIndexInformer + triggerInformer map[string]k8sCache.SharedIndexInformer functions []fv1.Function - funcInformer k8sCache.SharedIndexInformer + funcInformer map[string]k8sCache.SharedIndexInformer updateRouterRequestChannel chan struct{} tsRoundTripperParams *tsRoundTripperParams isDebugEnv bool @@ -76,18 +75,15 @@ func makeHTTPTriggerSet(logger *zap.Logger, fmap *functionServiceMap, fissionCli svcAddrUpdateThrottler: actionThrottler, unTapServiceTimeout: unTapServiceTimeout, } - - informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) - httpTriggerSet.triggerInformer = informerFactory.Core().V1().HTTPTriggers().Informer() - httpTriggerSet.funcInformer = informerFactory.Core().V1().Functions().Informer() - + httpTriggerSet.triggerInformer = utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.HttpTriggerResource) + httpTriggerSet.funcInformer = utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.FunctionResource) httpTriggerSet.addTriggerHandlers() httpTriggerSet.addFunctionHandlers() return httpTriggerSet } func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mr *mutableRouter) { - resolver := makeFunctionReferenceResolver(&ts.funcInformer) + resolver := makeFunctionReferenceResolver(ts.logger, ts.funcInformer) ts.resolver = resolver ts.mutableRouter = mr @@ -280,70 +276,75 @@ func (ts *HTTPTriggerSet) updateTriggerStatusFailed(ht *fv1.HTTPTrigger, err err } func (ts *HTTPTriggerSet) addTriggerHandlers() { - ts.triggerInformer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ - AddFunc: func(obj interface{}) { - trigger := obj.(*fv1.HTTPTrigger) - go createIngress(context.Background(), ts.logger, trigger, ts.kubeClient) - ts.syncTriggers() - }, - DeleteFunc: func(obj interface{}) { - ts.syncTriggers() - trigger := obj.(*fv1.HTTPTrigger) - go deleteIngress(context.Background(), ts.logger, trigger, ts.kubeClient) - }, - UpdateFunc: func(oldObj interface{}, newObj interface{}) { - oldTrigger := oldObj.(*fv1.HTTPTrigger) - newTrigger := newObj.(*fv1.HTTPTrigger) + for _, triggerInformer := range ts.triggerInformer { + triggerInformer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ + AddFunc: func(obj interface{}) { + trigger := obj.(*fv1.HTTPTrigger) + go createIngress(context.Background(), ts.logger, trigger, ts.kubeClient) + ts.syncTriggers() + }, + DeleteFunc: func(obj interface{}) { + ts.syncTriggers() + trigger := obj.(*fv1.HTTPTrigger) + go deleteIngress(context.Background(), 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(context.Background(), ts.logger, oldTrigger, newTrigger, ts.kubeClient) - ts.syncTriggers() - }, - }) + go updateIngress(context.Background(), ts.logger, oldTrigger, newTrigger, ts.kubeClient) + ts.syncTriggers() + }, + }) + } } 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) + for _, funcInformer := range ts.funcInformer { - if oldFn.ObjectMeta.ResourceVersion == fn.ObjectMeta.ResourceVersion { - return - } + 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) - // 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 + if oldFn.ObjectMeta.ResourceVersion == fn.ObjectMeta.ResourceVersion { + return } - } - ts.syncTriggers() - }, - }) + + // 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() + }, + }) + } } -func (ts *HTTPTriggerSet) runInformer(ctx context.Context, informer k8sCache.SharedIndexInformer) { - go func() { - informer.Run(ctx.Done()) - }() +func (ts *HTTPTriggerSet) runInformer(ctx context.Context, informer map[string]k8sCache.SharedIndexInformer) { + for _, inf := range informer { + go inf.Run(ctx.Done()) + } } func (ts *HTTPTriggerSet) syncTriggers() { @@ -353,23 +354,27 @@ func (ts *HTTPTriggerSet) syncTriggers() { func (ts *HTTPTriggerSet) updateRouter() { for range ts.updateRouterRequestChannel { // get triggers - latestTriggers := ts.triggerInformer.GetStore().List() - triggers := make([]fv1.HTTPTrigger, len(latestTriggers)) - for _, t := range latestTriggers { - triggers = append(triggers, *t.(*fv1.HTTPTrigger)) + alltriggers := make([]fv1.HTTPTrigger, 0) + for _, triggerInformer := range ts.triggerInformer { + latestTriggers := triggerInformer.GetStore().List() + for _, t := range latestTriggers { + alltriggers = append(alltriggers, *t.(*fv1.HTTPTrigger)) + } } - ts.triggers = triggers + ts.triggers = alltriggers // get functions - latestFunctions := ts.funcInformer.GetStore().List() - functionTimeout := make(map[types.UID]int, len(latestFunctions)) - functions := make([]fv1.Function, len(latestFunctions)) - for _, f := range latestFunctions { - fn := *f.(*fv1.Function) - functionTimeout[fn.ObjectMeta.UID] = fn.Spec.FunctionTimeout - functions = append(functions, *f.(*fv1.Function)) + allfunctions := make([]fv1.Function, 0) + functionTimeout := make(map[types.UID]int, 0) + for _, funcInformer := range ts.funcInformer { + latestFunctions := funcInformer.GetStore().List() + for _, f := range latestFunctions { + fn := *f.(*fv1.Function) + functionTimeout[fn.ObjectMeta.UID] = fn.Spec.FunctionTimeout + allfunctions = append(allfunctions, fn) + } } - ts.functions = functions + ts.functions = allfunctions // make a new router and use it ts.mutableRouter.updateRouter(ts.getRouter(functionTimeout)) diff --git a/pkg/utils/informer.go b/pkg/utils/informer.go index 77bd88ce..2c779423 100644 --- a/pkg/utils/informer.go +++ b/pkg/utils/informer.go @@ -1,6 +1,8 @@ package utils import ( + "os" + "strings" "time" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -8,11 +10,65 @@ import ( "k8s.io/apimachinery/pkg/selection" k8sInformers "k8s.io/client-go/informers" "k8s.io/client-go/kubernetes" + "k8s.io/client-go/tools/cache" metricsapi "k8s.io/metrics/pkg/apis/metrics" v1 "github.com/fission/fission/pkg/apis/core/v1" + + fv1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/generated/clientset/versioned" + genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" ) +const additionalNamespaces string = "FISSION_RESOURCE_NAMESPACES" + +func GetNamespaces() []string { + envValue := os.Getenv(additionalNamespaces) + if len(envValue) == 0 { + return []string{ + metav1.NamespaceAll, + } + } + + informerNS := make([]string, 0) + lstNamespaces := strings.Split(envValue, ",") + for _, namespace := range lstNamespaces { + //check to handle string with additional comma at the end of string. eg- ns1,ns2, + if namespace != "" { + informerNS = append(informerNS, namespace) + } + } + return informerNS +} + +func GetInformersForNamespaces(client versioned.Interface, defaultSync time.Duration, kind string) map[string]cache.SharedIndexInformer { + informers := make(map[string]cache.SharedIndexInformer) + for _, ns := range GetNamespaces() { + factory := genInformer.NewFilteredSharedInformerFactory(client, defaultSync, ns, nil).Core().V1() + switch kind { + case fv1.CanaryConfigResource: + informers[ns] = factory.CanaryConfigs().Informer() + case fv1.EnvironmentResource: + informers[ns] = factory.Environments().Informer() + case fv1.FunctionResource: + informers[ns] = factory.Functions().Informer() + case fv1.HttpTriggerResource: + informers[ns] = factory.HTTPTriggers().Informer() + case fv1.KubernetesWatchResource: + informers[ns] = factory.KubernetesWatchTriggers().Informer() + case fv1.MessageQueueResource: + informers[ns] = factory.MessageQueueTriggers().Informer() + case fv1.PackagesResource: + informers[ns] = factory.Packages().Informer() + case fv1.TimeTriggerResource: + informers[ns] = factory.TimeTriggers().Informer() + default: + panic("Unknown kind: " + kind) + } + } + return informers +} + func GetInformerFactoryByReadyPod(client kubernetes.Interface, namespace string, labelSelector *metav1.LabelSelector) (k8sInformers.SharedInformerFactory, error) { informerFactory := k8sInformers.NewSharedInformerFactoryWithOptions(client, 0, k8sInformers.WithNamespace(namespace),