From 5b296e01b3ad75de02457071c56097d4df959805 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sat, 26 Jan 2019 12:36:27 +0200 Subject: [PATCH] Detect changes in configs and trigger canary analysis - restart analysis if a ConfigMap or Secret changes during rollout - add tests for tracked changes --- pkg/controller/controller_test.go | 71 +++++++++++++------------------ pkg/controller/deployer.go | 23 +++++++++- pkg/controller/deployer_test.go | 9 ++++ pkg/controller/scheduler.go | 22 ++++++---- pkg/controller/scheduler_test.go | 29 +++++++------ pkg/controller/tracker.go | 48 +++++++++++++++++++++ 6 files changed, 138 insertions(+), 64 deletions(-) diff --git a/pkg/controller/controller_test.go b/pkg/controller/controller_test.go index 1129955f..cbb6b197 100644 --- a/pkg/controller/controller_test.go +++ b/pkg/controller/controller_test.go @@ -40,28 +40,26 @@ type Mocks struct { } func SetupMocks() Mocks { + // init canary canary := newTestCanary() - configMap := NewTestConfigMap() - configMapEnv := NewTestConfigMapEnv() - configMapVol := NewTestConfigMapVol() - secret := NewTestSecret() - secretEnv := NewTestSecretEnv() - secretVol := NewTestSecretVol() - dep := newTestDeployment() - hpa := newTestHPA() - - kubeClient := fake.NewSimpleClientset(secret, secretEnv, secretVol, configMap, configMapEnv, configMapVol, dep, hpa) - - istioClient := fakeIstio.NewSimpleClientset() - flaggerClient := fakeFlagger.NewSimpleClientset(canary) + // init kube clientset and register mock objects + kubeClient := fake.NewSimpleClientset( + newTestDeployment(), + newTestHPA(), + NewTestConfigMap(), + NewTestConfigMapEnv(), + NewTestConfigMapVol(), + NewTestSecret(), + NewTestSecretEnv(), + NewTestSecretVol(), + ) + + istioClient := fakeIstio.NewSimpleClientset() logger, _ := logging.NewLogger("debug") - observer := CanaryObserver{ - metricsServer: "fake", - } - + // init controller helpers deployer := CanaryDeployer{ flaggerClient: flaggerClient, kubeClient: kubeClient, @@ -72,38 +70,17 @@ func SetupMocks() Mocks { flaggerClient: flaggerClient, }, } - router := CanaryRouter{ flaggerClient: flaggerClient, kubeClient: kubeClient, istioClient: istioClient, logger: logger, } - - controller := newTestController(kubeClient, istioClient, flaggerClient, logger, deployer, router, observer) - - return Mocks{ - canary: canary, - observer: observer, - router: router, - deployer: deployer, - logger: logger, - flaggerClient: flaggerClient, - istioClient: istioClient, - kubeClient: kubeClient, - ctrl: controller, + observer := CanaryObserver{ + metricsServer: "fake", } -} -func newTestController( - kubeClient kubernetes.Interface, - istioClient istioclientset.Interface, - flaggerClient clientset.Interface, - logger *zap.SugaredLogger, - deployer CanaryDeployer, - router CanaryRouter, - observer CanaryObserver, -) *Controller { + // init controller flaggerInformerFactory := informers.NewSharedInformerFactory(flaggerClient, noResyncPeriodFunc()) flaggerInformer := flaggerInformerFactory.Flagger().V1alpha3().Canaries() @@ -125,7 +102,17 @@ func newTestController( } ctrl.flaggerSynced = alwaysReady - return ctrl + return Mocks{ + canary: canary, + observer: observer, + router: router, + deployer: deployer, + logger: logger, + flaggerClient: flaggerClient, + istioClient: istioClient, + kubeClient: kubeClient, + ctrl: ctrl, + } } func NewTestConfigMap() *corev1.ConfigMap { diff --git a/pkg/controller/deployer.go b/pkg/controller/deployer.go index 0019fd8b..7d177445 100644 --- a/pkg/controller/deployer.go +++ b/pkg/controller/deployer.go @@ -169,7 +169,22 @@ func (c *CanaryDeployer) ShouldAdvance(cd *flaggerv1.Canary) (bool, error) { if cd.Status.LastAppliedSpec == "" || cd.Status.Phase == flaggerv1.CanaryProgressing { return true, nil } - return c.IsNewSpec(cd) + + newDep, err := c.IsNewSpec(cd) + if err != nil { + return false, err + } + if newDep { + return newDep, nil + } + + newCfg, err := c.configTracker.HasConfigChanged(cd) + if err != nil { + return false, err + } + + return newCfg, nil + } // SetStatusFailedChecks updates the canary failed checks counter @@ -230,12 +245,18 @@ func (c *CanaryDeployer) SyncStatus(cd *flaggerv1.Canary, status flaggerv1.Canar return fmt.Errorf("deployment %s.%s marshal error %v", cd.Spec.TargetRef.Name, cd.Namespace, err) } + configs, err := c.configTracker.GetConfigRefs(cd) + if err != nil { + return fmt.Errorf("configs query error %v", err) + } + cdCopy := cd.DeepCopy() cdCopy.Status.Phase = status.Phase cdCopy.Status.CanaryWeight = status.CanaryWeight cdCopy.Status.FailedChecks = status.FailedChecks cdCopy.Status.LastAppliedSpec = base64.StdEncoding.EncodeToString(specJson) cdCopy.Status.LastTransitionTime = metav1.Now() + cdCopy.Status.TrackedConfigs = configs cd, err = c.flaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) if err != nil { diff --git a/pkg/controller/deployer_test.go b/pkg/controller/deployer_test.go index 08a62833..e7bdd0ba 100644 --- a/pkg/controller/deployer_test.go +++ b/pkg/controller/deployer_test.go @@ -251,6 +251,15 @@ func TestCanaryDeployer_SyncStatus(t *testing.T) { if res.Status.FailedChecks != status.FailedChecks { t.Errorf("Got failed checks %v wanted %v", res.Status.FailedChecks, status.FailedChecks) } + + if res.Status.TrackedConfigs == nil { + t.Fatalf("Status tracking configs are empty") + } + configs := *res.Status.TrackedConfigs + secret := NewTestSecret() + if _, exists := configs["secret/"+secret.GetName()]; !exists { + t.Errorf("Secret %s not found in status", secret.GetName()) + } } func TestCanaryDeployer_Scale(t *testing.T) { diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index ece90fb4..757af39e 100644 --- a/pkg/controller/scheduler.go +++ b/pkg/controller/scheduler.go @@ -91,10 +91,13 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh return } - if ok, err := c.deployer.ShouldAdvance(cd); !ok { - if err != nil { - c.recordEventWarningf(cd, "%v", err) - } + shouldAdvance, err := c.deployer.ShouldAdvance(cd) + if err != nil { + c.recordEventWarningf(cd, "%v", err) + return + } + + if !shouldAdvance { return } @@ -123,7 +126,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh 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); !ok { + if ok := c.checkCanaryStatus(cd, shouldAdvance); !ok { return } @@ -291,7 +294,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh } } -func (c *Controller) checkCanaryStatus(cd *flaggerv1.Canary) bool { +func (c *Controller) checkCanaryStatus(cd *flaggerv1.Canary, shouldAdvance bool) bool { c.recorder.SetStatus(cd) if cd.Status.Phase == flaggerv1.CanaryProgressing { return true @@ -309,11 +312,11 @@ func (c *Controller) checkCanaryStatus(cd *flaggerv1.Canary) bool { return false } - if diff, err := c.deployer.IsNewSpec(cd); diff { + if shouldAdvance { 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.Scale(cd, 1); err != nil { + if err := c.deployer.Scale(cd, 1); err != nil { c.recordEventErrorf(cd, "%v", err) return false } @@ -332,6 +335,9 @@ func (c *Controller) hasCanaryRevisionChanged(cd *flaggerv1.Canary) bool { if diff, _ := c.deployer.IsNewSpec(cd); diff { return true } + if diff, _ := c.deployer.configTracker.HasConfigChanged(cd); diff { + return true + } } return false } diff --git a/pkg/controller/scheduler_test.go b/pkg/controller/scheduler_test.go index e93b7e1e..900ddebf 100644 --- a/pkg/controller/scheduler_test.go +++ b/pkg/controller/scheduler_test.go @@ -130,7 +130,22 @@ func TestScheduler_Promotion(t *testing.T) { t.Fatal(err.Error()) } - // detect changes + // detect pod spec changes + mocks.ctrl.advanceCanary("podinfo", "default", true) + + config2 := NewTestConfigMapV2() + _, err = mocks.kubeClient.CoreV1().ConfigMaps("default").Update(config2) + if err != nil { + t.Fatal(err.Error()) + } + + secret2 := NewTestSecretV2() + _, err = mocks.kubeClient.CoreV1().Secrets("default").Update(secret2) + if err != nil { + t.Fatal(err.Error()) + } + + // detect configs changes mocks.ctrl.advanceCanary("podinfo", "default", true) primaryRoute, canaryRoute, err := mocks.router.GetRoutes(mocks.canary) @@ -145,18 +160,6 @@ func TestScheduler_Promotion(t *testing.T) { t.Fatal(err.Error()) } - config2 := NewTestConfigMapV2() - _, err = mocks.kubeClient.CoreV1().ConfigMaps("default").Update(config2) - if err != nil { - t.Fatal(err.Error()) - } - - secret2 := NewTestSecretV2() - _, err = mocks.kubeClient.CoreV1().Secrets("default").Update(secret2) - if err != nil { - t.Fatal(err.Error()) - } - // advance mocks.ctrl.advanceCanary("podinfo", "default", true) diff --git a/pkg/controller/tracker.go b/pkg/controller/tracker.go index f7eb767a..812b8a5a 100644 --- a/pkg/controller/tracker.go +++ b/pkg/controller/tracker.go @@ -182,6 +182,54 @@ func (ct *ConfigTracker) GetTargetConfigs(cd *flaggerv1.Canary) (map[string]Conf return res, nil } +// GetConfigRefs returns a map of configs and their checksum +func (ct *ConfigTracker) GetConfigRefs(cd *flaggerv1.Canary) (*map[string]string, error) { + res := make(map[string]string) + configs, err := ct.GetTargetConfigs(cd) + if err != nil { + return nil, err + } + + for _, cfg := range configs { + res[cfg.GetName()] = cfg.Checksum + } + + return &res, nil +} + +// HasConfigChanged checks for changes in ConfigMaps and Secretes by comparing +// the checksum for each ConfigRef stored in Canary.Status.TrackedConfigs +func (ct *ConfigTracker) HasConfigChanged(cd *flaggerv1.Canary) (bool, error) { + configs, err := ct.GetTargetConfigs(cd) + if err != nil { + return false, err + } + + if len(configs) == 0 && cd.Status.TrackedConfigs == nil { + return false, nil + } + + if len(configs) > 0 && cd.Status.TrackedConfigs == nil { + return true, nil + } + + trackedConfigs := *cd.Status.TrackedConfigs + + if len(configs) != len(trackedConfigs) { + return true, nil + } + + for _, cfg := range configs { + if trackedConfigs[cfg.GetName()] != cfg.Checksum { + ct.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)). + Infof("%s %s has changed", cfg.Type, cfg.Name) + return true, nil + } + } + + return false, nil +} + // CreatePrimaryConfigs syncs the primary Kubernetes ConfigMaps and Secretes // with those found in the target deployment func (ct *ConfigTracker) CreatePrimaryConfigs(cd *flaggerv1.Canary, refs map[string]ConfigRef) error {