From 3f5c22d863a07bfa8ae5c53438d28ff3eea5c569 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Tue, 5 Mar 2019 02:02:58 +0200 Subject: [PATCH] Extract routing to dedicated package - split routing management into Kubernetes service router and Istio Virtual service router --- pkg/router/factory.go | 44 ++++++++ pkg/router/istio.go | 218 +++++++++++++++++++++++++++++++++++++ pkg/router/istio_test.go | 151 +++++++++++++++++++++++++ pkg/router/router.go | 9 ++ pkg/router/router_test.go | 127 +++++++++++++++++++++ pkg/router/service.go | 152 ++++++++++++++++++++++++++ pkg/router/service_test.go | 46 ++++++++ 7 files changed, 747 insertions(+) create mode 100644 pkg/router/factory.go create mode 100644 pkg/router/istio.go create mode 100644 pkg/router/istio_test.go create mode 100644 pkg/router/router.go create mode 100644 pkg/router/router_test.go create mode 100644 pkg/router/service.go create mode 100644 pkg/router/service_test.go diff --git a/pkg/router/factory.go b/pkg/router/factory.go new file mode 100644 index 00000000..3c73abd2 --- /dev/null +++ b/pkg/router/factory.go @@ -0,0 +1,44 @@ +package router + +import ( + istioclientset "github.com/knative/pkg/client/clientset/versioned" + clientset "github.com/stefanprodan/flagger/pkg/client/clientset/versioned" + "go.uber.org/zap" + "k8s.io/client-go/kubernetes" +) + +type Factory struct { + kubeClient kubernetes.Interface + istioClient istioclientset.Interface + flaggerClient clientset.Interface + logger *zap.SugaredLogger +} + +func NewFactory(kubeClient kubernetes.Interface, + flaggerClient clientset.Interface, + logger *zap.SugaredLogger, + istioClient istioclientset.Interface) *Factory { + return &Factory{ + istioClient: istioClient, + kubeClient: kubeClient, + flaggerClient: flaggerClient, + logger: logger, + } +} + +func (factory *Factory) ServiceRouter() *ServiceRouter { + return &ServiceRouter{ + logger: factory.logger, + flaggerClient: factory.flaggerClient, + kubeClient: factory.kubeClient, + } +} + +func (factory *Factory) IstioRouter() *IstioRouter { + return &IstioRouter{ + logger: factory.logger, + flaggerClient: factory.flaggerClient, + kubeClient: factory.kubeClient, + istioClient: factory.istioClient, + } +} diff --git a/pkg/router/istio.go b/pkg/router/istio.go new file mode 100644 index 00000000..17e39259 --- /dev/null +++ b/pkg/router/istio.go @@ -0,0 +1,218 @@ +package router + +import ( + "fmt" + "github.com/google/go-cmp/cmp" + "github.com/google/go-cmp/cmp/cmpopts" + istiov1alpha3 "github.com/knative/pkg/apis/istio/v1alpha3" + istioclientset "github.com/knative/pkg/client/clientset/versioned" + flaggerv1 "github.com/stefanprodan/flagger/pkg/apis/flagger/v1alpha3" + clientset "github.com/stefanprodan/flagger/pkg/client/clientset/versioned" + "go.uber.org/zap" + "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/apis/meta/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/client-go/kubernetes" +) + +// IstioRouter is managing Istio virtual services +type IstioRouter struct { + kubeClient kubernetes.Interface + istioClient istioclientset.Interface + flaggerClient clientset.Interface + logger *zap.SugaredLogger +} + +// Sync creates or updates the the Istio virtual service. +func (ir *IstioRouter) Sync(cd *flaggerv1.Canary) error { + targetName := cd.Spec.TargetRef.Name + primaryName := fmt.Sprintf("%s-primary", targetName) + hosts := append(cd.Spec.Service.Hosts, targetName) + gateways := cd.Spec.Service.Gateways + var hasMeshGateway bool + for _, g := range gateways { + if g == "mesh" { + hasMeshGateway = true + } + } + if !hasMeshGateway { + gateways = append(gateways, "mesh") + } + + route := []istiov1alpha3.DestinationWeight{ + { + Destination: istiov1alpha3.Destination{ + Host: primaryName, + Port: istiov1alpha3.PortSelector{ + Number: uint32(cd.Spec.Service.Port), + }, + }, + Weight: 100, + }, + { + Destination: istiov1alpha3.Destination{ + Host: fmt.Sprintf("%s-canary", targetName), + Port: istiov1alpha3.PortSelector{ + Number: uint32(cd.Spec.Service.Port), + }, + }, + Weight: 0, + }, + } + newSpec := istiov1alpha3.VirtualServiceSpec{ + Hosts: hosts, + Gateways: gateways, + Http: []istiov1alpha3.HTTPRoute{ + { + Match: cd.Spec.Service.Match, + Rewrite: cd.Spec.Service.Rewrite, + Timeout: cd.Spec.Service.Timeout, + Retries: cd.Spec.Service.Retries, + AppendHeaders: cd.Spec.Service.AppendHeaders, + Route: route, + }, + }, + } + + virtualService, err := ir.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Get(targetName, metav1.GetOptions{}) + // insert + if errors.IsNotFound(err) { + virtualService = &istiov1alpha3.VirtualService{ + ObjectMeta: metav1.ObjectMeta{ + Name: targetName, + Namespace: cd.Namespace, + OwnerReferences: []metav1.OwnerReference{ + *metav1.NewControllerRef(cd, schema.GroupVersionKind{ + Group: flaggerv1.SchemeGroupVersion.Group, + Version: flaggerv1.SchemeGroupVersion.Version, + Kind: flaggerv1.CanaryKind, + }), + }, + }, + Spec: newSpec, + } + _, err = ir.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Create(virtualService) + if err != nil { + return fmt.Errorf("VirtualService %s.%s create error %v", targetName, cd.Namespace, err) + } + ir.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)). + Infof("VirtualService %s.%s created", virtualService.GetName(), cd.Namespace) + return nil + } + + if err != nil { + return fmt.Errorf("VirtualService %s.%s query error %v", targetName, cd.Namespace, err) + } + + // update service but keep the original destination weights + if virtualService != nil { + if diff := cmp.Diff(newSpec, virtualService.Spec, cmpopts.IgnoreTypes(istiov1alpha3.DestinationWeight{})); diff != "" { + vtClone := virtualService.DeepCopy() + vtClone.Spec = newSpec + if len(virtualService.Spec.Http) > 0 { + vtClone.Spec.Http[0].Route = virtualService.Spec.Http[0].Route + } + _, err = ir.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Update(vtClone) + if err != nil { + return fmt.Errorf("VirtualService %s.%s update error %v", targetName, cd.Namespace, err) + } + ir.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)). + Infof("VirtualService %s.%s updated", virtualService.GetName(), cd.Namespace) + } + } + + return nil +} + +// GetRoutes returns the destinations weight for primary and canary +func (ir *IstioRouter) GetRoutes(cd *flaggerv1.Canary) ( + primaryWeight int, + canaryWeight int, + err error, +) { + targetName := cd.Spec.TargetRef.Name + vs := &istiov1alpha3.VirtualService{} + vs, err = ir.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Get(targetName, v1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + err = fmt.Errorf("VirtualService %s.%s not found", targetName, cd.Namespace) + return + } + 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", targetName) { + primaryWeight = route.Weight + } + if route.Destination.Host == fmt.Sprintf("%s-canary", targetName) { + canaryWeight = route.Weight + } + } + } + + if primaryWeight == 0 && canaryWeight == 0 { + err = fmt.Errorf("VirtualService %s.%s does not contain routes for %s-primary and %s-canary", + targetName, cd.Namespace, targetName, targetName) + } + + return +} + +// SetRoutes updates the destinations weight for primary and canary +func (ir *IstioRouter) SetRoutes( + cd *flaggerv1.Canary, + primaryWeight int, + canaryWeight int, +) error { + targetName := cd.Spec.TargetRef.Name + vs, err := ir.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Get(targetName, v1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return fmt.Errorf("VirtualService %s.%s not found", targetName, cd.Namespace) + + } + return fmt.Errorf("VirtualService %s.%s query error %v", targetName, cd.Namespace, err) + } + + vsCopy := vs.DeepCopy() + vsCopy.Spec.Http = []istiov1alpha3.HTTPRoute{ + { + Match: cd.Spec.Service.Match, + Rewrite: cd.Spec.Service.Rewrite, + Timeout: cd.Spec.Service.Timeout, + Retries: cd.Spec.Service.Retries, + AppendHeaders: cd.Spec.Service.AppendHeaders, + Route: []istiov1alpha3.DestinationWeight{ + { + Destination: istiov1alpha3.Destination{ + Host: fmt.Sprintf("%s-primary", targetName), + Port: istiov1alpha3.PortSelector{ + Number: uint32(cd.Spec.Service.Port), + }, + }, + Weight: primaryWeight, + }, + { + Destination: istiov1alpha3.Destination{ + Host: fmt.Sprintf("%s-canary", targetName), + Port: istiov1alpha3.PortSelector{ + Number: uint32(cd.Spec.Service.Port), + }, + }, + Weight: canaryWeight, + }, + }, + }, + } + + vs, err = ir.istioClient.NetworkingV1alpha3().VirtualServices(cd.Namespace).Update(vsCopy) + if err != nil { + return fmt.Errorf("VirtualService %s.%s update failed: %v", targetName, cd.Namespace, err) + + } + return nil +} diff --git a/pkg/router/istio_test.go b/pkg/router/istio_test.go new file mode 100644 index 00000000..a123fd7c --- /dev/null +++ b/pkg/router/istio_test.go @@ -0,0 +1,151 @@ +package router + +import ( + "fmt" + istiov1alpha3 "github.com/knative/pkg/apis/istio/v1alpha3" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "testing" +) + +func TestIstioRouter_Sync(t *testing.T) { + mocks := setupfakeClients() + router := &IstioRouter{ + logger: mocks.logger, + flaggerClient: mocks.flaggerClient, + istioClient: mocks.istioClient, + kubeClient: mocks.kubeClient, + } + + err := router.Sync(mocks.canary) + if err != nil { + t.Fatal(err.Error()) + } + + // test insert + vs, err := mocks.istioClient.NetworkingV1alpha3().VirtualServices("default").Get("podinfo", metav1.GetOptions{}) + if err != nil { + t.Fatal(err.Error()) + } + + if len(vs.Spec.Http) != 1 { + t.Errorf("Got Istio VS Http %v wanted %v", len(vs.Spec.Http), 1) + } + + if len(vs.Spec.Http[0].Route) != 2 { + t.Errorf("Got Istio VS routes %v wanted %v", len(vs.Spec.Http[0].Route), 2) + } + + // test update + cd, err := mocks.flaggerClient.FlaggerV1alpha3().Canaries("default").Get("podinfo", metav1.GetOptions{}) + if err != nil { + t.Fatal(err.Error()) + } + + cdClone := cd.DeepCopy() + hosts := cdClone.Spec.Service.Hosts + hosts = append(hosts, "test.example.com") + cdClone.Spec.Service.Hosts = hosts + canary, err := mocks.flaggerClient.FlaggerV1alpha3().Canaries("default").Update(cdClone) + if err != nil { + t.Fatal(err.Error()) + } + + // apply change + err = router.Sync(canary) + if err != nil { + t.Fatal(err.Error()) + } + + // verify + vs, err = mocks.istioClient.NetworkingV1alpha3().VirtualServices("default").Get("podinfo", metav1.GetOptions{}) + if err != nil { + t.Fatal(err.Error()) + } + if len(vs.Spec.Hosts) != 2 { + t.Errorf("Got Istio VS hosts %v wanted %v", vs.Spec.Hosts, 2) + } + + // test drift + vsClone := vs.DeepCopy() + gateways := vsClone.Spec.Gateways + gateways = append(gateways, "test-gateway.istio-system") + vsClone.Spec.Gateways = gateways + + vsGateways, err := mocks.istioClient.NetworkingV1alpha3().VirtualServices("default").Update(vsClone) + if err != nil { + t.Fatal(err.Error()) + } + if len(vsGateways.Spec.Gateways) != 2 { + t.Errorf("Got Istio VS gateway %v wanted %v", vsGateways.Spec.Gateways, 2) + } + + // undo change + err = router.Sync(mocks.canary) + if err != nil { + t.Fatal(err.Error()) + } + + // verify + vs, err = mocks.istioClient.NetworkingV1alpha3().VirtualServices("default").Get("podinfo", metav1.GetOptions{}) + if err != nil { + t.Fatal(err.Error()) + } + if len(vs.Spec.Gateways) != 1 { + t.Errorf("Got Istio VS gateways %v wanted %v", vs.Spec.Gateways, 1) + } +} + +func TestIstioRouter_SetRoutes(t *testing.T) { + mocks := setupfakeClients() + router := &IstioRouter{ + logger: mocks.logger, + flaggerClient: mocks.flaggerClient, + istioClient: mocks.istioClient, + kubeClient: mocks.kubeClient, + } + + err := router.Sync(mocks.canary) + if err != nil { + t.Fatal(err.Error()) + } + + p, c, err := router.GetRoutes(mocks.canary) + if err != nil { + t.Fatal(err.Error()) + } + + p = 50 + c = 50 + + err = router.SetRoutes(mocks.canary, p, c) + if err != nil { + t.Fatal(err.Error()) + } + + vs, err := mocks.istioClient.NetworkingV1alpha3().VirtualServices("default").Get("podinfo", metav1.GetOptions{}) + if err != nil { + t.Fatal(err.Error()) + } + + pRoute := istiov1alpha3.DestinationWeight{} + cRoute := istiov1alpha3.DestinationWeight{} + + for _, http := range vs.Spec.Http { + for _, route := range http.Route { + if route.Destination.Host == fmt.Sprintf("%s-primary", mocks.canary.Spec.TargetRef.Name) { + pRoute = route + } + if route.Destination.Host == fmt.Sprintf("%s-canary", mocks.canary.Spec.TargetRef.Name) { + cRoute = route + } + } + } + + if pRoute.Weight != p { + t.Errorf("Got primary weight %v wanted %v", pRoute.Weight, p) + } + + if cRoute.Weight != c { + t.Errorf("Got canary weight %v wanted %v", cRoute.Weight, c) + } +} diff --git a/pkg/router/router.go b/pkg/router/router.go new file mode 100644 index 00000000..21686957 --- /dev/null +++ b/pkg/router/router.go @@ -0,0 +1,9 @@ +package router + +import flaggerv1 "github.com/stefanprodan/flagger/pkg/apis/flagger/v1alpha3" + +type Interface interface { + Sync(canary *flaggerv1.Canary) error + SetRoutes(canary *flaggerv1.Canary, primaryWeight int, canaryWeight int) error + GetRoutes(canary *flaggerv1.Canary) (primaryWeight int, canaryWeight int, err error) +} diff --git a/pkg/router/router_test.go b/pkg/router/router_test.go new file mode 100644 index 00000000..c89095df --- /dev/null +++ b/pkg/router/router_test.go @@ -0,0 +1,127 @@ +package router + +import ( + istioclientset "github.com/knative/pkg/client/clientset/versioned" + fakeIstio "github.com/knative/pkg/client/clientset/versioned/fake" + "github.com/stefanprodan/flagger/pkg/apis/flagger/v1alpha3" + clientset "github.com/stefanprodan/flagger/pkg/client/clientset/versioned" + fakeFlagger "github.com/stefanprodan/flagger/pkg/client/clientset/versioned/fake" + "github.com/stefanprodan/flagger/pkg/logging" + "go.uber.org/zap" + appsv1 "k8s.io/api/apps/v1" + hpav1 "k8s.io/api/autoscaling/v1" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/kubernetes/fake" +) + +type fakeClients struct { + canary *v1alpha3.Canary + kubeClient kubernetes.Interface + istioClient istioclientset.Interface + flaggerClient clientset.Interface + logger *zap.SugaredLogger +} + +func setupfakeClients() fakeClients { + canary := newMockCanary() + flaggerClient := fakeFlagger.NewSimpleClientset(canary) + + kubeClient := fake.NewSimpleClientset( + newMockDeployment(), + ) + + istioClient := fakeIstio.NewSimpleClientset() + logger, _ := logging.NewLogger("debug") + + return fakeClients{ + canary: canary, + kubeClient: kubeClient, + istioClient: istioClient, + flaggerClient: flaggerClient, + logger: logger, + } +} + +func newMockCanary() *v1alpha3.Canary { + cd := &v1alpha3.Canary{ + TypeMeta: metav1.TypeMeta{APIVersion: v1alpha3.SchemeGroupVersion.String()}, + ObjectMeta: metav1.ObjectMeta{ + Namespace: "default", + Name: "podinfo", + }, + Spec: v1alpha3.CanarySpec{ + TargetRef: hpav1.CrossVersionObjectReference{ + Name: "podinfo", + APIVersion: "apps/v1", + Kind: "Deployment", + }, + Service: v1alpha3.CanaryService{ + Port: 9898, + }, CanaryAnalysis: v1alpha3.CanaryAnalysis{ + Threshold: 10, + StepWeight: 10, + MaxWeight: 50, + Metrics: []v1alpha3.CanaryMetric{ + { + Name: "istio_requests_total", + Threshold: 99, + Interval: "1m", + }, + { + Name: "istio_request_duration_seconds_bucket", + Threshold: 500, + Interval: "1m", + }, + }, + }, + }, + } + return cd +} + +func newMockDeployment() *appsv1.Deployment { + d := &appsv1.Deployment{ + TypeMeta: metav1.TypeMeta{APIVersion: appsv1.SchemeGroupVersion.String()}, + ObjectMeta: metav1.ObjectMeta{ + Namespace: "default", + Name: "podinfo", + }, + Spec: appsv1.DeploymentSpec{ + Selector: &metav1.LabelSelector{ + MatchLabels: map[string]string{ + "app": "podinfo", + }, + }, + Template: corev1.PodTemplateSpec{ + ObjectMeta: metav1.ObjectMeta{ + Labels: map[string]string{ + "app": "podinfo", + }, + }, + Spec: corev1.PodSpec{ + Containers: []corev1.Container{ + { + Name: "podinfo", + Image: "quay.io/stefanprodan/podinfo:1.4.0", + Command: []string{ + "./podinfo", + "--port=9898", + }, + Ports: []corev1.ContainerPort{ + { + Name: "http", + ContainerPort: 9898, + Protocol: corev1.ProtocolTCP, + }, + }, + }, + }, + }, + }, + }, + } + + return d +} diff --git a/pkg/router/service.go b/pkg/router/service.go new file mode 100644 index 00000000..ab2ac01b --- /dev/null +++ b/pkg/router/service.go @@ -0,0 +1,152 @@ +package router + +import ( + "fmt" + flaggerv1 "github.com/stefanprodan/flagger/pkg/apis/flagger/v1alpha3" + clientset "github.com/stefanprodan/flagger/pkg/client/clientset/versioned" + "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/apimachinery/pkg/util/intstr" + "k8s.io/client-go/kubernetes" +) + +// ServiceRouter is managing ClusterIP services +type ServiceRouter struct { + kubeClient kubernetes.Interface + flaggerClient clientset.Interface + logger *zap.SugaredLogger +} + +// Sync creates to updates the primary and canary services +func (c *ServiceRouter) Sync(cd *flaggerv1.Canary) error { + 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: targetName, + Namespace: cd.Namespace, + OwnerReferences: []metav1.OwnerReference{ + *metav1.NewControllerRef(cd, schema.GroupVersionKind{ + Group: flaggerv1.SchemeGroupVersion.Group, + Version: flaggerv1.SchemeGroupVersion.Version, + Kind: flaggerv1.CanaryKind, + }), + }, + }, + Spec: corev1.ServiceSpec{ + Type: corev1.ServiceTypeClusterIP, + Selector: map[string]string{"app": targetName}, + Ports: []corev1.ServicePort{ + { + Name: "http", + Protocol: corev1.ProtocolTCP, + Port: cd.Spec.Service.Port, + TargetPort: intstr.IntOrString{ + Type: intstr.Int, + IntVal: cd.Spec.Service.Port, + }, + }, + }, + }, + } + + _, err = c.kubeClient.CoreV1().Services(cd.Namespace).Create(canaryService) + if err != nil { + return err + } + c.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("Service %s.%s created", canaryService.GetName(), cd.Namespace) + } + + canaryTestServiceName := fmt.Sprintf("%s-canary", cd.Spec.TargetRef.Name) + canaryTestService, err := c.kubeClient.CoreV1().Services(cd.Namespace).Get(canaryTestServiceName, metav1.GetOptions{}) + if errors.IsNotFound(err) { + canaryTestService = &corev1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: canaryTestServiceName, + Namespace: cd.Namespace, + OwnerReferences: []metav1.OwnerReference{ + *metav1.NewControllerRef(cd, schema.GroupVersionKind{ + Group: flaggerv1.SchemeGroupVersion.Group, + Version: flaggerv1.SchemeGroupVersion.Version, + Kind: flaggerv1.CanaryKind, + }), + }, + }, + Spec: corev1.ServiceSpec{ + Type: corev1.ServiceTypeClusterIP, + Selector: map[string]string{"app": targetName}, + Ports: []corev1.ServicePort{ + { + Name: "http", + Protocol: corev1.ProtocolTCP, + Port: cd.Spec.Service.Port, + TargetPort: intstr.IntOrString{ + Type: intstr.Int, + IntVal: cd.Spec.Service.Port, + }, + }, + }, + }, + } + + _, err = c.kubeClient.CoreV1().Services(cd.Namespace).Create(canaryTestService) + if err != nil { + return err + } + c.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("Service %s.%s created", canaryTestService.GetName(), cd.Namespace) + } + + primaryService, err := c.kubeClient.CoreV1().Services(cd.Namespace).Get(primaryName, metav1.GetOptions{}) + if errors.IsNotFound(err) { + primaryService = &corev1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: primaryName, + Namespace: cd.Namespace, + OwnerReferences: []metav1.OwnerReference{ + *metav1.NewControllerRef(cd, schema.GroupVersionKind{ + Group: flaggerv1.SchemeGroupVersion.Group, + Version: flaggerv1.SchemeGroupVersion.Version, + Kind: flaggerv1.CanaryKind, + }), + }, + }, + Spec: corev1.ServiceSpec{ + Type: corev1.ServiceTypeClusterIP, + Selector: map[string]string{"app": primaryName}, + Ports: []corev1.ServicePort{ + { + Name: "http", + Protocol: corev1.ProtocolTCP, + Port: cd.Spec.Service.Port, + TargetPort: intstr.IntOrString{ + Type: intstr.Int, + IntVal: cd.Spec.Service.Port, + }, + }, + }, + }, + } + + _, err = c.kubeClient.CoreV1().Services(cd.Namespace).Create(primaryService) + if err != nil { + return err + } + + c.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("Service %s.%s created", primaryService.GetName(), cd.Namespace) + } + + return nil +} + +func (c *ServiceRouter) SetRoutes(canary flaggerv1.Canary, primaryRoute int, canaryRoute int) error { + return nil +} + +func (c *ServiceRouter) GetRoutes(canary flaggerv1.Canary) (primaryRoute int, canaryRoute int, err error) { + return 0, 0, nil +} diff --git a/pkg/router/service_test.go b/pkg/router/service_test.go new file mode 100644 index 00000000..2cae1aa3 --- /dev/null +++ b/pkg/router/service_test.go @@ -0,0 +1,46 @@ +package router + +import ( + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "testing" +) + +func TestServiceRouter_Sync(t *testing.T) { + mocks := setupfakeClients() + router := &ServiceRouter{ + kubeClient: mocks.kubeClient, + flaggerClient: mocks.flaggerClient, + logger: mocks.logger, + } + + err := router.Sync(mocks.canary) + if err != nil { + t.Fatal(err.Error()) + } + + canarySvc, err := mocks.kubeClient.CoreV1().Services("default").Get("podinfo-canary", metav1.GetOptions{}) + if err != nil { + t.Fatal(err.Error()) + } + + if canarySvc.Spec.Ports[0].Name != "http" { + t.Errorf("Got svc port name %s wanted %s", canarySvc.Spec.Ports[0].Name, "http") + } + + if canarySvc.Spec.Ports[0].Port != 9898 { + t.Errorf("Got svc port %v wanted %v", canarySvc.Spec.Ports[0].Port, 9898) + } + + primarySvc, err := mocks.kubeClient.CoreV1().Services("default").Get("podinfo-primary", metav1.GetOptions{}) + if err != nil { + t.Fatal(err.Error()) + } + + if primarySvc.Spec.Ports[0].Name != "http" { + t.Errorf("Got primary svc port name %s wanted %s", primarySvc.Spec.Ports[0].Name, "http") + } + + if primarySvc.Spec.Ports[0].Port != 9898 { + t.Errorf("Got primary svc port %v wanted %v", primarySvc.Spec.Ports[0].Port, 9898) + } +}