From 1ad7ac2dcf33e91372d70b55cddce924039aa2db Mon Sep 17 00:00:00 2001 From: Ta-Ching Chen Date: Wed, 27 Nov 2019 23:08:45 +0800 Subject: [PATCH] Adopt existing orphan kubernetes resources when executor starts up (#1443) Previously, once the executor is deleted for reasons (like upgrade or cluster scale-in), the new executor deletes all existing resources created by the old executor and creates new one. This mechanism becomes a problem when there are requests connecting to the existing pods. Also in the worst case, the cluster may not have enough resources to create new pods and cause service downtime. This PR let each executor type adopts existing resources before starting the executor API services, and so the alive connections won't experience failure. However, the requests send to the function that doesn't have alive function pods will still fail due to the executor is in bootstrapping. --- pkg/executor/api.go | 16 +- pkg/executor/cms/cmscontroller.go | 2 +- pkg/executor/executor.go | 23 ++- pkg/executor/executortype/executortype.go | 3 + .../executortype/newdeploy/newdeploy.go | 76 ++++++--- .../executortype/newdeploy/newdeploymgr.go | 113 +++++++++++-- pkg/executor/executortype/poolmgr/gp.go | 95 ++++++++--- pkg/executor/executortype/poolmgr/gpm.go | 157 +++++++++++++++++- pkg/executor/reaper/reaper.go | 44 +++-- pkg/types/types.go | 19 ++- test/utils.sh | 4 +- 11 files changed, 437 insertions(+), 115 deletions(-) diff --git a/pkg/executor/api.go b/pkg/executor/api.go index 7f4ad86d..c1466615 100644 --- a/pkg/executor/api.go +++ b/pkg/executor/api.go @@ -17,7 +17,6 @@ limitations under the License. package executor import ( - "context" "encoding/json" "fmt" "io/ioutil" @@ -164,8 +163,8 @@ func (executor *Executor) tapServices(w http.ResponseWriter, r *http.Request) { err = et.TapService(svcHost) if err != nil { errs = multierror.Append(errs, - errors.Wrapf(err, "'%v' failed to tap function '%v/%v' with service url '%v'", - req.FnMetadata.Namespace, req.FnMetadata.Name, req.ServiceUrl, req.FnExecutorType)) + errors.Wrapf(err, "'%v' failed to tap function '%v' in '%v' with service url '%v'", + req.FnMetadata.Name, req.FnMetadata.Namespace, req.ServiceUrl, req.FnExecutorType)) } } @@ -192,16 +191,7 @@ func (executor *Executor) GetHandler() http.Handler { } func (executor *Executor) Serve(port int) { - executor.logger.Info("starting executor", zap.Int("port", port)) - - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() - - for _, et := range executor.executorTypes { - et.Run(ctx) - } - - executor.cms.Run(ctx) + executor.logger.Info("starting executor API", zap.Int("port", port)) address := fmt.Sprintf(":%v", port) err := http.ListenAndServe(address, &ochttp.Handler{ Handler: executor.GetHandler(), diff --git a/pkg/executor/cms/cmscontroller.go b/pkg/executor/cms/cmscontroller.go index 2df4b2f8..7bfb8f27 100644 --- a/pkg/executor/cms/cmscontroller.go +++ b/pkg/executor/cms/cmscontroller.go @@ -83,7 +83,7 @@ func initConfigmapController(logger *zap.Logger, fissionClient *crd.FissionClien } funcs, err := getConfigmapRelatedFuncs(logger, &newCm.ObjectMeta, fissionClient) if err != nil { - logger.Error("Failed to get functions related to secret", zap.String("secret_name", newCm.ObjectMeta.Name), zap.String("secret_namespace", newCm.ObjectMeta.Namespace)) + logger.Error("Failed to get functions related to configmap", zap.String("configmap_name", newCm.ObjectMeta.Name), zap.String("configmap_namespace", newCm.ObjectMeta.Namespace)) } recyclePods(logger, funcs, types) } diff --git a/pkg/executor/executor.go b/pkg/executor/executor.go index 622fc8a2..0350a02c 100644 --- a/pkg/executor/executor.go +++ b/pkg/executor/executor.go @@ -217,32 +217,45 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames } restClient := fissionClient.GetCrdClient() + executorInstanceID := strings.ToLower(uniuri.NewLen(8)) - poolID := strings.ToLower(uniuri.NewLen(8)) - reaper.CleanupOldExecutorObjects(logger, kubernetesClient, poolID) - go reaper.CleanupRoleBindings(logger, kubernetesClient, fissionClient, functionNamespace, envBuilderNamespace, time.Minute*30) + logger.Info("Starting executor", zap.String("instanceID", executorInstanceID)) gpm := poolmgr.MakeGenericPoolManager( logger, fissionClient, kubernetesClient, - functionNamespace, fetcherConfig, poolID) + functionNamespace, fetcherConfig, executorInstanceID) ndm := newdeploy.MakeNewDeploy( logger, fissionClient, kubernetesClient, restClient, - functionNamespace, fetcherConfig, poolID) + functionNamespace, fetcherConfig, executorInstanceID) executorTypes := make(map[fv1.ExecutorType]executortype.ExecutorType) executorTypes[gpm.GetTypeName()] = gpm executorTypes[ndm.GetTypeName()] = ndm + wg := &sync.WaitGroup{} + for _, et := range executorTypes { + wg.Add(1) + go func(et executortype.ExecutorType) { + defer wg.Done() + et.AdoptOrphanResources() + et.Run(context.Background()) + }(et) + } + wg.Wait() + cms := cms.MakeConfigSecretController(logger, fissionClient, kubernetesClient, executorTypes) + cms.Run(context.Background()) api, err := MakeExecutor(logger, cms, fissionClient, executorTypes) if err != nil { return err } + go reaper.CleanupRoleBindings(logger, kubernetesClient, fissionClient, functionNamespace, envBuilderNamespace, time.Minute*30) + go reaper.CleanupOldExecutorObjects(logger, kubernetesClient, executorInstanceID) go api.Serve(port) go serveMetric(logger) diff --git a/pkg/executor/executortype/executortype.go b/pkg/executor/executortype/executortype.go index d5879071..cb41d411 100644 --- a/pkg/executor/executortype/executortype.go +++ b/pkg/executor/executortype/executortype.go @@ -50,4 +50,7 @@ type ExecutorType interface { // RefreshFuncPods refreshes function pods if the secrets/configmaps pods reference to get updated. RefreshFuncPods(*zap.Logger, fv1.Function) error + + // AdoptOrphanResources adopts existing resources created by the deleted executor. + AdoptOrphanResources() } diff --git a/pkg/executor/executortype/newdeploy/newdeploy.go b/pkg/executor/executortype/newdeploy/newdeploy.go index 8220577f..9fd4fbaa 100644 --- a/pkg/executor/executortype/newdeploy/newdeploy.go +++ b/pkg/executor/executortype/newdeploy/newdeploy.go @@ -29,6 +29,7 @@ import ( k8s_err "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + k8sTypes "k8s.io/apimachinery/pkg/types" "k8s.io/apimachinery/pkg/util/intstr" fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" @@ -43,7 +44,7 @@ const ( ) func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Environment, - deployName string, deployLabels map[string]string, deployNamespace string, firstcreate bool) (*appsv1.Deployment, error) { + deployName string, deployLabels map[string]string, deployAnnotations map[string]string, deployNamespace string, firstcreate bool) (*appsv1.Deployment, error) { minScale := int32(fn.Spec.InvokeStrategy.ExecutionStrategy.MinScale) specializationTimeout := int(fn.Spec.InvokeStrategy.ExecutionStrategy.SpecializationTimeout) @@ -58,13 +59,23 @@ func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Enviro existingDepl, err := deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Get(deployName, metav1.GetOptions{}) if err == nil { - if waitForDeploy { - err = deploy.scaleDeployment(existingDepl.Namespace, existingDepl.Name, minScale) + if existingDepl.Labels[fv1.EXECUTOR_INSTANCEID_LABEL] != deploy.instanceID { + patch := fmt.Sprintf(`{"metadata":{"annotations":{"%v":"%v"}}}`, types.EXECUTOR_INSTANCEID_LABEL, deploy.instanceID) + existingDepl, err = deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Patch(deployName, k8sTypes.StrategicMergePatchType, []byte(patch)) if err != nil { - deploy.logger.Error("error scaling up function deployment", zap.Error(err), zap.String("function", fn.Metadata.Name)) + deploy.logger.Warn("error patching executor instance ID of deploy", zap.Error(err), + zap.String("deploy", deployName), zap.String("ns", deployNamespace)) return nil, err } - + } + if waitForDeploy { + if *existingDepl.Spec.Replicas < minScale { + err = deploy.scaleDeployment(existingDepl.Namespace, existingDepl.Name, minScale) + if err != nil { + deploy.logger.Error("error scaling up function deployment", zap.Error(err), zap.String("function", fn.Metadata.Name)) + return nil, err + } + } if existingDepl.Status.AvailableReplicas < minScale { existingDepl, err = deploy.waitForDeploy(existingDepl, minScale, specializationTimeout) } @@ -76,7 +87,7 @@ func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Enviro return nil, err } - deployment, err := deploy.getDeploymentSpec(fn, env, deployName, deployNamespace, deployLabels) + deployment, err := deploy.getDeploymentSpec(fn, env, deployName, deployNamespace, deployLabels, deployAnnotations) if err != nil { return nil, err } @@ -158,7 +169,7 @@ func (deploy *NewDeploy) deleteDeployment(ns string, name string) error { } func (deploy *NewDeploy) getDeploymentSpec(fn *fv1.Function, env *fv1.Environment, - deployName string, deployNamespace string, deployLabels map[string]string) (*appsv1.Deployment, error) { + deployName string, deployNamespace string, deployLabels map[string]string, deployAnnotations map[string]string) (*appsv1.Deployment, error) { replicas := int32(fn.Spec.InvokeStrategy.ExecutionStrategy.MinScale) @@ -174,6 +185,9 @@ func (deploy *NewDeploy) getDeploymentSpec(fn *fv1.Function, env *fv1.Environmen if deploy.useIstio && env.Spec.AllowAccessToExternalNetwork { podAnnotations["sidecar.istio.io/inject"] = "false" } + for k, v := range deployAnnotations { + podAnnotations[k] = v + } resources := deploy.getResources(env, fn) // Set maxUnavailable and maxSurge to 20% is because we want @@ -243,8 +257,9 @@ func (deploy *NewDeploy) getDeploymentSpec(fn *fv1.Function, env *fv1.Environmen deployment := &appsv1.Deployment{ ObjectMeta: metav1.ObjectMeta{ - Name: deployName, - Labels: deployLabels, + Name: deployName, + Labels: deployLabels, + Annotations: deployAnnotations, }, Spec: appsv1.DeploymentSpec{ Replicas: &replicas, @@ -319,7 +334,10 @@ func (deploy *NewDeploy) getResources(env *fv1.Environment, fn *fv1.Function) ap return resources } -func (deploy *NewDeploy) createOrGetHpa(hpaName string, execStrategy *fv1.ExecutionStrategy, depl *appsv1.Deployment) (*asv1.HorizontalPodAutoscaler, error) { +func (deploy *NewDeploy) createOrGetHpa(hpaName string, execStrategy *fv1.ExecutionStrategy, depl *appsv1.Deployment, deployLabels map[string]string, deployAnnotations map[string]string) (*asv1.HorizontalPodAutoscaler, error) { + if depl == nil { + return nil, errors.New("failed to create HPA, found empty deployment") + } minRepl := int32(execStrategy.MinScale) if minRepl == 0 { @@ -333,18 +351,22 @@ func (deploy *NewDeploy) createOrGetHpa(hpaName string, execStrategy *fv1.Execut existingHpa, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Get(hpaName, metav1.GetOptions{}) if err == nil { + if existingHpa.Labels[fv1.EXECUTOR_INSTANCEID_LABEL] != deploy.instanceID { + patch := fmt.Sprintf(`{"metadata":{"annotations":{"%v":"%v"}}}`, types.EXECUTOR_INSTANCEID_LABEL, deploy.instanceID) + existingHpa, err = deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Patch(hpaName, k8sTypes.StrategicMergePatchType, []byte(patch)) + if err != nil { + deploy.logger.Warn("error patching executor instance ID of HPA", zap.Error(err), + zap.String("HPA", hpaName), zap.String("ns", depl.ObjectMeta.Namespace)) + return nil, err + } + } return existingHpa, err - } - - if depl == nil { - return nil, errors.New("failed to create HPA, found empty deployment") - } - - if err != nil && k8s_err.IsNotFound(err) { + } else if k8s_err.IsNotFound(err) { hpa := asv1.HorizontalPodAutoscaler{ ObjectMeta: metav1.ObjectMeta{ - Name: hpaName, - Labels: depl.Labels, + Name: hpaName, + Labels: deployLabels, + Annotations: deployAnnotations, }, Spec: asv1.HorizontalPodAutoscalerSpec{ ScaleTargetRef: asv1.CrossVersionObjectReference{ @@ -382,15 +404,25 @@ func (deploy *NewDeploy) deleteHpa(ns string, name string) error { return deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Delete(name, &metav1.DeleteOptions{}) } -func (deploy *NewDeploy) createOrGetSvc(deployLabels map[string]string, svcName string, svcNamespace string) (*apiv1.Service, error) { +func (deploy *NewDeploy) createOrGetSvc(deployLabels map[string]string, deployAnnotations map[string]string, svcName string, svcNamespace string) (*apiv1.Service, error) { existingSvc, err := deploy.kubernetesClient.CoreV1().Services(svcNamespace).Get(svcName, metav1.GetOptions{}) if err == nil { + if existingSvc.Labels[fv1.EXECUTOR_INSTANCEID_LABEL] != deploy.instanceID { + patch := fmt.Sprintf(`{"metadata":{"annotations":{"%v":"%v"}}}`, types.EXECUTOR_INSTANCEID_LABEL, deploy.instanceID) + existingSvc, err = deploy.kubernetesClient.CoreV1().Services(svcNamespace).Patch(svcName, k8sTypes.StrategicMergePatchType, []byte(patch)) + if err != nil { + deploy.logger.Warn("error patching executor instance ID of service", zap.Error(err), + zap.String("service", svcName), zap.String("ns", svcNamespace)) + return nil, err + } + } return existingSvc, err } else if k8s_err.IsNotFound(err) { service := &apiv1.Service{ ObjectMeta: metav1.ObjectMeta{ - Name: svcName, - Labels: deployLabels, + Name: svcName, + Labels: deployLabels, + Annotations: deployAnnotations, }, Spec: apiv1.ServiceSpec{ Ports: []apiv1.ServicePort{ diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr.go b/pkg/executor/executortype/newdeploy/newdeploymgr.go index 38361657..aa81d22c 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr.go @@ -19,9 +19,11 @@ package newdeploy import ( "context" "fmt" + "math/rand" "os" "strconv" "strings" + "sync" "time" multierror "github.com/hashicorp/go-multierror" @@ -171,7 +173,9 @@ func (deploy *NewDeploy) IsValid(fsvc *fscache.FuncSvc) bool { _, err := deploy.kubernetesClient.CoreV1().Services(service[1]).Get(service[0], metav1.GetOptions{}) if err != nil { - deploy.logger.Error("error validating function service address", zap.String("function", fsvc.Function.Name), zap.Error(err)) + if !k8sErrs.IsNotFound(err) { + deploy.logger.Error("error validating function service address", zap.String("function", fsvc.Function.Name), zap.Error(err)) + } return false } @@ -184,7 +188,9 @@ func (deploy *NewDeploy) IsValid(fsvc *fscache.FuncSvc) bool { currentDeploy, err := deploy.kubernetesClient.AppsV1(). Deployments(deployObj.Namespace).Get(deployObj.Name, metav1.GetOptions{}) if err != nil { - deploy.logger.Error("error validating function deployment", zap.Error(err), zap.String("function", fsvc.Function.Name)) + if !k8sErrs.IsNotFound(err) { + deploy.logger.Error("error validating function deployment", zap.Error(err), zap.String("function", fsvc.Function.Name)) + } return false } @@ -235,6 +241,78 @@ func (deploy *NewDeploy) RefreshFuncPods(logger *zap.Logger, f fv1.Function) err return nil } +func (deploy *NewDeploy) AdoptOrphanResources() { + l := map[string]string{ + types.EXECUTOR_TYPE: string(fv1.ExecutorTypeNewdeploy), + } + + podList, err := deploy.kubernetesClient.CoreV1().Pods(metav1.NamespaceAll).List(metav1.ListOptions{ + LabelSelector: labels.Set(l).AsSelector().String(), + }) + + if err != nil { + deploy.logger.Error("error getting pod list", zap.Error(err)) + return + } + + podWG := &sync.WaitGroup{} + + for i := range podList.Items { + pod := &podList.Items[i] + if !utils.IsReadyPod(pod) { + continue + } + + podWG.Add(1) + go func() { + defer podWG.Done() + + // avoid too many requests arrive Kubernetes API server at the same time. + time.Sleep(time.Duration(rand.Intn(30)) * time.Millisecond) + + patch := fmt.Sprintf(`{"metadata":{"annotations":{"%v":"%v"}}}`, types.EXECUTOR_INSTANCEID_LABEL, deploy.instanceID) + pod, err = deploy.kubernetesClient.CoreV1().Pods(pod.Namespace).Patch(pod.Name, k8sTypes.StrategicMergePatchType, []byte(patch)) + if err != nil { + // just log the error since it won't affect the function serving + deploy.logger.Warn("error patching executor instance ID of pod", zap.Error(err), + zap.String("pod", pod.Name), zap.String("ns", pod.Namespace)) + return + } + + deploy.logger.Info("adopt newdeploy function pod", + zap.String("pod", pod.Name), zap.Any("labels", pod.Labels), zap.Any("annotations", pod.Annotations)) + }() + } + + fnList, err := deploy.fissionClient.Functions(metav1.NamespaceAll).List(metav1.ListOptions{}) + if err != nil { + deploy.logger.Error("error getting function list", zap.Error(err)) + return + } + + deployWG := &sync.WaitGroup{} + + for i := range fnList.Items { + fn := &fnList.Items[i] + if fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType == fv1.ExecutorTypeNewdeploy { + deployWG.Add(1) + go func() { + defer deployWG.Done() + + _, err = deploy.fnCreate(fn, true) + if err != nil { + deploy.logger.Warn("failed to adopt resources for function", zap.Error(err)) + return + } + deploy.logger.Info("adopt resources for function", zap.String("function", fn.Metadata.Name)) + }() + } + } + + podWG.Wait() + deployWG.Wait() +} + func (deploy *NewDeploy) initFuncController() (k8sCache.Store, k8sCache.Controller) { resyncPeriod := 30 * time.Second listWatch := k8sCache.NewListWatchFromClient(deploy.crdClient, "functions", metav1.NamespaceAll, fields.Everything()) @@ -378,6 +456,7 @@ func (deploy *NewDeploy) fnCreate(fn *fv1.Function, firstcreate bool) (*fscache. objName := deploy.getObjName(fn) deployLabels := deploy.getDeployLabels(fn.Metadata, env.Metadata) + deployAnnotations := deploy.getDeployAnnotations(fn.Metadata) // to support backward compatibility, if the function was created in default ns, we fall back to creating the // deployment of the function in fission-function ns @@ -391,7 +470,7 @@ func (deploy *NewDeploy) fnCreate(fn *fv1.Function, firstcreate bool) (*fscache. // Since newdeploy waits for pods of deployment to be ready, // change the order of kubeObject creation (create service first, // then deployment) to take advantage of waiting time. - svc, err := deploy.createOrGetSvc(deployLabels, objName, ns) + svc, err := deploy.createOrGetSvc(deployLabels, deployAnnotations, objName, ns) if err != nil { deploy.logger.Error("error creating service", zap.Error(err), zap.String("service", objName)) go deploy.cleanupNewdeploy(ns, objName) @@ -399,14 +478,14 @@ func (deploy *NewDeploy) fnCreate(fn *fv1.Function, firstcreate bool) (*fscache. } svcAddress := fmt.Sprintf("%v.%v", svc.Name, svc.Namespace) - depl, err := deploy.createOrGetDeployment(fn, env, objName, deployLabels, ns, firstcreate) + depl, err := deploy.createOrGetDeployment(fn, env, objName, deployLabels, deployAnnotations, ns, firstcreate) if err != nil { deploy.logger.Error("error creating deployment", zap.Error(err), zap.String("deployment", objName)) go deploy.cleanupNewdeploy(ns, objName) return nil, errors.Wrapf(err, "error creating deployment %v", objName) } - hpa, err := deploy.createOrGetHpa(objName, &fn.Spec.InvokeStrategy.ExecutionStrategy, depl) + hpa, err := deploy.createOrGetHpa(objName, &fn.Spec.InvokeStrategy.ExecutionStrategy, depl, deployLabels, deployAnnotations) if err != nil { deploy.logger.Error("error creating HPA", zap.Error(err), zap.String("hpa", objName)) go deploy.cleanupNewdeploy(ns, objName) @@ -607,7 +686,7 @@ func (deploy *NewDeploy) updateFuncDeployment(fn *fv1.Function, env *fv1.Environ ns = fn.Metadata.Namespace } - newDeployment, err := deploy.getDeploymentSpec(fn, env, fnObjName, ns, deployLabels) + newDeployment, err := deploy.getDeploymentSpec(fn, env, fnObjName, ns, deployLabels, deploy.getDeployAnnotations(fn.Metadata)) if err != nil { deploy.updateStatus(fn, err, "failed to get new deployment spec while updating function") return err @@ -665,15 +744,21 @@ func (deploy *NewDeploy) getObjName(fn *fv1.Function) string { } func (deploy *NewDeploy) getDeployLabels(fnMeta metav1.ObjectMeta, envMeta metav1.ObjectMeta) map[string]string { + return map[string]string{ + types.EXECUTOR_TYPE: string(fv1.ExecutorTypeNewdeploy), + types.ENVIRONMENT_NAME: envMeta.Name, + types.ENVIRONMENT_NAMESPACE: envMeta.Namespace, + types.ENVIRONMENT_UID: string(envMeta.UID), + types.FUNCTION_NAME: fnMeta.Name, + types.FUNCTION_NAMESPACE: fnMeta.Namespace, + types.FUNCTION_UID: string(fnMeta.UID), + } +} + +func (deploy *NewDeploy) getDeployAnnotations(fnMeta metav1.ObjectMeta) map[string]string { return map[string]string{ types.EXECUTOR_INSTANCEID_LABEL: deploy.instanceID, - types.EXECUTOR_TYPE: string(fv1.ExecutorTypeNewdeploy), - types.ENVIRONMENT_NAME: envMeta.Name, - types.ENVIRONMENT_NAMESPACE: envMeta.Namespace, - types.ENVIRONMENT_UID: string(envMeta.UID), - types.FUNCTION_NAME: fnMeta.Name, - types.FUNCTION_NAMESPACE: fnMeta.Namespace, - types.FUNCTION_UID: string(fnMeta.UID), + types.FUNCTION_RESOURCE_VERSION: fnMeta.ResourceVersion, } } @@ -739,7 +824,7 @@ func (deploy *NewDeploy) idleObjectReaper() { currentDeploy, err := deploy.kubernetesClient.AppsV1(). Deployments(deployObj.Namespace).Get(deployObj.Name, metav1.GetOptions{}) if err != nil { - deploy.logger.Error("error validating function deployment", zap.Error(err), zap.String("function", fsvc.Function.Name)) + deploy.logger.Error("error getting function deployment", zap.Error(err), zap.String("function", fsvc.Function.Name)) continue } diff --git a/pkg/executor/executortype/poolmgr/gp.go b/pkg/executor/executortype/poolmgr/gp.go index c9481f86..f983991a 100644 --- a/pkg/executor/executortype/poolmgr/gp.go +++ b/pkg/executor/executortype/poolmgr/gp.go @@ -33,8 +33,10 @@ import ( "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" "k8s.io/apimachinery/pkg/labels" + k8sTypes "k8s.io/apimachinery/pkg/types" "k8s.io/apimachinery/pkg/util/intstr" "k8s.io/client-go/kubernetes" @@ -64,9 +66,9 @@ type ( kubernetesClient *kubernetes.Clientset fissionClient *crd.FissionClient instanceId string // poolmgr instance id - labelsForPool map[string]string requestChannel chan *choosePodRequest fetcherConfig *fetcherConfig.Config + stopCh context.CancelFunc } // serialize the choosing of pods so that choices don't conflict @@ -97,6 +99,8 @@ func MakeGenericPool( gpLogger.Info("creating pool", zap.Any("environment", env.Metadata)) + ctx, stopCh := context.WithCancel(context.Background()) + // TODO: in general we need to provide the user a way to configure pools. Initial // replicas, autoscaling params, various timeouts, etc. gp := &GenericPool{ @@ -116,6 +120,7 @@ func MakeGenericPool( instanceId: instanceId, useSvc: false, // defaults off -- svc takes a second or more to become routable, slowing cold start useIstio: enableIstio, // defaults off -- istio integration requires pod relabeling and it takes a second or more to become routable, slowing cold start + stopCh: stopCh, } gp.runtimeImagePullPolicy = utils.GetImagePullPolicy(os.Getenv("RUNTIME_IMAGE_PULL_POLICY")) @@ -127,7 +132,7 @@ func MakeGenericPool( } // Labels for generic deployment/RS/pods. - gp.labelsForPool = gp.getDeployLabels() + //gp.labelsForPool = gp.getDeployLabels() // create the pool err = gp.createPool() @@ -136,31 +141,41 @@ func MakeGenericPool( } gpLogger.Info("deployment created", zap.Any("environment", env.Metadata)) - go gp.choosePodService() + go gp.choosePodService(ctx) return gp, nil } func (gp *GenericPool) getDeployLabels() map[string]string { + return map[string]string{ + types.EXECUTOR_TYPE: string(fv1.ExecutorTypePoolmgr), + types.ENVIRONMENT_NAME: gp.env.Metadata.Name, + types.ENVIRONMENT_NAMESPACE: gp.env.Metadata.Namespace, + types.ENVIRONMENT_UID: string(gp.env.Metadata.UID), + "managed": "true", // this allows us to easily find pods managed by the deployment + } +} + +func (gp *GenericPool) getDeployAnnotations() map[string]string { return map[string]string{ fv1.EXECUTOR_INSTANCEID_LABEL: gp.instanceId, - types.EXECUTOR_TYPE: string(fv1.ExecutorTypePoolmgr), - types.ENVIRONMENT_NAME: gp.env.Metadata.Name, - types.ENVIRONMENT_NAMESPACE: gp.env.Metadata.Namespace, - types.ENVIRONMENT_UID: string(gp.env.Metadata.UID), - "managed": "true", // this allows us to easily find pods managed by the deployment } } // choosePodService serializes the choosing of pods -func (gp *GenericPool) choosePodService() { - for req := range gp.requestChannel { - pod, err := gp._choosePod(req.newLabels) - if err != nil { - req.responseChannel <- &choosePodResponse{error: err} - continue +func (gp *GenericPool) choosePodService(ctx context.Context) { + for { + select { + case req := <-gp.requestChannel: + pod, err := gp._choosePod(req.newLabels) + if err != nil { + req.responseChannel <- &choosePodResponse{error: err} + continue + } + req.responseChannel <- &choosePodResponse{pod: pod} + case <-ctx.Done(): + return } - req.responseChannel <- &choosePodResponse{pod: pod} } } @@ -330,12 +345,15 @@ func (gp *GenericPool) specializePod(ctx context.Context, pod *apiv1.Pod, fn *fv // getPoolName returns a unique name of an environment func (gp *GenericPool) getPoolName() string { - return strings.ToLower(fmt.Sprintf("poolmgr-%v-%v-%v", gp.env.Metadata.Name, gp.env.Metadata.Namespace, uniuri.NewLen(8))) + return strings.ToLower(fmt.Sprintf("poolmgr-%v-%v-%v", gp.env.Metadata.Name, gp.env.Metadata.Namespace, gp.env.Metadata.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.getDeployLabels() + deployAnnotations := gp.getDeployAnnotations() + // Use long terminationGracePeriodSeconds for connection draining in case that // pod still runs user functions. gracePeriodSeconds := int64(6 * 60) @@ -350,6 +368,9 @@ func (gp *GenericPool) createPool() error { if gp.useIstio && gp.env.Spec.AllowAccessToExternalNetwork { podAnnotations["sidecar.istio.io/inject"] = "false" } + for k, v := range deployAnnotations { + podAnnotations[k] = v + } container, err := util.MergeContainer(&apiv1.Container{ Name: gp.env.Metadata.Name, @@ -391,7 +412,7 @@ func (gp *GenericPool) createPool() error { pod := apiv1.PodTemplateSpec{ ObjectMeta: metav1.ObjectMeta{ - Labels: gp.labelsForPool, + Labels: deployLabels, Annotations: podAnnotations, }, Spec: apiv1.PodSpec{ @@ -408,13 +429,14 @@ func (gp *GenericPool) createPool() error { deployment := &appsv1.Deployment{ ObjectMeta: metav1.ObjectMeta{ - Name: gp.getPoolName(), - Labels: gp.labelsForPool, + Name: gp.getPoolName(), + Labels: deployLabels, + Annotations: deployAnnotations, }, Spec: appsv1.DeploymentSpec{ Replicas: &gp.replicas, Selector: &metav1.LabelSelector{ - MatchLabels: gp.labelsForPool, + MatchLabels: deployLabels, }, Template: pod, }, @@ -434,11 +456,25 @@ func (gp *GenericPool) createPool() error { deployment.Spec.Template.Spec = *newPodSpec } + _, err = gp.kubernetesClient.AppsV1().Deployments(gp.namespace).Get(deployment.Name, metav1.GetOptions{}) + if err == nil { + patch := fmt.Sprintf(`{"metadata":{"annotations":{"%v":"%v"}}}`, fv1.EXECUTOR_INSTANCEID_LABEL, gp.instanceId) + depl, err := gp.kubernetesClient.AppsV1().Deployments(gp.namespace).Patch(deployment.Name, k8sTypes.StrategicMergePatchType, []byte(patch)) + if err == nil { + gp.deployment = depl + return nil + } + } 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(deployment) 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 } @@ -497,7 +533,7 @@ func (gp *GenericPool) createSvc(name string, labels map[string]string) (*apiv1. Ports: []apiv1.ServicePort{ { Protocol: apiv1.ProtocolTCP, - Port: 80, + Port: 8888, TargetPort: intstr.FromInt(8888), }, }, @@ -584,7 +620,7 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac // the fission router isn't in the same namespace, so return a // namespace-qualified hostname - svcHost = fmt.Sprintf("%v.%v", svcName, gp.namespace) + svcHost = fmt.Sprintf("%v.%v:8888", svcName, gp.namespace) } else if gp.useIstio { svc := utils.GetFunctionIstioServiceName(fn.Metadata.Name, fn.Metadata.Namespace) svcHost = fmt.Sprintf("%v.%v:8888", svc, gp.namespace) @@ -592,6 +628,18 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac svcHost = fmt.Sprintf("%v:8888", pod.Status.PodIP) } + // patch svc-host and resource version to the pod annotations for new executor to adopt the pod + patch := fmt.Sprintf(`{"metadata":{"annotations":{"%v":"%v","%v":"%v"}}}`, + types.ANNOTATION_SVC_HOST, svcHost, types.FUNCTION_RESOURCE_VERSION, fn.Metadata.ResourceVersion) + p, err := gp.kubernetesClient.CoreV1().Pods(pod.Namespace).Patch(pod.Name, k8sTypes.StrategicMergePatchType, []byte(patch)) + 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), + 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), @@ -634,10 +682,13 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac // destroys the pool -- the deployment, replicaset and pods func (gp *GenericPool) destroy() error { + gp.stopCh() + deletePropagation := metav1.DeletePropagationBackground delOpt := metav1.DeleteOptions{ PropagationPolicy: &deletePropagation, } + err := gp.kubernetesClient.AppsV1(). Deployments(gp.namespace).Delete(gp.deployment.ObjectMeta.Name, &delOpt) if err != nil { diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index 84c35070..b4c7cc41 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -18,13 +18,16 @@ package poolmgr import ( "context" + "fmt" + "math/rand" "os" "strconv" "strings" + "sync" "time" - "github.com/fission/fission/pkg/utils" "go.uber.org/zap" + apiv1 "k8s.io/api/core/v1" k8serrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" @@ -40,6 +43,7 @@ import ( "github.com/fission/fission/pkg/executor/reaper" fetcherConfig "github.com/fission/fission/pkg/fetcher/config" "github.com/fission/fission/pkg/types" + "github.com/fission/fission/pkg/utils" ) var _ executortype.ExecutorType = &GenericPoolManager{} @@ -110,6 +114,7 @@ func MakeGenericPoolManager( idlePodReapTime: 2 * time.Minute, fetcherConfig: fetcherConfig, } + go gpm.service() go gpm.eagerPoolCreator() @@ -238,6 +243,140 @@ func (gpm *GenericPoolManager) RefreshFuncPods(logger *zap.Logger, f fv1.Functio return nil } +func (gpm *GenericPoolManager) AdoptOrphanResources() { + envs, err := gpm.fissionClient.Environments(metav1.NamespaceAll).List(metav1.ListOptions{}) + if err != nil { + gpm.logger.Error("error getting environment list", zap.Error(err)) + return + } + + envMap := make(map[string]fv1.Environment, len(envs.Items)) + envWG := &sync.WaitGroup{} + + for i := range envs.Items { + env := envs.Items[i] + + if gpm.getEnvPoolsize(&env) > 0 { + envWG.Add(1) + go func() { + defer envWG.Done() + _, err := gpm.getPool(&env) + if err != nil { + gpm.logger.Error("adopt pool failed", zap.Error(err)) + } + }() + } + + // create environment map for later use + key := fmt.Sprintf("%v/%v", env.Metadata.Namespace, env.Metadata.Name) + envMap[key] = env + } + + l := map[string]string{ + types.EXECUTOR_TYPE: string(fv1.ExecutorTypePoolmgr), + } + + podList, err := gpm.kubernetesClient.CoreV1().Pods(metav1.NamespaceAll).List(metav1.ListOptions{ + LabelSelector: labels.Set(l).AsSelector().String(), + }) + + if err != nil { + gpm.logger.Error("error getting pod list", zap.Error(err)) + return + } + + podWG := &sync.WaitGroup{} + + for i := range podList.Items { + pod := &podList.Items[i] + if !utils.IsReadyPod(pod) { + continue + } + + podWG.Add(1) + go func() { + defer podWG.Done() + + // avoid too many requests arrive Kubernetes API server at the same time. + time.Sleep(time.Duration(rand.Intn(30)) * time.Millisecond) + + patch := fmt.Sprintf(`{"metadata":{"annotations":{"%v":"%v"}}}`, types.EXECUTOR_INSTANCEID_LABEL, gpm.instanceId) + pod, err = gpm.kubernetesClient.CoreV1().Pods(pod.Namespace).Patch(pod.Name, k8sTypes.StrategicMergePatchType, []byte(patch)) + if err != nil { + // just log the error since it won't affect the function serving + gpm.logger.Warn("error patching executor instance ID of pod", zap.Error(err), + zap.String("pod", pod.Name), zap.String("ns", pod.Namespace)) + return + } + + // for unspecialized pod, we only update its annotations + if pod.Labels["managed"] == "true" { + return + } + + fnName, ok1 := pod.Labels[types.FUNCTION_NAME] + fnNS, ok2 := pod.Labels[types.FUNCTION_NAMESPACE] + fnUID, ok3 := pod.Labels[types.FUNCTION_UID] + fnRV, ok4 := pod.Annotations[types.FUNCTION_RESOURCE_VERSION] + envName, ok5 := pod.Labels[types.ENVIRONMENT_NAME] + envNS, ok6 := pod.Labels[types.ENVIRONMENT_NAMESPACE] + svcHost, ok7 := pod.Annotations[types.ANNOTATION_SVC_HOST] + env, ok8 := envMap[fmt.Sprintf("%v/%v", envNS, envName)] + + if !(ok1 && ok2 && ok3 && ok4 && ok5 && ok6 && ok7 && ok8) { + gpm.logger.Warn("failed to adopt pod for function due to lack necessary information", + zap.String("pod", pod.Name), zap.Any("labels", pod.Labels), zap.Any("annotations", pod.Annotations), + zap.String("env", env.Metadata.Name)) + return + } + + fsvc := fscache.FuncSvc{ + Name: pod.Name, + Function: &metav1.ObjectMeta{ + Name: fnName, + Namespace: fnNS, + UID: k8sTypes.UID(fnUID), + ResourceVersion: fnRV, + }, + Environment: &env, + Address: svcHost, + KubernetesObjects: []apiv1.ObjectReference{ + { + Kind: "pod", + Name: pod.Name, + APIVersion: pod.TypeMeta.APIVersion, + Namespace: pod.ObjectMeta.Namespace, + ResourceVersion: pod.ObjectMeta.ResourceVersion, + UID: pod.ObjectMeta.UID, + }, + }, + Executor: fv1.ExecutorTypePoolmgr, + Ctime: time.Now(), + Atime: time.Now(), + } + + // If fsvc already exists we just skip the duplicate one. And let reaper to recycle the duplicate pod. + // This is for the case that there are multiple function pods for the same function due to unknown reason. + _, err := gpm.fsCache.GetByFunction(fsvc.Function) + if err == nil { + return + } + + _, err = gpm.fsCache.Add(fsvc) + if err != nil { + gpm.logger.Warn("failed to adopt pod for function", zap.Error(err), zap.String("pod", pod.Name)) + return + } + + gpm.logger.Info("adopt function pod", + zap.String("pod", pod.Name), zap.Any("labels", pod.Labels), zap.Any("annotations", pod.Annotations)) + }() + } + + envWG.Wait() + podWG.Wait() +} + func (gpm *GenericPoolManager) service() { for { req := <-gpm.requestChannel @@ -353,19 +492,27 @@ func (gpm *GenericPoolManager) eagerPoolCreator() { // creating pools for envs that are actually used by functions. Also we might want // to keep these eagerly created pools smaller than the ones created when there are // actual function calls. + + wg := &sync.WaitGroup{} + for i := range envs.Items { env := envs.Items[i] // Create pool only if poolsize greater than zero if gpm.getEnvPoolsize(&env) > 0 { - _, err := gpm.getPool(&envs.Items[i]) - if err != nil { - gpm.logger.Error("eager-create pool failed", zap.Error(err)) - } + wg.Add(1) + go func() { + defer wg.Done() + _, err := gpm.getPool(&env) + if err != nil { + gpm.logger.Error("eager-create pool failed", zap.Error(err)) + } + }() } } // Clean up pools whose env was deleted gpm.cleanupPools(envs.Items) + wg.Wait() time.Sleep(pollSleep) } } diff --git a/pkg/executor/reaper/reaper.go b/pkg/executor/reaper/reaper.go index 5f3ad061..a37b132b 100644 --- a/pkg/executor/reaper/reaper.go +++ b/pkg/executor/reaper/reaper.go @@ -38,16 +38,20 @@ var ( // CleanupOldExecutorObjects cleans up resources created by old executor instances func CleanupOldExecutorObjects(logger *zap.Logger, kubernetesClient *kubernetes.Clientset, instanceId string) { - go func() { - err := cleanup(logger, kubernetesClient, instanceId) - if err != nil { - // TODO retry reaper; logged and ignored for now - logger.Error("Failed to cleanup old executor objects", zap.Error(err)) - } - }() + err := cleanup(logger, kubernetesClient, instanceId) + if err != nil { + // TODO retry reaper; logged and ignored for now + logger.Error("Failed to cleanup old executor objects", zap.Error(err)) + } } func cleanup(logger *zap.Logger, client *kubernetes.Clientset, instanceId string) error { + // Pods might still be running user functions, so we give them + // a few minutes before terminating them. This time is the + // maximum function runtime, plus the time a router might + // still route to an old instance, i.e. router cache expiry + // time. + time.Sleep(6 * time.Minute) err := cleanupServices(logger, client, instanceId) if err != nil { @@ -67,13 +71,6 @@ func cleanup(logger *zap.Logger, client *kubernetes.Clientset, instanceId string return err } - // Pods might still be running user functions, so we give them - // a few minutes before terminating them. This time is the - // maximum function runtime, plus the time a router might - // still route to an old instance, i.e. router cache expiry - // time. - time.Sleep(6 * time.Minute) - err = cleanupPods(logger, client, instanceId) if err != nil { return err @@ -121,7 +118,7 @@ func cleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instan return err } for _, dep := range deploymentList.Items { - id, ok := dep.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL] + id, ok := dep.ObjectMeta.Annotations[types.EXECUTOR_INSTANCEID_LABEL] if ok && id != instanceId { logger.Debug("cleaning up deployment", zap.String("deployment", dep.ObjectMeta.Name)) err := client.AppsV1().Deployments(dep.ObjectMeta.Namespace).Delete(dep.ObjectMeta.Name, &delOpt) @@ -134,7 +131,7 @@ func cleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instan // ignore err } // Backward compatibility with older label name - pid, pok := dep.ObjectMeta.Labels[types.POOLMGR_INSTANCEID_LABEL] + pid, pok := dep.ObjectMeta.Annotations[types.POOLMGR_INSTANCEID_LABEL] if pok && pid != instanceId { logger.Debug("cleaning up deployment", zap.String("deployment", dep.ObjectMeta.Name)) err := client.AppsV1().Deployments(dep.ObjectMeta.Namespace).Delete(dep.ObjectMeta.Name, &delOpt) @@ -156,7 +153,7 @@ func cleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceId st return err } for _, pod := range podList.Items { - id, ok := pod.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL] + id, ok := pod.ObjectMeta.Annotations[types.EXECUTOR_INSTANCEID_LABEL] if ok && id != instanceId { logger.Debug("cleaning up pod", zap.String("pod", pod.ObjectMeta.Name)) err := client.CoreV1().Pods(pod.ObjectMeta.Namespace).Delete(pod.ObjectMeta.Name, nil) @@ -169,7 +166,7 @@ func cleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceId st // ignore err } // Backward compatibility with older label name - pid, pok := pod.ObjectMeta.Labels[types.POOLMGR_INSTANCEID_LABEL] + pid, pok := pod.ObjectMeta.Annotations[types.POOLMGR_INSTANCEID_LABEL] if pok && pid != instanceId { logger.Debug("cleaning up pod", zap.String("pod", pod.ObjectMeta.Name)) err := client.CoreV1().Pods(pod.ObjectMeta.Namespace).Delete(pod.ObjectMeta.Name, nil) @@ -190,7 +187,7 @@ func cleanupServices(logger *zap.Logger, client *kubernetes.Clientset, instanceI return err } for _, svc := range svcList.Items { - id, ok := svc.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL] + id, ok := svc.ObjectMeta.Annotations[types.EXECUTOR_INSTANCEID_LABEL] if ok && id != instanceId { logger.Debug("cleaning up service", zap.String("service", svc.ObjectMeta.Name)) err := client.CoreV1().Services(svc.ObjectMeta.Namespace).Delete(svc.ObjectMeta.Name, nil) @@ -213,7 +210,7 @@ func cleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceId str } for _, hpa := range hpaList.Items { - id, ok := hpa.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL] + id, ok := hpa.ObjectMeta.Annotations[types.EXECUTOR_INSTANCEID_LABEL] if ok && id != instanceId { logger.Debug("cleaning up HPA", zap.String("hpa", hpa.ObjectMeta.Name)) err := client.AutoscalingV1().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Delete(hpa.ObjectMeta.Name, nil) @@ -225,7 +222,6 @@ func cleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceId str } // ignore err } - } return nil @@ -235,6 +231,9 @@ func cleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceId str // deletes the rolebindings completely if there are no Service Accounts in a rolebinding object. func CleanupRoleBindings(logger *zap.Logger, client *kubernetes.Clientset, fissionClient *crd.FissionClient, functionNs, envBuilderNs string, cleanupRoleBindingInterval time.Duration) { for { + // some sleep before the next reaper iteration + time.Sleep(cleanupRoleBindingInterval) + logger.Debug("starting cleanupRoleBindings cycle") // get all rolebindings ( just to be efficient, one call to kubernetes ) rbList, err := client.RbacV1beta1().RoleBindings(meta_v1.NamespaceAll).List(meta_v1.ListOptions{}) @@ -363,8 +362,5 @@ func CleanupRoleBindings(logger *zap.Logger, client *kubernetes.Clientset, fissi } } } - - // some sleep before the next reaper iteration - time.Sleep(cleanupRoleBindingInterval) } } diff --git a/pkg/types/types.go b/pkg/types/types.go index 852fa152..7b5804ac 100644 --- a/pkg/types/types.go +++ b/pkg/types/types.go @@ -121,13 +121,18 @@ const ( // executor kubernetes object label key const ( - ENVIRONMENT_NAMESPACE = "environmentNamespace" - ENVIRONMENT_NAME = "environmentName" - ENVIRONMENT_UID = "environmentUid" - FUNCTION_NAMESPACE = "functionNamespace" - FUNCTION_NAME = "functionName" - FUNCTION_UID = "functionUid" - EXECUTOR_TYPE = "executorType" + ENVIRONMENT_NAMESPACE = "environmentNamespace" + ENVIRONMENT_NAME = "environmentName" + ENVIRONMENT_UID = "environmentUid" + FUNCTION_NAMESPACE = "functionNamespace" + FUNCTION_NAME = "functionName" + FUNCTION_UID = "functionUid" + FUNCTION_RESOURCE_VERSION = "functionResourceVersion" + EXECUTOR_TYPE = "executorType" +) + +const ( + ANNOTATION_SVC_HOST = "svcHost" ) const ( diff --git a/test/utils.sh b/test/utils.sh index bf37501a..6bb5cf2b 100755 --- a/test/utils.sh +++ b/test/utils.sh @@ -166,8 +166,8 @@ export FISSION_NATS_STREAMING_URL="http://defaultFissionAuthToken@$(kubectl -n $ ## To change the environment image setting for CI test, please refer run_all_tests() in test_utils.sh. export PYTHON_RUNTIME_IMAGE=${PYTHON_RUNTIME_IMAGE:-fission/python-env} export PYTHON_BUILDER_IMAGE=${PYTHON_BUILDER_IMAGE:-fission/python-builder} -export GO_RUNTIME_IMAGE=${GO_RUNTIME_IMAGE:-fission/go-env} -export GO_BUILDER_IMAGE=${GO_BUILDER_IMAGE:-fission/go-builder} +export GO_RUNTIME_IMAGE=${GO_RUNTIME_IMAGE:-fission/go-env-1.12} +export GO_BUILDER_IMAGE=${GO_BUILDER_IMAGE:-fission/go-builder-1.12} export JVM_RUNTIME_IMAGE=${JVM_RUNTIME_IMAGE:-fission/jvm-env} export JVM_BUILDER_IMAGE=${JVM_BUILDER_IMAGE:-fission/jvm-builder} export NODE_RUNTIME_IMAGE=${NODE_RUNTIME_IMAGE:-fission/node-env}