From aff8b117d4fc4213bb7309aa62d6218c4a28234c Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Thu, 17 Jan 2019 15:13:59 +0200 Subject: [PATCH] Restart validation if revision changes during analysis --- pkg/controller/scheduler.go | 52 ++++++++++++++++++++++++++++++------- 1 file changed, 43 insertions(+), 9 deletions(-) diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index f6975df3..40c7d38b 100644 --- a/pkg/controller/scheduler.go +++ b/pkg/controller/scheduler.go @@ -58,8 +58,7 @@ func (c *Controller) scheduleCanaries() { for canaryName, targetName := range current { for name, target := range current { if name != canaryName && target == targetName { - c.logger.Errorf("Bad things will happen! Found more than one canary with the same target %s", - targetName) + c.logger.With("canary", canaryName).Errorf("Bad things will happen! Found more than one canary with the same target %s", targetName) } } } @@ -121,7 +120,33 @@ func (c *Controller) advanceCanary(name string, namespace string) { c.recorder.SetWeight(cd, primaryRoute.Weight, canaryRoute.Weight) // check if canary analysis should start (canary revision has changes) or continue - if ok := c.checkCanaryStatus(cd, c.deployer); !ok { + if ok := c.checkCanaryStatus(cd); !ok { + return + } + + // check if canary revision changed during analysis + if restart := c.hasCanaryRevisionChanged(cd); restart { + c.recordEventInfof(cd, "New revision detected! Restarting analysis for %s.%s", + cd.Spec.TargetRef.Name, cd.Namespace) + + // route all traffic back to primary + primaryRoute.Weight = 100 + canaryRoute.Weight = 0 + if err := c.router.SetRoutes(cd, primaryRoute, canaryRoute); err != nil { + c.recordEventWarningf(cd, "%v", err) + return + } + + // reset status + status := flaggerv1.CanaryStatus{ + Phase: flaggerv1.CanaryProgressing, + CanaryWeight: 0, + FailedChecks: 0, + } + if err := c.deployer.SyncStatus(cd, status); err != nil { + c.recordEventWarningf(cd, "%v", err) + return + } return } @@ -185,7 +210,7 @@ func (c *Controller) advanceCanary(name string, namespace string) { // check if the canary success rate is above the threshold // skip check if no traffic is routed to canary if canaryRoute.Weight == 0 { - c.recordEventInfof(cd, "Starting canary deployment for %s.%s", cd.Name, cd.Namespace) + c.recordEventInfof(cd, "Starting canary analysis for %s.%s", cd.Spec.TargetRef.Name, cd.Namespace) } else { if ok := c.analyseCanary(cd); !ok { if err := c.deployer.SetStatusFailedChecks(cd, cd.Status.FailedChecks+1); err != nil { @@ -260,14 +285,14 @@ func (c *Controller) advanceCanary(name string, namespace string) { } } -func (c *Controller) checkCanaryStatus(cd *flaggerv1.Canary, deployer CanaryDeployer) bool { +func (c *Controller) checkCanaryStatus(cd *flaggerv1.Canary) bool { c.recorder.SetStatus(cd) if cd.Status.Phase == flaggerv1.CanaryProgressing { return true } if cd.Status.Phase == "" { - if err := deployer.SyncStatus(cd, flaggerv1.CanaryStatus{Phase: flaggerv1.CanaryInitialized}); err != nil { + if err := c.deployer.SyncStatus(cd, flaggerv1.CanaryStatus{Phase: flaggerv1.CanaryInitialized}); err != nil { c.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Errorf("%v", err) return false } @@ -278,15 +303,15 @@ func (c *Controller) checkCanaryStatus(cd *flaggerv1.Canary, deployer CanaryDepl return false } - if diff, err := deployer.IsNewSpec(cd); diff { + if diff, err := c.deployer.IsNewSpec(cd); diff { 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 = deployer.Scale(cd, 1); err != nil { + if err = c.deployer.Scale(cd, 1); err != nil { c.recordEventErrorf(cd, "%v", err) return false } - if err := deployer.SyncStatus(cd, flaggerv1.CanaryStatus{Phase: flaggerv1.CanaryProgressing}); err != nil { + if err := c.deployer.SyncStatus(cd, flaggerv1.CanaryStatus{Phase: flaggerv1.CanaryProgressing}); err != nil { c.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Errorf("%v", err) return false } @@ -296,6 +321,15 @@ func (c *Controller) checkCanaryStatus(cd *flaggerv1.Canary, deployer CanaryDepl return false } +func (c *Controller) hasCanaryRevisionChanged(cd *flaggerv1.Canary) bool { + if cd.Status.Phase == flaggerv1.CanaryProgressing { + if diff, _ := c.deployer.IsNewSpec(cd); diff { + return true + } + } + return false +} + func (c *Controller) analyseCanary(r *flaggerv1.Canary) bool { // run metrics checks for _, metric := range r.Spec.CanaryAnalysis.Metrics {