From 19ae7d5ac64b4cf85c066f005dd0a512861272fd Mon Sep 17 00:00:00 2001 From: Ta-Ching Chen Date: Fri, 29 Nov 2019 21:52:59 +0800 Subject: [PATCH] Fix adopted deployment uses old fetcher image (#1447) When a new executor starts up, it adopts the orphan kubernetes resources created by the old executor instance. However, the adopted resource won't reflect the changes come with the new executor, for example, the fetcher image inside won't be changed. To solve this, executor updates the resource spec (HPA/Deployment/Service) with the latest resources spec. By doing this, we can prevent the inconsistency between resources created by different executor instance, also minimizes the impact on users. --- .../executortype/newdeploy/newdeploy.go | 121 ++++++++++-------- pkg/executor/executortype/poolmgr/gp.go | 4 +- 2 files changed, 69 insertions(+), 56 deletions(-) diff --git a/pkg/executor/executortype/newdeploy/newdeploy.go b/pkg/executor/executortype/newdeploy/newdeploy.go index 7a634c92..44fd8e2f 100644 --- a/pkg/executor/executortype/newdeploy/newdeploy.go +++ b/pkg/executor/executortype/newdeploy/newdeploy.go @@ -29,7 +29,6 @@ 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" @@ -56,14 +55,23 @@ func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Enviro minScale = 1 } + deployment, err := deploy.getDeploymentSpec(fn, env, &minScale, deployName, deployNamespace, deployLabels, deployAnnotations) + if err != nil { + return nil, err + } + existingDepl, err := deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Get(deployName, metav1.GetOptions{}) if err == nil { // Try to adopt orphan deployment created by the old executor. if existingDepl.Annotations[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)) + existingDepl.Annotations = deployment.Annotations + existingDepl.Labels = deployment.Labels + existingDepl.Spec.Template.Spec.Containers = deployment.Spec.Template.Spec.Containers + existingDepl.Spec.Template.Spec.ServiceAccountName = deployment.Spec.Template.Spec.ServiceAccountName + existingDepl.Spec.Template.Spec.TerminationGracePeriodSeconds = deployment.Spec.Template.Spec.TerminationGracePeriodSeconds + existingDepl, err = deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Update(existingDepl) if err != nil { - deploy.logger.Warn("error patching executor instance ID of deploy", zap.Error(err), + deploy.logger.Warn("error adopting deploy", zap.Error(err), zap.String("deploy", deployName), zap.String("ns", deployNamespace)) return nil, err } @@ -89,11 +97,6 @@ func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Enviro return nil, err } - deployment, err := deploy.getDeploymentSpec(fn, env, &minScale, deployName, deployNamespace, deployLabels, deployAnnotations) - if err != nil { - return nil, err - } - depl, err := deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Create(deployment) if err != nil { deploy.logger.Error("error while creating function deployment", @@ -339,7 +342,9 @@ 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, deployLabels map[string]string, deployAnnotations map[string]string) (*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") } @@ -354,38 +359,41 @@ func (deploy *NewDeploy) createOrGetHpa(hpaName string, execStrategy *fv1.Execut } targetCPU := int32(execStrategy.TargetCPUPercent) + hpa := &asv1.HorizontalPodAutoscaler{ + ObjectMeta: metav1.ObjectMeta{ + Name: hpaName, + Labels: deployLabels, + Annotations: deployAnnotations, + }, + Spec: asv1.HorizontalPodAutoscalerSpec{ + ScaleTargetRef: asv1.CrossVersionObjectReference{ + Kind: DeploymentKind, + Name: depl.ObjectMeta.Name, + APIVersion: DeploymentVersion, + }, + MinReplicas: &minRepl, + MaxReplicas: maxRepl, + TargetCPUUtilizationPercentage: &targetCPU, + }, + } + existingHpa, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Get(hpaName, metav1.GetOptions{}) if err == nil { + // to adopt orphan service if existingHpa.Annotations[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)) + existingHpa.Annotations = hpa.Annotations + existingHpa.Labels = hpa.Labels + existingHpa.Spec = hpa.Spec + existingHpa, err = deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Update(existingHpa) if err != nil { - deploy.logger.Warn("error patching executor instance ID of HPA", zap.Error(err), + deploy.logger.Warn("error adopting HPA", zap.Error(err), zap.String("HPA", hpaName), zap.String("ns", depl.ObjectMeta.Namespace)) return nil, err } } return existingHpa, err } else if k8s_err.IsNotFound(err) { - hpa := asv1.HorizontalPodAutoscaler{ - ObjectMeta: metav1.ObjectMeta{ - Name: hpaName, - Labels: deployLabels, - Annotations: deployAnnotations, - }, - Spec: asv1.HorizontalPodAutoscalerSpec{ - ScaleTargetRef: asv1.CrossVersionObjectReference{ - Kind: DeploymentKind, - Name: depl.ObjectMeta.Name, - APIVersion: DeploymentVersion, - }, - MinReplicas: &minRepl, - MaxReplicas: maxRepl, - TargetCPUUtilizationPercentage: &targetCPU, - }, - } - - cHpa, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Create(&hpa) + cHpa, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Create(hpa) if err != nil { return nil, err } @@ -410,38 +418,43 @@ func (deploy *NewDeploy) deleteHpa(ns string, name string) error { } func (deploy *NewDeploy) createOrGetSvc(deployLabels map[string]string, deployAnnotations map[string]string, svcName string, svcNamespace string) (*apiv1.Service, error) { + service := &apiv1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: svcName, + Labels: deployLabels, + Annotations: deployAnnotations, + }, + Spec: apiv1.ServiceSpec{ + Ports: []apiv1.ServicePort{ + { + Name: "http-env", + Port: int32(80), + TargetPort: intstr.FromInt(8888), + }, + }, + Selector: deployLabels, + Type: apiv1.ServiceTypeClusterIP, + }, + } + existingSvc, err := deploy.kubernetesClient.CoreV1().Services(svcNamespace).Get(svcName, metav1.GetOptions{}) if err == nil { + // to adopt orphan service if existingSvc.Annotations[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)) + existingSvc.Annotations = service.Annotations + existingSvc.Labels = service.Labels + existingSvc.Spec.Ports = service.Spec.Ports + existingSvc.Spec.Selector = service.Spec.Selector + existingSvc.Spec.Type = service.Spec.Type + existingSvc, err = deploy.kubernetesClient.CoreV1().Services(svcNamespace).Update(existingSvc) if err != nil { - deploy.logger.Warn("error patching executor instance ID of service", zap.Error(err), + deploy.logger.Warn("error adopting 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, - Annotations: deployAnnotations, - }, - Spec: apiv1.ServiceSpec{ - Ports: []apiv1.ServicePort{ - { - Name: "http-env", - Port: int32(80), - TargetPort: intstr.FromInt(8888), - }, - }, - Selector: deployLabels, - Type: apiv1.ServiceTypeClusterIP, - }, - } - svc, err := deploy.kubernetesClient.CoreV1().Services(svcNamespace).Create(service) if err != nil { return nil, err diff --git a/pkg/executor/executortype/poolmgr/gp.go b/pkg/executor/executortype/poolmgr/gp.go index 028fe84b..8de8350d 100644 --- a/pkg/executor/executortype/poolmgr/gp.go +++ b/pkg/executor/executortype/poolmgr/gp.go @@ -459,8 +459,8 @@ func (gp *GenericPool) createPool() error { depl, err := gp.kubernetesClient.AppsV1().Deployments(gp.namespace).Get(deployment.Name, metav1.GetOptions{}) if err == nil { if depl.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] != gp.instanceId { - 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)) + deployment.Annotations[fv1.EXECUTOR_INSTANCEID_LABEL] = gp.instanceId + depl, err = gp.kubernetesClient.AppsV1().Deployments(gp.namespace).Update(deployment) } gp.deployment = depl return err