From a24934f1a0bbf44dadb5e985196c835427304bd4 Mon Sep 17 00:00:00 2001 From: Sanket Sudake Date: Wed, 18 Aug 2021 13:21:54 +0530 Subject: [PATCH] Organize pool manager code and few improvements (#2166) The patch adds few improvements in pool manager and adds better function composability by reorganizing code. 1. Added created status in get pool call 2. Improved logging in certain areas and having logger per component 3. Separated deployment-specific code in diff file for extensibility Signed-off-by: Sanket Sudake --- pkg/executor/executortype/poolmgr/common.go | 29 +++ .../executortype/poolmgr/funchandlers.go | 22 +- pkg/executor/executortype/poolmgr/gp.go | 200 ++--------------- .../executortype/poolmgr/gp_deployment.go | 203 ++++++++++++++++++ pkg/executor/executortype/poolmgr/gpm.go | 87 ++++---- .../executortype/poolmgr/packagehandlers.go | 19 +- 6 files changed, 323 insertions(+), 237 deletions(-) create mode 100644 pkg/executor/executortype/poolmgr/common.go create mode 100644 pkg/executor/executortype/poolmgr/gp_deployment.go diff --git a/pkg/executor/executortype/poolmgr/common.go b/pkg/executor/executortype/poolmgr/common.go new file mode 100644 index 00000000..5cf07bea --- /dev/null +++ b/pkg/executor/executortype/poolmgr/common.go @@ -0,0 +1,29 @@ +/* +Copyright 2018 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package poolmgr + +import fv1 "github.com/fission/fission/pkg/apis/core/v1" + +func getEnvPoolSize(env *fv1.Environment) int32 { + var poolsize int32 + if env.Spec.Version < 3 { + poolsize = 3 + } else { + poolsize = int32(env.Spec.Poolsize) + } + return poolsize +} diff --git a/pkg/executor/executortype/poolmgr/funchandlers.go b/pkg/executor/executortype/poolmgr/funchandlers.go index e16dea16..12e6f5e2 100644 --- a/pkg/executor/executortype/poolmgr/funchandlers.go +++ b/pkg/executor/executortype/poolmgr/funchandlers.go @@ -37,7 +37,11 @@ func getIstioServiceLabels(fnName string) map[string]string { } } -func (gpm *GenericPoolManager) FunctionEventHandlers(kubernetesClient *kubernetes.Clientset, fissionfnNamespace string, istioEnabled bool) k8sCache.ResourceEventHandlerFuncs { +// FunctionEventHandlers provides handlers for function resource events. +// Based on function create/update/delete event, we create role binding +// for the secret/configmap access which is used by fetcher component. +// If istio is enabled, we create a service for the function. +func FunctionEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clientset, fissionfnNamespace string, istioEnabled bool) k8sCache.ResourceEventHandlerFuncs { return k8sCache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { fn := obj.(*fv1.Function) @@ -65,11 +69,11 @@ func (gpm *GenericPoolManager) FunctionEventHandlers(kubernetesClient *kubernete // setup rolebinding is tried, if it fails, we don't return. we just log an error and move on, because : // 1. not all functions have secrets and/or configmaps, so things will work without this rolebinding in that case. // 2. on the contrary, when the route is tried, the env fetcher logs will show a 403 forbidden message and same will be relayed to executor. - err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.SecretConfigMapGetterRB, fn.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs) + err := utils.SetupRoleBinding(logger, kubernetesClient, fv1.SecretConfigMapGetterRB, fn.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs) if err != nil { - gpm.logger.Error("error creating rolebinding", zap.Error(err), zap.String("role_binding", fv1.SecretConfigMapGetterRB)) + logger.Error("error creating rolebinding", zap.Error(err), zap.String("role_binding", fv1.SecretConfigMapGetterRB)) } else { - gpm.logger.Debug("successfully set up rolebinding for fetcher service account for function", + logger.Debug("successfully set up rolebinding for fetcher service account for function", zap.String("service_account", fv1.FissionFetcherSA), zap.String("service_account_namepsace", envNs), zap.String("function_name", fn.ObjectMeta.Name), @@ -118,7 +122,7 @@ func (gpm *GenericPoolManager) FunctionEventHandlers(kubernetesClient *kubernete // create function istio service if it does not exist _, err = kubernetesClient.CoreV1().Services(envNs).Create(context.TODO(), &svc, metav1.CreateOptions{}) if err != nil && !kerrors.IsAlreadyExists(err) { - gpm.logger.Error("error creating istio service for function", + logger.Error("error creating istio service for function", zap.Error(err), zap.String("service_name", svcName), zap.String("function_name", fn.ObjectMeta.Name), @@ -145,7 +149,7 @@ func (gpm *GenericPoolManager) FunctionEventHandlers(kubernetesClient *kubernete // delete function istio service err := kubernetesClient.CoreV1().Services(envNs).Delete(context.TODO(), svcName, metav1.DeleteOptions{}) if err != nil && !kerrors.IsNotFound(err) { - gpm.logger.Error("error deleting istio service for function", + logger.Error("error deleting istio service for function", zap.Error(err), zap.String("service_name", svcName), zap.String("function_name", fn.ObjectMeta.Name)) @@ -178,14 +182,14 @@ func (gpm *GenericPoolManager) FunctionEventHandlers(kubernetesClient *kubernete envNs = newFunc.Spec.Environment.Namespace } - err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.SecretConfigMapGetterRB, + err := utils.SetupRoleBinding(logger, kubernetesClient, fv1.SecretConfigMapGetterRB, newFunc.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs) if err != nil { - gpm.logger.Error("error creating rolebinding", zap.Error(err), zap.String("role_binding", fv1.SecretConfigMapGetterRB)) + logger.Error("error creating rolebinding", zap.Error(err), zap.String("role_binding", fv1.SecretConfigMapGetterRB)) } else { - gpm.logger.Debug("successfully set up rolebinding for fetcher service account for function", + logger.Debug("successfully set up rolebinding for fetcher service account for function", zap.String("service_account", fv1.FissionFetcherSA), zap.String("service_account_namepsace", envNs), zap.String("function_name", newFunc.ObjectMeta.Name), diff --git a/pkg/executor/executortype/poolmgr/gp.go b/pkg/executor/executortype/poolmgr/gp.go index 37315fa4..aa0edf34 100644 --- a/pkg/executor/executortype/poolmgr/gp.go +++ b/pkg/executor/executortype/poolmgr/gp.go @@ -32,7 +32,6 @@ import ( "go.uber.org/zap" appsv1 "k8s.io/api/apps/v1" apiv1 "k8s.io/api/core/v1" - k8sErrs "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" @@ -46,7 +45,6 @@ import ( fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/executor/fscache" - "github.com/fission/fission/pkg/executor/util" fetcherClient "github.com/fission/fission/pkg/fetcher/client" fetcherConfig "github.com/fission/fission/pkg/fetcher/config" "github.com/fission/fission/pkg/utils" @@ -58,7 +56,6 @@ type ( GenericPool struct { logger *zap.Logger env *fv1.Environment - replicas int32 // num idle pods deployment *appsv1.Deployment // kubernetes deployment namespace string // namespace to keep our resources functionNamespace string // fallback namespace for fission functions @@ -88,7 +85,6 @@ func MakeGenericPool( kubernetesClient *kubernetes.Clientset, metricsClient *metricsclient.Clientset, env *fv1.Environment, - initialReplicas int32, namespace string, functionNamespace string, fsCache *fscache.FunctionServiceCache, @@ -115,7 +111,6 @@ func MakeGenericPool( gp := &GenericPool{ logger: gpLogger, env: env, - replicas: initialReplicas, // TODO make this an env param instead? fissionClient: fissionClient, kubernetesClient: kubernetesClient, metricsClient: metricsClient, @@ -144,7 +139,7 @@ func MakeGenericPool( //gp.labelsForPool = gp.getDeployLabels() // create the pool - err = gp.createPool() + err = gp.createPoolDeployment(context.Background(), env) if err != nil { return nil, err } @@ -155,18 +150,18 @@ func MakeGenericPool( return gp, nil } -func (gp *GenericPool) getEnvironmentPoolLabels() map[string]string { - envLabels := maps.CopyStringMap(gp.env.ObjectMeta.Labels) +func (gp *GenericPool) getEnvironmentPoolLabels(env *fv1.Environment) map[string]string { + envLabels := maps.CopyStringMap(env.ObjectMeta.Labels) envLabels[fv1.EXECUTOR_TYPE] = string(fv1.ExecutorTypePoolmgr) - envLabels[fv1.ENVIRONMENT_NAME] = gp.env.ObjectMeta.Name - envLabels[fv1.ENVIRONMENT_NAMESPACE] = gp.env.ObjectMeta.Namespace - envLabels[fv1.ENVIRONMENT_UID] = string(gp.env.ObjectMeta.UID) + envLabels[fv1.ENVIRONMENT_NAME] = env.ObjectMeta.Name + envLabels[fv1.ENVIRONMENT_NAMESPACE] = env.ObjectMeta.Namespace + envLabels[fv1.ENVIRONMENT_UID] = string(env.ObjectMeta.UID) envLabels["managed"] = "true" // this allows us to easily find pods managed by the deployment return envLabels } -func (gp *GenericPool) getDeployAnnotations() map[string]string { - deployAnnotations := maps.CopyStringMap(gp.env.Annotations) +func (gp *GenericPool) getDeployAnnotations(env *fv1.Environment) map[string]string { + deployAnnotations := maps.CopyStringMap(env.Annotations) deployAnnotations[fv1.EXECUTOR_INSTANCEID_LABEL] = gp.instanceID return deployAnnotations } @@ -277,7 +272,7 @@ func (gp *GenericPool) choosePod(newLabels map[string]string) (string, *apiv1.Po // Append executor instance id to pod annotations to // indicate this pod is managed by this executor. - annotations := gp.getDeployAnnotations() + annotations := gp.getDeployAnnotations(gp.env) annotationPatch, _ := json.Marshal(annotations) patch := fmt.Sprintf(`{"metadata":{"annotations":%v, "labels":%v}}`, string(annotationPatch), string(labelPatch)) @@ -316,7 +311,7 @@ func (gp *GenericPool) choosePod(newLabels map[string]string) (string, *apiv1.Po } func (gp *GenericPool) labelsForFunction(metadata *metav1.ObjectMeta) map[string]string { - label := gp.getEnvironmentPoolLabels() + label := gp.getEnvironmentPoolLabels(gp.env) label[fv1.FUNCTION_NAME] = metadata.Name label[fv1.FUNCTION_UID] = string(metadata.UID) label[fv1.FUNCTION_NAMESPACE] = metadata.Namespace // function CRD must stay within same namespace of environment CRD @@ -402,156 +397,6 @@ func (gp *GenericPool) specializePod(ctx context.Context, pod *apiv1.Pod, fn *fv return nil } -// getPoolName returns a unique name of an environment -func (gp *GenericPool) getPoolName() string { - return strings.ToLower(fmt.Sprintf("poolmgr-%v-%v-%v", gp.env.ObjectMeta.Name, gp.env.ObjectMeta.Namespace, gp.env.ObjectMeta.ResourceVersion)) -} - -// A pool is a deployment of generic containers for an env. This -// creates the pool but doesn't wait for any pods to be ready. -func (gp *GenericPool) createPool() error { - deployLabels := gp.getEnvironmentPoolLabels() - deployAnnotations := gp.getDeployAnnotations() - - // Use long terminationGracePeriodSeconds for connection draining in case that - // pod still runs user functions. - gracePeriodSeconds := int64(6 * 60) - if gp.env.Spec.TerminationGracePeriod > 0 { - gracePeriodSeconds = gp.env.Spec.TerminationGracePeriod - } - - podAnnotations := gp.env.ObjectMeta.Annotations - if podAnnotations == nil { - podAnnotations = make(map[string]string) - } - - // Here, we don't append executor instance-id to pod annotations - // to prevent unwanted rolling updates occur. Pool manager will - // append executor instance-id to pod annotations when a pod is chosen - // for function specialization. - - if gp.useIstio && gp.env.Spec.AllowAccessToExternalNetwork { - podAnnotations["sidecar.istio.io/inject"] = "false" - } - - podLabels := gp.env.ObjectMeta.Labels - if podLabels == nil { - podLabels = make(map[string]string) - } - - for k, v := range deployLabels { - podLabels[k] = v - } - - container, err := util.MergeContainer(&apiv1.Container{ - Name: gp.env.ObjectMeta.Name, - Image: gp.env.Spec.Runtime.Image, - ImagePullPolicy: gp.runtimeImagePullPolicy, - TerminationMessagePath: "/dev/termination-log", - Resources: gp.env.Spec.Resources, - // Pod is removed from endpoints list for service when it's - // state became "Termination". We used preStop hook as the - // workaround for connection draining since pod maybe shutdown - // before grace period expires. - // https://kubernetes.io/docs/concepts/workloads/pods/pod/#termination-of-pods - // https://github.com/kubernetes/kubernetes/issues/47576#issuecomment-308900172 - Lifecycle: &apiv1.Lifecycle{ - PreStop: &apiv1.Handler{ - Exec: &apiv1.ExecAction{ - Command: []string{ - "/bin/sleep", - fmt.Sprintf("%v", gracePeriodSeconds), - }, - }, - }, - }, - // https://istio.io/docs/setup/kubernetes/additional-setup/requirements/ - Ports: []apiv1.ContainerPort{ - { - Name: "http-fetcher", - ContainerPort: int32(8000), - }, - { - Name: "http-env", - ContainerPort: int32(8888), - }, - }, - }, gp.env.Spec.Runtime.Container) - if err != nil { - return err - } - - pod := apiv1.PodTemplateSpec{ - ObjectMeta: metav1.ObjectMeta{ - Labels: podLabels, - Annotations: podAnnotations, - }, - Spec: apiv1.PodSpec{ - Containers: []apiv1.Container{*container}, - ServiceAccountName: "fission-fetcher", - // TerminationGracePeriodSeconds should be equal to the - // sleep time of preStop to make sure that SIGTERM is sent - // to pod after 6 mins. - TerminationGracePeriodSeconds: &gracePeriodSeconds, - }, - } - - pod.Spec = *(util.ApplyImagePullSecret(gp.env.Spec.ImagePullSecret, pod.Spec)) - - deployment := &appsv1.Deployment{ - ObjectMeta: metav1.ObjectMeta{ - Name: gp.getPoolName(), - Labels: deployLabels, - Annotations: deployAnnotations, - }, - Spec: appsv1.DeploymentSpec{ - Replicas: &gp.replicas, - Selector: &metav1.LabelSelector{ - MatchLabels: deployLabels, - }, - Template: pod, - }, - } - - // Order of merging is important here - first fetcher, then containers and lastly pod spec - err = gp.fetcherConfig.AddFetcherToPodSpec(&deployment.Spec.Template.Spec, gp.env.ObjectMeta.Name) - if err != nil { - return err - } - - if gp.env.Spec.Runtime.PodSpec != nil { - newPodSpec, err := util.MergePodSpec(&deployment.Spec.Template.Spec, gp.env.Spec.Runtime.PodSpec) - if err != nil { - return err - } - deployment.Spec.Template.Spec = *newPodSpec - } - - depl, err := gp.kubernetesClient.AppsV1().Deployments(gp.namespace).Get(context.TODO(), deployment.Name, metav1.GetOptions{}) - if err == nil { - if depl.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] != gp.instanceID { - deployment.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] = gp.instanceID - // Update with the latest deployment spec. Kubernetes will trigger - // rolling update if spec is different from the one in the cluster. - depl, err = gp.kubernetesClient.AppsV1().Deployments(gp.namespace).Update(context.TODO(), deployment, metav1.UpdateOptions{}) - } - gp.deployment = depl - return err - } else if !k8sErrs.IsNotFound(err) { - gp.logger.Error("error getting deployment in kubernetes", zap.Error(err), zap.String("deployment", deployment.Name)) - return err - } - - depl, err = gp.kubernetesClient.AppsV1().Deployments(gp.namespace).Create(context.TODO(), deployment, metav1.CreateOptions{}) - if err != nil { - gp.logger.Error("error creating deployment in kubernetes", zap.Error(err), zap.String("deployment", deployment.Name)) - return err - } - - gp.deployment = depl - return nil -} - func (gp *GenericPool) createSvc(name string, labels map[string]string) (*apiv1.Service, error) { service := apiv1.Service{ ObjectMeta: metav1.ObjectMeta{ @@ -575,7 +420,9 @@ func (gp *GenericPool) createSvc(name string, labels map[string]string) (*apiv1. } func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscache.FuncSvc, error) { - gp.logger.Info("choosing pod from pool", zap.Any("function", fn.ObjectMeta)) + log := gp.logger.With(zap.String("function", fn.ObjectMeta.Name), zap.String("functionNamespace", fn.ObjectMeta.Namespace), + zap.String("env", fn.Spec.Environment.Name), zap.String("envNamespace", fn.Spec.Environment.Namespace)) + log.Info("choosing pod from pool") funcLabels := gp.labelsForFunction(&fn.ObjectMeta) if gp.useIstio { @@ -629,7 +476,7 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac gp.scheduleDeletePod(pod.ObjectMeta.Name) return nil, err } - gp.logger.Info("specialized pod", zap.String("pod", pod.ObjectMeta.Name), zap.Any("function", fn.ObjectMeta)) + log.Info("specialized pod", zap.String("pod", pod.ObjectMeta.Name), zap.String("podNamespace", pod.ObjectMeta.Namespace), zap.String("podIP", pod.Status.PodIP)) var svcHost string if gp.useSvc && !gp.useIstio { @@ -664,19 +511,12 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac p, err := gp.kubernetesClient.CoreV1().Pods(pod.Namespace).Patch(context.TODO(), pod.Name, k8sTypes.StrategicMergePatchType, []byte(patch), metav1.PatchOptions{}) if err != nil { // just log the error since it won't affect the function serving - gp.logger.Warn("error patching svc-host to pod", zap.Error(err), + log.Warn("error patching svc-host to pod", zap.Error(err), zap.String("pod", pod.Name), zap.String("ns", pod.Namespace)) } else { pod = p } - gp.logger.Info("specialized pod", - zap.String("pod", pod.ObjectMeta.Name), - zap.String("podNamespace", pod.ObjectMeta.Namespace), - zap.String("function", fn.ObjectMeta.Name), - zap.String("functionNamespace", fn.ObjectMeta.Namespace), - zap.String("specialization_host", svcHost)) - kubeObjRefs := []apiv1.ObjectReference{ { Kind: "pod", @@ -696,10 +536,10 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac // set cpuLimit to 85th percentage of the cpuUsage cpuLimit, err := gp.getPercent(cpuUsage, 0.85) if err != nil { - gp.logger.Error("failed to get 85 of CPU usage", zap.Error(err)) + log.Error("failed to get 85 of CPU usage", zap.Error(err)) cpuLimit = cpuUsage } - gp.logger.Debug("cpuLimit set to", zap.Any("cpulimit", cpuLimit)) + log.Debug("cpuLimit set to", zap.Any("cpulimit", cpuLimit)) m := fn.ObjectMeta // only cache necessary part fsvc := &fscache.FuncSvc{ @@ -720,6 +560,12 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac gp.fsCache.IncreaseColdStarts(fn.ObjectMeta.Name, string(fn.ObjectMeta.UID)) + log.Info("added function service", + zap.String("pod", pod.ObjectMeta.Name), + zap.String("podNamespace", pod.ObjectMeta.Namespace), + zap.String("serviceHost", svcHost), + zap.String("podIP", pod.Status.PodIP)) + return fsvc, nil } diff --git a/pkg/executor/executortype/poolmgr/gp_deployment.go b/pkg/executor/executortype/poolmgr/gp_deployment.go new file mode 100644 index 00000000..87881294 --- /dev/null +++ b/pkg/executor/executortype/poolmgr/gp_deployment.go @@ -0,0 +1,203 @@ +/* +Copyright 2016 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package poolmgr + +import ( + "context" + "fmt" + "strings" + + "go.uber.org/zap" + appsv1 "k8s.io/api/apps/v1" + apiv1 "k8s.io/api/core/v1" + k8sErrs "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + + fv1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/executor/util" +) + +// getPoolName returns a unique name of an environment +func (gp *GenericPool) getPoolName(env *fv1.Environment) string { + // TODO: get rid of resource version here + return strings.ToLower(fmt.Sprintf("poolmgr-%v-%v-%v", env.ObjectMeta.Name, env.ObjectMeta.Namespace, env.ObjectMeta.ResourceVersion)) +} + +func (gp *GenericPool) genDeploymentMeta(env *fv1.Environment) metav1.ObjectMeta { + deployLabels := gp.getEnvironmentPoolLabels(env) + deployAnnotations := gp.getDeployAnnotations(env) + return metav1.ObjectMeta{ + Name: gp.getPoolName(env), + Labels: deployLabels, + Annotations: deployAnnotations, + } +} + +func (gp *GenericPool) genDeploymentSpec(env *fv1.Environment) (*appsv1.DeploymentSpec, error) { + deployLabels := gp.getEnvironmentPoolLabels(env) + // Use long terminationGracePeriodSeconds for connection draining in case that + // pod still runs user functions. + gracePeriodSeconds := int64(6 * 60) + if env.Spec.TerminationGracePeriod > 0 { + gracePeriodSeconds = env.Spec.TerminationGracePeriod + } + + podAnnotations := env.ObjectMeta.Annotations + if podAnnotations == nil { + podAnnotations = make(map[string]string) + } + + // Here, we don't append executor instance-id to pod annotations + // to prevent unwanted rolling updates occur. Pool manager will + // append executor instance-id to pod annotations when a pod is chosen + // for function specialization. + + if gp.useIstio && env.Spec.AllowAccessToExternalNetwork { + podAnnotations["sidecar.istio.io/inject"] = "false" + } + + podLabels := env.ObjectMeta.Labels + if podLabels == nil { + podLabels = make(map[string]string) + } + + for k, v := range deployLabels { + podLabels[k] = v + } + + container, err := util.MergeContainer(&apiv1.Container{ + Name: env.ObjectMeta.Name, + Image: env.Spec.Runtime.Image, + ImagePullPolicy: gp.runtimeImagePullPolicy, + TerminationMessagePath: "/dev/termination-log", + Resources: env.Spec.Resources, + // Pod is removed from endpoints list for service when it's + // state became "Termination". We used preStop hook as the + // workaround for connection draining since pod maybe shutdown + // before grace period expires. + // https://kubernetes.io/docs/concepts/workloads/pods/pod/#termination-of-pods + // https://github.com/kubernetes/kubernetes/issues/47576#issuecomment-308900172 + Lifecycle: &apiv1.Lifecycle{ + PreStop: &apiv1.Handler{ + Exec: &apiv1.ExecAction{ + Command: []string{ + "/bin/sleep", + fmt.Sprintf("%v", gracePeriodSeconds), + }, + }, + }, + }, + // https://istio.io/docs/setup/kubernetes/additional-setup/requirements/ + Ports: []apiv1.ContainerPort{ + { + Name: "http-fetcher", + ContainerPort: int32(8000), + }, + { + Name: "http-env", + ContainerPort: int32(8888), + }, + }, + }, env.Spec.Runtime.Container) + if err != nil { + return nil, err + } + + pod := apiv1.PodTemplateSpec{ + ObjectMeta: metav1.ObjectMeta{ + Labels: podLabels, + Annotations: podAnnotations, + }, + Spec: apiv1.PodSpec{ + Containers: []apiv1.Container{*container}, + ServiceAccountName: "fission-fetcher", + // TerminationGracePeriodSeconds should be equal to the + // sleep time of preStop to make sure that SIGTERM is sent + // to pod after 6 mins. + TerminationGracePeriodSeconds: &gracePeriodSeconds, + }, + } + + pod.Spec = *(util.ApplyImagePullSecret(env.Spec.ImagePullSecret, pod.Spec)) + + poolsize := getEnvPoolSize(env) + switch env.Spec.AllowedFunctionsPerContainer { + case fv1.AllowedFunctionsPerContainerInfinite: + poolsize = 1 + } + + deploymentSpec := appsv1.DeploymentSpec{ + // TODO: fix this hardcoded value + Replicas: &poolsize, + Selector: &metav1.LabelSelector{ + MatchLabels: deployLabels, + }, + Template: pod, + } + + // Order of merging is important here - first fetcher, then containers and lastly pod spec + err = gp.fetcherConfig.AddFetcherToPodSpec(&deploymentSpec.Template.Spec, env.ObjectMeta.Name) + if err != nil { + return nil, err + } + + if env.Spec.Runtime.PodSpec != nil { + newPodSpec, err := util.MergePodSpec(&deploymentSpec.Template.Spec, env.Spec.Runtime.PodSpec) + if err != nil { + return nil, err + } + deploymentSpec.Template.Spec = *newPodSpec + } + return &deploymentSpec, nil +} + +// A pool is a deployment of generic containers for an env. This +// creates the pool but doesn't wait for any pods to be ready. +func (gp *GenericPool) createPoolDeployment(ctx context.Context, env *fv1.Environment) error { + deploymentMeta := gp.genDeploymentMeta(env) + deploymentSpec, err := gp.genDeploymentSpec(env) + if err != nil { + return err + } + deployment := &appsv1.Deployment{ + ObjectMeta: deploymentMeta, + Spec: *deploymentSpec, + } + depl, err := gp.kubernetesClient.AppsV1().Deployments(gp.namespace).Get(ctx, deployment.Name, metav1.GetOptions{}) + if err == nil { + if depl.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] != gp.instanceID { + deployment.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] = gp.instanceID + // Update with the latest deployment spec. Kubernetes will trigger + // rolling update if spec is different from the one in the cluster. + depl, err = gp.kubernetesClient.AppsV1().Deployments(gp.namespace).Update(ctx, deployment, metav1.UpdateOptions{}) + } + gp.deployment = depl + return err + } else if !k8sErrs.IsNotFound(err) { + gp.logger.Error("error getting deployment in kubernetes", zap.Error(err), zap.String("deployment", deployment.Name)) + return err + } + + depl, err = gp.kubernetesClient.AppsV1().Deployments(gp.namespace).Create(ctx, deployment, metav1.CreateOptions{}) + if err != nil { + gp.logger.Error("error creating deployment in kubernetes", zap.Error(err), zap.String("deployment", deployment.Name)) + return err + } + + gp.deployment = depl + return nil +} diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index 5c6eb696..1d7aabe4 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -91,7 +91,8 @@ type ( } response struct { error - pool *GenericPool + pool *GenericPool + created bool } ) @@ -109,6 +110,15 @@ func MakeGenericPoolManager( gpmLogger := logger.Named("generic_pool_manager") + enableIstio := false + if len(os.Getenv("ENABLE_ISTIO")) > 0 { + istio, err := strconv.ParseBool(os.Getenv("ENABLE_ISTIO")) + if err != nil { + gpmLogger.Error("failed to parse 'ENABLE_ISTIO', set to false", zap.Error(err)) + } + enableIstio = istio + } + gpm := &GenericPoolManager{ logger: gpmLogger, pools: make(map[string]*GenericPool), @@ -122,22 +132,15 @@ func MakeGenericPoolManager( requestChannel: make(chan *request), defaultIdlePodReapTime: 2 * time.Minute, fetcherConfig: fetcherConfig, + enableIstio: enableIstio, funcInformer: funcInformer, pkgInformer: pkgInformer, } go gpm.service() - if len(os.Getenv("ENABLE_ISTIO")) > 0 { - istio, err := strconv.ParseBool(os.Getenv("ENABLE_ISTIO")) - if err != nil { - gpmLogger.Error("failed to parse 'ENABLE_ISTIO', set to false", zap.Error(err)) - } - gpm.enableIstio = istio - } - - (*gpm.funcInformer).AddEventHandler(gpm.FunctionEventHandlers(gpm.kubernetesClient, gpm.namespace, gpm.enableIstio)) - (*gpm.pkgInformer).AddEventHandler(gpm.PackageEventHandlers(gpm.kubernetesClient, gpm.namespace)) + (*gpm.funcInformer).AddEventHandler(FunctionEventHandlers(gpm.logger, gpm.kubernetesClient, gpm.namespace, gpm.enableIstio)) + (*gpm.pkgInformer).AddEventHandler(PackageEventHandlers(gpm.logger, gpm.kubernetesClient, gpm.namespace)) kubeInformerFactory, err := utils.GetInformerFactoryByExecutor(gpm.kubernetesClient, fv1.ExecutorTypePoolmgr) if err != nil { @@ -152,6 +155,8 @@ func (gpm *GenericPoolManager) Run(ctx context.Context) { // Otherwise, the poolmanager may wrongly delete the deployment. go gpm.eagerPoolCreator() go gpm.podInformer.Run(ctx.Done()) + go gpm.WebsocketStartEventChecker(gpm.kubernetesClient) + go gpm.NoActiveConnectionEventChecker(gpm.kubernetesClient) go gpm.idleObjectReaper() } @@ -167,11 +172,15 @@ func (gpm *GenericPoolManager) GetFuncSvc(ctx context.Context, fn *fv1.Function) return nil, err } - pool, err := gpm.getPool(env) + pool, created, err := gpm.getPool(env) if err != nil { return nil, err } + if created { + gpm.logger.Info("created pool for the environment", zap.String("env", env.ObjectMeta.Name), zap.String("namespace", gpm.namespace)) + } + // from GenericPool -> get one function container // (this also adds to the cache) gpm.logger.Debug("getting function service from pool", zap.String("function", fn.ObjectMeta.Name)) @@ -253,11 +262,15 @@ func (gpm *GenericPoolManager) RefreshFuncPods(logger *zap.Logger, f fv1.Functio return err } - gp, err := gpm.getPool(env) + gp, created, err := gpm.getPool(env) if err != nil { return err } + if created { + gpm.logger.Info("created pool for the environment", zap.String("env", env.ObjectMeta.Name), zap.String("namespace", gpm.namespace)) + } + funcSvc, err := gp.fsCache.GetByFunction(&f.ObjectMeta) if err != nil { return err @@ -301,14 +314,17 @@ func (gpm *GenericPoolManager) AdoptExistingResources() { for i := range envs.Items { env := envs.Items[i] - if gpm.getEnvPoolsize(&env) > 0 { + if getEnvPoolSize(&env) > 0 { wg.Add(1) go func() { defer wg.Done() - _, err := gpm.getPool(&env) + _, created, err := gpm.getPool(&env) if err != nil { 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)) + } }() } @@ -448,14 +464,9 @@ func (gpm *GenericPoolManager) service() { case GET_POOL: // just because they are missing in the cache, we end up creating another duplicate pool. var err error + created := false pool, ok := gpm.pools[crd.CacheKey(&req.env.ObjectMeta)] if !ok { - poolsize := gpm.getEnvPoolsize(req.env) - switch req.env.Spec.AllowedFunctionsPerContainer { - case fv1.AllowedFunctionsPerContainerInfinite: - poolsize = 1 - } - // 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 @@ -464,19 +475,20 @@ func (gpm *GenericPoolManager) service() { } pool, err = MakeGenericPool(gpm.logger, - gpm.fissionClient, gpm.kubernetesClient, gpm.metricsClient, req.env, poolsize, - ns, gpm.namespace, gpm.fsCache, gpm.fetcherConfig, gpm.instanceID, gpm.enableIstio) + gpm.fissionClient, gpm.kubernetesClient, gpm.metricsClient, req.env, ns, + gpm.namespace, gpm.fsCache, gpm.fetcherConfig, gpm.instanceID, gpm.enableIstio) if err != nil { req.responseChannel <- &response{error: err} continue } gpm.pools[crd.CacheKey(&req.env.ObjectMeta)] = pool + created = true } - req.responseChannel <- &response{pool: pool} + req.responseChannel <- &response{pool: pool, created: created} case CLEANUP_POOLS: latestEnvPoolsize := make(map[string]int) for _, env := range req.envList { - latestEnvPoolsize[crd.CacheKey(&env.ObjectMeta)] = int(gpm.getEnvPoolsize(&env)) + latestEnvPoolsize[crd.CacheKey(&env.ObjectMeta)] = int(getEnvPoolSize(&env)) } for key, pool := range gpm.pools { poolsize, ok := latestEnvPoolsize[key] @@ -495,7 +507,7 @@ func (gpm *GenericPoolManager) service() { } } -func (gpm *GenericPoolManager) getPool(env *fv1.Environment) (*GenericPool, error) { +func (gpm *GenericPoolManager) getPool(env *fv1.Environment) (*GenericPool, bool, error) { c := make(chan *response) gpm.requestChannel <- &request{ requestType: GET_POOL, @@ -503,7 +515,7 @@ func (gpm *GenericPoolManager) getPool(env *fv1.Environment) (*GenericPool, erro responseChannel: c, } resp := <-c - return resp.pool, resp.error + return resp.pool, resp.created, resp.error } func (gpm *GenericPoolManager) cleanupPools(envs []fv1.Environment) { @@ -568,14 +580,17 @@ func (gpm *GenericPoolManager) eagerPoolCreator() { for i := range envs.Items { env := envs.Items[i] // Create pool only if poolsize greater than zero - if gpm.getEnvPoolsize(&env) > 0 { + if getEnvPoolSize(&env) > 0 { wg.Add(1) go func() { defer wg.Done() - _, err := gpm.getPool(&env) + _, created, err := gpm.getPool(&env) if err != nil { gpm.logger.Error("eager-create 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)) + } }() } } @@ -587,25 +602,11 @@ func (gpm *GenericPoolManager) eagerPoolCreator() { } } -func (gpm *GenericPoolManager) getEnvPoolsize(env *fv1.Environment) int32 { - var poolsize int32 - if env.Spec.Version < 3 { - poolsize = 3 - } else { - poolsize = int32(env.Spec.Poolsize) - } - return poolsize -} - // idleObjectReaper reaps objects after certain idle time func (gpm *GenericPoolManager) idleObjectReaper() { pollSleep := 5 * time.Second - go gpm.WebsocketStartEventChecker(gpm.kubernetesClient) - - go gpm.NoActiveConnectionEventChecker(gpm.kubernetesClient) - for { time.Sleep(pollSleep) diff --git a/pkg/executor/executortype/poolmgr/packagehandlers.go b/pkg/executor/executortype/poolmgr/packagehandlers.go index 0c6c025a..30d7b445 100644 --- a/pkg/executor/executortype/poolmgr/packagehandlers.go +++ b/pkg/executor/executortype/poolmgr/packagehandlers.go @@ -26,11 +26,14 @@ import ( "github.com/fission/fission/pkg/utils" ) -func (gpm *GenericPoolManager) PackageEventHandlers(kubernetesClient *kubernetes.Clientset, fissionfnNamespace string) k8sCache.ResourceEventHandlerFuncs { +// PackageEventHandlers provides handlers for package events. +// Based on package create/update event, we create role binding +// for the package which is used by fetcher component. +func PackageEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clientset, fissionfnNamespace string) k8sCache.ResourceEventHandlerFuncs { return k8sCache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { pkg := obj.(*fv1.Package) - gpm.logger.Debug("list watch for package reported a new package addition", + logger.Debug("list watch for package reported a new package addition", zap.String("package_name", pkg.ObjectMeta.Name), zap.String("package_namespace", pkg.ObjectMeta.Namespace)) @@ -42,9 +45,9 @@ func (gpm *GenericPoolManager) PackageEventHandlers(kubernetesClient *kubernetes // here, we return if we hit an error during rolebinding setup. this is because this rolebinding is mandatory for // every function's package to be loaded into its env. without that, there's no point to move forward. - err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.PackageGetterRB, pkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs) + err := utils.SetupRoleBinding(logger, kubernetesClient, fv1.PackageGetterRB, pkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs) if err != nil { - gpm.logger.Error("error creating rolebinding for package", + logger.Error("error creating rolebinding for package", zap.Error(err), zap.String("role_binding", fv1.PackageGetterRB), zap.String("package_name", pkg.ObjectMeta.Name), @@ -52,7 +55,7 @@ func (gpm *GenericPoolManager) PackageEventHandlers(kubernetesClient *kubernetes return } - gpm.logger.Debug("successfully set up rolebinding for fetcher service account", + logger.Debug("successfully set up rolebinding for fetcher service account", zap.String("service_account", fv1.FissionFetcherSA), zap.String("service_account_namespace", envNs), zap.String("package_name", pkg.ObjectMeta.Name), @@ -76,11 +79,11 @@ func (gpm *GenericPoolManager) PackageEventHandlers(kubernetesClient *kubernetes envNs = newPkg.Spec.Environment.Namespace } - err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.PackageGetterRB, + err := utils.SetupRoleBinding(logger, kubernetesClient, fv1.PackageGetterRB, newPkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs) if err != nil { - gpm.logger.Error("error updating rolebinding for package", + logger.Error("error updating rolebinding for package", zap.Error(err), zap.String("role_binding", fv1.PackageGetterRB), zap.String("package_name", newPkg.ObjectMeta.Name), @@ -88,7 +91,7 @@ func (gpm *GenericPoolManager) PackageEventHandlers(kubernetesClient *kubernetes return } - gpm.logger.Debug("successfully updated rolebinding for fetcher service account", + logger.Debug("successfully updated rolebinding for fetcher service account", zap.String("service_account", fv1.FissionFetcherSA), zap.String("service_account_namespace", envNs), zap.String("package_name", newPkg.ObjectMeta.Name),