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 <sanketsudake@gmail.com>
This commit is contained in:
Sanket Sudake
2021-08-18 13:21:54 +05:30
committed by GitHub
parent 9d54f5daac
commit a24934f1a0
6 changed files with 323 additions and 237 deletions
@@ -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
}
@@ -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),
+23 -177
View File
@@ -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
}
@@ -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
}
+44 -43
View File
@@ -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)
@@ -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),