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),