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}