diff --git a/pkg/canary/config_tracker.go b/pkg/canary/config_tracker.go index e3922fcb..16ad9333 100644 --- a/pkg/canary/config_tracker.go +++ b/pkg/canary/config_tracker.go @@ -54,7 +54,7 @@ func checksum(data interface{}) string { func (ct *ConfigTracker) getRefFromConfigMap(name string, namespace string) (*ConfigRef, error) { config, err := ct.KubeClient.CoreV1().ConfigMaps(namespace).Get(name, metav1.GetOptions{}) if err != nil { - return nil, err + return nil, fmt.Errorf("configmap %s.%s get query error: %w", name, namespace, err) } return &ConfigRef{ @@ -69,7 +69,7 @@ func (ct *ConfigTracker) getRefFromConfigMap(name string, namespace string) (*Co func (ct *ConfigTracker) getRefFromSecret(name string, namespace string) (*ConfigRef, error) { secret, err := ct.KubeClient.CoreV1().Secrets(namespace).Get(name, metav1.GetOptions{}) if err != nil { - return nil, err + return nil, fmt.Errorf("secret %s.%s get query error: %w", name, namespace, err) } // ignore registry secrets (those should be set via service account) @@ -100,20 +100,14 @@ func (ct *ConfigTracker) GetTargetConfigs(cd *flaggerv1.Canary) (map[string]Conf case "Deployment": targetDep, err := ct.KubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { - if errors.IsNotFound(err) { - return res, fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) - } - return res, fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) + return res, fmt.Errorf("deployment %s.%s get query error: %w", targetName, cd.Namespace, err) } vs = targetDep.Spec.Template.Spec.Volumes cs = targetDep.Spec.Template.Spec.Containers case "DaemonSet": targetDae, err := ct.KubeClient.AppsV1().DaemonSets(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { - if errors.IsNotFound(err) { - return res, fmt.Errorf("daemonset %s.%s not found", targetName, cd.Namespace) - } - return res, fmt.Errorf("daemonset %s.%s query error %v", targetName, cd.Namespace, err) + return res, fmt.Errorf("daemonset %s.%s get query error: %w", targetName, cd.Namespace, err) } vs = targetDae.Spec.Template.Spec.Volumes cs = targetDae.Spec.Template.Spec.Containers @@ -126,18 +120,16 @@ func (ct *ConfigTracker) GetTargetConfigs(cd *flaggerv1.Canary) (map[string]Conf if cmv := volume.ConfigMap; cmv != nil { config, err := ct.getRefFromConfigMap(cmv.Name, cd.Namespace) if err != nil { - ct.Logger.Errorf("configMap %s.%s query error %v", cmv.Name, cd.Namespace, err) + ct.Logger.Errorf("getRefFromConfigMap failed: %v", err) continue } - if config != nil { - res[config.GetName()] = *config - } + res[config.GetName()] = *config } if sv := volume.Secret; sv != nil { secret, err := ct.getRefFromSecret(sv.SecretName, cd.Namespace) if err != nil { - ct.Logger.Errorf("secret %s.%s query error %v", sv.SecretName, cd.Namespace, err) + ct.Logger.Errorf("getRefFromSecret failed: %v", err) continue } if secret != nil { @@ -150,18 +142,16 @@ func (ct *ConfigTracker) GetTargetConfigs(cd *flaggerv1.Canary) (map[string]Conf if cmv := source.ConfigMap; cmv != nil { config, err := ct.getRefFromConfigMap(cmv.Name, cd.Namespace) if err != nil { - ct.Logger.Errorf("configMap %s.%s query error %v", cmv.Name, cd.Namespace, err) + ct.Logger.Errorf("getRefFromConfigMap failed: %v", err) continue } - if config != nil { - res[config.GetName()] = *config - } + res[config.GetName()] = *config } if sv := source.Secret; sv != nil { secret, err := ct.getRefFromSecret(sv.Name, cd.Namespace) if err != nil { - ct.Logger.Errorf("secret %s.%s query error %v", sv.Name, cd.Namespace, err) + ct.Logger.Errorf("getRefFromSecret failed: %v", err) continue } if secret != nil { @@ -181,17 +171,15 @@ func (ct *ConfigTracker) GetTargetConfigs(cd *flaggerv1.Canary) (map[string]Conf name := env.ValueFrom.ConfigMapKeyRef.LocalObjectReference.Name config, err := ct.getRefFromConfigMap(name, cd.Namespace) if err != nil { - ct.Logger.Errorf("configMap %s.%s query error %v", name, cd.Namespace, err) + ct.Logger.Errorf("getRefFromConfigMap failed: %v", err) continue } - if config != nil { - res[config.GetName()] = *config - } + res[config.GetName()] = *config case env.ValueFrom.SecretKeyRef != nil: name := env.ValueFrom.SecretKeyRef.LocalObjectReference.Name secret, err := ct.getRefFromSecret(name, cd.Namespace) if err != nil { - ct.Logger.Errorf("secret %s.%s query error %v", name, cd.Namespace, err) + ct.Logger.Errorf("getRefFromSecret failed: %v", err) continue } if secret != nil { @@ -207,17 +195,15 @@ func (ct *ConfigTracker) GetTargetConfigs(cd *flaggerv1.Canary) (map[string]Conf name := envFrom.ConfigMapRef.LocalObjectReference.Name config, err := ct.getRefFromConfigMap(name, cd.Namespace) if err != nil { - ct.Logger.Errorf("configMap %s.%s query error %v", name, cd.Namespace, err) + ct.Logger.Errorf("getRefFromConfigMap failed %v", err) continue } - if config != nil { - res[config.GetName()] = *config - } + res[config.GetName()] = *config case envFrom.SecretRef != nil: name := envFrom.SecretRef.LocalObjectReference.Name secret, err := ct.getRefFromSecret(name, cd.Namespace) if err != nil { - ct.Logger.Errorf("secret %s.%s query error %v", name, cd.Namespace, err) + ct.Logger.Errorf("getRefFromSecret failed %v", err) continue } if secret != nil { @@ -235,7 +221,7 @@ func (ct *ConfigTracker) GetConfigRefs(cd *flaggerv1.Canary) (*map[string]string res := make(map[string]string) configs, err := ct.GetTargetConfigs(cd) if err != nil { - return nil, err + return nil, fmt.Errorf("GetTargetConfigs failed: %w", err) } for _, cfg := range configs { @@ -250,7 +236,7 @@ func (ct *ConfigTracker) GetConfigRefs(cd *flaggerv1.Canary) (*map[string]string func (ct *ConfigTracker) HasConfigChanged(cd *flaggerv1.Canary) (bool, error) { configs, err := ct.GetTargetConfigs(cd) if err != nil { - return false, err + return false, fmt.Errorf("GetTargetConfigs failed: %w", err) } if len(configs) == 0 && cd.Status.TrackedConfigs == nil { @@ -286,7 +272,7 @@ func (ct *ConfigTracker) CreatePrimaryConfigs(cd *flaggerv1.Canary, refs map[str case ConfigRefMap: config, err := ct.KubeClient.CoreV1().ConfigMaps(cd.Namespace).Get(ref.Name, metav1.GetOptions{}) if err != nil { - return err + return fmt.Errorf("configmap %s.%s query failed : %w", ref.Name, cd.Name, err) } primaryName := fmt.Sprintf("%s-primary", config.GetName()) primaryConfigMap := &corev1.ConfigMap{ @@ -311,10 +297,10 @@ func (ct *ConfigTracker) CreatePrimaryConfigs(cd *flaggerv1.Canary, refs map[str if errors.IsNotFound(err) { _, err = ct.KubeClient.CoreV1().ConfigMaps(cd.Namespace).Create(primaryConfigMap) if err != nil { - return err + return fmt.Errorf("creating configmap %s.%s failed: %w", primaryConfigMap.Name, cd.Namespace, err) } } else { - return err + return fmt.Errorf("updating configmap %s.%s failed: %w", primaryConfigMap.Name, cd.Namespace, err) } } @@ -349,10 +335,10 @@ func (ct *ConfigTracker) CreatePrimaryConfigs(cd *flaggerv1.Canary, refs map[str if errors.IsNotFound(err) { _, err = ct.KubeClient.CoreV1().Secrets(cd.Namespace).Create(primarySecret) if err != nil { - return err + return fmt.Errorf("creating secret %s.%s failed: %w", primarySecret.Name, cd.Namespace, err) } } else { - return err + return fmt.Errorf("updating secret %s.%s failed: %w", primarySecret.Name, cd.Namespace, err) } } diff --git a/pkg/canary/controller.go b/pkg/canary/controller.go index c615cc5e..d54c1a20 100644 --- a/pkg/canary/controller.go +++ b/pkg/canary/controller.go @@ -5,7 +5,7 @@ import ( ) type Controller interface { - IsPrimaryReady(canary *flaggerv1.Canary) (bool, error) + IsPrimaryReady(canary *flaggerv1.Canary) error IsCanaryReady(canary *flaggerv1.Canary) (bool, error) GetMetadata(canary *flaggerv1.Canary) (string, map[string]int32, error) SyncStatus(canary *flaggerv1.Canary, status flaggerv1.CanaryStatus) error @@ -17,6 +17,6 @@ type Controller interface { Promote(canary *flaggerv1.Canary) error HasTargetChanged(canary *flaggerv1.Canary) (bool, error) HaveDependenciesChanged(canary *flaggerv1.Canary) (bool, error) - Scale(canary *flaggerv1.Canary, replicas int32) error + ScaleToZero(canary *flaggerv1.Canary) error ScaleFromZero(canary *flaggerv1.Canary) error } diff --git a/pkg/canary/daemonset_controller.go b/pkg/canary/daemonset_controller.go index 0041247b..074d343a 100644 --- a/pkg/canary/daemonset_controller.go +++ b/pkg/canary/daemonset_controller.go @@ -28,33 +28,27 @@ type DaemonSetController struct { labels []string } -func (c *DaemonSetController) Scale(cd *flaggerv1.Canary, v int32) error { - // there's no concept `replicas` for DaemonSet - if v == 0 { - targetName := cd.Spec.TargetRef.Name - dae, err := c.kubeClient.AppsV1().DaemonSets(cd.Namespace).Get(targetName, metav1.GetOptions{}) - if err != nil { - if errors.IsNotFound(err) { - return fmt.Errorf("daemonset %s.%s not found", targetName, cd.Namespace) - } - return fmt.Errorf("daemonset %s.%s query error %v", targetName, cd.Namespace, err) - } +func (c *DaemonSetController) ScaleToZero(cd *flaggerv1.Canary) error { + targetName := cd.Spec.TargetRef.Name + dae, err := c.kubeClient.AppsV1().DaemonSets(cd.Namespace).Get(targetName, metav1.GetOptions{}) + if err != nil { + return fmt.Errorf("daemonset %s.%s query error: %w", targetName, cd.Namespace, err) + } - daeCopy := dae.DeepCopy() - daeCopy.Spec.Template.Spec.NodeSelector = make(map[string]string, - len(dae.Spec.Template.Spec.NodeSelector)+len(daemonSetScaleDownNodeSelector)) - for k, v := range dae.Spec.Template.Spec.NodeSelector { - daeCopy.Spec.Template.Spec.NodeSelector[k] = v - } + daeCopy := dae.DeepCopy() + daeCopy.Spec.Template.Spec.NodeSelector = make(map[string]string, + len(dae.Spec.Template.Spec.NodeSelector)+len(daemonSetScaleDownNodeSelector)) + for k, v := range dae.Spec.Template.Spec.NodeSelector { + daeCopy.Spec.Template.Spec.NodeSelector[k] = v + } - for k, v := range daemonSetScaleDownNodeSelector { - daeCopy.Spec.Template.Spec.NodeSelector[k] = v - } + for k, v := range daemonSetScaleDownNodeSelector { + daeCopy.Spec.Template.Spec.NodeSelector[k] = v + } - _, err = c.kubeClient.AppsV1().DaemonSets(dae.Namespace).Update(daeCopy) - if err != nil { - return fmt.Errorf("scaling down daemonset %s.%s failed: %v", daeCopy.GetName(), daeCopy.Namespace, err) - } + _, err = c.kubeClient.AppsV1().DaemonSets(dae.Namespace).Update(daeCopy) + if err != nil { + return fmt.Errorf("updating daemonset %s.%s failed: %w", daeCopy.GetName(), daeCopy.Namespace, err) } return nil } @@ -63,10 +57,7 @@ func (c *DaemonSetController) ScaleFromZero(cd *flaggerv1.Canary) error { targetName := cd.Spec.TargetRef.Name dep, err := c.kubeClient.AppsV1().DaemonSets(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { - if errors.IsNotFound(err) { - return fmt.Errorf("daemonset %s.%s not found", targetName, cd.Namespace) - } - return fmt.Errorf("daemonset %s.%s query error %v", targetName, cd.Namespace, err) + return fmt.Errorf("daemonset %s.%s query error: %w", targetName, cd.Namespace, err) } depCopy := dep.DeepCopy() @@ -76,7 +67,7 @@ func (c *DaemonSetController) ScaleFromZero(cd *flaggerv1.Canary) error { _, err = c.kubeClient.AppsV1().DaemonSets(dep.Namespace).Update(depCopy) if err != nil { - return fmt.Errorf("scaling up daemonset %s.%s failed: %v", depCopy.GetName(), depCopy.Namespace, err) + return fmt.Errorf("scaling up daemonset %s.%s failed: %w", depCopy.GetName(), depCopy.Namespace, err) } return nil } @@ -84,23 +75,21 @@ func (c *DaemonSetController) ScaleFromZero(cd *flaggerv1.Canary) error { // Initialize creates the primary DaemonSet and // delete the canary DaemonSet and returns the pod selector label and container ports func (c *DaemonSetController) Initialize(cd *flaggerv1.Canary, skipLivenessChecks bool) (err error) { - primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) err = c.createPrimaryDaemonSet(cd) if err != nil { - return fmt.Errorf("creating daemonset %s.%s failed: %v", primaryName, cd.Namespace, err) + return fmt.Errorf("createPrimaryDaemonSet failed: %w", err) } if cd.Status.Phase == "" || cd.Status.Phase == flaggerv1.CanaryPhaseInitializing { if !skipLivenessChecks && !cd.SkipAnalysis() { - _, readyErr := c.IsPrimaryReady(cd) - if readyErr != nil { - return readyErr + if err := c.IsPrimaryReady(cd); err != nil { + return fmt.Errorf("IsPrimaryReady failed: %w", err) } } 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 err + if err := c.ScaleToZero(cd); err != nil { + return fmt.Errorf("ScaleToZero failed: %w", err) } } return nil @@ -113,33 +102,26 @@ func (c *DaemonSetController) Promote(cd *flaggerv1.Canary) error { canary, err := c.kubeClient.AppsV1().DaemonSets(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { - if errors.IsNotFound(err) { - return fmt.Errorf("damonset %s.%s not found", targetName, cd.Namespace) - } - return fmt.Errorf("damonset %s.%s query error %v", targetName, cd.Namespace, err) + return fmt.Errorf("damonset %s.%s get query error %v", targetName, cd.Namespace, err) } label, err := c.getSelectorLabel(canary) if err != nil { - return fmt.Errorf("invalid label selector! DaemonSet %s.%s spec.selector.matchLabels must contain selector 'app: %s'", - targetName, cd.Namespace, targetName) + return fmt.Errorf("getSelectorLabel failed: %w", err) } primary, err := c.kubeClient.AppsV1().DaemonSets(cd.Namespace).Get(primaryName, metav1.GetOptions{}) if err != nil { - if errors.IsNotFound(err) { - return fmt.Errorf("daemonset %s.%s not found", primaryName, cd.Namespace) - } - return fmt.Errorf("daemonset %s.%s query error %v", primaryName, cd.Namespace, err) + return fmt.Errorf("daemonset %s.%s get query error %w", primaryName, cd.Namespace, err) } // promote secrets and config maps configRefs, err := c.configTracker.GetTargetConfigs(cd) if err != nil { - return err + return fmt.Errorf("GetTargetConfigs failed: %w", err) } if err := c.configTracker.CreatePrimaryConfigs(cd, configRefs); err != nil { - return err + return fmt.Errorf("CreatePrimaryConfigs failed: %w", err) } primaryCopy := primary.DeepCopy() @@ -158,16 +140,16 @@ func (c *DaemonSetController) Promote(cd *flaggerv1.Canary) error { // update pod annotations to ensure a rolling update annotations, err := makeAnnotations(canary.Spec.Template.Annotations) if err != nil { - return err + return fmt.Errorf("makeAnnotations failed: %w", err) } - primaryCopy.Spec.Template.Annotations = annotations + primaryCopy.Spec.Template.Annotations = annotations primaryCopy.Spec.Template.Labels = makePrimaryLabels(canary.Spec.Template.Labels, primaryName, label) // apply update _, err = c.kubeClient.AppsV1().DaemonSets(cd.Namespace).Update(primaryCopy) if err != nil { - return fmt.Errorf("updating deployment %s.%s template spec failed: %v", + return fmt.Errorf("updating daemonset %s.%s template spec failed: %w", primaryCopy.GetName(), primaryCopy.Namespace, err) } return nil @@ -178,10 +160,7 @@ func (c *DaemonSetController) HasTargetChanged(cd *flaggerv1.Canary) (bool, erro targetName := cd.Spec.TargetRef.Name canary, err := c.kubeClient.AppsV1().DaemonSets(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { - if errors.IsNotFound(err) { - return false, fmt.Errorf("daemonset %s.%s not found", targetName, cd.Namespace) - } - return false, fmt.Errorf("daemonset %s.%s query error %v", targetName, cd.Namespace, err) + return false, fmt.Errorf("daemonset %s.%s get query error %w", targetName, cd.Namespace, err) } // ignore `daemonSetScaleDownNodeSelector` node selector @@ -203,27 +182,18 @@ func (c *DaemonSetController) GetMetadata(cd *flaggerv1.Canary) (string, map[str canaryDae, err := c.kubeClient.AppsV1().DaemonSets(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { - if errors.IsNotFound(err) { - return "", nil, fmt.Errorf("daemonset %s.%s not found, retrying", targetName, cd.Namespace) - } - return "", nil, err + return "", nil, fmt.Errorf("daemonset %s.%s get query error %w", targetName, cd.Namespace, err) } label, err := c.getSelectorLabel(canaryDae) if err != nil { - return "", nil, fmt.Errorf("invalid label selector! DaemonSet %s.%s spec.selector.matchLabels must contain selector 'app: %s'", - targetName, cd.Namespace, targetName) + return "", nil, fmt.Errorf("getSelectorLabel failed: %w", err) } var ports map[string]int32 if cd.Spec.Service.PortDiscovery { - p, err := getPorts(cd, canaryDae.Spec.Template.Spec.Containers) - if err != nil { - return "", nil, fmt.Errorf("port discovery failed with error: %v", err) - } - ports = p + ports = getPorts(cd, canaryDae.Spec.Template.Spec.Containers) } - return label, ports, nil } @@ -233,10 +203,7 @@ func (c *DaemonSetController) createPrimaryDaemonSet(cd *flaggerv1.Canary) error canaryDae, err := c.kubeClient.AppsV1().DaemonSets(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { - if errors.IsNotFound(err) { - return fmt.Errorf("daemonset %s.%s not found, retrying", targetName, cd.Namespace) - } - return err + return fmt.Errorf("daemonset %s.%s get query error %w", targetName, cd.Namespace, err) } if canaryDae.Spec.UpdateStrategy.Type != "" && @@ -247,27 +214,26 @@ func (c *DaemonSetController) createPrimaryDaemonSet(cd *flaggerv1.Canary) error label, err := c.getSelectorLabel(canaryDae) if err != nil { - return fmt.Errorf("invalid label selector! DaemonSet %s.%s spec.selector.matchLabels must contain selector 'app: %s'", - targetName, cd.Namespace, targetName) + return fmt.Errorf("getSelectorLabel failed: %w", err) } - primaryDep, err := c.kubeClient.AppsV1().DaemonSets(cd.Namespace).Get(primaryName, metav1.GetOptions{}) + primaryDae, err := c.kubeClient.AppsV1().DaemonSets(cd.Namespace).Get(primaryName, metav1.GetOptions{}) if errors.IsNotFound(err) { // create primary secrets and config maps configRefs, err := c.configTracker.GetTargetConfigs(cd) if err != nil { - return err + return fmt.Errorf("GetTargetConfigs failed: %w", err) } if err := c.configTracker.CreatePrimaryConfigs(cd, configRefs); err != nil { - return err + return fmt.Errorf("CreatePrimaryConfigs failed: %w", err) } annotations, err := makeAnnotations(canaryDae.Spec.Template.Annotations) if err != nil { - return err + return fmt.Errorf("makeAnnotations failed: %w", err) } - // create primary deployment - primaryDep = &appsv1.DaemonSet{ + // create primary daemonset + primaryDae = &appsv1.DaemonSet{ ObjectMeta: metav1.ObjectMeta{ Name: primaryName, Namespace: cd.Namespace, @@ -302,12 +268,12 @@ func (c *DaemonSetController) createPrimaryDaemonSet(cd *flaggerv1.Canary) error }, } - _, err = c.kubeClient.AppsV1().DaemonSets(cd.Namespace).Create(primaryDep) + _, err = c.kubeClient.AppsV1().DaemonSets(cd.Namespace).Create(primaryDae) if err != nil { - return err + return fmt.Errorf("creating daemonset %s.%s failed: %w", primaryDae.Name, cd.Namespace, err) } - c.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("DaemonSet %s.%s created", primaryDep.GetName(), cd.Namespace) + c.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("DaemonSet %s.%s created", primaryDae.GetName(), cd.Namespace) } return nil } @@ -320,7 +286,10 @@ func (c *DaemonSetController) getSelectorLabel(daemonSet *appsv1.DaemonSet) (str } } - return "", fmt.Errorf("selector not found") + return "", fmt.Errorf( + "daemonset %s.%s spec.selector.matchLabels must contain one of %v'", + c.labels, daemonSet.Name, daemonSet.Namespace, + ) } func (c *DaemonSetController) HaveDependenciesChanged(cd *flaggerv1.Canary) (bool, error) { diff --git a/pkg/canary/daemonset_controller_test.go b/pkg/canary/daemonset_controller_test.go index 9bd5bc53..0222a710 100644 --- a/pkg/canary/daemonset_controller_test.go +++ b/pkg/canary/daemonset_controller_test.go @@ -161,12 +161,12 @@ func TestDaemonSetController_HasTargetChanged(t *testing.T) { } func TestDaemonSetController_Scale(t *testing.T) { - t.Run("Scale", func(t *testing.T) { + t.Run("ScaleToZero", func(t *testing.T) { mocks := newDaemonSetFixture() err := mocks.controller.Initialize(mocks.canary, true) require.NoError(t, err) - err = mocks.controller.Scale(mocks.canary, 0) + err = mocks.controller.ScaleToZero(mocks.canary) require.NoError(t, err) c, err := mocks.kubeClient.AppsV1().DaemonSets("default").Get("podinfo", metav1.GetOptions{}) diff --git a/pkg/canary/daemonset_ready.go b/pkg/canary/daemonset_ready.go index 4120082a..81b892b1 100644 --- a/pkg/canary/daemonset_ready.go +++ b/pkg/canary/daemonset_ready.go @@ -5,7 +5,6 @@ import ( "time" appsv1 "k8s.io/api/apps/v1" - "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1beta1" @@ -13,21 +12,18 @@ import ( // IsPrimaryReady checks the primary daemonset status and returns an error if // the daemonset is in the middle of a rolling update -func (c *DaemonSetController) IsPrimaryReady(cd *flaggerv1.Canary) (bool, error) { +func (c *DaemonSetController) IsPrimaryReady(cd *flaggerv1.Canary) error { primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) primary, err := c.kubeClient.AppsV1().DaemonSets(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) - } - return true, fmt.Errorf("deployment %s.%s query error %v", primaryName, cd.Namespace, err) + return fmt.Errorf("daemonset %s.%s get query error %w", primaryName, cd.Namespace, err) } - retriable, err := c.isDaemonSetReady(cd, primary) + _, err = c.isDaemonSetReady(cd, primary) if err != nil { - return retriable, fmt.Errorf("halt advancement %s.%s %s", primaryName, cd.Namespace, err.Error()) + return fmt.Errorf("primary daemonset %s.%s not ready: %w", primaryName, cd.Namespace, err) } - return true, nil + return nil } // IsCanaryReady checks the primary daemonset and returns an error if @@ -36,15 +32,13 @@ func (c *DaemonSetController) IsCanaryReady(cd *flaggerv1.Canary) (bool, error) targetName := cd.Spec.TargetRef.Name canary, err := c.kubeClient.AppsV1().DaemonSets(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { - if errors.IsNotFound(err) { - return true, fmt.Errorf("daemonset %s.%s not found", targetName, cd.Namespace) - } - return true, fmt.Errorf("daemonset %s.%s query error %v", targetName, cd.Namespace, err) + return true, fmt.Errorf("daemonset %s.%s get query error %w", targetName, cd.Namespace, err) } retriable, err := c.isDaemonSetReady(cd, canary) if err != nil { - return retriable, fmt.Errorf("halt advancement %s.%s %s", targetName, cd.Namespace, err.Error()) + return retriable, fmt.Errorf("canary damonset %s.%s not ready with retryablility: %v: %w", + targetName, cd.Namespace, retriable, err) } return true, nil } @@ -56,7 +50,7 @@ func (c *DaemonSetController) isDaemonSetReady(cd *flaggerv1.Canary, daemonSet * delta := time.Duration(cd.GetProgressDeadlineSeconds()) * time.Second dl := from.Add(delta) if dl.Before(time.Now()) { - return false, fmt.Errorf("daemonset %s exceeded its progress deadline", cd.GetName()) + return false, fmt.Errorf("exceeded its progressDeadlineSeconds: %d", cd.GetProgressDeadlineSeconds()) } else { return true, fmt.Errorf( "waiting for rollout to finish: desiredNumberScheduled=%d, updatedNumberScheduled=%d, numberUnavailable=%d", diff --git a/pkg/canary/daemonset_ready_test.go b/pkg/canary/daemonset_ready_test.go index 9754b866..76a810cf 100644 --- a/pkg/canary/daemonset_ready_test.go +++ b/pkg/canary/daemonset_ready_test.go @@ -16,7 +16,7 @@ func TestDaemonSetController_IsReady(t *testing.T) { err := mocks.controller.Initialize(mocks.canary, true) assert.NoError(t, err, "Expected primary readiness check to fail") - _, err = mocks.controller.IsPrimaryReady(mocks.canary) + err = mocks.controller.IsPrimaryReady(mocks.canary) require.NoError(t, err) _, err = mocks.controller.IsCanaryReady(mocks.canary) diff --git a/pkg/canary/daemonset_status.go b/pkg/canary/daemonset_status.go index 49329866..dd558f96 100644 --- a/pkg/canary/daemonset_status.go +++ b/pkg/canary/daemonset_status.go @@ -3,8 +3,6 @@ package canary import ( "fmt" - ex "github.com/pkg/errors" - "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1beta1" @@ -14,10 +12,7 @@ import ( func (c *DaemonSetController) SyncStatus(cd *flaggerv1.Canary, status flaggerv1.CanaryStatus) error { dae, err := c.kubeClient.AppsV1().DaemonSets(cd.Namespace).Get(cd.Spec.TargetRef.Name, metav1.GetOptions{}) if err != nil { - if errors.IsNotFound(err) { - return fmt.Errorf("daemonset %s.%s not found", cd.Spec.TargetRef.Name, cd.Namespace) - } - return ex.Wrap(err, "SyncStatus daemonset query error") + return fmt.Errorf("daemonset %s.%s get query error: %w", cd.Spec.TargetRef.Name, cd.Namespace, err) } // ignore `daemonSetScaleDownNodeSelector` node selector @@ -32,7 +27,7 @@ func (c *DaemonSetController) SyncStatus(cd *flaggerv1.Canary, status flaggerv1. configs, err := c.configTracker.GetConfigRefs(cd) if err != nil { - return ex.Wrap(err, "SyncStatus configs query error") + return fmt.Errorf("GetConfigRefs failed: %w", err) } return syncCanaryStatus(c.flaggerClient, cd, status, dae.Spec.Template, func(cdCopy *flaggerv1.Canary) { diff --git a/pkg/canary/deployment_controller.go b/pkg/canary/deployment_controller.go index 96b6b77f..8364845d 100644 --- a/pkg/canary/deployment_controller.go +++ b/pkg/canary/deployment_controller.go @@ -33,26 +33,31 @@ func (c *DeploymentController) Initialize(cd *flaggerv1.Canary, skipLivenessChec err = c.createPrimaryDeployment(cd) if err != nil { - return fmt.Errorf("creating deployment %s.%s failed: %v", primaryName, cd.Namespace, err) + return fmt.Errorf("createPrimaryDeployment %s.%s failed: %w", primaryName, cd.Namespace, err) } if cd.Status.Phase == "" || cd.Status.Phase == flaggerv1.CanaryPhaseInitializing { if !skipLivenessChecks && !cd.SkipAnalysis() { - _, readyErr := c.IsPrimaryReady(cd) - if readyErr != nil { - return readyErr + if err := c.IsPrimaryReady(cd); err != nil { + return fmt.Errorf("primary deployment %s.%s not ready: %w", primaryName, cd.Namespace, err) } } 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 err + if err := c.ScaleToZero(cd); err != nil { + return fmt.Errorf("scaling down canary daemon set %s.%s failed: %w", cd.Spec.TargetRef.Name, cd.Namespace, err) } } - if cd.Spec.AutoscalerRef != nil && cd.Spec.AutoscalerRef.Kind == "HorizontalPodAutoscaler" { - if err := c.reconcilePrimaryHpa(cd, true); err != nil { - return fmt.Errorf("creating HorizontalPodAutoscaler %s.%s failed: %v", primaryName, cd.Namespace, err) + if cd.Spec.AutoscalerRef != nil { + switch cd.Spec.AutoscalerRef.Kind { + case "HorizontalPodAutoscaler": + if err := c.reconcilePrimaryHpa(cd, true); err != nil { + return fmt.Errorf( + "initial reconcilePrimaryHpa for %s.%s failed: %w", primaryName, cd.Namespace, err) + } + default: + return fmt.Errorf("cd.Spec.AutoscalerRef.Kind is invalid: %s", cd.Spec.AutoscalerRef.Kind) } } return nil @@ -65,33 +70,26 @@ func (c *DeploymentController) Promote(cd *flaggerv1.Canary) error { 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) - } return fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) } label, err := c.getSelectorLabel(canary) if err != nil { - return fmt.Errorf("invalid label selector! Deployment %s.%s spec.selector.matchLabels must contain selector 'app: %s'", - targetName, cd.Namespace, targetName) + return fmt.Errorf("getSelectorLabel failed: %w", err) } 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) - } return fmt.Errorf("deployment %s.%s query error %v", primaryName, cd.Namespace, err) } // promote secrets and config maps configRefs, err := c.configTracker.GetTargetConfigs(cd) if err != nil { - return err + return fmt.Errorf("GetTargetConfigs failed: %w", err) } if err := c.configTracker.CreatePrimaryConfigs(cd, configRefs); err != nil { - return err + return fmt.Errorf("CreatePrimaryConfigs failed: %w", err) } primaryCopy := primary.DeepCopy() @@ -106,26 +104,31 @@ func (c *DeploymentController) Promote(cd *flaggerv1.Canary) error { // update pod annotations to ensure a rolling update annotations, err := makeAnnotations(canary.Spec.Template.Annotations) if err != nil { - return err + return fmt.Errorf("makeAnnotations failed: %w", err) } - primaryCopy.Spec.Template.Annotations = annotations + primaryCopy.Spec.Template.Annotations = annotations primaryCopy.Spec.Template.Labels = makePrimaryLabels(canary.Spec.Template.Labels, primaryName, label) // apply update _, err = c.kubeClient.AppsV1().Deployments(cd.Namespace).Update(primaryCopy) if err != nil { - return fmt.Errorf("updating deployment %s.%s template spec failed: %v", + return fmt.Errorf("updating deployment %s.%s template spec failed: %w", primaryCopy.GetName(), primaryCopy.Namespace, err) } // update HPA - if cd.Spec.AutoscalerRef != nil && cd.Spec.AutoscalerRef.Kind == "HorizontalPodAutoscaler" { - if err := c.reconcilePrimaryHpa(cd, false); err != nil { - return fmt.Errorf("updating HorizontalPodAutoscaler %s.%s failed: %v", primaryName, cd.Namespace, err) + if cd.Spec.AutoscalerRef != nil { + switch cd.Spec.AutoscalerRef.Kind { + case "HorizontalPodAutoscaler": + if err := c.reconcilePrimaryHpa(cd, false); err != nil { + return fmt.Errorf( + "reconcilePrimaryHpa for %s.%s failed: %w", primaryName, cd.Namespace, err) + } + default: + return fmt.Errorf("cd.Spec.AutoscalerRef.Kind is invalid: %s", cd.Spec.AutoscalerRef.Kind) } } - return nil } @@ -134,32 +137,26 @@ func (c *DeploymentController) HasTargetChanged(cd *flaggerv1.Canary) (bool, err targetName := cd.Spec.TargetRef.Name 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) - } - return false, fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) + return false, fmt.Errorf("deployment %s.%s query error: %w", targetName, cd.Namespace, err) } return hasSpecChanged(cd, canary.Spec.Template) } // Scale sets the canary deployment replicas -func (c *DeploymentController) Scale(cd *flaggerv1.Canary, replicas int32) error { +func (c *DeploymentController) ScaleToZero(cd *flaggerv1.Canary) error { targetName := cd.Spec.TargetRef.Name 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) - } - return fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) + return fmt.Errorf("deployment %s.%s get query error: %w", targetName, cd.Namespace, err) } depCopy := dep.DeepCopy() - depCopy.Spec.Replicas = int32p(replicas) + depCopy.Spec.Replicas = int32p(0) _, 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 fmt.Errorf("deployment %s.%s update query error: %w", targetName, cd.Namespace, err) } return nil } @@ -168,10 +165,7 @@ 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{}) if err != nil { - if errors.IsNotFound(err) { - return fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) - } - return fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) + return fmt.Errorf("deployment %s.%s get query error: %w", targetName, cd.Namespace, err) } replicas := int32p(1) @@ -183,7 +177,7 @@ func (c *DeploymentController) ScaleFromZero(cd *flaggerv1.Canary) error { _, 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 fmt.Errorf("scaling up %s.%s to %v failed: %v", depCopy.GetName(), depCopy.Namespace, replicas, err) } return nil } @@ -194,25 +188,17 @@ func (c *DeploymentController) GetMetadata(cd *flaggerv1.Canary) (string, map[st 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) - } - return "", nil, err + return "", nil, fmt.Errorf("deployment %s.%s get query error %w", targetName, cd.Namespace, err) } label, err := c.getSelectorLabel(canaryDep) if err != nil { - return "", nil, fmt.Errorf("invalid label selector! Deployment %s.%s spec.selector.matchLabels must contain selector 'app: %s'", - targetName, cd.Namespace, targetName) + return "", nil, fmt.Errorf("getSelectorLabel failed: %w", err) } var ports map[string]int32 if cd.Spec.Service.PortDiscovery { - p, err := getPorts(cd, canaryDep.Spec.Template.Spec.Containers) - if err != nil { - return "", nil, fmt.Errorf("port discovery failed with error: %v", err) - } - ports = p + ports = getPorts(cd, canaryDep.Spec.Template.Spec.Containers) } return label, ports, nil @@ -223,16 +209,12 @@ func (c *DeploymentController) createPrimaryDeployment(cd *flaggerv1.Canary) err canaryDep, 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, retrying", targetName, cd.Namespace) - } - return err + return fmt.Errorf("deplyoment %s.%s get query error %w", targetName, cd.Namespace, err) } label, err := c.getSelectorLabel(canaryDep) if err != nil { - return fmt.Errorf("invalid label selector! Deployment %s.%s spec.selector.matchLabels must contain selector 'app: %s'", - targetName, cd.Namespace, targetName) + return fmt.Errorf("getSelectorLabel failed: %w", err) } primaryDep, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(primaryName, metav1.GetOptions{}) @@ -240,14 +222,14 @@ func (c *DeploymentController) createPrimaryDeployment(cd *flaggerv1.Canary) err // create primary secrets and config maps configRefs, err := c.configTracker.GetTargetConfigs(cd) if err != nil { - return err + return fmt.Errorf("GetTargetConfigs failed: %w", err) } if err := c.configTracker.CreatePrimaryConfigs(cd, configRefs); err != nil { - return err + return fmt.Errorf("CreatePrimaryConfigs failed: %w", err) } annotations, err := makeAnnotations(canaryDep.Spec.Template.Annotations) if err != nil { - return err + return fmt.Errorf("makeAnnotations failed: %w", err) } replicas := int32(1) @@ -295,7 +277,7 @@ func (c *DeploymentController) createPrimaryDeployment(cd *flaggerv1.Canary) err _, err = c.kubeClient.AppsV1().Deployments(cd.Namespace).Create(primaryDep) if err != nil { - return err + return fmt.Errorf("creating deployment %s.%s failed: %w", primaryDep.Name, cd.Namespace, err) } c.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("Deployment %s.%s created", primaryDep.GetName(), cd.Namespace) @@ -308,11 +290,8 @@ func (c *DeploymentController) reconcilePrimaryHpa(cd *flaggerv1.Canary, init bo primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) 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", - cd.Spec.AutoscalerRef.Name, cd.Namespace) - } - return err + return fmt.Errorf("HorizontalPodAutoscaler %s.%s get query error: %w", + cd.Spec.AutoscalerRef.Name, cd.Namespace, err) } hpaSpec := hpav1.HorizontalPodAutoscalerSpec{ @@ -349,14 +328,15 @@ func (c *DeploymentController) reconcilePrimaryHpa(cd *flaggerv1.Canary, init bo _, err = c.kubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Create(primaryHpa) if err != nil { - return err + return fmt.Errorf("creating HorizontalPodAutoscaler %s.%s failed: %w", + primaryHpa.Name, primaryHpa.Namespace, 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 - } - - if err != nil { - return err + } else if err != nil { + return fmt.Errorf("HorizontalPodAutoscaler %s.%s exists but get query failed: %w", + primaryHpa.Name, primaryHpa.Namespace, err) } // update HPA @@ -369,14 +349,15 @@ func (c *DeploymentController) reconcilePrimaryHpa(cd *flaggerv1.Canary, init bo hpaClone.Spec.MinReplicas = hpaSpec.MinReplicas hpaClone.Spec.Metrics = hpaSpec.Metrics - _, upErr := c.kubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Update(hpaClone) - if upErr != nil { - return upErr + _, err := c.kubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Update(hpaClone) + if err != nil { + return fmt.Errorf("updating HorizontalPodAutoscaler %s.%s failed: %w", + hpaClone.Name, hpaClone.Namespace, err) } - 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) } } - return nil } @@ -388,7 +369,10 @@ func (c *DeploymentController) getSelectorLabel(deployment *appsv1.Deployment) ( } } - return "", fmt.Errorf("selector not found") + return "", fmt.Errorf( + "deployment %s.%s spec.selector.matchLabels must contain one of %v", + c.labels, deployment.Name, deployment.Namespace, + ) } func (c *DeploymentController) HaveDependenciesChanged(cd *flaggerv1.Canary) (bool, error) { diff --git a/pkg/canary/deployment_controller_test.go b/pkg/canary/deployment_controller_test.go index 06958e3f..fb8d574c 100644 --- a/pkg/canary/deployment_controller_test.go +++ b/pkg/canary/deployment_controller_test.go @@ -77,7 +77,7 @@ func TestDeploymentController_IsReady(t *testing.T) { err := mocks.controller.Initialize(mocks.canary, true) require.NoError(t, err, "Expected primary readiness check to fail") - _, err = mocks.controller.IsPrimaryReady(mocks.canary) + err = mocks.controller.IsPrimaryReady(mocks.canary) require.Error(t, err) _, err = mocks.controller.IsCanaryReady(mocks.canary) @@ -134,17 +134,17 @@ func TestDeploymentController_SyncStatus(t *testing.T) { assert.True(t, exists, "Secret %s not found in status", secret.GetName()) } -func TestDeploymentController_Scale(t *testing.T) { +func TestDeploymentController_ScaleToZero(t *testing.T) { mocks := newDeploymentFixture() err := mocks.controller.Initialize(mocks.canary, true) require.NoError(t, err) - err = mocks.controller.Scale(mocks.canary, 2) + err = mocks.controller.ScaleToZero(mocks.canary) require.NoError(t, err) c, err := mocks.kubeClient.AppsV1().Deployments("default").Get("podinfo", metav1.GetOptions{}) require.NoError(t, err) - assert.Equal(t, int32(2), *c.Spec.Replicas) + assert.Equal(t, int32(0), *c.Spec.Replicas) } func TestDeploymentController_NoConfigTracking(t *testing.T) { diff --git a/pkg/canary/deployment_ready.go b/pkg/canary/deployment_ready.go index 083633f9..6fdd9d9a 100644 --- a/pkg/canary/deployment_ready.go +++ b/pkg/canary/deployment_ready.go @@ -5,7 +5,6 @@ import ( "time" appsv1 "k8s.io/api/apps/v1" - "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1beta1" @@ -14,26 +13,23 @@ 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 *DeploymentController) IsPrimaryReady(cd *flaggerv1.Canary) (bool, error) { +func (c *DeploymentController) IsPrimaryReady(cd *flaggerv1.Canary) error { primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) 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) - } - return true, fmt.Errorf("deployment %s.%s query error %v", primaryName, cd.Namespace, err) + return fmt.Errorf("deployment %s.%s get query error %w", primaryName, cd.Namespace, err) } - retriable, err := c.isDeploymentReady(primary, cd.GetProgressDeadlineSeconds()) + _, err = c.isDeploymentReady(primary, cd.GetProgressDeadlineSeconds()) if err != nil { - return retriable, fmt.Errorf("Halt advancement %s.%s %s", primaryName, cd.Namespace, err.Error()) + return fmt.Errorf("primary daemonset %s.%s not ready: %w", primaryName, cd.Namespace, err) } if primary.Spec.Replicas == int32p(0) { - return true, fmt.Errorf("Halt %s.%s advancement primary deployment is scaled to zero", + return fmt.Errorf("halt %s.%s advancement: primary deployment is scaled to zero", cd.Name, cd.Namespace) } - return true, nil + return nil } // IsCanaryReady checks the canary deployment status and returns an error if @@ -43,22 +39,16 @@ 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{}) if err != nil { - if errors.IsNotFound(err) { - return true, fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) - } - return true, fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) + return true, fmt.Errorf("deployment %s.%s get query error: %w", targetName, cd.Namespace, err) } retriable, err := c.isDeploymentReady(canary, cd.GetProgressDeadlineSeconds()) if err != nil { - if retriable { - return retriable, fmt.Errorf("Halt advancement %s.%s %s", targetName, cd.Namespace, err.Error()) - } else { - return retriable, fmt.Errorf("deployment does not have minimum availability for more than %vs", - cd.GetProgressDeadlineSeconds()) - } + return retriable, fmt.Errorf( + "canary deployment %s.%s not ready with retriablility %v: %w", + targetName, cd.Namespace, retriable, err, + ) } - return true, nil } @@ -93,9 +83,9 @@ func (c *DeploymentController) isDeploymentReady(deployment *appsv1.Deployment, } } else { - return true, fmt.Errorf("waiting for rollout to finish: observed deployment generation less then desired generation") + return true, fmt.Errorf( + "waiting for rollout to finish: observed deployment generation less then desired generation") } - return true, nil } diff --git a/pkg/canary/deployment_status.go b/pkg/canary/deployment_status.go index e5577b51..6c2babe6 100644 --- a/pkg/canary/deployment_status.go +++ b/pkg/canary/deployment_status.go @@ -3,8 +3,6 @@ package canary import ( "fmt" - ex "github.com/pkg/errors" - "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1beta1" @@ -14,15 +12,12 @@ import ( 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) - } - return ex.Wrap(err, "SyncStatus deployment query error") + return fmt.Errorf("deployment %s.%s get query error: %w", cd.Spec.TargetRef.Name, cd.Namespace, err) } configs, err := c.configTracker.GetConfigRefs(cd) if err != nil { - return ex.Wrap(err, "SyncStatus configs query error") + return fmt.Errorf("GetConfigRefs failed: %w", err) } return syncCanaryStatus(c.flaggerClient, cd, status, dep.Spec.Template, func(cdCopy *flaggerv1.Canary) { diff --git a/pkg/canary/factory.go b/pkg/canary/factory.go index 416a0c91..95f243b8 100644 --- a/pkg/canary/factory.go +++ b/pkg/canary/factory.go @@ -50,12 +50,12 @@ func (factory *Factory) Controller(kind string) Controller { flaggerClient: factory.flaggerClient, } - switch { - case kind == "DaemonSet": + switch kind { + case "DaemonSet": return daemonSetCtrl - case kind == "Deployment": + case "Deployment": return deploymentCtrl - case kind == "Service": + case "Service": return serviceCtrl default: return deploymentCtrl diff --git a/pkg/canary/nop_tracker.go b/pkg/canary/nop_tracker.go index d2234a0b..866c08c4 100644 --- a/pkg/canary/nop_tracker.go +++ b/pkg/canary/nop_tracker.go @@ -6,8 +6,7 @@ import ( ) // NopTracker no-operation tracker -type NopTracker struct { -} +type NopTracker struct{} func (nt *NopTracker) GetTargetConfigs(*flaggerv1.Canary) (map[string]ConfigRef, error) { res := make(map[string]ConfigRef) @@ -26,6 +25,6 @@ func (nt *NopTracker) CreatePrimaryConfigs(*flaggerv1.Canary, map[string]ConfigR return nil } -func (nt *NopTracker) ApplyPrimaryConfigs(spec corev1.PodSpec, refs map[string]ConfigRef) corev1.PodSpec { +func (nt *NopTracker) ApplyPrimaryConfigs(spec corev1.PodSpec, _ map[string]ConfigRef) corev1.PodSpec { return spec } diff --git a/pkg/canary/service_controller.go b/pkg/canary/service_controller.go index 3e3724d2..cfd1703f 100644 --- a/pkg/canary/service_controller.go +++ b/pkg/canary/service_controller.go @@ -3,7 +3,6 @@ package canary import ( "fmt" - ex "github.com/pkg/errors" "go.uber.org/zap" corev1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/api/errors" @@ -43,31 +42,27 @@ func (c *ServiceController) SetStatusPhase(cd *flaggerv1.Canary, phase flaggerv1 } // GetMetadata returns the pod label selector and svc ports -func (c *ServiceController) GetMetadata(cd *flaggerv1.Canary) (string, map[string]int32, error) { +func (c *ServiceController) GetMetadata(_ *flaggerv1.Canary) (string, map[string]int32, error) { return "", nil, nil } // Initialize creates or updates the primary and canary services to prepare for the canary release process targeted on the K8s service -func (c *ServiceController) Initialize(cd *flaggerv1.Canary, skipLivenessChecks bool) (err error) { +func (c *ServiceController) Initialize(cd *flaggerv1.Canary, _ bool) (err error) { targetName := cd.Spec.TargetRef.Name primaryName := fmt.Sprintf("%s-primary", targetName) canaryName := fmt.Sprintf("%s-canary", targetName) svc, err := c.kubeClient.CoreV1().Services(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { - return err + return fmt.Errorf("service %s.%s get query error: %w", primaryName, cd.Namespace, err) } - // canary svc - err = c.reconcileCanaryService(cd, canaryName, svc) - if err != nil { - return err + if err = c.reconcileCanaryService(cd, canaryName, svc); err != nil { + return fmt.Errorf("reconcileCanaryService failed: %w", err) } - // primary svc - err = c.reconcilePrimaryService(cd, primaryName, svc) - if err != nil { - return err + if err = c.reconcilePrimaryService(cd, primaryName, svc); err != nil { + return fmt.Errorf("reconcilePrimaryService failed: %w", err) } return nil @@ -77,31 +72,29 @@ func (c *ServiceController) reconcileCanaryService(canary *flaggerv1.Canary, nam current, err := c.kubeClient.CoreV1().Services(canary.Namespace).Get(name, metav1.GetOptions{}) if errors.IsNotFound(err) { return c.createService(canary, name, src) + } else if err != nil { + return fmt.Errorf("service %s.%s get query error: %w", name, canary.Namespace, err) } - if err != nil { - return fmt.Errorf("service %s query error %v", name, err) - } + ns := buildService(canary, name, src) - new := buildService(canary, name, src) - - if new.Spec.Type == "ClusterIP" { + if ns.Spec.Type == "ClusterIP" { // We can't change this immutable field - new.Spec.ClusterIP = current.Spec.ClusterIP + ns.Spec.ClusterIP = current.Spec.ClusterIP } // We can't change this immutable field - new.ObjectMeta.UID = current.ObjectMeta.UID + ns.ObjectMeta.UID = current.ObjectMeta.UID - new.ObjectMeta.ResourceVersion = current.ObjectMeta.ResourceVersion + ns.ObjectMeta.ResourceVersion = current.ObjectMeta.ResourceVersion - _, err = c.kubeClient.CoreV1().Services(canary.Namespace).Update(new) + _, err = c.kubeClient.CoreV1().Services(canary.Namespace).Update(ns) if err != nil { - return err + return fmt.Errorf("updating service %s.%s failed: %w", name, canary.Namespace, err) } c.logger.With("canary", fmt.Sprintf("%s.%s", canary.Name, canary.Namespace)). - Infof("Service %s.%s updated", new.GetName(), canary.Namespace) + Infof("Service %s.%s updated", ns.GetName(), canary.Namespace) return nil } @@ -109,12 +102,9 @@ func (c *ServiceController) reconcilePrimaryService(canary *flaggerv1.Canary, na _, err := c.kubeClient.CoreV1().Services(canary.Namespace).Get(name, metav1.GetOptions{}) if errors.IsNotFound(err) { return c.createService(canary, name, src) + } else if err != nil { + return fmt.Errorf("service %s.%s get query error: %w", name, canary.Namespace, err) } - - if err != nil { - return fmt.Errorf("service %s query error %v", name, err) - } - return nil } @@ -131,7 +121,7 @@ func (c *ServiceController) createService(canary *flaggerv1.Canary, name string, _, err := c.kubeClient.CoreV1().Services(canary.Namespace).Create(svc) if err != nil { - return err + return fmt.Errorf("creating service %s.%s query error: %w", canary.Name, canary.Namespace, err) } c.logger.With("canary", fmt.Sprintf("%s.%s", canary.Name, canary.Namespace)). @@ -156,7 +146,6 @@ func buildService(canary *flaggerv1.Canary, name string, src *corev1.Service) *c // Operation cannot be fulfilled on services "mysvc-canary": the object has been modified; please apply your changes to the latest version and try again delete(svc.ObjectMeta.Annotations, "kubectl.kubernetes.io/last-applied-configuration") } - return svc } @@ -167,18 +156,12 @@ func (c *ServiceController) Promote(cd *flaggerv1.Canary) error { canary, err := c.kubeClient.CoreV1().Services(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { - if errors.IsNotFound(err) { - return fmt.Errorf("service %s.%s not found", targetName, cd.Namespace) - } - return fmt.Errorf("service %s.%s query error %v", targetName, cd.Namespace, err) + return fmt.Errorf("service %s.%s get query error %w", targetName, cd.Namespace, err) } primary, err := c.kubeClient.CoreV1().Services(cd.Namespace).Get(primaryName, metav1.GetOptions{}) if err != nil { - if errors.IsNotFound(err) { - return fmt.Errorf("service %s.%s not found", primaryName, cd.Namespace) - } - return fmt.Errorf("service %s.%s query error %v", primaryName, cd.Namespace, err) + return fmt.Errorf("service %s.%s get query error %w", primaryName, cd.Namespace, err) } primaryCopy := canary.DeepCopy() @@ -192,7 +175,7 @@ func (c *ServiceController) Promote(cd *flaggerv1.Canary) error { // apply update _, err = c.kubeClient.CoreV1().Services(cd.Namespace).Update(primaryCopy) if err != nil { - return fmt.Errorf("updating service %s.%s spec failed: %v", + return fmt.Errorf("updating service %s.%s spec failed: %w", primaryCopy.GetName(), primaryCopy.Namespace, err) } @@ -204,44 +187,37 @@ func (c *ServiceController) HasTargetChanged(cd *flaggerv1.Canary) (bool, error) targetName := cd.Spec.TargetRef.Name canary, err := c.kubeClient.CoreV1().Services(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { - if errors.IsNotFound(err) { - return false, fmt.Errorf("service %s.%s not found", targetName, cd.Namespace) - } - return false, fmt.Errorf("service %s.%s query error %v", targetName, cd.Namespace, err) + return false, fmt.Errorf("service %s.%s get query error %w", targetName, cd.Namespace, err) } - return hasSpecChanged(cd, canary.Spec) } // Scale sets the canary deployment replicas -func (c *ServiceController) Scale(cd *flaggerv1.Canary, replicas int32) error { +func (c *ServiceController) ScaleToZero(_ *flaggerv1.Canary) error { return nil } -func (c *ServiceController) ScaleFromZero(cd *flaggerv1.Canary) error { +func (c *ServiceController) ScaleFromZero(_ *flaggerv1.Canary) error { return nil } func (c *ServiceController) SyncStatus(cd *flaggerv1.Canary, status flaggerv1.CanaryStatus) error { dep, err := c.kubeClient.CoreV1().Services(cd.Namespace).Get(cd.Spec.TargetRef.Name, metav1.GetOptions{}) if err != nil { - if errors.IsNotFound(err) { - return fmt.Errorf("service %s.%s not found", cd.Spec.TargetRef.Name, cd.Namespace) - } - return ex.Wrap(err, "SyncStatus service query error") + return fmt.Errorf("service %s.%s get query error %w", cd.Spec.TargetRef.Name, cd.Namespace, err) } return syncCanaryStatus(c.flaggerClient, cd, status, dep.Spec, func(cdCopy *flaggerv1.Canary) {}) } -func (c *ServiceController) HaveDependenciesChanged(cd *flaggerv1.Canary) (bool, error) { +func (c *ServiceController) HaveDependenciesChanged(_ *flaggerv1.Canary) (bool, error) { return false, nil } -func (c *ServiceController) IsPrimaryReady(cd *flaggerv1.Canary) (bool, error) { - return true, nil +func (c *ServiceController) IsPrimaryReady(_ *flaggerv1.Canary) error { + return nil } -func (c *ServiceController) IsCanaryReady(cd *flaggerv1.Canary) (bool, error) { +func (c *ServiceController) IsCanaryReady(_ *flaggerv1.Canary) (bool, error) { return true, nil } diff --git a/pkg/canary/status.go b/pkg/canary/status.go index d65fd979..6d3f2e34 100644 --- a/pkg/canary/status.go +++ b/pkg/canary/status.go @@ -4,7 +4,6 @@ import ( "fmt" "strings" - ex "github.com/pkg/errors" corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/client-go/util/retry" @@ -17,12 +16,12 @@ func syncCanaryStatus(flaggerClient clientset.Interface, cd *flaggerv1.Canary, s hash := computeHash(canaryResource) firstTry := true + name, ns := cd.GetName(), cd.GetNamespace() err := retry.RetryOnConflict(retry.DefaultBackoff, func() (err error) { - var selErr error if !firstTry { - cd, selErr = flaggerClient.FlaggerV1beta1().Canaries(cd.Namespace).Get(cd.GetName(), metav1.GetOptions{}) - if selErr != nil { - return selErr + cd, err = flaggerClient.FlaggerV1beta1().Canaries(ns).Get(name, metav1.GetOptions{}) + if err != nil { + return fmt.Errorf("canary %s.%s get query failed: %w", name, ns, err) } } @@ -39,72 +38,79 @@ func syncCanaryStatus(flaggerClient clientset.Interface, cd *flaggerv1.Canary, s cdCopy.Status.Conditions = conditions } - err = updateStatusWithUpgrade(flaggerClient, cdCopy) + if err = updateStatusWithUpgrade(flaggerClient, cdCopy); err != nil { + return fmt.Errorf("updateStatusWithUpgrade failed: %w", err) + } firstTry = false return }) + if err != nil { - return ex.Wrap(err, "SyncStatus") + return fmt.Errorf("failed after retries: %w", err) } return nil } func setStatusFailedChecks(flaggerClient clientset.Interface, cd *flaggerv1.Canary, val int) error { firstTry := true + name, ns := cd.GetName(), cd.GetNamespace() err := retry.RetryOnConflict(retry.DefaultBackoff, func() (err error) { - var selErr error if !firstTry { - cd, selErr = flaggerClient.FlaggerV1beta1().Canaries(cd.Namespace).Get(cd.GetName(), metav1.GetOptions{}) - if selErr != nil { - return selErr + cd, err = flaggerClient.FlaggerV1beta1().Canaries(name).Get(ns, metav1.GetOptions{}) + if err != nil { + return fmt.Errorf("canary %s.%s get query failed: %w", name, ns, err) } } cdCopy := cd.DeepCopy() cdCopy.Status.FailedChecks = val cdCopy.Status.LastTransitionTime = metav1.Now() - err = updateStatusWithUpgrade(flaggerClient, cdCopy) + if err = updateStatusWithUpgrade(flaggerClient, cdCopy); err != nil { + return fmt.Errorf("updateStatusWithUpgrade failed: %w", err) + } firstTry = false return }) if err != nil { - return ex.Wrap(err, "SetStatusFailedChecks") + return fmt.Errorf("failed after retries: %w", err) } return nil } func setStatusWeight(flaggerClient clientset.Interface, cd *flaggerv1.Canary, val int) error { firstTry := true + name, ns := cd.GetName(), cd.GetNamespace() err := retry.RetryOnConflict(retry.DefaultBackoff, func() (err error) { - var selErr error if !firstTry { - cd, selErr = flaggerClient.FlaggerV1beta1().Canaries(cd.Namespace).Get(cd.GetName(), metav1.GetOptions{}) - if selErr != nil { - return selErr + cd, err = flaggerClient.FlaggerV1beta1().Canaries(cd.Namespace).Get(cd.GetName(), metav1.GetOptions{}) + if err != nil { + return fmt.Errorf("canary %s.%s get query failed: %w", name, ns, err) } } cdCopy := cd.DeepCopy() cdCopy.Status.CanaryWeight = val cdCopy.Status.LastTransitionTime = metav1.Now() - err = updateStatusWithUpgrade(flaggerClient, cdCopy) + if err = updateStatusWithUpgrade(flaggerClient, cdCopy); err != nil { + return fmt.Errorf("updateStatusWithUpgrade failed: %w", err) + } firstTry = false return }) if err != nil { - return ex.Wrap(err, "SetStatusWeight") + return fmt.Errorf("failed after retries: %w", err) } return nil } func setStatusIterations(flaggerClient clientset.Interface, cd *flaggerv1.Canary, val int) error { firstTry := true + name, ns := cd.GetName(), cd.GetNamespace() err := retry.RetryOnConflict(retry.DefaultBackoff, func() (err error) { - var selErr error if !firstTry { - cd, selErr = flaggerClient.FlaggerV1beta1().Canaries(cd.Namespace).Get(cd.GetName(), metav1.GetOptions{}) - if selErr != nil { - return selErr + cd, err = flaggerClient.FlaggerV1beta1().Canaries(cd.Namespace).Get(cd.GetName(), metav1.GetOptions{}) + if err != nil { + return fmt.Errorf("canary %s.%s get query failed: %w", name, ns, err) } } @@ -112,25 +118,27 @@ func setStatusIterations(flaggerClient clientset.Interface, cd *flaggerv1.Canary cdCopy.Status.Iterations = val cdCopy.Status.LastTransitionTime = metav1.Now() - err = updateStatusWithUpgrade(flaggerClient, cdCopy) + if err = updateStatusWithUpgrade(flaggerClient, cdCopy); err != nil { + return fmt.Errorf("updateStatusWithUpgrade failed: %w", err) + } firstTry = false return }) if err != nil { - return ex.Wrap(err, "SetStatusIterations") + return fmt.Errorf("failed after retries: %w", err) } return nil } func setStatusPhase(flaggerClient clientset.Interface, cd *flaggerv1.Canary, phase flaggerv1.CanaryPhase) error { firstTry := true + name, ns := cd.GetName(), cd.GetNamespace() err := retry.RetryOnConflict(retry.DefaultBackoff, func() (err error) { - var selErr error if !firstTry { - cd, selErr = flaggerClient.FlaggerV1beta1().Canaries(cd.Namespace).Get(cd.GetName(), metav1.GetOptions{}) - if selErr != nil { - return selErr + cd, err = flaggerClient.FlaggerV1beta1().Canaries(cd.Namespace).Get(cd.GetName(), metav1.GetOptions{}) + if err != nil { + return fmt.Errorf("canary %s.%s get query failed: %w", name, ns, err) } } @@ -152,12 +160,14 @@ func setStatusPhase(flaggerClient clientset.Interface, cd *flaggerv1.Canary, pha cdCopy.Status.Conditions = conditions } - err = updateStatusWithUpgrade(flaggerClient, cdCopy) + if err = updateStatusWithUpgrade(flaggerClient, cdCopy); err != nil { + return fmt.Errorf("updateStatusWithUpgrade failed: %w", err) + } firstTry = false return }) if err != nil { - return ex.Wrap(err, "SetStatusPhase") + return fmt.Errorf("failed after retries: %w", err) } return nil } @@ -239,11 +249,14 @@ func updateStatusWithUpgrade(flaggerClient clientset.Interface, cd *flaggerv1.Ca // upgrade alpha resource _, updateErr := flaggerClient.FlaggerV1beta1().Canaries(cd.Namespace).Update(cd) if updateErr != nil { - return updateErr + return fmt.Errorf("updating canary %s.%s from v1alpha to v1beta failed: %w", cd.Name, cd.Namespace, updateErr) } - // retry status update _, err = flaggerClient.FlaggerV1beta1().Canaries(cd.Namespace).UpdateStatus(cd) } + + if err != nil { + return fmt.Errorf("updating canary %s.%s status failed: %w", cd.Name, cd.Namespace, err) + } return err } diff --git a/pkg/canary/util.go b/pkg/canary/util.go index cbe9c0e8..f4f58304 100644 --- a/pkg/canary/util.go +++ b/pkg/canary/util.go @@ -16,7 +16,7 @@ var sidecars = map[string]bool{ "envoy": true, } -func getPorts(cd *flaggerv1.Canary, cs []corev1.Container) (map[string]int32, error) { +func getPorts(cd *flaggerv1.Canary, cs []corev1.Container) map[string]int32 { ports := make(map[string]int32, len(cs)) for _, container := range cs { // exclude service mesh proxies based on container name @@ -49,7 +49,7 @@ func getPorts(cd *flaggerv1.Canary, cs []corev1.Container) (map[string]int32, er ports[name] = p.ContainerPort } } - return ports, nil + return ports } // makeAnnotations appends an unique ID to annotations map @@ -59,7 +59,7 @@ func makeAnnotations(annotations map[string]string) (map[string]string, error) { uuid := make([]byte, 16) n, err := io.ReadFull(rand.Reader, uuid) if n != len(uuid) || err != nil { - return res, err + return res, fmt.Errorf("%w", err) } uuid[8] = uuid[8]&^0xc0 | 0x80 uuid[6] = uuid[6]&^0xf0 | 0x40 diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index 6d393b06..050fbdd1 100644 --- a/pkg/controller/scheduler.go +++ b/pkg/controller/scheduler.go @@ -163,7 +163,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh // check primary status if !skipLivenessChecks && !cd.SkipAnalysis() { - if _, err := canaryController.IsPrimaryReady(cd); err != nil { + if err := canaryController.IsPrimaryReady(cd); err != nil { c.recordEventWarningf(cd, "%v", err) return } @@ -258,7 +258,7 @@ 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 := canaryController.Scale(cd, 0); err != nil { + if err := canaryController.ScaleToZero(cd); err != nil { c.recordEventWarningf(cd, "%v", err) return } @@ -554,7 +554,7 @@ func (c *Controller) shouldSkipAnalysis(canary *flaggerv1.Canary, canaryControll } // shutdown canary - if err := canaryController.Scale(canary, 0); err != nil { + if err := canaryController.ScaleToZero(canary); err != nil { c.recordEventWarningf(canary, "%v", err) return false } @@ -1030,7 +1030,7 @@ func (c *Controller) rollback(canary *flaggerv1.Canary, canaryController canary. c.recorder.SetWeight(canary, primaryWeight, canaryWeight) // shutdown canary - if err := canaryController.Scale(canary, 0); err != nil { + if err := canaryController.ScaleToZero(canary); err != nil { c.recordEventWarningf(canary, "%v", err) return }