From 177dc824e3f4229ad9796b224bb042928209313f Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Mon, 6 May 2019 18:42:02 +0300 Subject: [PATCH] Implement nginx ingress router --- pkg/router/factory.go | 7 +- pkg/router/ingress.go | 165 ++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 171 insertions(+), 1 deletion(-) create mode 100644 pkg/router/ingress.go diff --git a/pkg/router/factory.go b/pkg/router/factory.go index 9d8f66ec..ab836d31 100644 --- a/pkg/router/factory.go +++ b/pkg/router/factory.go @@ -41,9 +41,14 @@ func (factory *Factory) KubernetesRouter(label string) *KubernetesRouter { } } -// MeshRouter returns a service mesh router (Istio or AppMesh) +// MeshRouter returns a service mesh router func (factory *Factory) MeshRouter(provider string) Interface { switch { + case provider == "nginx": + return &IngressRouter{ + logger: factory.logger, + kubeClient: factory.kubeClient, + } case provider == "appmesh": return &AppMeshRouter{ logger: factory.logger, diff --git a/pkg/router/ingress.go b/pkg/router/ingress.go new file mode 100644 index 00000000..7e9513e0 --- /dev/null +++ b/pkg/router/ingress.go @@ -0,0 +1,165 @@ +package router + +import ( + "fmt" + "github.com/google/go-cmp/cmp" + flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1alpha3" + "go.uber.org/zap" + "k8s.io/api/extensions/v1beta1" + "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" + "strconv" + "strings" +) + +type IngressRouter struct { + kubeClient kubernetes.Interface + logger *zap.SugaredLogger +} + +func (i *IngressRouter) Reconcile(canary *flaggerv1.Canary) error { + if canary.Spec.IngressRef == nil || canary.Spec.IngressRef.Name == "" { + return fmt.Errorf("ingress selector is empty") + } + + targetName := canary.Spec.TargetRef.Name + canaryName := fmt.Sprintf("%s-canary", targetName) + canaryIngressName := fmt.Sprintf("%s-canary", canary.Spec.IngressRef.Name) + + ingress, err := i.kubeClient.ExtensionsV1beta1().Ingresses(canary.Namespace).Get(canary.Spec.IngressRef.Name, metav1.GetOptions{}) + if err != nil { + return err + } + + ingressClone := ingress.DeepCopy() + + // change backend to -canary + backendExists := false + for k, v := range ingressClone.Spec.Rules { + for x, y := range v.HTTP.Paths { + if y.Backend.ServiceName == targetName { + ingressClone.Spec.Rules[k].HTTP.Paths[x].Backend.ServiceName = canaryName + backendExists = true + break + } + } + } + + if !backendExists { + return fmt.Errorf("backend %s not found in ingress %s", targetName, canary.Spec.IngressRef.Name) + } + + canaryIngress, err := i.kubeClient.ExtensionsV1beta1().Ingresses(canary.Namespace).Get(canaryIngressName, metav1.GetOptions{}) + + if errors.IsNotFound(err) { + ing := &v1beta1.Ingress{ + ObjectMeta: metav1.ObjectMeta{ + Name: canaryIngressName, + Namespace: canary.Namespace, + OwnerReferences: []metav1.OwnerReference{ + *metav1.NewControllerRef(canary, schema.GroupVersionKind{ + Group: flaggerv1.SchemeGroupVersion.Group, + Version: flaggerv1.SchemeGroupVersion.Version, + Kind: flaggerv1.CanaryKind, + }), + }, + Annotations: i.makeAnnotations(ingressClone.Annotations), + Labels: ingressClone.Labels, + }, + Spec: ingressClone.Spec, + } + + _, err := i.kubeClient.ExtensionsV1beta1().Ingresses(canary.Namespace).Create(ing) + if err != nil { + return err + } + + i.logger.With("canary", fmt.Sprintf("%s.%s", canary.Name, canary.Namespace)). + Infof("Ingress %s.%s created", ing.GetName(), canary.Namespace) + return nil + } + + if err != nil { + return fmt.Errorf("ingress %s query error %v", canaryIngressName, err) + } + + if diff := cmp.Diff(ingressClone.Spec, canaryIngress.Spec); diff != "" { + iClone := canaryIngress.DeepCopy() + iClone.Spec = ingressClone.Spec + + _, err := i.kubeClient.ExtensionsV1beta1().Ingresses(canary.Namespace).Update(iClone) + if err != nil { + return fmt.Errorf("ingress %s update error %v", canaryIngressName, err) + } + + i.logger.With("canary", fmt.Sprintf("%s.%s", canary.Name, canary.Namespace)). + Infof("Ingress %s updated", canaryIngressName) + } + + return nil +} + +func (i *IngressRouter) GetRoutes(canary *flaggerv1.Canary) ( + primaryWeight int, + canaryWeight int, + err error, +) { + canaryIngressName := fmt.Sprintf("%s-canary", canary.Spec.IngressRef.Name) + canaryIngress, err := i.kubeClient.ExtensionsV1beta1().Ingresses(canary.Namespace).Get(canaryIngressName, metav1.GetOptions{}) + if err != nil { + return 0, 0, err + } + + for k, v := range canaryIngress.Annotations { + if k == "nginx.ingress.kubernetes.io/canary-weight" { + val, err := strconv.Atoi(v) + if err != nil { + return 0, 0, err + } + + canaryWeight = val + break + } + } + + primaryWeight = 100 - canaryWeight + return +} + +func (i *IngressRouter) SetRoutes( + canary *flaggerv1.Canary, + primaryWeight int, + canaryWeight int, +) error { + canaryIngressName := fmt.Sprintf("%s-canary", canary.Spec.IngressRef.Name) + canaryIngress, err := i.kubeClient.ExtensionsV1beta1().Ingresses(canary.Namespace).Get(canaryIngressName, metav1.GetOptions{}) + if err != nil { + return err + } + + iClone := canaryIngress.DeepCopy() + iClone.Annotations["nginx.ingress.kubernetes.io/canary-weight"] = fmt.Sprintf("%v", canaryWeight) + + _, err = i.kubeClient.ExtensionsV1beta1().Ingresses(canary.Namespace).Update(iClone) + if err != nil { + return fmt.Errorf("ingress %s update error %v", canaryIngressName, err) + } + + return nil +} + +func (i *IngressRouter) makeAnnotations(annotations map[string]string) map[string]string { + res := make(map[string]string) + for k, v := range annotations { + if !strings.Contains(v, "nginx.ingress.kubernetes.io/canary") { + res[k] = v + } + } + + res["nginx.ingress.kubernetes.io/canary"] = "true" + res["nginx.ingress.kubernetes.io/canary-weight"] = "0" + + return res +}