From e1d8703a151e590b9742148bc2b86af86925a953 Mon Sep 17 00:00:00 2001 From: Yusuke Kuoka Date: Wed, 27 Nov 2019 22:39:49 +0900 Subject: [PATCH] Refactor to merge KubernetesServiceRouter into ServiceController The current design is that everything related to managing the targeted resource should go into the respective implementation of `canary.Controller`. In the service-canary use-case our target is Service so rather than splitting and scattering the logics over Controller and Router, everything should naturally go to `ServiceController`. Maybe at the time of writing the first implementation, I was confusing the target service vs the router. --- pkg/canary/service_controller.go | 113 +++++++++++++++++++++++++++- pkg/router/factory.go | 11 +-- pkg/router/kubernetes_noop.go | 14 ++++ pkg/router/kubernetes_service.go | 124 ------------------------------- 4 files changed, 127 insertions(+), 135 deletions(-) create mode 100644 pkg/router/kubernetes_noop.go delete mode 100644 pkg/router/kubernetes_service.go diff --git a/pkg/canary/service_controller.go b/pkg/canary/service_controller.go index fc1c0450..bd8a505b 100644 --- a/pkg/canary/service_controller.go +++ b/pkg/canary/service_controller.go @@ -5,8 +5,10 @@ import ( ex "github.com/pkg/errors" "go.uber.org/zap" + corev1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/client-go/kubernetes" flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1alpha3" @@ -42,12 +44,119 @@ func (c *ServiceController) SetStatusPhase(cd *flaggerv1.Canary, phase flaggerv1 var _ Controller = &ServiceController{} -// Initialize creates the primary deployment, hpa, -// scales to zero the canary deployment and returns the pod selector label and container ports +// 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) (label string, ports map[string]int32, 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 "", nil, err + } + + // canary svc + err = c.reconcileCanaryService(cd, canaryName, svc) + if err != nil { + return "", nil, err + } + + // primary svc + err = c.reconcilePrimaryService(cd, primaryName, svc) + if err != nil { + return "", nil, err + } + return "", nil, nil } +func (c *ServiceController) reconcileCanaryService(canary *flaggerv1.Canary, name string, src *corev1.Service) error { + current, err := c.kubeClient.CoreV1().Services(canary.Namespace).Get(name, metav1.GetOptions{}) + if errors.IsNotFound(err) { + return c.createService(canary, name, src) + } + + if err != nil { + return fmt.Errorf("service %s query error %v", name, err) + } + + new := buildService(canary, name, src) + + if new.Spec.Type == "ClusterIP" { + // We can't change this immutable field + new.Spec.ClusterIP = current.Spec.ClusterIP + } + + // We can't change this immutable field + new.ObjectMeta.UID = current.ObjectMeta.UID + + new.ObjectMeta.ResourceVersion = current.ObjectMeta.ResourceVersion + + _, err = c.kubeClient.CoreV1().Services(canary.Namespace).Update(new) + if err != nil { + return err + } + + c.logger.With("canary", fmt.Sprintf("%s.%s", canary.Name, canary.Namespace)). + Infof("Service %s.%s updated", new.GetName(), canary.Namespace) + return nil +} + +func (c *ServiceController) reconcilePrimaryService(canary *flaggerv1.Canary, name string, src *corev1.Service) error { + _, err := c.kubeClient.CoreV1().Services(canary.Namespace).Get(name, metav1.GetOptions{}) + if errors.IsNotFound(err) { + return c.createService(canary, name, src) + } + + if err != nil { + return fmt.Errorf("service %s query error %v", name, err) + } + + return nil +} + +func (c *ServiceController) createService(canary *flaggerv1.Canary, name string, src *corev1.Service) error { + svc := buildService(canary, name, src) + + if svc.Spec.Type == "ClusterIP" { + // Reset and let K8s assign the IP. Otherwise we get an error due to the IP is already assigned + svc.Spec.ClusterIP = "" + } + + // Let K8s set this. Otherwise K8s API complains with "resourceVersion should not be set on objects to be created" + svc.ObjectMeta.ResourceVersion = "" + + _, err := c.kubeClient.CoreV1().Services(canary.Namespace).Create(svc) + if err != nil { + return err + } + + c.logger.With("canary", fmt.Sprintf("%s.%s", canary.Name, canary.Namespace)). + Infof("Service %s.%s created", svc.GetName(), canary.Namespace) + return nil +} + +func buildService(canary *flaggerv1.Canary, name string, src *corev1.Service) *corev1.Service { + svc := src.DeepCopy() + svc.ObjectMeta.Name = name + svc.ObjectMeta.Namespace = canary.Namespace + svc.ObjectMeta.OwnerReferences = []metav1.OwnerReference{ + *metav1.NewControllerRef(canary, schema.GroupVersionKind{ + Group: flaggerv1.SchemeGroupVersion.Group, + Version: flaggerv1.SchemeGroupVersion.Version, + Kind: flaggerv1.CanaryKind, + }), + } + _, exists := svc.ObjectMeta.Annotations["kubectl.kubernetes.io/last-applied-configuration"] + if exists { + // Leaving this results in updates from flagger to this svc never succeed due to resourceVersion mismatch: + // 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 +} + // Promote copies target's spec from canary to primary func (c *ServiceController) Promote(cd *flaggerv1.Canary) error { targetName := cd.Spec.TargetRef.Name diff --git a/pkg/router/factory.go b/pkg/router/factory.go index a01cb09f..2e2803e5 100644 --- a/pkg/router/factory.go +++ b/pkg/router/factory.go @@ -44,20 +44,13 @@ func (factory *Factory) KubernetesRouter(kind string, labelSelector string, anno annotations: annotations, ports: ports, } - serviceRouter := &KubernetesServiceRouter{ - logger: factory.logger, - flaggerClient: factory.flaggerClient, - kubeClient: factory.kubeClient, - labelSelector: labelSelector, - annotations: annotations, - ports: ports, - } + noopRouter := &KubernetesNoopRouter{} switch { case kind == "Deployment": return deploymentRouter case kind == "Service": - return serviceRouter + return noopRouter default: return deploymentRouter } diff --git a/pkg/router/kubernetes_noop.go b/pkg/router/kubernetes_noop.go new file mode 100644 index 00000000..fa97362a --- /dev/null +++ b/pkg/router/kubernetes_noop.go @@ -0,0 +1,14 @@ +package router + +import ( + flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1alpha3" +) + +// KubernetesNoopRouter manages nothing. This is useful when one uses Flagger for progressive delivery of +// services that are not load-balanced by a Kubernetes service +type KubernetesNoopRouter struct { +} + +func (c *KubernetesNoopRouter) Reconcile(canary *flaggerv1.Canary) error { + return nil +} diff --git a/pkg/router/kubernetes_service.go b/pkg/router/kubernetes_service.go deleted file mode 100644 index 8ac0baf7..00000000 --- a/pkg/router/kubernetes_service.go +++ /dev/null @@ -1,124 +0,0 @@ -package router - -import ( - "fmt" - - "go.uber.org/zap" - corev1 "k8s.io/api/core/v1" - "k8s.io/apimachinery/pkg/api/errors" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/client-go/kubernetes" - - flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1alpha3" - clientset "github.com/weaveworks/flagger/pkg/client/clientset/versioned" -) - -// KubernetesServiceRouter manages ClusterIP services -type KubernetesServiceRouter struct { - kubeClient kubernetes.Interface - flaggerClient clientset.Interface - logger *zap.SugaredLogger - labelSelector string - annotations map[string]string - ports map[string]int32 -} - -// Reconcile creates or updates the primary and canary services to prepare for the canary release process targeted on the K8s service -func (c *KubernetesServiceRouter) Reconcile(canary *flaggerv1.Canary) error { - targetName := canary.Spec.TargetRef.Name - primaryName := fmt.Sprintf("%s-primary", targetName) - canaryName := fmt.Sprintf("%s-canary", targetName) - - svc, err := c.kubeClient.CoreV1().Services(canary.Namespace).Get(targetName, metav1.GetOptions{}) - if err != nil { - return err - } - - // canary svc - err = c.reconcileCanaryService(canary, canaryName, svc) - if err != nil { - return err - } - - // primary svc - err = c.reconcilePrimaryService(canary, primaryName, svc) - if err != nil { - return err - } - - return nil -} - -func (c *KubernetesServiceRouter) SetRoutes(canary *flaggerv1.Canary, primaryRoute int, canaryRoute int) error { - return nil -} - -func (c *KubernetesServiceRouter) GetRoutes(canary *flaggerv1.Canary) (primaryRoute int, canaryRoute int, err error) { - return 0, 0, nil -} - -func (c *KubernetesServiceRouter) reconcileCanaryService(canary *flaggerv1.Canary, name string, src *corev1.Service) error { - current, err := c.kubeClient.CoreV1().Services(canary.Namespace).Get(name, metav1.GetOptions{}) - if errors.IsNotFound(err) { - return c.createService(canary, name, src) - } - - if err != nil { - return fmt.Errorf("service %s query error %v", name, err) - } - - new := buildService(canary, name, src) - - if new.Spec.Type == "ClusterIP" { - // We can't change this immutable field - new.Spec.ClusterIP = current.Spec.ClusterIP - } - - // We can't change this immutable field - new.ObjectMeta.UID = current.ObjectMeta.UID - - new.ObjectMeta.ResourceVersion = current.ObjectMeta.ResourceVersion - - _, err = c.kubeClient.CoreV1().Services(canary.Namespace).Update(new) - if err != nil { - return err - } - - c.logger.With("canary", fmt.Sprintf("%s.%s", canary.Name, canary.Namespace)). - Infof("Service %s.%s updated", new.GetName(), canary.Namespace) - return nil -} - -func (c *KubernetesServiceRouter) reconcilePrimaryService(canary *flaggerv1.Canary, name string, src *corev1.Service) error { - _, err := c.kubeClient.CoreV1().Services(canary.Namespace).Get(name, metav1.GetOptions{}) - if errors.IsNotFound(err) { - return c.createService(canary, name, src) - } - - if err != nil { - return fmt.Errorf("service %s query error %v", name, err) - } - - return nil -} - -func (c *KubernetesServiceRouter) createService(canary *flaggerv1.Canary, name string, src *corev1.Service) error { - svc := buildService(canary, name, src) - - if svc.Spec.Type == "ClusterIP" { - // Reset and let K8s assign the IP. Otherwise we get an error due to the IP is already assigned - svc.Spec.ClusterIP = "" - } - - // Let K8s set this. Otherwise K8s API complains with "resourceVersion should not be set on objects to be created" - svc.ObjectMeta.ResourceVersion = "" - - _, err := c.kubeClient.CoreV1().Services(canary.Namespace).Create(svc) - if err != nil { - return err - } - - c.logger.With("canary", fmt.Sprintf("%s.%s", canary.Name, canary.Namespace)). - Infof("Service %s.%s created", svc.GetName(), canary.Namespace) - return nil -}