From 6661406b75c7de71971714151fd411e447c1469e Mon Sep 17 00:00:00 2001 From: Yusuke Kuoka Date: Sat, 30 Nov 2019 12:33:34 +0900 Subject: [PATCH] Metrics provider for deployments and services behind Envoy Assumes `envoy:smi` as the mesh provider name as I've successfully tested the progressive delivery for Envoy + Crossover with it. This enhances Flagger to translate it to the metrics provider name of `envoy` for deployment targets, or `envoy:service` for service targets. The `envoy` metrics provider is equivalent to `appmesh`, as both relies on the same set of standard metrics exposed by Envoy itself. The `envoy:service` is almost the same as the `envoy` provider, but removing the condition on pod name, as we only need to filter on the backing service name = envoy_cluster_name. We don't consider other Envoy xDS implementations that uses anything that is different to original servicen ames as `envoy_cluster_name`, for now. Ref #385 --- cmd/flagger/main.go | 2 +- pkg/controller/controller_test.go | 2 +- pkg/controller/scheduler.go | 17 ++++++- pkg/metrics/envoy_service.go | 73 ++++++++++++++++++++++++++++++ pkg/metrics/envoy_service_test.go | 74 +++++++++++++++++++++++++++++++ pkg/metrics/factory.go | 14 +++--- 6 files changed, 170 insertions(+), 12 deletions(-) create mode 100644 pkg/metrics/envoy_service.go create mode 100644 pkg/metrics/envoy_service_test.go diff --git a/cmd/flagger/main.go b/cmd/flagger/main.go index 8548de56..7d9df913 100644 --- a/cmd/flagger/main.go +++ b/cmd/flagger/main.go @@ -161,7 +161,7 @@ func main() { logger.Infof("Watching namespace %s", namespace) } - observerFactory, err := metrics.NewFactory(metricsServer, meshProvider, 5*time.Second) + observerFactory, err := metrics.NewFactory(metricsServer, 5*time.Second) if err != nil { logger.Fatalf("Error building prometheus client: %s", err.Error()) } diff --git a/pkg/controller/controller_test.go b/pkg/controller/controller_test.go index 8554c436..6945ae2a 100644 --- a/pkg/controller/controller_test.go +++ b/pkg/controller/controller_test.go @@ -73,7 +73,7 @@ func SetupMocks(c *flaggerv1.Canary) Mocks { rf := router.NewFactory(nil, kubeClient, flaggerClient, "annotationsPrefix", logger, flaggerClient) // init observer - observerFactory, _ := metrics.NewFactory("fake", "istio", 5*time.Second) + observerFactory, _ := metrics.NewFactory("fake", 5*time.Second) // init canary factory configTracker := canary.ConfigTracker{ diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index 25cdabb6..ec16b602 100644 --- a/pkg/controller/scheduler.go +++ b/pkg/controller/scheduler.go @@ -13,6 +13,10 @@ import ( "github.com/weaveworks/flagger/pkg/router" ) +const ( + MetricsProviderServiceSuffix = ":service" +) + // scheduleCanaries synchronises the canary map with the jobs map, // for new canaries new jobs are created and started // for the removed canaries the jobs are stopped and deleted @@ -747,10 +751,19 @@ func (c *Controller) analyseCanary(r *flaggerv1.Canary) bool { if r.Spec.Provider != "" { metricsProvider = r.Spec.Provider - // set the metrics server to Linkerd Prometheus when Linkerd is the default mesh provider + // set the metrics provider to Linkerd Prometheus when Linkerd is the default mesh provider if strings.Contains(c.meshProvider, "linkerd") { metricsProvider = "linkerd" } + + // set the metrics provider to Envoy Prometheus when Envoy is the default mesh provider + if strings.Contains(c.meshProvider, "envoy") { + metricsProvider = "envoy" + } + } + // set the metrics provider to query Prometheus for the canary Kubernetes service if the canary target is Service + if r.Spec.TargetRef.Kind == "Service" { + metricsProvider = metricsProvider + MetricsProviderServiceSuffix } // create observer based on the mesh provider @@ -761,7 +774,7 @@ func (c *Controller) analyseCanary(r *flaggerv1.Canary) bool { if r.Spec.MetricsServer != "" { metricsServer = r.Spec.MetricsServer var err error - observerFactory, err = metrics.NewFactory(metricsServer, metricsProvider, 5*time.Second) + observerFactory, err = metrics.NewFactory(metricsServer, 5*time.Second) if err != nil { c.recordEventErrorf(r, "Error building Prometheus client for %s %v", r.Spec.MetricsServer, err) return false diff --git a/pkg/metrics/envoy_service.go b/pkg/metrics/envoy_service.go new file mode 100644 index 00000000..5e9d26c3 --- /dev/null +++ b/pkg/metrics/envoy_service.go @@ -0,0 +1,73 @@ +package metrics + +import ( + "time" +) + +var envoyServiceQueries = map[string]string{ + "request-success-rate": ` + sum( + rate( + envoy_cluster_upstream_rq{ + kubernetes_namespace="{{ .Namespace }}", + envoy_cluster_name="{{ .Name }}-canary", + envoy_response_code!~"5.*" + }[{{ .Interval }}] + ) + ) + / + sum( + rate( + envoy_cluster_upstream_rq{ + kubernetes_namespace="{{ .Namespace }}", + envoy_cluster_name="{{ .Name }}-canary" + }[{{ .Interval }}] + ) + ) + * 100`, + "request-duration": ` + histogram_quantile( + 0.99, + sum( + rate( + envoy_cluster_upstream_rq_time_bucket{ + kubernetes_namespace="{{ .Namespace }}", + envoy_cluster_name="{{ .Name }}-canary" + }[{{ .Interval }}] + ) + ) by (le) + )`, +} + +type EnvoyServiceObserver struct { + client *PrometheusClient +} + +func (ob *EnvoyServiceObserver) GetRequestSuccessRate(name string, namespace string, interval string) (float64, error) { + query, err := ob.client.RenderQuery(name, namespace, interval, envoyServiceQueries["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 *EnvoyServiceObserver) GetRequestDuration(name string, namespace string, interval string) (time.Duration, error) { + query, err := ob.client.RenderQuery(name, namespace, interval, envoyServiceQueries["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/envoy_service_test.go b/pkg/metrics/envoy_service_test.go new file mode 100644 index 00000000..fca854bd --- /dev/null +++ b/pkg/metrics/envoy_service_test.go @@ -0,0 +1,74 @@ +package metrics + +import ( + "net/http" + "net/http/httptest" + "testing" + "time" +) + +func TestEnvoyServiceObserver_GetRequestSuccessRate(t *testing.T) { + expected := ` sum( rate( envoy_cluster_upstream_rq{ kubernetes_namespace="default", envoy_cluster_name="podinfo-canary", envoy_response_code!~"5.*" }[1m] ) ) / sum( rate( envoy_cluster_upstream_rq{ kubernetes_namespace="default", envoy_cluster_name="podinfo-canary" }[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 := &EnvoyServiceObserver{ + 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 TestEnvoyServiceObserver_GetRequestDuration(t *testing.T) { + expected := ` histogram_quantile( 0.99, sum( rate( envoy_cluster_upstream_rq_time_bucket{ kubernetes_namespace="default", envoy_cluster_name="podinfo-canary" }[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 := &EnvoyServiceObserver{ + 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) + } +} diff --git a/pkg/metrics/factory.go b/pkg/metrics/factory.go index e2b69b8d..ce6d85d0 100644 --- a/pkg/metrics/factory.go +++ b/pkg/metrics/factory.go @@ -6,19 +6,17 @@ import ( ) type Factory struct { - MeshProvider string - Client *PrometheusClient + Client *PrometheusClient } -func NewFactory(metricsServer string, meshProvider string, timeout time.Duration) (*Factory, error) { +func NewFactory(metricsServer string, timeout time.Duration) (*Factory, error) { client, err := NewPrometheusClient(metricsServer, timeout) if err != nil { return nil, err } return &Factory{ - MeshProvider: meshProvider, - Client: client, + Client: client, }, nil } @@ -32,7 +30,7 @@ func (factory Factory) Observer(provider string) Interface { return &HttpObserver{ client: factory.Client, } - case provider == "appmesh": + case provider == "appmesh", provider == "envoy": return &EnvoyObserver{ client: factory.Client, } @@ -44,8 +42,8 @@ func (factory Factory) Observer(provider string) Interface { return &GlooObserver{ client: factory.Client, } - case provider == "smi:linkerd": - return &LinkerdObserver{ + case provider == "appmesh:service", provider == "envoy:service": + return &EnvoyServiceObserver{ client: factory.Client, } case provider == "linkerd":