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..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 CanaryObserver - recorder metrics.CanaryRecorder + observer metrics.Observer + recorder metrics.Recorder notifier *notifier.Slack meshProvider string } @@ -81,11 +81,7 @@ func NewController( }, } - observer := CanaryObserver{ - metricsServer: metricServer, - } - - recorder := metrics.NewCanaryRecorder(controllerAgentName, true) + recorder := metrics.NewRecorder(controllerAgentName, true) recorder.SetInfo(version, meshProvider) ctrl := &Controller{ @@ -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..dd605d7c 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.Observer ctrl *Controller logger *zap.SugaredLogger router router.Interface @@ -74,9 +74,7 @@ func SetupMocks(abtest bool) Mocks { flaggerClient: flaggerClient, }, } - observer := CanaryObserver{ - metricsServer: "fake", - } + observer := metrics.NewObserver("fake") // init controller flaggerInformerFactory := informers.NewSharedInformerFactory(flaggerClient, noResyncPeriodFunc()) @@ -95,7 +93,7 @@ func SetupMocks(abtest bool) Mocks { flaggerWindow: time.Second, deployer: deployer, observer: observer, - recorder: metrics.NewCanaryRecorder(controllerAgentName, false), + recorder: metrics.NewRecorder(controllerAgentName, false), } ctrl.flaggerSynced = alwaysReady diff --git a/pkg/controller/observer.go b/pkg/controller/observer.go deleted file mode 100644 index 19eabff0..00000000 --- a/pkg/controller/observer.go +++ /dev/null @@ -1,257 +0,0 @@ -package controller - -import ( - "context" - "encoding/json" - "fmt" - "io/ioutil" - "net/http" - "net/url" - "strconv" - "strings" - "time" -) - -// CanaryObserver is used to query the Istio Prometheus db -type CanaryObserver 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"` - } - } -} - -func (c *CanaryObserver) 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 *CanaryObserver) 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 - result, err := c.queryMetric(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 - } - value = &f - } - } - if value == nil { - return 0, fmt.Errorf("no values found for query %s", query) - } - return *value, nil -} - -func (c *CanaryObserver) 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 *CanaryObserver) 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 *CanaryObserver) 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) { - 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 -} diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index ac058b1c..4c699f73 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 } @@ -513,13 +513,13 @@ 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", 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 } @@ -531,9 +531,9 @@ 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.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/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 new file mode 100644 index 00000000..ddb87119 --- /dev/null +++ b/pkg/metrics/observer.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" +) + +// 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 + result, err := c.queryMetric(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 + } + 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 +} diff --git a/pkg/controller/observer_test.go b/pkg/metrics/observer_test.go similarity index 62% rename from pkg/controller/observer_test.go rename to pkg/metrics/observer_test.go index 1ef00982..2bcc0ab5 100644 --- a/pkg/controller/observer_test.go +++ b/pkg/metrics/observer_test.go @@ -1,4 +1,4 @@ -package controller +package metrics import ( "net/http" @@ -7,18 +7,16 @@ 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 := CanaryObserver{ - metricsServer: ts.URL, - } + observer := NewObserver(ts.URL) - val, err := observer.GetDeploymentCounter("podinfo", "default", "istio_requests_total", "1m") + val, err := observer.GetEnvoySuccessRate("podinfo", "default", "envoy_cluster_upstream_rq", "1m") if err != nil { t.Fatal(err.Error()) } @@ -29,18 +27,36 @@ func TestCanaryObserver_GetDeploymentCounter(t *testing.T) { } -func TestCanaryObserver_GetDeploymentHistogram(t *testing.T) { +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 := CanaryObserver{ - metricsServer: ts.URL, - } + observer := NewObserver(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()) } 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)) } diff --git a/test/e2e-istio-values.yaml b/test/e2e-istio-values.yaml index 131065ba..7386ca3a 100644 --- a/test/e2e-istio-values.yaml +++ b/test/e2e-istio-values.yaml @@ -6,16 +6,16 @@ pilot: enabled: true sidecar: true + resources: + requests: + cpu: 10m + memory: 128Mi gateways: enabled: false istio-ingressgateway: autoscaleMax: 1 -# citadel configuration -security: - enabled: true - # sidecar-injector webhook configuration sidecarInjectorWebhook: enabled: true @@ -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 - memory: 128Mi + cpu: 10m + memory: 64Mi limits: - cpu: 2000m - memory: 128Mi + cpu: 1000m + memory: 256Mi useMCP: false diff --git a/test/e2e-istio.sh b/test/e2e-istio.sh index 71e2791d..cf6bb992 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")" @@ -25,4 +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 --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 --wait --namespace istio-system -f ${REPO_ROOT}/test/e2e-istio-values.yaml + +kubectl -n istio-system get all 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")"