From 673b6102a7379b99b32296dfa6574255054833cd Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Wed, 9 Oct 2019 10:59:50 +0300 Subject: [PATCH] Add the name label to ClusterIP services and primary deployment --- pkg/canary/deployer.go | 11 +++++++---- pkg/controller/scheduler.go | 8 ++++---- pkg/router/factory.go | 5 +++-- pkg/router/kubernetes.go | 37 +++++++++++++++++++------------------ 4 files changed, 33 insertions(+), 28 deletions(-) diff --git a/pkg/canary/deployer.go b/pkg/canary/deployer.go index 9f35254d..ef976637 100644 --- a/pkg/canary/deployer.go +++ b/pkg/canary/deployer.go @@ -32,7 +32,7 @@ type Deployer struct { // 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 *Deployer) 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 { @@ -185,7 +185,7 @@ func (c *Deployer) Scale(cd *flaggerv1.Canary, replicas int32) error { return nil } -func (c *Deployer) createPrimaryDeployment(cd *flaggerv1.Canary) (string, *map[string]int32, error) { +func (c *Deployer) createPrimaryDeployment(cd *flaggerv1.Canary) (string, map[string]int32, error) { targetName := cd.Spec.TargetRef.Name primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) @@ -203,13 +203,13 @@ func (c *Deployer) createPrimaryDeployment(cd *flaggerv1.Canary) (string, *map[s targetName, cd.Namespace, targetName) } - var ports *map[string]int32 + var ports map[string]int32 if cd.Spec.Service.PortDiscovery { p, err := c.getPorts(cd, canaryDep) if err != nil { return "", nil, fmt.Errorf("port discovery failed with error: %v", err) } - ports = &p + ports = p } primaryDep, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(primaryName, metav1.GetOptions{}) @@ -237,6 +237,9 @@ func (c *Deployer) createPrimaryDeployment(cd *flaggerv1.Canary) (string, *map[s ObjectMeta: metav1.ObjectMeta{ Name: primaryName, Namespace: cd.Namespace, + Labels: map[string]string{ + label: primaryName, + }, OwnerReferences: []metav1.OwnerReference{ *metav1.NewControllerRef(cd, schema.GroupVersionKind{ Group: flaggerv1.SchemeGroupVersion.Group, diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index fa9b0172..8eeb4c61 100644 --- a/pkg/controller/scheduler.go +++ b/pkg/controller/scheduler.go @@ -102,7 +102,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh if skipLivenessChecks || strings.Contains(provider, "istio") || strings.Contains(provider, "appmesh") { skipPrimaryCheck = true } - label, ports, err := c.deployer.Initialize(cd, skipPrimaryCheck) + labelSelector, ports, err := c.deployer.Initialize(cd, skipPrimaryCheck) if err != nil { c.recordEventWarningf(cd, "%v", err) return @@ -112,7 +112,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh meshRouter := c.routerFactory.MeshRouter(provider) // create or update ClusterIP services - if err := c.routerFactory.KubernetesRouter(label, ports).Reconcile(cd); err != nil { + if err := c.routerFactory.KubernetesRouter(labelSelector, map[string]string{}, ports).Reconcile(cd); err != nil { c.recordEventWarningf(cd, "%v", err) return } @@ -123,6 +123,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh return } + // check for deployment spec or configs changes shouldAdvance, err := c.shouldAdvance(cd) if err != nil { c.recordEventWarningf(cd, "%v", err) @@ -153,8 +154,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh } } - // check if virtual service exists - // and if it contains weighted destination routes to the primary and canary services + // get the routing settings primaryWeight, canaryWeight, mirrored, err := meshRouter.GetRoutes(cd) if err != nil { c.recordEventWarningf(cd, "%v", err) diff --git a/pkg/router/factory.go b/pkg/router/factory.go index 332095bc..80ab2471 100644 --- a/pkg/router/factory.go +++ b/pkg/router/factory.go @@ -36,12 +36,13 @@ func NewFactory(kubeConfig *restclient.Config, kubeClient kubernetes.Interface, } // KubernetesRouter returns a ClusterIP service router -func (factory *Factory) KubernetesRouter(label string, ports *map[string]int32) *KubernetesRouter { +func (factory *Factory) KubernetesRouter(labelSelector string, annotations map[string]string, ports map[string]int32) *KubernetesRouter { return &KubernetesRouter{ logger: factory.logger, flaggerClient: factory.flaggerClient, kubeClient: factory.kubeClient, - label: label, + labelSelector: labelSelector, + annotations: annotations, ports: ports, } } diff --git a/pkg/router/kubernetes.go b/pkg/router/kubernetes.go index c64f5270..5fc29ded 100644 --- a/pkg/router/kubernetes.go +++ b/pkg/router/kubernetes.go @@ -22,8 +22,9 @@ type KubernetesRouter struct { kubeClient kubernetes.Interface flaggerClient clientset.Interface logger *zap.SugaredLogger - label string - ports *map[string]int32 + labelSelector string + annotations map[string]string + ports map[string]int32 } // Reconcile creates or updates the primary and canary services @@ -78,7 +79,7 @@ func (c *KubernetesRouter) reconcileService(canary *flaggerv1.Canary, name strin svcSpec := corev1.ServiceSpec{ Type: corev1.ServiceTypeClusterIP, - Selector: map[string]string{c.label: target}, + Selector: map[string]string{c.labelSelector: target}, Ports: []corev1.ServicePort{ { Name: portName, @@ -89,28 +90,28 @@ func (c *KubernetesRouter) reconcileService(canary *flaggerv1.Canary, name strin }, } - if c.ports != nil { - for n, p := range *c.ports { - cp := corev1.ServicePort{ - Name: n, - Protocol: corev1.ProtocolTCP, - Port: p, - TargetPort: intstr.IntOrString{ - Type: intstr.Int, - IntVal: p, - }, - } - - svcSpec.Ports = append(svcSpec.Ports, cp) + for n, p := range c.ports { + cp := corev1.ServicePort{ + Name: n, + Protocol: corev1.ProtocolTCP, + Port: p, + TargetPort: intstr.IntOrString{ + Type: intstr.Int, + IntVal: p, + }, } + + svcSpec.Ports = append(svcSpec.Ports, cp) } svc, err := c.kubeClient.CoreV1().Services(canary.Namespace).Get(name, metav1.GetOptions{}) if errors.IsNotFound(err) { svc = &corev1.Service{ ObjectMeta: metav1.ObjectMeta{ - Name: name, - Namespace: canary.Namespace, + Name: name, + Namespace: canary.Namespace, + Labels: map[string]string{c.labelSelector: name}, + Annotations: c.annotations, OwnerReferences: []metav1.OwnerReference{ *metav1.NewControllerRef(canary, schema.GroupVersionKind{ Group: flaggerv1.SchemeGroupVersion.Group,