diff --git a/pkg/controller/deployer.go b/pkg/controller/deployer.go index 402d6370..f590c781 100644 --- a/pkg/controller/deployer.go +++ b/pkg/controller/deployer.go @@ -32,15 +32,17 @@ type CanaryDeployer struct { // Promote copies the pod spec from canary to primary func (c *CanaryDeployer) Promote(cd *flaggerv1.Canary) error { - canary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(cd.Spec.TargetRef.Name, metav1.GetOptions{}) + targetName := cd.Spec.TargetRef.Name + primaryName := fmt.Sprintf("%s-primary", targetName) + + 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", cd.Spec.TargetRef.Name, cd.Namespace) + return fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) } - return fmt.Errorf("deployment %s.%s query error %v", cd.Spec.TargetRef.Name, cd.Namespace, err) + return fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) } - 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) { @@ -98,12 +100,13 @@ func (c *CanaryDeployer) IsPrimaryReady(cd *flaggerv1.Canary) (bool, error) { // 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 *CanaryDeployer) IsCanaryReady(cd *flaggerv1.Canary) (bool, error) { - canary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(cd.Spec.TargetRef.Name, metav1.GetOptions{}) + 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", cd.Spec.TargetRef.Name, cd.Namespace) + return true, fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) } - return true, fmt.Errorf("deployment %s.%s query error %v", cd.Spec.TargetRef.Name, cd.Namespace, err) + return true, fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) } retriable, err := c.isDeploymentReady(canary, cd.GetProgressDeadlineSeconds()) @@ -121,12 +124,13 @@ func (c *CanaryDeployer) IsCanaryReady(cd *flaggerv1.Canary) (bool, error) { // IsNewSpec returns true if the canary deployment pod spec has changed func (c *CanaryDeployer) IsNewSpec(cd *flaggerv1.Canary) (bool, error) { - canary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(cd.Spec.TargetRef.Name, metav1.GetOptions{}) + 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", cd.Spec.TargetRef.Name, cd.Namespace) + return false, fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) } - return false, fmt.Errorf("deployment %s.%s query error %v", cd.Spec.TargetRef.Name, cd.Namespace, err) + return false, fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) } if cd.Status.CanaryRevision == "" { @@ -208,12 +212,13 @@ func (c *CanaryDeployer) SyncStatus(cd *flaggerv1.Canary, status flaggerv1.Canar // Scale sets the canary deployment replicas func (c *CanaryDeployer) Scale(cd *flaggerv1.Canary, replicas int32) error { - dep, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(cd.Spec.TargetRef.Name, metav1.GetOptions{}) + 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", cd.Spec.TargetRef.Name, cd.Namespace) + return fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) } - return fmt.Errorf("deployment %s.%s query error %v", cd.Spec.TargetRef.Name, cd.Namespace, err) + return fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) } depCopy := dep.DeepCopy() @@ -250,13 +255,13 @@ func (c *CanaryDeployer) Sync(cd *flaggerv1.Canary) error { } func (c *CanaryDeployer) createPrimaryDeployment(cd *flaggerv1.Canary) error { - canaryName := cd.Spec.TargetRef.Name + targetName := cd.Spec.TargetRef.Name primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) - canaryDep, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(canaryName, metav1.GetOptions{}) + 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", canaryName, cd.Namespace) + return fmt.Errorf("deployment %s.%s not found, retrying", targetName, cd.Namespace) } return err } diff --git a/pkg/controller/router.go b/pkg/controller/router.go index 1bf27525..ce8f787b 100644 --- a/pkg/controller/router.go +++ b/pkg/controller/router.go @@ -42,13 +42,13 @@ func (c *CanaryRouter) Sync(cd *flaggerv1.Canary) error { } func (c *CanaryRouter) createServices(cd *flaggerv1.Canary) error { - canaryName := cd.Spec.TargetRef.Name - primaryName := fmt.Sprintf("%s-primary", canaryName) - canaryService, err := c.kubeClient.CoreV1().Services(cd.Namespace).Get(canaryName, metav1.GetOptions{}) + targetName := cd.Spec.TargetRef.Name + primaryName := fmt.Sprintf("%s-primary", targetName) + canaryService, err := c.kubeClient.CoreV1().Services(cd.Namespace).Get(targetName, metav1.GetOptions{}) if errors.IsNotFound(err) { canaryService = &corev1.Service{ ObjectMeta: metav1.ObjectMeta{ - Name: canaryName, + Name: targetName, Namespace: cd.Namespace, OwnerReferences: []metav1.OwnerReference{ *metav1.NewControllerRef(cd, schema.GroupVersionKind{ @@ -60,7 +60,7 @@ func (c *CanaryRouter) createServices(cd *flaggerv1.Canary) error { }, Spec: corev1.ServiceSpec{ Type: corev1.ServiceTypeClusterIP, - Selector: map[string]string{"app": canaryName}, + Selector: map[string]string{"app": targetName}, Ports: []corev1.ServicePort{ { Name: "http", @@ -99,7 +99,7 @@ func (c *CanaryRouter) createServices(cd *flaggerv1.Canary) error { }, Spec: corev1.ServiceSpec{ Type: corev1.ServiceTypeClusterIP, - Selector: map[string]string{"app": canaryName}, + Selector: map[string]string{"app": targetName}, Ports: []corev1.ServicePort{ { Name: "http", @@ -164,16 +164,16 @@ func (c *CanaryRouter) createServices(cd *flaggerv1.Canary) error { } func (c *CanaryRouter) createVirtualService(cd *flaggerv1.Canary) error { - canaryName := cd.Spec.TargetRef.Name - primaryName := fmt.Sprintf("%s-primary", canaryName) - hosts := append(cd.Spec.Service.Hosts, canaryName) + targetName := cd.Spec.TargetRef.Name + primaryName := fmt.Sprintf("%s-primary", targetName) + hosts := append(cd.Spec.Service.Hosts, targetName) gateways := append(cd.Spec.Service.Gateways, "mesh") - virtualService, err := c.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Get(canaryName, metav1.GetOptions{}) + virtualService, err := c.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Get(targetName, metav1.GetOptions{}) if errors.IsNotFound(err) { - c.logger.Debugf("VirtualService %s.%s not found", canaryName, cd.Namespace) + c.logger.Debugf("VirtualService %s.%s not found", targetName, cd.Namespace) virtualService = &istiov1alpha3.VirtualService{ ObjectMeta: metav1.ObjectMeta{ - Name: cd.Name, + Name: targetName, Namespace: cd.Namespace, OwnerReferences: []metav1.OwnerReference{ *metav1.NewControllerRef(cd, schema.GroupVersionKind{ @@ -200,7 +200,7 @@ func (c *CanaryRouter) createVirtualService(cd *flaggerv1.Canary) error { }, { Destination: istiov1alpha3.Destination{ - Host: canaryName, + Host: targetName, Port: istiov1alpha3.PortSelector{ Number: uint32(cd.Spec.Service.Port), }, @@ -216,7 +216,7 @@ func (c *CanaryRouter) createVirtualService(cd *flaggerv1.Canary) error { c.logger.Debugf("Creating VirtualService %s.%s", virtualService.GetName(), cd.Namespace) _, err = c.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Create(virtualService) if err != nil { - return fmt.Errorf("VirtualService %s.%s create error %v", cd.Name, cd.Namespace, err) + return fmt.Errorf("VirtualService %s.%s create error %v", targetName, cd.Namespace, err) } c.logger.Infof("VirtualService %s.%s created", virtualService.GetName(), cd.Namespace) } @@ -230,23 +230,24 @@ func (c *CanaryRouter) GetRoutes(cd *flaggerv1.Canary) ( canary istiov1alpha3.DestinationWeight, err error, ) { + targetName := cd.Spec.TargetRef.Name vs := &istiov1alpha3.VirtualService{} - vs, err = c.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Get(cd.Name, v1.GetOptions{}) + vs, err = c.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Get(targetName, v1.GetOptions{}) if err != nil { if errors.IsNotFound(err) { - err = fmt.Errorf("VirtualService %s.%s not found", cd.Name, cd.Namespace) + err = fmt.Errorf("VirtualService %s.%s not found", targetName, cd.Namespace) return } - err = fmt.Errorf("VirtualService %s.%s query error %v", cd.Name, cd.Namespace, err) + err = fmt.Errorf("VirtualService %s.%s query error %v", targetName, cd.Namespace, err) return } for _, http := range vs.Spec.Http { for _, route := range http.Route { - if route.Destination.Host == fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) { + if route.Destination.Host == fmt.Sprintf("%s-primary", targetName) { primary = route } - if route.Destination.Host == cd.Spec.TargetRef.Name { + if route.Destination.Host == targetName { canary = route } } @@ -254,7 +255,7 @@ func (c *CanaryRouter) GetRoutes(cd *flaggerv1.Canary) ( if primary.Weight == 0 && canary.Weight == 0 { err = fmt.Errorf("VirtualService %s.%s does not contain routes for %s and %s", - cd.Name, cd.Namespace, fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name), cd.Spec.TargetRef.Name) + targetName, cd.Namespace, fmt.Sprintf("%s-primary", targetName), targetName) } return @@ -266,13 +267,14 @@ func (c *CanaryRouter) SetRoutes( primary istiov1alpha3.DestinationWeight, canary istiov1alpha3.DestinationWeight, ) error { - vs, err := c.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Get(cd.Name, v1.GetOptions{}) + targetName := cd.Spec.TargetRef.Name + vs, err := c.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Get(targetName, v1.GetOptions{}) if err != nil { if errors.IsNotFound(err) { - return fmt.Errorf("VirtualService %s.%s not found", cd.Name, cd.Namespace) + return fmt.Errorf("VirtualService %s.%s not found", targetName, cd.Namespace) } - return fmt.Errorf("VirtualService %s.%s query error %v", cd.Name, cd.Namespace, err) + return fmt.Errorf("VirtualService %s.%s query error %v", targetName, cd.Namespace, err) } vs.Spec.Http = []istiov1alpha3.HTTPRoute{ { @@ -282,7 +284,7 @@ func (c *CanaryRouter) SetRoutes( vs, err = c.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Update(vs) if err != nil { - return fmt.Errorf("VirtualService %s.%s update failed: %v", cd.Name, cd.Namespace, err) + return fmt.Errorf("VirtualService %s.%s update failed: %v", targetName, cd.Namespace, err) } return nil diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index 8608d2bb..25a95892 100644 --- a/pkg/controller/scheduler.go +++ b/pkg/controller/scheduler.go @@ -12,7 +12,7 @@ import ( // for new canaries new jobs are created and started // for the removed canaries the jobs are stopped and deleted func (c *Controller) scheduleCanaries() { - current := make(map[string]bool) + current := make(map[string]string) stats := make(map[string]int) c.canaries.Range(func(key interface{}, value interface{}) bool { @@ -20,7 +20,7 @@ func (c *Controller) scheduleCanaries() { // format: . name := key.(string) - current[name] = true + current[name] = fmt.Sprintf("%s.%s", canary.Spec.TargetRef.Name, canary.Namespace) // schedule new jobs if _, exists := c.jobs[name]; !exists { @@ -54,6 +54,16 @@ func (c *Controller) scheduleCanaries() { } } + // check if multiple canaries have the same target + 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) + } + } + } + // set total canaries per namespace metric for k, v := range stats { c.recorder.SetTotal(k, v)