From 0e8b1ef20f2d3176b0ed8c764c47b0e70a5c80f8 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sat, 11 May 2019 12:49:23 +0300 Subject: [PATCH 01/14] Generate the SMI TrafficSplit clientset --- hack/update-codegen.sh | 2 +- pkg/apis/smi/register.go | 5 + pkg/apis/smi/v1alpha1/doc.go | 4 + pkg/apis/smi/v1alpha1/register.go | 48 +++++ pkg/apis/smi/v1alpha1/traffic_split.go | 56 ++++++ .../smi/v1alpha1/zz_generated.deepcopy.go | 129 +++++++++++++ pkg/client/clientset/versioned/clientset.go | 22 +++ .../versioned/fake/clientset_generated.go | 12 ++ .../clientset/versioned/fake/register.go | 2 + .../clientset/versioned/scheme/register.go | 2 + .../versioned/typed/smi/v1alpha1/doc.go | 20 ++ .../versioned/typed/smi/v1alpha1/fake/doc.go | 20 ++ .../smi/v1alpha1/fake/fake_smi_client.go | 40 ++++ .../smi/v1alpha1/fake/fake_trafficsplit.go | 128 +++++++++++++ .../typed/smi/v1alpha1/generated_expansion.go | 21 +++ .../typed/smi/v1alpha1/smi_client.go | 90 +++++++++ .../typed/smi/v1alpha1/trafficsplit.go | 174 ++++++++++++++++++ .../informers/externalversions/factory.go | 6 + .../informers/externalversions/generic.go | 5 + .../externalversions/smi/interface.go | 46 +++++ .../smi/v1alpha1/interface.go | 45 +++++ .../smi/v1alpha1/trafficsplit.go | 89 +++++++++ .../smi/v1alpha1/expansion_generated.go | 27 +++ .../listers/smi/v1alpha1/trafficsplit.go | 94 ++++++++++ 24 files changed, 1086 insertions(+), 1 deletion(-) create mode 100644 pkg/apis/smi/register.go create mode 100644 pkg/apis/smi/v1alpha1/doc.go create mode 100644 pkg/apis/smi/v1alpha1/register.go create mode 100644 pkg/apis/smi/v1alpha1/traffic_split.go create mode 100644 pkg/apis/smi/v1alpha1/zz_generated.deepcopy.go create mode 100644 pkg/client/clientset/versioned/typed/smi/v1alpha1/doc.go create mode 100644 pkg/client/clientset/versioned/typed/smi/v1alpha1/fake/doc.go create mode 100644 pkg/client/clientset/versioned/typed/smi/v1alpha1/fake/fake_smi_client.go create mode 100644 pkg/client/clientset/versioned/typed/smi/v1alpha1/fake/fake_trafficsplit.go create mode 100644 pkg/client/clientset/versioned/typed/smi/v1alpha1/generated_expansion.go create mode 100644 pkg/client/clientset/versioned/typed/smi/v1alpha1/smi_client.go create mode 100644 pkg/client/clientset/versioned/typed/smi/v1alpha1/trafficsplit.go create mode 100644 pkg/client/informers/externalversions/smi/interface.go create mode 100644 pkg/client/informers/externalversions/smi/v1alpha1/interface.go create mode 100644 pkg/client/informers/externalversions/smi/v1alpha1/trafficsplit.go create mode 100644 pkg/client/listers/smi/v1alpha1/expansion_generated.go create mode 100644 pkg/client/listers/smi/v1alpha1/trafficsplit.go diff --git a/hack/update-codegen.sh b/hack/update-codegen.sh index 4cdcbb36..5ca27f13 100755 --- a/hack/update-codegen.sh +++ b/hack/update-codegen.sh @@ -23,5 +23,5 @@ CODEGEN_PKG=${CODEGEN_PKG:-$(cd ${SCRIPT_ROOT}; ls -d -1 ./vendor/k8s.io/code-ge ${CODEGEN_PKG}/generate-groups.sh "deepcopy,client,informer,lister" \ github.com/weaveworks/flagger/pkg/client github.com/weaveworks/flagger/pkg/apis \ - "appmesh:v1beta1 istio:v1alpha3 flagger:v1alpha3" \ + "appmesh:v1beta1 istio:v1alpha3 flagger:v1alpha3 smi:v1alpha1" \ --go-header-file ${SCRIPT_ROOT}/hack/boilerplate.go.txt diff --git a/pkg/apis/smi/register.go b/pkg/apis/smi/register.go new file mode 100644 index 00000000..67cbbf79 --- /dev/null +++ b/pkg/apis/smi/register.go @@ -0,0 +1,5 @@ +package smi + +const ( + GroupName = "split.smi-spec.io" +) diff --git a/pkg/apis/smi/v1alpha1/doc.go b/pkg/apis/smi/v1alpha1/doc.go new file mode 100644 index 00000000..7792b7f6 --- /dev/null +++ b/pkg/apis/smi/v1alpha1/doc.go @@ -0,0 +1,4 @@ +// +k8s:deepcopy-gen=package +// +groupName=split.smi-spec.io + +package v1alpha1 diff --git a/pkg/apis/smi/v1alpha1/register.go b/pkg/apis/smi/v1alpha1/register.go new file mode 100644 index 00000000..c9868ec5 --- /dev/null +++ b/pkg/apis/smi/v1alpha1/register.go @@ -0,0 +1,48 @@ +package v1alpha1 + +import ( + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" + + ts "github.com/weaveworks/flagger/pkg/apis/smi" +) + +// SchemeGroupVersion is the identifier for the API which includes +// the name of the group and the version of the API +var SchemeGroupVersion = schema.GroupVersion{ + Group: ts.GroupName, + Version: "v1alpha1", +} + +// Kind takes an unqualified kind and returns back a Group qualified GroupKind +func Kind(kind string) schema.GroupKind { + return SchemeGroupVersion.WithKind(kind).GroupKind() +} + +// Resource takes an unqualified resource and returns a Group qualified GroupResource +func Resource(resource string) schema.GroupResource { + return SchemeGroupVersion.WithResource(resource).GroupResource() +} + +var ( + // SchemeBuilder collects functions that add things to a scheme. It's to allow + // code to compile without explicitly referencing generated types. You should + // declare one in each package that will have generated deep copy or conversion + // functions. + SchemeBuilder = runtime.NewSchemeBuilder(addKnownTypes) + + // AddToScheme applies all the stored functions to the scheme. A non-nil error + // indicates that one function failed and the attempt was abandoned. + AddToScheme = SchemeBuilder.AddToScheme +) + +// Adds the list of known types to Scheme. +func addKnownTypes(scheme *runtime.Scheme) error { + scheme.AddKnownTypes(SchemeGroupVersion, + &TrafficSplit{}, + &TrafficSplitList{}, + ) + metav1.AddToGroupVersion(scheme, SchemeGroupVersion) + return nil +} diff --git a/pkg/apis/smi/v1alpha1/traffic_split.go b/pkg/apis/smi/v1alpha1/traffic_split.go new file mode 100644 index 00000000..72574832 --- /dev/null +++ b/pkg/apis/smi/v1alpha1/traffic_split.go @@ -0,0 +1,56 @@ +package v1alpha1 + +import ( + "k8s.io/apimachinery/pkg/api/resource" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +// +genclient +// +genclient:noStatus +// +k8s:deepcopy-gen:interfaces=k8s.io/apimachinery/pkg/runtime.Object + +// TrafficSplit allows users to incrementally direct percentages of traffic +// between various services. It will be used by clients such as ingress +// controllers or service mesh sidecars to split the outgoing traffic to +// different destinations. +type TrafficSplit struct { + metav1.TypeMeta `json:",inline"` + // Standard object's metadata. + // More info: https://git.k8s.io/community/contributors/devel/api-conventions.md#metadata + // +optional + metav1.ObjectMeta `json:"metadata,omitempty" protobuf:"bytes,1,opt,name=metadata"` + + // Specification of the desired behavior of the traffic split. + // More info: https://git.k8s.io/community/contributors/devel/api-conventions.md#spec-and-status + // +optional + Spec TrafficSplitSpec `json:"spec,omitempty" protobuf:"bytes,2,opt,name=spec"` + + // Most recently observed status of the pod. + // This data may not be up to date. + // Populated by the system. + // Read-only. + // More info: https://git.k8s.io/community/contributors/devel/api-conventions.md#spec-and-status + // +optional + //Status Status `json:"status,omitempty" protobuf:"bytes,3,opt,name=status"` +} + +// TrafficSplitSpec is the specification for a TrafficSplit +type TrafficSplitSpec struct { + Service string `json:"service,omitempty"` + Backends []TrafficSplitBackend `json:"backends,omitempty"` +} + +// TrafficSplitBackend defines a backend +type TrafficSplitBackend struct { + Service string `json:"service,omitempty"` + Weight *resource.Quantity `json:"weight,omitempty"` +} + +// +k8s:deepcopy-gen:interfaces=k8s.io/apimachinery/pkg/runtime.Object + +type TrafficSplitList struct { + metav1.TypeMeta `json:",inline"` + metav1.ListMeta `json:"metadata"` + + Items []TrafficSplit `json:"items"` +} diff --git a/pkg/apis/smi/v1alpha1/zz_generated.deepcopy.go b/pkg/apis/smi/v1alpha1/zz_generated.deepcopy.go new file mode 100644 index 00000000..96694d48 --- /dev/null +++ b/pkg/apis/smi/v1alpha1/zz_generated.deepcopy.go @@ -0,0 +1,129 @@ +// +build !ignore_autogenerated + +/* +Copyright The Flagger Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +// Code generated by deepcopy-gen. DO NOT EDIT. + +package v1alpha1 + +import ( + runtime "k8s.io/apimachinery/pkg/runtime" +) + +// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. +func (in *TrafficSplit) DeepCopyInto(out *TrafficSplit) { + *out = *in + out.TypeMeta = in.TypeMeta + in.ObjectMeta.DeepCopyInto(&out.ObjectMeta) + in.Spec.DeepCopyInto(&out.Spec) + return +} + +// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new TrafficSplit. +func (in *TrafficSplit) DeepCopy() *TrafficSplit { + if in == nil { + return nil + } + out := new(TrafficSplit) + in.DeepCopyInto(out) + return out +} + +// DeepCopyObject is an autogenerated deepcopy function, copying the receiver, creating a new runtime.Object. +func (in *TrafficSplit) DeepCopyObject() runtime.Object { + if c := in.DeepCopy(); c != nil { + return c + } + return nil +} + +// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. +func (in *TrafficSplitBackend) DeepCopyInto(out *TrafficSplitBackend) { + *out = *in + if in.Weight != nil { + in, out := &in.Weight, &out.Weight + x := (*in).DeepCopy() + *out = &x + } + return +} + +// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new TrafficSplitBackend. +func (in *TrafficSplitBackend) DeepCopy() *TrafficSplitBackend { + if in == nil { + return nil + } + out := new(TrafficSplitBackend) + in.DeepCopyInto(out) + return out +} + +// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. +func (in *TrafficSplitList) DeepCopyInto(out *TrafficSplitList) { + *out = *in + out.TypeMeta = in.TypeMeta + out.ListMeta = in.ListMeta + if in.Items != nil { + in, out := &in.Items, &out.Items + *out = make([]TrafficSplit, len(*in)) + for i := range *in { + (*in)[i].DeepCopyInto(&(*out)[i]) + } + } + return +} + +// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new TrafficSplitList. +func (in *TrafficSplitList) DeepCopy() *TrafficSplitList { + if in == nil { + return nil + } + out := new(TrafficSplitList) + in.DeepCopyInto(out) + return out +} + +// DeepCopyObject is an autogenerated deepcopy function, copying the receiver, creating a new runtime.Object. +func (in *TrafficSplitList) DeepCopyObject() runtime.Object { + if c := in.DeepCopy(); c != nil { + return c + } + return nil +} + +// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. +func (in *TrafficSplitSpec) DeepCopyInto(out *TrafficSplitSpec) { + *out = *in + if in.Backends != nil { + in, out := &in.Backends, &out.Backends + *out = make([]TrafficSplitBackend, len(*in)) + for i := range *in { + (*in)[i].DeepCopyInto(&(*out)[i]) + } + } + return +} + +// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new TrafficSplitSpec. +func (in *TrafficSplitSpec) DeepCopy() *TrafficSplitSpec { + if in == nil { + return nil + } + out := new(TrafficSplitSpec) + in.DeepCopyInto(out) + return out +} diff --git a/pkg/client/clientset/versioned/clientset.go b/pkg/client/clientset/versioned/clientset.go index c6039c8f..1a3e17bc 100644 --- a/pkg/client/clientset/versioned/clientset.go +++ b/pkg/client/clientset/versioned/clientset.go @@ -22,6 +22,7 @@ import ( appmeshv1beta1 "github.com/weaveworks/flagger/pkg/client/clientset/versioned/typed/appmesh/v1beta1" flaggerv1alpha3 "github.com/weaveworks/flagger/pkg/client/clientset/versioned/typed/flagger/v1alpha3" networkingv1alpha3 "github.com/weaveworks/flagger/pkg/client/clientset/versioned/typed/istio/v1alpha3" + splitv1alpha1 "github.com/weaveworks/flagger/pkg/client/clientset/versioned/typed/smi/v1alpha1" discovery "k8s.io/client-go/discovery" rest "k8s.io/client-go/rest" flowcontrol "k8s.io/client-go/util/flowcontrol" @@ -38,6 +39,9 @@ type Interface interface { NetworkingV1alpha3() networkingv1alpha3.NetworkingV1alpha3Interface // Deprecated: please explicitly pick a version if possible. Networking() networkingv1alpha3.NetworkingV1alpha3Interface + SplitV1alpha1() splitv1alpha1.SplitV1alpha1Interface + // Deprecated: please explicitly pick a version if possible. + Split() splitv1alpha1.SplitV1alpha1Interface } // Clientset contains the clients for groups. Each group has exactly one @@ -47,6 +51,7 @@ type Clientset struct { appmeshV1beta1 *appmeshv1beta1.AppmeshV1beta1Client flaggerV1alpha3 *flaggerv1alpha3.FlaggerV1alpha3Client networkingV1alpha3 *networkingv1alpha3.NetworkingV1alpha3Client + splitV1alpha1 *splitv1alpha1.SplitV1alpha1Client } // AppmeshV1beta1 retrieves the AppmeshV1beta1Client @@ -82,6 +87,17 @@ func (c *Clientset) Networking() networkingv1alpha3.NetworkingV1alpha3Interface return c.networkingV1alpha3 } +// SplitV1alpha1 retrieves the SplitV1alpha1Client +func (c *Clientset) SplitV1alpha1() splitv1alpha1.SplitV1alpha1Interface { + return c.splitV1alpha1 +} + +// Deprecated: Split retrieves the default version of SplitClient. +// Please explicitly pick a version. +func (c *Clientset) Split() splitv1alpha1.SplitV1alpha1Interface { + return c.splitV1alpha1 +} + // Discovery retrieves the DiscoveryClient func (c *Clientset) Discovery() discovery.DiscoveryInterface { if c == nil { @@ -110,6 +126,10 @@ func NewForConfig(c *rest.Config) (*Clientset, error) { if err != nil { return nil, err } + cs.splitV1alpha1, err = splitv1alpha1.NewForConfig(&configShallowCopy) + if err != nil { + return nil, err + } cs.DiscoveryClient, err = discovery.NewDiscoveryClientForConfig(&configShallowCopy) if err != nil { @@ -125,6 +145,7 @@ func NewForConfigOrDie(c *rest.Config) *Clientset { cs.appmeshV1beta1 = appmeshv1beta1.NewForConfigOrDie(c) cs.flaggerV1alpha3 = flaggerv1alpha3.NewForConfigOrDie(c) cs.networkingV1alpha3 = networkingv1alpha3.NewForConfigOrDie(c) + cs.splitV1alpha1 = splitv1alpha1.NewForConfigOrDie(c) cs.DiscoveryClient = discovery.NewDiscoveryClientForConfigOrDie(c) return &cs @@ -136,6 +157,7 @@ func New(c rest.Interface) *Clientset { cs.appmeshV1beta1 = appmeshv1beta1.New(c) cs.flaggerV1alpha3 = flaggerv1alpha3.New(c) cs.networkingV1alpha3 = networkingv1alpha3.New(c) + cs.splitV1alpha1 = splitv1alpha1.New(c) cs.DiscoveryClient = discovery.NewDiscoveryClient(c) return &cs diff --git a/pkg/client/clientset/versioned/fake/clientset_generated.go b/pkg/client/clientset/versioned/fake/clientset_generated.go index abb76981..57e467ca 100644 --- a/pkg/client/clientset/versioned/fake/clientset_generated.go +++ b/pkg/client/clientset/versioned/fake/clientset_generated.go @@ -26,6 +26,8 @@ import ( fakeflaggerv1alpha3 "github.com/weaveworks/flagger/pkg/client/clientset/versioned/typed/flagger/v1alpha3/fake" networkingv1alpha3 "github.com/weaveworks/flagger/pkg/client/clientset/versioned/typed/istio/v1alpha3" fakenetworkingv1alpha3 "github.com/weaveworks/flagger/pkg/client/clientset/versioned/typed/istio/v1alpha3/fake" + splitv1alpha1 "github.com/weaveworks/flagger/pkg/client/clientset/versioned/typed/smi/v1alpha1" + fakesplitv1alpha1 "github.com/weaveworks/flagger/pkg/client/clientset/versioned/typed/smi/v1alpha1/fake" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/watch" "k8s.io/client-go/discovery" @@ -104,3 +106,13 @@ func (c *Clientset) NetworkingV1alpha3() networkingv1alpha3.NetworkingV1alpha3In func (c *Clientset) Networking() networkingv1alpha3.NetworkingV1alpha3Interface { return &fakenetworkingv1alpha3.FakeNetworkingV1alpha3{Fake: &c.Fake} } + +// SplitV1alpha1 retrieves the SplitV1alpha1Client +func (c *Clientset) SplitV1alpha1() splitv1alpha1.SplitV1alpha1Interface { + return &fakesplitv1alpha1.FakeSplitV1alpha1{Fake: &c.Fake} +} + +// Split retrieves the SplitV1alpha1Client +func (c *Clientset) Split() splitv1alpha1.SplitV1alpha1Interface { + return &fakesplitv1alpha1.FakeSplitV1alpha1{Fake: &c.Fake} +} diff --git a/pkg/client/clientset/versioned/fake/register.go b/pkg/client/clientset/versioned/fake/register.go index 29214ba6..b5910b2e 100644 --- a/pkg/client/clientset/versioned/fake/register.go +++ b/pkg/client/clientset/versioned/fake/register.go @@ -22,6 +22,7 @@ import ( appmeshv1beta1 "github.com/weaveworks/flagger/pkg/apis/appmesh/v1beta1" flaggerv1alpha3 "github.com/weaveworks/flagger/pkg/apis/flagger/v1alpha3" networkingv1alpha3 "github.com/weaveworks/flagger/pkg/apis/istio/v1alpha3" + splitv1alpha1 "github.com/weaveworks/flagger/pkg/apis/smi/v1alpha1" v1 "k8s.io/apimachinery/pkg/apis/meta/v1" runtime "k8s.io/apimachinery/pkg/runtime" schema "k8s.io/apimachinery/pkg/runtime/schema" @@ -36,6 +37,7 @@ var localSchemeBuilder = runtime.SchemeBuilder{ appmeshv1beta1.AddToScheme, flaggerv1alpha3.AddToScheme, networkingv1alpha3.AddToScheme, + splitv1alpha1.AddToScheme, } // AddToScheme adds all types of this clientset into the given scheme. This allows composition diff --git a/pkg/client/clientset/versioned/scheme/register.go b/pkg/client/clientset/versioned/scheme/register.go index c1e0adcc..7337f87a 100644 --- a/pkg/client/clientset/versioned/scheme/register.go +++ b/pkg/client/clientset/versioned/scheme/register.go @@ -22,6 +22,7 @@ import ( appmeshv1beta1 "github.com/weaveworks/flagger/pkg/apis/appmesh/v1beta1" flaggerv1alpha3 "github.com/weaveworks/flagger/pkg/apis/flagger/v1alpha3" networkingv1alpha3 "github.com/weaveworks/flagger/pkg/apis/istio/v1alpha3" + splitv1alpha1 "github.com/weaveworks/flagger/pkg/apis/smi/v1alpha1" v1 "k8s.io/apimachinery/pkg/apis/meta/v1" runtime "k8s.io/apimachinery/pkg/runtime" schema "k8s.io/apimachinery/pkg/runtime/schema" @@ -36,6 +37,7 @@ var localSchemeBuilder = runtime.SchemeBuilder{ appmeshv1beta1.AddToScheme, flaggerv1alpha3.AddToScheme, networkingv1alpha3.AddToScheme, + splitv1alpha1.AddToScheme, } // AddToScheme adds all types of this clientset into the given scheme. This allows composition diff --git a/pkg/client/clientset/versioned/typed/smi/v1alpha1/doc.go b/pkg/client/clientset/versioned/typed/smi/v1alpha1/doc.go new file mode 100644 index 00000000..20b3d7fd --- /dev/null +++ b/pkg/client/clientset/versioned/typed/smi/v1alpha1/doc.go @@ -0,0 +1,20 @@ +/* +Copyright The Flagger Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +// Code generated by client-gen. DO NOT EDIT. + +// This package has the automatically generated typed clients. +package v1alpha1 diff --git a/pkg/client/clientset/versioned/typed/smi/v1alpha1/fake/doc.go b/pkg/client/clientset/versioned/typed/smi/v1alpha1/fake/doc.go new file mode 100644 index 00000000..7a3b19cb --- /dev/null +++ b/pkg/client/clientset/versioned/typed/smi/v1alpha1/fake/doc.go @@ -0,0 +1,20 @@ +/* +Copyright The Flagger Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +// Code generated by client-gen. DO NOT EDIT. + +// Package fake has the automatically generated clients. +package fake diff --git a/pkg/client/clientset/versioned/typed/smi/v1alpha1/fake/fake_smi_client.go b/pkg/client/clientset/versioned/typed/smi/v1alpha1/fake/fake_smi_client.go new file mode 100644 index 00000000..e3cf89ed --- /dev/null +++ b/pkg/client/clientset/versioned/typed/smi/v1alpha1/fake/fake_smi_client.go @@ -0,0 +1,40 @@ +/* +Copyright The Flagger Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +// Code generated by client-gen. DO NOT EDIT. + +package fake + +import ( + v1alpha1 "github.com/weaveworks/flagger/pkg/client/clientset/versioned/typed/smi/v1alpha1" + rest "k8s.io/client-go/rest" + testing "k8s.io/client-go/testing" +) + +type FakeSplitV1alpha1 struct { + *testing.Fake +} + +func (c *FakeSplitV1alpha1) TrafficSplits(namespace string) v1alpha1.TrafficSplitInterface { + return &FakeTrafficSplits{c, namespace} +} + +// RESTClient returns a RESTClient that is used to communicate +// with API server by this client implementation. +func (c *FakeSplitV1alpha1) RESTClient() rest.Interface { + var ret *rest.RESTClient + return ret +} diff --git a/pkg/client/clientset/versioned/typed/smi/v1alpha1/fake/fake_trafficsplit.go b/pkg/client/clientset/versioned/typed/smi/v1alpha1/fake/fake_trafficsplit.go new file mode 100644 index 00000000..ac6b67cc --- /dev/null +++ b/pkg/client/clientset/versioned/typed/smi/v1alpha1/fake/fake_trafficsplit.go @@ -0,0 +1,128 @@ +/* +Copyright The Flagger Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +// Code generated by client-gen. DO NOT EDIT. + +package fake + +import ( + v1alpha1 "github.com/weaveworks/flagger/pkg/apis/smi/v1alpha1" + v1 "k8s.io/apimachinery/pkg/apis/meta/v1" + labels "k8s.io/apimachinery/pkg/labels" + schema "k8s.io/apimachinery/pkg/runtime/schema" + types "k8s.io/apimachinery/pkg/types" + watch "k8s.io/apimachinery/pkg/watch" + testing "k8s.io/client-go/testing" +) + +// FakeTrafficSplits implements TrafficSplitInterface +type FakeTrafficSplits struct { + Fake *FakeSplitV1alpha1 + ns string +} + +var trafficsplitsResource = schema.GroupVersionResource{Group: "split.smi-spec.io", Version: "v1alpha1", Resource: "trafficsplits"} + +var trafficsplitsKind = schema.GroupVersionKind{Group: "split.smi-spec.io", Version: "v1alpha1", Kind: "TrafficSplit"} + +// Get takes name of the trafficSplit, and returns the corresponding trafficSplit object, and an error if there is any. +func (c *FakeTrafficSplits) Get(name string, options v1.GetOptions) (result *v1alpha1.TrafficSplit, err error) { + obj, err := c.Fake. + Invokes(testing.NewGetAction(trafficsplitsResource, c.ns, name), &v1alpha1.TrafficSplit{}) + + if obj == nil { + return nil, err + } + return obj.(*v1alpha1.TrafficSplit), err +} + +// List takes label and field selectors, and returns the list of TrafficSplits that match those selectors. +func (c *FakeTrafficSplits) List(opts v1.ListOptions) (result *v1alpha1.TrafficSplitList, err error) { + obj, err := c.Fake. + Invokes(testing.NewListAction(trafficsplitsResource, trafficsplitsKind, c.ns, opts), &v1alpha1.TrafficSplitList{}) + + if obj == nil { + return nil, err + } + + label, _, _ := testing.ExtractFromListOptions(opts) + if label == nil { + label = labels.Everything() + } + list := &v1alpha1.TrafficSplitList{ListMeta: obj.(*v1alpha1.TrafficSplitList).ListMeta} + for _, item := range obj.(*v1alpha1.TrafficSplitList).Items { + if label.Matches(labels.Set(item.Labels)) { + list.Items = append(list.Items, item) + } + } + return list, err +} + +// Watch returns a watch.Interface that watches the requested trafficSplits. +func (c *FakeTrafficSplits) Watch(opts v1.ListOptions) (watch.Interface, error) { + return c.Fake. + InvokesWatch(testing.NewWatchAction(trafficsplitsResource, c.ns, opts)) + +} + +// Create takes the representation of a trafficSplit and creates it. Returns the server's representation of the trafficSplit, and an error, if there is any. +func (c *FakeTrafficSplits) Create(trafficSplit *v1alpha1.TrafficSplit) (result *v1alpha1.TrafficSplit, err error) { + obj, err := c.Fake. + Invokes(testing.NewCreateAction(trafficsplitsResource, c.ns, trafficSplit), &v1alpha1.TrafficSplit{}) + + if obj == nil { + return nil, err + } + return obj.(*v1alpha1.TrafficSplit), err +} + +// Update takes the representation of a trafficSplit and updates it. Returns the server's representation of the trafficSplit, and an error, if there is any. +func (c *FakeTrafficSplits) Update(trafficSplit *v1alpha1.TrafficSplit) (result *v1alpha1.TrafficSplit, err error) { + obj, err := c.Fake. + Invokes(testing.NewUpdateAction(trafficsplitsResource, c.ns, trafficSplit), &v1alpha1.TrafficSplit{}) + + if obj == nil { + return nil, err + } + return obj.(*v1alpha1.TrafficSplit), err +} + +// Delete takes name of the trafficSplit and deletes it. Returns an error if one occurs. +func (c *FakeTrafficSplits) Delete(name string, options *v1.DeleteOptions) error { + _, err := c.Fake. + Invokes(testing.NewDeleteAction(trafficsplitsResource, c.ns, name), &v1alpha1.TrafficSplit{}) + + return err +} + +// DeleteCollection deletes a collection of objects. +func (c *FakeTrafficSplits) DeleteCollection(options *v1.DeleteOptions, listOptions v1.ListOptions) error { + action := testing.NewDeleteCollectionAction(trafficsplitsResource, c.ns, listOptions) + + _, err := c.Fake.Invokes(action, &v1alpha1.TrafficSplitList{}) + return err +} + +// Patch applies the patch and returns the patched trafficSplit. +func (c *FakeTrafficSplits) Patch(name string, pt types.PatchType, data []byte, subresources ...string) (result *v1alpha1.TrafficSplit, err error) { + obj, err := c.Fake. + Invokes(testing.NewPatchSubresourceAction(trafficsplitsResource, c.ns, name, pt, data, subresources...), &v1alpha1.TrafficSplit{}) + + if obj == nil { + return nil, err + } + return obj.(*v1alpha1.TrafficSplit), err +} diff --git a/pkg/client/clientset/versioned/typed/smi/v1alpha1/generated_expansion.go b/pkg/client/clientset/versioned/typed/smi/v1alpha1/generated_expansion.go new file mode 100644 index 00000000..4cc5c42a --- /dev/null +++ b/pkg/client/clientset/versioned/typed/smi/v1alpha1/generated_expansion.go @@ -0,0 +1,21 @@ +/* +Copyright The Flagger Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +// Code generated by client-gen. DO NOT EDIT. + +package v1alpha1 + +type TrafficSplitExpansion interface{} diff --git a/pkg/client/clientset/versioned/typed/smi/v1alpha1/smi_client.go b/pkg/client/clientset/versioned/typed/smi/v1alpha1/smi_client.go new file mode 100644 index 00000000..745325bd --- /dev/null +++ b/pkg/client/clientset/versioned/typed/smi/v1alpha1/smi_client.go @@ -0,0 +1,90 @@ +/* +Copyright The Flagger Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +// Code generated by client-gen. DO NOT EDIT. + +package v1alpha1 + +import ( + v1alpha1 "github.com/weaveworks/flagger/pkg/apis/smi/v1alpha1" + "github.com/weaveworks/flagger/pkg/client/clientset/versioned/scheme" + serializer "k8s.io/apimachinery/pkg/runtime/serializer" + rest "k8s.io/client-go/rest" +) + +type SplitV1alpha1Interface interface { + RESTClient() rest.Interface + TrafficSplitsGetter +} + +// SplitV1alpha1Client is used to interact with features provided by the split.smi-spec.io group. +type SplitV1alpha1Client struct { + restClient rest.Interface +} + +func (c *SplitV1alpha1Client) TrafficSplits(namespace string) TrafficSplitInterface { + return newTrafficSplits(c, namespace) +} + +// NewForConfig creates a new SplitV1alpha1Client for the given config. +func NewForConfig(c *rest.Config) (*SplitV1alpha1Client, error) { + config := *c + if err := setConfigDefaults(&config); err != nil { + return nil, err + } + client, err := rest.RESTClientFor(&config) + if err != nil { + return nil, err + } + return &SplitV1alpha1Client{client}, nil +} + +// NewForConfigOrDie creates a new SplitV1alpha1Client for the given config and +// panics if there is an error in the config. +func NewForConfigOrDie(c *rest.Config) *SplitV1alpha1Client { + client, err := NewForConfig(c) + if err != nil { + panic(err) + } + return client +} + +// New creates a new SplitV1alpha1Client for the given RESTClient. +func New(c rest.Interface) *SplitV1alpha1Client { + return &SplitV1alpha1Client{c} +} + +func setConfigDefaults(config *rest.Config) error { + gv := v1alpha1.SchemeGroupVersion + config.GroupVersion = &gv + config.APIPath = "/apis" + config.NegotiatedSerializer = serializer.DirectCodecFactory{CodecFactory: scheme.Codecs} + + if config.UserAgent == "" { + config.UserAgent = rest.DefaultKubernetesUserAgent() + } + + return nil +} + +// RESTClient returns a RESTClient that is used to communicate +// with API server by this client implementation. +func (c *SplitV1alpha1Client) RESTClient() rest.Interface { + if c == nil { + return nil + } + return c.restClient +} diff --git a/pkg/client/clientset/versioned/typed/smi/v1alpha1/trafficsplit.go b/pkg/client/clientset/versioned/typed/smi/v1alpha1/trafficsplit.go new file mode 100644 index 00000000..18f7bdc6 --- /dev/null +++ b/pkg/client/clientset/versioned/typed/smi/v1alpha1/trafficsplit.go @@ -0,0 +1,174 @@ +/* +Copyright The Flagger Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +// Code generated by client-gen. DO NOT EDIT. + +package v1alpha1 + +import ( + "time" + + v1alpha1 "github.com/weaveworks/flagger/pkg/apis/smi/v1alpha1" + scheme "github.com/weaveworks/flagger/pkg/client/clientset/versioned/scheme" + v1 "k8s.io/apimachinery/pkg/apis/meta/v1" + types "k8s.io/apimachinery/pkg/types" + watch "k8s.io/apimachinery/pkg/watch" + rest "k8s.io/client-go/rest" +) + +// TrafficSplitsGetter has a method to return a TrafficSplitInterface. +// A group's client should implement this interface. +type TrafficSplitsGetter interface { + TrafficSplits(namespace string) TrafficSplitInterface +} + +// TrafficSplitInterface has methods to work with TrafficSplit resources. +type TrafficSplitInterface interface { + Create(*v1alpha1.TrafficSplit) (*v1alpha1.TrafficSplit, error) + Update(*v1alpha1.TrafficSplit) (*v1alpha1.TrafficSplit, error) + Delete(name string, options *v1.DeleteOptions) error + DeleteCollection(options *v1.DeleteOptions, listOptions v1.ListOptions) error + Get(name string, options v1.GetOptions) (*v1alpha1.TrafficSplit, error) + List(opts v1.ListOptions) (*v1alpha1.TrafficSplitList, error) + Watch(opts v1.ListOptions) (watch.Interface, error) + Patch(name string, pt types.PatchType, data []byte, subresources ...string) (result *v1alpha1.TrafficSplit, err error) + TrafficSplitExpansion +} + +// trafficSplits implements TrafficSplitInterface +type trafficSplits struct { + client rest.Interface + ns string +} + +// newTrafficSplits returns a TrafficSplits +func newTrafficSplits(c *SplitV1alpha1Client, namespace string) *trafficSplits { + return &trafficSplits{ + client: c.RESTClient(), + ns: namespace, + } +} + +// Get takes name of the trafficSplit, and returns the corresponding trafficSplit object, and an error if there is any. +func (c *trafficSplits) Get(name string, options v1.GetOptions) (result *v1alpha1.TrafficSplit, err error) { + result = &v1alpha1.TrafficSplit{} + err = c.client.Get(). + Namespace(c.ns). + Resource("trafficsplits"). + Name(name). + VersionedParams(&options, scheme.ParameterCodec). + Do(). + Into(result) + return +} + +// List takes label and field selectors, and returns the list of TrafficSplits that match those selectors. +func (c *trafficSplits) List(opts v1.ListOptions) (result *v1alpha1.TrafficSplitList, err error) { + var timeout time.Duration + if opts.TimeoutSeconds != nil { + timeout = time.Duration(*opts.TimeoutSeconds) * time.Second + } + result = &v1alpha1.TrafficSplitList{} + err = c.client.Get(). + Namespace(c.ns). + Resource("trafficsplits"). + VersionedParams(&opts, scheme.ParameterCodec). + Timeout(timeout). + Do(). + Into(result) + return +} + +// Watch returns a watch.Interface that watches the requested trafficSplits. +func (c *trafficSplits) Watch(opts v1.ListOptions) (watch.Interface, error) { + var timeout time.Duration + if opts.TimeoutSeconds != nil { + timeout = time.Duration(*opts.TimeoutSeconds) * time.Second + } + opts.Watch = true + return c.client.Get(). + Namespace(c.ns). + Resource("trafficsplits"). + VersionedParams(&opts, scheme.ParameterCodec). + Timeout(timeout). + Watch() +} + +// Create takes the representation of a trafficSplit and creates it. Returns the server's representation of the trafficSplit, and an error, if there is any. +func (c *trafficSplits) Create(trafficSplit *v1alpha1.TrafficSplit) (result *v1alpha1.TrafficSplit, err error) { + result = &v1alpha1.TrafficSplit{} + err = c.client.Post(). + Namespace(c.ns). + Resource("trafficsplits"). + Body(trafficSplit). + Do(). + Into(result) + return +} + +// Update takes the representation of a trafficSplit and updates it. Returns the server's representation of the trafficSplit, and an error, if there is any. +func (c *trafficSplits) Update(trafficSplit *v1alpha1.TrafficSplit) (result *v1alpha1.TrafficSplit, err error) { + result = &v1alpha1.TrafficSplit{} + err = c.client.Put(). + Namespace(c.ns). + Resource("trafficsplits"). + Name(trafficSplit.Name). + Body(trafficSplit). + Do(). + Into(result) + return +} + +// Delete takes name of the trafficSplit and deletes it. Returns an error if one occurs. +func (c *trafficSplits) Delete(name string, options *v1.DeleteOptions) error { + return c.client.Delete(). + Namespace(c.ns). + Resource("trafficsplits"). + Name(name). + Body(options). + Do(). + Error() +} + +// DeleteCollection deletes a collection of objects. +func (c *trafficSplits) DeleteCollection(options *v1.DeleteOptions, listOptions v1.ListOptions) error { + var timeout time.Duration + if listOptions.TimeoutSeconds != nil { + timeout = time.Duration(*listOptions.TimeoutSeconds) * time.Second + } + return c.client.Delete(). + Namespace(c.ns). + Resource("trafficsplits"). + VersionedParams(&listOptions, scheme.ParameterCodec). + Timeout(timeout). + Body(options). + Do(). + Error() +} + +// Patch applies the patch and returns the patched trafficSplit. +func (c *trafficSplits) Patch(name string, pt types.PatchType, data []byte, subresources ...string) (result *v1alpha1.TrafficSplit, err error) { + result = &v1alpha1.TrafficSplit{} + err = c.client.Patch(pt). + Namespace(c.ns). + Resource("trafficsplits"). + SubResource(subresources...). + Name(name). + Body(data). + Do(). + Into(result) + return +} diff --git a/pkg/client/informers/externalversions/factory.go b/pkg/client/informers/externalversions/factory.go index 62626dd9..c19c5ca2 100644 --- a/pkg/client/informers/externalversions/factory.go +++ b/pkg/client/informers/externalversions/factory.go @@ -28,6 +28,7 @@ import ( flagger "github.com/weaveworks/flagger/pkg/client/informers/externalversions/flagger" internalinterfaces "github.com/weaveworks/flagger/pkg/client/informers/externalversions/internalinterfaces" istio "github.com/weaveworks/flagger/pkg/client/informers/externalversions/istio" + smi "github.com/weaveworks/flagger/pkg/client/informers/externalversions/smi" v1 "k8s.io/apimachinery/pkg/apis/meta/v1" runtime "k8s.io/apimachinery/pkg/runtime" schema "k8s.io/apimachinery/pkg/runtime/schema" @@ -177,6 +178,7 @@ type SharedInformerFactory interface { Appmesh() appmesh.Interface Flagger() flagger.Interface Networking() istio.Interface + Split() smi.Interface } func (f *sharedInformerFactory) Appmesh() appmesh.Interface { @@ -190,3 +192,7 @@ func (f *sharedInformerFactory) Flagger() flagger.Interface { func (f *sharedInformerFactory) Networking() istio.Interface { return istio.New(f, f.namespace, f.tweakListOptions) } + +func (f *sharedInformerFactory) Split() smi.Interface { + return smi.New(f, f.namespace, f.tweakListOptions) +} diff --git a/pkg/client/informers/externalversions/generic.go b/pkg/client/informers/externalversions/generic.go index f938d50b..b374b9be 100644 --- a/pkg/client/informers/externalversions/generic.go +++ b/pkg/client/informers/externalversions/generic.go @@ -24,6 +24,7 @@ import ( v1beta1 "github.com/weaveworks/flagger/pkg/apis/appmesh/v1beta1" v1alpha3 "github.com/weaveworks/flagger/pkg/apis/flagger/v1alpha3" istiov1alpha3 "github.com/weaveworks/flagger/pkg/apis/istio/v1alpha3" + v1alpha1 "github.com/weaveworks/flagger/pkg/apis/smi/v1alpha1" schema "k8s.io/apimachinery/pkg/runtime/schema" cache "k8s.io/client-go/tools/cache" ) @@ -70,6 +71,10 @@ func (f *sharedInformerFactory) ForResource(resource schema.GroupVersionResource case istiov1alpha3.SchemeGroupVersion.WithResource("virtualservices"): return &genericInformer{resource: resource.GroupResource(), informer: f.Networking().V1alpha3().VirtualServices().Informer()}, nil + // Group=split.smi-spec.io, Version=v1alpha1 + case v1alpha1.SchemeGroupVersion.WithResource("trafficsplits"): + return &genericInformer{resource: resource.GroupResource(), informer: f.Split().V1alpha1().TrafficSplits().Informer()}, nil + } return nil, fmt.Errorf("no informer found for %v", resource) diff --git a/pkg/client/informers/externalversions/smi/interface.go b/pkg/client/informers/externalversions/smi/interface.go new file mode 100644 index 00000000..cfc8f719 --- /dev/null +++ b/pkg/client/informers/externalversions/smi/interface.go @@ -0,0 +1,46 @@ +/* +Copyright The Flagger Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +// Code generated by informer-gen. DO NOT EDIT. + +package split + +import ( + internalinterfaces "github.com/weaveworks/flagger/pkg/client/informers/externalversions/internalinterfaces" + v1alpha1 "github.com/weaveworks/flagger/pkg/client/informers/externalversions/smi/v1alpha1" +) + +// Interface provides access to each of this group's versions. +type Interface interface { + // V1alpha1 provides access to shared informers for resources in V1alpha1. + V1alpha1() v1alpha1.Interface +} + +type group struct { + factory internalinterfaces.SharedInformerFactory + namespace string + tweakListOptions internalinterfaces.TweakListOptionsFunc +} + +// New returns a new Interface. +func New(f internalinterfaces.SharedInformerFactory, namespace string, tweakListOptions internalinterfaces.TweakListOptionsFunc) Interface { + return &group{factory: f, namespace: namespace, tweakListOptions: tweakListOptions} +} + +// V1alpha1 returns a new v1alpha1.Interface. +func (g *group) V1alpha1() v1alpha1.Interface { + return v1alpha1.New(g.factory, g.namespace, g.tweakListOptions) +} diff --git a/pkg/client/informers/externalversions/smi/v1alpha1/interface.go b/pkg/client/informers/externalversions/smi/v1alpha1/interface.go new file mode 100644 index 00000000..dabfa548 --- /dev/null +++ b/pkg/client/informers/externalversions/smi/v1alpha1/interface.go @@ -0,0 +1,45 @@ +/* +Copyright The Flagger Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +// Code generated by informer-gen. DO NOT EDIT. + +package v1alpha1 + +import ( + internalinterfaces "github.com/weaveworks/flagger/pkg/client/informers/externalversions/internalinterfaces" +) + +// Interface provides access to all the informers in this group version. +type Interface interface { + // TrafficSplits returns a TrafficSplitInformer. + TrafficSplits() TrafficSplitInformer +} + +type version struct { + factory internalinterfaces.SharedInformerFactory + namespace string + tweakListOptions internalinterfaces.TweakListOptionsFunc +} + +// New returns a new Interface. +func New(f internalinterfaces.SharedInformerFactory, namespace string, tweakListOptions internalinterfaces.TweakListOptionsFunc) Interface { + return &version{factory: f, namespace: namespace, tweakListOptions: tweakListOptions} +} + +// TrafficSplits returns a TrafficSplitInformer. +func (v *version) TrafficSplits() TrafficSplitInformer { + return &trafficSplitInformer{factory: v.factory, namespace: v.namespace, tweakListOptions: v.tweakListOptions} +} diff --git a/pkg/client/informers/externalversions/smi/v1alpha1/trafficsplit.go b/pkg/client/informers/externalversions/smi/v1alpha1/trafficsplit.go new file mode 100644 index 00000000..feb6b703 --- /dev/null +++ b/pkg/client/informers/externalversions/smi/v1alpha1/trafficsplit.go @@ -0,0 +1,89 @@ +/* +Copyright The Flagger Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +// Code generated by informer-gen. DO NOT EDIT. + +package v1alpha1 + +import ( + time "time" + + smiv1alpha1 "github.com/weaveworks/flagger/pkg/apis/smi/v1alpha1" + versioned "github.com/weaveworks/flagger/pkg/client/clientset/versioned" + internalinterfaces "github.com/weaveworks/flagger/pkg/client/informers/externalversions/internalinterfaces" + v1alpha1 "github.com/weaveworks/flagger/pkg/client/listers/smi/v1alpha1" + v1 "k8s.io/apimachinery/pkg/apis/meta/v1" + runtime "k8s.io/apimachinery/pkg/runtime" + watch "k8s.io/apimachinery/pkg/watch" + cache "k8s.io/client-go/tools/cache" +) + +// TrafficSplitInformer provides access to a shared informer and lister for +// TrafficSplits. +type TrafficSplitInformer interface { + Informer() cache.SharedIndexInformer + Lister() v1alpha1.TrafficSplitLister +} + +type trafficSplitInformer struct { + factory internalinterfaces.SharedInformerFactory + tweakListOptions internalinterfaces.TweakListOptionsFunc + namespace string +} + +// NewTrafficSplitInformer constructs a new informer for TrafficSplit type. +// Always prefer using an informer factory to get a shared informer instead of getting an independent +// one. This reduces memory footprint and number of connections to the server. +func NewTrafficSplitInformer(client versioned.Interface, namespace string, resyncPeriod time.Duration, indexers cache.Indexers) cache.SharedIndexInformer { + return NewFilteredTrafficSplitInformer(client, namespace, resyncPeriod, indexers, nil) +} + +// NewFilteredTrafficSplitInformer constructs a new informer for TrafficSplit type. +// Always prefer using an informer factory to get a shared informer instead of getting an independent +// one. This reduces memory footprint and number of connections to the server. +func NewFilteredTrafficSplitInformer(client versioned.Interface, namespace string, resyncPeriod time.Duration, indexers cache.Indexers, tweakListOptions internalinterfaces.TweakListOptionsFunc) cache.SharedIndexInformer { + return cache.NewSharedIndexInformer( + &cache.ListWatch{ + ListFunc: func(options v1.ListOptions) (runtime.Object, error) { + if tweakListOptions != nil { + tweakListOptions(&options) + } + return client.SplitV1alpha1().TrafficSplits(namespace).List(options) + }, + WatchFunc: func(options v1.ListOptions) (watch.Interface, error) { + if tweakListOptions != nil { + tweakListOptions(&options) + } + return client.SplitV1alpha1().TrafficSplits(namespace).Watch(options) + }, + }, + &smiv1alpha1.TrafficSplit{}, + resyncPeriod, + indexers, + ) +} + +func (f *trafficSplitInformer) defaultInformer(client versioned.Interface, resyncPeriod time.Duration) cache.SharedIndexInformer { + return NewFilteredTrafficSplitInformer(client, f.namespace, resyncPeriod, cache.Indexers{cache.NamespaceIndex: cache.MetaNamespaceIndexFunc}, f.tweakListOptions) +} + +func (f *trafficSplitInformer) Informer() cache.SharedIndexInformer { + return f.factory.InformerFor(&smiv1alpha1.TrafficSplit{}, f.defaultInformer) +} + +func (f *trafficSplitInformer) Lister() v1alpha1.TrafficSplitLister { + return v1alpha1.NewTrafficSplitLister(f.Informer().GetIndexer()) +} diff --git a/pkg/client/listers/smi/v1alpha1/expansion_generated.go b/pkg/client/listers/smi/v1alpha1/expansion_generated.go new file mode 100644 index 00000000..271ee243 --- /dev/null +++ b/pkg/client/listers/smi/v1alpha1/expansion_generated.go @@ -0,0 +1,27 @@ +/* +Copyright The Flagger Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +// Code generated by lister-gen. DO NOT EDIT. + +package v1alpha1 + +// TrafficSplitListerExpansion allows custom methods to be added to +// TrafficSplitLister. +type TrafficSplitListerExpansion interface{} + +// TrafficSplitNamespaceListerExpansion allows custom methods to be added to +// TrafficSplitNamespaceLister. +type TrafficSplitNamespaceListerExpansion interface{} diff --git a/pkg/client/listers/smi/v1alpha1/trafficsplit.go b/pkg/client/listers/smi/v1alpha1/trafficsplit.go new file mode 100644 index 00000000..23c52eba --- /dev/null +++ b/pkg/client/listers/smi/v1alpha1/trafficsplit.go @@ -0,0 +1,94 @@ +/* +Copyright The Flagger Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +// Code generated by lister-gen. DO NOT EDIT. + +package v1alpha1 + +import ( + v1alpha1 "github.com/weaveworks/flagger/pkg/apis/smi/v1alpha1" + "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/labels" + "k8s.io/client-go/tools/cache" +) + +// TrafficSplitLister helps list TrafficSplits. +type TrafficSplitLister interface { + // List lists all TrafficSplits in the indexer. + List(selector labels.Selector) (ret []*v1alpha1.TrafficSplit, err error) + // TrafficSplits returns an object that can list and get TrafficSplits. + TrafficSplits(namespace string) TrafficSplitNamespaceLister + TrafficSplitListerExpansion +} + +// trafficSplitLister implements the TrafficSplitLister interface. +type trafficSplitLister struct { + indexer cache.Indexer +} + +// NewTrafficSplitLister returns a new TrafficSplitLister. +func NewTrafficSplitLister(indexer cache.Indexer) TrafficSplitLister { + return &trafficSplitLister{indexer: indexer} +} + +// List lists all TrafficSplits in the indexer. +func (s *trafficSplitLister) List(selector labels.Selector) (ret []*v1alpha1.TrafficSplit, err error) { + err = cache.ListAll(s.indexer, selector, func(m interface{}) { + ret = append(ret, m.(*v1alpha1.TrafficSplit)) + }) + return ret, err +} + +// TrafficSplits returns an object that can list and get TrafficSplits. +func (s *trafficSplitLister) TrafficSplits(namespace string) TrafficSplitNamespaceLister { + return trafficSplitNamespaceLister{indexer: s.indexer, namespace: namespace} +} + +// TrafficSplitNamespaceLister helps list and get TrafficSplits. +type TrafficSplitNamespaceLister interface { + // List lists all TrafficSplits in the indexer for a given namespace. + List(selector labels.Selector) (ret []*v1alpha1.TrafficSplit, err error) + // Get retrieves the TrafficSplit from the indexer for a given namespace and name. + Get(name string) (*v1alpha1.TrafficSplit, error) + TrafficSplitNamespaceListerExpansion +} + +// trafficSplitNamespaceLister implements the TrafficSplitNamespaceLister +// interface. +type trafficSplitNamespaceLister struct { + indexer cache.Indexer + namespace string +} + +// List lists all TrafficSplits in the indexer for a given namespace. +func (s trafficSplitNamespaceLister) List(selector labels.Selector) (ret []*v1alpha1.TrafficSplit, err error) { + err = cache.ListAllByNamespace(s.indexer, s.namespace, selector, func(m interface{}) { + ret = append(ret, m.(*v1alpha1.TrafficSplit)) + }) + return ret, err +} + +// Get retrieves the TrafficSplit from the indexer for a given namespace and name. +func (s trafficSplitNamespaceLister) Get(name string) (*v1alpha1.TrafficSplit, error) { + obj, exists, err := s.indexer.GetByKey(s.namespace + "/" + name) + if err != nil { + return nil, err + } + if !exists { + return nil, errors.NewNotFound(v1alpha1.Resource("trafficsplit"), name) + } + return obj.(*v1alpha1.TrafficSplit), nil +} From 95b8840bf253765f5c860d0c39179c6042b80ed0 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sat, 11 May 2019 13:05:19 +0300 Subject: [PATCH 02/14] Add SMI traffic split to router --- pkg/router/factory.go | 9 ++ pkg/router/smi.go | 190 ++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 199 insertions(+) create mode 100644 pkg/router/smi.go diff --git a/pkg/router/factory.go b/pkg/router/factory.go index ab836d31..09ed940b 100644 --- a/pkg/router/factory.go +++ b/pkg/router/factory.go @@ -56,6 +56,15 @@ func (factory *Factory) MeshRouter(provider string) Interface { kubeClient: factory.kubeClient, appmeshClient: factory.meshClient, } + case strings.HasPrefix(provider, "smi:"): + mesh := strings.TrimPrefix(provider, "smi:") + return &SmiRouter{ + logger: factory.logger, + flaggerClient: factory.flaggerClient, + kubeClient: factory.kubeClient, + smiClient: factory.meshClient, + targetMesh: mesh, + } case strings.HasPrefix(provider, "supergloo"): supergloo, err := NewSuperglooRouter(context.TODO(), provider, factory.flaggerClient, factory.logger, factory.kubeConfig) if err != nil { diff --git a/pkg/router/smi.go b/pkg/router/smi.go new file mode 100644 index 00000000..4bdb2fbd --- /dev/null +++ b/pkg/router/smi.go @@ -0,0 +1,190 @@ +package router + +import ( + "encoding/json" + "fmt" + + "github.com/google/go-cmp/cmp" + "github.com/google/go-cmp/cmp/cmpopts" + flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1alpha3" + smiv1 "github.com/weaveworks/flagger/pkg/apis/smi/v1alpha1" + clientset "github.com/weaveworks/flagger/pkg/client/clientset/versioned" + "go.uber.org/zap" + "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/api/resource" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/client-go/kubernetes" +) + +type SmiRouter struct { + kubeClient kubernetes.Interface + flaggerClient clientset.Interface + smiClient clientset.Interface + logger *zap.SugaredLogger + targetMesh string +} + +// Reconcile creates or updates the SMI traffic split +func (sr *SmiRouter) Reconcile(canary *flaggerv1.Canary) error { + targetName := canary.Spec.TargetRef.Name + canaryName := fmt.Sprintf("%s-canary", targetName) + primaryName := fmt.Sprintf("%s-primary", targetName) + + var host string + if len(canary.Spec.Service.Hosts) > 0 { + host = canary.Spec.Service.Hosts[0] + } else { + host = targetName + } + + tsSpec := smiv1.TrafficSplitSpec{ + Service: host, + Backends: []smiv1.TrafficSplitBackend{ + { + Service: canaryName, + Weight: resource.NewQuantity(0, resource.DecimalExponent), + }, + { + Service: primaryName, + Weight: resource.NewQuantity(100, resource.DecimalExponent), + }, + }, + } + + ts, err := sr.smiClient.SplitV1alpha1().TrafficSplits(canary.Namespace).Get(targetName, metav1.GetOptions{}) + // create traffic split + if errors.IsNotFound(err) { + t := &smiv1.TrafficSplit{ + ObjectMeta: metav1.ObjectMeta{ + Name: targetName, + Namespace: canary.Namespace, + OwnerReferences: []metav1.OwnerReference{ + *metav1.NewControllerRef(canary, schema.GroupVersionKind{ + Group: flaggerv1.SchemeGroupVersion.Group, + Version: flaggerv1.SchemeGroupVersion.Version, + Kind: flaggerv1.CanaryKind, + }), + }, + Annotations: sr.makeAnnotations(canary.Spec.Service.Gateways), + }, + Spec: tsSpec, + } + + _, err := sr.smiClient.SplitV1alpha1().TrafficSplits(canary.Namespace).Create(t) + if err != nil { + return err + } + + sr.logger.With("canary", fmt.Sprintf("%s.%s", canary.Name, canary.Namespace)). + Infof("TrafficSplit %s.%s created", t.GetName(), canary.Namespace) + return nil + } + + if err != nil { + return fmt.Errorf("traffic split %s query error %v", targetName, err) + } + + // update traffic split + if diff := cmp.Diff(tsSpec, ts.Spec, cmpopts.IgnoreTypes(resource.Quantity{})); diff != "" { + tsClone := ts.DeepCopy() + tsClone.Spec = tsSpec + + _, err := sr.smiClient.SplitV1alpha1().TrafficSplits(canary.Namespace).Update(tsClone) + if err != nil { + return fmt.Errorf("TrafficSplit %s update error %v", targetName, err) + } + + sr.logger.With("canary", fmt.Sprintf("%s.%s", canary.Name, canary.Namespace)). + Infof("TrafficSplit %s.%s updated", targetName, canary.Namespace) + return nil + } + + return nil +} + +// GetRoutes returns the destinations weight for primary and canary +func (sr *SmiRouter) GetRoutes(canary *flaggerv1.Canary) ( + primaryWeight int, + canaryWeight int, + err error, +) { + targetName := canary.Spec.TargetRef.Name + canaryName := fmt.Sprintf("%s-canary", targetName) + primaryName := fmt.Sprintf("%s-primary", targetName) + ts, err := sr.smiClient.SplitV1alpha1().TrafficSplits(canary.Namespace).Get(targetName, metav1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + err = fmt.Errorf("TrafficSplit %s.%s not found", targetName, canary.Namespace) + return + } + err = fmt.Errorf("TrafficSplit %s.%s query error %v", targetName, canary.Namespace, err) + return + } + + for _, r := range ts.Spec.Backends { + w, _ := r.Weight.AsInt64() + if r.Service == primaryName { + primaryWeight = int(w) + } + if r.Service == canaryName { + canaryWeight = int(w) + } + } + + if primaryWeight == 0 && canaryWeight == 0 { + err = fmt.Errorf("TrafficSplit %s.%s does not contain routes for %s and %s", + targetName, canary.Namespace, primaryName, canaryName) + } + + return +} + +// SetRoutes updates the destinations weight for primary and canary +func (sr *SmiRouter) SetRoutes( + canary *flaggerv1.Canary, + primaryWeight int, + canaryWeight int, +) error { + targetName := canary.Spec.TargetRef.Name + canaryName := fmt.Sprintf("%s-canary", targetName) + primaryName := fmt.Sprintf("%s-primary", targetName) + ts, err := sr.smiClient.SplitV1alpha1().TrafficSplits(canary.Namespace).Get(targetName, metav1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return fmt.Errorf("TrafficSplit %s.%s not found", targetName, canary.Namespace) + + } + return fmt.Errorf("TrafficSplit %s.%s query error %v", targetName, canary.Namespace, err) + } + + backends := []smiv1.TrafficSplitBackend{ + { + Service: canaryName, + Weight: resource.NewQuantity(int64(canaryWeight), resource.DecimalExponent), + }, + { + Service: primaryName, + Weight: resource.NewQuantity(int64(primaryWeight), resource.DecimalExponent), + }, + } + + tsClone := ts.DeepCopy() + tsClone.Spec.Backends = backends + + _, err = sr.smiClient.SplitV1alpha1().TrafficSplits(canary.Namespace).Update(tsClone) + if err != nil { + return fmt.Errorf("TrafficSplit %s update error %v", targetName, err) + } + + return nil +} + +func (sr *SmiRouter) makeAnnotations(gateways []string) map[string]string { + res := make(map[string]string) + if sr.targetMesh == "istio" && len(gateways) > 0 { + g, _ := json.Marshal(gateways) + res["VirtualService.v1alpha3.networking.istio.io/spec.gateways"] = string(g) + } + return res +} From 81481204214dff20afa1ab77d013bb95a1b33a37 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sat, 11 May 2019 13:06:06 +0300 Subject: [PATCH 03/14] Enable Istio checks for SMI-Istio adapter --- pkg/controller/scheduler.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index 0fc06978..8fa792e4 100644 --- a/pkg/controller/scheduler.go +++ b/pkg/controller/scheduler.go @@ -598,7 +598,7 @@ func (c *Controller) analyseCanary(r *flaggerv1.Canary) bool { } // Istio checks - if c.meshProvider == "istio" { + if strings.Contains(c.meshProvider, "istio") { if metric.Name == "request-success-rate" || metric.Name == "istio_requests_total" { val, err := c.observer.GetIstioSuccessRate(r.Spec.TargetRef.Name, r.Namespace, metric.Name, metric.Interval) if err != nil { From 8fde6bdb8a3ab49b449e87307780ed7d39ec803e Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sat, 11 May 2019 13:35:36 +0300 Subject: [PATCH 04/14] Add SMI Istio adapter deployment --- Makefile | 6 ++ artifacts/smi/istio-adapter.yaml | 131 +++++++++++++++++++++++++++++++ 2 files changed, 137 insertions(+) create mode 100644 artifacts/smi/istio-adapter.yaml diff --git a/Makefile b/Makefile index 262f98b9..15771dcc 100644 --- a/Makefile +++ b/Makefile @@ -24,6 +24,12 @@ run-nginx: -slack-url=https://hooks.slack.com/services/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 \ -slack-channel="devops-alerts" +run-smi: + go run cmd/flagger/* -kubeconfig=$$HOME/.kube/config -log-level=info -mesh-provider=smi:istio -namespace=smi \ + -metrics-server=https://prometheus.istio.weavedx.com \ + -slack-url=https://hooks.slack.com/services/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 \ + -slack-channel="devops-alerts" + build: docker build -t weaveworks/flagger:$(TAG) . -f Dockerfile diff --git a/artifacts/smi/istio-adapter.yaml b/artifacts/smi/istio-adapter.yaml new file mode 100644 index 00000000..ce40d9bf --- /dev/null +++ b/artifacts/smi/istio-adapter.yaml @@ -0,0 +1,131 @@ +apiVersion: apiextensions.k8s.io/v1beta1 +kind: CustomResourceDefinition +metadata: + name: trafficsplits.split.smi-spec.io +spec: + additionalPrinterColumns: + - JSONPath: .spec.service + description: The service + name: Service + type: string + group: split.smi-spec.io + names: + kind: TrafficSplit + listKind: TrafficSplitList + plural: trafficsplits + singular: trafficsplit + scope: Namespaced + subresources: + status: {} + version: v1alpha1 + versions: + - name: v1alpha1 + served: true + storage: true +--- +apiVersion: v1 +kind: ServiceAccount +metadata: + name: smi-adapter-istio + namespace: istio-system +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRole +metadata: + name: smi-adapter-istio +rules: + - apiGroups: + - "" + resources: + - pods + - services + - endpoints + - persistentvolumeclaims + - events + - configmaps + - secrets + verbs: + - '*' + - apiGroups: + - apps + resources: + - deployments + - daemonsets + - replicasets + - statefulsets + verbs: + - '*' + - apiGroups: + - monitoring.coreos.com + resources: + - servicemonitors + verbs: + - get + - create + - apiGroups: + - apps + resourceNames: + - smi-adapter-istio + resources: + - deployments/finalizers + verbs: + - update + - apiGroups: + - split.smi-spec.io + resources: + - '*' + verbs: + - '*' + - apiGroups: + - networking.istio.io + resources: + - '*' + verbs: + - '*' +--- +kind: ClusterRoleBinding +apiVersion: rbac.authorization.k8s.io/v1 +metadata: + name: smi-adapter-istio +subjects: + - kind: ServiceAccount + name: smi-adapter-istio + namespace: istio-system +roleRef: + kind: Role + name: smi-adapter-istio + apiGroup: rbac.authorization.k8s.io +--- +apiVersion: apps/v1 +kind: Deployment +metadata: + name: smi-adapter-istio + namespace: istio-system +spec: + replicas: 1 + selector: + matchLabels: + name: smi-adapter-istio + template: + metadata: + labels: + name: smi-adapter-istio + annotations: + sidecar.istio.io/inject: "false" + spec: + serviceAccountName: smi-adapter-istio + containers: + - name: smi-adapter-istio + image: docker.io/stefanprodan/smi-adapter-istio:0.0.2-beta.1 + command: + - smi-adapter-istio + imagePullPolicy: Always + env: + - name: WATCH_NAMESPACE + value: "" + - name: POD_NAME + valueFrom: + fieldRef: + fieldPath: metadata.name + - name: OPERATOR_NAME + value: "smi-adapter-istio" From d63f05c92e096eaeb9d8fb6cdc4c7e1def419895 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sat, 11 May 2019 13:45:32 +0300 Subject: [PATCH 05/14] Add SMI group to RBAC --- artifacts/flagger/account.yaml | 5 +++++ charts/flagger/templates/rbac.yaml | 5 +++++ 2 files changed, 10 insertions(+) diff --git a/artifacts/flagger/account.yaml b/artifacts/flagger/account.yaml index 0fac89ad..d31e7568 100644 --- a/artifacts/flagger/account.yaml +++ b/artifacts/flagger/account.yaml @@ -59,6 +59,11 @@ rules: - virtualservices - virtualservices/status verbs: ["*"] + - apiGroups: + - split.smi-spec.io + resources: + - trafficsplits + verbs: ["*"] - nonResourceURLs: - /version verbs: diff --git a/charts/flagger/templates/rbac.yaml b/charts/flagger/templates/rbac.yaml index e03755c0..782e1df1 100644 --- a/charts/flagger/templates/rbac.yaml +++ b/charts/flagger/templates/rbac.yaml @@ -55,6 +55,11 @@ rules: - virtualservices - virtualservices/status verbs: ["*"] + - apiGroups: + - split.smi-spec.io + resources: + - trafficsplits + verbs: ["*"] - nonResourceURLs: - /version verbs: From eb856fda13410baeaac638ece1f714460844f48e Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sat, 11 May 2019 13:46:24 +0300 Subject: [PATCH 06/14] Add SMI Istio e2e tests --- .circleci/config.yml | 9 +++++++++ test/e2e-smi-istio-build.sh | 26 ++++++++++++++++++++++++++ 2 files changed, 35 insertions(+) create mode 100755 test/e2e-smi-istio-build.sh diff --git a/.circleci/config.yml b/.circleci/config.yml index f9164539..9b73ca51 100644 --- a/.circleci/config.yml +++ b/.circleci/config.yml @@ -9,6 +9,15 @@ jobs: - run: test/e2e-build.sh - run: test/e2e-tests.sh + e2e-smi-istio-testing: + machine: true + steps: + - checkout + - run: test/e2e-kind.sh + - run: test/e2e-istio.sh + - run: test/e2e-smi-istio-build.sh + - run: test/e2e-tests.sh canary + e2e-supergloo-testing: machine: true steps: diff --git a/test/e2e-smi-istio-build.sh b/test/e2e-smi-istio-build.sh new file mode 100755 index 00000000..06b28ff1 --- /dev/null +++ b/test/e2e-smi-istio-build.sh @@ -0,0 +1,26 @@ +#!/usr/bin/env bash + +set -o errexit + +REPO_ROOT=$(git rev-parse --show-toplevel) +export KUBECONFIG="$(kind get kubeconfig-path --name="kind")" + +echo '>>> Building Flagger' +cd ${REPO_ROOT} && docker build -t test/flagger:latest . -f Dockerfile + +kind load docker-image test/flagger:latest + +echo '>>> Installing Flagger' +helm upgrade -i flagger ${REPO_ROOT}/charts/flagger \ +--wait \ +--namespace istio-system \ +--set meshProvider=smi:istio + +kubectl -n istio-system set image deployment/flagger flagger=test/flagger:latest + +kubectl -n istio-system rollout status deployment/flagger + +echo '>>> Installing the SMI Istio adapter' +kubectl apply -f ${REPO_ROOT}/artifacts/smi/istio-adapter.yaml + +kubectl -n istio-system rollout status deployment/smi-adapter-istio \ No newline at end of file From bd817cc520334c8ae64f8bc7de5a1f657134fc0b Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sat, 11 May 2019 14:00:53 +0300 Subject: [PATCH 07/14] Run SMI Istio e2e tests --- .circleci/config.yml | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/.circleci/config.yml b/.circleci/config.yml index 9b73ca51..9d79a90b 100644 --- a/.circleci/config.yml +++ b/.circleci/config.yml @@ -47,6 +47,13 @@ workflows: - /gh-pages.*/ - /docs-.*/ - /release-.*/ + - e2e-smi-istio-testing: + filters: + branches: + ignore: + - /gh-pages.*/ + - /docs-.*/ + - /release-.*/ - e2e-supergloo-testing: filters: branches: From 7fe273a21d27750a07ddda84a24c8d53c01c05fc Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sat, 11 May 2019 14:08:58 +0300 Subject: [PATCH 08/14] Fix SMI cluster role binding --- artifacts/smi/istio-adapter.yaml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/artifacts/smi/istio-adapter.yaml b/artifacts/smi/istio-adapter.yaml index ce40d9bf..eaebdcb8 100644 --- a/artifacts/smi/istio-adapter.yaml +++ b/artifacts/smi/istio-adapter.yaml @@ -92,7 +92,7 @@ subjects: name: smi-adapter-istio namespace: istio-system roleRef: - kind: Role + kind: ClusterRole name: smi-adapter-istio apiGroup: rbac.authorization.k8s.io --- From 24a74d3589a70b1e1af31375a6f87155cc3dc54b Mon Sep 17 00:00:00 2001 From: Carlos Sanchez Date: Sat, 11 May 2019 13:42:08 +0200 Subject: [PATCH 09/14] Fix #177 Do not copy labels from canary to primary deployment --- pkg/canary/deployer.go | 1 - 1 file changed, 1 deletion(-) diff --git a/pkg/canary/deployer.go b/pkg/canary/deployer.go index bb8aa121..e6d86c35 100644 --- a/pkg/canary/deployer.go +++ b/pkg/canary/deployer.go @@ -214,7 +214,6 @@ func (c *Deployer) createPrimaryDeployment(cd *flaggerv1.Canary) (string, error) primaryDep = &appsv1.Deployment{ ObjectMeta: metav1.ObjectMeta{ Name: primaryName, - Labels: canaryDep.Labels, Namespace: cd.Namespace, OwnerReferences: []metav1.OwnerReference{ *metav1.NewControllerRef(cd, schema.GroupVersionKind{ From 1902884b56d586ddec09e1c4148fdbb8c276779d Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sat, 11 May 2019 15:16:31 +0300 Subject: [PATCH 10/14] Release v0.13.2 --- CHANGELOG.md | 13 ++++++++++++- artifacts/flagger/deployment.yaml | 2 +- charts/flagger/Chart.yaml | 4 ++-- charts/flagger/values.yaml | 2 +- pkg/version/version.go | 2 +- 5 files changed, 17 insertions(+), 6 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index db761266..4d31ae5c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,18 @@ All notable changes to this project are documented in this file. +## 0.13.2 (2019-04-11) + +Fixes for Jenkins X deployments (prevent the jx GC from removing the primary instance) + +#### Fixes + +- Do not copy labels from canary to primary deployment [#178](https://github.com/weaveworks/flagger/pull/178) + +#### Improvements + +- Add NGINX ingress controller e2e and unit tests [#176](https://github.com/weaveworks/flagger/pull/176) + ## 0.13.1 (2019-04-09) Fixes for custom metrics checks and NGINX Prometheus queries @@ -10,7 +22,6 @@ Fixes for custom metrics checks and NGINX Prometheus queries - Fix promql queries for custom checks and NGINX [#174](https://github.com/weaveworks/flagger/pull/174) - ## 0.13.0 (2019-04-08) Adds support for [NGINX](https://docs.flagger.app/usage/nginx-progressive-delivery) ingress controller diff --git a/artifacts/flagger/deployment.yaml b/artifacts/flagger/deployment.yaml index 833be438..7e5fe182 100644 --- a/artifacts/flagger/deployment.yaml +++ b/artifacts/flagger/deployment.yaml @@ -22,7 +22,7 @@ spec: serviceAccountName: flagger containers: - name: flagger - image: weaveworks/flagger:0.13.1 + image: weaveworks/flagger:0.13.2 imagePullPolicy: IfNotPresent ports: - name: http diff --git a/charts/flagger/Chart.yaml b/charts/flagger/Chart.yaml index c7b1ca67..24d5f2b5 100644 --- a/charts/flagger/Chart.yaml +++ b/charts/flagger/Chart.yaml @@ -1,7 +1,7 @@ apiVersion: v1 name: flagger -version: 0.13.1 -appVersion: 0.13.1 +version: 0.13.2 +appVersion: 0.13.2 kubeVersion: ">=1.11.0-0" engine: gotpl description: Flagger is a Kubernetes operator that automates the promotion of canary deployments using Istio, App Mesh or NGINX routing for traffic shifting and Prometheus metrics for canary analysis. diff --git a/charts/flagger/values.yaml b/charts/flagger/values.yaml index 6b6dafdd..3b9e823a 100644 --- a/charts/flagger/values.yaml +++ b/charts/flagger/values.yaml @@ -2,7 +2,7 @@ image: repository: weaveworks/flagger - tag: 0.13.1 + tag: 0.13.2 pullPolicy: IfNotPresent metricsServer: "http://prometheus:9090" diff --git a/pkg/version/version.go b/pkg/version/version.go index d9f5035d..b04a3033 100644 --- a/pkg/version/version.go +++ b/pkg/version/version.go @@ -1,4 +1,4 @@ package version -var VERSION = "0.13.1" +var VERSION = "0.13.2" var REVISION = "unknown" From 0032c14a78a24b0ee6d0aadb1b494e2d1ace62ea Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Mon, 13 May 2019 17:34:08 +0300 Subject: [PATCH 11/14] Refactor metrics - add observer interface with builtin metrics functions - add metrics observer factory - add prometheus client - implement the observer interface for istio, envoy and nginx - remove deprecated istio and app mesh metric aliases (istio_requests_total, istio_request_duration_seconds_bucket, envoy_cluster_upstream_rq, envoy_cluster_upstream_rq_time_bucket) --- cmd/flagger/main.go | 22 ++-- pkg/controller/controller.go | 69 +++++------ pkg/controller/controller_test.go | 34 +++--- pkg/controller/scheduler.go | 148 +++++++----------------- pkg/metrics/client.go | 184 ++++++++++++++++++++++++++++++ pkg/metrics/client_test.go | 85 ++++++++++++++ pkg/metrics/envoy.go | 150 +++++++++--------------- pkg/metrics/envoy_test.go | 79 ++++++++----- pkg/metrics/factory.go | 39 +++++++ pkg/metrics/istio.go | 157 +++++++++---------------- pkg/metrics/istio_test.go | 79 ++++++++----- pkg/metrics/nginx.go | 160 ++++++++++---------------- pkg/metrics/nginx_test.go | 79 ++++++++----- pkg/metrics/observer.go | 182 +---------------------------- pkg/metrics/observer_test.go | 119 ------------------- 15 files changed, 735 insertions(+), 851 deletions(-) create mode 100644 pkg/metrics/client.go create mode 100644 pkg/metrics/client_test.go create mode 100644 pkg/metrics/factory.go delete mode 100644 pkg/metrics/observer_test.go diff --git a/cmd/flagger/main.go b/cmd/flagger/main.go index 0c09b7dc..da6a61d6 100644 --- a/cmd/flagger/main.go +++ b/cmd/flagger/main.go @@ -45,8 +45,8 @@ var ( func init() { flag.StringVar(&kubeconfig, "kubeconfig", "", "Path to a kubeconfig. Only required if out-of-cluster.") flag.StringVar(&masterURL, "master", "", "The address of the Kubernetes API server. Overrides any value in kubeconfig. Only required if out-of-cluster.") - flag.StringVar(&metricsServer, "metrics-server", "http://prometheus:9090", "Prometheus URL") - flag.DurationVar(&controlLoopInterval, "control-loop-interval", 10*time.Second, "Kubernetes API sync interval") + flag.StringVar(&metricsServer, "metrics-server", "http://prometheus:9090", "Prometheus URL.") + flag.DurationVar(&controlLoopInterval, "control-loop-interval", 10*time.Second, "Kubernetes API sync interval.") flag.StringVar(&logLevel, "log-level", "debug", "Log level can be: debug, info, warning, error.") flag.StringVar(&port, "port", "8080", "Port to listen on.") flag.StringVar(&slackURL, "slack-url", "", "Slack hook URL.") @@ -55,9 +55,9 @@ func init() { flag.IntVar(&threadiness, "threadiness", 2, "Worker concurrency.") flag.BoolVar(&zapReplaceGlobals, "zap-replace-globals", false, "Whether to change the logging level of the global zap logger.") flag.StringVar(&zapEncoding, "zap-encoding", "json", "Zap logger encoding.") - flag.StringVar(&namespace, "namespace", "", "Namespace that flagger would watch canary object") - flag.StringVar(&meshProvider, "mesh-provider", "istio", "Service mesh provider, can be istio or appmesh") - flag.StringVar(&selectorLabels, "selector-labels", "app,name,app.kubernetes.io/name", "List of pod labels that Flagger uses to create pod selectors") + flag.StringVar(&namespace, "namespace", "", "Namespace that flagger would watch canary object.") + flag.StringVar(&meshProvider, "mesh-provider", "istio", "Service mesh provider, can be istio, appmesh, supergloo, nginx or smi.") + flag.StringVar(&selectorLabels, "selector-labels", "app,name,app.kubernetes.io/name", "List of pod labels that Flagger uses to create pod selectors.") } func main() { @@ -87,12 +87,12 @@ func main() { meshClient, err := clientset.NewForConfig(cfg) if err != nil { - logger.Fatalf("Error building istio clientset: %v", err) + logger.Fatalf("Error building mesh clientset: %v", err) } flaggerClient, err := clientset.NewForConfig(cfg) if err != nil { - logger.Fatalf("Error building example clientset: %s", err.Error()) + logger.Fatalf("Error building flagger clientset: %s", err.Error()) } flaggerInformerFactory := informers.NewSharedInformerFactoryWithOptions(flaggerClient, time.Second*30, informers.WithNamespace(namespace)) @@ -116,7 +116,12 @@ func main() { logger.Infof("Watching namespace %s", namespace) } - ok, err := metrics.CheckMetricsServer(metricsServer) + observerFactory, err := metrics.NewFactory(metricsServer, meshProvider, 5*time.Second) + if err != nil { + logger.Fatalf("Error building prometheus client: %s", err.Error()) + } + + ok, err := observerFactory.Client.IsOnline() if ok { logger.Infof("Connected to metrics server %s", metricsServer) } else { @@ -148,6 +153,7 @@ func main() { logger, slack, routerFactory, + observerFactory, meshProvider, version.VERSION, labels, diff --git a/pkg/controller/controller.go b/pkg/controller/controller.go index 662f4f01..4f44adee 100644 --- a/pkg/controller/controller.go +++ b/pkg/controller/controller.go @@ -33,23 +33,23 @@ const controllerAgentName = "flagger" // Controller is managing the canary objects and schedules canary deployments type Controller struct { - kubeClient kubernetes.Interface - istioClient clientset.Interface - flaggerClient clientset.Interface - flaggerLister flaggerlisters.CanaryLister - flaggerSynced cache.InformerSynced - flaggerWindow time.Duration - workqueue workqueue.RateLimitingInterface - eventRecorder record.EventRecorder - logger *zap.SugaredLogger - canaries *sync.Map - jobs map[string]CanaryJob - deployer canary.Deployer - observer metrics.Observer - recorder metrics.Recorder - notifier *notifier.Slack - routerFactory *router.Factory - meshProvider string + kubeClient kubernetes.Interface + istioClient clientset.Interface + flaggerClient clientset.Interface + flaggerLister flaggerlisters.CanaryLister + flaggerSynced cache.InformerSynced + flaggerWindow time.Duration + workqueue workqueue.RateLimitingInterface + eventRecorder record.EventRecorder + logger *zap.SugaredLogger + canaries *sync.Map + jobs map[string]CanaryJob + deployer canary.Deployer + recorder metrics.Recorder + notifier *notifier.Slack + routerFactory *router.Factory + observerFactory *metrics.Factory + meshProvider string } func NewController( @@ -62,6 +62,7 @@ func NewController( logger *zap.SugaredLogger, notifier *notifier.Slack, routerFactory *router.Factory, + observerFactory *metrics.Factory, meshProvider string, version string, labels []string, @@ -92,23 +93,23 @@ func NewController( recorder.SetInfo(version, meshProvider) ctrl := &Controller{ - kubeClient: kubeClient, - istioClient: istioClient, - flaggerClient: flaggerClient, - flaggerLister: flaggerInformer.Lister(), - flaggerSynced: flaggerInformer.Informer().HasSynced, - workqueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), controllerAgentName), - eventRecorder: eventRecorder, - logger: logger, - canaries: new(sync.Map), - jobs: map[string]CanaryJob{}, - flaggerWindow: flaggerWindow, - deployer: deployer, - observer: metrics.NewObserver(metricServer), - recorder: recorder, - notifier: notifier, - routerFactory: routerFactory, - meshProvider: meshProvider, + kubeClient: kubeClient, + istioClient: istioClient, + flaggerClient: flaggerClient, + flaggerLister: flaggerInformer.Lister(), + flaggerSynced: flaggerInformer.Informer().HasSynced, + workqueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), controllerAgentName), + eventRecorder: eventRecorder, + logger: logger, + canaries: new(sync.Map), + jobs: map[string]CanaryJob{}, + flaggerWindow: flaggerWindow, + deployer: deployer, + observerFactory: observerFactory, + recorder: recorder, + notifier: notifier, + routerFactory: routerFactory, + meshProvider: meshProvider, } flaggerInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ diff --git a/pkg/controller/controller_test.go b/pkg/controller/controller_test.go index 28ac9fcd..423fb27f 100644 --- a/pkg/controller/controller_test.go +++ b/pkg/controller/controller_test.go @@ -37,7 +37,6 @@ type Mocks struct { meshClient clientset.Interface flaggerClient clientset.Interface deployer canary.Deployer - observer metrics.Observer ctrl *Controller logger *zap.SugaredLogger router router.Interface @@ -77,7 +76,6 @@ func SetupMocks(abtest bool) Mocks { FlaggerClient: flaggerClient, }, } - observer := metrics.NewObserver("fake") // init controller flaggerInformerFactory := informers.NewSharedInformerFactory(flaggerClient, noResyncPeriodFunc()) @@ -86,21 +84,24 @@ func SetupMocks(abtest bool) Mocks { // init router rf := router.NewFactory(nil, kubeClient, flaggerClient, logger, flaggerClient) + // init observer + observerFactory, _ := metrics.NewFactory("fake", "istio", 5*time.Second) + ctrl := &Controller{ - kubeClient: kubeClient, - istioClient: flaggerClient, - flaggerClient: flaggerClient, - flaggerLister: flaggerInformer.Lister(), - flaggerSynced: flaggerInformer.Informer().HasSynced, - workqueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), controllerAgentName), - eventRecorder: &record.FakeRecorder{}, - logger: logger, - canaries: new(sync.Map), - flaggerWindow: time.Second, - deployer: deployer, - observer: observer, - recorder: metrics.NewRecorder(controllerAgentName, false), - routerFactory: rf, + kubeClient: kubeClient, + istioClient: flaggerClient, + flaggerClient: flaggerClient, + flaggerLister: flaggerInformer.Lister(), + flaggerSynced: flaggerInformer.Informer().HasSynced, + workqueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), controllerAgentName), + eventRecorder: &record.FakeRecorder{}, + logger: logger, + canaries: new(sync.Map), + flaggerWindow: time.Second, + deployer: deployer, + observerFactory: observerFactory, + recorder: metrics.NewRecorder(controllerAgentName, false), + routerFactory: rf, } ctrl.flaggerSynced = alwaysReady @@ -108,7 +109,6 @@ func SetupMocks(abtest bool) Mocks { return Mocks{ canary: c, - observer: observer, deployer: deployer, logger: logger, flaggerClient: flaggerClient, diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index 0fc06978..e15b4642 100644 --- a/pkg/controller/scheduler.go +++ b/pkg/controller/scheduler.go @@ -556,126 +556,56 @@ func (c *Controller) analyseCanary(r *flaggerv1.Canary) bool { } } + // create observer based on the mesh provider + observer := c.observerFactory.Observer() + // run metrics checks for _, metric := range r.Spec.CanaryAnalysis.Metrics { if metric.Interval == "" { metric.Interval = r.GetMetricInterval() } - // App Mesh checks - if c.meshProvider == "appmesh" { - if metric.Name == "request-success-rate" || metric.Name == "envoy_cluster_upstream_rq" { - val, err := c.observer.GetEnvoySuccessRate(r.Spec.TargetRef.Name, r.Namespace, metric.Name, metric.Interval) - if err != nil { - if strings.Contains(err.Error(), "no values found") { - c.recordEventWarningf(r, "Halt advancement no values found for metric %s probably %s.%s is not receiving traffic", - metric.Name, r.Spec.TargetRef.Name, r.Namespace) - } else { - c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.observer.GetMetricsServer(), err) - } - return false - } - if float64(metric.Threshold) > val { - c.recordEventWarningf(r, "Halt %s.%s advancement success rate %.2f%% < %v%%", - r.Name, r.Namespace, val, metric.Threshold) - return false - } - } - - if metric.Name == "request-duration" || metric.Name == "envoy_cluster_upstream_rq_time_bucket" { - val, err := c.observer.GetEnvoyRequestDuration(r.Spec.TargetRef.Name, r.Namespace, metric.Name, metric.Interval) - if err != nil { - c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.observer.GetMetricsServer(), err) - return false - } - t := time.Duration(metric.Threshold) * time.Millisecond - if val > t { - c.recordEventWarningf(r, "Halt %s.%s advancement request duration %v > %v", - r.Name, r.Namespace, val, t) - return false - } - } - } - - // Istio checks - if c.meshProvider == "istio" { - if metric.Name == "request-success-rate" || metric.Name == "istio_requests_total" { - val, err := c.observer.GetIstioSuccessRate(r.Spec.TargetRef.Name, r.Namespace, metric.Name, metric.Interval) - if err != nil { - if strings.Contains(err.Error(), "no values found") { - c.recordEventWarningf(r, "Halt advancement no values found for metric %s probably %s.%s is not receiving traffic", - metric.Name, r.Spec.TargetRef.Name, r.Namespace) - } else { - c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.observer.GetMetricsServer(), err) - } - return false - } - if float64(metric.Threshold) > val { - c.recordEventWarningf(r, "Halt %s.%s advancement success rate %.2f%% < %v%%", - r.Name, r.Namespace, val, metric.Threshold) - return false - } - } - - if metric.Name == "request-duration" || metric.Name == "istio_request_duration_seconds_bucket" { - val, err := c.observer.GetIstioRequestDuration(r.Spec.TargetRef.Name, r.Namespace, metric.Name, metric.Interval) - if err != nil { - c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.observer.GetMetricsServer(), err) - return false - } - t := time.Duration(metric.Threshold) * time.Millisecond - if val > t { - c.recordEventWarningf(r, "Halt %s.%s advancement request duration %v > %v", - r.Name, r.Namespace, val, t) - return false - } - } - } - - // NGINX checks - if c.meshProvider == "nginx" { - if metric.Name == "request-success-rate" { - val, err := c.observer.GetNginxSuccessRate(r.Spec.IngressRef.Name, r.Namespace, metric.Name, metric.Interval) - if err != nil { - if strings.Contains(err.Error(), "no values found") { - c.recordEventWarningf(r, "Halt advancement no values found for metric %s probably %s.%s is not receiving traffic", - metric.Name, r.Spec.TargetRef.Name, r.Namespace) - } else { - c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.observer.GetMetricsServer(), err) - } - return false - } - if float64(metric.Threshold) > val { - c.recordEventWarningf(r, "Halt %s.%s advancement success rate %.2f%% < %v%%", - r.Name, r.Namespace, val, metric.Threshold) - return false - } - } - - if metric.Name == "request-duration" { - val, err := c.observer.GetNginxRequestDuration(r.Spec.IngressRef.Name, r.Namespace, metric.Name, metric.Interval) - if err != nil { - c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.observer.GetMetricsServer(), err) - return false - } - t := time.Duration(metric.Threshold) * time.Millisecond - if val > t { - c.recordEventWarningf(r, "Halt %s.%s advancement request duration %v > %v", - r.Name, r.Namespace, val, t) - return false - } - } - } - - // custom checks - if metric.Query != "" { - val, err := c.observer.GetScalar(metric.Query) + if metric.Name == "request-success-rate" { + val, err := observer.GetRequestSuccessRate(r.Spec.TargetRef.Name, r.Namespace, metric.Interval) if err != nil { if strings.Contains(err.Error(), "no values found") { c.recordEventWarningf(r, "Halt advancement no values found for metric %s probably %s.%s is not receiving traffic", metric.Name, r.Spec.TargetRef.Name, r.Namespace) } else { - c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.observer.GetMetricsServer(), err) + c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.observerFactory.Client.GetMetricsServer(), err) + } + return false + } + if float64(metric.Threshold) > val { + c.recordEventWarningf(r, "Halt %s.%s advancement success rate %.2f%% < %v%%", + r.Name, r.Namespace, val, metric.Threshold) + return false + } + } + + if metric.Name == "request-duration" { + val, err := observer.GetRequestDuration(r.Spec.TargetRef.Name, r.Namespace, metric.Interval) + if err != nil { + c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.observerFactory.Client.GetMetricsServer(), err) + return false + } + t := time.Duration(metric.Threshold) * time.Millisecond + if val > t { + c.recordEventWarningf(r, "Halt %s.%s advancement request duration %v > %v", + r.Name, r.Namespace, val, t) + return false + } + } + + // custom checks + if metric.Query != "" { + val, err := c.observerFactory.Client.RunQuery(metric.Query) + if err != nil { + if strings.Contains(err.Error(), "no values found") { + c.recordEventWarningf(r, "Halt advancement no values found for metric %s probably %s.%s is not receiving traffic", + metric.Name, r.Spec.TargetRef.Name, r.Namespace) + } else { + c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.observerFactory.Client.GetMetricsServer(), err) } return false } diff --git a/pkg/metrics/client.go b/pkg/metrics/client.go new file mode 100644 index 00000000..60541aff --- /dev/null +++ b/pkg/metrics/client.go @@ -0,0 +1,184 @@ +package metrics + +import ( + "bufio" + "bytes" + "context" + "encoding/json" + "fmt" + "io/ioutil" + "net/http" + "net/url" + "strconv" + "strings" + "text/template" + "time" +) + +// PrometheusClient is executing promql queries +type PrometheusClient struct { + timeout time.Duration + url url.URL +} + +type prometheusResponse struct { + Data struct { + Result []struct { + Metric struct { + Name string `json:"name"` + } + Value []interface{} `json:"value"` + } + } +} + +// NewPrometheusClient creates a Prometheus client for the provided URL address +func NewPrometheusClient(address string, timeout time.Duration) (*PrometheusClient, error) { + promURL, err := url.Parse(address) + if err != nil { + return nil, err + } + + return &PrometheusClient{timeout: timeout, url: *promURL}, nil +} + +// RenderQuery renders the promql query using the provided text template +func (p *PrometheusClient) RenderQuery(name string, namespace string, interval string, tmpl string) (string, error) { + meta := struct { + Name string + Namespace string + Interval string + }{ + name, + namespace, + interval, + } + + t, err := template.New("tmpl").Parse(tmpl) + if err != nil { + return "", err + } + var data bytes.Buffer + b := bufio.NewWriter(&data) + + if err := t.Execute(b, meta); err != nil { + return "", err + } + + err = b.Flush() + if err != nil { + return "", err + } + + return data.String(), nil +} + +// RunQuery executes the promql and converts the result to float64 +func (p *PrometheusClient) RunQuery(query string) (float64, error) { + if p.url.Host == "fake" { + return 100, nil + } + + query = url.QueryEscape(p.TrimQuery(query)) + u, err := url.Parse(fmt.Sprintf("./api/v1/query?query=%s", query)) + if err != nil { + return 0, err + } + + u = p.url.ResolveReference(u) + + req, err := http.NewRequest("GET", u.String(), nil) + if err != nil { + return 0, err + } + + ctx, cancel := context.WithTimeout(req.Context(), p.timeout) + defer cancel() + + r, err := http.DefaultClient.Do(req.WithContext(ctx)) + if err != nil { + return 0, err + } + defer r.Body.Close() + + b, err := ioutil.ReadAll(r.Body) + if err != nil { + return 0, fmt.Errorf("error reading body: %s", err.Error()) + } + + if 400 <= r.StatusCode { + return 0, fmt.Errorf("error response: %s", string(b)) + } + + var result prometheusResponse + err = json.Unmarshal(b, &result) + if err != nil { + return 0, fmt.Errorf("error unmarshaling result: %s, '%s'", err.Error(), string(b)) + } + + var value *float64 + for _, v := range result.Data.Result { + metricValue := v.Value[1] + switch metricValue.(type) { + case string: + f, err := strconv.ParseFloat(metricValue.(string), 64) + if err != nil { + return 0, err + } + value = &f + } + } + if value == nil { + return 0, fmt.Errorf("no values found") + } + + return *value, nil +} + +// TrimQuery takes a promql query and removes spaces, tabs and new lines +func (p *PrometheusClient) TrimQuery(query string) string { + query = strings.Replace(query, "\n", "", -1) + query = strings.Replace(query, "\t", "", -1) + query = strings.Replace(query, " ", "", -1) + + return query +} + +// IsOnline call Prometheus status endpoint and returns an error if the API is unreachable +func (p *PrometheusClient) IsOnline() (bool, error) { + u, err := url.Parse("./api/v1/status/flags") + if err != nil { + return false, err + } + + u = p.url.ResolveReference(u) + + req, err := http.NewRequest("GET", u.String(), nil) + if err != nil { + return false, err + } + + ctx, cancel := context.WithTimeout(req.Context(), p.timeout) + defer cancel() + + r, err := http.DefaultClient.Do(req.WithContext(ctx)) + if err != nil { + return false, err + } + defer r.Body.Close() + + b, err := ioutil.ReadAll(r.Body) + if err != nil { + return false, fmt.Errorf("error reading body: %s", err.Error()) + } + + if 400 <= r.StatusCode { + return false, fmt.Errorf("error response: %s", string(b)) + } + + return true, nil +} + +func (p *PrometheusClient) GetMetricsServer() string { + return p.url.RawQuery +} diff --git a/pkg/metrics/client_test.go b/pkg/metrics/client_test.go new file mode 100644 index 00000000..9fdbe85a --- /dev/null +++ b/pkg/metrics/client_test.go @@ -0,0 +1,85 @@ +package metrics + +import ( + "net/http" + "net/http/httptest" + "testing" + "time" +) + +func TestPrometheusClient_RunQuery(t *testing.T) { + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + json := `{"status":"success","data":{"resultType":"vector","result":[{"metric":{},"value":[1545905245.458,"100"]}]}}` + w.Write([]byte(json)) + })) + defer ts.Close() + + client, err := NewPrometheusClient(ts.URL, time.Second) + if err != nil { + t.Fatal(err) + } + + query := ` + histogram_quantile(0.99, + sum( + rate( + http_request_duration_seconds_bucket{ + kubernetes_namespace="test", + kubernetes_pod_name=~"podinfo-[0-9a-zA-Z]+(-[0-9a-zA-Z]+)" + }[1m] + ) + ) by (le) + )` + + val, err := client.RunQuery(query) + if err != nil { + t.Fatal(err.Error()) + } + + if val != 100 { + t.Errorf("Got %v wanted %v", val, 100) + } +} + +func TestPrometheusClient_IsOnline(t *testing.T) { + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + json := `{"status":"success","data":{"config.file":"/etc/prometheus/prometheus.yml"}}` + w.Write([]byte(json)) + })) + defer ts.Close() + + client, err := NewPrometheusClient(ts.URL, time.Second) + if err != nil { + t.Fatal(err) + } + + ok, err := client.IsOnline() + if err != nil { + t.Fatal(err.Error()) + } + + if !ok { + t.Errorf("Got %v wanted %v", ok, true) + } +} + +func TestPrometheusClient_IsOffline(t *testing.T) { + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusBadGateway) + })) + defer ts.Close() + + client, err := NewPrometheusClient(ts.URL, time.Second) + if err != nil { + t.Fatal(err) + } + + ok, err := client.IsOnline() + if err == nil { + t.Errorf("Got no error wanted %v", http.StatusBadGateway) + } + + if ok { + t.Errorf("Got %v wanted %v", ok, false) + } +} diff --git a/pkg/metrics/envoy.go b/pkg/metrics/envoy.go index 5373f532..ce88b540 100644 --- a/pkg/metrics/envoy.go +++ b/pkg/metrics/envoy.go @@ -1,119 +1,73 @@ package metrics import ( - "fmt" - "net/url" - "strconv" "time" ) -const envoySuccessRateQuery = ` -sum(rate( -envoy_cluster_upstream_rq{kubernetes_namespace="{{ .Namespace }}", -kubernetes_pod_name=~"{{ .Name }}-[0-9a-zA-Z]+(-[0-9a-zA-Z]+)", -envoy_response_code!~"5.*"} -[{{ .Interval }}])) -/ -sum(rate( -envoy_cluster_upstream_rq{kubernetes_namespace="{{ .Namespace }}", -kubernetes_pod_name=~"{{ .Name }}-[0-9a-zA-Z]+(-[0-9a-zA-Z]+)"} -[{{ .Interval }}])) -* 100 -` - -func (c *Observer) GetEnvoySuccessRate(name string, namespace string, metric string, interval string) (float64, error) { - if c.metricsServer == "fake" { - return 100, nil - } - - meta := struct { - Name string - Namespace string - Interval string - }{ - name, - namespace, - interval, - } - - query, err := render(meta, envoySuccessRateQuery) - if err != nil { - return 0, err - } - - var rate *float64 - querySt := url.QueryEscape(query) - result, err := c.queryMetric(querySt) - if err != nil { - return 0, err - } - - for _, v := range result.Data.Result { - metricValue := v.Value[1] - switch metricValue.(type) { - case string: - f, err := strconv.ParseFloat(metricValue.(string), 64) - if err != nil { - return 0, err - } - rate = &f - } - } - if rate == nil { - return 0, fmt.Errorf("no values found for metric %s", metric) - } - return *rate, nil +var envoyQueries = map[string]string{ + "request-success-rate": ` + sum( + rate( + envoy_cluster_upstream_rq{ + kubernetes_namespace="{{ .Namespace }}", + kubernetes_pod_name=~"{{ .Name }}-[0-9a-zA-Z]+(-[0-9a-zA-Z]+)", + envoy_response_code!~"5.*" + }[{{ .Interval }}] + ) + ) + / + sum( + rate( + envoy_cluster_upstream_rq{ + kubernetes_namespace="{{ .Namespace }}", + kubernetes_pod_name=~"{{ .Name }}-[0-9a-zA-Z]+(-[0-9a-zA-Z]+)" + }[{{ .Interval }}] + ) + ) + * 100`, + "request-duration": ` + histogram_quantile( + 0.99, + sum( + rate( + envoy_cluster_upstream_rq_time_bucket{ + kubernetes_namespace="{{ .Namespace }}", + kubernetes_pod_name=~"{{ .Name }}-[0-9a-zA-Z]+(-[0-9a-zA-Z]+)" + }[{{ .Interval }}] + ) + ) by (le) + )`, } -const envoyRequestDurationQuery = ` -histogram_quantile(0.99, sum(rate( -envoy_cluster_upstream_rq_time_bucket{kubernetes_namespace="{{ .Namespace }}", -kubernetes_pod_name=~"{{ .Name }}-[0-9a-zA-Z]+(-[0-9a-zA-Z]+)"} -[{{ .Interval }}])) by (le)) -` +type EnvoyObserver struct { + client *PrometheusClient +} -// GetEnvoyRequestDuration returns the 99P requests delay using envoy_cluster_upstream_rq_time_bucket metrics -func (c *Observer) GetEnvoyRequestDuration(name string, namespace string, metric string, interval string) (time.Duration, error) { - if c.metricsServer == "fake" { - return 1, nil - } - - meta := struct { - Name string - Namespace string - Interval string - }{ - name, - namespace, - interval, - } - - query, err := render(meta, envoyRequestDurationQuery) +func (ob *EnvoyObserver) GetRequestSuccessRate(name string, namespace string, interval string) (float64, error) { + query, err := ob.client.RenderQuery(name, namespace, interval, envoyQueries["request-success-rate"]) if err != nil { return 0, err } - var rate *float64 - querySt := url.QueryEscape(query) - result, err := c.queryMetric(querySt) + value, err := ob.client.RunQuery(query) if err != nil { return 0, err } - for _, v := range result.Data.Result { - metricValue := v.Value[1] - switch metricValue.(type) { - case string: - f, err := strconv.ParseFloat(metricValue.(string), 64) - if err != nil { - return 0, err - } - rate = &f - } + return value, nil +} + +func (ob *EnvoyObserver) GetRequestDuration(name string, namespace string, interval string) (time.Duration, error) { + query, err := ob.client.RenderQuery(name, namespace, interval, envoyQueries["request-duration"]) + if err != nil { + return 0, err } - if rate == nil { - return 0, fmt.Errorf("no values found for metric %s", metric) + + value, err := ob.client.RunQuery(query) + if err != nil { + return 0, err } - ms := time.Duration(int64(*rate)) * time.Millisecond + + ms := time.Duration(int64(value)) * time.Millisecond return ms, nil } diff --git a/pkg/metrics/envoy_test.go b/pkg/metrics/envoy_test.go index 0ae9a96c..85f27905 100644 --- a/pkg/metrics/envoy_test.go +++ b/pkg/metrics/envoy_test.go @@ -1,51 +1,74 @@ package metrics import ( + "net/http" + "net/http/httptest" "testing" + "time" ) -func Test_EnvoySuccessRateQueryRender(t *testing.T) { - meta := struct { - Name string - Namespace string - Interval string - }{ - "podinfo", - "default", - "1m", - } +func TestEnvoyObserver_GetRequestSuccessRate(t *testing.T) { + expected := `sum(rate(envoy_cluster_upstream_rq{kubernetes_namespace="default",kubernetes_pod_name=~"podinfo-[0-9a-zA-Z]+(-[0-9a-zA-Z]+)",envoy_response_code!~"5.*"}[1m]))/sum(rate(envoy_cluster_upstream_rq{kubernetes_namespace="default",kubernetes_pod_name=~"podinfo-[0-9a-zA-Z]+(-[0-9a-zA-Z]+)"}[1m]))*100` - query, err := render(meta, envoySuccessRateQuery) + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + promql := r.URL.Query()["query"][0] + if promql != expected { + t.Errorf("\nGot %s \nWanted %s", promql, expected) + } + + json := `{"status":"success","data":{"resultType":"vector","result":[{"metric":{},"value":[1,"100"]}]}}` + w.Write([]byte(json)) + })) + defer ts.Close() + + client, err := NewPrometheusClient(ts.URL, time.Second) if err != nil { t.Fatal(err) } - expected := `sum(rate(envoy_cluster_upstream_rq{kubernetes_namespace="default",kubernetes_pod_name=~"podinfo-[0-9a-zA-Z]+(-[0-9a-zA-Z]+)",envoy_response_code!~"5.*"}[1m])) / sum(rate(envoy_cluster_upstream_rq{kubernetes_namespace="default",kubernetes_pod_name=~"podinfo-[0-9a-zA-Z]+(-[0-9a-zA-Z]+)"}[1m])) * 100` + observer := &EnvoyObserver{ + client: client, + } - if query != expected { - t.Errorf("\nGot %s \nWanted %s", query, expected) + val, err := observer.GetRequestSuccessRate("podinfo", "default", "1m") + if err != nil { + t.Fatal(err.Error()) + } + + if val != 100 { + t.Errorf("Got %v wanted %v", val, 100) } } -func Test_EnvoyRequestDurationQueryRender(t *testing.T) { - meta := struct { - Name string - Namespace string - Interval string - }{ - "podinfo", - "default", - "1m", - } +func TestEnvoyObserver_GetRequestDuration(t *testing.T) { + expected := `histogram_quantile(0.99,sum(rate(envoy_cluster_upstream_rq_time_bucket{kubernetes_namespace="default",kubernetes_pod_name=~"podinfo-[0-9a-zA-Z]+(-[0-9a-zA-Z]+)"}[1m]))by(le))` - query, err := render(meta, envoyRequestDurationQuery) + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + promql := r.URL.Query()["query"][0] + if promql != expected { + t.Errorf("\nGot %s \nWanted %s", promql, expected) + } + + json := `{"status":"success","data":{"resultType":"vector","result":[{"metric":{},"value":[1,"100"]}]}}` + w.Write([]byte(json)) + })) + defer ts.Close() + + client, err := NewPrometheusClient(ts.URL, time.Second) if err != nil { t.Fatal(err) } - expected := `histogram_quantile(0.99, sum(rate(envoy_cluster_upstream_rq_time_bucket{kubernetes_namespace="default",kubernetes_pod_name=~"podinfo-[0-9a-zA-Z]+(-[0-9a-zA-Z]+)"}[1m])) by (le))` + observer := &EnvoyObserver{ + client: client, + } - if query != expected { - t.Errorf("\nGot %s \nWanted %s", query, expected) + val, err := observer.GetRequestDuration("podinfo", "default", "1m") + if err != nil { + t.Fatal(err.Error()) + } + + if val != 100*time.Millisecond { + t.Errorf("Got %v wanted %v", val, 100*time.Millisecond) } } diff --git a/pkg/metrics/factory.go b/pkg/metrics/factory.go new file mode 100644 index 00000000..c717b3a6 --- /dev/null +++ b/pkg/metrics/factory.go @@ -0,0 +1,39 @@ +package metrics + +import ( + "time" +) + +type Factory struct { + MeshProvider string + Client *PrometheusClient +} + +func NewFactory(metricsServer string, meshProvider string, timeout time.Duration) (*Factory, error) { + client, err := NewPrometheusClient(metricsServer, timeout) + if err != nil { + return nil, err + } + + return &Factory{ + MeshProvider: meshProvider, + Client: client, + }, nil +} + +func (factory Factory) Observer() Interface { + switch { + case factory.MeshProvider == "appmesh": + return &EnvoyObserver{ + client: factory.Client, + } + case factory.MeshProvider == "nginx": + return &NginxObserver{ + client: factory.Client, + } + default: + return &IstioObserver{ + client: factory.Client, + } + } +} diff --git a/pkg/metrics/istio.go b/pkg/metrics/istio.go index 5c8e9855..f6fc9c9a 100644 --- a/pkg/metrics/istio.go +++ b/pkg/metrics/istio.go @@ -1,123 +1,76 @@ package metrics import ( - "fmt" - "net/url" - "strconv" "time" ) -const istioSuccessRateQuery = ` -sum(rate( -istio_requests_total{reporter="destination", -destination_workload_namespace="{{ .Namespace }}", -destination_workload=~"{{ .Name }}", -response_code!~"5.*"} -[{{ .Interval }}])) -/ -sum(rate( -istio_requests_total{reporter="destination", -destination_workload_namespace="{{ .Namespace }}", -destination_workload=~"{{ .Name }}"} -[{{ .Interval }}])) -* 100 -` - -// GetIstioSuccessRate returns the requests success rate (non 5xx) using istio_requests_total metric -func (c *Observer) GetIstioSuccessRate(name string, namespace string, metric string, interval string) (float64, error) { - if c.metricsServer == "fake" { - return 100, nil - } - - meta := struct { - Name string - Namespace string - Interval string - }{ - name, - namespace, - interval, - } - - query, err := render(meta, istioSuccessRateQuery) - if err != nil { - return 0, err - } - - var rate *float64 - querySt := url.QueryEscape(query) - result, err := c.queryMetric(querySt) - if err != nil { - return 0, err - } - - for _, v := range result.Data.Result { - metricValue := v.Value[1] - switch metricValue.(type) { - case string: - f, err := strconv.ParseFloat(metricValue.(string), 64) - if err != nil { - return 0, err - } - rate = &f - } - } - if rate == nil { - return 0, fmt.Errorf("no values found for metric %s", metric) - } - return *rate, nil +var istioQueries = map[string]string{ + "request-success-rate": ` + sum( + rate( + istio_requests_total{ + reporter="destination", + destination_workload_namespace="{{ .Namespace }}", + destination_workload=~"{{ .Name }}", + response_code!~"5.*" + }[{{ .Interval }}] + ) + ) + / + sum( + rate( + istio_requests_total{ + reporter="destination", + destination_workload_namespace="{{ .Namespace }}", + destination_workload=~"{{ .Name }}" + }[{{ .Interval }}] + ) + ) + * 100`, + "request-duration": ` + histogram_quantile( + 0.99, + sum( + rate( + istio_request_duration_seconds_bucket{ + reporter="destination", + destination_workload_namespace="{{ .Namespace }}", + destination_workload=~"{{ .Name }}" + }[{{ .Interval }}] + ) + ) by (le) + )`, } -const istioRequestDurationQuery = ` -histogram_quantile(0.99, sum(rate( -istio_request_duration_seconds_bucket{reporter="destination", -destination_workload_namespace="{{ .Namespace }}", -destination_workload=~"{{ .Name }}"} -[{{ .Interval }}])) by (le)) -` +type IstioObserver struct { + client *PrometheusClient +} -// GetIstioRequestDuration returns the 99P requests delay using istio_request_duration_seconds_bucket metrics -func (c *Observer) GetIstioRequestDuration(name string, namespace string, metric string, interval string) (time.Duration, error) { - if c.metricsServer == "fake" { - return 1, nil - } - - meta := struct { - Name string - Namespace string - Interval string - }{ - name, - namespace, - interval, - } - - query, err := render(meta, istioRequestDurationQuery) +func (ob *IstioObserver) GetRequestSuccessRate(name string, namespace string, interval string) (float64, error) { + query, err := ob.client.RenderQuery(name, namespace, interval, istioQueries["request-success-rate"]) if err != nil { return 0, err } - var rate *float64 - querySt := url.QueryEscape(query) - result, err := c.queryMetric(querySt) + value, err := ob.client.RunQuery(query) if err != nil { return 0, err } - for _, v := range result.Data.Result { - metricValue := v.Value[1] - switch metricValue.(type) { - case string: - f, err := strconv.ParseFloat(metricValue.(string), 64) - if err != nil { - return 0, err - } - rate = &f - } + return value, nil +} + +func (ob *IstioObserver) GetRequestDuration(name string, namespace string, interval string) (time.Duration, error) { + query, err := ob.client.RenderQuery(name, namespace, interval, istioQueries["request-duration"]) + if err != nil { + return 0, err } - if rate == nil { - return 0, fmt.Errorf("no values found for metric %s", metric) + + value, err := ob.client.RunQuery(query) + if err != nil { + return 0, err } - ms := time.Duration(int64(*rate*1000)) * time.Millisecond + + ms := time.Duration(int64(value)) * time.Millisecond return ms, nil } diff --git a/pkg/metrics/istio_test.go b/pkg/metrics/istio_test.go index 28826a77..dfba3c05 100644 --- a/pkg/metrics/istio_test.go +++ b/pkg/metrics/istio_test.go @@ -1,51 +1,74 @@ package metrics import ( + "net/http" + "net/http/httptest" "testing" + "time" ) -func Test_IstioSuccessRateQueryRender(t *testing.T) { - meta := struct { - Name string - Namespace string - Interval string - }{ - "podinfo", - "default", - "1m", - } +func TestIstioObserver_GetRequestSuccessRate(t *testing.T) { + expected := `sum(rate(istio_requests_total{reporter="destination",destination_workload_namespace="default",destination_workload=~"podinfo",response_code!~"5.*"}[1m]))/sum(rate(istio_requests_total{reporter="destination",destination_workload_namespace="default",destination_workload=~"podinfo"}[1m]))*100` - query, err := render(meta, istioSuccessRateQuery) + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + promql := r.URL.Query()["query"][0] + if promql != expected { + t.Errorf("\nGot %s \nWanted %s", promql, expected) + } + + json := `{"status":"success","data":{"resultType":"vector","result":[{"metric":{},"value":[1,"100"]}]}}` + w.Write([]byte(json)) + })) + defer ts.Close() + + client, err := NewPrometheusClient(ts.URL, time.Second) if err != nil { t.Fatal(err) } - expected := `sum(rate(istio_requests_total{reporter="destination",destination_workload_namespace="default",destination_workload=~"podinfo",response_code!~"5.*"}[1m])) / sum(rate(istio_requests_total{reporter="destination",destination_workload_namespace="default",destination_workload=~"podinfo"}[1m])) * 100` + observer := &IstioObserver{ + client: client, + } - if query != expected { - t.Errorf("\nGot %s \nWanted %s", query, expected) + val, err := observer.GetRequestSuccessRate("podinfo", "default", "1m") + if err != nil { + t.Fatal(err.Error()) + } + + if val != 100 { + t.Errorf("Got %v wanted %v", val, 100) } } -func Test_IstioRequestDurationQueryRender(t *testing.T) { - meta := struct { - Name string - Namespace string - Interval string - }{ - "podinfo", - "default", - "1m", - } +func TestIstioObserver_GetRequestDuration(t *testing.T) { + expected := `histogram_quantile(0.99,sum(rate(istio_request_duration_seconds_bucket{reporter="destination",destination_workload_namespace="default",destination_workload=~"podinfo"}[1m]))by(le))` - query, err := render(meta, istioRequestDurationQuery) + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + promql := r.URL.Query()["query"][0] + if promql != expected { + t.Errorf("\nGot %s \nWanted %s", promql, expected) + } + + json := `{"status":"success","data":{"resultType":"vector","result":[{"metric":{},"value":[1,"100"]}]}}` + w.Write([]byte(json)) + })) + defer ts.Close() + + client, err := NewPrometheusClient(ts.URL, time.Second) if err != nil { t.Fatal(err) } - expected := `histogram_quantile(0.99, sum(rate(istio_request_duration_seconds_bucket{reporter="destination",destination_workload_namespace="default",destination_workload=~"podinfo"}[1m])) by (le))` + observer := &IstioObserver{ + client: client, + } - if query != expected { - t.Errorf("\nGot %s \nWanted %s", query, expected) + val, err := observer.GetRequestDuration("podinfo", "default", "1m") + if err != nil { + t.Fatal(err.Error()) + } + + if val != 100*time.Millisecond { + t.Errorf("Got %v wanted %v", val, 100*time.Millisecond) } } diff --git a/pkg/metrics/nginx.go b/pkg/metrics/nginx.go index bdd1c8f3..593d2306 100644 --- a/pkg/metrics/nginx.go +++ b/pkg/metrics/nginx.go @@ -1,122 +1,80 @@ package metrics import ( - "fmt" - "net/url" - "strconv" "time" ) -const nginxSuccessRateQuery = ` -sum(rate( -nginx_ingress_controller_requests{namespace="{{ .Namespace }}", -ingress="{{ .Name }}", -status!~"5.*"} -[{{ .Interval }}])) -/ -sum(rate( -nginx_ingress_controller_requests{namespace="{{ .Namespace }}", -ingress="{{ .Name }}"} -[{{ .Interval }}])) -* 100 -` - -// GetNginxSuccessRate returns the requests success rate (non 5xx) using nginx_ingress_controller_requests metric -func (c *Observer) GetNginxSuccessRate(name string, namespace string, metric string, interval string) (float64, error) { - if c.metricsServer == "fake" { - return 100, nil - } - - meta := struct { - Name string - Namespace string - Interval string - }{ - name, - namespace, - interval, - } - - query, err := render(meta, nginxSuccessRateQuery) - if err != nil { - return 0, err - } - - var rate *float64 - querySt := url.QueryEscape(query) - result, err := c.queryMetric(querySt) - if err != nil { - return 0, err - } - - for _, v := range result.Data.Result { - metricValue := v.Value[1] - switch metricValue.(type) { - case string: - f, err := strconv.ParseFloat(metricValue.(string), 64) - if err != nil { - return 0, err - } - rate = &f - } - } - if rate == nil { - return 0, fmt.Errorf("no values found for metric %s", metric) - } - return *rate, nil +var nginxQueries = map[string]string{ + "request-success-rate": ` + sum( + rate( + nginx_ingress_controller_requests{ + namespace="{{ .Namespace }}", + ingress="{{ .Name }}", + status!~"5.*" + }[{{ .Interval }}] + ) + ) + / + sum( + rate( + nginx_ingress_controller_requests{ + namespace="{{ .Namespace }}", + ingress="{{ .Name }}" + }[{{ .Interval }}] + ) + ) + * 100`, + "request-duration": ` + sum( + rate( + nginx_ingress_controller_ingress_upstream_latency_seconds_sum{ + namespace="{{ .Namespace }}", + ingress="{{ .Name }}" + }[{{ .Interval }}] + ) + ) + / + sum( + rate( + nginx_ingress_controller_ingress_upstream_latency_seconds_count{ + namespace="{{ .Namespace }}", + ingress="{{ .Name }}" + }[{{ .Interval }}] + ) + ) + * 1000`, } -const nginxRequestDurationQuery = ` -sum(rate( -nginx_ingress_controller_ingress_upstream_latency_seconds_sum{namespace="{{ .Namespace }}", -ingress="{{ .Name }}"}[{{ .Interval }}])) -/ -sum(rate(nginx_ingress_controller_ingress_upstream_latency_seconds_count{namespace="{{ .Namespace }}", -ingress="{{ .Name }}"}[{{ .Interval }}])) * 1000 -` +type NginxObserver struct { + client *PrometheusClient +} -// GetNginxRequestDuration returns the avg requests latency using nginx_ingress_controller_ingress_upstream_latency_seconds_sum metric -func (c *Observer) GetNginxRequestDuration(name string, namespace string, metric string, interval string) (time.Duration, error) { - if c.metricsServer == "fake" { - return 1, nil - } - - meta := struct { - Name string - Namespace string - Interval string - }{ - name, - namespace, - interval, - } - - query, err := render(meta, nginxRequestDurationQuery) +func (ob *NginxObserver) GetRequestSuccessRate(name string, namespace string, interval string) (float64, error) { + query, err := ob.client.RenderQuery(name, namespace, interval, nginxQueries["request-success-rate"]) if err != nil { return 0, err } - var rate *float64 - querySt := url.QueryEscape(query) - result, err := c.queryMetric(querySt) + value, err := ob.client.RunQuery(query) if err != nil { return 0, err } - for _, v := range result.Data.Result { - metricValue := v.Value[1] - switch metricValue.(type) { - case string: - f, err := strconv.ParseFloat(metricValue.(string), 64) - if err != nil { - return 0, err - } - rate = &f - } + return value, nil +} + +func (ob *NginxObserver) GetRequestDuration(name string, namespace string, interval string) (time.Duration, error) { + query, err := ob.client.RenderQuery(name, namespace, interval, nginxQueries["request-duration"]) + if err != nil { + return 0, err } - if rate == nil { - return 0, fmt.Errorf("no values found for metric %s", metric) + + value, err := ob.client.RunQuery(query) + if err != nil { + return 0, err } - ms := time.Duration(int64(*rate)) * time.Millisecond + + ms := time.Duration(int64(value)) * time.Millisecond return ms, nil } diff --git a/pkg/metrics/nginx_test.go b/pkg/metrics/nginx_test.go index 00e9ef90..30b51e31 100644 --- a/pkg/metrics/nginx_test.go +++ b/pkg/metrics/nginx_test.go @@ -1,51 +1,74 @@ package metrics import ( + "net/http" + "net/http/httptest" "testing" + "time" ) -func Test_NginxSuccessRateQueryRender(t *testing.T) { - meta := struct { - Name string - Namespace string - Interval string - }{ - "podinfo", - "nginx", - "1m", - } +func TestNginxObserver_GetRequestSuccessRate(t *testing.T) { + expected := `sum(rate(nginx_ingress_controller_requests{namespace="nginx",ingress="podinfo",status!~"5.*"}[1m]))/sum(rate(nginx_ingress_controller_requests{namespace="nginx",ingress="podinfo"}[1m]))*100` - query, err := render(meta, nginxSuccessRateQuery) + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + promql := r.URL.Query()["query"][0] + if promql != expected { + t.Errorf("\nGot %s \nWanted %s", promql, expected) + } + + json := `{"status":"success","data":{"resultType":"vector","result":[{"metric":{},"value":[1,"100"]}]}}` + w.Write([]byte(json)) + })) + defer ts.Close() + + client, err := NewPrometheusClient(ts.URL, time.Second) if err != nil { t.Fatal(err) } - expected := `sum(rate(nginx_ingress_controller_requests{namespace="nginx",ingress="podinfo",status!~"5.*"}[1m])) / sum(rate(nginx_ingress_controller_requests{namespace="nginx",ingress="podinfo"}[1m])) * 100` + observer := &NginxObserver{ + client: client, + } - if query != expected { - t.Errorf("\nGot %s \nWanted %s", query, expected) + val, err := observer.GetRequestSuccessRate("podinfo", "nginx", "1m") + if err != nil { + t.Fatal(err.Error()) + } + + if val != 100 { + t.Errorf("Got %v wanted %v", val, 100) } } -func Test_NginxRequestDurationQueryRender(t *testing.T) { - meta := struct { - Name string - Namespace string - Interval string - }{ - "podinfo", - "nginx", - "1m", - } +func TestNginxObserver_GetRequestDuration(t *testing.T) { + expected := `sum(rate(nginx_ingress_controller_ingress_upstream_latency_seconds_sum{namespace="nginx",ingress="podinfo"}[1m]))/sum(rate(nginx_ingress_controller_ingress_upstream_latency_seconds_count{namespace="nginx",ingress="podinfo"}[1m]))*1000` - query, err := render(meta, nginxRequestDurationQuery) + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + promql := r.URL.Query()["query"][0] + if promql != expected { + t.Errorf("\nGot %s \nWanted %s", promql, expected) + } + + json := `{"status":"success","data":{"resultType":"vector","result":[{"metric":{},"value":[1,"100"]}]}}` + w.Write([]byte(json)) + })) + defer ts.Close() + + client, err := NewPrometheusClient(ts.URL, time.Second) if err != nil { t.Fatal(err) } - expected := `sum(rate(nginx_ingress_controller_ingress_upstream_latency_seconds_sum{namespace="nginx",ingress="podinfo"}[1m])) /sum(rate(nginx_ingress_controller_ingress_upstream_latency_seconds_count{namespace="nginx",ingress="podinfo"}[1m])) * 1000` + observer := &NginxObserver{ + client: client, + } - if query != expected { - t.Errorf("\nGot %s \nWanted %s", query, expected) + val, err := observer.GetRequestDuration("podinfo", "nginx", "1m") + if err != nil { + t.Fatal(err.Error()) + } + + if val != 100*time.Millisecond { + t.Errorf("Got %v wanted %v", val, 100*time.Millisecond) } } diff --git a/pkg/metrics/observer.go b/pkg/metrics/observer.go index ef0a5cfd..906285e7 100644 --- a/pkg/metrics/observer.go +++ b/pkg/metrics/observer.go @@ -1,186 +1,10 @@ package metrics import ( - "bufio" - "bytes" - "context" - "encoding/json" - "fmt" - "io/ioutil" - "net/http" - "net/url" - "strconv" - "strings" - "text/template" "time" ) -// Observer is used to query Prometheus -type Observer struct { - metricsServer string -} - -type vectorQueryResponse struct { - Data struct { - Result []struct { - Metric struct { - Code string `json:"response_code"` - Name string `json:"destination_workload"` - } - Value []interface{} `json:"value"` - } - } -} - -// NewObserver creates a new observer -func NewObserver(metricsServer string) Observer { - return Observer{ - metricsServer: metricsServer, - } -} - -// GetMetricsServer returns the Prometheus URL -func (c *Observer) GetMetricsServer() string { - return c.metricsServer -} - -func (c *Observer) queryMetric(query string) (*vectorQueryResponse, error) { - promURL, err := url.Parse(c.metricsServer) - if err != nil { - return nil, err - } - - u, err := url.Parse(fmt.Sprintf("./api/v1/query?query=%s", query)) - if err != nil { - return nil, err - } - - u = promURL.ResolveReference(u) - - req, err := http.NewRequest("GET", u.String(), nil) - if err != nil { - return nil, err - } - - ctx, cancel := context.WithTimeout(req.Context(), 5*time.Second) - defer cancel() - - r, err := http.DefaultClient.Do(req.WithContext(ctx)) - if err != nil { - return nil, err - } - defer r.Body.Close() - - b, err := ioutil.ReadAll(r.Body) - if err != nil { - return nil, fmt.Errorf("error reading body: %s", err.Error()) - } - - if 400 <= r.StatusCode { - return nil, fmt.Errorf("error response: %s", string(b)) - } - - var values vectorQueryResponse - err = json.Unmarshal(b, &values) - if err != nil { - return nil, fmt.Errorf("error unmarshaling result: %s, '%s'", err.Error(), string(b)) - } - - return &values, nil -} - -// GetScalar runs the promql query and returns the first value found -func (c *Observer) GetScalar(query string) (float64, error) { - if c.metricsServer == "fake" { - return 100, nil - } - - query = strings.Replace(query, "\n", "", -1) - query = strings.Replace(query, " ", "", -1) - - var value *float64 - - querySt := url.QueryEscape(query) - result, err := c.queryMetric(querySt) - if err != nil { - return 0, err - } - - for _, v := range result.Data.Result { - metricValue := v.Value[1] - switch metricValue.(type) { - case string: - f, err := strconv.ParseFloat(metricValue.(string), 64) - if err != nil { - return 0, err - } - value = &f - } - } - if value == nil { - return 0, fmt.Errorf("no values found for query %s", query) - } - return *value, nil -} - -// CheckMetricsServer call Prometheus status endpoint and returns an error if -// the API is unreachable -func CheckMetricsServer(address string) (bool, error) { - promURL, err := url.Parse(address) - if err != nil { - return false, err - } - - u, err := url.Parse("./api/v1/status/flags") - if err != nil { - return false, err - } - - u = promURL.ResolveReference(u) - - req, err := http.NewRequest("GET", u.String(), nil) - if err != nil { - return false, err - } - - ctx, cancel := context.WithTimeout(req.Context(), 5*time.Second) - defer cancel() - - r, err := http.DefaultClient.Do(req.WithContext(ctx)) - if err != nil { - return false, err - } - defer r.Body.Close() - - b, err := ioutil.ReadAll(r.Body) - if err != nil { - return false, fmt.Errorf("error reading body: %s", err.Error()) - } - - if 400 <= r.StatusCode { - return false, fmt.Errorf("error response: %s", string(b)) - } - - return true, nil -} - -func render(meta interface{}, tmpl string) (string, error) { - t, err := template.New("tmpl").Parse(tmpl) - if err != nil { - return "", err - } - var data bytes.Buffer - b := bufio.NewWriter(&data) - - if err := t.Execute(b, meta); err != nil { - return "", err - } - err = b.Flush() - if err != nil { - return "", err - } - - res := strings.ReplaceAll(data.String(), "\n", "") - - return res, nil +type Interface interface { + GetRequestSuccessRate(name string, namespace string, interval string) (float64, error) + GetRequestDuration(name string, namespace string, interval string) (time.Duration, error) } diff --git a/pkg/metrics/observer_test.go b/pkg/metrics/observer_test.go deleted file mode 100644 index fcd12eb4..00000000 --- a/pkg/metrics/observer_test.go +++ /dev/null @@ -1,119 +0,0 @@ -package metrics - -import ( - "net/http" - "net/http/httptest" - "testing" - "time" -) - -func TestCanaryObserver_GetEnvoySuccessRate(t *testing.T) { - ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - json := `{"status":"success","data":{"resultType":"vector","result":[{"metric":{},"value":[1545905245.458,"100"]}]}}` - w.Write([]byte(json)) - })) - defer ts.Close() - - observer := NewObserver(ts.URL) - - val, err := observer.GetEnvoySuccessRate("podinfo", "default", "envoy_cluster_upstream_rq", "1m") - if err != nil { - t.Fatal(err.Error()) - } - - if val != 100 { - t.Errorf("Got %v wanted %v", val, 100) - } - -} - -func TestCanaryObserver_GetEnvoyRequestDuration(t *testing.T) { - ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - json := `{"status":"success","data":{"resultType":"vector","result":[{"metric":{},"value":[1545905245.596,"200"]}]}}` - w.Write([]byte(json)) - })) - defer ts.Close() - - observer := NewObserver(ts.URL) - - val, err := observer.GetEnvoyRequestDuration("podinfo", "default", "envoy_cluster_upstream_rq_time_bucket", "1m") - if err != nil { - t.Fatal(err.Error()) - } - - if val != 200*time.Millisecond { - t.Errorf("Got %v wanted %v", val, 200*time.Millisecond) - } -} - -func TestCanaryObserver_GetIstioSuccessRate(t *testing.T) { - ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - json := `{"status":"success","data":{"resultType":"vector","result":[{"metric":{},"value":[1545905245.458,"100"]}]}}` - w.Write([]byte(json)) - })) - defer ts.Close() - - observer := NewObserver(ts.URL) - - val, err := observer.GetIstioSuccessRate("podinfo", "default", "istio_requests_total", "1m") - if err != nil { - t.Fatal(err.Error()) - } - - if val != 100 { - t.Errorf("Got %v wanted %v", val, 100) - } - -} - -func TestCanaryObserver_GetIstioRequestDuration(t *testing.T) { - ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - json := `{"status":"success","data":{"resultType":"vector","result":[{"metric":{},"value":[1545905245.596,"0.2"]}]}}` - w.Write([]byte(json)) - })) - defer ts.Close() - - observer := NewObserver(ts.URL) - - val, err := observer.GetIstioRequestDuration("podinfo", "default", "istio_request_duration_seconds_bucket", "1m") - if err != nil { - t.Fatal(err.Error()) - } - - if val != 200*time.Millisecond { - t.Errorf("Got %v wanted %v", val, 200*time.Millisecond) - } -} - -func TestCheckMetricsServer(t *testing.T) { - ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - json := `{"status":"success","data":{"config.file":"/etc/prometheus/prometheus.yml"}}` - w.Write([]byte(json)) - })) - defer ts.Close() - - ok, err := CheckMetricsServer(ts.URL) - if err != nil { - t.Fatal(err.Error()) - } - - if !ok { - t.Errorf("Got %v wanted %v", ok, true) - } -} - -func TestCheckMetricsServer_Offline(t *testing.T) { - ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - w.WriteHeader(http.StatusBadGateway) - })) - defer ts.Close() - - ok, err := CheckMetricsServer(ts.URL) - if err == nil { - t.Errorf("Got no error wanted %v", http.StatusBadGateway) - } - - if ok { - t.Errorf("Got %v wanted %v", ok, false) - } -} From ee500d83aca1376ba2eec731d83a16fee170add0 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Mon, 13 May 2019 17:51:39 +0300 Subject: [PATCH 12/14] Add Linkerd observer implementation --- pkg/metrics/factory.go | 4 ++ pkg/metrics/linkerd.go | 73 ++++++++++++++++++++++++++++++++++++ pkg/metrics/linkerd_test.go | 74 +++++++++++++++++++++++++++++++++++++ 3 files changed, 151 insertions(+) create mode 100644 pkg/metrics/linkerd.go create mode 100644 pkg/metrics/linkerd_test.go diff --git a/pkg/metrics/factory.go b/pkg/metrics/factory.go index c717b3a6..b351a44d 100644 --- a/pkg/metrics/factory.go +++ b/pkg/metrics/factory.go @@ -31,6 +31,10 @@ func (factory Factory) Observer() Interface { return &NginxObserver{ client: factory.Client, } + case factory.MeshProvider == "smi:linkerd": + return &LinkerdObserver{ + client: factory.Client, + } default: return &IstioObserver{ client: factory.Client, diff --git a/pkg/metrics/linkerd.go b/pkg/metrics/linkerd.go new file mode 100644 index 00000000..50bcd00f --- /dev/null +++ b/pkg/metrics/linkerd.go @@ -0,0 +1,73 @@ +package metrics + +import ( + "time" +) + +var linkerdQueries = map[string]string{ + "request-success-rate": ` + sum( + rate( + response_total{ + namespace="{{ .Namespace }}", + dst_deployment=~"{{ .Name }}", + classification="failure" + }[{{ .Interval }}] + ) + ) + / + sum( + rate( + response_total{ + namespace="{{ .Namespace }}", + dst_deployment=~"{{ .Name }}" + }[{{ .Interval }}] + ) + ) + * 100`, + "request-duration": ` + histogram_quantile( + 0.99, + sum( + rate( + response_latency_ms_bucket{ + namespace="{{ .Namespace }}", + dst_deployment=~"{{ .Name }}" + }[{{ .Interval }}] + ) + ) by (le) + )`, +} + +type LinkerdObserver struct { + client *PrometheusClient +} + +func (ob *LinkerdObserver) GetRequestSuccessRate(name string, namespace string, interval string) (float64, error) { + query, err := ob.client.RenderQuery(name, namespace, interval, linkerdQueries["request-success-rate"]) + if err != nil { + return 0, err + } + + value, err := ob.client.RunQuery(query) + if err != nil { + return 0, err + } + + return value, nil +} + +func (ob *LinkerdObserver) GetRequestDuration(name string, namespace string, interval string) (time.Duration, error) { + query, err := ob.client.RenderQuery(name, namespace, interval, linkerdQueries["request-duration"]) + if err != nil { + return 0, err + } + + value, err := ob.client.RunQuery(query) + if err != nil { + return 0, err + } + + ms := time.Duration(int64(value)) * time.Millisecond + return ms, nil +} diff --git a/pkg/metrics/linkerd_test.go b/pkg/metrics/linkerd_test.go new file mode 100644 index 00000000..502109a1 --- /dev/null +++ b/pkg/metrics/linkerd_test.go @@ -0,0 +1,74 @@ +package metrics + +import ( + "net/http" + "net/http/httptest" + "testing" + "time" +) + +func TestLinkerdObserver_GetRequestSuccessRate(t *testing.T) { + expected := `sum(rate(response_total{namespace="default",dst_deployment=~"podinfo",classification="failure"}[1m]))/sum(rate(response_total{namespace="default",dst_deployment=~"podinfo"}[1m]))*100` + + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + promql := r.URL.Query()["query"][0] + if promql != expected { + t.Errorf("\nGot %s \nWanted %s", promql, expected) + } + + json := `{"status":"success","data":{"resultType":"vector","result":[{"metric":{},"value":[1,"100"]}]}}` + w.Write([]byte(json)) + })) + defer ts.Close() + + client, err := NewPrometheusClient(ts.URL, time.Second) + if err != nil { + t.Fatal(err) + } + + observer := &LinkerdObserver{ + client: client, + } + + val, err := observer.GetRequestSuccessRate("podinfo", "default", "1m") + if err != nil { + t.Fatal(err.Error()) + } + + if val != 100 { + t.Errorf("Got %v wanted %v", val, 100) + } +} + +func TestLinkerdObserver_GetRequestDuration(t *testing.T) { + expected := `histogram_quantile(0.99,sum(rate(response_latency_ms_bucket{namespace="default",dst_deployment=~"podinfo"}[1m]))by(le))` + + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + promql := r.URL.Query()["query"][0] + if promql != expected { + t.Errorf("\nGot %s \nWanted %s", promql, expected) + } + + json := `{"status":"success","data":{"resultType":"vector","result":[{"metric":{},"value":[1,"100"]}]}}` + w.Write([]byte(json)) + })) + defer ts.Close() + + client, err := NewPrometheusClient(ts.URL, time.Second) + if err != nil { + t.Fatal(err) + } + + observer := &LinkerdObserver{ + client: client, + } + + val, err := observer.GetRequestDuration("podinfo", "default", "1m") + if err != nil { + t.Fatal(err.Error()) + } + + if val != 100*time.Millisecond { + t.Errorf("Got %v wanted %v", val, 100*time.Millisecond) + } +} From 674c79da9425f0b63b4bd05af71f1931a4488ddf Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Tue, 14 May 2019 12:14:47 +0300 Subject: [PATCH 13/14] Fix Linkerd promql queries - include all inbound traffic stats --- pkg/metrics/linkerd.go | 11 +++++++---- pkg/metrics/linkerd_test.go | 4 ++-- 2 files changed, 9 insertions(+), 6 deletions(-) diff --git a/pkg/metrics/linkerd.go b/pkg/metrics/linkerd.go index 50bcd00f..4a9ee294 100644 --- a/pkg/metrics/linkerd.go +++ b/pkg/metrics/linkerd.go @@ -10,8 +10,9 @@ var linkerdQueries = map[string]string{ rate( response_total{ namespace="{{ .Namespace }}", - dst_deployment=~"{{ .Name }}", - classification="failure" + deployment=~"{{ .Name }}", + classification="failure", + direction="inbound" }[{{ .Interval }}] ) ) @@ -20,7 +21,8 @@ var linkerdQueries = map[string]string{ rate( response_total{ namespace="{{ .Namespace }}", - dst_deployment=~"{{ .Name }}" + deployment=~"{{ .Name }}", + direction="inbound" }[{{ .Interval }}] ) ) @@ -32,7 +34,8 @@ var linkerdQueries = map[string]string{ rate( response_latency_ms_bucket{ namespace="{{ .Namespace }}", - dst_deployment=~"{{ .Name }}" + deployment=~"{{ .Name }}", + direction="inbound" }[{{ .Interval }}] ) ) by (le) diff --git a/pkg/metrics/linkerd_test.go b/pkg/metrics/linkerd_test.go index 502109a1..6dbfef5c 100644 --- a/pkg/metrics/linkerd_test.go +++ b/pkg/metrics/linkerd_test.go @@ -8,7 +8,7 @@ import ( ) func TestLinkerdObserver_GetRequestSuccessRate(t *testing.T) { - expected := `sum(rate(response_total{namespace="default",dst_deployment=~"podinfo",classification="failure"}[1m]))/sum(rate(response_total{namespace="default",dst_deployment=~"podinfo"}[1m]))*100` + expected := `sum(rate(response_total{namespace="default",deployment=~"podinfo",classification="failure",direction="inbound"}[1m]))/sum(rate(response_total{namespace="default",deployment=~"podinfo",direction="inbound"}[1m]))*100` ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { promql := r.URL.Query()["query"][0] @@ -41,7 +41,7 @@ func TestLinkerdObserver_GetRequestSuccessRate(t *testing.T) { } func TestLinkerdObserver_GetRequestDuration(t *testing.T) { - expected := `histogram_quantile(0.99,sum(rate(response_latency_ms_bucket{namespace="default",dst_deployment=~"podinfo"}[1m]))by(le))` + expected := `histogram_quantile(0.99,sum(rate(response_latency_ms_bucket{namespace="default",deployment=~"podinfo",direction="inbound"}[1m]))by(le))` ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { promql := r.URL.Query()["query"][0] From 5a490abfdd9294d856e2f412c1cae9f47deb1bdc Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Tue, 14 May 2019 13:06:52 +0300 Subject: [PATCH 14/14] Remove the mesh gateway from docs examples --- README.md | 10 ++-------- docs/gitbook/how-it-works.md | 1 - docs/gitbook/usage/ab-testing.md | 1 - docs/gitbook/usage/progressive-delivery.md | 1 - 4 files changed, 2 insertions(+), 11 deletions(-) diff --git a/README.md b/README.md index 60e382c8..9e20ec4a 100644 --- a/README.md +++ b/README.md @@ -82,7 +82,6 @@ spec: # Istio gateways (optional) gateways: - public-gateway.istio-system.svc.cluster.local - - mesh # Istio virtual service host names (optional) hosts: - podinfo.example.com @@ -93,17 +92,12 @@ spec: # HTTP rewrite (optional) rewrite: uri: / - # Envoy timeout and retry policy (optional) - headers: - request: - add: - x-envoy-upstream-rq-timeout-ms: "15000" - x-envoy-max-retries: "10" - x-envoy-retry-on: "gateway-error,connect-failure,refused-stream" # cross-origin resource sharing policy (optional) corsPolicy: allowOrigin: - example.com + # request timeout (optional) + timeout: 5s # promote the canary without analysing it (default false) skipAnalysis: false # define the canary analysis timing and KPIs diff --git a/docs/gitbook/how-it-works.md b/docs/gitbook/how-it-works.md index d5c9160f..3a08f8eb 100644 --- a/docs/gitbook/how-it-works.md +++ b/docs/gitbook/how-it-works.md @@ -38,7 +38,6 @@ spec: # Istio gateways (optional) gateways: - public-gateway.istio-system.svc.cluster.local - - mesh # Istio virtual service host names (optional) hosts: - podinfo.example.com diff --git a/docs/gitbook/usage/ab-testing.md b/docs/gitbook/usage/ab-testing.md index 3014a794..32bdc9e3 100644 --- a/docs/gitbook/usage/ab-testing.md +++ b/docs/gitbook/usage/ab-testing.md @@ -60,7 +60,6 @@ spec: # Istio gateways (optional) gateways: - public-gateway.istio-system.svc.cluster.local - - mesh # Istio virtual service host names (optional) hosts: - app.example.com diff --git a/docs/gitbook/usage/progressive-delivery.md b/docs/gitbook/usage/progressive-delivery.md index 2bb25dcf..0c693657 100644 --- a/docs/gitbook/usage/progressive-delivery.md +++ b/docs/gitbook/usage/progressive-delivery.md @@ -54,7 +54,6 @@ spec: # Istio gateways (optional) gateways: - public-gateway.istio-system.svc.cluster.local - - mesh # Istio virtual service host names (optional) hosts: - app.example.com