From b2c12c11316b3ea6a2de0de206d0f3fc63ae490c Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sat, 30 Mar 2019 11:45:39 +0200 Subject: [PATCH 01/14] Move observer to metrics package --- cmd/flagger/main.go | 3 ++- pkg/controller/controller.go | 8 ++------ pkg/controller/controller_test.go | 7 ++----- pkg/controller/scheduler.go | 8 ++++---- pkg/{controller => metrics}/observer.go | 12 +++++++++++- pkg/{controller => metrics}/observer_test.go | 2 +- 6 files changed, 22 insertions(+), 18 deletions(-) rename pkg/{controller => metrics}/observer.go (96%) rename pkg/{controller => metrics}/observer_test.go (99%) diff --git a/cmd/flagger/main.go b/cmd/flagger/main.go index 463e06db..ed2fa659 100644 --- a/cmd/flagger/main.go +++ b/cmd/flagger/main.go @@ -7,6 +7,7 @@ import ( informers "github.com/weaveworks/flagger/pkg/client/informers/externalversions" "github.com/weaveworks/flagger/pkg/controller" "github.com/weaveworks/flagger/pkg/logging" + "github.com/weaveworks/flagger/pkg/metrics" "github.com/weaveworks/flagger/pkg/notifier" "github.com/weaveworks/flagger/pkg/server" "github.com/weaveworks/flagger/pkg/signals" @@ -105,7 +106,7 @@ func main() { logger.Infof("Watching namespace %s", namespace) } - ok, err := controller.CheckMetricsServer(metricsServer) + ok, err := metrics.CheckMetricsServer(metricsServer) if ok { logger.Infof("Connected to metrics server %s", metricsServer) } else { diff --git a/pkg/controller/controller.go b/pkg/controller/controller.go index c887083a..28415540 100644 --- a/pkg/controller/controller.go +++ b/pkg/controller/controller.go @@ -42,7 +42,7 @@ type Controller struct { canaries *sync.Map jobs map[string]CanaryJob deployer CanaryDeployer - observer CanaryObserver + observer metrics.CanaryObserver recorder metrics.CanaryRecorder notifier *notifier.Slack meshProvider string @@ -81,10 +81,6 @@ func NewController( }, } - observer := CanaryObserver{ - metricsServer: metricServer, - } - recorder := metrics.NewCanaryRecorder(controllerAgentName, true) recorder.SetInfo(version, meshProvider) @@ -101,7 +97,7 @@ func NewController( jobs: map[string]CanaryJob{}, flaggerWindow: flaggerWindow, deployer: deployer, - observer: observer, + observer: metrics.NewObserver(metricServer), recorder: recorder, notifier: notifier, meshProvider: meshProvider, diff --git a/pkg/controller/controller_test.go b/pkg/controller/controller_test.go index d4f3e418..ade7a935 100644 --- a/pkg/controller/controller_test.go +++ b/pkg/controller/controller_test.go @@ -35,7 +35,7 @@ type Mocks struct { meshClient clientset.Interface flaggerClient clientset.Interface deployer CanaryDeployer - observer CanaryObserver + observer metrics.CanaryObserver ctrl *Controller logger *zap.SugaredLogger router router.Interface @@ -74,9 +74,6 @@ func SetupMocks(abtest bool) Mocks { flaggerClient: flaggerClient, }, } - observer := CanaryObserver{ - metricsServer: "fake", - } // init controller flaggerInformerFactory := informers.NewSharedInformerFactory(flaggerClient, noResyncPeriodFunc()) @@ -94,7 +91,7 @@ func SetupMocks(abtest bool) Mocks { canaries: new(sync.Map), flaggerWindow: time.Second, deployer: deployer, - observer: observer, + observer: metrics.NewObserver("fake"), recorder: metrics.NewCanaryRecorder(controllerAgentName, false), } ctrl.flaggerSynced = alwaysReady diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index ac058b1c..892fc41f 100644 --- a/pkg/controller/scheduler.go +++ b/pkg/controller/scheduler.go @@ -501,7 +501,7 @@ func (c *Controller) analyseCanary(r *flaggerv1.Canary) bool { 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.metricsServer, err) + c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.observer.GetMetricsServer(), err) } return false } @@ -519,7 +519,7 @@ func (c *Controller) analyseCanary(r *flaggerv1.Canary) bool { 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.metricsServer, err) + c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.observer.GetMetricsServer(), err) } return false } @@ -533,7 +533,7 @@ func (c *Controller) analyseCanary(r *flaggerv1.Canary) bool { if metric.Name == "istio_request_duration_seconds_bucket" { val, err := c.observer.GetDeploymentHistogram(r.Spec.TargetRef.Name, r.Namespace, metric.Name, metric.Interval) if err != nil { - c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.observer.metricsServer, err) + c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.observer.GetMetricsServer(), err) return false } t := time.Duration(metric.Threshold) * time.Millisecond @@ -551,7 +551,7 @@ func (c *Controller) analyseCanary(r *flaggerv1.Canary) bool { 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.metricsServer, err) + c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.observer.GetMetricsServer(), err) } return false } diff --git a/pkg/controller/observer.go b/pkg/metrics/observer.go similarity index 96% rename from pkg/controller/observer.go rename to pkg/metrics/observer.go index 19eabff0..12dc24e3 100644 --- a/pkg/controller/observer.go +++ b/pkg/metrics/observer.go @@ -1,4 +1,4 @@ -package controller +package metrics import ( "context" @@ -29,6 +29,16 @@ type vectorQueryResponse struct { } } +func NewObserver(metricsServer string) CanaryObserver { + return CanaryObserver{ + metricsServer: metricsServer, + } +} + +func (c *CanaryObserver) GetMetricsServer() string { + return c.metricsServer +} + func (c *CanaryObserver) queryMetric(query string) (*vectorQueryResponse, error) { promURL, err := url.Parse(c.metricsServer) if err != nil { diff --git a/pkg/controller/observer_test.go b/pkg/metrics/observer_test.go similarity index 99% rename from pkg/controller/observer_test.go rename to pkg/metrics/observer_test.go index 1ef00982..5d4073fe 100644 --- a/pkg/controller/observer_test.go +++ b/pkg/metrics/observer_test.go @@ -1,4 +1,4 @@ -package controller +package metrics import ( "net/http" From 6a080f303212ce13dc17424800a17860f6f7ea5a Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sat, 30 Mar 2019 11:49:43 +0200 Subject: [PATCH 02/14] Rename observer and recorder --- pkg/controller/controller.go | 6 +++--- pkg/controller/controller_test.go | 4 ++-- pkg/metrics/observer.go | 20 ++++++++++---------- pkg/metrics/observer_test.go | 4 ++-- pkg/metrics/recorder.go | 20 ++++++++++---------- 5 files changed, 27 insertions(+), 27 deletions(-) diff --git a/pkg/controller/controller.go b/pkg/controller/controller.go index 28415540..6fbc2b91 100644 --- a/pkg/controller/controller.go +++ b/pkg/controller/controller.go @@ -42,8 +42,8 @@ type Controller struct { canaries *sync.Map jobs map[string]CanaryJob deployer CanaryDeployer - observer metrics.CanaryObserver - recorder metrics.CanaryRecorder + observer metrics.Observer + recorder metrics.Recorder notifier *notifier.Slack meshProvider string } @@ -81,7 +81,7 @@ func NewController( }, } - recorder := metrics.NewCanaryRecorder(controllerAgentName, true) + recorder := metrics.NewRecorder(controllerAgentName, true) recorder.SetInfo(version, meshProvider) ctrl := &Controller{ diff --git a/pkg/controller/controller_test.go b/pkg/controller/controller_test.go index ade7a935..a636f7f8 100644 --- a/pkg/controller/controller_test.go +++ b/pkg/controller/controller_test.go @@ -35,7 +35,7 @@ type Mocks struct { meshClient clientset.Interface flaggerClient clientset.Interface deployer CanaryDeployer - observer metrics.CanaryObserver + observer metrics.Observer ctrl *Controller logger *zap.SugaredLogger router router.Interface @@ -92,7 +92,7 @@ func SetupMocks(abtest bool) Mocks { flaggerWindow: time.Second, deployer: deployer, observer: metrics.NewObserver("fake"), - recorder: metrics.NewCanaryRecorder(controllerAgentName, false), + recorder: metrics.NewRecorder(controllerAgentName, false), } ctrl.flaggerSynced = alwaysReady diff --git a/pkg/metrics/observer.go b/pkg/metrics/observer.go index 12dc24e3..7eaacd4d 100644 --- a/pkg/metrics/observer.go +++ b/pkg/metrics/observer.go @@ -12,8 +12,8 @@ import ( "time" ) -// CanaryObserver is used to query the Istio Prometheus db -type CanaryObserver struct { +// Observer is used to query Prometheus +type Observer struct { metricsServer string } @@ -29,17 +29,17 @@ type vectorQueryResponse struct { } } -func NewObserver(metricsServer string) CanaryObserver { - return CanaryObserver{ +func NewObserver(metricsServer string) Observer { + return Observer{ metricsServer: metricsServer, } } -func (c *CanaryObserver) GetMetricsServer() string { +func (c *Observer) GetMetricsServer() string { return c.metricsServer } -func (c *CanaryObserver) queryMetric(query string) (*vectorQueryResponse, error) { +func (c *Observer) queryMetric(query string) (*vectorQueryResponse, error) { promURL, err := url.Parse(c.metricsServer) if err != nil { return nil, err @@ -85,7 +85,7 @@ func (c *CanaryObserver) queryMetric(query string) (*vectorQueryResponse, error) } // GetScalar runs the promql query and returns the first value found -func (c *CanaryObserver) GetScalar(query string) (float64, error) { +func (c *Observer) GetScalar(query string) (float64, error) { if c.metricsServer == "fake" { return 100, nil } @@ -116,7 +116,7 @@ func (c *CanaryObserver) GetScalar(query string) (float64, error) { return *value, nil } -func (c *CanaryObserver) GetEnvoySuccessRate(name string, namespace string, metric string, interval string) (float64, error) { +func (c *Observer) GetEnvoySuccessRate(name string, namespace string, metric string, interval string) (float64, error) { if c.metricsServer == "fake" { return 100, nil } @@ -154,7 +154,7 @@ func (c *CanaryObserver) GetEnvoySuccessRate(name string, namespace string, metr } // GetDeploymentCounter returns the requests success rate using istio_requests_total metric -func (c *CanaryObserver) GetDeploymentCounter(name string, namespace string, metric string, interval string) (float64, error) { +func (c *Observer) GetDeploymentCounter(name string, namespace string, metric string, interval string) (float64, error) { if c.metricsServer == "fake" { return 100, nil } @@ -192,7 +192,7 @@ func (c *CanaryObserver) GetDeploymentCounter(name string, namespace string, met } // GetDeploymentHistogram returns the 99P requests delay using istio_request_duration_seconds_bucket metrics -func (c *CanaryObserver) GetDeploymentHistogram(name string, namespace string, metric string, interval string) (time.Duration, error) { +func (c *Observer) GetDeploymentHistogram(name string, namespace string, metric string, interval string) (time.Duration, error) { if c.metricsServer == "fake" { return 1, nil } diff --git a/pkg/metrics/observer_test.go b/pkg/metrics/observer_test.go index 5d4073fe..ac961a31 100644 --- a/pkg/metrics/observer_test.go +++ b/pkg/metrics/observer_test.go @@ -14,7 +14,7 @@ func TestCanaryObserver_GetDeploymentCounter(t *testing.T) { })) defer ts.Close() - observer := CanaryObserver{ + observer := Observer{ metricsServer: ts.URL, } @@ -36,7 +36,7 @@ func TestCanaryObserver_GetDeploymentHistogram(t *testing.T) { })) defer ts.Close() - observer := CanaryObserver{ + observer := Observer{ metricsServer: ts.URL, } diff --git a/pkg/metrics/recorder.go b/pkg/metrics/recorder.go index 699db590..6b8d50ed 100644 --- a/pkg/metrics/recorder.go +++ b/pkg/metrics/recorder.go @@ -8,8 +8,8 @@ import ( flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1alpha3" ) -// CanaryRecorder records the canary analysis as Prometheus metrics -type CanaryRecorder struct { +// Recorder records the canary analysis as Prometheus metrics +type Recorder struct { info *prometheus.GaugeVec duration *prometheus.HistogramVec total *prometheus.GaugeVec @@ -17,8 +17,8 @@ type CanaryRecorder struct { weight *prometheus.GaugeVec } -// NewCanaryRecorder creates a new recorder and registers the Prometheus metrics -func NewCanaryRecorder(controller string, register bool) CanaryRecorder { +// NewRecorder creates a new recorder and registers the Prometheus metrics +func NewRecorder(controller string, register bool) Recorder { info := prometheus.NewGaugeVec(prometheus.GaugeOpts{ Subsystem: controller, Name: "info", @@ -59,7 +59,7 @@ func NewCanaryRecorder(controller string, register bool) CanaryRecorder { prometheus.MustRegister(weight) } - return CanaryRecorder{ + return Recorder{ info: info, duration: duration, total: total, @@ -69,22 +69,22 @@ func NewCanaryRecorder(controller string, register bool) CanaryRecorder { } // SetInfo sets the version and mesh provider labels -func (cr *CanaryRecorder) SetInfo(version string, meshProvider string) { +func (cr *Recorder) SetInfo(version string, meshProvider string) { cr.info.WithLabelValues(version, meshProvider).Set(1) } // SetDuration sets the time spent in seconds performing canary analysis -func (cr *CanaryRecorder) SetDuration(cd *flaggerv1.Canary, duration time.Duration) { +func (cr *Recorder) SetDuration(cd *flaggerv1.Canary, duration time.Duration) { cr.duration.WithLabelValues(cd.Spec.TargetRef.Name, cd.Namespace).Observe(duration.Seconds()) } // SetTotal sets the total number of canaries per namespace -func (cr *CanaryRecorder) SetTotal(namespace string, total int) { +func (cr *Recorder) SetTotal(namespace string, total int) { cr.total.WithLabelValues(namespace).Set(float64(total)) } // SetStatus sets the last known canary analysis status -func (cr *CanaryRecorder) SetStatus(cd *flaggerv1.Canary, phase flaggerv1.CanaryPhase) { +func (cr *Recorder) SetStatus(cd *flaggerv1.Canary, phase flaggerv1.CanaryPhase) { status := 1 switch phase { case flaggerv1.CanaryProgressing: @@ -98,7 +98,7 @@ func (cr *CanaryRecorder) SetStatus(cd *flaggerv1.Canary, phase flaggerv1.Canary } // SetWeight sets the weight values for primary and canary destinations -func (cr *CanaryRecorder) SetWeight(cd *flaggerv1.Canary, primary int, canary int) { +func (cr *Recorder) SetWeight(cd *flaggerv1.Canary, primary int, canary int) { cr.weight.WithLabelValues(fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name), cd.Namespace).Set(float64(primary)) cr.weight.WithLabelValues(cd.Spec.TargetRef.Name, cd.Namespace).Set(float64(canary)) } From c91a128b65aed4c532c4b2782780cc4b9a84fb85 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sat, 30 Mar 2019 11:55:41 +0200 Subject: [PATCH 03/14] Fix observer mock init --- pkg/controller/controller_test.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pkg/controller/controller_test.go b/pkg/controller/controller_test.go index a636f7f8..dd605d7c 100644 --- a/pkg/controller/controller_test.go +++ b/pkg/controller/controller_test.go @@ -74,6 +74,7 @@ func SetupMocks(abtest bool) Mocks { flaggerClient: flaggerClient, }, } + observer := metrics.NewObserver("fake") // init controller flaggerInformerFactory := informers.NewSharedInformerFactory(flaggerClient, noResyncPeriodFunc()) @@ -91,7 +92,7 @@ func SetupMocks(abtest bool) Mocks { canaries: new(sync.Map), flaggerWindow: time.Second, deployer: deployer, - observer: metrics.NewObserver("fake"), + observer: observer, recorder: metrics.NewRecorder(controllerAgentName, false), } ctrl.flaggerSynced = alwaysReady From f211e0fe31716f7a88570171e4cc939f90bf6762 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sun, 31 Mar 2019 13:55:14 +0300 Subject: [PATCH 04/14] Use go templates to render the builtin promql queries --- pkg/controller/scheduler.go | 4 +- pkg/metrics/envoy.go | 65 +++++++++++++++++ pkg/metrics/envoy_test.go | 28 ++++++++ pkg/metrics/istio.go | 123 +++++++++++++++++++++++++++++++ pkg/metrics/istio_test.go | 51 +++++++++++++ pkg/metrics/observer.go | 135 +++++++---------------------------- pkg/metrics/observer_test.go | 4 +- 7 files changed, 297 insertions(+), 113 deletions(-) create mode 100644 pkg/metrics/envoy.go create mode 100644 pkg/metrics/envoy_test.go create mode 100644 pkg/metrics/istio.go create mode 100644 pkg/metrics/istio_test.go diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index 892fc41f..4c699f73 100644 --- a/pkg/controller/scheduler.go +++ b/pkg/controller/scheduler.go @@ -513,7 +513,7 @@ func (c *Controller) analyseCanary(r *flaggerv1.Canary) bool { } if metric.Name == "istio_requests_total" { - val, err := c.observer.GetDeploymentCounter(r.Spec.TargetRef.Name, r.Namespace, metric.Name, metric.Interval) + 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", @@ -531,7 +531,7 @@ func (c *Controller) analyseCanary(r *flaggerv1.Canary) bool { } if metric.Name == "istio_request_duration_seconds_bucket" { - val, err := c.observer.GetDeploymentHistogram(r.Spec.TargetRef.Name, r.Namespace, metric.Name, metric.Interval) + 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 diff --git a/pkg/metrics/envoy.go b/pkg/metrics/envoy.go new file mode 100644 index 00000000..04002bb7 --- /dev/null +++ b/pkg/metrics/envoy.go @@ -0,0 +1,65 @@ +package metrics + +import ( + "fmt" + "net/url" + "strconv" +) + +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 +} diff --git a/pkg/metrics/envoy_test.go b/pkg/metrics/envoy_test.go new file mode 100644 index 00000000..12d93570 --- /dev/null +++ b/pkg/metrics/envoy_test.go @@ -0,0 +1,28 @@ +package metrics + +import ( + "testing" +) + +func Test_EnvoySuccessRateQueryRender(t *testing.T) { + meta := struct { + Name string + Namespace string + Interval string + }{ + "podinfo", + "default", + "1m", + } + + query, err := render(meta, envoySuccessRateQuery) + 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` + + if query != expected { + t.Errorf("\nGot %s \nWanted %s", query, expected) + } +} diff --git a/pkg/metrics/istio.go b/pkg/metrics/istio.go new file mode 100644 index 00000000..5c8e9855 --- /dev/null +++ b/pkg/metrics/istio.go @@ -0,0 +1,123 @@ +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 +} + +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)) +` + +// 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) + 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) + } + ms := time.Duration(int64(*rate*1000)) * time.Millisecond + return ms, nil +} diff --git a/pkg/metrics/istio_test.go b/pkg/metrics/istio_test.go new file mode 100644 index 00000000..28826a77 --- /dev/null +++ b/pkg/metrics/istio_test.go @@ -0,0 +1,51 @@ +package metrics + +import ( + "testing" +) + +func Test_IstioSuccessRateQueryRender(t *testing.T) { + meta := struct { + Name string + Namespace string + Interval string + }{ + "podinfo", + "default", + "1m", + } + + query, err := render(meta, istioSuccessRateQuery) + 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` + + if query != expected { + t.Errorf("\nGot %s \nWanted %s", query, expected) + } +} + +func Test_IstioRequestDurationQueryRender(t *testing.T) { + meta := struct { + Name string + Namespace string + Interval string + }{ + "podinfo", + "default", + "1m", + } + + query, err := render(meta, istioRequestDurationQuery) + 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))` + + if query != expected { + t.Errorf("\nGot %s \nWanted %s", query, expected) + } +} diff --git a/pkg/metrics/observer.go b/pkg/metrics/observer.go index 7eaacd4d..ddb87119 100644 --- a/pkg/metrics/observer.go +++ b/pkg/metrics/observer.go @@ -1,6 +1,8 @@ package metrics import ( + "bufio" + "bytes" "context" "encoding/json" "fmt" @@ -9,6 +11,7 @@ import ( "net/url" "strconv" "strings" + "text/template" "time" ) @@ -29,12 +32,14 @@ type vectorQueryResponse struct { } } +// 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 } @@ -116,115 +121,6 @@ func (c *Observer) GetScalar(query string) (float64, error) { return *value, nil } -func (c *Observer) GetEnvoySuccessRate(name string, namespace string, metric string, interval string) (float64, error) { - if c.metricsServer == "fake" { - return 100, nil - } - - var rate *float64 - querySt := url.QueryEscape(`sum(rate(` + - metric + `{kubernetes_namespace="` + - namespace + `",kubernetes_pod_name=~"` + - name + `-[0-9a-zA-Z]+(-[0-9a-zA-Z]+)",envoy_response_code!~"5.*"}[` + - interval + `])) / sum(rate(` + - metric + `{kubernetes_namespace="` + - namespace + `",kubernetes_pod_name=~"` + - name + `-[0-9a-zA-Z]+(-[0-9a-zA-Z]+)"}[` + - interval + `])) * 100 `) - 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 -} - -// GetDeploymentCounter returns the requests success rate using istio_requests_total metric -func (c *Observer) GetDeploymentCounter(name string, namespace string, metric string, interval string) (float64, error) { - if c.metricsServer == "fake" { - return 100, nil - } - - var rate *float64 - querySt := url.QueryEscape(`sum(rate(` + - metric + `{reporter="destination",destination_workload_namespace=~"` + - namespace + `",destination_workload=~"` + - name + `",response_code!~"5.*"}[` + - interval + `])) / sum(rate(` + - metric + `{reporter="destination",destination_workload_namespace=~"` + - namespace + `",destination_workload=~"` + - name + `"}[` + - interval + `])) * 100 `) - 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 -} - -// GetDeploymentHistogram returns the 99P requests delay using istio_request_duration_seconds_bucket metrics -func (c *Observer) GetDeploymentHistogram(name string, namespace string, metric string, interval string) (time.Duration, error) { - if c.metricsServer == "fake" { - return 1, nil - } - var rate *float64 - querySt := url.QueryEscape(`histogram_quantile(0.99, sum(rate(` + - metric + `{reporter="destination",destination_workload=~"` + - name + `", destination_workload_namespace=~"` + - namespace + `"}[` + - interval + `])) by (le))`) - 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) - } - ms := time.Duration(int64(*rate*1000)) * time.Millisecond - return ms, nil -} - // CheckMetricsServer call Prometheus status endpoint and returns an error if // the API is unreachable func CheckMetricsServer(address string) (bool, error) { @@ -265,3 +161,24 @@ func CheckMetricsServer(address string) (bool, error) { 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 +} diff --git a/pkg/metrics/observer_test.go b/pkg/metrics/observer_test.go index ac961a31..1cd4cb87 100644 --- a/pkg/metrics/observer_test.go +++ b/pkg/metrics/observer_test.go @@ -18,7 +18,7 @@ func TestCanaryObserver_GetDeploymentCounter(t *testing.T) { metricsServer: ts.URL, } - val, err := observer.GetDeploymentCounter("podinfo", "default", "istio_requests_total", "1m") + val, err := observer.GetIstioSuccessRate("podinfo", "default", "istio_requests_total", "1m") if err != nil { t.Fatal(err.Error()) } @@ -40,7 +40,7 @@ func TestCanaryObserver_GetDeploymentHistogram(t *testing.T) { metricsServer: ts.URL, } - val, err := observer.GetDeploymentHistogram("podinfo", "default", "istio_request_duration_seconds_bucket", "1m") + val, err := observer.GetIstioRequestDuration("podinfo", "default", "istio_request_duration_seconds_bucket", "1m") if err != nil { t.Fatal(err.Error()) } From ec759ce467c16824cfbaa570c38e0f0e98a3f1b5 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sun, 31 Mar 2019 14:17:39 +0300 Subject: [PATCH 05/14] Add Envoy success rate test --- pkg/metrics/observer_test.go | 30 +++++++++++++++++++++++------- 1 file changed, 23 insertions(+), 7 deletions(-) diff --git a/pkg/metrics/observer_test.go b/pkg/metrics/observer_test.go index 1cd4cb87..2bcc0ab5 100644 --- a/pkg/metrics/observer_test.go +++ b/pkg/metrics/observer_test.go @@ -7,17 +7,35 @@ import ( "time" ) -func TestCanaryObserver_GetDeploymentCounter(t *testing.T) { +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 := Observer{ - metricsServer: ts.URL, + 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_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()) @@ -29,16 +47,14 @@ func TestCanaryObserver_GetDeploymentCounter(t *testing.T) { } -func TestCanaryObserver_GetDeploymentHistogram(t *testing.T) { +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 := Observer{ - metricsServer: ts.URL, - } + observer := NewObserver(ts.URL) val, err := observer.GetIstioRequestDuration("podinfo", "default", "istio_request_duration_seconds_bucket", "1m") if err != nil { From 347cfd06def9e65346320ba81a7af64cf80c8220 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sun, 31 Mar 2019 14:30:27 +0300 Subject: [PATCH 06/14] Upgrade Kubernetes Kind to v0.2.1 --- test/e2e-kind.sh | 12 +++++------- 1 file changed, 5 insertions(+), 7 deletions(-) diff --git a/test/e2e-kind.sh b/test/e2e-kind.sh index 2312aebe..e617c551 100755 --- a/test/e2e-kind.sh +++ b/test/e2e-kind.sh @@ -3,21 +3,19 @@ set -o errexit REPO_ROOT=$(git rev-parse --show-toplevel) +KIND_VERSION=0.2.1 echo ">>> Installing kubectl" curl -LO https://storage.googleapis.com/kubernetes-release/release/$(curl -s https://storage.googleapis.com/kubernetes-release/release/stable.txt)/bin/linux/amd64/kubectl && \ chmod +x kubectl && \ sudo mv kubectl /usr/local/bin/ -echo ">>> Building sigs.k8s.io/kind" -docker build -t kind:src . -f ${REPO_ROOT}/test/Dockerfile.kind -docker create -ti --name dummy kind:src sh -docker cp dummy:/go/bin/kind ./kind -docker rm -f dummy - echo ">>> Installing kind" +curl -sSLo kind "https://github.com/kubernetes-sigs/kind/releases/download/$KIND_VERSION/kind-linux-amd64" chmod +x kind -sudo mv kind /usr/local/bin/ +sudo mv kind /usr/local/bin/kind + +echo ">>> Creating kind cluster" kind create cluster --wait 5m export KUBECONFIG="$(kind get kubeconfig-path --name="kind")" From fc10745a1ad99f0bac5959e61f25fd05ce0cd845 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sun, 31 Mar 2019 14:43:08 +0300 Subject: [PATCH 07/14] Upgrade Istio e2e to v1.1.1 --- test/e2e-istio.sh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/test/e2e-istio.sh b/test/e2e-istio.sh index 71e2791d..f444673c 100755 --- a/test/e2e-istio.sh +++ b/test/e2e-istio.sh @@ -2,7 +2,7 @@ set -o errexit -ISTIO_VER="1.1.0-rc.0" +ISTIO_VER="1.1.1" REPO_ROOT=$(git rev-parse --show-toplevel) export KUBECONFIG="$(kind get kubeconfig-path --name="kind")" From 3a1018cff6475ac567003c60e50e354d37352953 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sun, 31 Mar 2019 18:00:03 +0300 Subject: [PATCH 08/14] Change Istio e2e limits --- test/e2e-istio-values.yaml | 8 ++------ 1 file changed, 2 insertions(+), 6 deletions(-) diff --git a/test/e2e-istio-values.yaml b/test/e2e-istio-values.yaml index 131065ba..97076f40 100644 --- a/test/e2e-istio-values.yaml +++ b/test/e2e-istio-values.yaml @@ -12,10 +12,6 @@ gateways: istio-ingressgateway: autoscaleMax: 1 -# citadel configuration -security: - enabled: true - # sidecar-injector webhook configuration sidecarInjectorWebhook: enabled: true @@ -55,8 +51,8 @@ global: resources: requests: cpu: 100m - memory: 128Mi + memory: 64Mi limits: cpu: 2000m - memory: 128Mi + memory: 256Mi useMCP: false From dbf36082b21adc9608317075a905af050f7b5a8a Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Mon, 1 Apr 2019 12:55:00 +0300 Subject: [PATCH 09/14] Revert Istio e2e to 1.1.0 --- test/e2e-istio.sh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/test/e2e-istio.sh b/test/e2e-istio.sh index f444673c..71e2791d 100755 --- a/test/e2e-istio.sh +++ b/test/e2e-istio.sh @@ -2,7 +2,7 @@ set -o errexit -ISTIO_VER="1.1.1" +ISTIO_VER="1.1.0-rc.0" REPO_ROOT=$(git rev-parse --show-toplevel) export KUBECONFIG="$(kind get kubeconfig-path --name="kind")" From 36dfd4dd3501866f4a6ccd43a58310c2420dd002 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Mon, 1 Apr 2019 13:09:52 +0300 Subject: [PATCH 10/14] Bring down Istio Pilot memory requests --- test/e2e-istio-values.yaml | 4 ++++ test/e2e-istio.sh | 2 +- 2 files changed, 5 insertions(+), 1 deletion(-) diff --git a/test/e2e-istio-values.yaml b/test/e2e-istio-values.yaml index 97076f40..68350462 100644 --- a/test/e2e-istio-values.yaml +++ b/test/e2e-istio-values.yaml @@ -6,6 +6,10 @@ pilot: enabled: true sidecar: true + resources: + requests: + cpu: 100m + memory: 128Mi gateways: enabled: false diff --git a/test/e2e-istio.sh b/test/e2e-istio.sh index 71e2791d..f444673c 100755 --- a/test/e2e-istio.sh +++ b/test/e2e-istio.sh @@ -2,7 +2,7 @@ set -o errexit -ISTIO_VER="1.1.0-rc.0" +ISTIO_VER="1.1.1" REPO_ROOT=$(git rev-parse --show-toplevel) export KUBECONFIG="$(kind get kubeconfig-path --name="kind")" From ee4a009a0619655b4bd6ec5aa2d45635e8dad7d2 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Mon, 1 Apr 2019 14:09:19 +0300 Subject: [PATCH 11/14] Print Istio e2e status --- test/e2e-istio.sh | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/test/e2e-istio.sh b/test/e2e-istio.sh index f444673c..918646fc 100755 --- a/test/e2e-istio.sh +++ b/test/e2e-istio.sh @@ -25,4 +25,13 @@ kubectl -n istio-system wait --for=condition=complete job/istio-init-crd-10 kubectl -n istio-system wait --for=condition=complete job/istio-init-crd-11 echo '>>> Installing Istio control plane' -helm upgrade -i istio istio.io/istio --wait --namespace istio-system -f ${REPO_ROOT}/test/e2e-istio-values.yaml \ No newline at end of file +helm upgrade -i istio istio.io/istio --namespace istio-system -f ${REPO_ROOT}/test/e2e-istio-values.yaml + +kubectl -n istio-system get all + +kubectl -n istio-system describe deployment/istio-pilot +kubectl -n istio-system describe deployment/istio-telemetry +kubectl -n istio-system describe deployment/istio-citadel +kubectl -n istio-system describe deployment/prometheus + +kubectl -n istio-system rollout status deployment/prometheus From 664e7ad555e06d70cf133724f21981e72d5472ef Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Mon, 1 Apr 2019 14:33:47 +0300 Subject: [PATCH 12/14] Debug e2e load test --- test/e2e-istio.sh | 4 ++-- test/e2e-tests.sh | 5 +++++ 2 files changed, 7 insertions(+), 2 deletions(-) diff --git a/test/e2e-istio.sh b/test/e2e-istio.sh index 918646fc..16b3bda0 100755 --- a/test/e2e-istio.sh +++ b/test/e2e-istio.sh @@ -27,11 +27,11 @@ kubectl -n istio-system wait --for=condition=complete job/istio-init-crd-11 echo '>>> Installing Istio control plane' helm upgrade -i istio istio.io/istio --namespace istio-system -f ${REPO_ROOT}/test/e2e-istio-values.yaml +kubectl -n istio-system rollout status deployment/prometheus + kubectl -n istio-system get all kubectl -n istio-system describe deployment/istio-pilot kubectl -n istio-system describe deployment/istio-telemetry kubectl -n istio-system describe deployment/istio-citadel kubectl -n istio-system describe deployment/prometheus - -kubectl -n istio-system rollout status deployment/prometheus diff --git a/test/e2e-tests.sh b/test/e2e-tests.sh index 3756a10f..fdbb97a8 100755 --- a/test/e2e-tests.sh +++ b/test/e2e-tests.sh @@ -14,6 +14,11 @@ kubectl label namespace test istio-injection=enabled echo '>>> Installing the load tester' kubectl -n test apply -f ${REPO_ROOT}/artifacts/loadtester/ + +sleep 30 + +kubectl -n test describe deployment/flagger-loadtester + kubectl -n test rollout status deployment/flagger-loadtester echo '>>> Initialising canary' From 69a6e260f5e6d8b823f49adb963918a927c85a8b Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Mon, 1 Apr 2019 14:40:11 +0300 Subject: [PATCH 13/14] Bring down the Istio e2e CPU requests --- test/e2e-istio-values.yaml | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/test/e2e-istio-values.yaml b/test/e2e-istio-values.yaml index 68350462..7386ca3a 100644 --- a/test/e2e-istio-values.yaml +++ b/test/e2e-istio-values.yaml @@ -8,7 +8,7 @@ pilot: sidecar: true resources: requests: - cpu: 100m + cpu: 10m memory: 128Mi gateways: @@ -36,7 +36,7 @@ mixer: autoscaleEnabled: false resources: requests: - cpu: 100m + cpu: 10m memory: 128Mi # addon prometheus configuration @@ -54,9 +54,9 @@ global: # Resources for the sidecar. resources: requests: - cpu: 100m + cpu: 10m memory: 64Mi limits: - cpu: 2000m + cpu: 1000m memory: 256Mi useMCP: false From 3e43963daaeab185fb4b8d6deb9959b044ae6244 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Mon, 1 Apr 2019 14:52:03 +0300 Subject: [PATCH 14/14] Wait for Istio pods to be ready --- test/e2e-istio.sh | 9 +-------- test/e2e-tests.sh | 5 ----- 2 files changed, 1 insertion(+), 13 deletions(-) diff --git a/test/e2e-istio.sh b/test/e2e-istio.sh index 16b3bda0..cf6bb992 100755 --- a/test/e2e-istio.sh +++ b/test/e2e-istio.sh @@ -25,13 +25,6 @@ kubectl -n istio-system wait --for=condition=complete job/istio-init-crd-10 kubectl -n istio-system wait --for=condition=complete job/istio-init-crd-11 echo '>>> Installing Istio control plane' -helm upgrade -i istio istio.io/istio --namespace istio-system -f ${REPO_ROOT}/test/e2e-istio-values.yaml - -kubectl -n istio-system rollout status deployment/prometheus +helm upgrade -i istio istio.io/istio --wait --namespace istio-system -f ${REPO_ROOT}/test/e2e-istio-values.yaml kubectl -n istio-system get all - -kubectl -n istio-system describe deployment/istio-pilot -kubectl -n istio-system describe deployment/istio-telemetry -kubectl -n istio-system describe deployment/istio-citadel -kubectl -n istio-system describe deployment/prometheus diff --git a/test/e2e-tests.sh b/test/e2e-tests.sh index fdbb97a8..3756a10f 100755 --- a/test/e2e-tests.sh +++ b/test/e2e-tests.sh @@ -14,11 +14,6 @@ kubectl label namespace test istio-injection=enabled echo '>>> Installing the load tester' kubectl -n test apply -f ${REPO_ROOT}/artifacts/loadtester/ - -sleep 30 - -kubectl -n test describe deployment/flagger-loadtester - kubectl -n test rollout status deployment/flagger-loadtester echo '>>> Initialising canary'