From baeee62a26e626b7a5019cb6aa6162399df4b675 Mon Sep 17 00:00:00 2001 From: Stefan Prodan Date: Thu, 11 Oct 2018 19:59:40 +0300 Subject: [PATCH] Controller refactoring - split controller logic into components (deployer, observer, router and scheduler) - set the canary analysis final state (failed or finished) in a single run --- pkg/controller/controller.go | 20 +- pkg/controller/deployer.go | 494 +++++++++++++++++++---------------- pkg/controller/deployment.go | 422 ------------------------------ pkg/controller/observer.go | 10 +- pkg/controller/router.go | 285 ++++++++++++++++++++ pkg/controller/scheduler.go | 266 +++++++++++++++++++ pkg/controller/utils.go | 79 ------ 7 files changed, 845 insertions(+), 731 deletions(-) delete mode 100644 pkg/controller/deployment.go create mode 100644 pkg/controller/router.go create mode 100644 pkg/controller/scheduler.go delete mode 100644 pkg/controller/utils.go diff --git a/pkg/controller/controller.go b/pkg/controller/controller.go index 72273006..460b7f40 100644 --- a/pkg/controller/controller.go +++ b/pkg/controller/controller.go @@ -123,7 +123,7 @@ func (c *Controller) Run(threadiness int, stopCh <-chan struct{}) error { for { select { case <-tickChan: - c.doRollouts() + c.scheduleCanaries() case <-stopCh: c.logger.Info("Shutting down operator workers") return nil @@ -182,13 +182,13 @@ func (c *Controller) syncHandler(key string) error { c.rollouts.Store(fmt.Sprintf("%s.%s", cd.Name, cd.Namespace), cd) - if cd.Spec.TargetRef.Kind == "Deployment" { - err = c.bootstrapDeployment(cd) - if err != nil { - c.logger.Warnf("%s.%s bootstrap error %v", cd.Name, cd.Namespace, err) - return err - } - } + //if cd.Spec.TargetRef.Kind == "Deployment" { + // err = c.bootstrapDeployment(cd) + // if err != nil { + // c.logger.Warnf("%s.%s bootstrap error %v", cd.Name, cd.Namespace, err) + // return err + // } + //} c.logger.Infof("Synced %s", key) @@ -229,3 +229,7 @@ func (c *Controller) recordEventWarningf(r *flaggerv1.Canary, template string, a c.logger.Infof(template, args...) c.recorder.Event(r, corev1.EventTypeWarning, "Synced", fmt.Sprintf(template, args...)) } + +func int32p(i int32) *int32 { + return &i +} diff --git a/pkg/controller/deployer.go b/pkg/controller/deployer.go index a3976a9e..3adbd2c1 100644 --- a/pkg/controller/deployer.go +++ b/pkg/controller/deployer.go @@ -1,29 +1,219 @@ package controller import ( + "encoding/base64" + "encoding/json" "fmt" - istiov1alpha3 "github.com/knative/pkg/apis/istio/v1alpha3" + "github.com/google/go-cmp/cmp" + "github.com/google/go-cmp/cmp/cmpopts" + istioclientset "github.com/knative/pkg/client/clientset/versioned" flaggerv1 "github.com/stefanprodan/flagger/pkg/apis/flagger/v1alpha1" + clientset "github.com/stefanprodan/flagger/pkg/client/clientset/versioned" + "go.uber.org/zap" appsv1 "k8s.io/api/apps/v1" hpav1 "k8s.io/api/autoscaling/v2beta1" corev1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime/schema" - "k8s.io/apimachinery/pkg/util/intstr" + "k8s.io/client-go/kubernetes" ) -func (c *Controller) bootstrapDeployment(cd *flaggerv1.Canary) error { +type CanaryDeployer struct { + kubeClient kubernetes.Interface + istioClient istioclientset.Interface + flaggerClient clientset.Interface + logger *zap.SugaredLogger +} +// Promote copies the pod spec from canary to primary +func (c *CanaryDeployer) Promote(cd *flaggerv1.Canary) error { + canary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(cd.Spec.TargetRef.Name, metav1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return fmt.Errorf("deployment %s.%s not found", cd.Spec.TargetRef.Name, cd.Namespace) + } + return fmt.Errorf("deployment %s.%s query error %v", cd.Spec.TargetRef.Name, cd.Namespace, err) + } + + primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) + primary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(primaryName, metav1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return fmt.Errorf("deployment %s.%s not found", primaryName, cd.Namespace) + } + return fmt.Errorf("deployment %s.%s query error %v", primaryName, cd.Namespace, err) + } + + primary.Spec.Template.Spec = canary.Spec.Template.Spec + _, err = c.kubeClient.AppsV1().Deployments(primary.Namespace).Update(primary) + if err != nil { + return fmt.Errorf("updating template spec %s.%s failed: %v", primary.GetName(), primary.Namespace, err) + + } + + return nil +} + +// IsDeploymentHealthy checks the primary and canary deployment status and returns an error if +// the deployment is in the middle of a rolling update or if the pods are unhealthy +func (c *CanaryDeployer) IsDeploymentHealthy(cd *flaggerv1.Canary) error { + canary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(cd.Spec.TargetRef.Name, metav1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return fmt.Errorf("deployment %s.%s not found", cd.Spec.TargetRef.Name, cd.Namespace) + } + return fmt.Errorf("deployment %s.%s query error %v", cd.Spec.TargetRef.Name, cd.Namespace, err) + } + if msg, healthy := c.getDeploymentStatus(canary); !healthy { + return fmt.Errorf("Halt %s.%s advancement %s", cd.Name, cd.Namespace, msg) + } + + primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) + primary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(primaryName, metav1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return fmt.Errorf("deployment %s.%s not found", primaryName, cd.Namespace) + } + return fmt.Errorf("deployment %s.%s query error %v", primaryName, cd.Namespace, err) + } + if msg, healthy := c.getDeploymentStatus(primary); !healthy { + return fmt.Errorf("Halt %s.%s advancement %s", cd.Name, cd.Namespace, msg) + } + + if primary.Spec.Replicas == int32p(0) { + return fmt.Errorf("halt %s.%s advancement %s", + cd.Name, cd.Namespace, "primary deployment is scaled to zero") + } + return nil +} + +// IsNewSpec returns true if the canary deployment pod spec has changed +func (c *CanaryDeployer) IsNewSpec(cd *flaggerv1.Canary) (bool, error) { + canary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(cd.Spec.TargetRef.Name, metav1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return false, fmt.Errorf("deployment %s.%s not found", cd.Spec.TargetRef.Name, cd.Namespace) + } + return false, fmt.Errorf("deployment %s.%s query error %v", cd.Spec.TargetRef.Name, cd.Namespace, err) + } + + if cd.Status.CanaryRevision == "" { + return true, nil + } + + newSpec := &canary.Spec.Template.Spec + oldSpecJson, err := base64.StdEncoding.DecodeString(cd.Status.CanaryRevision) + if err != nil { + return false, err + } + oldSpec := &corev1.PodSpec{} + err = json.Unmarshal(oldSpecJson, oldSpec) + if err != nil { + return false, fmt.Errorf("%s.%s unmarshal error %v", cd.Name, cd.Namespace, err) + } + + if diff := cmp.Diff(*newSpec, *oldSpec, cmpopts.IgnoreUnexported(resource.Quantity{})); diff != "" { + //fmt.Println(diff) + return true, nil + } + + return false, nil +} + +// SyncStatus updates the canary status state +func (c *CanaryDeployer) SetFailedChecks(cd *flaggerv1.Canary, val int) error { + cd.Status.FailedChecks = val + cd, err := c.flaggerClient.FlaggerV1alpha1().Canaries(cd.Namespace).Update(cd) + if err != nil { + return fmt.Errorf("deployment %s.%s update error %v", cd.Spec.TargetRef.Name, cd.Namespace, err) + } + return nil +} + +// SyncStatus updates the canary status state +func (c *CanaryDeployer) SetState(cd *flaggerv1.Canary, state string) error { + cd.Status.State = state + cd, err := c.flaggerClient.FlaggerV1alpha1().Canaries(cd.Namespace).Update(cd) + if err != nil { + return fmt.Errorf("deployment %s.%s update error %v", cd.Spec.TargetRef.Name, cd.Namespace, err) + } + return nil +} + +// SyncStatus encodes the canary pod spec and updates the canary status +func (c *CanaryDeployer) SyncStatus(cd *flaggerv1.Canary, status flaggerv1.CanaryStatus) error { + canary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(cd.Spec.TargetRef.Name, metav1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return fmt.Errorf("deployment %s.%s not found", cd.Spec.TargetRef.Name, cd.Namespace) + } + return fmt.Errorf("deployment %s.%s query error %v", cd.Spec.TargetRef.Name, cd.Namespace, err) + } + + specJson, err := json.Marshal(canary.Spec.Template.Spec) + if err != nil { + return fmt.Errorf("deployment %s.%s marshal error %v", cd.Spec.TargetRef.Name, cd.Namespace, err) + } + + specEnc := base64.StdEncoding.EncodeToString(specJson) + cd.Status.State = status.State + cd.Status.FailedChecks = status.FailedChecks + cd.Status.CanaryRevision = specEnc + cd, err = c.flaggerClient.FlaggerV1alpha1().Canaries(cd.Namespace).Update(cd) + if err != nil { + return fmt.Errorf("deployment %s.%s update error %v", cd.Spec.TargetRef.Name, cd.Namespace, err) + } + return nil +} + +// Scale sets the canary deployment replicas +func (c *CanaryDeployer) Scale(cd *flaggerv1.Canary, replicas int32) error { + canary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(cd.Spec.TargetRef.Name, metav1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return fmt.Errorf("deployment %s.%s not found", cd.Spec.TargetRef.Name, cd.Namespace) + } + return fmt.Errorf("deployment %s.%s query error %v", cd.Spec.TargetRef.Name, cd.Namespace, err) + } + canary.Spec.Replicas = int32p(replicas) + canary, err = c.kubeClient.AppsV1().Deployments(canary.Namespace).Update(canary) + if err != nil { + return fmt.Errorf("scaling %s.%s to %v failed: %v", canary.GetName(), canary.Namespace, replicas, err) + + } + + return nil +} + +// Sync creates the primary deployment and hpa +func (c *CanaryDeployer) Sync(cd *flaggerv1.Canary) error { + if err := c.createPrimaryDeployment(cd); err != nil { + return err + } + + if cd.Status.State == "" { + c.Scale(cd, 0) + } + + if cd.Spec.AutoscalerRef.Kind == "HorizontalPodAutoscaler" { + if err := c.createPrimaryHpa(cd); err != nil { + return err + } + } + return nil +} + +func (c *CanaryDeployer) createPrimaryDeployment(cd *flaggerv1.Canary) error { canaryName := cd.Spec.TargetRef.Name primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) canaryDep, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(canaryName, metav1.GetOptions{}) if err != nil { if errors.IsNotFound(err) { - return fmt.Errorf("deployment %s.%s not found, retrying in %v", - canaryName, cd.Namespace, c.rolloutWindow) + return fmt.Errorf("deployment %s.%s not found, retrying", canaryName, cd.Namespace) } else { return err } @@ -67,225 +257,91 @@ func (c *Controller) bootstrapDeployment(cd *flaggerv1.Canary) error { return err } - c.recordEventInfof(cd, "Deployment %s.%s created", primaryDep.GetName(), cd.Namespace) + c.logger.Infof("Deployment %s.%s created", primaryDep.GetName(), cd.Namespace) } - if cd.Status.State == "" { - c.scaleToZeroCanary(cd) - } - - canaryService, err := c.kubeClient.CoreV1().Services(cd.Namespace).Get(canaryName, metav1.GetOptions{}) - if errors.IsNotFound(err) { - canaryService = &corev1.Service{ - ObjectMeta: metav1.ObjectMeta{ - Name: canaryName, - Namespace: cd.Namespace, - OwnerReferences: []metav1.OwnerReference{ - *metav1.NewControllerRef(cd, schema.GroupVersionKind{ - Group: flaggerv1.SchemeGroupVersion.Group, - Version: flaggerv1.SchemeGroupVersion.Version, - Kind: flaggerv1.CanaryKind, - }), - }, - }, - Spec: corev1.ServiceSpec{ - Type: corev1.ServiceTypeClusterIP, - Selector: map[string]string{"app": canaryName}, - Ports: []corev1.ServicePort{ - { - Name: "http", - Protocol: corev1.ProtocolTCP, - Port: cd.Spec.Service.Port, - TargetPort: intstr.IntOrString{ - Type: intstr.Int, - IntVal: cd.Spec.Service.Port, - }, - }, - }, - }, - } - - _, err = c.kubeClient.CoreV1().Services(cd.Namespace).Create(canaryService) - if err != nil { - return err - } - c.recordEventInfof(cd, "Service %s.%s created", canaryService.GetName(), cd.Namespace) - } - - canaryTestServiceName := fmt.Sprintf("%s-canary", cd.Spec.TargetRef.Name) - canaryTestService, err := c.kubeClient.CoreV1().Services(cd.Namespace).Get(canaryTestServiceName, metav1.GetOptions{}) - if errors.IsNotFound(err) { - canaryTestService = &corev1.Service{ - ObjectMeta: metav1.ObjectMeta{ - Name: canaryTestServiceName, - Namespace: cd.Namespace, - OwnerReferences: []metav1.OwnerReference{ - *metav1.NewControllerRef(cd, schema.GroupVersionKind{ - Group: flaggerv1.SchemeGroupVersion.Group, - Version: flaggerv1.SchemeGroupVersion.Version, - Kind: flaggerv1.CanaryKind, - }), - }, - }, - Spec: corev1.ServiceSpec{ - Type: corev1.ServiceTypeClusterIP, - Selector: map[string]string{"app": canaryName}, - Ports: []corev1.ServicePort{ - { - Name: "http", - Protocol: corev1.ProtocolTCP, - Port: cd.Spec.Service.Port, - TargetPort: intstr.IntOrString{ - Type: intstr.Int, - IntVal: cd.Spec.Service.Port, - }, - }, - }, - }, - } - - _, err = c.kubeClient.CoreV1().Services(cd.Namespace).Create(canaryTestService) - if err != nil { - return err - } - c.recordEventInfof(cd, "Service %s.%s created", canaryTestService.GetName(), cd.Namespace) - } - - primaryService, err := c.kubeClient.CoreV1().Services(cd.Namespace).Get(primaryName, metav1.GetOptions{}) - if errors.IsNotFound(err) { - primaryService = &corev1.Service{ - ObjectMeta: metav1.ObjectMeta{ - Name: primaryName, - Namespace: cd.Namespace, - OwnerReferences: []metav1.OwnerReference{ - *metav1.NewControllerRef(cd, schema.GroupVersionKind{ - Group: flaggerv1.SchemeGroupVersion.Group, - Version: flaggerv1.SchemeGroupVersion.Version, - Kind: flaggerv1.CanaryKind, - }), - }, - }, - Spec: corev1.ServiceSpec{ - Type: corev1.ServiceTypeClusterIP, - Selector: map[string]string{"app": primaryName}, - Ports: []corev1.ServicePort{ - { - Name: "http", - Protocol: corev1.ProtocolTCP, - Port: cd.Spec.Service.Port, - TargetPort: intstr.IntOrString{ - Type: intstr.Int, - IntVal: cd.Spec.Service.Port, - }, - }, - }, - }, - } - - _, err = c.kubeClient.CoreV1().Services(cd.Namespace).Create(primaryService) - if err != nil { - return err - } - - c.recordEventInfof(cd, "Service %s.%s created", primaryService.GetName(), cd.Namespace) - } - - hosts := append(cd.Spec.Service.Hosts, canaryName) - gateways := append(cd.Spec.Service.Gateways, "mesh") - virtualService, err := c.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Get(canaryName, metav1.GetOptions{}) - if errors.IsNotFound(err) { - virtualService = &istiov1alpha3.VirtualService{ - ObjectMeta: metav1.ObjectMeta{ - Name: canaryName, - Namespace: cd.Namespace, - OwnerReferences: []metav1.OwnerReference{ - *metav1.NewControllerRef(cd, schema.GroupVersionKind{ - Group: flaggerv1.SchemeGroupVersion.Group, - Version: flaggerv1.SchemeGroupVersion.Version, - Kind: flaggerv1.CanaryKind, - }), - }, - }, - Spec: istiov1alpha3.VirtualServiceSpec{ - Hosts: hosts, - Gateways: gateways, - Http: []istiov1alpha3.HTTPRoute{ - { - Route: []istiov1alpha3.DestinationWeight{ - { - Destination: istiov1alpha3.Destination{ - Host: primaryName, - Port: istiov1alpha3.PortSelector{ - Number: uint32(cd.Spec.Service.Port), - }, - }, - Weight: 100, - }, - { - Destination: istiov1alpha3.Destination{ - Host: canaryName, - Port: istiov1alpha3.PortSelector{ - Number: uint32(cd.Spec.Service.Port), - }, - }, - Weight: 0, - }, - }, - }, - }, - }, - } - - _, err = c.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Create(virtualService) - if err != nil { - return err - } - c.recordEventInfof(cd, "VirtualService %s.%s created", virtualService.GetName(), cd.Namespace) - } - - if cd.Spec.AutoscalerRef.Kind == "HorizontalPodAutoscaler" { - hpa, err := c.kubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Get(cd.Spec.AutoscalerRef.Name, metav1.GetOptions{}) - if err != nil { - if errors.IsNotFound(err) { - return fmt.Errorf("HorizontalPodAutoscaler %s.%s not found, retrying in %v", - cd.Spec.AutoscalerRef.Name, cd.Namespace, c.rolloutWindow) - } else { - return err - } - } - primaryHpaName := fmt.Sprintf("%s-primary", cd.Spec.AutoscalerRef.Name) - primaryHpa, err := c.kubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Get(primaryHpaName, metav1.GetOptions{}) + return nil +} +func (c *CanaryDeployer) createPrimaryHpa(cd *flaggerv1.Canary) error { + primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) + hpa, err := c.kubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Get(cd.Spec.AutoscalerRef.Name, metav1.GetOptions{}) + if err != nil { if errors.IsNotFound(err) { - primaryHpa = &hpav1.HorizontalPodAutoscaler{ - ObjectMeta: metav1.ObjectMeta{ - Name: primaryHpaName, - Namespace: cd.Namespace, - OwnerReferences: []metav1.OwnerReference{ - *metav1.NewControllerRef(cd, schema.GroupVersionKind{ - Group: flaggerv1.SchemeGroupVersion.Group, - Version: flaggerv1.SchemeGroupVersion.Version, - Kind: flaggerv1.CanaryKind, - }), - }, - }, - Spec: hpav1.HorizontalPodAutoscalerSpec{ - ScaleTargetRef: hpav1.CrossVersionObjectReference{ - Name: primaryName, - Kind: hpa.Spec.ScaleTargetRef.Kind, - APIVersion: hpa.Spec.ScaleTargetRef.APIVersion, - }, - MinReplicas: hpa.Spec.MinReplicas, - MaxReplicas: hpa.Spec.MaxReplicas, - Metrics: hpa.Spec.Metrics, - }, - } + return fmt.Errorf("HorizontalPodAutoscaler %s.%s not found, retrying", + cd.Spec.AutoscalerRef.Name, cd.Namespace) + } else { + return err + } + } + primaryHpaName := fmt.Sprintf("%s-primary", cd.Spec.AutoscalerRef.Name) + primaryHpa, err := c.kubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Get(primaryHpaName, metav1.GetOptions{}) - _, err = c.kubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Create(primaryHpa) - if err != nil { - return err - } - c.recordEventInfof(cd, "HorizontalPodAutoscaler %s.%s created", primaryHpa.GetName(), cd.Namespace) + if errors.IsNotFound(err) { + primaryHpa = &hpav1.HorizontalPodAutoscaler{ + ObjectMeta: metav1.ObjectMeta{ + Name: primaryHpaName, + Namespace: cd.Namespace, + OwnerReferences: []metav1.OwnerReference{ + *metav1.NewControllerRef(cd, schema.GroupVersionKind{ + Group: flaggerv1.SchemeGroupVersion.Group, + Version: flaggerv1.SchemeGroupVersion.Version, + Kind: flaggerv1.CanaryKind, + }), + }, + }, + Spec: hpav1.HorizontalPodAutoscalerSpec{ + ScaleTargetRef: hpav1.CrossVersionObjectReference{ + Name: primaryName, + Kind: hpa.Spec.ScaleTargetRef.Kind, + APIVersion: hpa.Spec.ScaleTargetRef.APIVersion, + }, + MinReplicas: hpa.Spec.MinReplicas, + MaxReplicas: hpa.Spec.MaxReplicas, + Metrics: hpa.Spec.Metrics, + }, + } + + _, err = c.kubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Create(primaryHpa) + if err != nil { + return err + } + c.logger.Infof("HorizontalPodAutoscaler %s.%s created", primaryHpa.GetName(), cd.Namespace) + } + + return nil +} + +func (c *CanaryDeployer) getDeploymentStatus(deployment *appsv1.Deployment) (string, bool) { + if deployment.Generation <= deployment.Status.ObservedGeneration { + cond := c.getDeploymentCondition(deployment.Status, appsv1.DeploymentProgressing) + if cond != nil && cond.Reason == "ProgressDeadlineExceeded" { + return fmt.Sprintf("deployment %q exceeded its progress deadline", deployment.GetName()), false + } else if deployment.Spec.Replicas != nil && deployment.Status.UpdatedReplicas < *deployment.Spec.Replicas { + return fmt.Sprintf("waiting for rollout to finish: %d out of %d new replicas have been updated", + deployment.Status.UpdatedReplicas, *deployment.Spec.Replicas), false + } else if deployment.Status.Replicas > deployment.Status.UpdatedReplicas { + return fmt.Sprintf("waiting for rollout to finish: %d old replicas are pending termination", + deployment.Status.Replicas-deployment.Status.UpdatedReplicas), false + } else if deployment.Status.AvailableReplicas < deployment.Status.UpdatedReplicas { + return fmt.Sprintf("waiting for rollout to finish: %d of %d updated replicas are available", + deployment.Status.AvailableReplicas, deployment.Status.UpdatedReplicas), false + } + } else { + return "waiting for rollout to finish: observed deployment generation less then desired generation", false + } + + return "ready", true +} + +func (c *CanaryDeployer) getDeploymentCondition( + status appsv1.DeploymentStatus, + conditionType appsv1.DeploymentConditionType, +) *appsv1.DeploymentCondition { + for i := range status.Conditions { + c := status.Conditions[i] + if c.Type == conditionType { + return &c } } return nil diff --git a/pkg/controller/deployment.go b/pkg/controller/deployment.go deleted file mode 100644 index 778f4609..00000000 --- a/pkg/controller/deployment.go +++ /dev/null @@ -1,422 +0,0 @@ -package controller - -import ( - "fmt" - "time" - - istiov1alpha3 "github.com/knative/pkg/apis/istio/v1alpha3" - flaggerv1 "github.com/stefanprodan/flagger/pkg/apis/flagger/v1alpha1" - appsv1 "k8s.io/api/apps/v1" - "k8s.io/apimachinery/pkg/apis/meta/v1" -) - -func (c *Controller) doRollouts() { - c.rollouts.Range(func(key interface{}, value interface{}) bool { - r := value.(*flaggerv1.Canary) - if r.Spec.TargetRef.Kind == "Deployment" { - go c.advanceDeploymentRollout(r.Name, r.Namespace) - } - return true - }) -} - -func (c *Controller) advanceDeploymentRollout(name string, namespace string) { - // gate stage: check if the rollout exists - r, ok := c.getRollout(name, namespace) - if !ok { - return - } - - err := c.bootstrapDeployment(r) - if err != nil { - c.recordEventWarningf(r, "%v", err) - return - } - - // set max weight default value to 100% - maxWeight := 100 - if r.Spec.CanaryAnalysis.MaxWeight > 0 { - maxWeight = r.Spec.CanaryAnalysis.MaxWeight - } - - // gate stage: check if canary deployment exists and is healthy - canary, ok := c.getCanaryDeployment(r, r.Spec.TargetRef.Name, r.Namespace) - if !ok { - return - } - - // gate stage: check if primary deployment exists and is healthy - primary, ok := c.getDeployment(r, fmt.Sprintf("%s-primary", r.Spec.TargetRef.Name), r.Namespace) - if !ok { - return - } - - // gate stage: check if virtual service exists - // and if it contains weighted destination routes to the primary and canary services - vs, primaryRoute, canaryRoute, ok := c.getVirtualService(r) - if !ok { - return - } - - // gate stage: check if rollout should start (canary revision has changes) or continue - if ok := c.checkRolloutStatus(r, canary); !ok { - return - } - - // gate stage: check if the number of failed checks reached the threshold - if r.Status.State == "running" && r.Status.FailedChecks >= r.Spec.CanaryAnalysis.Threshold { - c.recordEventWarningf(r, "Rolling back %s.%s failed checks threshold reached %v", - r.Name, r.Namespace, r.Status.FailedChecks) - - // route all traffic back to primary - primaryRoute.Weight = 100 - canaryRoute.Weight = 0 - if ok := c.updateVirtualServiceRoutes(r, vs, primaryRoute, canaryRoute); !ok { - return - } - - c.recordEventWarningf(r, "Canary failed! Scaling down %s.%s", - canary.GetName(), canary.Namespace) - - // shutdown canary - c.scaleToZeroCanary(r) - - // mark rollout as failed - c.updateRolloutStatus(r, "promotion-failed") - return - } - - // gate stage: check if the canary success rate is above the threshold - // skip check if no traffic is routed to canary - if canaryRoute.Weight == 0 { - c.recordEventInfof(r, "Starting canary deployment for %s.%s", r.Name, r.Namespace) - } else { - if ok := c.checkDeploymentMetrics(r); !ok { - c.updateRolloutFailedChecks(r, r.Status.FailedChecks+1) - return - } - } - - // routing stage: increase canary traffic percentage - if canaryRoute.Weight < maxWeight { - primaryRoute.Weight -= r.Spec.CanaryAnalysis.StepWeight - if primaryRoute.Weight < 0 { - primaryRoute.Weight = 0 - } - canaryRoute.Weight += r.Spec.CanaryAnalysis.StepWeight - if primaryRoute.Weight > 100 { - primaryRoute.Weight = 100 - } - - if ok := c.updateVirtualServiceRoutes(r, vs, primaryRoute, canaryRoute); !ok { - return - } - - c.recordEventInfof(r, "Advance %s.%s canary weight %v", r.Name, r.Namespace, canaryRoute.Weight) - - // promotion stage: override primary.template.spec with the canary spec - if canaryRoute.Weight == maxWeight { - c.recordEventInfof(r, "Copying %s.%s template spec to %s.%s", - canary.GetName(), canary.Namespace, primary.GetName(), primary.Namespace) - - primary.Spec.Template.Spec = canary.Spec.Template.Spec - _, err := c.kubeClient.AppsV1().Deployments(primary.Namespace).Update(primary) - if err != nil { - c.recordEventErrorf(r, "Updating template spec %s.%s failed: %v", primary.GetName(), primary.Namespace, err) - return - } - } - } else { - // final stage: route all traffic back to primary - primaryRoute.Weight = 100 - canaryRoute.Weight = 0 - if ok := c.updateVirtualServiceRoutes(r, vs, primaryRoute, canaryRoute); !ok { - return - } - - // final stage: mark rollout as finished and scale canary to zero replicas - c.recordEventInfof(r, "Scaling down %s.%s", canary.GetName(), canary.Namespace) - c.scaleToZeroCanary(r) - c.updateRolloutStatus(r, "promotion-finished") - } -} - -func (c *Controller) getRollout(name string, namespace string) (*flaggerv1.Canary, bool) { - r, err := c.rolloutClient.FlaggerV1alpha1().Canaries(namespace).Get(name, v1.GetOptions{}) - if err != nil { - c.logger.Errorf("Canary %s.%s not found", name, namespace) - return nil, false - } - - return r, true -} - -func (c *Controller) checkRolloutStatus(r *flaggerv1.Canary, canary *appsv1.Deployment) bool { - canaryRevision, err := c.getDeploymentSpecEnc(canary) - if err != nil { - c.logger.Errorf("Canary %s.%s not found: %v", r.Name, r.Namespace, err) - return false - } - - if r.Status.State == "" { - r.Status = flaggerv1.CanaryStatus{ - State: "initialized", - CanaryRevision: canaryRevision, - FailedChecks: 0, - } - r, err = c.rolloutClient.FlaggerV1alpha1().Canaries(r.Namespace).Update(r) - if err != nil { - c.logger.Errorf("Canary %s.%s status update failed: %v", r.Name, r.Namespace, err) - return false - } - - c.recordEventInfof(r, "Initialization done! %s.%s", canary.GetName(), canary.Namespace) - return false - } - - if r.Status.State == "running" { - return true - } - - if r.Status.State == "promotion-finished" { - c.setCanaryRevision(r, canary, "finished") - c.logger.Infof("Promotion completed! %s.%s", r.Name, r.Namespace) - return false - } - - if r.Status.State == "promotion-failed" { - c.setCanaryRevision(r, canary, "failed") - c.logger.Infof("Promotion failed! %s.%s", r.Name, r.Namespace) - return false - } - - if diff, err := c.diffDeploymentSpec(r, canary); diff { - c.recordEventInfof(r, "New revision detected %s.%s", - canary.GetName(), canary.Namespace) - canary.Spec.Replicas = int32p(1) - canary, err = c.kubeClient.AppsV1().Deployments(canary.Namespace).Update(canary) - if err != nil { - c.recordEventErrorf(r, "Scaling up %s.%s failed: %v", canary.GetName(), canary.Namespace, err) - return false - } - - r.Status = flaggerv1.CanaryStatus{ - FailedChecks: 0, - } - c.setCanaryRevision(r, canary, "running") - c.recordEventInfof(r, "Scaling up %s.%s", canary.GetName(), canary.Namespace) - - return false - } - - return false -} - -func (c *Controller) updateRolloutStatus(r *flaggerv1.Canary, status string) bool { - var err error - r.Status.State = status - r, err = c.rolloutClient.FlaggerV1alpha1().Canaries(r.Namespace).Update(r) - if err != nil { - c.logger.Errorf("Canary %s.%s status update failed: %v", r.Name, r.Namespace, err) - return false - } - return true -} - -func (c *Controller) updateRolloutFailedChecks(r *flaggerv1.Canary, val int) bool { - var err error - r.Status.FailedChecks = val - r, err = c.rolloutClient.FlaggerV1alpha1().Canaries(r.Namespace).Update(r) - if err != nil { - c.logger.Errorf("Canary %s.%s status update failed: %v", r.Name, r.Namespace, err) - return false - } - return true -} - -func (c *Controller) getDeployment(r *flaggerv1.Canary, name string, namespace string) (*appsv1.Deployment, bool) { - dep, err := c.kubeClient.AppsV1().Deployments(namespace).Get(name, v1.GetOptions{}) - if err != nil { - c.recordEventErrorf(r, "Deployment %s.%s not found", name, namespace) - return nil, false - } - - if msg, healthy := getDeploymentStatus(dep); !healthy { - c.recordEventWarningf(r, "Halt %s.%s advancement %s", dep.GetName(), dep.Namespace, msg) - return nil, false - } - - if dep.Spec.Replicas == nil || *dep.Spec.Replicas == 0 { - return nil, false - } - - return dep, true -} - -func (c *Controller) getCanaryDeployment(r *flaggerv1.Canary, name string, namespace string) (*appsv1.Deployment, bool) { - dep, err := c.kubeClient.AppsV1().Deployments(namespace).Get(name, v1.GetOptions{}) - if err != nil { - c.recordEventErrorf(r, "Deployment %s.%s not found", name, namespace) - return nil, false - } - - if msg, healthy := getDeploymentStatus(dep); !healthy { - c.recordEventWarningf(r, "Halt %s.%s advancement %s", dep.GetName(), dep.Namespace, msg) - return nil, false - } - - return dep, true -} - -func (c *Controller) checkDeploymentMetrics(r *flaggerv1.Canary) bool { - for _, metric := range r.Spec.CanaryAnalysis.Metrics { - if metric.Name == "istio_requests_total" { - val, err := c.getDeploymentCounter(r.Spec.TargetRef.Name, r.Namespace, metric.Name, metric.Interval) - if err != nil { - c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.metricsServer, err) - return false - } - if float64(metric.Threshold) > val { - c.recordEventWarningf(r, "Halt %s.%s advancement success rate %.2f%% < %v%%", - r.Name, r.Namespace, val, metric.Threshold) - return false - } - } - - if metric.Name == "istio_request_duration_seconds_bucket" { - val, err := c.GetDeploymentHistogram(r.Spec.TargetRef.Name, r.Namespace, metric.Name, metric.Interval) - if err != nil { - c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.metricsServer, err) - return false - } - t := time.Duration(metric.Threshold) * time.Millisecond - if val > t { - c.recordEventWarningf(r, "Halt %s.%s advancement request duration %v > %v", - r.Name, r.Namespace, val, t) - return false - } - } - } - - return true -} - -func (c *Controller) scaleToZeroCanary(r *flaggerv1.Canary) { - canary, err := c.kubeClient.AppsV1().Deployments(r.Namespace).Get(r.Spec.TargetRef.Name, v1.GetOptions{}) - if err != nil { - c.recordEventErrorf(r, "Deployment %s.%s not found", r.Spec.TargetRef.Name, r.Namespace) - return - } - //HPA https://github.com/kubernetes/kubernetes/pull/29212 - canary.Spec.Replicas = int32p(0) - canary, err = c.kubeClient.AppsV1().Deployments(canary.Namespace).Update(canary) - if err != nil { - c.recordEventErrorf(r, "Scaling down %s.%s failed: %v", canary.GetName(), canary.Namespace, err) - return - } -} - -func (c *Controller) setCanaryRevision(r *flaggerv1.Canary, canary *appsv1.Deployment, status string) { - r.Status = flaggerv1.CanaryStatus{ - State: status, - FailedChecks: r.Status.FailedChecks, - } - err := c.saveDeploymentSpec(r, canary) - if err != nil { - c.logger.Errorf("Canary %s.%s status update failed: %v", r.Name, r.Namespace, err) - } -} - -func (c *Controller) getVirtualService(r *flaggerv1.Canary) ( - vs *istiov1alpha3.VirtualService, - primary istiov1alpha3.DestinationWeight, - canary istiov1alpha3.DestinationWeight, - ok bool, -) { - var err error - vs, err = c.istioClient.NetworkingV1alpha3().VirtualServices(r.Namespace).Get(r.Name, v1.GetOptions{}) - if err != nil { - c.recordEventErrorf(r, "VirtualService %s.%s not found", r.Name, r.Namespace) - return - } - - for _, http := range vs.Spec.Http { - for _, route := range http.Route { - if route.Destination.Host == fmt.Sprintf("%s-primary", r.Spec.TargetRef.Name) { - primary = route - } - if route.Destination.Host == r.Spec.TargetRef.Name { - canary = route - } - } - } - - if primary.Weight == 0 && canary.Weight == 0 { - c.recordEventErrorf(r, "VirtualService %s.%s does not contain routes for %s and %s", - r.Name, r.Namespace, fmt.Sprintf("%s-primary", r.Spec.TargetRef.Name), r.Spec.TargetRef.Name) - return - } - - ok = true - return -} - -func (c *Controller) updateVirtualServiceRoutes( - r *flaggerv1.Canary, - vs *istiov1alpha3.VirtualService, - primary istiov1alpha3.DestinationWeight, - canary istiov1alpha3.DestinationWeight, -) bool { - vs.Spec.Http = []istiov1alpha3.HTTPRoute{ - { - Route: []istiov1alpha3.DestinationWeight{primary, canary}, - }, - } - - var err error - vs, err = c.istioClient.NetworkingV1alpha3().VirtualServices(r.Namespace).Update(vs) - if err != nil { - c.recordEventErrorf(r, "VirtualService %s.%s update failed: %v", r.Name, r.Namespace, err) - return false - } - return true -} - -func getDeploymentStatus(deployment *appsv1.Deployment) (string, bool) { - if deployment.Generation <= deployment.Status.ObservedGeneration { - cond := getDeploymentCondition(deployment.Status, appsv1.DeploymentProgressing) - if cond != nil && cond.Reason == "ProgressDeadlineExceeded" { - return fmt.Sprintf("deployment %q exceeded its progress deadline", deployment.GetName()), false - } else if deployment.Spec.Replicas != nil && deployment.Status.UpdatedReplicas < *deployment.Spec.Replicas { - return fmt.Sprintf("waiting for rollout to finish: %d out of %d new replicas have been updated", - deployment.Status.UpdatedReplicas, *deployment.Spec.Replicas), false - } else if deployment.Status.Replicas > deployment.Status.UpdatedReplicas { - return fmt.Sprintf("waiting for rollout to finish: %d old replicas are pending termination", - deployment.Status.Replicas-deployment.Status.UpdatedReplicas), false - } else if deployment.Status.AvailableReplicas < deployment.Status.UpdatedReplicas { - return fmt.Sprintf("waiting for rollout to finish: %d of %d updated replicas are available", - deployment.Status.AvailableReplicas, deployment.Status.UpdatedReplicas), false - } - } else { - return "waiting for rollout to finish: observed deployment generation less then desired generation", false - } - - return "ready", true -} - -func getDeploymentCondition( - status appsv1.DeploymentStatus, - conditionType appsv1.DeploymentConditionType, -) *appsv1.DeploymentCondition { - for i := range status.Conditions { - c := status.Conditions[i] - if c.Type == conditionType { - return &c - } - } - return nil -} - -func int32p(i int32) *int32 { - return &i -} diff --git a/pkg/controller/observer.go b/pkg/controller/observer.go index 98e598eb..76f1d012 100644 --- a/pkg/controller/observer.go +++ b/pkg/controller/observer.go @@ -11,6 +11,10 @@ import ( "time" ) +type CanaryObserver struct { + metricsServer string +} + type VectorQueryResponse struct { Data struct { Result []struct { @@ -23,7 +27,7 @@ type VectorQueryResponse struct { } } -func (c *Controller) queryMetric(query string) (*VectorQueryResponse, error) { +func (c *CanaryObserver) queryMetric(query string) (*VectorQueryResponse, error) { promURL, err := url.Parse(c.metricsServer) if err != nil { return nil, err @@ -69,7 +73,7 @@ func (c *Controller) queryMetric(query string) (*VectorQueryResponse, error) { } // istio_requests_total -func (c *Controller) getDeploymentCounter(name string, namespace string, metric string, interval string) (float64, error) { +func (c *CanaryObserver) GetDeploymentCounter(name string, namespace string, metric string, interval string) (float64, error) { var rate *float64 querySt := url.QueryEscape(`sum(rate(` + metric + `{reporter="destination",destination_workload_namespace=~"` + @@ -102,7 +106,7 @@ func (c *Controller) getDeploymentCounter(name string, namespace string, metric } // istio_request_duration_seconds_bucket -func (c *Controller) GetDeploymentHistogram(name string, namespace string, metric string, interval string) (time.Duration, error) { +func (c *CanaryObserver) GetDeploymentHistogram(name string, namespace string, metric string, interval string) (time.Duration, error) { var rate *float64 querySt := url.QueryEscape(`histogram_quantile(0.99, sum(rate(` + metric + `{reporter="destination",destination_workload=~"` + diff --git a/pkg/controller/router.go b/pkg/controller/router.go new file mode 100644 index 00000000..609a5eba --- /dev/null +++ b/pkg/controller/router.go @@ -0,0 +1,285 @@ +package controller + +import ( + "fmt" + + istiov1alpha3 "github.com/knative/pkg/apis/istio/v1alpha3" + istioclientset "github.com/knative/pkg/client/clientset/versioned" + flaggerv1 "github.com/stefanprodan/flagger/pkg/apis/flagger/v1alpha1" + clientset "github.com/stefanprodan/flagger/pkg/client/clientset/versioned" + "go.uber.org/zap" + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/apis/meta/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/util/intstr" + "k8s.io/client-go/kubernetes" +) + +type CanaryRouter struct { + kubeClient kubernetes.Interface + istioClient istioclientset.Interface + flaggerClient clientset.Interface + logger *zap.SugaredLogger +} + +// Sync creates the primary and canary ClusterIP services +// and sets up a virtual service with routes for the two services +// all traffic goes to primary +func (c *CanaryRouter) Sync(cd *flaggerv1.Canary) error { + err := c.createServices(cd) + if err != nil { + return err + } + err = c.createVirtualService(cd) + if err != nil { + return err + } + return nil +} + +func (c *CanaryRouter) createServices(cd *flaggerv1.Canary) error { + canaryName := cd.Spec.TargetRef.Name + primaryName := fmt.Sprintf("%s-primary", canaryName) + canaryService, err := c.kubeClient.CoreV1().Services(cd.Namespace).Get(canaryName, metav1.GetOptions{}) + if errors.IsNotFound(err) { + canaryService = &corev1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: canaryName, + Namespace: cd.Namespace, + OwnerReferences: []metav1.OwnerReference{ + *metav1.NewControllerRef(cd, schema.GroupVersionKind{ + Group: flaggerv1.SchemeGroupVersion.Group, + Version: flaggerv1.SchemeGroupVersion.Version, + Kind: flaggerv1.CanaryKind, + }), + }, + }, + Spec: corev1.ServiceSpec{ + Type: corev1.ServiceTypeClusterIP, + Selector: map[string]string{"app": canaryName}, + Ports: []corev1.ServicePort{ + { + Name: "http", + Protocol: corev1.ProtocolTCP, + Port: cd.Spec.Service.Port, + TargetPort: intstr.IntOrString{ + Type: intstr.Int, + IntVal: cd.Spec.Service.Port, + }, + }, + }, + }, + } + + _, err = c.kubeClient.CoreV1().Services(cd.Namespace).Create(canaryService) + if err != nil { + return err + } + c.logger.Infof("Service %s.%s created", canaryService.GetName(), cd.Namespace) + } + + canaryTestServiceName := fmt.Sprintf("%s-canary", cd.Spec.TargetRef.Name) + canaryTestService, err := c.kubeClient.CoreV1().Services(cd.Namespace).Get(canaryTestServiceName, metav1.GetOptions{}) + if errors.IsNotFound(err) { + canaryTestService = &corev1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: canaryTestServiceName, + Namespace: cd.Namespace, + OwnerReferences: []metav1.OwnerReference{ + *metav1.NewControllerRef(cd, schema.GroupVersionKind{ + Group: flaggerv1.SchemeGroupVersion.Group, + Version: flaggerv1.SchemeGroupVersion.Version, + Kind: flaggerv1.CanaryKind, + }), + }, + }, + Spec: corev1.ServiceSpec{ + Type: corev1.ServiceTypeClusterIP, + Selector: map[string]string{"app": canaryName}, + Ports: []corev1.ServicePort{ + { + Name: "http", + Protocol: corev1.ProtocolTCP, + Port: cd.Spec.Service.Port, + TargetPort: intstr.IntOrString{ + Type: intstr.Int, + IntVal: cd.Spec.Service.Port, + }, + }, + }, + }, + } + + _, err = c.kubeClient.CoreV1().Services(cd.Namespace).Create(canaryTestService) + if err != nil { + return err + } + c.logger.Infof("Service %s.%s created", canaryTestService.GetName(), cd.Namespace) + } + + primaryService, err := c.kubeClient.CoreV1().Services(cd.Namespace).Get(primaryName, metav1.GetOptions{}) + if errors.IsNotFound(err) { + primaryService = &corev1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: primaryName, + Namespace: cd.Namespace, + OwnerReferences: []metav1.OwnerReference{ + *metav1.NewControllerRef(cd, schema.GroupVersionKind{ + Group: flaggerv1.SchemeGroupVersion.Group, + Version: flaggerv1.SchemeGroupVersion.Version, + Kind: flaggerv1.CanaryKind, + }), + }, + }, + Spec: corev1.ServiceSpec{ + Type: corev1.ServiceTypeClusterIP, + Selector: map[string]string{"app": primaryName}, + Ports: []corev1.ServicePort{ + { + Name: "http", + Protocol: corev1.ProtocolTCP, + Port: cd.Spec.Service.Port, + TargetPort: intstr.IntOrString{ + Type: intstr.Int, + IntVal: cd.Spec.Service.Port, + }, + }, + }, + }, + } + + _, err = c.kubeClient.CoreV1().Services(cd.Namespace).Create(primaryService) + if err != nil { + return err + } + + c.logger.Infof("Service %s.%s created", primaryService.GetName(), cd.Namespace) + } + + return nil +} + +func (c *CanaryRouter) createVirtualService(cd *flaggerv1.Canary) error { + canaryName := cd.Spec.TargetRef.Name + primaryName := fmt.Sprintf("%s-primary", canaryName) + hosts := append(cd.Spec.Service.Hosts, canaryName) + gateways := append(cd.Spec.Service.Gateways, "mesh") + virtualService, err := c.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Get(canaryName, metav1.GetOptions{}) + if errors.IsNotFound(err) { + virtualService = &istiov1alpha3.VirtualService{ + ObjectMeta: metav1.ObjectMeta{ + Name: cd.Name, + Namespace: cd.Namespace, + OwnerReferences: []metav1.OwnerReference{ + *metav1.NewControllerRef(cd, schema.GroupVersionKind{ + Group: flaggerv1.SchemeGroupVersion.Group, + Version: flaggerv1.SchemeGroupVersion.Version, + Kind: flaggerv1.CanaryKind, + }), + }, + }, + Spec: istiov1alpha3.VirtualServiceSpec{ + Hosts: hosts, + Gateways: gateways, + Http: []istiov1alpha3.HTTPRoute{ + { + Route: []istiov1alpha3.DestinationWeight{ + { + Destination: istiov1alpha3.Destination{ + Host: primaryName, + Port: istiov1alpha3.PortSelector{ + Number: uint32(cd.Spec.Service.Port), + }, + }, + Weight: 100, + }, + { + Destination: istiov1alpha3.Destination{ + Host: canaryName, + Port: istiov1alpha3.PortSelector{ + Number: uint32(cd.Spec.Service.Port), + }, + }, + Weight: 0, + }, + }, + }, + }, + }, + } + + _, err = c.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Create(virtualService) + if err != nil { + return fmt.Errorf("VirtualService %s.%s create error %v", cd.Name, cd.Namespace, err) + } + c.logger.Infof("VirtualService %s.%s created", virtualService.GetName(), cd.Namespace) + } + + return nil +} + +// GetRoutes returns the destinations weight for primary and canary +func (c *CanaryRouter) GetRoutes(cd *flaggerv1.Canary) ( + primary istiov1alpha3.DestinationWeight, + canary istiov1alpha3.DestinationWeight, + err error, +) { + vs := &istiov1alpha3.VirtualService{} + vs, err = c.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Get(cd.Name, v1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + err = fmt.Errorf("VirtualService %s.%s not found", cd.Name, cd.Namespace) + return + } + err = fmt.Errorf("VirtualService %s.%s query error %v", cd.Name, cd.Namespace, err) + return + } + + for _, http := range vs.Spec.Http { + for _, route := range http.Route { + if route.Destination.Host == fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) { + primary = route + } + if route.Destination.Host == cd.Spec.TargetRef.Name { + canary = route + } + } + } + + if primary.Weight == 0 && canary.Weight == 0 { + err = fmt.Errorf("VirtualService %s.%s does not contain routes for %s and %s", + cd.Name, cd.Namespace, fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name), cd.Spec.TargetRef.Name) + } + + return +} + +// SetRoutes updates the destinations weight for primary and canary +func (c *CanaryRouter) SetRoutes( + cd *flaggerv1.Canary, + primary istiov1alpha3.DestinationWeight, + canary istiov1alpha3.DestinationWeight, +) error { + vs, err := c.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Get(cd.Name, v1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return fmt.Errorf("VirtualService %s.%s not found", cd.Name, cd.Namespace) + + } + return fmt.Errorf("VirtualService %s.%s query error %v", cd.Name, cd.Namespace, err) + } + vs.Spec.Http = []istiov1alpha3.HTTPRoute{ + { + Route: []istiov1alpha3.DestinationWeight{primary, canary}, + }, + } + + vs, err = c.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Update(vs) + if err != nil { + return fmt.Errorf("VirtualService %s.%s update failed: %v", cd.Name, cd.Namespace, err) + + } + return nil +} diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go new file mode 100644 index 00000000..afd128d1 --- /dev/null +++ b/pkg/controller/scheduler.go @@ -0,0 +1,266 @@ +package controller + +import ( + "fmt" + "time" + + flaggerv1 "github.com/stefanprodan/flagger/pkg/apis/flagger/v1alpha1" + "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +func (c *Controller) scheduleCanaries() { + c.rollouts.Range(func(key interface{}, value interface{}) bool { + r := value.(*flaggerv1.Canary) + if r.Spec.TargetRef.Kind == "Deployment" { + go c.advanceCanary(r.Name, r.Namespace) + } + return true + }) +} + +func (c *Controller) advanceCanary(name string, namespace string) { + // check if the rollout exists + r, err := c.rolloutClient.FlaggerV1alpha1().Canaries(namespace).Get(name, v1.GetOptions{}) + if err != nil { + c.logger.Errorf("Canary %s.%s not found", name, namespace) + return + } + + deployer := CanaryDeployer{ + logger: c.logger, + kubeClient: c.kubeClient, + istioClient: c.istioClient, + flaggerClient: c.rolloutClient, + } + + // create primary deployment and hpa if needed + err = deployer.Sync(r) + if err != nil { + c.recordEventWarningf(r, "%v", err) + return + } + + router := CanaryRouter{ + logger: c.logger, + kubeClient: c.kubeClient, + istioClient: c.istioClient, + flaggerClient: c.rolloutClient, + } + + // create ClusterIP services and virtual service if needed + err = router.Sync(r) + if err != nil { + c.recordEventWarningf(r, "%v", err) + return + } + + // set max weight default value to 100% + maxWeight := 100 + if r.Spec.CanaryAnalysis.MaxWeight > 0 { + maxWeight = r.Spec.CanaryAnalysis.MaxWeight + } + + // check primary and canary deployments status + err = deployer.IsDeploymentHealthy(r) + if err != nil { + c.recordEventWarningf(r, "%v", err) + return + } + + // check if virtual service exists + // and if it contains weighted destination routes to the primary and canary services + primaryRoute, canaryRoute, err := router.GetRoutes(r) + if err != nil { + c.recordEventWarningf(r, "%v", err) + return + } + + // check if canary analysis should start (canary revision has changes) or continue + if ok := c.checkCanaryStatus(r, deployer); !ok { + return + } + + // check if the number of failed checks reached the threshold + if r.Status.State == "running" && r.Status.FailedChecks >= r.Spec.CanaryAnalysis.Threshold { + c.recordEventWarningf(r, "Rolling back %s.%s failed checks threshold reached %v", + r.Name, r.Namespace, r.Status.FailedChecks) + + // route all traffic back to primary + primaryRoute.Weight = 100 + canaryRoute.Weight = 0 + if err := router.SetRoutes(r, primaryRoute, canaryRoute); err != nil { + c.recordEventWarningf(r, "%v", err) + return + } + + c.recordEventWarningf(r, "Canary failed! Scaling down %s.%s", + r.Spec.TargetRef.Name, r.Namespace) + + // shutdown canary + err = deployer.Scale(r, 0) + if err != nil { + c.recordEventWarningf(r, "%v", err) + return + } + + // mark canary as failed + err := deployer.SetState(r, "failed") + if err != nil { + c.logger.Errorf("%v", err) + return + } + return + } + + // check if the canary success rate is above the threshold + // skip check if no traffic is routed to canary + if canaryRoute.Weight == 0 { + c.recordEventInfof(r, "Starting canary deployment for %s.%s", r.Name, r.Namespace) + } else { + if ok := c.checkCanaryMetrics(r); !ok { + if err = deployer.SetFailedChecks(r, r.Status.FailedChecks+1); err != nil { + c.recordEventWarningf(r, "%v", err) + return + } + return + } + } + + // increase canary traffic percentage + if canaryRoute.Weight < maxWeight { + primaryRoute.Weight -= r.Spec.CanaryAnalysis.StepWeight + if primaryRoute.Weight < 0 { + primaryRoute.Weight = 0 + } + canaryRoute.Weight += r.Spec.CanaryAnalysis.StepWeight + if primaryRoute.Weight > 100 { + primaryRoute.Weight = 100 + } + + if err = router.SetRoutes(r, primaryRoute, canaryRoute); err != nil { + c.recordEventWarningf(r, "%v", err) + return + } + + c.recordEventInfof(r, "Advance %s.%s canary weight %v", r.Name, r.Namespace, canaryRoute.Weight) + + // promote canary + primaryName := fmt.Sprintf("%s-primary", r.Spec.TargetRef.Name) + if canaryRoute.Weight == maxWeight { + c.recordEventInfof(r, "Copying %s.%s template spec to %s.%s", + r.Spec.TargetRef.Name, r.Namespace, primaryName, r.Namespace) + + err := deployer.Promote(r) + if err != nil { + c.recordEventWarningf(r, "%v", err) + return + } + } + } else { + // route all traffic back to primary + primaryRoute.Weight = 100 + canaryRoute.Weight = 0 + if err = router.SetRoutes(r, primaryRoute, canaryRoute); err != nil { + c.recordEventWarningf(r, "%v", err) + return + } + + c.recordEventInfof(r, "Promotion completed! Scaling down %s.%s", r.Spec.TargetRef.Name, r.Namespace) + + // shutdown canary + err = deployer.Scale(r, 0) + if err != nil { + c.recordEventWarningf(r, "%v", err) + return + } + + // update status + err = deployer.SetState(r, "finished") + if err != nil { + c.recordEventWarningf(r, "%v", err) + return + } + } +} + +func (c *Controller) checkCanaryStatus(r *flaggerv1.Canary, deployer CanaryDeployer) bool { + if r.Status.State == "running" { + return true + } + + if r.Status.State == "" { + status := flaggerv1.CanaryStatus{ + State: "initialized", + FailedChecks: 0, + } + + err := deployer.SyncStatus(r, status) + if err != nil { + c.logger.Errorf("%v", err) + return false + } + + c.recordEventInfof(r, "Initialization done! %s.%s", r.Name, r.Namespace) + return false + } + + if diff, err := deployer.IsNewSpec(r); diff { + c.recordEventInfof(r, "New revision detected %s.%s", r.Spec.TargetRef.Name, r.Namespace) + err = deployer.Scale(r, 1) + if err != nil { + c.recordEventErrorf(r, "%v", err) + return false + } + + status := flaggerv1.CanaryStatus{ + State: "running", + FailedChecks: 0, + } + err := deployer.SyncStatus(r, status) + if err != nil { + c.logger.Errorf("%v", err) + return false + } + c.recordEventInfof(r, "Scaling up %s.%s", r.Spec.TargetRef.Name, r.Namespace) + + return false + } + + return false +} + +func (c *Controller) checkCanaryMetrics(r *flaggerv1.Canary) bool { + observer := &CanaryObserver{ + metricsServer: c.metricsServer, + } + for _, metric := range r.Spec.CanaryAnalysis.Metrics { + if metric.Name == "istio_requests_total" { + val, err := observer.GetDeploymentCounter(r.Spec.TargetRef.Name, r.Namespace, metric.Name, metric.Interval) + if err != nil { + c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.metricsServer, err) + return false + } + if float64(metric.Threshold) > val { + c.recordEventWarningf(r, "Halt %s.%s advancement success rate %.2f%% < %v%%", + r.Name, r.Namespace, val, metric.Threshold) + return false + } + } + + if metric.Name == "istio_request_duration_seconds_bucket" { + val, err := observer.GetDeploymentHistogram(r.Spec.TargetRef.Name, r.Namespace, metric.Name, metric.Interval) + if err != nil { + c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.metricsServer, err) + return false + } + t := time.Duration(metric.Threshold) * time.Millisecond + if val > t { + c.recordEventWarningf(r, "Halt %s.%s advancement request duration %v > %v", + r.Name, r.Namespace, val, t) + return false + } + } + } + + return true +} diff --git a/pkg/controller/utils.go b/pkg/controller/utils.go deleted file mode 100644 index 735696be..00000000 --- a/pkg/controller/utils.go +++ /dev/null @@ -1,79 +0,0 @@ -package controller - -import ( - "encoding/base64" - "encoding/json" - "fmt" - - "github.com/google/go-cmp/cmp" - "github.com/google/go-cmp/cmp/cmpopts" - flaggerv1 "github.com/stefanprodan/flagger/pkg/apis/flagger/v1alpha1" - appsv1 "k8s.io/api/apps/v1" - corev1 "k8s.io/api/core/v1" - "k8s.io/apimachinery/pkg/api/resource" - "k8s.io/apimachinery/pkg/apis/meta/v1" -) - -func (c *Controller) saveDeploymentSpec(cd *flaggerv1.Canary, dep *appsv1.Deployment) error { - specJson, err := json.Marshal(dep.Spec.Template.Spec) - if err != nil { - return err - } - - specEnc := base64.StdEncoding.EncodeToString(specJson) - cd.Status.CanaryRevision = specEnc - cd, err = c.rolloutClient.FlaggerV1alpha1().Canaries(cd.Namespace).Update(cd) - if err != nil { - return err - } - return nil -} - -func (c *Controller) diffDeploymentSpec(cd *flaggerv1.Canary, dep *appsv1.Deployment) (bool, error) { - if cd.Status.CanaryRevision == "" { - return true, nil - } - - newSpec := &dep.Spec.Template.Spec - oldSpecJson, err := base64.StdEncoding.DecodeString(cd.Status.CanaryRevision) - if err != nil { - return false, err - } - oldSpec := &corev1.PodSpec{} - err = json.Unmarshal(oldSpecJson, oldSpec) - if err != nil { - return false, err - } - - if diff := cmp.Diff(*newSpec, *oldSpec, cmpopts.IgnoreUnexported(resource.Quantity{})); diff != "" { - fmt.Println(diff) - return true, nil - } - - return false, nil -} - -func (c *Controller) getDeploymentSpec(name string, namespace string) (string, error) { - dep, err := c.kubeClient.AppsV1().Deployments(namespace).Get(name, v1.GetOptions{}) - if err != nil { - return "", err - } - - specJson, err := json.Marshal(dep.Spec.Template.Spec) - if err != nil { - return "", err - } - - specEnc := base64.StdEncoding.EncodeToString(specJson) - return specEnc, nil -} - -func (c *Controller) getDeploymentSpecEnc(dep *appsv1.Deployment) (string, error) { - specJson, err := json.Marshal(dep.Spec.Template.Spec) - if err != nil { - return "", err - } - - specEnc := base64.StdEncoding.EncodeToString(specJson) - return specEnc, nil -}