diff --git a/cmd/flagger/main.go b/cmd/flagger/main.go index 529c9b32..8548de56 100644 --- a/cmd/flagger/main.go +++ b/cmd/flagger/main.go @@ -10,16 +10,6 @@ import ( "time" "github.com/Masterminds/semver" - clientset "github.com/weaveworks/flagger/pkg/client/clientset/versioned" - informers "github.com/weaveworks/flagger/pkg/client/informers/externalversions" - "github.com/weaveworks/flagger/pkg/controller" - "github.com/weaveworks/flagger/pkg/logger" - "github.com/weaveworks/flagger/pkg/metrics" - "github.com/weaveworks/flagger/pkg/notifier" - "github.com/weaveworks/flagger/pkg/router" - "github.com/weaveworks/flagger/pkg/server" - "github.com/weaveworks/flagger/pkg/signals" - "github.com/weaveworks/flagger/pkg/version" "go.uber.org/zap" "k8s.io/apimachinery/pkg/util/uuid" "k8s.io/client-go/kubernetes" @@ -30,6 +20,18 @@ import ( "k8s.io/client-go/tools/leaderelection/resourcelock" "k8s.io/client-go/transport" _ "k8s.io/code-generator/cmd/client-gen/generators" + + "github.com/weaveworks/flagger/pkg/canary" + clientset "github.com/weaveworks/flagger/pkg/client/clientset/versioned" + informers "github.com/weaveworks/flagger/pkg/client/informers/externalversions" + "github.com/weaveworks/flagger/pkg/controller" + "github.com/weaveworks/flagger/pkg/logger" + "github.com/weaveworks/flagger/pkg/metrics" + "github.com/weaveworks/flagger/pkg/notifier" + "github.com/weaveworks/flagger/pkg/router" + "github.com/weaveworks/flagger/pkg/server" + "github.com/weaveworks/flagger/pkg/signals" + "github.com/weaveworks/flagger/pkg/version" ) var ( @@ -178,6 +180,12 @@ func main() { go server.ListenAndServe(port, 3*time.Second, logger, stopCh) routerFactory := router.NewFactory(cfg, kubeClient, flaggerClient, ingressAnnotationsPrefix, logger, meshClient) + configTracker := canary.ConfigTracker{ + Logger: logger, + KubeClient: kubeClient, + FlaggerClient: flaggerClient, + } + canaryFactory := canary.NewFactory(kubeClient, flaggerClient, configTracker, labels, logger) c := controller.NewController( kubeClient, @@ -187,11 +195,11 @@ func main() { controlLoopInterval, logger, notifierClient, + canaryFactory, routerFactory, observerFactory, meshProvider, version.VERSION, - labels, ) flaggerInformerFactory.Start(stopCh) diff --git a/pkg/canary/controller.go b/pkg/canary/controller.go new file mode 100644 index 00000000..a76e2603 --- /dev/null +++ b/pkg/canary/controller.go @@ -0,0 +1,19 @@ +package canary + +import "github.com/weaveworks/flagger/pkg/apis/flagger/v1alpha3" + +type Controller interface { + IsPrimaryReady(canary *v1alpha3.Canary) (bool, error) + IsCanaryReady(canary *v1alpha3.Canary) (bool, error) + SyncStatus(canary *v1alpha3.Canary, status v1alpha3.CanaryStatus) error + SetStatusFailedChecks(canary *v1alpha3.Canary, val int) error + SetStatusWeight(canary *v1alpha3.Canary, val int) error + SetStatusIterations(canary *v1alpha3.Canary, val int) error + SetStatusPhase(canary *v1alpha3.Canary, phase v1alpha3.CanaryPhase) error + Initialize(canary *v1alpha3.Canary, skipLivenessChecks bool) (label string, ports map[string]int32, err error) + Promote(canary *v1alpha3.Canary) error + HasTargetChanged(canary *v1alpha3.Canary) (bool, error) + HaveDependenciesChanged(canary *v1alpha3.Canary) (bool, error) + Scale(canary *v1alpha3.Canary, replicas int32) error + ScaleFromZero(canary *v1alpha3.Canary) error +} diff --git a/pkg/canary/deployer.go b/pkg/canary/deployment.go similarity index 80% rename from pkg/canary/deployer.go rename to pkg/canary/deployment.go index a11b80ff..f098dc9d 100644 --- a/pkg/canary/deployer.go +++ b/pkg/canary/deployment.go @@ -21,18 +21,18 @@ import ( clientset "github.com/weaveworks/flagger/pkg/client/clientset/versioned" ) -// Deployer is managing the operations for Kubernetes deployment kind -type Deployer struct { - KubeClient kubernetes.Interface - FlaggerClient clientset.Interface - Logger *zap.SugaredLogger - ConfigTracker ConfigTracker - Labels []string +// DeploymentController is managing the operations for Kubernetes Deployment kind +type DeploymentController struct { + kubeClient kubernetes.Interface + flaggerClient clientset.Interface + logger *zap.SugaredLogger + configTracker ConfigTracker + labels []string } // Initialize creates the primary deployment, hpa, // scales to zero the canary deployment and returns the pod selector label and container ports -func (c *Deployer) Initialize(cd *flaggerv1.Canary, skipLivenessChecks bool) (label string, ports map[string]int32, err error) { +func (c *DeploymentController) Initialize(cd *flaggerv1.Canary, skipLivenessChecks bool) (label string, ports map[string]int32, err error) { primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) label, ports, err = c.createPrimaryDeployment(cd) if err != nil { @@ -47,7 +47,7 @@ func (c *Deployer) Initialize(cd *flaggerv1.Canary, skipLivenessChecks bool) (la } } - c.Logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("Scaling down %s.%s", cd.Spec.TargetRef.Name, cd.Namespace) + c.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("Scaling down %s.%s", cd.Spec.TargetRef.Name, cd.Namespace) if err := c.Scale(cd, 0); err != nil { return "", ports, err } @@ -62,11 +62,11 @@ func (c *Deployer) Initialize(cd *flaggerv1.Canary, skipLivenessChecks bool) (la } // Promote copies the pod spec, secrets and config maps from canary to primary -func (c *Deployer) Promote(cd *flaggerv1.Canary) error { +func (c *DeploymentController) Promote(cd *flaggerv1.Canary) error { targetName := cd.Spec.TargetRef.Name primaryName := fmt.Sprintf("%s-primary", targetName) - canary, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) + canary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { if errors.IsNotFound(err) { return fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) @@ -80,7 +80,7 @@ func (c *Deployer) Promote(cd *flaggerv1.Canary) error { targetName, cd.Namespace, targetName) } - primary, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(primaryName, metav1.GetOptions{}) + 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) @@ -89,11 +89,11 @@ func (c *Deployer) Promote(cd *flaggerv1.Canary) error { } // promote secrets and config maps - configRefs, err := c.ConfigTracker.GetTargetConfigs(cd) + configRefs, err := c.configTracker.GetTargetConfigs(cd) if err != nil { return err } - if err := c.ConfigTracker.CreatePrimaryConfigs(cd, configRefs); err != nil { + if err := c.configTracker.CreatePrimaryConfigs(cd, configRefs); err != nil { return err } @@ -104,7 +104,7 @@ func (c *Deployer) Promote(cd *flaggerv1.Canary) error { primaryCopy.Spec.Strategy = canary.Spec.Strategy // update spec with primary secrets and config maps - primaryCopy.Spec.Template.Spec = c.ConfigTracker.ApplyPrimaryConfigs(canary.Spec.Template.Spec, configRefs) + primaryCopy.Spec.Template.Spec = c.configTracker.ApplyPrimaryConfigs(canary.Spec.Template.Spec, configRefs) // update pod annotations to ensure a rolling update annotations, err := c.makeAnnotations(canary.Spec.Template.Annotations) @@ -116,7 +116,7 @@ func (c *Deployer) Promote(cd *flaggerv1.Canary) error { primaryCopy.Spec.Template.Labels = makePrimaryLabels(canary.Spec.Template.Labels, primaryName, label) // apply update - _, err = c.KubeClient.AppsV1().Deployments(cd.Namespace).Update(primaryCopy) + _, err = c.kubeClient.AppsV1().Deployments(cd.Namespace).Update(primaryCopy) if err != nil { return fmt.Errorf("updating deployment %s.%s template spec failed: %v", primaryCopy.GetName(), primaryCopy.Namespace, err) @@ -132,10 +132,10 @@ func (c *Deployer) Promote(cd *flaggerv1.Canary) error { return nil } -// HasDeploymentChanged returns true if the canary deployment pod spec has changed -func (c *Deployer) HasDeploymentChanged(cd *flaggerv1.Canary) (bool, error) { +// HasTargetChanged returns true if the canary deployment pod spec has changed +func (c *DeploymentController) HasTargetChanged(cd *flaggerv1.Canary) (bool, error) { targetName := cd.Spec.TargetRef.Name - canary, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) + canary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { if errors.IsNotFound(err) { return false, fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) @@ -164,10 +164,15 @@ func (c *Deployer) HasDeploymentChanged(cd *flaggerv1.Canary) (bool, error) { return false, nil } +// HaveDependenciesChanged returns true if the canary configmaps or secrets have changed +func (c *DeploymentController) HaveDependenciesChanged(cd *flaggerv1.Canary) (bool, error) { + return c.configTracker.HasConfigChanged(cd) +} + // Scale sets the canary deployment replicas -func (c *Deployer) Scale(cd *flaggerv1.Canary, replicas int32) error { +func (c *DeploymentController) Scale(cd *flaggerv1.Canary, replicas int32) error { targetName := cd.Spec.TargetRef.Name - dep, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) + dep, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { if errors.IsNotFound(err) { return fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) @@ -178,16 +183,16 @@ func (c *Deployer) Scale(cd *flaggerv1.Canary, replicas int32) error { depCopy := dep.DeepCopy() depCopy.Spec.Replicas = int32p(replicas) - _, err = c.KubeClient.AppsV1().Deployments(dep.Namespace).Update(depCopy) + _, err = c.kubeClient.AppsV1().Deployments(dep.Namespace).Update(depCopy) if err != nil { return fmt.Errorf("scaling %s.%s to %v failed: %v", depCopy.GetName(), depCopy.Namespace, replicas, err) } return nil } -func (c *Deployer) ScaleUp(cd *flaggerv1.Canary) error { +func (c *DeploymentController) ScaleFromZero(cd *flaggerv1.Canary) error { targetName := cd.Spec.TargetRef.Name - dep, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) + dep, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { if errors.IsNotFound(err) { return fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) @@ -202,18 +207,18 @@ func (c *Deployer) ScaleUp(cd *flaggerv1.Canary) error { depCopy := dep.DeepCopy() depCopy.Spec.Replicas = replicas - _, err = c.KubeClient.AppsV1().Deployments(dep.Namespace).Update(depCopy) + _, err = c.kubeClient.AppsV1().Deployments(dep.Namespace).Update(depCopy) if err != nil { return fmt.Errorf("scaling %s.%s to %v failed: %v", depCopy.GetName(), depCopy.Namespace, replicas, err) } return nil } -func (c *Deployer) createPrimaryDeployment(cd *flaggerv1.Canary) (string, map[string]int32, error) { +func (c *DeploymentController) createPrimaryDeployment(cd *flaggerv1.Canary) (string, map[string]int32, error) { targetName := cd.Spec.TargetRef.Name primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) - canaryDep, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) + canaryDep, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { if errors.IsNotFound(err) { return "", nil, fmt.Errorf("deployment %s.%s not found, retrying", targetName, cd.Namespace) @@ -236,14 +241,14 @@ func (c *Deployer) createPrimaryDeployment(cd *flaggerv1.Canary) (string, map[st ports = p } - primaryDep, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(primaryName, metav1.GetOptions{}) + primaryDep, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(primaryName, metav1.GetOptions{}) if errors.IsNotFound(err) { // create primary secrets and config maps - configRefs, err := c.ConfigTracker.GetTargetConfigs(cd) + configRefs, err := c.configTracker.GetTargetConfigs(cd) if err != nil { return "", nil, err } - if err := c.ConfigTracker.CreatePrimaryConfigs(cd, configRefs); err != nil { + if err := c.configTracker.CreatePrimaryConfigs(cd, configRefs); err != nil { return "", nil, err } annotations, err := c.makeAnnotations(canaryDep.Spec.Template.Annotations) @@ -289,25 +294,25 @@ func (c *Deployer) createPrimaryDeployment(cd *flaggerv1.Canary) (string, map[st Annotations: annotations, }, // update spec with the primary secrets and config maps - Spec: c.ConfigTracker.ApplyPrimaryConfigs(canaryDep.Spec.Template.Spec, configRefs), + Spec: c.configTracker.ApplyPrimaryConfigs(canaryDep.Spec.Template.Spec, configRefs), }, }, } - _, err = c.KubeClient.AppsV1().Deployments(cd.Namespace).Create(primaryDep) + _, err = c.kubeClient.AppsV1().Deployments(cd.Namespace).Create(primaryDep) if err != nil { return "", nil, err } - c.Logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("Deployment %s.%s created", primaryDep.GetName(), cd.Namespace) + c.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("Deployment %s.%s created", primaryDep.GetName(), cd.Namespace) } return label, ports, nil } -func (c *Deployer) reconcilePrimaryHpa(cd *flaggerv1.Canary, init bool) error { +func (c *DeploymentController) reconcilePrimaryHpa(cd *flaggerv1.Canary, init bool) 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{}) + 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", @@ -328,7 +333,7 @@ func (c *Deployer) reconcilePrimaryHpa(cd *flaggerv1.Canary, init bool) error { } primaryHpaName := fmt.Sprintf("%s-primary", cd.Spec.AutoscalerRef.Name) - primaryHpa, err := c.KubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Get(primaryHpaName, metav1.GetOptions{}) + primaryHpa, err := c.kubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Get(primaryHpaName, metav1.GetOptions{}) // create HPA if errors.IsNotFound(err) { @@ -348,11 +353,11 @@ func (c *Deployer) reconcilePrimaryHpa(cd *flaggerv1.Canary, init bool) error { Spec: hpaSpec, } - _, err = c.KubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Create(primaryHpa) + _, err = c.kubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Create(primaryHpa) if err != nil { return err } - c.Logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("HorizontalPodAutoscaler %s.%s created", primaryHpa.GetName(), cd.Namespace) + c.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("HorizontalPodAutoscaler %s.%s created", primaryHpa.GetName(), cd.Namespace) return nil } @@ -370,11 +375,11 @@ func (c *Deployer) reconcilePrimaryHpa(cd *flaggerv1.Canary, init bool) error { hpaClone.Spec.MinReplicas = hpaSpec.MinReplicas hpaClone.Spec.Metrics = hpaSpec.Metrics - _, upErr := c.KubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Update(hpaClone) + _, upErr := c.kubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Update(hpaClone) if upErr != nil { return upErr } - c.Logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("HorizontalPodAutoscaler %s.%s updated", primaryHpa.GetName(), cd.Namespace) + c.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("HorizontalPodAutoscaler %s.%s updated", primaryHpa.GetName(), cd.Namespace) } } @@ -382,7 +387,7 @@ func (c *Deployer) reconcilePrimaryHpa(cd *flaggerv1.Canary, init bool) error { } // makeAnnotations appends an unique ID to annotations map -func (c *Deployer) makeAnnotations(annotations map[string]string) (map[string]string, error) { +func (c *DeploymentController) makeAnnotations(annotations map[string]string) (map[string]string, error) { idKey := "flagger-id" res := make(map[string]string) uuid := make([]byte, 16) @@ -405,8 +410,8 @@ func (c *Deployer) makeAnnotations(annotations map[string]string) (map[string]st } // getSelectorLabel returns the selector match label -func (c *Deployer) getSelectorLabel(deployment *appsv1.Deployment) (string, error) { - for _, l := range c.Labels { +func (c *DeploymentController) getSelectorLabel(deployment *appsv1.Deployment) (string, error) { + for _, l := range c.labels { if _, ok := deployment.Spec.Selector.MatchLabels[l]; ok { return l, nil } @@ -421,7 +426,7 @@ var sidecars = map[string]bool{ } // getPorts returns a list of all container ports -func (c *Deployer) getPorts(cd *flaggerv1.Canary, deployment *appsv1.Deployment) (map[string]int32, error) { +func (c *DeploymentController) getPorts(cd *flaggerv1.Canary, deployment *appsv1.Deployment) (map[string]int32, error) { ports := make(map[string]int32) for _, container := range deployment.Spec.Template.Spec.Containers { diff --git a/pkg/canary/deployer_test.go b/pkg/canary/deployment_test.go similarity index 99% rename from pkg/canary/deployer_test.go rename to pkg/canary/deployment_test.go index efffe937..66914f31 100644 --- a/pkg/canary/deployer_test.go +++ b/pkg/canary/deployment_test.go @@ -107,7 +107,7 @@ func TestCanaryDeployer_IsNewSpec(t *testing.T) { t.Fatal(err.Error()) } - isNew, err := mocks.deployer.HasDeploymentChanged(mocks.canary) + isNew, err := mocks.deployer.HasTargetChanged(mocks.canary) if err != nil { t.Fatal(err.Error()) } diff --git a/pkg/canary/factory.go b/pkg/canary/factory.go new file mode 100644 index 00000000..db608d5f --- /dev/null +++ b/pkg/canary/factory.go @@ -0,0 +1,52 @@ +package canary + +import ( + "go.uber.org/zap" + "k8s.io/client-go/kubernetes" + + clientset "github.com/weaveworks/flagger/pkg/client/clientset/versioned" +) + +type Factory struct { + kubeClient kubernetes.Interface + flaggerClient clientset.Interface + logger *zap.SugaredLogger + configTracker ConfigTracker + labels []string +} + +func NewFactory(kubeClient kubernetes.Interface, + flaggerClient clientset.Interface, + configTracker ConfigTracker, + labels []string, + logger *zap.SugaredLogger) *Factory { + return &Factory{ + kubeClient: kubeClient, + flaggerClient: flaggerClient, + logger: logger, + configTracker: configTracker, + labels: labels, + } +} + +func (factory *Factory) Controller(kind string) Controller { + deploymentCtrl := &DeploymentController{ + logger: factory.logger, + kubeClient: factory.kubeClient, + flaggerClient: factory.flaggerClient, + labels: factory.labels, + configTracker: ConfigTracker{ + Logger: factory.logger, + KubeClient: factory.kubeClient, + FlaggerClient: factory.flaggerClient, + }, + } + + switch { + case kind == "Deployment": + return deploymentCtrl + default: + return deploymentCtrl + } + +} diff --git a/pkg/canary/mock.go b/pkg/canary/mock.go index ac381de3..f80361e3 100644 --- a/pkg/canary/mock.go +++ b/pkg/canary/mock.go @@ -20,7 +20,7 @@ type Mocks struct { canary *flaggerv1.Canary kubeClient kubernetes.Interface flaggerClient clientset.Interface - deployer Deployer + deployer DeploymentController logger *zap.SugaredLogger } @@ -43,12 +43,12 @@ func SetupMocks() Mocks { logger, _ := logger.NewLogger("debug") - deployer := Deployer{ - FlaggerClient: flaggerClient, - KubeClient: kubeClient, - Logger: logger, - Labels: []string{"app", "name"}, - ConfigTracker: ConfigTracker{ + deployer := DeploymentController{ + flaggerClient: flaggerClient, + kubeClient: kubeClient, + logger: logger, + labels: []string{"app", "name"}, + configTracker: ConfigTracker{ Logger: logger, KubeClient: kubeClient, FlaggerClient: flaggerClient, diff --git a/pkg/canary/ready.go b/pkg/canary/ready.go index 9d47744e..cfd2377e 100644 --- a/pkg/canary/ready.go +++ b/pkg/canary/ready.go @@ -14,9 +14,9 @@ import ( // IsPrimaryReady checks the primary deployment status and returns an error if // the deployment is in the middle of a rolling update or if the pods are unhealthy // it will return a non retriable error if the rolling update is stuck -func (c *Deployer) IsPrimaryReady(cd *flaggerv1.Canary) (bool, error) { +func (c *DeploymentController) IsPrimaryReady(cd *flaggerv1.Canary) (bool, error) { primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) - primary, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(primaryName, metav1.GetOptions{}) + primary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(primaryName, metav1.GetOptions{}) if err != nil { if errors.IsNotFound(err) { return true, fmt.Errorf("deployment %s.%s not found", primaryName, cd.Namespace) @@ -39,9 +39,9 @@ func (c *Deployer) IsPrimaryReady(cd *flaggerv1.Canary) (bool, error) { // IsCanaryReady checks the primary deployment status and returns an error if // the deployment is in the middle of a rolling update or if the pods are unhealthy // it will return a non retriable error if the rolling update is stuck -func (c *Deployer) IsCanaryReady(cd *flaggerv1.Canary) (bool, error) { +func (c *DeploymentController) IsCanaryReady(cd *flaggerv1.Canary) (bool, error) { targetName := cd.Spec.TargetRef.Name - canary, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) + canary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { if errors.IsNotFound(err) { return true, fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) @@ -64,7 +64,7 @@ func (c *Deployer) IsCanaryReady(cd *flaggerv1.Canary) (bool, error) { // isDeploymentReady determines if a deployment is ready by checking the status conditions // if a deployment has exceeded the progress deadline it returns a non retriable error -func (c *Deployer) isDeploymentReady(deployment *appsv1.Deployment, deadline int) (bool, error) { +func (c *DeploymentController) isDeploymentReady(deployment *appsv1.Deployment, deadline int) (bool, error) { retriable := true if deployment.Generation <= deployment.Status.ObservedGeneration { progress := c.getDeploymentCondition(deployment.Status, appsv1.DeploymentProgressing) @@ -99,7 +99,7 @@ func (c *Deployer) isDeploymentReady(deployment *appsv1.Deployment, deadline int return true, nil } -func (c *Deployer) getDeploymentCondition( +func (c *DeploymentController) getDeploymentCondition( status appsv1.DeploymentStatus, conditionType appsv1.DeploymentConditionType, ) *appsv1.DeploymentCondition { diff --git a/pkg/canary/status.go b/pkg/canary/status.go index 26f0f693..530f9504 100644 --- a/pkg/canary/status.go +++ b/pkg/canary/status.go @@ -14,8 +14,8 @@ import ( ) // SyncStatus encodes the canary pod spec and updates the canary status -func (c *Deployer) SyncStatus(cd *flaggerv1.Canary, status flaggerv1.CanaryStatus) error { - dep, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(cd.Spec.TargetRef.Name, metav1.GetOptions{}) +func (c *DeploymentController) SyncStatus(cd *flaggerv1.Canary, status flaggerv1.CanaryStatus) error { + dep, 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) @@ -23,7 +23,7 @@ func (c *Deployer) SyncStatus(cd *flaggerv1.Canary, status flaggerv1.CanaryStatu return ex.Wrap(err, "SyncStatus deployment query error") } - configs, err := c.ConfigTracker.GetConfigRefs(cd) + configs, err := c.configTracker.GetConfigRefs(cd) if err != nil { return ex.Wrap(err, "SyncStatus configs query error") } @@ -37,7 +37,7 @@ func (c *Deployer) SyncStatus(cd *flaggerv1.Canary, status flaggerv1.CanaryStatu err = retry.RetryOnConflict(retry.DefaultBackoff, func() (err error) { var selErr error if !firstTry { - cd, selErr = c.FlaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).Get(cd.GetName(), metav1.GetOptions{}) + cd, selErr = c.flaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).Get(cd.GetName(), metav1.GetOptions{}) if selErr != nil { return selErr } @@ -51,11 +51,11 @@ func (c *Deployer) SyncStatus(cd *flaggerv1.Canary, status flaggerv1.CanaryStatu cdCopy.Status.LastTransitionTime = metav1.Now() cdCopy.Status.TrackedConfigs = configs - if ok, conditions := c.MakeStatusConditions(cd.Status, status.Phase); ok { + if ok, conditions := MakeStatusConditions(cd.Status, status.Phase); ok { cdCopy.Status.Conditions = conditions } - _, err = c.FlaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) + _, err = c.flaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) firstTry = false return }) @@ -66,12 +66,12 @@ func (c *Deployer) SyncStatus(cd *flaggerv1.Canary, status flaggerv1.CanaryStatu } // SetStatusFailedChecks updates the canary failed checks counter -func (c *Deployer) SetStatusFailedChecks(cd *flaggerv1.Canary, val int) error { +func (c *DeploymentController) SetStatusFailedChecks(cd *flaggerv1.Canary, val int) error { firstTry := true err := retry.RetryOnConflict(retry.DefaultBackoff, func() (err error) { var selErr error if !firstTry { - cd, selErr = c.FlaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).Get(cd.GetName(), metav1.GetOptions{}) + cd, selErr = c.flaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).Get(cd.GetName(), metav1.GetOptions{}) if selErr != nil { return selErr } @@ -80,7 +80,7 @@ func (c *Deployer) SetStatusFailedChecks(cd *flaggerv1.Canary, val int) error { cdCopy.Status.FailedChecks = val cdCopy.Status.LastTransitionTime = metav1.Now() - _, err = c.FlaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) + _, err = c.flaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) firstTry = false return }) @@ -91,12 +91,12 @@ func (c *Deployer) SetStatusFailedChecks(cd *flaggerv1.Canary, val int) error { } // SetStatusWeight updates the canary status weight value -func (c *Deployer) SetStatusWeight(cd *flaggerv1.Canary, val int) error { +func (c *DeploymentController) SetStatusWeight(cd *flaggerv1.Canary, val int) error { firstTry := true err := retry.RetryOnConflict(retry.DefaultBackoff, func() (err error) { var selErr error if !firstTry { - cd, selErr = c.FlaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).Get(cd.GetName(), metav1.GetOptions{}) + cd, selErr = c.flaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).Get(cd.GetName(), metav1.GetOptions{}) if selErr != nil { return selErr } @@ -105,7 +105,7 @@ func (c *Deployer) SetStatusWeight(cd *flaggerv1.Canary, val int) error { cdCopy.Status.CanaryWeight = val cdCopy.Status.LastTransitionTime = metav1.Now() - _, err = c.FlaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) + _, err = c.flaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) firstTry = false return }) @@ -116,12 +116,12 @@ func (c *Deployer) SetStatusWeight(cd *flaggerv1.Canary, val int) error { } // SetStatusIterations updates the canary status iterations value -func (c *Deployer) SetStatusIterations(cd *flaggerv1.Canary, val int) error { +func (c *DeploymentController) SetStatusIterations(cd *flaggerv1.Canary, val int) error { firstTry := true err := retry.RetryOnConflict(retry.DefaultBackoff, func() (err error) { var selErr error if !firstTry { - cd, selErr = c.FlaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).Get(cd.GetName(), metav1.GetOptions{}) + cd, selErr = c.flaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).Get(cd.GetName(), metav1.GetOptions{}) if selErr != nil { return selErr } @@ -131,7 +131,7 @@ func (c *Deployer) SetStatusIterations(cd *flaggerv1.Canary, val int) error { cdCopy.Status.Iterations = val cdCopy.Status.LastTransitionTime = metav1.Now() - _, err = c.FlaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) + _, err = c.flaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) firstTry = false return }) @@ -143,12 +143,12 @@ func (c *Deployer) SetStatusIterations(cd *flaggerv1.Canary, val int) error { } // SetStatusPhase updates the canary status phase -func (c *Deployer) SetStatusPhase(cd *flaggerv1.Canary, phase flaggerv1.CanaryPhase) error { +func (c *DeploymentController) SetStatusPhase(cd *flaggerv1.Canary, phase flaggerv1.CanaryPhase) error { firstTry := true err := retry.RetryOnConflict(retry.DefaultBackoff, func() (err error) { var selErr error if !firstTry { - cd, selErr = c.FlaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).Get(cd.GetName(), metav1.GetOptions{}) + cd, selErr = c.flaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).Get(cd.GetName(), metav1.GetOptions{}) if selErr != nil { return selErr } @@ -167,11 +167,11 @@ func (c *Deployer) SetStatusPhase(cd *flaggerv1.Canary, phase flaggerv1.CanaryPh cdCopy.Status.LastPromotedSpec = cd.Status.LastAppliedSpec } - if ok, conditions := c.MakeStatusConditions(cdCopy.Status, phase); ok { + if ok, conditions := MakeStatusConditions(cdCopy.Status, phase); ok { cdCopy.Status.Conditions = conditions } - _, err = c.FlaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) + _, err = c.flaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) firstTry = false return }) @@ -181,8 +181,8 @@ func (c *Deployer) SetStatusPhase(cd *flaggerv1.Canary, phase flaggerv1.CanaryPh return nil } -// GetStatusCondition returns a condition based on type -func (c *Deployer) getStatusCondition(status flaggerv1.CanaryStatus, conditionType flaggerv1.CanaryConditionType) *flaggerv1.CanaryCondition { +// getStatusCondition returns a condition based on type +func getStatusCondition(status flaggerv1.CanaryStatus, conditionType flaggerv1.CanaryConditionType) *flaggerv1.CanaryCondition { for i := range status.Conditions { c := status.Conditions[i] if c.Type == conditionType { @@ -193,9 +193,9 @@ func (c *Deployer) getStatusCondition(status flaggerv1.CanaryStatus, conditionTy } // MakeStatusCondition updates the canary status conditions based on canary phase -func (c *Deployer) MakeStatusConditions(canaryStatus flaggerv1.CanaryStatus, +func MakeStatusConditions(canaryStatus flaggerv1.CanaryStatus, phase flaggerv1.CanaryPhase) (bool, []flaggerv1.CanaryCondition) { - currentCondition := c.getStatusCondition(canaryStatus, flaggerv1.PromotedType) + currentCondition := getStatusCondition(canaryStatus, flaggerv1.PromotedType) message := "New deployment detected, starting initialization." status := corev1.ConditionUnknown diff --git a/pkg/canary/tracker.go b/pkg/canary/tracker.go index 57fb1f78..d94567bd 100644 --- a/pkg/canary/tracker.go +++ b/pkg/canary/tracker.go @@ -16,7 +16,7 @@ import ( clientset "github.com/weaveworks/flagger/pkg/client/clientset/versioned" ) -// ConfigTracker is managing the operations for Kubernetes ConfigMaps and Secrets +// configTracker is managing the operations for Kubernetes ConfigMaps and Secrets type ConfigTracker struct { KubeClient kubernetes.Interface FlaggerClient clientset.Interface diff --git a/pkg/controller/controller.go b/pkg/controller/controller.go index 948c611e..574e6dbe 100644 --- a/pkg/controller/controller.go +++ b/pkg/controller/controller.go @@ -45,9 +45,9 @@ type Controller struct { logger *zap.SugaredLogger canaries *sync.Map jobs map[string]CanaryJob - deployer canary.Deployer recorder metrics.Recorder notifier notifier.Interface + canaryFactory *canary.Factory routerFactory *router.Factory observerFactory *metrics.Factory meshProvider string @@ -61,11 +61,11 @@ func NewController( flaggerWindow time.Duration, logger *zap.SugaredLogger, notifier notifier.Interface, + canaryFactory *canary.Factory, routerFactory *router.Factory, observerFactory *metrics.Factory, meshProvider string, version string, - labels []string, ) *Controller { logger.Debug("Creating event broadcaster") flaggerscheme.AddToScheme(scheme.Scheme) @@ -76,19 +76,6 @@ func NewController( }) eventRecorder := eventBroadcaster.NewRecorder( scheme.Scheme, corev1.EventSource{Component: controllerAgentName}) - - deployer := canary.Deployer{ - Logger: logger, - KubeClient: kubeClient, - FlaggerClient: flaggerClient, - Labels: labels, - ConfigTracker: canary.ConfigTracker{ - Logger: logger, - KubeClient: kubeClient, - FlaggerClient: flaggerClient, - }, - } - recorder := metrics.NewRecorder(controllerAgentName, true) recorder.SetInfo(version, meshProvider) @@ -104,10 +91,10 @@ func NewController( canaries: new(sync.Map), jobs: map[string]CanaryJob{}, flaggerWindow: flaggerWindow, - deployer: deployer, observerFactory: observerFactory, recorder: recorder, notifier: notifier, + canaryFactory: canaryFactory, routerFactory: routerFactory, meshProvider: meshProvider, } @@ -218,7 +205,7 @@ func (c *Controller) syncHandler(key string) error { // set status condition for new canaries if cd.Status.Conditions == nil { - if ok, conditions := c.deployer.MakeStatusConditions(cd.Status, flaggerv1.CanaryPhaseInitializing); ok { + if ok, conditions := canary.MakeStatusConditions(cd.Status, flaggerv1.CanaryPhaseInitializing); ok { cdCopy := cd.DeepCopy() cdCopy.Status.Conditions = conditions cdCopy.Status.LastTransitionTime = metav1.Now() diff --git a/pkg/controller/controller_test.go b/pkg/controller/controller_test.go index 19c42435..df91b713 100644 --- a/pkg/controller/controller_test.go +++ b/pkg/controller/controller_test.go @@ -37,7 +37,7 @@ type Mocks struct { kubeClient kubernetes.Interface meshClient clientset.Interface flaggerClient clientset.Interface - deployer canary.Deployer + deployer canary.Controller ctrl *Controller logger *zap.SugaredLogger router router.Interface @@ -63,19 +63,6 @@ func SetupMocks(c *flaggerv1.Canary) Mocks { logger, _ := logger.NewLogger("debug") - // init controller helpers - deployer := canary.Deployer{ - Logger: logger, - KubeClient: kubeClient, - FlaggerClient: flaggerClient, - Labels: []string{"app", "name"}, - ConfigTracker: canary.ConfigTracker{ - Logger: logger, - KubeClient: kubeClient, - FlaggerClient: flaggerClient, - }, - } - // init controller flaggerInformerFactory := informers.NewSharedInformerFactory(flaggerClient, noResyncPeriodFunc()) flaggerInformer := flaggerInformerFactory.Flagger().V1alpha3().Canaries() @@ -86,6 +73,14 @@ func SetupMocks(c *flaggerv1.Canary) Mocks { // init observer observerFactory, _ := metrics.NewFactory("fake", "istio", 5*time.Second) + // init canary factory + configTracker := canary.ConfigTracker{ + Logger: logger, + KubeClient: kubeClient, + FlaggerClient: flaggerClient, + } + canaryFactory := canary.NewFactory(kubeClient, flaggerClient, configTracker, []string{"app", "name"}, logger) + ctrl := &Controller{ kubeClient: kubeClient, istioClient: flaggerClient, @@ -97,7 +92,7 @@ func SetupMocks(c *flaggerv1.Canary) Mocks { logger: logger, canaries: new(sync.Map), flaggerWindow: time.Second, - deployer: deployer, + canaryFactory: canaryFactory, observerFactory: observerFactory, recorder: metrics.NewRecorder(controllerAgentName, false), routerFactory: rf, @@ -108,7 +103,7 @@ func SetupMocks(c *flaggerv1.Canary) Mocks { return Mocks{ canary: c, - deployer: deployer, + deployer: canaryFactory.Controller("Deployment"), logger: logger, flaggerClient: flaggerClient, meshClient: flaggerClient, diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index fda8e418..489d166c 100644 --- a/pkg/controller/scheduler.go +++ b/pkg/controller/scheduler.go @@ -2,13 +2,14 @@ package controller import ( "fmt" - "github.com/weaveworks/flagger/pkg/metrics" "strings" "time" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1alpha3" + "github.com/weaveworks/flagger/pkg/canary" + "github.com/weaveworks/flagger/pkg/metrics" "github.com/weaveworks/flagger/pkg/router" ) @@ -97,13 +98,16 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh provider = cd.Spec.Provider } + // init controller based on target kind + canaryController := c.canaryFactory.Controller(cd.Spec.TargetRef.Kind) + // create primary deployment and hpa if needed // skip primary check for Istio since the deployment will become ready after the ClusterIP are created skipPrimaryCheck := false if skipLivenessChecks || strings.Contains(provider, "istio") || strings.Contains(provider, "appmesh") { skipPrimaryCheck = true } - labelSelector, ports, err := c.deployer.Initialize(cd, skipPrimaryCheck) + labelSelector, ports, err := canaryController.Initialize(cd, skipPrimaryCheck) if err != nil { c.recordEventWarningf(cd, "%v", err) return @@ -125,7 +129,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh } // check for deployment spec or configs changes - shouldAdvance, err := c.shouldAdvance(cd) + shouldAdvance, err := c.shouldAdvance(cd, canaryController) if err != nil { c.recordEventWarningf(cd, "%v", err) return @@ -137,7 +141,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh } // check gates - if isApproved := c.runConfirmRolloutHooks(cd); !isApproved { + if isApproved := c.runConfirmRolloutHooks(cd, canaryController); !isApproved { return } @@ -149,7 +153,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh // check primary deployment status if !skipLivenessChecks { - if _, err := c.deployer.IsPrimaryReady(cd); err != nil { + if _, err := canaryController.IsPrimaryReady(cd); err != nil { c.recordEventWarningf(cd, "%v", err) return } @@ -165,12 +169,12 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh c.recorder.SetWeight(cd, primaryWeight, canaryWeight) // check if canary analysis should start (canary revision has changes) or continue - if ok := c.checkCanaryStatus(cd, shouldAdvance); !ok { + if ok := c.checkCanaryStatus(cd, canaryController, shouldAdvance); !ok { return } // check if canary revision changed during analysis - if restart := c.hasCanaryRevisionChanged(cd); restart { + if restart := c.hasCanaryRevisionChanged(cd, canaryController); restart { c.recordEventInfof(cd, "New revision detected! Restarting analysis for %s.%s", cd.Spec.TargetRef.Name, cd.Namespace) @@ -189,7 +193,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh FailedChecks: 0, Iterations: 0, } - if err := c.deployer.SyncStatus(cd, status); err != nil { + if err := canaryController.SyncStatus(cd, status); err != nil { c.recordEventWarningf(cd, "%v", err) return } @@ -203,7 +207,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh // check canary deployment status var retriable = true if !skipLivenessChecks { - retriable, err = c.deployer.IsCanaryReady(cd) + retriable, err = canaryController.IsCanaryReady(cd) if err != nil && retriable { c.recordEventWarningf(cd, "%v", err) return @@ -211,7 +215,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh } // check if analysis should be skipped - if skip := c.shouldSkipAnalysis(cd, meshRouter, primaryWeight, canaryWeight); skip { + if skip := c.shouldSkipAnalysis(cd, canaryController, meshRouter, primaryWeight, canaryWeight); skip { return } @@ -227,7 +231,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh } // update status phase - if err := c.deployer.SetStatusPhase(cd, flaggerv1.CanaryPhaseFinalising); err != nil { + if err := canaryController.SetStatusPhase(cd, flaggerv1.CanaryPhaseFinalising); err != nil { c.recordEventWarningf(cd, "%v", err) return } @@ -237,13 +241,13 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh // scale canary to zero if promotion has finished if cd.Status.Phase == flaggerv1.CanaryPhaseFinalising { - if err := c.deployer.Scale(cd, 0); err != nil { + if err := canaryController.Scale(cd, 0); err != nil { c.recordEventWarningf(cd, "%v", err) return } // set status to succeeded - if err := c.deployer.SetStatusPhase(cd, flaggerv1.CanaryPhaseSucceeded); err != nil { + if err := canaryController.SetStatusPhase(cd, flaggerv1.CanaryPhaseSucceeded); err != nil { c.recordEventWarningf(cd, "%v", err) return } @@ -286,13 +290,13 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh cd.Name, cd.Namespace) // shutdown canary - if err := c.deployer.Scale(cd, 0); err != nil { + if err := canaryController.Scale(cd, 0); err != nil { c.recordEventWarningf(cd, "%v", err) return } // mark canary as failed - if err := c.deployer.SyncStatus(cd, flaggerv1.CanaryStatus{Phase: flaggerv1.CanaryPhaseFailed, CanaryWeight: 0}); err != nil { + if err := canaryController.SyncStatus(cd, flaggerv1.CanaryStatus{Phase: flaggerv1.CanaryPhaseFailed, CanaryWeight: 0}); err != nil { c.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Errorf("%v", err) return } @@ -310,7 +314,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh // run pre-rollout web hooks if ok := c.runPreRolloutHooks(cd); !ok { - if err := c.deployer.SetStatusFailedChecks(cd, cd.Status.FailedChecks+1); err != nil { + if err := canaryController.SetStatusFailedChecks(cd, cd.Status.FailedChecks+1); err != nil { c.recordEventWarningf(cd, "%v", err) return } @@ -318,7 +322,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh } } else { if ok := c.analyseCanary(cd); !ok { - if err := c.deployer.SetStatusFailedChecks(cd, cd.Status.FailedChecks+1); err != nil { + if err := canaryController.SetStatusFailedChecks(cd, cd.Status.FailedChecks+1); err != nil { c.recordEventWarningf(cd, "%v", err) return } @@ -349,7 +353,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh } c.recorder.SetWeight(cd, 0, 100) - if err := c.deployer.SetStatusIterations(cd, cd.Status.Iterations+1); err != nil { + if err := canaryController.SetStatusIterations(cd, cd.Status.Iterations+1); err != nil { c.recordEventWarningf(cd, "%v", err) return } @@ -367,13 +371,13 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh if cd.Spec.CanaryAnalysis.Iterations == cd.Status.Iterations { c.recordEventInfof(cd, "Copying %s.%s template spec to %s.%s", cd.Spec.TargetRef.Name, cd.Namespace, primaryName, cd.Namespace) - if err := c.deployer.Promote(cd); err != nil { + if err := canaryController.Promote(cd); err != nil { c.recordEventWarningf(cd, "%v", err) return } // update status phase - if err := c.deployer.SetStatusPhase(cd, flaggerv1.CanaryPhasePromoting); err != nil { + if err := canaryController.SetStatusPhase(cd, flaggerv1.CanaryPhasePromoting); err != nil { c.recordEventWarningf(cd, "%v", err) return } @@ -396,7 +400,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh c.logger.With("canary", fmt.Sprintf("%s.%s", name, namespace)). Infof("Start traffic mirroring") } - if err := c.deployer.SetStatusIterations(cd, cd.Status.Iterations+1); err != nil { + if err := canaryController.SetStatusIterations(cd, cd.Status.Iterations+1); err != nil { c.recordEventWarningf(cd, "%v", err) return } @@ -426,7 +430,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh } // increment iterations - if err := c.deployer.SetStatusIterations(cd, cd.Status.Iterations+1); err != nil { + if err := canaryController.SetStatusIterations(cd, cd.Status.Iterations+1); err != nil { c.recordEventWarningf(cd, "%v", err) return } @@ -437,13 +441,13 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh if cd.Spec.CanaryAnalysis.Iterations < cd.Status.Iterations { c.recordEventInfof(cd, "Copying %s.%s template spec to %s.%s", cd.Spec.TargetRef.Name, cd.Namespace, primaryName, cd.Namespace) - if err := c.deployer.Promote(cd); err != nil { + if err := canaryController.Promote(cd); err != nil { c.recordEventWarningf(cd, "%v", err) return } // update status phase - if err := c.deployer.SetStatusPhase(cd, flaggerv1.CanaryPhasePromoting); err != nil { + if err := canaryController.SetStatusPhase(cd, flaggerv1.CanaryPhasePromoting); err != nil { c.recordEventWarningf(cd, "%v", err) return } @@ -489,7 +493,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh return } - if err := c.deployer.SetStatusWeight(cd, canaryWeight); err != nil { + if err := canaryController.SetStatusWeight(cd, canaryWeight); err != nil { c.recordEventWarningf(cd, "%v", err) return } @@ -509,13 +513,13 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh // update primary spec c.recordEventInfof(cd, "Copying %s.%s template spec to %s.%s", cd.Spec.TargetRef.Name, cd.Namespace, primaryName, cd.Namespace) - if err := c.deployer.Promote(cd); err != nil { + if err := canaryController.Promote(cd); err != nil { c.recordEventWarningf(cd, "%v", err) return } // update status phase - if err := c.deployer.SetStatusPhase(cd, flaggerv1.CanaryPhasePromoting); err != nil { + if err := canaryController.SetStatusPhase(cd, flaggerv1.CanaryPhasePromoting); err != nil { c.recordEventWarningf(cd, "%v", err) return } @@ -527,7 +531,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh } -func (c *Controller) shouldSkipAnalysis(cd *flaggerv1.Canary, meshRouter router.Interface, primaryWeight int, canaryWeight int) bool { +func (c *Controller) shouldSkipAnalysis(cd *flaggerv1.Canary, canaryController canary.Controller, meshRouter router.Interface, primaryWeight int, canaryWeight int) bool { if !cd.Spec.SkipAnalysis { return false } @@ -544,19 +548,19 @@ func (c *Controller) shouldSkipAnalysis(cd *flaggerv1.Canary, meshRouter router. // copy spec and configs from canary to primary c.recordEventInfof(cd, "Copying %s.%s template spec to %s-primary.%s", cd.Spec.TargetRef.Name, cd.Namespace, cd.Spec.TargetRef.Name, cd.Namespace) - if err := c.deployer.Promote(cd); err != nil { + if err := canaryController.Promote(cd); err != nil { c.recordEventWarningf(cd, "%v", err) return false } // shutdown canary - if err := c.deployer.Scale(cd, 0); err != nil { + if err := canaryController.Scale(cd, 0); err != nil { c.recordEventWarningf(cd, "%v", err) return false } // update status phase - if err := c.deployer.SetStatusPhase(cd, flaggerv1.CanaryPhaseSucceeded); err != nil { + if err := canaryController.SetStatusPhase(cd, flaggerv1.CanaryPhaseSucceeded); err != nil { c.recordEventWarningf(cd, "%v", err) return false } @@ -571,7 +575,7 @@ func (c *Controller) shouldSkipAnalysis(cd *flaggerv1.Canary, meshRouter router. return true } -func (c *Controller) shouldAdvance(cd *flaggerv1.Canary) (bool, error) { +func (c *Controller) shouldAdvance(cd *flaggerv1.Canary, canaryController canary.Controller) (bool, error) { if cd.Status.LastAppliedSpec == "" || cd.Status.Phase == flaggerv1.CanaryPhaseInitializing || cd.Status.Phase == flaggerv1.CanaryPhaseProgressing || @@ -581,7 +585,7 @@ func (c *Controller) shouldAdvance(cd *flaggerv1.Canary) (bool, error) { return true, nil } - newDep, err := c.deployer.HasDeploymentChanged(cd) + newDep, err := canaryController.HasTargetChanged(cd) if err != nil { return false, err } @@ -589,7 +593,7 @@ func (c *Controller) shouldAdvance(cd *flaggerv1.Canary) (bool, error) { return newDep, nil } - newCfg, err := c.deployer.ConfigTracker.HasConfigChanged(cd) + newCfg, err := canaryController.HaveDependenciesChanged(cd) if err != nil { return false, err } @@ -598,7 +602,7 @@ func (c *Controller) shouldAdvance(cd *flaggerv1.Canary) (bool, error) { } -func (c *Controller) checkCanaryStatus(cd *flaggerv1.Canary, shouldAdvance bool) bool { +func (c *Controller) checkCanaryStatus(cd *flaggerv1.Canary, canaryController canary.Controller, shouldAdvance bool) bool { c.recorder.SetStatus(cd, cd.Status.Phase) if cd.Status.Phase == flaggerv1.CanaryPhaseProgressing || cd.Status.Phase == flaggerv1.CanaryPhasePromoting || @@ -607,7 +611,7 @@ func (c *Controller) checkCanaryStatus(cd *flaggerv1.Canary, shouldAdvance bool) } if cd.Status.Phase == "" || cd.Status.Phase == flaggerv1.CanaryPhaseInitializing { - if err := c.deployer.SyncStatus(cd, flaggerv1.CanaryStatus{Phase: flaggerv1.CanaryPhaseInitialized}); err != nil { + if err := canaryController.SyncStatus(cd, flaggerv1.CanaryStatus{Phase: flaggerv1.CanaryPhaseInitialized}); err != nil { c.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Errorf("%v", err) return false } @@ -622,11 +626,11 @@ func (c *Controller) checkCanaryStatus(cd *flaggerv1.Canary, shouldAdvance bool) c.recordEventInfof(cd, "New revision detected! Scaling up %s.%s", cd.Spec.TargetRef.Name, cd.Namespace) c.sendNotification(cd, "New revision detected, starting canary analysis.", true, false) - if err := c.deployer.ScaleUp(cd); err != nil { + if err := canaryController.ScaleFromZero(cd); err != nil { c.recordEventErrorf(cd, "%v", err) return false } - if err := c.deployer.SyncStatus(cd, flaggerv1.CanaryStatus{Phase: flaggerv1.CanaryPhaseProgressing}); err != nil { + if err := canaryController.SyncStatus(cd, flaggerv1.CanaryStatus{Phase: flaggerv1.CanaryPhaseProgressing}); err != nil { c.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Errorf("%v", err) return false } @@ -636,25 +640,25 @@ func (c *Controller) checkCanaryStatus(cd *flaggerv1.Canary, shouldAdvance bool) return false } -func (c *Controller) hasCanaryRevisionChanged(cd *flaggerv1.Canary) bool { +func (c *Controller) hasCanaryRevisionChanged(cd *flaggerv1.Canary, canaryController canary.Controller) bool { if cd.Status.Phase == flaggerv1.CanaryPhaseProgressing { - if diff, _ := c.deployer.HasDeploymentChanged(cd); diff { + if diff, _ := canaryController.HasTargetChanged(cd); diff { return true } - if diff, _ := c.deployer.ConfigTracker.HasConfigChanged(cd); diff { + if diff, _ := canaryController.HaveDependenciesChanged(cd); diff { return true } } return false } -func (c *Controller) runConfirmRolloutHooks(canary *flaggerv1.Canary) bool { +func (c *Controller) runConfirmRolloutHooks(canary *flaggerv1.Canary, canaryController canary.Controller) bool { for _, webhook := range canary.Spec.CanaryAnalysis.Webhooks { if webhook.Type == flaggerv1.ConfirmRolloutHook { err := CallWebhook(canary.Name, canary.Namespace, flaggerv1.CanaryPhaseProgressing, webhook) if err != nil { if canary.Status.Phase != flaggerv1.CanaryPhaseWaiting { - if err := c.deployer.SetStatusPhase(canary, flaggerv1.CanaryPhaseWaiting); err != nil { + if err := canaryController.SetStatusPhase(canary, flaggerv1.CanaryPhaseWaiting); err != nil { c.logger.With("canary", fmt.Sprintf("%s.%s", canary.Name, canary.Namespace)).Errorf("%v", err) } c.recordEventWarningf(canary, "Halt %s.%s advancement waiting for approval %s", @@ -664,7 +668,7 @@ func (c *Controller) runConfirmRolloutHooks(canary *flaggerv1.Canary) bool { return false } else { if canary.Status.Phase == flaggerv1.CanaryPhaseWaiting { - if err := c.deployer.SetStatusPhase(canary, flaggerv1.CanaryPhaseProgressing); err != nil { + if err := canaryController.SetStatusPhase(canary, flaggerv1.CanaryPhaseProgressing); err != nil { c.logger.With("canary", fmt.Sprintf("%s.%s", canary.Name, canary.Namespace)).Errorf("%v", err) return false }