Merge pull request #1211 from aryan9600/scaler-reconciler

Introduce `ScalerReconciler` and refactor HPA reconciliation
This commit is contained in:
Stefan Prodan
2022-06-08 10:21:02 +03:00
committed by GitHub
11 changed files with 763 additions and 85 deletions
-22
View File
@@ -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
}
-29
View File
@@ -50,10 +50,6 @@ func TestDeploymentController_Sync_ConsistentNaming(t *testing.T) {
annotation := depPrimary.Annotations["kustomize.toolkit.fluxcd.io/checksum"]
assert.Equal(t, "", annotation)
hpaPrimary, err := mocks.kubeClient.AutoscalingV2beta2().HorizontalPodAutoscalers("default").Get(context.TODO(), "podinfo-primary", metav1.GetOptions{})
require.NoError(t, err)
assert.Equal(t, depPrimary.Name, hpaPrimary.Spec.ScaleTargetRef.Name)
}
func TestDeploymentController_Sync_InconsistentNaming(t *testing.T) {
@@ -71,10 +67,6 @@ func TestDeploymentController_Sync_InconsistentNaming(t *testing.T) {
primarySelectorValue := depPrimary.Spec.Selector.MatchLabels[dc.label]
assert.Equal(t, primarySelectorValue, fmt.Sprintf("%s-primary", dc.labelValue))
hpaPrimary, err := mocks.kubeClient.AutoscalingV2beta2().HorizontalPodAutoscalers("default").Get(context.TODO(), "podinfo-primary", metav1.GetOptions{})
require.NoError(t, err)
assert.Equal(t, depPrimary.Name, hpaPrimary.Spec.ScaleTargetRef.Name)
}
func TestDeploymentController_Promote(t *testing.T) {
@@ -90,15 +82,6 @@ func TestDeploymentController_Promote(t *testing.T) {
_, err = mocks.kubeClient.CoreV1().ConfigMaps("default").Update(context.TODO(), config2, metav1.UpdateOptions{})
require.NoError(t, err)
hpa, err := mocks.kubeClient.AutoscalingV2beta2().HorizontalPodAutoscalers("default").Get(context.TODO(), "podinfo", metav1.GetOptions{})
require.NoError(t, err)
hpaClone := hpa.DeepCopy()
hpaClone.Spec.MaxReplicas = 2
_, err = mocks.kubeClient.AutoscalingV2beta2().HorizontalPodAutoscalers("default").Update(context.TODO(), hpaClone, metav1.UpdateOptions{})
require.NoError(t, err)
err = mocks.controller.Promote(mocks.canary)
require.NoError(t, err)
@@ -121,18 +104,6 @@ func TestDeploymentController_Promote(t *testing.T) {
require.NoError(t, err)
assert.Equal(t, config2.Data["color"], configPrimary.Data["color"])
hpaPrimary, err := mocks.kubeClient.AutoscalingV2beta2().HorizontalPodAutoscalers("default").Get(context.TODO(), "podinfo-primary", metav1.GetOptions{})
require.NoError(t, err)
assert.Equal(t, int32(2), hpaPrimary.Spec.MaxReplicas)
hpaPrimaryLabels := hpaPrimary.ObjectMeta.Labels
hpaSourceLabels := hpa.ObjectMeta.Labels
assert.Equal(t, hpaSourceLabels["app.kubernetes.io/test-label-1"], hpaPrimaryLabels["app.kubernetes.io/test-label-1"])
hpaPrimaryAnnotations := hpaPrimary.ObjectMeta.Annotations
hpaSourceAnnotations := hpa.ObjectMeta.Annotations
assert.Equal(t, hpaSourceAnnotations["app.kubernetes.io/test-annotation-1"], hpaPrimaryAnnotations["app.kubernetes.io/test-annotation-1"])
value := depPrimary.Spec.Template.Spec.Affinity.PodAntiAffinity.PreferredDuringSchedulingIgnoredDuringExecution[0].PodAffinityTerm.LabelSelector.MatchExpressions[0].Values[0]
assert.Equal(t, "podinfo-primary", value)
-32
View File
@@ -24,7 +24,6 @@ import (
"github.com/stretchr/testify/require"
"go.uber.org/zap"
appsv1 "k8s.io/api/apps/v1"
hpav2 "k8s.io/api/autoscaling/v2beta2"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/resource"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -91,7 +90,6 @@ func newCustomizableFixture(dc deploymentConfigs) (deploymentControllerFixture,
// init kube clientset and register mock objects
kubeClient := fake.NewSimpleClientset(
newDeploymentControllerTest(dc),
newDeploymentControllerTestHPA(),
newDeploymentControllerTestConfigMap(),
newDeploymentControllerTestConfigMapEnv(),
newDeploymentControllerTestConfigMapVol(),
@@ -1000,33 +998,3 @@ func newDeploymentControllerTestV2() *appsv1.Deployment {
return d
}
func newDeploymentControllerTestHPA() *hpav2.HorizontalPodAutoscaler {
h := &hpav2.HorizontalPodAutoscaler{
TypeMeta: metav1.TypeMeta{APIVersion: hpav2.SchemeGroupVersion.String()},
ObjectMeta: metav1.ObjectMeta{
Namespace: "default",
Name: "podinfo",
},
Spec: hpav2.HorizontalPodAutoscalerSpec{
ScaleTargetRef: hpav2.CrossVersionObjectReference{
Name: "podinfo",
APIVersion: "apps/v1",
Kind: "Deployment",
},
Metrics: []hpav2.MetricSpec{
{
Type: "Resource",
Resource: &hpav2.ResourceMetricSource{
Name: "cpu",
Target: hpav2.MetricTarget{
AverageUtilization: int32p(99),
},
},
},
},
},
}
return h
}
+16
View File
@@ -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
}
}
+289
View File
@@ -0,0 +1,289 @@
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: makeObjectMeta(primaryHpaName, hpa.Labels, cd),
Spec: hpaSpec,
}
_, err = hr.kubeClient.AutoscalingV2().HorizontalPodAutoscalers(cd.Namespace).Create(context.TODO(), primaryHpa, metav1.CreateOptions{})
if err != nil {
return fmt.Errorf("creating HorizontalPodAutoscaler v2 %s.%s failed: %w",
primaryHpa.Name, primaryHpa.Namespace, err)
}
hr.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof(
"HorizontalPodAutoscaler v2 %s.%s created", primaryHpa.GetName(), cd.Namespace)
return nil
} else if err != nil {
return fmt.Errorf("HorizontalPodAutoscaler v2 %s.%s get query failed: %w",
primaryHpa.Name, primaryHpa.Namespace, err)
}
// update HPA
if !init && primaryHpa != nil {
targetFields := hpaFields{
metrics: hpaSpec.Metrics,
behavior: hpaSpec.Behavior,
annotations: hpa.Annotations,
labels: hpa.Labels,
min: hpaSpec.MinReplicas,
max: hpaSpec.MaxReplicas,
}
primaryFields := hpaFields{
metrics: primaryHpa.Spec.Metrics,
behavior: primaryHpa.Spec.Behavior,
annotations: primaryHpa.Annotations,
labels: primaryHpa.Labels,
min: primaryHpa.Spec.MinReplicas,
max: primaryHpa.Spec.MaxReplicas,
}
if hasHPAChanged(targetFields, primaryFields) {
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
hr.updateObjectMeta(hpaClone.ObjectMeta)
_, err = hr.kubeClient.AutoscalingV2().HorizontalPodAutoscalers(cd.Namespace).Update(context.TODO(), hpaClone, metav1.UpdateOptions{})
return err
})
if err != nil {
return fmt.Errorf("updating HorizontalPodAutoscaler v2 %s.%s failed: %w",
primaryHpa.Name, primaryHpa.Namespace, err)
}
hr.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).
Infof("HorizontalPodAutoscaler v2 %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: makeObjectMeta(primaryHpaName, hpa.Labels, cd),
Spec: hpaSpec,
}
_, err = hr.kubeClient.AutoscalingV2beta2().HorizontalPodAutoscalers(cd.Namespace).Create(context.TODO(), primaryHpa, metav1.CreateOptions{})
if err != nil {
return fmt.Errorf("creating HorizontalPodAutoscaler v2beta2 %s.%s failed: %w",
primaryHpa.Name, primaryHpa.Namespace, err)
}
hr.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof(
"HorizontalPodAutoscaler v2beta2 %s.%s created", primaryHpa.GetName(), cd.Namespace)
return nil
} else if err != nil {
return fmt.Errorf("HorizontalPodAutoscaler v2beta2 %s.%s get query failed: %w",
primaryHpa.Name, primaryHpa.Namespace, err)
}
// update HPA
if !init && primaryHpa != nil {
targetFields := hpaFields{
metrics: hpaSpec.Metrics,
behavior: hpaSpec.Behavior,
annotations: hpa.Annotations,
labels: hpa.Labels,
min: hpaSpec.MinReplicas,
max: hpaSpec.MaxReplicas,
}
primaryFields := hpaFields{
metrics: primaryHpa.Spec.Metrics,
behavior: primaryHpa.Spec.Behavior,
annotations: primaryHpa.Annotations,
labels: primaryHpa.Labels,
min: primaryHpa.Spec.MinReplicas,
max: primaryHpa.Spec.MaxReplicas,
}
if hasHPAChanged(targetFields, primaryFields) {
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
hr.updateObjectMeta(hpaClone.ObjectMeta)
_, err = hr.kubeClient.AutoscalingV2beta2().HorizontalPodAutoscalers(cd.Namespace).Update(context.TODO(), hpaClone, metav1.UpdateOptions{})
return err
})
if err != nil {
return fmt.Errorf("updating HorizontalPodAutoscaler v2beta2 %s.%s failed: %w",
primaryHpa.Name, primaryHpa.Namespace, err)
}
hr.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).
Infof("HorizontalPodAutoscaler v2beta2 %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
}
func (hr *HPAReconciler) updateObjectMeta(meta metav1.ObjectMeta) {
// update hpa annotations
meta.Annotations = make(map[string]string)
filteredAnnotations := includeLabelsByPrefix(meta.Annotations, hr.includeLabelPrefix)
for k, v := range filteredAnnotations {
meta.Annotations[k] = v
}
// update hpa labels
meta.Labels = make(map[string]string)
filteredLabels := includeLabelsByPrefix(meta.Labels, hr.includeLabelPrefix)
for k, v := range filteredLabels {
meta.Labels[k] = v
}
}
type hpaFields struct {
metrics interface{}
behavior interface{}
annotations map[string]string
labels map[string]string
min *int32
max int32
}
func hasHPAChanged(target, primary hpaFields) bool {
diffMetrics := cmp.Diff(target.metrics, primary.metrics)
diffBehavior := cmp.Diff(target.behavior, primary.behavior)
diffLabels := cmp.Diff(target.labels, primary.labels)
diffAnnotations := cmp.Diff(target.annotations, primary.annotations)
if diffMetrics != "" || diffBehavior != "" || diffLabels != "" || diffAnnotations != "" ||
int32Default(target.min) != int32Default(primary.min) || target.max != primary.max {
return true
}
return false
}
func makeObjectMeta(name string, labels map[string]string, cd *flaggerv1.Canary) metav1.ObjectMeta {
return metav1.ObjectMeta{
Name: name,
Namespace: cd.Namespace,
Labels: filterMetadata(labels),
OwnerReferences: []metav1.OwnerReference{
*metav1.NewControllerRef(cd, schema.GroupVersionKind{
Group: flaggerv1.SchemeGroupVersion.Group,
Version: flaggerv1.SchemeGroupVersion.Version,
Kind: flaggerv1.CanaryKind,
}),
},
}
}
+110
View File
@@ -0,0 +1,110 @@
package canary
import (
"context"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
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"
)
func Test_reconcilePrimaryHpa(t *testing.T) {
mocks := newScalerReconcilerFixture(scalerConfig{
targetName: "podinfo",
scaler: "HorizontalPodAutoscaler",
// avoid creating a v2 HPA.
excludeObjs: []string{"HPAV2"},
})
hpaReconciler := mocks.scalerReconciler.(*HPAReconciler)
err := hpaReconciler.reconcilePrimaryHpa(mocks.canary, true)
require.NoError(t, err)
// assert that we fallback to v2beta2, when HPAv2 fails.
_, err = mocks.kubeClient.AutoscalingV2().HorizontalPodAutoscalers("default").Get(context.TODO(), "podinfo-primary", metav1.GetOptions{})
assert.True(t, errors.IsNotFound(err))
hpa, err := mocks.kubeClient.AutoscalingV2beta2().HorizontalPodAutoscalers("default").Get(context.TODO(), "podinfo-primary", metav1.GetOptions{})
require.NoError(t, err)
require.NotNil(t, hpa)
mocks = newScalerReconcilerFixture(scalerConfig{
targetName: "podinfo",
scaler: "HorizontalPodAutoscaler",
// avoid creating _any_ HPAs.
excludeObjs: []string{"HPAV2", "HPAV2Beta2"},
})
hpaReconciler = mocks.scalerReconciler.(*HPAReconciler)
// assert that we return an error if no HPAs are found.
err = hpaReconciler.reconcilePrimaryHpa(mocks.canary, true)
require.Error(t, err)
}
func Test_reconcilePrimaryHpaV2(t *testing.T) {
mocks := newScalerReconcilerFixture(scalerConfig{
targetName: "podinfo",
scaler: "HorizontalPodAutoscaler",
})
hpaReconciler := mocks.scalerReconciler.(*HPAReconciler)
hpa, err := mocks.kubeClient.AutoscalingV2().HorizontalPodAutoscalers("default").Get(context.TODO(), "podinfo", metav1.GetOptions{})
require.NoError(t, err)
err = hpaReconciler.reconcilePrimaryHpaV2(mocks.canary, hpa, true)
require.NoError(t, err)
primaryHPA, err := mocks.kubeClient.AutoscalingV2().HorizontalPodAutoscalers("default").Get(context.TODO(), "podinfo-primary", metav1.GetOptions{})
require.NoError(t, err)
assert.Equal(t, primaryHPA.Spec.ScaleTargetRef.Name, "podinfo-primary")
assert.Equal(t, int(*primaryHPA.Spec.Metrics[0].Resource.Target.AverageUtilization), 99)
hpa.Spec.Metrics[0].Resource.Target = hpav2.MetricTarget{AverageUtilization: int32p(50)}
hpa.Spec.MaxReplicas = 10
_, err = mocks.kubeClient.AutoscalingV2().HorizontalPodAutoscalers("default").Update(context.TODO(), hpa, metav1.UpdateOptions{})
require.NoError(t, err)
err = hpaReconciler.reconcilePrimaryHpaV2(mocks.canary, hpa, false)
require.NoError(t, err)
primaryHPA, err = mocks.kubeClient.AutoscalingV2().HorizontalPodAutoscalers("default").Get(context.TODO(), "podinfo-primary", metav1.GetOptions{})
require.NoError(t, err)
assert.Equal(t, int(*primaryHPA.Spec.Metrics[0].Resource.Target.AverageUtilization), 50)
assert.Equal(t, int(primaryHPA.Spec.MaxReplicas), 10)
}
func Test_reconcilePrimaryHpaV2Beta2(t *testing.T) {
mocks := newScalerReconcilerFixture(scalerConfig{
targetName: "podinfo",
scaler: "HorizontalPodAutoscaler",
})
hpaReconciler := mocks.scalerReconciler.(*HPAReconciler)
hpa, err := mocks.kubeClient.AutoscalingV2beta2().HorizontalPodAutoscalers("default").Get(context.TODO(), "podinfo", metav1.GetOptions{})
require.NoError(t, err)
err = hpaReconciler.reconcilePrimaryHpaV2Beta2(mocks.canary, hpa, true)
require.NoError(t, err)
primaryHPA, err := mocks.kubeClient.AutoscalingV2beta2().HorizontalPodAutoscalers("default").Get(context.TODO(), "podinfo-primary", metav1.GetOptions{})
require.NoError(t, err)
assert.Equal(t, primaryHPA.Spec.ScaleTargetRef.Name, "podinfo-primary")
assert.Equal(t, int(*primaryHPA.Spec.Metrics[0].Resource.Target.AverageUtilization), 99)
hpa.Spec.Metrics[0].Resource.Target = hpav2beta2.MetricTarget{AverageUtilization: int32p(50)}
hpa.Spec.MaxReplicas = 10
_, err = mocks.kubeClient.AutoscalingV2beta2().HorizontalPodAutoscalers("default").Update(context.TODO(), hpa, metav1.UpdateOptions{})
require.NoError(t, err)
err = hpaReconciler.reconcilePrimaryHpaV2Beta2(mocks.canary, hpa, false)
require.NoError(t, err)
primaryHPA, err = mocks.kubeClient.AutoscalingV2beta2().HorizontalPodAutoscalers("default").Get(context.TODO(), "podinfo-primary", metav1.GetOptions{})
require.NoError(t, err)
assert.Equal(t, int(*primaryHPA.Spec.Metrics[0].Resource.Target.AverageUtilization), 50)
assert.Equal(t, int(primaryHPA.Spec.MaxReplicas), 10)
}
+13
View File
@@ -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
}
@@ -0,0 +1,136 @@
package canary
import (
"go.uber.org/zap"
hpav2 "k8s.io/api/autoscaling/v2"
hpav2beta2 "k8s.io/api/autoscaling/v2beta2"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime/schema"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/kubernetes/fake"
flaggerv1 "github.com/fluxcd/flagger/pkg/apis/flagger/v1beta1"
clientset "github.com/fluxcd/flagger/pkg/client/clientset/versioned"
fakeFlagger "github.com/fluxcd/flagger/pkg/client/clientset/versioned/fake"
"github.com/fluxcd/flagger/pkg/logger"
)
type scalerReconcilerFixture struct {
canary *flaggerv1.Canary
kubeClient kubernetes.Interface
flaggerClient clientset.Interface
scalerReconciler ScalerReconciler
logger *zap.SugaredLogger
}
type scalerConfig struct {
targetName string
excludeObjs []string
scaler string
}
func newScalerReconcilerFixture(cfg scalerConfig) scalerReconcilerFixture {
canary := newDeploymentControllerTestCanary(canaryConfigs{targetName: cfg.targetName})
flaggerClient := fakeFlagger.NewSimpleClientset(canary)
kubeClient := fake.NewSimpleClientset(
newScalerReconcilerTestHPAV2(),
newScalerReconcilerTestHPAV2Beta2(),
)
for _, obj := range cfg.excludeObjs {
if obj == "HPAV2" {
kubeClient.Tracker().Delete(schema.GroupVersionResource{
Group: "autoscaling",
Version: "v2",
Resource: "horizontalpodautoscalers",
}, "default", "podinfo")
}
if obj == "HPAV2Beta2" {
kubeClient.Tracker().Delete(schema.GroupVersionResource{
Group: "autoscaling",
Version: "v2beta2",
Resource: "horizontalpodautoscalers",
}, "default", "podinfo")
}
}
logger, _ := logger.NewLogger("debug")
var hpaReconciler HPAReconciler
if cfg.scaler == "HorizontalPodAutoscaler" {
hpaReconciler = HPAReconciler{
kubeClient: kubeClient,
flaggerClient: flaggerClient,
logger: logger,
includeLabelPrefix: []string{"app.kubernetes.io"},
}
}
return scalerReconcilerFixture{
canary: canary,
kubeClient: kubeClient,
flaggerClient: flaggerClient,
scalerReconciler: &hpaReconciler,
logger: logger,
}
}
func newScalerReconcilerTestHPAV2Beta2() *hpav2beta2.HorizontalPodAutoscaler {
h := &hpav2beta2.HorizontalPodAutoscaler{
TypeMeta: metav1.TypeMeta{APIVersion: hpav2beta2.SchemeGroupVersion.String()},
ObjectMeta: metav1.ObjectMeta{
Namespace: "default",
Name: "podinfo",
},
Spec: hpav2beta2.HorizontalPodAutoscalerSpec{
ScaleTargetRef: hpav2beta2.CrossVersionObjectReference{
Name: "podinfo",
APIVersion: "apps/v1",
Kind: "Deployment",
},
Metrics: []hpav2beta2.MetricSpec{
{
Type: "Resource",
Resource: &hpav2beta2.ResourceMetricSource{
Name: "cpu",
Target: hpav2beta2.MetricTarget{
AverageUtilization: int32p(99),
},
},
},
},
},
}
return h
}
func newScalerReconcilerTestHPAV2() *hpav2.HorizontalPodAutoscaler {
h := &hpav2.HorizontalPodAutoscaler{
TypeMeta: metav1.TypeMeta{APIVersion: hpav2.SchemeGroupVersion.String()},
ObjectMeta: metav1.ObjectMeta{
Namespace: "default",
Name: "podinfo",
},
Spec: hpav2.HorizontalPodAutoscalerSpec{
ScaleTargetRef: hpav2.CrossVersionObjectReference{
Name: "podinfo",
APIVersion: "apps/v1",
Kind: "Deployment",
},
Metrics: []hpav2.MetricSpec{
{
Type: "Resource",
Resource: &hpav2.ResourceMetricSource{
Name: "cpu",
Target: hpav2.MetricTarget{
AverageUtilization: int32p(99),
},
},
},
},
},
}
return h
}
+29 -2
View File
@@ -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)
+5
View File
@@ -10,3 +10,8 @@ DIR="$(cd "$(dirname "$0")" && pwd)"
"$REPO_ROOT"/test/workloads/init.sh
"$DIR"/test-deployment.sh
"$DIR"/test-daemonset.sh
kubectl -n test delete deploy podinfo
kubectl -n test delete svc podinfo-svc
kubectl apply -f ${REPO_ROOT}/test/workloads/deployment.yaml -n test
"$DIR"/test-hpa.sh
+165
View File
@@ -0,0 +1,165 @@
#!/usr/bin/env bash
set -o errexit
REPO_ROOT=$(git rev-parse --show-toplevel)
cat <<EOF | kubectl apply -f -
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
name: podinfo
namespace: test
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: podinfo
minReplicas: 2
maxReplicas: 4
metrics:
- type: Resource
resource:
name: cpu
target:
type: Utilization
averageUtilization: 99
EOF
cat <<EOF | kubectl apply -f -
apiVersion: v1
kind: Service
metadata:
name: podinfo-svc
namespace: test
spec:
type: ClusterIP
selector:
app: podinfo
ports:
- name: http
port: 9898
protocol: TCP
targetPort: http
---
apiVersion: flagger.app/v1beta1
kind: Canary
metadata:
name: podinfo
namespace: test
spec:
provider: kubernetes
targetRef:
apiVersion: apps/v1
kind: Deployment
name: podinfo
autoscalerRef:
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
name: podinfo
progressDeadlineSeconds: 60
service:
port: 80
targetPort: 9898
name: podinfo-svc
portDiscovery: true
analysis:
interval: 15s
threshold: 10
iterations: 5
metrics:
- name: request-success-rate
interval: 1m
thresholdRange:
min: 99
- name: request-duration
interval: 30s
thresholdRange:
max: 500
webhooks:
- name: "gate"
type: confirm-rollout
url: http://flagger-loadtester.test/gate/approve
- name: acceptance-test
type: pre-rollout
url: http://flagger-loadtester.test/
timeout: 10s
metadata:
type: bash
cmd: "curl -sd 'test' http://podinfo-svc-canary/token | grep token"
- name: load-test
url: http://flagger-loadtester.test/
timeout: 5s
metadata:
type: cmd
cmd: "hey -z 10m -q 10 -c 2 http://podinfo-svc-canary.test/"
EOF
echo '>>> Waiting for primary to be ready'
retries=50
count=0
ok=false
until ${ok}; do
kubectl -n test get canary/podinfo | grep 'Initialized' && ok=true || ok=false
sleep 5
count=$(($count + 1))
if [[ ${count} -eq ${retries} ]]; then
kubectl -n flagger-system logs deployment/flagger
echo "No more retries left"
exit 1
fi
done
echo '✔ Canary initialization test passed'
echo '>>> Waiting for primary HPA to be created'
retries=10
count=0
ok=false
until ${ok}; do
kubectl -n test get hpa podinfo-primary && ok=true || ok=false
sleep 2
count=$(($count + 1))
if [[ ${count} -eq ${retries} ]]; then
kubectl -n flagger-system logs deployment/flagger
echo "No more retries left"
echo ' Primary HPA not found'
exit 1
fi
done
echo '✔ Primary HPA successfully reconciled'
# update the target hpa resource utilization
kubectl -n test patch hpa podinfo --type json --patch='[ { "op": "replace", "path": "/spec/metrics/0/resource/target/averageUtilization", "value": 50 } ]'
echo '>>> Triggering canary deployment'
kubectl -n test set image deployment/podinfo podinfod=ghcr.io/stefanprodan/podinfo:6.0.1
echo '>>> Waiting for canary promotion'
retries=50
count=0
ok=false
until ${ok}; do
kubectl -n test describe deployment/podinfo-primary | grep '6.0.1' && ok=true || ok=false
sleep 10
kubectl -n flagger-system logs deployment/flagger --tail 1
count=$(($count + 1))
if [[ ${count} -eq ${retries} ]]; then
kubectl -n test describe deployment/podinfo
kubectl -n test describe deployment/podinfo-primary
kubectl -n flagger-system logs deployment/flagger
echo "No more retries left"
exit 1
fi
done
echo '✔ Canary promotion test passed'
util=$(kubectl -n test get hpa podinfo -ojsonpath='{.spec.metrics[0].resource.target.averageUtilization}' | xargs)
if [[ ${util} -eq 50 ]]; then
echo '✔ Primary HPA successfully reconciled'
else
echo " Unexpected primary HPA resource target average utilization value: ${util}, expected: 50"
exit 1
fi