From 6d117ad43aa83df8f0ce0fc774385c870000c351 Mon Sep 17 00:00:00 2001 From: Shubham Bansal <62992590+shubham-bansal96@users.noreply.github.com> Date: Wed, 16 Nov 2022 22:16:05 +0530 Subject: [PATCH] Allow empty namespace for fission function and builder (#2621) Currently, we create Fission resources in the default namespace, function-related resources are created in the fission-function namespace, whereas builder resources are created in the fission-builder namespace. This causes confusion for a lot of users. In this fix, we allow the user to set the function and builder namespace empty so that function and builder resources are created in the same namespace as the function resource always. If the user desires older behaviour they can functionNamespace and builderNamespace the same previous before the upgrade. * use default namespace for fission function and builder * support for existing fission namespaces * Replace builder and function namespace with template * Fix namespace creation template Signed-off-by: Sanket Sudake --- .../templates/_function-access-role.tpl | 4 +- charts/fission-all/templates/_helpers.tpl | 18 ++- .../templates/buildermgr/deployment.yaml | 8 +- .../templates/controller/deployment.yaml | 6 +- .../templates/executor/deployment.yaml | 8 +- .../templates/misc-functions/namespace.yaml | 13 +- .../templates/misc-functions/role.yaml | 2 +- .../templates/misc-functions/rolebinding.yaml | 4 +- .../misc-functions/serviceaccount.yaml | 4 +- .../pre-upgrade-checks/pre-upgrade-job.yaml | 1 - charts/fission-all/values.yaml | 22 ++-- cmd/fission-bundle/main.go | 15 +-- cmd/preupgradechecks/checks.go | 6 +- cmd/preupgradechecks/main.go | 27 +--- pkg/buildermgr/buildermgr.go | 6 +- pkg/buildermgr/envwatcher.go | 27 ++-- pkg/buildermgr/pkgwatcher.go | 41 +++--- pkg/controller/api.go | 9 +- pkg/executor/executor.go | 11 +- pkg/executor/executor_test.go | 2 +- .../executortype/container/containermgr.go | 30 ++--- .../executortype/newdeploy/newdeploymgr.go | 34 ++--- .../newdeploy/newdeploymgr_test.go | 9 +- pkg/executor/executortype/poolmgr/gp.go | 5 +- pkg/executor/executortype/poolmgr/gpm.go | 26 ++-- .../executortype/poolmgr/poolpodcontroller.go | 12 +- .../poolmgr/poolpodcontroller_test.go | 5 +- pkg/executor/reaper/reaper.go | 10 +- pkg/utils/namespace.go | 63 +++++++++ pkg/utils/namespace_test.go | 124 ++++++++++++++++++ test/kind_CI.sh | 4 +- 31 files changed, 357 insertions(+), 199 deletions(-) create mode 100644 pkg/utils/namespace.go create mode 100644 pkg/utils/namespace_test.go diff --git a/charts/fission-all/templates/_function-access-role.tpl b/charts/fission-all/templates/_function-access-role.tpl index 2485136f..a5235a61 100644 --- a/charts/fission-all/templates/_function-access-role.tpl +++ b/charts/fission-all/templates/_function-access-role.tpl @@ -61,7 +61,7 @@ roleRef: subjects: - kind: ServiceAccount name: fission-fetcher - namespace: {{ .Values.functionNamespace }} + namespace: {{template "fission-function-ns" . }} --- apiVersion: rbac.authorization.k8s.io/v1 kind: RoleBinding @@ -75,5 +75,5 @@ roleRef: subjects: - kind: ServiceAccount name: fission-builder - namespace: {{ .Values.builderNamespace }} + namespace: {{ template "fission-builder-ns" . }} {{- end -}} diff --git a/charts/fission-all/templates/_helpers.tpl b/charts/fission-all/templates/_helpers.tpl index 908b035f..60f4fa4e 100644 --- a/charts/fission-all/templates/_helpers.tpl +++ b/charts/fission-all/templates/_helpers.tpl @@ -86,4 +86,20 @@ Define the svc's name */}} {{- define "fission-webhook.svc" -}} {{- printf "webhook-service" -}} -{{- end -}} \ No newline at end of file +{{- end -}} + +{{- define "fission-function-ns" -}} +{{- if .Values.functionNamespace -}} +{{- printf "%s" .Values.functionNamespace -}} +{{- else -}} +{{- printf "%s" .Values.defaultNamespace -}} +{{- end -}} +{{- end -}} + +{{- define "fission-builder-ns" -}} +{{- if .Values.builderNamespace -}} +{{- printf "%s" .Values.builderNamespace -}} +{{- else -}} +{{- printf "%s" .Values.builderNamespace -}} +{{- end -}} +{{- end -}} diff --git a/charts/fission-all/templates/buildermgr/deployment.yaml b/charts/fission-all/templates/buildermgr/deployment.yaml index f6b8f761..274c7bb3 100644 --- a/charts/fission-all/templates/buildermgr/deployment.yaml +++ b/charts/fission-all/templates/buildermgr/deployment.yaml @@ -27,7 +27,7 @@ spec: image: {{ include "fission-bundleImage" . | quote }} imagePullPolicy: {{ .Values.pullPolicy }} command: ["/fission-bundle"] - args: ["--builderMgr", "--storageSvcUrl", "http://storagesvc.{{ .Release.Namespace }}", "--envbuilder-namespace", "{{ .Values.builderNamespace }}"] + args: ["--builderMgr", "--storageSvcUrl", "http://storagesvc.{{ .Release.Namespace }}"] env: - name: FETCHER_IMAGE {{- if eq .Values.fetcher.imageTag "" }} @@ -39,6 +39,12 @@ spec: value: "{{ .Values.pullPolicy }}" - name: BUILDER_IMAGE_PULL_POLICY value: "{{ .Values.pullPolicy }}" + - name: FISSION_BUILDER_NAMESPACE + value: "{{ .Values.builderNamespace }}" + - name: FISSION_FUNCTION_NAMESPACE + value: "{{ .Values.functionNamespace }}" + - name: FISSION_DEFAULT_NAMESPACE + value: "{{ .Values.defaultNamespace }}" - name: ENABLE_ISTIO value: "{{ .Values.enableIstio }}" - name: FETCHER_MINCPU diff --git a/charts/fission-all/templates/controller/deployment.yaml b/charts/fission-all/templates/controller/deployment.yaml index 69aee70f..8df20d76 100644 --- a/charts/fission-all/templates/controller/deployment.yaml +++ b/charts/fission-all/templates/controller/deployment.yaml @@ -33,8 +33,12 @@ spec: command: ["/fission-bundle"] args: ["--controllerPort", "8888"] env: + - name: FISSION_DEFAULT_NAMESPACE + value: "{{ .Values.defaultNamespace }}" + - name: FISSION_BUILDER_NAMESPACE + value: "{{ .Values.builderNamespace }}" - name: FISSION_FUNCTION_NAMESPACE - value: "{{ .Values.functionNamespace }}" + value: "{{ .Values.functionNamespace }}" - name: DEBUG_ENV value: {{ .Values.debugEnv | quote }} - name: PPROF_ENABLED diff --git a/charts/fission-all/templates/executor/deployment.yaml b/charts/fission-all/templates/executor/deployment.yaml index ed2f24c0..14a11178 100644 --- a/charts/fission-all/templates/executor/deployment.yaml +++ b/charts/fission-all/templates/executor/deployment.yaml @@ -27,7 +27,7 @@ spec: image: {{ include "fission-bundleImage" . | quote }} imagePullPolicy: {{ .Values.pullPolicy }} command: ["/fission-bundle"] - args: ["--executorPort", "8888", "--namespace", "{{ .Values.functionNamespace }}"] + args: ["--executorPort", "8888"] env: - name: FETCHER_IMAGE {{- if eq .Values.fetcher.imageTag "" }} @@ -37,6 +37,12 @@ spec: {{- end }} - name: FETCHER_IMAGE_PULL_POLICY value: "{{ .Values.pullPolicy }}" + - name: FISSION_BUILDER_NAMESPACE + value: "{{ .Values.builderNamespace }}" + - name: FISSION_FUNCTION_NAMESPACE + value: "{{ .Values.functionNamespace }}" + - name: FISSION_DEFAULT_NAMESPACE + value: "{{ .Values.defaultNamespace }}" - name: RUNTIME_IMAGE_PULL_POLICY value: "{{ .Values.pullPolicy }}" - name: ADOPT_EXISTING_RESOURCES diff --git a/charts/fission-all/templates/misc-functions/namespace.yaml b/charts/fission-all/templates/misc-functions/namespace.yaml index 61a5f6bb..bf9bf0e7 100644 --- a/charts/fission-all/templates/misc-functions/namespace.yaml +++ b/charts/fission-all/templates/misc-functions/namespace.yaml @@ -1,24 +1,29 @@ {{- if .Values.createNamespace }} +{{- if and (ne .Values.functionNamespace "default") (ne .Values.functionNamespace "") }} apiVersion: v1 kind: Namespace metadata: - name: {{ .Values.functionNamespace }} + name: {{ template "fission-function-ns" . }} labels: name: fission-function - chart: "{{ .Chart.Name }}-{{ .Chart.Version }}" + chart: {{ .Chart.Name }}-{{ .Chart.Version }} {{- if .Values.enableIstio }} istio-injection: enabled {{- end }} +{{- end}} --- + +{{- if and (ne .Values.builderNamespace "default") (ne .Values.builderNamespace "") }} apiVersion: v1 kind: Namespace metadata: - name: {{ .Values.builderNamespace }} + name: {{ template "fission-builder-ns" . }} labels: name: fission-builder - chart: "{{ .Chart.Name }}-{{ .Chart.Version }}" + chart: {{ .Chart.Name }}-{{ .Chart.Version }} {{- if .Values.enableIstio }} istio-injection: enabled {{- end }} +{{- end }} {{- end }} \ No newline at end of file diff --git a/charts/fission-all/templates/misc-functions/role.yaml b/charts/fission-all/templates/misc-functions/role.yaml index 46241b11..5ea35b12 100644 --- a/charts/fission-all/templates/misc-functions/role.yaml +++ b/charts/fission-all/templates/misc-functions/role.yaml @@ -8,7 +8,7 @@ can be used apiVersion: rbac.authorization.k8s.io/v1 kind: Role metadata: - namespace: {{ .Values.functionNamespace }} + namespace: {{ template "fission-function-ns" . }} name: {{ .Release.Name }}-event-fetcher rules: - apiGroups: diff --git a/charts/fission-all/templates/misc-functions/rolebinding.yaml b/charts/fission-all/templates/misc-functions/rolebinding.yaml index 019c7b1a..235701cd 100644 --- a/charts/fission-all/templates/misc-functions/rolebinding.yaml +++ b/charts/fission-all/templates/misc-functions/rolebinding.yaml @@ -9,7 +9,7 @@ apiVersion: rbac.authorization.k8s.io/v1 kind: RoleBinding metadata: name: {{ .Release.Name }}-fission-fetcher-pod-reader - namespace: {{ .Values.functionNamespace }} + namespace: {{ template "fission-function-ns" . }} roleRef: apiGroup: rbac.authorization.k8s.io kind: Role @@ -17,7 +17,7 @@ roleRef: subjects: - kind: ServiceAccount name: fission-fetcher - namespace: {{ .Values.functionNamespace }} + namespace: {{ template "fission-function-ns" . }} {{- if not .Values.singleDefaultNamespace }} {{- range $namespace := $.Values.additionalFissionNamespaces }} diff --git a/charts/fission-all/templates/misc-functions/serviceaccount.yaml b/charts/fission-all/templates/misc-functions/serviceaccount.yaml index 6a550352..f54b7764 100644 --- a/charts/fission-all/templates/misc-functions/serviceaccount.yaml +++ b/charts/fission-all/templates/misc-functions/serviceaccount.yaml @@ -2,11 +2,11 @@ apiVersion: v1 kind: ServiceAccount metadata: name: fission-fetcher - namespace: {{ .Values.functionNamespace }} + namespace: {{ template "fission-function-ns" . }} --- apiVersion: v1 kind: ServiceAccount metadata: name: fission-builder - namespace: {{ .Values.builderNamespace }} + namespace: {{ template "fission-builder-ns" . }} diff --git a/charts/fission-all/templates/pre-upgrade-checks/pre-upgrade-job.yaml b/charts/fission-all/templates/pre-upgrade-checks/pre-upgrade-job.yaml index 350f6419..49e20980 100644 --- a/charts/fission-all/templates/pre-upgrade-checks/pre-upgrade-job.yaml +++ b/charts/fission-all/templates/pre-upgrade-checks/pre-upgrade-job.yaml @@ -35,7 +35,6 @@ spec: {{- end }} imagePullPolicy: {{ .Values.pullPolicy }} command: [ "/pre-upgrade-checks" ] - args: ["--fn-pod-namespace", "{{ .Values.functionNamespace }}", "--envbuilder-namespace", "{{ .Values.builderNamespace }}"] env: {{- include "fission-resource-namespace.envs" . | indent 8 }} {{- if .Values.terminationMessagePath }} diff --git a/charts/fission-all/values.yaml b/charts/fission-all/values.yaml index 03a2997c..8f3f49b8 100644 --- a/charts/fission-all/values.yaml +++ b/charts/fission-all/values.yaml @@ -61,16 +61,6 @@ controllerPort: 31313 ## routerPort: 31314 -## functionNamespace represents the namespace in which Fission Function resources will be created. -## This is different from the release namespace. -## -functionNamespace: fission-function - -## builderNamespace represents the namespace in which Fission Builder resources will be created. -## This is different from the release namespace. -## -builderNamespace: fission-builder - ## 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 @@ -79,6 +69,18 @@ builderNamespace: fission-builder ## defaultNamespace: default +## builderNamespace represents the namespace in which Fission Builder resources will be created. +## if builderNamespace is set to empty then builder resources will be created in the same namespace as the Fission resources. +## This is different from the release namespace. +## +builderNamespace: "" + +## functionNamespace represents the namespace in which Fission Function resources will be created. +## if functionNamespace is set to empty then function resources will be created in the same namespace as the Fission resources. +## This is different from the release namespace. +## +functionNamespace: "" + ## If true, fission will only watch for fission custom resources created in the `defaultNamespace` above. ## singleDefaultNamespace: true diff --git a/cmd/fission-bundle/main.go b/cmd/fission-bundle/main.go index 53491918..1d345d27 100644 --- a/cmd/fission-bundle/main.go +++ b/cmd/fission-bundle/main.go @@ -62,8 +62,8 @@ func runRouter(ctx context.Context, logger *zap.Logger, port int, executorUrl st router.Start(ctx, logger, port, executorUrl) } -func runExecutor(ctx context.Context, logger *zap.Logger, port int, functionNamespace, envBuilderNamespace string) error { - return executor.StartExecutor(ctx, logger, functionNamespace, envBuilderNamespace, port) +func runExecutor(ctx context.Context, logger *zap.Logger, port int) error { + return executor.StartExecutor(ctx, logger, port) } func runKubeWatcher(ctx context.Context, logger *zap.Logger, routerUrl string) error { @@ -87,8 +87,8 @@ func runStorageSvc(ctx context.Context, logger *zap.Logger, port int, storage st return storagesvc.Start(ctx, logger, storage, port) } -func runBuilderMgr(ctx context.Context, logger *zap.Logger, storageSvcUrl string, envBuilderNamespace string) error { - return buildermgr.Start(ctx, logger, storageSvcUrl, envBuilderNamespace) +func runBuilderMgr(ctx context.Context, logger *zap.Logger, storageSvcUrl string) error { + return buildermgr.Start(ctx, logger, storageSvcUrl) } func runLogger(ctx context.Context, logger *zap.Logger) { @@ -232,9 +232,6 @@ Options: defer shutdown(ctx) } - functionNs := getStringArgWithDefault(arguments["--namespace"], "fission-function") - envBuilderNs := getStringArgWithDefault(arguments["--envbuilder-namespace"], "fission-builder") - executorUrl := getStringArgWithDefault(arguments["--executorUrl"], "http://executor.fission") routerUrl := getStringArgWithDefault(arguments["--routerUrl"], "http://router.fission") storageSvcUrl := getStringArgWithDefault(arguments["--storageSvcUrl"], "http://storagesvc.fission") @@ -270,7 +267,7 @@ Options: if arguments["--executorPort"] != nil { port := getPort(logger, arguments["--executorPort"]) - err = runExecutor(ctx, logger, port, functionNs, envBuilderNs) + err = runExecutor(ctx, logger, port) if err != nil { logger.Error("executor exited", zap.Error(err)) return @@ -310,7 +307,7 @@ Options: } if arguments["--builderMgr"] == true { - err = runBuilderMgr(ctx, logger, storageSvcUrl, envBuilderNs) + err = runBuilderMgr(ctx, logger, storageSvcUrl) if err != nil { logger.Error("builder manager exited", zap.Error(err)) return diff --git a/cmd/preupgradechecks/checks.go b/cmd/preupgradechecks/checks.go index ee76d68c..f3474b9e 100644 --- a/cmd/preupgradechecks/checks.go +++ b/cmd/preupgradechecks/checks.go @@ -42,8 +42,6 @@ type ( fissionClient versioned.Interface k8sClient kubernetes.Interface apiExtClient apiextensionsclient.Interface - fnPodNs string - envBuilderNs string } ) @@ -53,7 +51,7 @@ const ( MqtCRD = "messagequeuetriggers.fission.io" ) -func makePreUpgradeTaskClient(logger *zap.Logger, fnPodNs, envBuilderNs string) (*PreUpgradeTaskClient, error) { +func makePreUpgradeTaskClient(logger *zap.Logger) (*PreUpgradeTaskClient, error) { fissionClient, k8sClient, apiExtClient, _, err := crd.MakeFissionClient() if err != nil { return nil, errors.Wrap(err, "error making fission client") @@ -63,8 +61,6 @@ func makePreUpgradeTaskClient(logger *zap.Logger, fnPodNs, envBuilderNs string) logger: logger.Named("pre_upgrade_task_client"), fissionClient: fissionClient, k8sClient: k8sClient, - fnPodNs: fnPodNs, - envBuilderNs: envBuilderNs, apiExtClient: apiExtClient, }, nil } diff --git a/cmd/preupgradechecks/main.go b/cmd/preupgradechecks/main.go index 39115378..624265f4 100644 --- a/cmd/preupgradechecks/main.go +++ b/cmd/preupgradechecks/main.go @@ -17,42 +17,17 @@ limitations under the License. package main import ( - "github.com/docopt/docopt-go" "go.uber.org/zap" "sigs.k8s.io/controller-runtime/pkg/manager/signals" - "github.com/fission/fission/pkg/info" "github.com/fission/fission/pkg/utils/loggerfactory" ) -func getStringArgWithDefault(arg interface{}, defaultValue string) string { - if arg != nil { - return arg.(string) - } else { - return defaultValue - } -} - func main() { logger := loggerfactory.GetLogger() defer logger.Sync() - usage := `Package to perform operations needed prior to fission installation -Usage: - pre-upgrade-checks --fn-pod-namespace= --envbuilder-namespace= -Options: - --fn-pod-namespace= Namespace where function pods get deployed. - --envbuilder-namespace= Namespace where builder env pods are deployed.` - - arguments, err := docopt.ParseArgs(usage, nil, info.BuildInfo().String()) - if err != nil { - logger.Fatal("Could not parse command line arguments", zap.Error(err)) - } - - functionPodNs := getStringArgWithDefault(arguments["--fn-pod-namespace"], "fission-function") - envBuilderNs := getStringArgWithDefault(arguments["--envbuilder-namespace"], "fission-builder") - - crdBackedClient, err := makePreUpgradeTaskClient(logger, functionPodNs, envBuilderNs) + crdBackedClient, err := makePreUpgradeTaskClient(logger) if err != nil { logger.Fatal("error creating a crd client, please retry helm upgrade", zap.Error(err)) diff --git a/pkg/buildermgr/buildermgr.go b/pkg/buildermgr/buildermgr.go index 7244657f..38ea40f5 100644 --- a/pkg/buildermgr/buildermgr.go +++ b/pkg/buildermgr/buildermgr.go @@ -33,7 +33,7 @@ import ( ) // Start the buildermgr service. -func Start(ctx context.Context, logger *zap.Logger, storageSvcUrl string, envBuilderNamespace string) error { +func Start(ctx context.Context, logger *zap.Logger, storageSvcUrl string) error { bmLogger := logger.Named("builder_manager") fissionClient, kubernetesClient, _, _, err := crd.MakeFissionClient() @@ -62,13 +62,13 @@ func Start(ctx context.Context, logger *zap.Logger, storageSvcUrl string, envBui } } - envWatcher := makeEnvironmentWatcher(ctx, bmLogger, fissionClient, kubernetesClient, fetcherConfig, envBuilderNamespace, podSpecPatch) + envWatcher := makeEnvironmentWatcher(ctx, bmLogger, fissionClient, kubernetesClient, fetcherConfig, podSpecPatch) envWatcher.Run(ctx) k8sInformerFactory := k8sInformers.NewSharedInformerFactory(kubernetesClient, time.Minute*30) podInformer := k8sInformerFactory.Core().V1().Pods().Informer() pkgWatcher := makePackageWatcher(bmLogger, fissionClient, - kubernetesClient, envBuilderNamespace, storageSvcUrl, podInformer, + kubernetesClient, storageSvcUrl, podInformer, utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.PackagesResource)) pkgWatcher.Run(ctx) return nil diff --git a/pkg/buildermgr/envwatcher.go b/pkg/buildermgr/envwatcher.go index bba2ed90..1e82c288 100644 --- a/pkg/buildermgr/envwatcher.go +++ b/pkg/buildermgr/envwatcher.go @@ -64,9 +64,9 @@ type ( environmentWatcher struct { logger *zap.Logger cache map[string]*builderInfo - builderNamespace string fissionClient versioned.Interface kubernetesClient kubernetes.Interface + nsResolver *utils.NamespaceResolver fetcherConfig *fetcherConfig.Config builderImagePullPolicy apiv1.PullPolicy useIstio bool @@ -81,7 +81,6 @@ func makeEnvironmentWatcher( fissionClient versioned.Interface, kubernetesClient kubernetes.Interface, fetcherConfig *fetcherConfig.Config, - builderNamespace string, podSpecPatch *apiv1.PodSpec) *environmentWatcher { useIstio := false @@ -99,9 +98,9 @@ func makeEnvironmentWatcher( envWatcher := &environmentWatcher{ logger: logger.Named("environment_watcher"), cache: make(map[string]*builderInfo), - builderNamespace: builderNamespace, fissionClient: fissionClient, kubernetesClient: kubernetesClient, + nsResolver: utils.DefaultNSResolver(), builderImagePullPolicy: builderImagePullPolicy, useIstio: useIstio, fetcherConfig: fetcherConfig, @@ -161,7 +160,7 @@ func (envw *environmentWatcher) AddUpdateBuilder(ctx context.Context, env *fv1.E //builder is not supported with v1 interface and ignore env without builder image if env.Spec.Version != 1 && len(env.Spec.Builder.Image) != 0 { if _, ok := envw.cache[crd.CacheKeyUID(&env.ObjectMeta)]; !ok { - builderInfo, err := envw.createBuilder(ctx, env, envw.getNamespace(env)) + builderInfo, err := envw.createBuilder(ctx, env, envw.nsResolver.GetBuilderNS(env.ObjectMeta.Namespace)) if err != nil { envw.logger.Error("error creating builder service", zap.Error(err)) return @@ -170,7 +169,7 @@ func (envw *environmentWatcher) AddUpdateBuilder(ctx context.Context, env *fv1.E } else { envw.DeleteBuilder(ctx, env) // once older builder deleted then add new builder service - builderInfo, err := envw.createBuilder(ctx, env, envw.getNamespace(env)) + builderInfo, err := envw.createBuilder(ctx, env, envw.nsResolver.GetBuilderNS(env.ObjectMeta.Namespace)) if err != nil { envw.logger.Error("error updating builder service", zap.Error(err)) return @@ -185,24 +184,14 @@ func (envw *environmentWatcher) DeleteBuilder(ctx context.Context, env *fv1.Envi envw.DeleteBuilderService(ctx, env) envw.DeleteBuilderDeployment(ctx, env) delete(envw.cache, crd.CacheKeyUID(&env.ObjectMeta)) - envw.logger.Info("builder service deleted", zap.String("env_name", env.ObjectMeta.Name), zap.String("namespace", envw.getNamespace(env))) + envw.logger.Info("builder service deleted", zap.String("env_name", env.ObjectMeta.Name), zap.String("namespace", envw.nsResolver.GetBuilderNS(env.ObjectMeta.Namespace))) } else { - envw.logger.Debug("builder service not found", zap.String("env_name", env.ObjectMeta.Name), zap.String("namespace", envw.getNamespace(env))) + envw.logger.Debug("builder service not found", zap.String("env_name", env.ObjectMeta.Name), zap.String("namespace", envw.nsResolver.GetBuilderNS(env.ObjectMeta.Namespace))) } } -func (envw *environmentWatcher) getNamespace(env *fv1.Environment) string { - // In order to support backward compatibility, for all environments with builder image created in default env, - // the pods will be created in fission-builder namespace - ns := envw.builderNamespace - if env.ObjectMeta.Namespace != metav1.NamespaceDefault { - ns = env.ObjectMeta.Namespace - } - return ns -} - func (envw *environmentWatcher) DeleteBuilderService(ctx context.Context, env *fv1.Environment) { - ns := envw.getNamespace(env) + ns := envw.nsResolver.GetBuilderNS(env.ObjectMeta.Namespace) svcList, err := envw.getBuilderServiceList(ctx, envw.getDeploymentLabels(env.ObjectMeta.Name), ns) if err != nil { envw.logger.Error("error getting the builder service list", zap.Error(err)) @@ -227,7 +216,7 @@ func (envw *environmentWatcher) DeleteBuilderService(ctx context.Context, env *f } func (envw *environmentWatcher) DeleteBuilderDeployment(ctx context.Context, env *fv1.Environment) { - ns := envw.getNamespace(env) + ns := envw.nsResolver.GetBuilderNS(env.ObjectMeta.Namespace) deployList, err := envw.getBuilderDeploymentList(ctx, envw.getDeploymentLabels(env.ObjectMeta.Name), ns) if err != nil { envw.logger.Error("error getting the builder deployment list", zap.Error(err)) diff --git a/pkg/buildermgr/pkgwatcher.go b/pkg/buildermgr/pkgwatcher.go index 8d4d0e5a..3e3b94c7 100644 --- a/pkg/buildermgr/pkgwatcher.go +++ b/pkg/buildermgr/pkgwatcher.go @@ -37,29 +37,29 @@ import ( type ( packageWatcher struct { - logger *zap.Logger - fissionClient versioned.Interface - k8sClient kubernetes.Interface - podInformer k8sCache.SharedIndexInformer - pkgInformer map[string]k8sCache.SharedIndexInformer - builderNamespace string - storageSvcUrl string - buildCache *cache.Cache + logger *zap.Logger + fissionClient versioned.Interface + nsResolver *utils.NamespaceResolver + k8sClient kubernetes.Interface + podInformer k8sCache.SharedIndexInformer + pkgInformer map[string]k8sCache.SharedIndexInformer + storageSvcUrl string + buildCache *cache.Cache } ) func makePackageWatcher(logger *zap.Logger, fissionClient versioned.Interface, k8sClientSet kubernetes.Interface, - builderNamespace string, storageSvcUrl string, podInformer k8sCache.SharedIndexInformer, + storageSvcUrl string, podInformer k8sCache.SharedIndexInformer, pkgInformer map[string]k8sCache.SharedIndexInformer) *packageWatcher { pkgw := &packageWatcher{ - logger: logger.Named("package_watcher"), - fissionClient: fissionClient, - k8sClient: k8sClientSet, - podInformer: podInformer, - pkgInformer: pkgInformer, - builderNamespace: builderNamespace, - storageSvcUrl: storageSvcUrl, - buildCache: cache.MakeCache(0, 0), + logger: logger.Named("package_watcher"), + fissionClient: fissionClient, + k8sClient: k8sClientSet, + nsResolver: utils.DefaultNSResolver(), + podInformer: podInformer, + pkgInformer: pkgInformer, + storageSvcUrl: storageSvcUrl, + buildCache: cache.MakeCache(0, 0), } return pkgw } @@ -137,12 +137,7 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) { for _, item := range items { pod := item.(*apiv1.Pod) - // In order to support backward compatibility, for all builder images created in default env, - // the pods will be created in fission-builder namespace - builderNs := pkgw.builderNamespace - if env.ObjectMeta.Namespace != metav1.NamespaceDefault { - builderNs = env.ObjectMeta.Namespace - } + builderNs := pkgw.nsResolver.GetBuilderNS(env.ObjectMeta.Namespace) // Filter non-matching pods if pod.ObjectMeta.Labels[LABEL_ENV_NAME] != env.ObjectMeta.Name || diff --git a/pkg/controller/api.go b/pkg/controller/api.go index d799f280..07f63384 100644 --- a/pkg/controller/api.go +++ b/pkg/controller/api.go @@ -34,6 +34,7 @@ import ( "github.com/fission/fission/pkg/fission-cli/logdb" "github.com/fission/fission/pkg/generated/clientset/versioned" "github.com/fission/fission/pkg/info" + "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/httpserver" "github.com/fission/fission/pkg/utils/metrics" "github.com/fission/fission/pkg/utils/otel" @@ -91,12 +92,8 @@ func MakeAPI(logger *zap.Logger) (*API, error) { api.workflowApiUrl = "http://workflows-apiserver" } - fnNs := os.Getenv("FISSION_FUNCTION_NAMESPACE") - if len(fnNs) > 0 { - api.functionNamespace = fnNs - } else { - api.functionNamespace = "fission-function" - } + nsResolver := utils.DefaultNSResolver() + api.functionNamespace = nsResolver.ResolveNamespace(os.Getenv(utils.ENV_FUNCTION_NAMESPACE)) return api, err } diff --git a/pkg/executor/executor.go b/pkg/executor/executor.go index fba860dd..5ac425e5 100644 --- a/pkg/executor/executor.go +++ b/pkg/executor/executor.go @@ -255,7 +255,7 @@ func (executor *Executor) getFunctionServiceFromCache(ctx context.Context, fn *f // StartExecutor Starts executor and the executor components such as Poolmgr, // deploymgr and potential future executor types -func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace string, envBuilderNamespace string, port int) error { +func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error { fissionClient, kubernetesClient, _, metricsClient, err := crd.MakeFissionClient() if err != nil { return errors.Wrap(err, "failed to get kubernetes client") @@ -306,7 +306,7 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace st gpm, err := poolmgr.MakeGenericPoolManager(ctx, logger, fissionClient, kubernetesClient, metricsClient, - functionNamespace, fetcherConfig, executorInstanceID, + fetcherConfig, executorInstanceID, funcInformer, pkgInformer, envInformer, gpmPodInformer, gpmRsInformer, podSpecPatch) if err != nil { @@ -322,7 +322,7 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace st ndm, err := newdeploy.MakeNewDeploy(ctx, logger, fissionClient, kubernetesClient, - functionNamespace, fetcherConfig, executorInstanceID, + fetcherConfig, executorInstanceID, funcInformer, envInformer, ndmDeplInformer, ndmSvcInformer, podSpecPatch) if err != nil { @@ -338,7 +338,7 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace st cnm, err := container.MakeContainer( ctx, logger, fissionClient, kubernetesClient, - functionNamespace, executorInstanceID, funcInformer, + executorInstanceID, funcInformer, cnmDeplInformer, cnmSvcInformer) if err != nil { return errors.Wrap(err, "container manager creation failed") @@ -408,7 +408,8 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace st if err != nil { return err } - go reaper.CleanupRoleBindings(ctx, logger, kubernetesClient, fissionClient, functionNamespace, envBuilderNamespace, time.Minute*30) + + go reaper.CleanupRoleBindings(ctx, logger, kubernetesClient, fissionClient, time.Minute*30) go metrics.ServeMetrics(ctx, logger) go api.Serve(ctx, port) diff --git a/pkg/executor/executor_test.go b/pkg/executor/executor_test.go index 69479338..8af11808 100644 --- a/pkg/executor/executor_test.go +++ b/pkg/executor/executor_test.go @@ -174,7 +174,7 @@ func TestExecutor(t *testing.T) { // create poolmgr port := 9999 - err = StartExecutor(ctx, logger, functionNs, "fission-builder", port) + err = StartExecutor(ctx, logger, port) if err != nil { log.Panicf("failed to start poolmgr: %v", err) } diff --git a/pkg/executor/executortype/container/containermgr.go b/pkg/executor/executortype/container/containermgr.go index 692a08ba..d84b260f 100644 --- a/pkg/executor/executortype/container/containermgr.go +++ b/pkg/executor/executortype/container/containermgr.go @@ -57,7 +57,9 @@ import ( otelUtils "github.com/fission/fission/pkg/utils/otel" ) -var _ executortype.ExecutorType = &Container{} +var ( + _ executortype.ExecutorType = &Container{} +) type ( // Container represents an executor type @@ -67,10 +69,10 @@ type ( kubernetesClient kubernetes.Interface fissionClient versioned.Interface instanceID string + nsResolver *utils.NamespaceResolver // fetcherConfig *fetcherConfig.Config runtimeImagePullPolicy apiv1.PullPolicy - namespace string useIstio bool fsCache *fscache.FunctionServiceCache // cache funcSvc's by function, address and pod name @@ -96,7 +98,6 @@ func MakeContainer( logger *zap.Logger, fissionClient versioned.Interface, kubernetesClient kubernetes.Interface, - namespace string, instanceID string, funcInformer map[string]finformerv1.FunctionInformer, deplInformer appsinformers.DeploymentInformer, @@ -117,8 +118,8 @@ func MakeContainer( fissionClient: fissionClient, kubernetesClient: kubernetesClient, instanceID: instanceID, + nsResolver: utils.DefaultNSResolver(), - namespace: namespace, fsCache: fscache.MakeFunctionServiceCache(logger), throttler: throttler.MakeThrottler(1 * time.Minute), @@ -129,6 +130,7 @@ func MakeContainer( objectReaperIntervalSecond: time.Duration(executorUtils.GetObjectReaperInterval(logger, fv1.ExecutorTypeContainer, 5)) * time.Second, hpaops: hpautils.NewHpaOperations(logger, kubernetesClient, instanceID), } + caaf.deplLister = deplInformer.Lister() caaf.deplListerSynced = deplInformer.Informer().HasSynced @@ -382,10 +384,7 @@ func (caaf *Container) fnCreate(ctx context.Context, fn *fv1.Function) (*fscache // to support backward compatibility, if the function was created in default ns, we fall back to creating the // deployment of the function in fission-function ns - ns := caaf.namespace - if fn.ObjectMeta.Namespace != metav1.NamespaceDefault { - ns = fn.ObjectMeta.Namespace - } + ns := caaf.nsResolver.GetFunctionNS(fn.ObjectMeta.Namespace) // Envoy(istio-proxy) returns 404 directly before istio pilot // propagates latest Envoy-specific configuration. @@ -499,10 +498,7 @@ func (caaf *Container) updateFunction(ctx context.Context, oldFn *fv1.Function, if !reflect.DeepEqual(oldFn.Spec.InvokeStrategy, newFn.Spec.InvokeStrategy) { // to support backward compatibility, if the function was created in default ns, we fall back to creating the // deployment of the function in fission-function ns, so cleaning up resources there - ns := caaf.namespace - if newFn.ObjectMeta.Namespace != metav1.NamespaceDefault { - ns = newFn.ObjectMeta.Namespace - } + ns := caaf.nsResolver.GetFunctionNS(newFn.ObjectMeta.Namespace) fsvc, err := caaf.fsCache.GetByFunctionUID(newFn.ObjectMeta.UID) if err != nil { @@ -598,10 +594,7 @@ func (caaf *Container) updateFuncDeployment(ctx context.Context, fn *fv1.Functio // to support backward compatibility, if the function was created in default ns, we fall back to creating the // deployment of the function in fission-function ns - ns := caaf.namespace - if fn.ObjectMeta.Namespace != metav1.NamespaceDefault { - ns = fn.ObjectMeta.Namespace - } + ns := caaf.nsResolver.GetFunctionNS(fn.ObjectMeta.Namespace) existingDepl, err := caaf.kubernetesClient.AppsV1().Deployments(ns).Get(ctx, fnObjName, metav1.GetOptions{}) if err != nil { @@ -651,10 +644,7 @@ func (caaf *Container) fnDelete(ctx context.Context, fn *fv1.Function) error { // to support backward compatibility, if the function was created in default ns, we fall back to creating the // deployment of the function in fission-function ns, so cleaning up resources there - ns := caaf.namespace - if fn.ObjectMeta.Namespace != metav1.NamespaceDefault { - ns = fn.ObjectMeta.Namespace - } + ns := caaf.nsResolver.GetFunctionNS(fn.ObjectMeta.Namespace) err = caaf.cleanupContainer(ctx, ns, objName) multierr = multierror.Append(multierr, err) diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr.go b/pkg/executor/executortype/newdeploy/newdeploymgr.go index e8e50384..cf5f7971 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr.go @@ -59,7 +59,9 @@ import ( otelUtils "github.com/fission/fission/pkg/utils/otel" ) -var _ executortype.ExecutorType = &NewDeploy{} +var ( + _ executortype.ExecutorType = &NewDeploy{} +) type ( // NewDeploy represents an ExecutorType @@ -70,9 +72,9 @@ type ( fissionClient versioned.Interface instanceID string fetcherConfig *fetcherConfig.Config + nsResolver *utils.NamespaceResolver runtimeImagePullPolicy apiv1.PullPolicy - namespace string useIstio bool fsCache *fscache.FunctionServiceCache // cache funcSvc's by function, address and pod name @@ -100,7 +102,6 @@ func MakeNewDeploy( logger *zap.Logger, fissionClient versioned.Interface, kubernetesClient kubernetes.Interface, - namespace string, fetcherConfig *fetcherConfig.Config, instanceID string, funcInformer map[string]finformerv1.FunctionInformer, @@ -124,10 +125,9 @@ func MakeNewDeploy( fissionClient: fissionClient, kubernetesClient: kubernetesClient, instanceID: instanceID, - - namespace: namespace, - fsCache: fscache.MakeFunctionServiceCache(logger), - throttler: throttler.MakeThrottler(1 * time.Minute), + fsCache: fscache.MakeFunctionServiceCache(logger), + throttler: throttler.MakeThrottler(1 * time.Minute), + nsResolver: utils.DefaultNSResolver(), fetcherConfig: fetcherConfig, runtimeImagePullPolicy: utils.GetImagePullPolicy(os.Getenv("RUNTIME_IMAGE_PULL_POLICY")), @@ -428,10 +428,7 @@ func (deploy *NewDeploy) fnCreate(ctx context.Context, fn *fv1.Function) (*fscac // to support backward compatibility, if the function was created in default ns, we fall back to creating the // deployment of the function in fission-function ns - ns := deploy.namespace - if fn.ObjectMeta.Namespace != metav1.NamespaceDefault { - ns = fn.ObjectMeta.Namespace - } + ns := deploy.nsResolver.GetFunctionNS(fn.ObjectMeta.Namespace) // Envoy(istio-proxy) returns 404 directly before istio pilot // propagates latest Envoy-specific configuration. @@ -549,10 +546,7 @@ func (deploy *NewDeploy) updateFunction(ctx context.Context, oldFn *fv1.Function // to support backward compatibility, if the function was created in default ns, we fall back to creating the // deployment of the function in fission-function ns, so cleaning up resources there - ns := deploy.namespace - if newFn.ObjectMeta.Namespace != metav1.NamespaceDefault { - ns = newFn.ObjectMeta.Namespace - } + ns := deploy.nsResolver.GetFunctionNS(newFn.ObjectMeta.Namespace) fsvc, err := deploy.fsCache.GetByFunctionUID(newFn.ObjectMeta.UID) if err != nil { @@ -654,10 +648,7 @@ func (deploy *NewDeploy) updateFuncDeployment(ctx context.Context, fn *fv1.Funct // to support backward compatibility, if the function was created in default ns, we fall back to creating the // deployment of the function in fission-function ns - ns := deploy.namespace - if fn.ObjectMeta.Namespace != metav1.NamespaceDefault { - ns = fn.ObjectMeta.Namespace - } + ns := deploy.nsResolver.GetFunctionNS(fn.ObjectMeta.Namespace) existingDepl, err := deploy.kubernetesClient.AppsV1().Deployments(ns).Get(ctx, fnObjName, metav1.GetOptions{}) if err != nil { @@ -708,10 +699,7 @@ func (deploy *NewDeploy) fnDelete(ctx context.Context, fn *fv1.Function) error { // to support backward compatibility, if the function was created in default ns, we fall back to creating the // deployment of the function in fission-function ns, so cleaning up resources there - ns := deploy.namespace - if fn.ObjectMeta.Namespace != metav1.NamespaceDefault { - ns = fn.ObjectMeta.Namespace - } + ns := deploy.nsResolver.GetFunctionNS(fn.ObjectMeta.Namespace) err = deploy.cleanupNewdeploy(ctx, ns, objName) multierr = multierror.Append(multierr, err) diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr_test.go b/pkg/executor/executortype/newdeploy/newdeploymgr_test.go index 2252e1eb..1a1a0059 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr_test.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr_test.go @@ -30,6 +30,7 @@ import ( const ( defaultNamespace string = "default" functionNamespace string = "fission-function" + builderNamespace string = "fission-builder" envName string = "newdeploy-test-env" functionName string = "newdeploy-test-func" configmapName string = "newdeploy-test-configmap" @@ -80,7 +81,7 @@ func TestRefreshFuncPods(t *testing.T) { t.Fatalf("Error creating fetcher config: %s", err) } - executor, err := MakeNewDeploy(ctx, logger, fissionClient, kubernetesClient, functionNamespace, fetcherConfig, "test", + executor, err := MakeNewDeploy(ctx, logger, fissionClient, kubernetesClient, fetcherConfig, "test", funcInformer, envInformer, deployInformer, svcInformer, podSpecPatch) if err != nil { t.Fatalf("new deploy manager creation failed: %s", err) @@ -88,6 +89,12 @@ func TestRefreshFuncPods(t *testing.T) { ndm := executor.(*NewDeploy) + nsResolver := utils.NamespaceResolver{ + FunctionNamespace: functionNamespace, + BuiderNamespace: builderNamespace, + } + ndm.nsResolver = &nsResolver + go ndm.Run(ctx) t.Log("New deploy manager started") diff --git a/pkg/executor/executortype/poolmgr/gp.go b/pkg/executor/executortype/poolmgr/gp.go index 2e63f51a..a6a18288 100644 --- a/pkg/executor/executortype/poolmgr/gp.go +++ b/pkg/executor/executortype/poolmgr/gp.go @@ -62,7 +62,6 @@ type ( env *fv1.Environment deployment *appsv1.Deployment // kubernetes deployment namespace string // namespace to keep our resources - functionNamespace string // fallback namespace for fission functions podReadyTimeout time.Duration // timeout for generic pods to become ready fsCache *fscache.FunctionServiceCache // cache funcSvc's by function, address and podname useSvc bool // create k8s service for specialized pods @@ -92,7 +91,6 @@ func MakeGenericPool( metricsClient metricsclient.Interface, env *fv1.Environment, namespace string, - functionNamespace string, fsCache *fscache.FunctionServiceCache, fetcherConfig *fetcherConfig.Config, instanceID string, @@ -122,7 +120,6 @@ func MakeGenericPool( kubernetesClient: kubernetesClient, metricsClient: metricsClient, namespace: namespace, - functionNamespace: functionNamespace, podReadyTimeout: podReadyTimeout, fsCache: fsCache, fetcherConfig: fetcherConfig, @@ -446,7 +443,7 @@ func (gp *GenericPool) createSvc(ctx context.Context, name string, labels map[st } func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscache.FuncSvc, error) { - logger := otelUtils.LoggerWithTraceID(ctx, gp.logger).With(zap.String("function", fn.ObjectMeta.Name), zap.String("functionNamespace", fn.ObjectMeta.Namespace), + logger := otelUtils.LoggerWithTraceID(ctx, gp.logger).With(zap.String("function", fn.ObjectMeta.Name), zap.String("namespace", fn.ObjectMeta.Namespace), zap.String("env", fn.Spec.Environment.Name), zap.String("envNamespace", fn.Spec.Environment.Namespace)) logger.Info("choosing pod from pool") diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index 901f4095..ba8acf88 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -58,7 +58,9 @@ import ( otelUtils "github.com/fission/fission/pkg/utils/otel" ) -var _ executortype.ExecutorType = &GenericPoolManager{} +var ( + _ executortype.ExecutorType = &GenericPoolManager{} +) type requestType int @@ -74,7 +76,7 @@ type ( pools map[string]*GenericPool kubernetesClient kubernetes.Interface metricsClient metricsclient.Interface - namespace string + nsResolver *utils.NamespaceResolver fissionClient versioned.Interface functionEnv *cache.Cache @@ -116,7 +118,6 @@ func MakeGenericPoolManager(ctx context.Context, fissionClient versioned.Interface, kubernetesClient kubernetes.Interface, metricsClient metricsclient.Interface, - functionNamespace string, fetcherConfig *fetcherConfig.Config, instanceID string, funcInformer map[string]finformerv1.FunctionInformer, @@ -138,15 +139,15 @@ func MakeGenericPoolManager(ctx context.Context, enableIstio = istio } - poolPodC := NewPoolPodController(ctx, gpmLogger, kubernetesClient, functionNamespace, + poolPodC := NewPoolPodController(ctx, gpmLogger, kubernetesClient, enableIstio, funcInformer, pkgInformer, envInformer, rsInformer, podInformer) gpm := &GenericPoolManager{ logger: gpmLogger, pools: make(map[string]*GenericPool), kubernetesClient: kubernetesClient, + nsResolver: utils.DefaultNSResolver(), metricsClient: metricsClient, - namespace: functionNamespace, fissionClient: fissionClient, functionEnv: cache.MakeCache(10*time.Second, 0), fsCache: fscache.MakeFunctionServiceCache(gpmLogger), @@ -162,6 +163,8 @@ func MakeGenericPoolManager(ctx context.Context, gpm.podLister = podInformer.Lister() gpm.podListerSynced = podInformer.Informer().HasSynced + gpm.logger.Debug("inside MakeGenericPoolManager") + return gpm, nil } @@ -197,7 +200,7 @@ func (gpm *GenericPoolManager) GetFuncSvc(ctx context.Context, fn *fv1.Function) } if created { - logger.Info("created pool for the environment", zap.String("env", env.ObjectMeta.Name), zap.String("namespace", gpm.namespace)) + logger.Info("created pool for the environment", zap.String("env", env.ObjectMeta.Name), zap.String("namespace", gpm.nsResolver.ResolveNamespace(gpm.nsResolver.FunctionNamespace))) } // from GenericPool -> get one function container @@ -277,7 +280,7 @@ func (gpm *GenericPoolManager) RefreshFuncPods(ctx context.Context, logger *zap. } if created { - gpm.logger.Info("created pool for the environment", zap.String("env", env.ObjectMeta.Name), zap.String("namespace", gpm.namespace)) + gpm.logger.Info("created pool for the environment", zap.String("env", env.ObjectMeta.Name), zap.String("namespace", gpm.nsResolver.ResolveNamespace(gpm.nsResolver.FunctionNamespace))) } funcSvc, err := gp.fsCache.GetByFunction(&f.ObjectMeta) @@ -333,7 +336,7 @@ func (gpm *GenericPoolManager) AdoptExistingResources(ctx context.Context) { gpm.logger.Error("adopt pool failed", zap.Error(err)) } if created { - gpm.logger.Info("created pool for the environment", zap.String("env", env.ObjectMeta.Name), zap.String("namespace", gpm.namespace)) + gpm.logger.Info("created pool for the environment", zap.String("env", env.ObjectMeta.Name), zap.String("namespace", gpm.nsResolver.ResolveNamespace(gpm.nsResolver.FunctionNamespace))) } }() } @@ -482,12 +485,9 @@ func (gpm *GenericPoolManager) service() { if !ok { // To support backward compatibility, if envs are created in default ns, we go ahead // and create pools in fission-function ns as earlier. - ns := gpm.namespace - if req.env.ObjectMeta.Namespace != metav1.NamespaceDefault { - ns = req.env.ObjectMeta.Namespace - } + ns := gpm.nsResolver.GetFunctionNS(req.env.ObjectMeta.Namespace) pool = MakeGenericPool(gpm.logger, gpm.fissionClient, gpm.kubernetesClient, - gpm.metricsClient, req.env, ns, gpm.namespace, gpm.fsCache, + gpm.metricsClient, req.env, ns, gpm.fsCache, gpm.fetcherConfig, gpm.instanceID, gpm.enableIstio, gpm.podSpecPatch) err = pool.setup(req.ctx) if err != nil { diff --git a/pkg/executor/executortype/poolmgr/poolpodcontroller.go b/pkg/executor/executortype/poolmgr/poolpodcontroller.go index d0a28fa8..f0fa4dbe 100644 --- a/pkg/executor/executortype/poolmgr/poolpodcontroller.go +++ b/pkg/executor/executortype/poolmgr/poolpodcontroller.go @@ -40,14 +40,15 @@ import ( "github.com/fission/fission/pkg/executor/fscache" finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1" flisterv1 "github.com/fission/fission/pkg/generated/listers/core/v1" + "github.com/fission/fission/pkg/utils" ) type ( PoolPodController struct { logger *zap.Logger kubernetesClient kubernetes.Interface - namespace string enableIstio bool + nsResolver *utils.NamespaceResolver envLister map[string]flisterv1.EnvironmentLister envListerSynced map[string]k8sCache.InformerSynced @@ -69,7 +70,6 @@ type ( func NewPoolPodController(ctx context.Context, logger *zap.Logger, kubernetesClient kubernetes.Interface, - namespace string, enableIstio bool, funcInformer map[string]finformerv1.FunctionInformer, pkgInformer map[string]finformerv1.PackageInformer, @@ -79,8 +79,8 @@ func NewPoolPodController(ctx context.Context, logger *zap.Logger, logger = logger.Named("pool_pod_controller") p := &PoolPodController{ logger: logger, + nsResolver: utils.DefaultNSResolver(), kubernetesClient: kubernetesClient, - namespace: namespace, enableIstio: enableIstio, envLister: make(map[string]flisterv1.EnvironmentLister, 0), envListerSynced: make(map[string]k8sCache.InformerSynced, 0), @@ -89,10 +89,10 @@ func NewPoolPodController(ctx context.Context, logger *zap.Logger, spCleanupPodQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "SpecializedPodCleanupQueue"), } for _, informer := range funcInformer { - informer.Informer().AddEventHandler(FunctionEventHandlers(ctx, p.logger, p.kubernetesClient, p.namespace, p.enableIstio)) + informer.Informer().AddEventHandler(FunctionEventHandlers(ctx, p.logger, p.kubernetesClient, p.nsResolver.ResolveNamespace(p.nsResolver.FunctionNamespace), p.enableIstio)) } for _, informer := range pkgInformer { - informer.Informer().AddEventHandler(PackageEventHandlers(ctx, p.logger, p.kubernetesClient, p.namespace)) + informer.Informer().AddEventHandler(PackageEventHandlers(ctx, p.logger, p.kubernetesClient, p.nsResolver.ResolveNamespace(p.nsResolver.FunctionNamespace))) } for ns, informer := range envInformer { informer.Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ @@ -373,7 +373,7 @@ func (p *PoolPodController) envDeleteQueueProcessFunc(ctx context.Context) bool p.logger.Debug("env delete request processing") p.gpm.cleanupPool(ctx, env) specializePodLables := getSpecializedPodLabels(env) - specializedPods, err := p.podLister.Pods(p.gpm.namespace).List(labels.SelectorFromSet(specializePodLables)) + specializedPods, err := p.podLister.Pods(p.nsResolver.ResolveNamespace(p.nsResolver.FunctionNamespace)).List(labels.SelectorFromSet(specializePodLables)) if err != nil { p.logger.Error("failed to list specialized pods", zap.Error(err)) p.envDeleteQueue.Forget(obj) diff --git a/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go b/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go index b7f5870d..4c45a657 100644 --- a/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go +++ b/pkg/executor/executortype/poolmgr/poolpodcontroller_test.go @@ -68,8 +68,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) { gpmPodInformer := gpmInformerFactory.Core().V1().Pods() gpmRsInformer := gpmInformerFactory.Apps().V1().ReplicaSets() - fnNamespace := "fission-function" - ppc := NewPoolPodController(ctx, logger, kubernetesClient, fnNamespace, false, + ppc := NewPoolPodController(ctx, logger, kubernetesClient, false, funcInformer, pkgInformer, envInformer, @@ -85,7 +84,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) { executor, err := MakeGenericPoolManager(ctx, logger, fissionClient, kubernetesClient, metricsClient, - fnNamespace, fetcherConfig, executorInstanceID, + fetcherConfig, executorInstanceID, funcInformer, pkgInformer, envInformer, gpmPodInformer, gpmRsInformer, nil) if err != nil { diff --git a/pkg/executor/reaper/reaper.go b/pkg/executor/reaper/reaper.go index bdc76695..08c52117 100644 --- a/pkg/executor/reaper/reaper.go +++ b/pkg/executor/reaper/reaper.go @@ -184,7 +184,8 @@ func CleanupHpa(ctx context.Context, logger *zap.Logger, client kubernetes.Inter // CleanupRoleBindings periodically lists rolebindings across all namespaces and removes Service Accounts from them or // deletes the rolebindings completely if there are no Service Accounts in a rolebinding object. -func CleanupRoleBindings(ctx context.Context, logger *zap.Logger, client kubernetes.Interface, fissionClient versioned.Interface, functionNs, envBuilderNs string, cleanupRoleBindingInterval time.Duration) { +func CleanupRoleBindings(ctx context.Context, logger *zap.Logger, client kubernetes.Interface, fissionClient versioned.Interface, cleanupRoleBindingInterval time.Duration) { + nsResolver := utils.DefaultNSResolver() for { // some sleep before the next reaper iteration time.Sleep(cleanupRoleBindingInterval) @@ -238,8 +239,8 @@ func CleanupRoleBindings(ctx context.Context, logger *zap.Logger, client kuberne // so now we need to look for the objects in default namespace. saNs := subj.Namespace isInReservedNS := false - if subj.Namespace == functionNs || - subj.Namespace == envBuilderNs { + if subj.Namespace == nsResolver.FunctionNamespace || + subj.Namespace == nsResolver.BuiderNamespace { saNs = metav1.NamespaceDefault isInReservedNS = true } @@ -249,7 +250,8 @@ func CleanupRoleBindings(ctx context.Context, logger *zap.Logger, client kuberne for _, fn := range funcList.Items { if fn.Spec.Environment.Namespace == saNs || // For the case that the environment is created in the reserved namespace. - (isInReservedNS && (fn.Spec.Environment.Namespace == functionNs || fn.Spec.Environment.Namespace == envBuilderNs)) { + (isInReservedNS && (fn.Spec.Environment.Namespace == nsResolver.FunctionNamespace || + fn.Spec.Environment.Namespace == nsResolver.BuiderNamespace)) { funcEnvReference = true break } diff --git a/pkg/utils/namespace.go b/pkg/utils/namespace.go new file mode 100644 index 00000000..9c27a9f5 --- /dev/null +++ b/pkg/utils/namespace.go @@ -0,0 +1,63 @@ +package utils + +import ( + "os" + + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +const ( + ENV_FUNCTION_NAMESPACE string = "FISSION_FUNCTION_NAMESPACE" + ENV_BUILDER_NAMESPACE string = "FISSION_BUILDER_NAMESPACE" + ENV_DEFAULT_NAMESPACE string = "FISSION_DEFAULT_NAMESPACE" +) + +type NamespaceResolver struct { + FunctionNamespace string + BuiderNamespace string + DefaultNamespace string +} + +var nsResolver *NamespaceResolver + +func init() { + nsResolver = &NamespaceResolver{ + FunctionNamespace: os.Getenv(ENV_FUNCTION_NAMESPACE), + BuiderNamespace: os.Getenv(ENV_BUILDER_NAMESPACE), + DefaultNamespace: os.Getenv(ENV_DEFAULT_NAMESPACE), + } +} + +func (nsr *NamespaceResolver) GetBuilderNS(namespace string) string { + if nsr.FunctionNamespace == "" || nsr.BuiderNamespace == "" { + return namespace + } + + if namespace != metav1.NamespaceDefault { + return namespace + } + return nsr.BuiderNamespace +} + +func (nsr *NamespaceResolver) GetFunctionNS(namespace string) string { + if nsr.FunctionNamespace == "" || nsr.BuiderNamespace == "" { + return namespace + } + + if namespace != metav1.NamespaceDefault { + return namespace + } + return nsr.FunctionNamespace +} + +func (nsr *NamespaceResolver) ResolveNamespace(namespace string) string { + if nsr.FunctionNamespace == "" || nsr.BuiderNamespace == "" { + return nsr.DefaultNamespace + } + return namespace +} + +// GetFissionNamespaces => return all fission core component namespaces +func DefaultNSResolver() *NamespaceResolver { + return nsResolver +} diff --git a/pkg/utils/namespace_test.go b/pkg/utils/namespace_test.go new file mode 100644 index 00000000..ad150dd5 --- /dev/null +++ b/pkg/utils/namespace_test.go @@ -0,0 +1,124 @@ +package utils + +import "testing" + +func TestNamespaceResolver(t *testing.T) { + t.Run("GetBuilderNS", func(t *testing.T) { + for _, test := range []struct { + name string + namespaceResolver *NamespaceResolver + namespace string + expected string + }{ + { + name: "should return fission-builder namespace", + namespaceResolver: getFissionNamespaces("fission-builder", "fission-function", "default"), + namespace: "default", + expected: "fission-builder", + }, + { + name: "should return testns2 namespace", + namespaceResolver: getFissionNamespaces("fission-builder", "fission-function", "testns"), + namespace: "testns2", + expected: "testns2", + }, + { + name: "should return testns3 namespace", + namespaceResolver: getFissionNamespaces("", "", "testns"), + namespace: "testns3", + expected: "testns3", + }, + { + name: "should return default namespace", + namespaceResolver: getFissionNamespaces("fission-builder", "", "default"), + namespace: "default", + expected: "default", + }, + } { + t.Run(test.name, func(t *testing.T) { + ns := test.namespaceResolver.GetBuilderNS(test.namespace) + if ns != test.expected { + t.Errorf("expected builder namespace %s, got %s", test.expected, ns) + } + }) + } + }) + + t.Run("GetFunctionNS", func(t *testing.T) { + for _, test := range []struct { + name string + namespaceResolver *NamespaceResolver + namespace string + expected string + }{ + { + name: "should return fission-function namespace", + namespaceResolver: getFissionNamespaces("fission-builder", "fission-function", "default"), + namespace: "default", + expected: "fission-function", + }, + { + name: "should return testns2 namespace", + namespaceResolver: getFissionNamespaces("fission-builder", "fission-function", "testns"), + namespace: "testns2", + expected: "testns2", + }, + { + name: "should return testns3 namespace", + namespaceResolver: getFissionNamespaces("", "", "testns"), + namespace: "testns3", + expected: "testns3", + }, + { + name: "should return default namespace", + namespaceResolver: getFissionNamespaces("", "fission-function", "default"), + namespace: "default", + expected: "default", + }, + } { + t.Run(test.name, func(t *testing.T) { + ns := test.namespaceResolver.GetFunctionNS(test.namespace) + if ns != test.expected { + t.Errorf("expected function namespace %s, got %s", test.expected, ns) + } + }) + } + }) + + t.Run("ResolveNamespace", func(t *testing.T) { + for _, test := range []struct { + name string + namespaceResolver *NamespaceResolver + namespace string + expected string + }{ + { + name: "should return testns namespace", + namespaceResolver: getFissionNamespaces("fission-builder", "fission-function", "default"), + namespace: "testns", + expected: "testns", + }, + { + name: "should return default namespace", + namespaceResolver: getFissionNamespaces("fission-builder", "", "default"), + namespace: "testns", + expected: "default", + }, + } { + t.Run(test.name, func(t *testing.T) { + ns := test.namespaceResolver.ResolveNamespace(test.namespace) + if ns != test.expected { + t.Errorf("expected function namespace %s, got %s", test.expected, ns) + } + }) + } + }) +} + +func getFissionNamespaces(builderNS, functionNS, defaultNS string) *NamespaceResolver { + return &NamespaceResolver{ + FunctionNamespace: functionNS, + BuiderNamespace: builderNS, + DefaultNamespace: defaultNS, + } +} diff --git a/test/kind_CI.sh b/test/kind_CI.sh index 418714d7..e857848e 100755 --- a/test/kind_CI.sh +++ b/test/kind_CI.sh @@ -14,8 +14,8 @@ echo "source test_utils done" dump_system_info -export FUNCTION_NAMESPACE=fission-function -export BUILDER_NAMESPACE=fission-builder +export FUNCTION_NAMESPACE=default +export BUILDER_NAMESPACE=default export FISSION_NAMESPACE=fission export FISSION_ROUTER=127.0.0.1:8888 export NODE_RUNTIME_IMAGE=fission/node-env-14