diff --git a/pkg/canary/deployment_controller.go b/pkg/canary/deployment_controller.go index 62abd669..0d3fdd75 100644 --- a/pkg/canary/deployment_controller.go +++ b/pkg/canary/deployment_controller.go @@ -48,7 +48,6 @@ type DeploymentController struct { // Initialize creates the primary deployment, hpa, // scales to zero the canary deployment and returns the pod selector label and container ports func (c *DeploymentController) Initialize(cd *flaggerv1.Canary) (err error) { - primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) if err := c.createPrimaryDeployment(cd, c.includeLabelPrefix); err != nil { return fmt.Errorf("createPrimaryDeployment failed: %w", err) } @@ -67,16 +66,6 @@ func (c *DeploymentController) Initialize(cd *flaggerv1.Canary) (err error) { } } - if cd.Spec.AutoscalerRef != nil { - if cd.Spec.AutoscalerRef.Kind == "HorizontalPodAutoscaler" { - if err := c.reconcilePrimaryHpa(cd, true); err != nil { - return fmt.Errorf( - "initial reconcilePrimaryHpa for %s.%s failed: %w", primaryName, cd.Namespace, err) - } - } else { - return fmt.Errorf("cd.Spec.AutoscalerRef.Kind is invalid: %s", cd.Spec.AutoscalerRef.Kind) - } - } return nil } @@ -155,17 +144,6 @@ func (c *DeploymentController) Promote(cd *flaggerv1.Canary) error { primaryName, cd.Namespace, err) } - // update HPA - if cd.Spec.AutoscalerRef != nil { - if cd.Spec.AutoscalerRef.Kind == "HorizontalPodAutoscaler" { - if err := c.reconcilePrimaryHpa(cd, false); err != nil { - return fmt.Errorf( - "reconcilePrimaryHpa for %s.%s failed: %w", primaryName, cd.Namespace, err) - } - } else { - return fmt.Errorf("cd.Spec.AutoscalerRef.Kind is invalid: %s", cd.Spec.AutoscalerRef.Kind) - } - } return nil } diff --git a/pkg/canary/factory.go b/pkg/canary/factory.go index 49265b98..b0ba97c6 100644 --- a/pkg/canary/factory.go +++ b/pkg/canary/factory.go @@ -83,3 +83,19 @@ func (factory *Factory) Controller(kind string) Controller { return deploymentCtrl } } + +func (factory *Factory) ScalerReconciler(kind string) ScalerReconciler { + hpaReconciler := &HPAReconciler{ + logger: factory.logger, + kubeClient: factory.kubeClient, + flaggerClient: factory.flaggerClient, + includeLabelPrefix: factory.includeLabelPrefix, + } + + switch kind { + case "HorizontalPodAutoscaler": + return hpaReconciler + default: + return nil + } +} diff --git a/pkg/canary/hpa_reconciler.go b/pkg/canary/hpa_reconciler.go new file mode 100644 index 00000000..e4a2e3e2 --- /dev/null +++ b/pkg/canary/hpa_reconciler.go @@ -0,0 +1,258 @@ +package canary + +import ( + "context" + "fmt" + + flaggerv1 "github.com/fluxcd/flagger/pkg/apis/flagger/v1beta1" + clientset "github.com/fluxcd/flagger/pkg/client/clientset/versioned" + "github.com/google/go-cmp/cmp" + "go.uber.org/zap" + hpav2 "k8s.io/api/autoscaling/v2" + hpav2beta2 "k8s.io/api/autoscaling/v2beta2" + "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/util/retry" +) + +// HPAReconciler is a ScalerReconciler that reconciles HPAs. +type HPAReconciler struct { + kubeClient kubernetes.Interface + flaggerClient clientset.Interface + logger *zap.SugaredLogger + includeLabelPrefix []string +} + +func (hr *HPAReconciler) ReconcilePrimaryScaler(cd *flaggerv1.Canary, init bool) error { + if cd.Spec.AutoscalerRef != nil { + if err := hr.reconcilePrimaryHpa(cd, init); err != nil { + return err + } + } + return nil +} + +func (hr *HPAReconciler) reconcilePrimaryHpa(cd *flaggerv1.Canary, init bool) error { + var betaHpa *hpav2beta2.HorizontalPodAutoscaler + hpa, err := hr.kubeClient.AutoscalingV2().HorizontalPodAutoscalers(cd.Namespace).Get(context.TODO(), cd.Spec.AutoscalerRef.Name, metav1.GetOptions{}) + if err != nil { + hr.logger.Debugf("v2 HorizontalPodAutoscaler %s.%s get query error: %w; falling back to v2beta2", + cd.Namespace, cd.Spec.AutoscalerRef.Name, err) + var betaErr error + betaHpa, betaErr = hr.kubeClient.AutoscalingV2beta2().HorizontalPodAutoscalers(cd.Namespace).Get(context.TODO(), cd.Spec.AutoscalerRef.Name, metav1.GetOptions{}) + if betaErr != nil { + return fmt.Errorf("HorizontalPodAutoscaler %s.%s get query error for both v2beta2: %s and v2: %s", + cd.Spec.AutoscalerRef.Name, cd.Namespace, betaErr, err) + } + } + + if hpa != nil { + if err = hr.reconcilePrimaryHpaV2(cd, hpa, init); err != nil { + return err + } + } else if betaHpa != nil { + if err = hr.reconcilePrimaryHpaV2Beta2(cd, betaHpa, init); err != nil { + return err + } + } + + return nil +} + +func (hr *HPAReconciler) reconcilePrimaryHpaV2(cd *flaggerv1.Canary, hpa *hpav2.HorizontalPodAutoscaler, init bool) error { + primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) + + hpaSpec := hpav2.HorizontalPodAutoscalerSpec{ + ScaleTargetRef: hpav2.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, + Behavior: hpa.Spec.Behavior, + } + + primaryHpaName := fmt.Sprintf("%s-primary", cd.Spec.AutoscalerRef.Name) + primaryHpa, err := hr.kubeClient.AutoscalingV2().HorizontalPodAutoscalers(cd.Namespace).Get(context.TODO(), primaryHpaName, metav1.GetOptions{}) + + // create HPA + if errors.IsNotFound(err) { + primaryHpa = &hpav2.HorizontalPodAutoscaler{ + ObjectMeta: metav1.ObjectMeta{ + Name: primaryHpaName, + Namespace: cd.Namespace, + Labels: filterMetadata(hpa.Labels), + OwnerReferences: []metav1.OwnerReference{ + *metav1.NewControllerRef(cd, schema.GroupVersionKind{ + Group: flaggerv1.SchemeGroupVersion.Group, + Version: flaggerv1.SchemeGroupVersion.Version, + Kind: flaggerv1.CanaryKind, + }), + }, + }, + Spec: hpaSpec, + } + + _, err = hr.kubeClient.AutoscalingV2().HorizontalPodAutoscalers(cd.Namespace).Create(context.TODO(), primaryHpa, metav1.CreateOptions{}) + if err != nil { + return fmt.Errorf("creating HorizontalPodAutoscaler %s.%s failed: %w", + primaryHpa.Name, primaryHpa.Namespace, err) + } + hr.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof( + "HorizontalPodAutoscaler %s.%s created", primaryHpa.GetName(), cd.Namespace) + return nil + } else if err != nil { + return fmt.Errorf("HorizontalPodAutoscaler %s.%s get query failed: %w", + primaryHpa.Name, primaryHpa.Namespace, err) + } + + // update HPA + if !init && primaryHpa != nil { + diffMetrics := cmp.Diff(hpaSpec.Metrics, primaryHpa.Spec.Metrics) + diffBehavior := cmp.Diff(hpaSpec.Behavior, primaryHpa.Spec.Behavior) + diffLabels := cmp.Diff(hpa.ObjectMeta.Labels, primaryHpa.ObjectMeta.Labels) + diffAnnotations := cmp.Diff(hpa.ObjectMeta.Annotations, primaryHpa.ObjectMeta.Annotations) + if diffMetrics != "" || diffBehavior != "" || diffLabels != "" || diffAnnotations != "" || int32Default(hpaSpec.MinReplicas) != int32Default(primaryHpa.Spec.MinReplicas) || hpaSpec.MaxReplicas != primaryHpa.Spec.MaxReplicas { + err = retry.RetryOnConflict(retry.DefaultRetry, func() error { + primaryHpa, err := hr.kubeClient.AutoscalingV2().HorizontalPodAutoscalers(cd.Namespace).Get(context.TODO(), primaryHpaName, metav1.GetOptions{}) + if err != nil { + return err + } + hpaClone := primaryHpa.DeepCopy() + hpaClone.Spec.MaxReplicas = hpaSpec.MaxReplicas + hpaClone.Spec.MinReplicas = hpaSpec.MinReplicas + hpaClone.Spec.Metrics = hpaSpec.Metrics + hpaClone.Spec.Behavior = hpaSpec.Behavior + + // update hpa annotations + hpaClone.ObjectMeta.Annotations = make(map[string]string) + filteredAnnotations := includeLabelsByPrefix(hpa.ObjectMeta.Annotations, hr.includeLabelPrefix) + for k, v := range filteredAnnotations { + hpaClone.ObjectMeta.Annotations[k] = v + } + // update hpa labels + hpaClone.ObjectMeta.Labels = make(map[string]string) + filteredLabels := includeLabelsByPrefix(hpa.ObjectMeta.Labels, hr.includeLabelPrefix) + for k, v := range filteredLabels { + hpaClone.ObjectMeta.Labels[k] = v + } + + _, err = hr.kubeClient.AutoscalingV2().HorizontalPodAutoscalers(cd.Namespace).Update(context.TODO(), hpaClone, metav1.UpdateOptions{}) + return err + }) + if err != nil { + return fmt.Errorf("updating HorizontalPodAutoscaler %s.%s failed: %w", + primaryHpa.Name, primaryHpa.Namespace, err) + } + hr.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)). + Infof("HorizontalPodAutoscaler %s.%s updated", primaryHpa.GetName(), cd.Namespace) + } + } + return nil +} + +func (hr *HPAReconciler) reconcilePrimaryHpaV2Beta2(cd *flaggerv1.Canary, hpa *hpav2beta2.HorizontalPodAutoscaler, init bool) error { + primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) + + hpaSpec := hpav2beta2.HorizontalPodAutoscalerSpec{ + ScaleTargetRef: hpav2beta2.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, + Behavior: hpa.Spec.Behavior, + } + + primaryHpaName := fmt.Sprintf("%s-primary", cd.Spec.AutoscalerRef.Name) + primaryHpa, err := hr.kubeClient.AutoscalingV2beta2().HorizontalPodAutoscalers(cd.Namespace).Get(context.TODO(), primaryHpaName, metav1.GetOptions{}) + + // create HPA + if errors.IsNotFound(err) { + primaryHpa = &hpav2beta2.HorizontalPodAutoscaler{ + ObjectMeta: metav1.ObjectMeta{ + Name: primaryHpaName, + Namespace: cd.Namespace, + Labels: filterMetadata(hpa.Labels), + OwnerReferences: []metav1.OwnerReference{ + *metav1.NewControllerRef(cd, schema.GroupVersionKind{ + Group: flaggerv1.SchemeGroupVersion.Group, + Version: flaggerv1.SchemeGroupVersion.Version, + Kind: flaggerv1.CanaryKind, + }), + }, + }, + Spec: hpaSpec, + } + + _, err = hr.kubeClient.AutoscalingV2beta2().HorizontalPodAutoscalers(cd.Namespace).Create(context.TODO(), primaryHpa, metav1.CreateOptions{}) + if err != nil { + return fmt.Errorf("creating HorizontalPodAutoscaler %s.%s failed: %w", + primaryHpa.Name, primaryHpa.Namespace, err) + } + hr.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof( + "HorizontalPodAutoscaler %s.%s created", primaryHpa.GetName(), cd.Namespace) + return nil + } else if err != nil { + return fmt.Errorf("HorizontalPodAutoscaler %s.%s get query failed: %w", + primaryHpa.Name, primaryHpa.Namespace, err) + } + + // update HPA + if !init && primaryHpa != nil { + diffMetrics := cmp.Diff(hpaSpec.Metrics, primaryHpa.Spec.Metrics) + diffBehavior := cmp.Diff(hpaSpec.Behavior, primaryHpa.Spec.Behavior) + diffLabels := cmp.Diff(hpa.ObjectMeta.Labels, primaryHpa.ObjectMeta.Labels) + diffAnnotations := cmp.Diff(hpa.ObjectMeta.Annotations, primaryHpa.ObjectMeta.Annotations) + if diffMetrics != "" || diffBehavior != "" || diffLabels != "" || diffAnnotations != "" || int32Default(hpaSpec.MinReplicas) != int32Default(primaryHpa.Spec.MinReplicas) || hpaSpec.MaxReplicas != primaryHpa.Spec.MaxReplicas { + err = retry.RetryOnConflict(retry.DefaultRetry, func() error { + primaryHpa, err := hr.kubeClient.AutoscalingV2beta2().HorizontalPodAutoscalers(cd.Namespace).Get(context.TODO(), primaryHpaName, metav1.GetOptions{}) + if err != nil { + return err + } + hpaClone := primaryHpa.DeepCopy() + hpaClone.Spec.MaxReplicas = hpaSpec.MaxReplicas + hpaClone.Spec.MinReplicas = hpaSpec.MinReplicas + hpaClone.Spec.Metrics = hpaSpec.Metrics + hpaClone.Spec.Behavior = hpaSpec.Behavior + + // update hpa annotations + hpaClone.ObjectMeta.Annotations = make(map[string]string) + filteredAnnotations := includeLabelsByPrefix(hpa.ObjectMeta.Annotations, hr.includeLabelPrefix) + for k, v := range filteredAnnotations { + hpaClone.ObjectMeta.Annotations[k] = v + } + // update hpa labels + hpaClone.ObjectMeta.Labels = make(map[string]string) + filteredLabels := includeLabelsByPrefix(hpa.ObjectMeta.Labels, hr.includeLabelPrefix) + for k, v := range filteredLabels { + hpaClone.ObjectMeta.Labels[k] = v + } + + _, err = hr.kubeClient.AutoscalingV2beta2().HorizontalPodAutoscalers(cd.Namespace).Update(context.TODO(), hpaClone, metav1.UpdateOptions{}) + return err + }) + if err != nil { + return fmt.Errorf("updating HorizontalPodAutoscaler %s.%s failed: %w", + primaryHpa.Name, primaryHpa.Namespace, err) + } + hr.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)). + Infof("HorizontalPodAutoscaler %s.%s updated", primaryHpa.GetName(), cd.Namespace) + } + } + return nil +} + +func (hr *HPAReconciler) PauseTargetScaler(cd *flaggerv1.Canary) error { + return nil +} + +func (hr *HPAReconciler) ResumeTargetScaler(cd *flaggerv1.Canary) error { + return nil +} diff --git a/pkg/canary/scaler_reconciler.go b/pkg/canary/scaler_reconciler.go new file mode 100644 index 00000000..a768bab8 --- /dev/null +++ b/pkg/canary/scaler_reconciler.go @@ -0,0 +1,13 @@ +package canary + +import ( + flaggerv1 "github.com/fluxcd/flagger/pkg/apis/flagger/v1beta1" +) + +// ScalerReconciler represents a reconciler that can reconcile resources +// that can scale other resources. +type ScalerReconciler interface { + ReconcilePrimaryScaler(cd *flaggerv1.Canary, init bool) error + PauseTargetScaler(cd *flaggerv1.Canary) error + ResumeTargetScaler(cd *flaggerv1.Canary) error +} diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index 54f142d6..eac33064 100644 --- a/pkg/controller/scheduler.go +++ b/pkg/controller/scheduler.go @@ -178,6 +178,11 @@ func (c *Controller) advanceCanary(name string, namespace string) { return } + var scalerReconciler canary.ScalerReconciler + if cd.Spec.AutoscalerRef != nil { + scalerReconciler = c.canaryFactory.ScalerReconciler(cd.Spec.AutoscalerRef.Kind) + } + // init Kubernetes router kubeRouter := c.routerFactory.KubernetesRouter(cd.Spec.TargetRef.Kind, labelSelector, labelValue, ports) @@ -213,6 +218,14 @@ func (c *Controller) advanceCanary(name string, namespace string) { return } + if scalerReconciler != nil { + err = scalerReconciler.ReconcilePrimaryScaler(cd, true) + if err != nil { + c.recordEventWarningf(cd, "%v", err) + return + } + } + // change the apex service pod selector to primary if err := kubeRouter.Reconcile(cd); err != nil { c.recordEventWarningf(cd, "%v", err) @@ -304,7 +317,7 @@ func (c *Controller) advanceCanary(name string, namespace string) { } // check if analysis should be skipped - if skip := c.shouldSkipAnalysis(cd, canaryController, meshRouter, err, retriable); skip { + if skip := c.shouldSkipAnalysis(cd, canaryController, meshRouter, scalerReconciler, err, retriable); skip { return } @@ -322,6 +335,13 @@ func (c *Controller) advanceCanary(name string, namespace string) { // route traffic back to primary if analysis has succeeded if cd.Status.Phase == flaggerv1.CanaryPhasePromoting { + if scalerReconciler != nil { + err = scalerReconciler.ReconcilePrimaryScaler(cd, false) + if err != nil { + c.recordEventWarningf(cd, "%v", err) + return + } + } c.runPromotionTrafficShift(cd, canaryController, meshRouter, provider, canaryWeight, primaryWeight) return } @@ -694,7 +714,7 @@ func (c *Controller) runAnalysis(canary *flaggerv1.Canary) bool { return true } -func (c *Controller) shouldSkipAnalysis(canary *flaggerv1.Canary, canaryController canary.Controller, meshRouter router.Interface, err error, retriable bool) bool { +func (c *Controller) shouldSkipAnalysis(canary *flaggerv1.Canary, canaryController canary.Controller, meshRouter router.Interface, scalerReconciler canary.ScalerReconciler, err error, retriable bool) bool { if !canary.SkipAnalysis() { return false } @@ -725,6 +745,13 @@ func (c *Controller) shouldSkipAnalysis(canary *flaggerv1.Canary, canaryControll return true } + if scalerReconciler != nil { + if err := scalerReconciler.ReconcilePrimaryScaler(canary, false); err != nil { + c.recordEventWarningf(canary, "%v", err) + return true + } + } + // shutdown canary if err := canaryController.ScaleToZero(canary); err != nil { c.recordEventWarningf(canary, "%v", err)