From 7c96e8b081d60d2f90441c9c37b308ee62bd5d36 Mon Sep 17 00:00:00 2001 From: Stefan Prodan Date: Fri, 28 Sep 2018 13:27:25 +0300 Subject: [PATCH] Add req duration metric to CRD spec --- pkg/apis/rollout/v1beta1/types.go | 3 +- .../rollout/v1beta1/zz_generated.deepcopy.go | 8 ++- pkg/controller/controller.go | 36 ------------- pkg/controller/deployment.go | 48 +++++++++++------ pkg/controller/observer.go | 51 ++++++++++++++++--- pkg/version/version.go | 2 +- 6 files changed, 84 insertions(+), 64 deletions(-) diff --git a/pkg/apis/rollout/v1beta1/types.go b/pkg/apis/rollout/v1beta1/types.go index 9901ce6b..2cef2f27 100755 --- a/pkg/apis/rollout/v1beta1/types.go +++ b/pkg/apis/rollout/v1beta1/types.go @@ -38,7 +38,7 @@ type RolloutSpec struct { Primary Target `json:"primary"` Canary Target `json:"canary"` VirtualService VirtualService `json:"virtualService"` - Metric Metric `json:"metric"` + Metrics []Metric `json:"metrics"` } type Target struct { @@ -52,7 +52,6 @@ type VirtualService struct { } type Metric struct { - Type string `json:"type"` Name string `json:"name"` Interval string `json:"interval"` Threshold int `json:"threshold"` diff --git a/pkg/apis/rollout/v1beta1/zz_generated.deepcopy.go b/pkg/apis/rollout/v1beta1/zz_generated.deepcopy.go index ff20d572..01f457a2 100644 --- a/pkg/apis/rollout/v1beta1/zz_generated.deepcopy.go +++ b/pkg/apis/rollout/v1beta1/zz_generated.deepcopy.go @@ -45,7 +45,7 @@ func (in *Rollout) DeepCopyInto(out *Rollout) { *out = *in out.TypeMeta = in.TypeMeta in.ObjectMeta.DeepCopyInto(&out.ObjectMeta) - out.Spec = in.Spec + in.Spec.DeepCopyInto(&out.Spec) out.Status = in.Status return } @@ -107,7 +107,11 @@ func (in *RolloutSpec) DeepCopyInto(out *RolloutSpec) { out.Primary = in.Primary out.Canary = in.Canary out.VirtualService = in.VirtualService - out.Metric = in.Metric + if in.Metrics != nil { + in, out := &in.Metrics, &out.Metrics + *out = make([]Metric, len(*in)) + copy(*out, *in) + } return } diff --git a/pkg/controller/controller.go b/pkg/controller/controller.go index 12f43cc3..76e620e9 100644 --- a/pkg/controller/controller.go +++ b/pkg/controller/controller.go @@ -15,7 +15,6 @@ import ( "go.uber.org/zap" corev1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/api/errors" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" utilruntime "k8s.io/apimachinery/pkg/util/runtime" "k8s.io/apimachinery/pkg/util/wait" "k8s.io/client-go/kubernetes" @@ -175,7 +174,6 @@ func (c *Controller) syncHandler(key string) error { utilruntime.HandleError(fmt.Errorf("invalid resource key: %s", key)) return nil } - rollout, err := c.rolloutLister.Rollouts(namespace).Get(name) if errors.IsNotFound(err) { utilruntime.HandleError(fmt.Errorf("rollout '%s' in work queue no longer exists", key)) @@ -198,40 +196,6 @@ func (c *Controller) enqueueRollout(obj interface{}) { c.workqueue.AddRateLimited(key) } -func (c *Controller) handleObject(obj interface{}) { - var object metav1.Object - var ok bool - if object, ok = obj.(metav1.Object); !ok { - tombstone, ok := obj.(cache.DeletedFinalStateUnknown) - if !ok { - utilruntime.HandleError(fmt.Errorf("error decoding object, invalid type")) - return - } - object, ok = tombstone.Obj.(metav1.Object) - if !ok { - utilruntime.HandleError(fmt.Errorf("error decoding object tombstone, invalid type")) - return - } - c.logger.Debugf("Recovered deleted object '%s' from tombstone", object.GetName()) - } - c.logger.Debugf("Processing object: %s", object.GetName()) - if ownerRef := metav1.GetControllerOf(object); ownerRef != nil { - if ownerRef.Kind != "Rollout" { - return - } - - vs, err := c.rolloutLister.Rollouts(object.GetNamespace()).Get(ownerRef.Name) - if err != nil { - c.logger.Debugf("ignoring orphaned object '%s' of '%s'", object.GetSelfLink(), ownerRef.Name) - return - } - - c.enqueueRollout(vs) - return - } - -} - func (c *Controller) recordEventInfof(r *rolloutv1.Rollout, template string, args ...interface{}) { c.logger.Infof(template, args...) c.recorder.Event(r, corev1.EventTypeNormal, "Synced", fmt.Sprintf(template, args...)) diff --git a/pkg/controller/deployment.go b/pkg/controller/deployment.go index 4987656e..2c08401b 100644 --- a/pkg/controller/deployment.go +++ b/pkg/controller/deployment.go @@ -2,6 +2,7 @@ package controller import ( "fmt" + "time" istiov1alpha3 "github.com/knative/pkg/apis/istio/v1alpha3" rolloutv1 "github.com/stefanprodan/steerer/pkg/apis/rollout/v1beta1" @@ -32,13 +33,13 @@ func (c *Controller) advanceDeploymentRollout(name string, namespace string) { } // gate stage: check if primary deployment exists and is healthy - primary, ok := c.getDeployment(r.Spec.Primary.Name, r.Namespace) + primary, ok := c.getDeployment(r, r.Spec.Primary.Name, r.Namespace) if !ok { return } // gate stage: check if canary deployment exists and is healthy - canary, ok := c.getDeployment(r.Spec.Canary.Name, r.Namespace) + canary, ok := c.getDeployment(r, r.Spec.Canary.Name, r.Namespace) if !ok { return } @@ -60,7 +61,7 @@ func (c *Controller) advanceDeploymentRollout(name string, namespace string) { if canaryRoute.Weight == 0 { c.recordEventInfof(r, "Starting rollout for %s.%s", r.Name, r.Namespace) } else { - if ok := c.checkDeploymentSuccessRate(r); !ok { + if ok := c.checkDeploymentMetrics(r); !ok { return } } @@ -161,10 +162,10 @@ func (c *Controller) updateRolloutStatus(r *rolloutv1.Rollout, status string) bo } -func (c *Controller) getDeployment(name string, namespace string) (*appsv1.Deployment, bool) { +func (c *Controller) getDeployment(r *rolloutv1.Rollout, name string, namespace string) (*appsv1.Deployment, bool) { dep, err := c.kubeClient.AppsV1().Deployments(namespace).Get(name, v1.GetOptions{}) if err != nil { - c.logger.Errorf("Deployment %s.%s not found", name, namespace) + c.recordEventErrorf(r, "Deployment %s.%s not found", name, namespace) return nil, false } @@ -180,17 +181,34 @@ func (c *Controller) getDeployment(name string, namespace string) (*appsv1.Deplo return dep, true } -func (c *Controller) checkDeploymentSuccessRate(r *rolloutv1.Rollout) bool { - val, err := c.getDeploymentMetric(r.Spec.Canary.Name, r.Namespace, r.Spec.Metric.Name, r.Spec.Metric.Interval) - if err != nil { - c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.metricsServer, err) - return false - } +func (c *Controller) checkDeploymentMetrics(r *rolloutv1.Rollout) bool { + for _, metric := range r.Spec.Metrics { + if metric.Name == "istio_requests_total" { + val, err := c.getDeploymentCounter(r.Spec.Canary.Name, r.Namespace, metric.Name, metric.Interval) + if err != nil { + c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.metricsServer, err) + return false + } + if float64(metric.Threshold) > val { + c.recordEventErrorf(r, "Halt rollout %s.%s success rate %.2f%% < %v%%", + r.Name, r.Namespace, val, metric.Threshold) + return false + } + } - if float64(r.Spec.Metric.Threshold) > val { - c.recordEventErrorf(r, "Halt rollout %s.%s success rate %.2f%% < %v%%", - r.Name, r.Namespace, val, r.Spec.Metric.Threshold) - return false + if metric.Name == "istio_request_duration_seconds_bucket" { + val, err := c.GetDeploymentHistogram(r.Spec.Canary.Name, r.Namespace, metric.Name, metric.Interval) + if err != nil { + c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.metricsServer, err) + return false + } + t := time.Duration(metric.Threshold) * time.Millisecond + if val > t { + c.recordEventErrorf(r, "Halt rollout %s.%s request duration %v > %v", + r.Name, r.Namespace, val, t) + return false + } + } } return true diff --git a/pkg/controller/observer.go b/pkg/controller/observer.go index 8956c534..6358369b 100644 --- a/pkg/controller/observer.go +++ b/pkg/controller/observer.go @@ -68,19 +68,20 @@ func (c *Controller) queryMetric(query string) (*VectorQueryResponse, error) { return &values, nil } -func (c *Controller) getDeploymentMetric(name string, namespace string, counter string, interval string) (float64, error) { - var rate float64 +// istio_requests_total +func (c *Controller) getDeploymentCounter(name string, namespace string, metric string, interval string) (float64, error) { + var rate *float64 querySt := url.QueryEscape(`sum(rate(` + - counter + `{reporter="destination",destination_workload_namespace=~"` + + metric + `{reporter="destination",destination_workload_namespace=~"` + namespace + `",destination_workload=~"` + name + `",response_code!~"5.*"}[1m])) / sum(rate(` + - counter + `{reporter="destination",destination_workload_namespace=~"` + + metric + `{reporter="destination",destination_workload_namespace=~"` + namespace + `",destination_workload=~"` + name + `"}[` + interval + `])) * 100 `) result, err := c.queryMetric(querySt) if err != nil { - return rate, err + return 0, err } for _, v := range result.Data.Result { @@ -89,12 +90,46 @@ func (c *Controller) getDeploymentMetric(name string, namespace string, counter case string: f, err := strconv.ParseFloat(metricValue.(string), 64) if err != nil { - return rate, err + return 0, err } - rate = f + rate = &f } } - return rate, nil + if rate == nil { + return 0, fmt.Errorf("no values found for metric %s", metric) + } + return *rate, nil +} + +// istio_request_duration_seconds_bucket +func (c *Controller) GetDeploymentHistogram(name string, namespace string, metric string, interval string) (time.Duration, error) { + var rate *float64 + querySt := url.QueryEscape(`histogram_quantile(0.99, sum(irate(` + + 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 } func CheckMetricsServer(address string) (bool, error) { diff --git a/pkg/version/version.go b/pkg/version/version.go index dfeb4d61..05c1159f 100644 --- a/pkg/version/version.go +++ b/pkg/version/version.go @@ -1,4 +1,4 @@ package version -var VERSION = "0.0.1-beta.5" +var VERSION = "0.0.1-beta.9" var REVISION = "unknown"