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 <sanketsudake@gmail.com>
Co-authored-by: shubham bansal <shubhambansaliimtgn@gmail.com>
This commit is contained in:
Sanket Sudake
2022-10-19 15:48:47 +05:30
committed by GitHub
co-authored by shubham bansal
parent facd14de90
commit 827baea974
27 changed files with 430 additions and 219 deletions
@@ -71,3 +71,12 @@ This template generates the image name for the deployment depending on the value
- name: OTEL_PROPAGATORS - name: OTEL_PROPAGATORS
value: "{{ .Values.openTelemetry.propagators }}" value: "{{ .Values.openTelemetry.propagators }}"
{{- end }} {{- 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 }}
@@ -55,6 +55,7 @@ spec:
value: {{ .Values.pprof.enabled | quote }} value: {{ .Values.pprof.enabled | quote }}
- name: HELM_RELEASE_NAME - name: HELM_RELEASE_NAME
value: {{ .Release.Name | quote }} value: {{ .Release.Name | quote }}
{{- include "fission-resource-namespace.envs" . | indent 8 }}
{{- include "opentelemtry.envs" . | indent 8 }} {{- include "opentelemtry.envs" . | indent 8 }}
ports: ports:
- containerPort: 8080 - containerPort: 8080
@@ -38,6 +38,7 @@ spec:
value: {{ .Values.debugEnv | quote }} value: {{ .Values.debugEnv | quote }}
- name: PPROF_ENABLED - name: PPROF_ENABLED
value: {{ .Values.pprof.enabled | quote }} value: {{ .Values.pprof.enabled | quote }}
{{- include "fission-resource-namespace.envs" . | indent 8 }}
- name: POD_NAMESPACE - name: POD_NAMESPACE
valueFrom: valueFrom:
fieldRef: fieldRef:
@@ -71,6 +71,7 @@ spec:
- name: CONTAINER_OBJECT_REAPER_INTERVAL - name: CONTAINER_OBJECT_REAPER_INTERVAL
value: {{ .Values.executor.container.objectReaperInterval | quote }} value: {{ .Values.executor.container.objectReaperInterval | quote }}
{{- end}} {{- end}}
{{- include "fission-resource-namespace.envs" . | indent 8 }}
- name: HELM_RELEASE_NAME - name: HELM_RELEASE_NAME
value: {{ .Release.Name | quote }} value: {{ .Release.Name | quote }}
{{- include "opentelemtry.envs" . | indent 8 }} {{- include "opentelemtry.envs" . | indent 8 }}
@@ -29,6 +29,7 @@ spec:
value: {{ .Values.debugEnv | quote }} value: {{ .Values.debugEnv | quote }}
- name: PPROF_ENABLED - name: PPROF_ENABLED
value: {{ .Values.pprof.enabled | quote }} value: {{ .Values.pprof.enabled | quote }}
{{- include "fission-resource-namespace.envs" . | indent 8 }}
{{- include "opentelemtry.envs" . | indent 8 }} {{- include "opentelemtry.envs" . | indent 8 }}
resources: resources:
{{- toYaml .Values.kubewatcher.resources | nindent 10 }} {{- toYaml .Values.kubewatcher.resources | nindent 10 }}
@@ -47,6 +47,7 @@ spec:
value: {{ .Values.debugEnv | quote }} value: {{ .Values.debugEnv | quote }}
- name: PPROF_ENABLED - name: PPROF_ENABLED
value: {{ .Values.pprof.enabled | quote }} value: {{ .Values.pprof.enabled | quote }}
{{- include "fission-resource-namespace.envs" . | indent 8 }}
{{- include "opentelemtry.envs" . | indent 8 }} {{- include "opentelemtry.envs" . | indent 8 }}
# TLS authentication is TLS with authentication (2 way) # TLS authentication is TLS with authentication (2 way)
# More info: https://docs.confluent.io/current/kafka/authentication_ssl.html#ssl-overview # More info: https://docs.confluent.io/current/kafka/authentication_ssl.html#ssl-overview
@@ -46,6 +46,7 @@ spec:
value: "{{ .Values.mqt_keda.connector_images.gcp_pubsub.image }}:{{ .Values.mqt_keda.connector_images.gcp_pubsub.tag }}" value: "{{ .Values.mqt_keda.connector_images.gcp_pubsub.image }}:{{ .Values.mqt_keda.connector_images.gcp_pubsub.tag }}"
- name: REDIS_IMAGE - name: REDIS_IMAGE
value: "{{ .Values.mqt_keda.connector_images.redis.image }}:{{ .Values.mqt_keda.connector_images.redis.tag }}" 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 }} {{- include "opentelemtry.envs" . | indent 8 }}
resources: resources:
{{- toYaml .Values.mqt_keda.resources | nindent 10 }} {{- toYaml .Values.mqt_keda.resources | nindent 10 }}
@@ -83,6 +83,7 @@ spec:
value: {{ .Values.pprof.enabled | quote }} value: {{ .Values.pprof.enabled | quote }}
- name: DISPLAY_ACCESS_LOG - name: DISPLAY_ACCESS_LOG
value: {{ .Values.router.displayAccessLog | default false | quote }} value: {{ .Values.router.displayAccessLog | default false | quote }}
{{- include "fission-resource-namespace.envs" . | indent 8 }}
{{- include "opentelemtry.envs" . | indent 8 }} {{- include "opentelemtry.envs" . | indent 8 }}
resources: resources:
{{- toYaml .Values.router.resources | nindent 10 }} {{- toYaml .Values.router.resources | nindent 10 }}
@@ -60,6 +60,7 @@ spec:
- name: STORAGE_S3_REGION - name: STORAGE_S3_REGION
value: {{ .Values.persistence.s3.region }} value: {{ .Values.persistence.s3.region }}
{{- end }} {{- end }}
{{- include "fission-resource-namespace.envs" . | indent 8 }}
{{- include "opentelemtry.envs" . | indent 8 }} {{- include "opentelemtry.envs" . | indent 8 }}
resources: resources:
{{- toYaml .Values.storagesvc.resources | nindent 10 }} {{- toYaml .Values.storagesvc.resources | nindent 10 }}
@@ -29,6 +29,7 @@ spec:
value: {{ .Values.debugEnv | quote }} value: {{ .Values.debugEnv | quote }}
- name: PPROF_ENABLED - name: PPROF_ENABLED
value: {{ .Values.pprof.enabled | quote }} value: {{ .Values.pprof.enabled | quote }}
{{- include "fission-resource-namespace.envs" . | indent 8 }}
{{- include "opentelemtry.envs" . | indent 8 }} {{- include "opentelemtry.envs" . | indent 8 }}
resources: resources:
{{- toYaml .Values.timer.resources | nindent 10 }} {{- toYaml .Values.timer.resources | nindent 10 }}
+23 -7
View File
@@ -71,10 +71,26 @@ functionNamespace: fission-function
## ##
builderNamespace: fission-builder 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 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. ## createNamespace decides to create namespaces by the chart.
## If set to true, functionNamespace and builderNamespace namespaces mentioned above will be created 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. ## Set to false if you want to create the namespaces manually.
@@ -251,7 +267,7 @@ router:
maxRetries: 10 maxRetries: 10
## Extend the container specs for the core fission pods. ## 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: ## For example:
## extraCoreComponentPodConfig: ## extraCoreComponentPodConfig:
## affinity: ## affinity:
@@ -479,8 +495,8 @@ serviceMonitor:
##namespace in which you want to deploy servicemonitor ##namespace in which you want to deploy servicemonitor
## ##
namespace: "" namespace: ""
## Map of additional lables to add to the ServiceMonitor resources ## Map of additional labels to add to the ServiceMonitor resources
# to allow selecting sepcific ServiceMonitors # to allow selecting specific ServiceMonitors
# in case of multiple prometheus deployments # in case of multiple prometheus deployments
additionalServiceMonitorLabels: {} additionalServiceMonitorLabels: {}
# release: "monitoring" # release: "monitoring"
@@ -493,8 +509,8 @@ podMonitor:
##namespace in which you want to deploy podmonitor ##namespace in which you want to deploy podmonitor
## ##
namespace: "" namespace: ""
## Map of additional lables to add to the PodMonitor resources ## Map of additional labels to add to the PodMonitor resources
# to allow selecting sepcific PodMonitor # to allow selecting specific PodMonitor
# in case of multiple prometheus deployments # in case of multiple prometheus deployments
additionalPodMonitorLabels: {} additionalPodMonitorLabels: {}
# release: "monitoring" # release: "monitoring"
@@ -542,7 +558,7 @@ persistence:
size: 8Gi size: 8Gi
## Extend the container specs for the core fission pods. ## 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: ## For example:
## extraCoreComponentPodConfig: ## extraCoreComponentPodConfig:
## affinity: ## affinity:
+11
View File
@@ -153,3 +153,14 @@ const (
ClusterRole = "ClusterRole" ClusterRole = "ClusterRole"
) )
const (
CanaryConfigResource = "canaryconfigs"
EnvironmentResource = "environments"
FunctionResource = "functions"
HttpTriggerResource = "httptriggers"
KubernetesWatchResource = "kuberneteswatchtriggers"
MessageQueueResource = "messagequeuetriggers"
PackagesResource = "packages"
TimeTriggerResource = "timetriggers"
)
+2 -4
View File
@@ -29,7 +29,6 @@ import (
"github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/crd"
"github.com/fission/fission/pkg/executor/util" "github.com/fission/fission/pkg/executor/util"
fetcherConfig "github.com/fission/fission/pkg/fetcher/config" fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions"
"github.com/fission/fission/pkg/utils" "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) go envWatcher.watchEnvironments(ctx)
k8sInformerFactory := k8sInformers.NewSharedInformerFactory(kubernetesClient, time.Minute*30) k8sInformerFactory := k8sInformers.NewSharedInformerFactory(kubernetesClient, time.Minute*30)
informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30)
podInformer := k8sInformerFactory.Core().V1().Pods().Informer() podInformer := k8sInformerFactory.Core().V1().Pods().Informer()
pkgInformer := informerFactory.Core().V1().Packages().Informer()
pkgWatcher := makePackageWatcher(bmLogger, fissionClient, pkgWatcher := makePackageWatcher(bmLogger, fissionClient,
kubernetesClient, envBuilderNamespace, storageSvcUrl, &podInformer, &pkgInformer) kubernetesClient, envBuilderNamespace, storageSvcUrl, podInformer,
utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.PackagesResource))
pkgWatcher.Run(ctx) pkgWatcher.Run(ctx)
return nil return nil
} }
+10 -8
View File
@@ -40,8 +40,8 @@ type (
logger *zap.Logger logger *zap.Logger
fissionClient versioned.Interface fissionClient versioned.Interface
k8sClient kubernetes.Interface k8sClient kubernetes.Interface
podInformer *k8sCache.SharedIndexInformer podInformer k8sCache.SharedIndexInformer
pkgInformer *k8sCache.SharedIndexInformer pkgInformer map[string]k8sCache.SharedIndexInformer
builderNamespace string builderNamespace string
storageSvcUrl string storageSvcUrl string
buildCache *cache.Cache buildCache *cache.Cache
@@ -49,8 +49,8 @@ type (
) )
func makePackageWatcher(logger *zap.Logger, fissionClient versioned.Interface, k8sClientSet kubernetes.Interface, func makePackageWatcher(logger *zap.Logger, fissionClient versioned.Interface, k8sClientSet kubernetes.Interface,
builderNamespace string, storageSvcUrl string, podInformer *k8sCache.SharedIndexInformer, builderNamespace string, storageSvcUrl string, podInformer k8sCache.SharedIndexInformer,
pkgInformer *k8sCache.SharedIndexInformer) *packageWatcher { pkgInformer map[string]k8sCache.SharedIndexInformer) *packageWatcher {
pkgw := &packageWatcher{ pkgw := &packageWatcher{
logger: logger.Named("package_watcher"), logger: logger.Named("package_watcher"),
fissionClient: fissionClient, fissionClient: fissionClient,
@@ -122,7 +122,7 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) {
for healthCheckBackOff.NextExists() { for healthCheckBackOff.NextExists() {
// Informer store is not able to use label to find the pod, // Informer store is not able to use label to find the pod,
// iterate all available environment builders. // iterate all available environment builders.
items := (*pkgw.podInformer).GetStore().List() items := pkgw.podInformer.GetStore().List()
if err != nil { if err != nil {
pkgw.logger.Error("error retrieving pod information for environment", zap.Error(err), zap.String("environment", env.ObjectMeta.Name)) pkgw.logger.Error("error retrieving pod information for environment", zap.Error(err), zap.String("environment", env.ObjectMeta.Name))
return return
@@ -323,9 +323,11 @@ func (pkgw *packageWatcher) packageInformerHandler(ctx context.Context) k8sCache
func (pkgw *packageWatcher) Run(ctx context.Context) { func (pkgw *packageWatcher) Run(ctx context.Context) {
go metrics.ServeMetrics(ctx, pkgw.logger) go metrics.ServeMetrics(ctx, pkgw.logger)
go (*pkgw.podInformer).Run(ctx.Done()) go pkgw.podInformer.Run(ctx.Done())
(*pkgw.pkgInformer).AddEventHandler(pkgw.packageInformerHandler(ctx)) for _, pkgInformer := range pkgw.pkgInformer {
(*pkgw.pkgInformer).Run(ctx.Done()) pkgInformer.AddEventHandler(pkgw.packageInformerHandler(ctx))
pkgInformer.Run(ctx.Done())
}
} }
// setInitialBuildStatus sets initial build status to a package if it is empty. // setInitialBuildStatus sets initial build status to a package if it is empty.
+45 -42
View File
@@ -33,7 +33,7 @@ import (
fv1 "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" "github.com/fission/fission/pkg/generated/clientset/versioned"
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" "github.com/fission/fission/pkg/utils"
) )
const ( const (
@@ -44,7 +44,7 @@ type canaryConfigMgr struct {
logger *zap.Logger logger *zap.Logger
fissionClient versioned.Interface fissionClient versioned.Interface
kubeClient kubernetes.Interface kubeClient kubernetes.Interface
canaryConfigInformer *k8sCache.SharedIndexInformer canaryConfigInformer map[string]k8sCache.SharedIndexInformer
promClient *PrometheusApiClient promClient *PrometheusApiClient
canaryCfgCancelFuncMap *canaryConfigCancelFuncMap canaryCfgCancelFuncMap *canaryConfigCancelFuncMap
} }
@@ -92,45 +92,46 @@ func MakeCanaryConfigMgr(ctx context.Context, logger *zap.Logger, fissionClient
promClient: promClient, promClient: promClient,
canaryCfgCancelFuncMap: makecanaryConfigCancelFuncMap(), canaryCfgCancelFuncMap: makecanaryConfigCancelFuncMap(),
} }
configMgr.canaryConfigInformer = utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.CanaryConfigResource)
informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30)
informer := informerFactory.Core().V1().CanaryConfigs().Informer()
configMgr.canaryConfigInformer = &informer
configMgr.CanaryConfigEventHandlers(ctx) configMgr.CanaryConfigEventHandlers(ctx)
return configMgr, nil return configMgr, nil
} }
func (canaryCfgMgr *canaryConfigMgr) CanaryConfigEventHandlers(ctx context.Context) { func (canaryCfgMgr *canaryConfigMgr) CanaryConfigEventHandlers(ctx context.Context) {
(*canaryCfgMgr.canaryConfigInformer).AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ for _, informer := range canaryCfgMgr.canaryConfigInformer {
AddFunc: func(obj interface{}) { informer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
canaryConfig := obj.(*fv1.CanaryConfig) AddFunc: func(obj interface{}) {
if canaryConfig.Status.Status == fv1.CanaryConfigStatusPending { canaryConfig := obj.(*fv1.CanaryConfig)
go canaryCfgMgr.addCanaryConfig(ctx, canaryConfig) if canaryConfig.Status.Status == fv1.CanaryConfigStatusPending {
} go canaryCfgMgr.addCanaryConfig(ctx, canaryConfig)
}, }
DeleteFunc: func(obj interface{}) { },
canaryConfig := obj.(*fv1.CanaryConfig) DeleteFunc: func(obj interface{}) {
go canaryCfgMgr.deleteCanaryConfig(canaryConfig) canaryConfig := obj.(*fv1.CanaryConfig)
}, go canaryCfgMgr.deleteCanaryConfig(canaryConfig)
UpdateFunc: func(oldObj interface{}, newObj interface{}) { },
oldConfig := oldObj.(*fv1.CanaryConfig) UpdateFunc: func(oldObj interface{}, newObj interface{}) {
newConfig := newObj.(*fv1.CanaryConfig) oldConfig := oldObj.(*fv1.CanaryConfig)
if oldConfig.ObjectMeta.ResourceVersion != newConfig.ObjectMeta.ResourceVersion && newConfig := newObj.(*fv1.CanaryConfig)
newConfig.Status.Status == fv1.CanaryConfigStatusPending { if oldConfig.ObjectMeta.ResourceVersion != newConfig.ObjectMeta.ResourceVersion &&
canaryCfgMgr.logger.Info("update canary config invoked", newConfig.Status.Status == fv1.CanaryConfigStatusPending {
zap.String("name", newConfig.ObjectMeta.Name), canaryCfgMgr.logger.Info("update canary config invoked",
zap.String("namespace", newConfig.ObjectMeta.Namespace), zap.String("name", newConfig.ObjectMeta.Name),
zap.String("version", newConfig.ObjectMeta.ResourceVersion)) zap.String("namespace", newConfig.ObjectMeta.Namespace),
go canaryCfgMgr.updateCanaryConfig(ctx, oldConfig, newConfig) zap.String("version", newConfig.ObjectMeta.ResourceVersion))
} go canaryCfgMgr.updateCanaryConfig(ctx, oldConfig, newConfig)
go canaryCfgMgr.reSyncCanaryConfigs(ctx) }
go canaryCfgMgr.reSyncCanaryConfigs(ctx)
}, },
}) })
}
} }
func (canaryCfgMgr *canaryConfigMgr) Run(ctx context.Context) { 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") 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) { func (canaryCfgMgr *canaryConfigMgr) reSyncCanaryConfigs(ctx context.Context) {
for _, obj := range (*canaryCfgMgr.canaryConfigInformer).GetStore().List() { for _, informer := range canaryCfgMgr.canaryConfigInformer {
canaryConfig := obj.(*fv1.CanaryConfig) for _, obj := range informer.GetStore().List() {
_, err := canaryCfgMgr.canaryCfgCancelFuncMap.lookup(&canaryConfig.ObjectMeta) canaryConfig := obj.(*fv1.CanaryConfig)
if err != nil && canaryConfig.Status.Status == fv1.CanaryConfigStatusPending { _, err := canaryCfgMgr.canaryCfgCancelFuncMap.lookup(&canaryConfig.ObjectMeta)
canaryCfgMgr.logger.Debug("adding canary config from resync loop", if err != nil && canaryConfig.Status.Status == fv1.CanaryConfigStatusPending {
zap.String("name", canaryConfig.ObjectMeta.Name), canaryCfgMgr.logger.Debug("adding canary config from resync loop",
zap.String("namespace", canaryConfig.ObjectMeta.Namespace), zap.String("name", canaryConfig.ObjectMeta.Name),
zap.String("version", canaryConfig.ObjectMeta.ResourceVersion)) zap.String("namespace", canaryConfig.ObjectMeta.Namespace),
zap.String("version", canaryConfig.ObjectMeta.ResourceVersion))
// new canaryConfig detected, add it to our cache and start processing it // new canaryConfig detected, add it to our cache and start processing it
go canaryCfgMgr.addCanaryConfig(ctx, canaryConfig) go canaryCfgMgr.addCanaryConfig(ctx, canaryConfig)
}
} }
} }
} }
+42 -18
View File
@@ -45,6 +45,7 @@ import (
fetcherConfig "github.com/fission/fission/pkg/fetcher/config" fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
"github.com/fission/fission/pkg/generated/clientset/versioned" "github.com/fission/fission/pkg/generated/clientset/versioned"
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" 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"
"github.com/fission/fission/pkg/utils/metrics" "github.com/fission/fission/pkg/utils/metrics"
otelUtils "github.com/fission/fission/pkg/utils/otel" otelUtils "github.com/fission/fission/pkg/utils/otel"
@@ -79,7 +80,7 @@ type (
// MakeExecutor returns an Executor for given ExecutorType(s). // MakeExecutor returns an Executor for given ExecutorType(s).
func MakeExecutor(ctx context.Context, logger *zap.Logger, cms *cms.ConfigSecretController, func MakeExecutor(ctx context.Context, logger *zap.Logger, cms *cms.ConfigSecretController,
fissionClient versioned.Interface, types map[fv1.ExecutorType]executortype.ExecutorType, fissionClient versioned.Interface, types map[fv1.ExecutorType]executortype.ExecutorType,
informers []k8sCache.SharedIndexInformer) (*Executor, error) { informers ...k8sCache.SharedIndexInformer) (*Executor, error) {
executor := &Executor{ executor := &Executor{
logger: logger.Named("executor"), logger: logger.Named("executor"),
cms: cms, cms: cms,
@@ -284,10 +285,24 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace st
logger.Info("Starting executor", zap.String("instanceID", executorInstanceID)) logger.Info("Starting executor", zap.String("instanceID", executorInstanceID))
informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) funcInformer := make(map[string]finformerv1.FunctionInformer, 0)
funcInformer := informerFactory.Core().V1().Functions() envInformer := make(map[string]finformerv1.EnvironmentInformer, 0)
pkgInformer := informerFactory.Core().V1().Packages() pkgInformer := make(map[string]finformerv1.PackageInformer, 0)
envInformer := informerFactory.Core().V1().Environments()
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) gpmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypePoolmgr, time.Minute*30)
if err != nil { 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) 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, api, err := MakeExecutor(ctx, logger, cms, fissionClient, executorTypes,
[]k8sCache.SharedIndexInformer{ fissionInformers...,
funcInformer.Informer(), )
pkgInformer.Informer(),
envInformer.Informer(),
configmapInformer.Informer(),
secretInformer.Informer(),
gpmPodInformer.Informer(),
gpmRsInformer.Informer(),
ndmDeplInformer.Informer(),
ndmSvcInformer.Informer(),
cnmDeplInformer.Informer(),
cnmSvcInformer.Informer(),
})
if err != nil { if err != nil {
return err return err
} }
@@ -98,7 +98,7 @@ func MakeContainer(
kubernetesClient kubernetes.Interface, kubernetesClient kubernetes.Interface,
namespace string, namespace string,
instanceID string, instanceID string,
funcInformer finformerv1.FunctionInformer, funcInformer map[string]finformerv1.FunctionInformer,
deplInformer appsinformers.DeploymentInformer, deplInformer appsinformers.DeploymentInformer,
svcInformer coreinformers.ServiceInformer, svcInformer coreinformers.ServiceInformer,
) (executortype.ExecutorType, error) { ) (executortype.ExecutorType, error) {
@@ -135,7 +135,9 @@ func MakeContainer(
caaf.svcLister = svcInformer.Lister() caaf.svcLister = svcInformer.Lister()
caaf.svcListerSynced = svcInformer.Informer().HasSynced caaf.svcListerSynced = svcInformer.Informer().HasSynced
funcInformer.Informer().AddEventHandler(caaf.FuncInformerHandler(ctx)) for _, informer := range funcInformer {
informer.Informer().AddEventHandler(caaf.FuncInformerHandler(ctx))
}
return caaf, nil return caaf, nil
} }
@@ -103,8 +103,8 @@ func MakeNewDeploy(
namespace string, namespace string,
fetcherConfig *fetcherConfig.Config, fetcherConfig *fetcherConfig.Config,
instanceID string, instanceID string,
funcInformer finformerv1.FunctionInformer, funcInformer map[string]finformerv1.FunctionInformer,
envInformer finformerv1.EnvironmentInformer, envInformer map[string]finformerv1.EnvironmentInformer,
deplInformer appsinformers.DeploymentInformer, deplInformer appsinformers.DeploymentInformer,
svcInformer coreinformers.ServiceInformer, svcInformer coreinformers.ServiceInformer,
podSpecPatch *apiv1.PodSpec, podSpecPatch *apiv1.PodSpec,
@@ -146,9 +146,12 @@ func MakeNewDeploy(
nd.svcLister = svcInformer.Lister() nd.svcLister = svcInformer.Lister()
nd.svcListerSynced = svcInformer.Informer().HasSynced nd.svcListerSynced = svcInformer.Informer().HasSynced
funcInformer.Informer().AddEventHandler(nd.FunctionEventHandlers(ctx)) for _, fnInformer := range funcInformer {
envInformer.Informer().AddEventHandler(nd.EnvEventHandlers(ctx)) fnInformer.Informer().AddEventHandler(nd.FunctionEventHandlers(ctx))
}
for _, envInformer := range envInformer {
envInformer.Informer().AddEventHandler(nd.EnvEventHandlers(ctx))
}
return nd, nil return nd, nil
} }
@@ -22,6 +22,7 @@ import (
fetcherConfig "github.com/fission/fission/pkg/fetcher/config" fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
fClient "github.com/fission/fission/pkg/generated/clientset/versioned/fake" fClient "github.com/fission/fission/pkg/generated/clientset/versioned/fake"
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" 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"
"github.com/fission/fission/pkg/utils/loggerfactory" "github.com/fission/fission/pkg/utils/loggerfactory"
) )
@@ -47,9 +48,12 @@ func TestRefreshFuncPods(t *testing.T) {
kubernetesClient := fake.NewSimpleClientset() kubernetesClient := fake.NewSimpleClientset()
fissionClient := fClient.NewSimpleClientset() fissionClient := fClient.NewSimpleClientset()
informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30)
funcInformer := informerFactory.Core().V1().Functions() funcInformer := map[string]finformerv1.FunctionInformer{
envInformer := informerFactory.Core().V1().Environments() 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) newDeployInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypeNewdeploy, time.Minute*30)
if err != nil { if err != nil {
t.Fatalf("Error creating informer factory: %s", err) t.Fatalf("Error creating informer factory: %s", err)
@@ -88,8 +92,8 @@ func TestRefreshFuncPods(t *testing.T) {
t.Log("New deploy manager started") t.Log("New deploy manager started")
runInformers(ctx, []k8sCache.SharedIndexInformer{ runInformers(ctx, []k8sCache.SharedIndexInformer{
envInformer.Informer(), envInformer[metav1.NamespaceAll].Informer(),
funcInformer.Informer(), funcInformer[metav1.NamespaceAll].Informer(),
deployInformer.Informer(), deployInformer.Informer(),
svcInformer.Informer(), svcInformer.Informer(),
}) })
+8 -4
View File
@@ -119,9 +119,9 @@ func MakeGenericPoolManager(ctx context.Context,
functionNamespace string, functionNamespace string,
fetcherConfig *fetcherConfig.Config, fetcherConfig *fetcherConfig.Config,
instanceID string, instanceID string,
funcInformer finformerv1.FunctionInformer, funcInformer map[string]finformerv1.FunctionInformer,
pkgInformer finformerv1.PackageInformer, pkgInformer map[string]finformerv1.PackageInformer,
envInformer finformerv1.EnvironmentInformer, envInformer map[string]finformerv1.EnvironmentInformer,
podInformer coreinformers.PodInformer, podInformer coreinformers.PodInformer,
rsInformer appsinformers.ReplicaSetInformer, rsInformer appsinformers.ReplicaSetInformer,
podSpecPatch *apiv1.PodSpec, podSpecPatch *apiv1.PodSpec,
@@ -554,7 +554,11 @@ func (gpm *GenericPoolManager) getFunctionEnv(ctx context.Context, fn *fv1.Funct
} }
// Get env from controller // 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 { if err != nil {
return nil, err return nil, err
} }
@@ -17,6 +17,7 @@ package poolmgr
import ( import (
"context" "context"
"fmt"
"strings" "strings"
"time" "time"
@@ -48,8 +49,8 @@ type (
namespace string namespace string
enableIstio bool enableIstio bool
envLister flisterv1.EnvironmentLister envLister map[string]flisterv1.EnvironmentLister
envListerSynced k8sCache.InformerSynced envListerSynced map[string]k8sCache.InformerSynced
// podLister can list/get pods from the shared informer's store // podLister can list/get pods from the shared informer's store
podLister corelisters.PodLister podLister corelisters.PodLister
@@ -70,37 +71,44 @@ func NewPoolPodController(ctx context.Context, logger *zap.Logger,
kubernetesClient kubernetes.Interface, kubernetesClient kubernetes.Interface,
namespace string, namespace string,
enableIstio bool, enableIstio bool,
funcInformer finformerv1.FunctionInformer, funcInformer map[string]finformerv1.FunctionInformer,
pkgInformer finformerv1.PackageInformer, pkgInformer map[string]finformerv1.PackageInformer,
envInformer finformerv1.EnvironmentInformer, envInformer map[string]finformerv1.EnvironmentInformer,
rsInformer appsinformers.ReplicaSetInformer, rsInformer appsinformers.ReplicaSetInformer,
podInformer coreinformers.PodInformer) *PoolPodController { podInformer coreinformers.PodInformer) *PoolPodController {
logger = logger.Named("pool_pod_controller") logger = logger.Named("pool_pod_controller")
p := &PoolPodController{ p := &PoolPodController{
logger: logger, logger: logger,
kubernetesClient: kubernetesClient, kubernetesClient: kubernetesClient,
namespace: namespace, namespace: namespace,
enableIstio: enableIstio, enableIstio: enableIstio,
envLister: make(map[string]flisterv1.EnvironmentLister, 0),
envListerSynced: make(map[string]k8sCache.InformerSynced, 0),
envCreateUpdateQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "EnvAddUpdateQueue"), envCreateUpdateQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "EnvAddUpdateQueue"),
envDeleteQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "EnvDeleteQueue"), envDeleteQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "EnvDeleteQueue"),
spCleanupPodQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "SpecializedPodCleanupQueue"), spCleanupPodQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "SpecializedPodCleanupQueue"),
} }
funcInformer.Informer().AddEventHandler(FunctionEventHandlers(ctx, p.logger, p.kubernetesClient, p.namespace, p.enableIstio)) for _, informer := range funcInformer {
pkgInformer.Informer().AddEventHandler(PackageEventHandlers(ctx, p.logger, p.kubernetesClient, p.namespace)) informer.Informer().AddEventHandler(FunctionEventHandlers(ctx, p.logger, p.kubernetesClient, p.namespace, p.enableIstio))
envInformer.Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ }
AddFunc: p.enqueueEnvAdd, for _, informer := range pkgInformer {
UpdateFunc: p.enqueueEnvUpdate, informer.Informer().AddEventHandler(PackageEventHandlers(ctx, p.logger, p.kubernetesClient, p.namespace))
DeleteFunc: p.enqueueEnvDelete, }
}) 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{ rsInformer.Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
AddFunc: p.handleRSAdd, AddFunc: p.handleRSAdd,
UpdateFunc: p.handleRSUpdate, UpdateFunc: p.handleRSUpdate,
DeleteFunc: p.handleRSDelete, DeleteFunc: p.handleRSDelete,
}) })
p.envLister = envInformer.Lister()
p.envListerSynced = envInformer.Informer().HasSynced
p.podLister = podInformer.Lister() p.podLister = podInformer.Lister()
p.podListerSynced = podInformer.Informer().HasSynced p.podListerSynced = podInformer.Informer().HasSynced
p.logger.Info("pool pod controller handlers registered") 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.envCreateUpdateQueue.ShutDown()
defer p.envDeleteQueue.ShutDown() defer p.envDeleteQueue.ShutDown()
defer p.spCleanupPodQueue.ShutDown() defer p.spCleanupPodQueue.ShutDown()
// Wait for the caches to be synced before starting workers // Wait for the caches to be synced before starting workers
p.logger.Info("Waiting for informer caches to sync") 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") p.logger.Fatal("failed to wait for caches to sync")
} }
for i := 0; i < 4; i++ { 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 { func (p *PoolPodController) envCreateUpdateQueueProcessFunc(ctx context.Context) bool {
maxRetries := 3 maxRetries := 3
handleEnv := func(ctx context.Context, env *fv1.Environment) error { handleEnv := func(ctx context.Context, env *fv1.Environment) error {
@@ -292,7 +319,13 @@ func (p *PoolPodController) envCreateUpdateQueueProcessFunc(ctx context.Context)
p.envCreateUpdateQueue.Forget(key) p.envCreateUpdateQueue.Forget(key)
return false 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) { if apierrors.IsNotFound(err) {
p.logger.Info("env not found", zap.String("key", key)) p.logger.Info("env not found", zap.String("key", key))
p.envCreateUpdateQueue.Forget(key) p.envCreateUpdateQueue.Forget(key)
@@ -32,6 +32,7 @@ import (
fetcherConfig "github.com/fission/fission/pkg/fetcher/config" fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
fClient "github.com/fission/fission/pkg/generated/clientset/versioned/fake" fClient "github.com/fission/fission/pkg/generated/clientset/versioned/fake"
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions" 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"
"github.com/fission/fission/pkg/utils/loggerfactory" "github.com/fission/fission/pkg/utils/loggerfactory"
) )
@@ -50,9 +51,15 @@ func TestPoolPodControllerPodCleanup(t *testing.T) {
kubernetesClient := fake.NewSimpleClientset() kubernetesClient := fake.NewSimpleClientset()
fissionClient := fClient.NewSimpleClientset() fissionClient := fClient.NewSimpleClientset()
informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30)
funcInformer := informerFactory.Core().V1().Functions() funcInformer := map[string]finformerv1.FunctionInformer{
pkgInformer := informerFactory.Core().V1().Packages() metav1.NamespaceAll: informerFactory.Core().V1().Functions(),
envInformer := informerFactory.Core().V1().Environments() }
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) gpmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypePoolmgr, time.Minute*30)
if err != nil { if err != nil {
@@ -92,9 +99,9 @@ func TestPoolPodControllerPodCleanup(t *testing.T) {
podInformer := gpmPodInformer.Informer() podInformer := gpmPodInformer.Informer()
runInformers(ctx, []k8sCache.SharedIndexInformer{ runInformers(ctx, []k8sCache.SharedIndexInformer{
funcInformer.Informer(), funcInformer[metav1.NamespaceAll].Informer(),
pkgInformer.Informer(), pkgInformer[metav1.NamespaceAll].Informer(),
envInformer.Informer(), envInformer[metav1.NamespaceAll].Informer(),
podInformer, podInformer,
gpmRsInformer.Informer(), gpmRsInformer.Informer(),
}) })
+7 -7
View File
@@ -26,8 +26,8 @@ import (
fv1 "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" "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/mqtrigger/messageQueue"
"github.com/fission/fission/pkg/utils"
"github.com/fission/fission/pkg/utils/metrics" "github.com/fission/fission/pkg/utils/metrics"
) )
@@ -80,12 +80,12 @@ func MakeMessageQueueTriggerManager(logger *zap.Logger,
func (mqt *MessageQueueTriggerManager) Run(ctx context.Context) { func (mqt *MessageQueueTriggerManager) Run(ctx context.Context) {
go mqt.service() go mqt.service()
informerFactory := genInformer.NewSharedInformerFactory(mqt.fissionClient, time.Minute*30) for _, informer := range utils.GetInformersForNamespaces(mqt.fissionClient, time.Minute*30, fv1.MessageQueueResource) {
mqTriggerInformer := informerFactory.Core().V1().MessageQueueTriggers().Informer() informer.AddEventHandler(mqt.mqtInformerHandlers())
mqTriggerInformer.AddEventHandler(mqt.mqtInformerHandlers()) go informer.Run(ctx.Done())
go mqTriggerInformer.Run(ctx.Done()) if ok := k8sCache.WaitForCacheSync(ctx.Done(), informer.HasSynced); !ok {
if ok := k8sCache.WaitForCacheSync(ctx.Done(), mqTriggerInformer.HasSynced); !ok { mqt.logger.Fatal("failed to wait for caches to sync")
mqt.logger.Fatal("failed to wait for caches to sync") }
} }
go metrics.ServeMetrics(ctx, mqt.logger) go metrics.ServeMetrics(ctx, mqt.logger)
} }
+8 -5
View File
@@ -24,7 +24,6 @@ import (
fv1 "github.com/fission/fission/pkg/apis/core/v1" fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/crd"
"github.com/fission/fission/pkg/executor/util" "github.com/fission/fission/pkg/executor/util"
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions"
"github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils"
) )
@@ -158,10 +157,14 @@ func StartScalerManager(ctx context.Context, logger *zap.Logger, routerURL strin
if err != nil { if err != nil {
return errors.Wrap(err, "error waiting for CRDs") return errors.Wrap(err, "error waiting for CRDs")
} }
informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30)
mqTriggerInformer := informerFactory.Core().V1().MessageQueueTriggers().Informer() for _, informer := range utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.MessageQueueResource) {
mqTriggerInformer.AddEventHandler(mqTriggerEventHandlers(ctx, logger, kubeClient, routerURL)) informer.AddEventHandler(mqTriggerEventHandlers(ctx, logger, kubeClient, routerURL))
mqTriggerInformer.Run(ctx.Done()) 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 return nil
} }
+30 -8
View File
@@ -21,6 +21,7 @@ import (
"time" "time"
"github.com/pkg/errors" "github.com/pkg/errors"
"go.uber.org/zap"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
k8sCache "k8s.io/client-go/tools/cache" k8sCache "k8s.io/client-go/tools/cache"
@@ -34,7 +35,8 @@ type (
functionReferenceResolver struct { functionReferenceResolver struct {
// FunctionReference -> function metadata // FunctionReference -> function metadata
refCache *cache.Cache refCache *cache.Cache
funcInformer *k8sCache.SharedIndexInformer funcInformer map[string]k8sCache.SharedIndexInformer
logger *zap.Logger
// store k8sCache.Store // store k8sCache.Store
} }
@@ -69,10 +71,11 @@ const (
resolveResultMultipleFunctions resolveResultMultipleFunctions
) )
func makeFunctionReferenceResolver(funcInformer *k8sCache.SharedIndexInformer) *functionReferenceResolver { func makeFunctionReferenceResolver(logger *zap.Logger, funcInformer map[string]k8sCache.SharedIndexInformer) *functionReferenceResolver {
frr := &functionReferenceResolver{ frr := &functionReferenceResolver{
refCache: cache.MakeCache(time.Minute, 0), refCache: cache.MakeCache(time.Minute, 0),
funcInformer: funcInformer, funcInformer: funcInformer,
logger: logger.Named("function_ref_resolver"),
} }
return frr return frr
} }
@@ -118,10 +121,24 @@ func (frr *functionReferenceResolver) resolve(trigger fv1.HTTPTrigger) (*resolve
return rr, nil 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. // resolveByName simply looks up function by name in a namespace.
func (frr *functionReferenceResolver) resolveByName(namespace, name string) (*resolveResult, error) { func (frr *functionReferenceResolver) resolveByName(namespace, name string) (*resolveResult, error) {
// get function from cache // 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{ ObjectMeta: metav1.ObjectMeta{
Namespace: namespace, Namespace: namespace,
Name: name, Name: name,
@@ -131,10 +148,11 @@ func (frr *functionReferenceResolver) resolveByName(namespace, name string) (*re
return nil, err return nil, err
} }
if !isExist { 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) f := obj.(*fv1.Function)
functionMap := map[string]*fv1.Function{ functionMap := map[string]*fv1.Function{
f.ObjectMeta.Name: f, f.ObjectMeta.Name: f,
} }
@@ -155,7 +173,11 @@ func (frr *functionReferenceResolver) resolveByFunctionWeights(namespace string,
for functionName, functionWeight := range fr.FunctionWeights { for functionName, functionWeight := range fr.FunctionWeights {
// get function from cache // 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{ ObjectMeta: metav1.ObjectMeta{
Namespace: namespace, Namespace: namespace,
Name: functionName, Name: functionName,
@@ -165,9 +187,9 @@ func (frr *functionReferenceResolver) resolveByFunctionWeights(namespace string,
return nil, err return nil, err
} }
if !isExist { 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) f := obj.(*fv1.Function)
functionMap[f.ObjectMeta.Name] = f functionMap[f.ObjectMeta.Name] = f
sumPrefix = sumPrefix + functionWeight sumPrefix = sumPrefix + functionWeight
+81 -76
View File
@@ -32,7 +32,6 @@ import (
executorClient "github.com/fission/fission/pkg/executor/client" executorClient "github.com/fission/fission/pkg/executor/client"
config "github.com/fission/fission/pkg/featureconfig" config "github.com/fission/fission/pkg/featureconfig"
"github.com/fission/fission/pkg/generated/clientset/versioned" "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/throttler"
"github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils"
"github.com/fission/fission/pkg/utils/metrics" "github.com/fission/fission/pkg/utils/metrics"
@@ -50,9 +49,9 @@ type HTTPTriggerSet struct {
executor *executorClient.Client executor *executorClient.Client
resolver *functionReferenceResolver resolver *functionReferenceResolver
triggers []fv1.HTTPTrigger triggers []fv1.HTTPTrigger
triggerInformer k8sCache.SharedIndexInformer triggerInformer map[string]k8sCache.SharedIndexInformer
functions []fv1.Function functions []fv1.Function
funcInformer k8sCache.SharedIndexInformer funcInformer map[string]k8sCache.SharedIndexInformer
updateRouterRequestChannel chan struct{} updateRouterRequestChannel chan struct{}
tsRoundTripperParams *tsRoundTripperParams tsRoundTripperParams *tsRoundTripperParams
isDebugEnv bool isDebugEnv bool
@@ -76,18 +75,15 @@ func makeHTTPTriggerSet(logger *zap.Logger, fmap *functionServiceMap, fissionCli
svcAddrUpdateThrottler: actionThrottler, svcAddrUpdateThrottler: actionThrottler,
unTapServiceTimeout: unTapServiceTimeout, unTapServiceTimeout: unTapServiceTimeout,
} }
httpTriggerSet.triggerInformer = utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.HttpTriggerResource)
informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30) httpTriggerSet.funcInformer = utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.FunctionResource)
httpTriggerSet.triggerInformer = informerFactory.Core().V1().HTTPTriggers().Informer()
httpTriggerSet.funcInformer = informerFactory.Core().V1().Functions().Informer()
httpTriggerSet.addTriggerHandlers() httpTriggerSet.addTriggerHandlers()
httpTriggerSet.addFunctionHandlers() httpTriggerSet.addFunctionHandlers()
return httpTriggerSet return httpTriggerSet
} }
func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mr *mutableRouter) { func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mr *mutableRouter) {
resolver := makeFunctionReferenceResolver(&ts.funcInformer) resolver := makeFunctionReferenceResolver(ts.logger, ts.funcInformer)
ts.resolver = resolver ts.resolver = resolver
ts.mutableRouter = mr ts.mutableRouter = mr
@@ -280,70 +276,75 @@ func (ts *HTTPTriggerSet) updateTriggerStatusFailed(ht *fv1.HTTPTrigger, err err
} }
func (ts *HTTPTriggerSet) addTriggerHandlers() { func (ts *HTTPTriggerSet) addTriggerHandlers() {
ts.triggerInformer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ for _, triggerInformer := range ts.triggerInformer {
AddFunc: func(obj interface{}) { triggerInformer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
trigger := obj.(*fv1.HTTPTrigger) AddFunc: func(obj interface{}) {
go createIngress(context.Background(), ts.logger, trigger, ts.kubeClient) trigger := obj.(*fv1.HTTPTrigger)
ts.syncTriggers() go createIngress(context.Background(), ts.logger, trigger, ts.kubeClient)
}, ts.syncTriggers()
DeleteFunc: func(obj interface{}) { },
ts.syncTriggers() DeleteFunc: func(obj interface{}) {
trigger := obj.(*fv1.HTTPTrigger) ts.syncTriggers()
go deleteIngress(context.Background(), ts.logger, trigger, ts.kubeClient) trigger := obj.(*fv1.HTTPTrigger)
}, go deleteIngress(context.Background(), ts.logger, trigger, ts.kubeClient)
UpdateFunc: func(oldObj interface{}, newObj interface{}) { },
oldTrigger := oldObj.(*fv1.HTTPTrigger) UpdateFunc: func(oldObj interface{}, newObj interface{}) {
newTrigger := newObj.(*fv1.HTTPTrigger) oldTrigger := oldObj.(*fv1.HTTPTrigger)
newTrigger := newObj.(*fv1.HTTPTrigger)
if oldTrigger.ObjectMeta.ResourceVersion == newTrigger.ObjectMeta.ResourceVersion { if oldTrigger.ObjectMeta.ResourceVersion == newTrigger.ObjectMeta.ResourceVersion {
return return
} }
go updateIngress(context.Background(), ts.logger, oldTrigger, newTrigger, ts.kubeClient) go updateIngress(context.Background(), ts.logger, oldTrigger, newTrigger, ts.kubeClient)
ts.syncTriggers() ts.syncTriggers()
}, },
}) })
}
} }
func (ts *HTTPTriggerSet) addFunctionHandlers() { func (ts *HTTPTriggerSet) addFunctionHandlers() {
ts.funcInformer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ for _, funcInformer := range ts.funcInformer {
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 { funcInformer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
return 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 if oldFn.ObjectMeta.ResourceVersion == fn.ObjectMeta.ResourceVersion {
for key, rr := range ts.resolver.copy() { return
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() // 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) { func (ts *HTTPTriggerSet) runInformer(ctx context.Context, informer map[string]k8sCache.SharedIndexInformer) {
go func() { for _, inf := range informer {
informer.Run(ctx.Done()) go inf.Run(ctx.Done())
}() }
} }
func (ts *HTTPTriggerSet) syncTriggers() { func (ts *HTTPTriggerSet) syncTriggers() {
@@ -353,23 +354,27 @@ func (ts *HTTPTriggerSet) syncTriggers() {
func (ts *HTTPTriggerSet) updateRouter() { func (ts *HTTPTriggerSet) updateRouter() {
for range ts.updateRouterRequestChannel { for range ts.updateRouterRequestChannel {
// get triggers // get triggers
latestTriggers := ts.triggerInformer.GetStore().List() alltriggers := make([]fv1.HTTPTrigger, 0)
triggers := make([]fv1.HTTPTrigger, len(latestTriggers)) for _, triggerInformer := range ts.triggerInformer {
for _, t := range latestTriggers { latestTriggers := triggerInformer.GetStore().List()
triggers = append(triggers, *t.(*fv1.HTTPTrigger)) for _, t := range latestTriggers {
alltriggers = append(alltriggers, *t.(*fv1.HTTPTrigger))
}
} }
ts.triggers = triggers ts.triggers = alltriggers
// get functions // get functions
latestFunctions := ts.funcInformer.GetStore().List() allfunctions := make([]fv1.Function, 0)
functionTimeout := make(map[types.UID]int, len(latestFunctions)) functionTimeout := make(map[types.UID]int, 0)
functions := make([]fv1.Function, len(latestFunctions)) for _, funcInformer := range ts.funcInformer {
for _, f := range latestFunctions { latestFunctions := funcInformer.GetStore().List()
fn := *f.(*fv1.Function) for _, f := range latestFunctions {
functionTimeout[fn.ObjectMeta.UID] = fn.Spec.FunctionTimeout fn := *f.(*fv1.Function)
functions = append(functions, *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 // make a new router and use it
ts.mutableRouter.updateRouter(ts.getRouter(functionTimeout)) ts.mutableRouter.updateRouter(ts.getRouter(functionTimeout))
+56
View File
@@ -1,6 +1,8 @@
package utils package utils
import ( import (
"os"
"strings"
"time" "time"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -8,11 +10,65 @@ import (
"k8s.io/apimachinery/pkg/selection" "k8s.io/apimachinery/pkg/selection"
k8sInformers "k8s.io/client-go/informers" k8sInformers "k8s.io/client-go/informers"
"k8s.io/client-go/kubernetes" "k8s.io/client-go/kubernetes"
"k8s.io/client-go/tools/cache"
metricsapi "k8s.io/metrics/pkg/apis/metrics" metricsapi "k8s.io/metrics/pkg/apis/metrics"
v1 "github.com/fission/fission/pkg/apis/core/v1" 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) { func GetInformerFactoryByReadyPod(client kubernetes.Interface, namespace string, labelSelector *metav1.LabelSelector) (k8sInformers.SharedInformerFactory, error) {
informerFactory := k8sInformers.NewSharedInformerFactoryWithOptions(client, 0, informerFactory := k8sInformers.NewSharedInformerFactoryWithOptions(client, 0,
k8sInformers.WithNamespace(namespace), k8sInformers.WithNamespace(namespace),