diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index 1dd9a101..bba6b869 100644 --- a/pkg/controller/scheduler.go +++ b/pkg/controller/scheduler.go @@ -1,17 +1,13 @@ package controller import ( - "errors" "fmt" - "strings" "time" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1beta1" "github.com/weaveworks/flagger/pkg/canary" - "github.com/weaveworks/flagger/pkg/metrics/observers" - "github.com/weaveworks/flagger/pkg/metrics/providers" "github.com/weaveworks/flagger/pkg/router" ) @@ -222,7 +218,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh } // check if analysis should be skipped - if skip := c.shouldSkipAnalysis(cd, canaryController, meshRouter, primaryWeight, canaryWeight); skip { + if skip := c.shouldSkipAnalysis(cd, canaryController, meshRouter); skip { return } @@ -533,14 +529,40 @@ func (c *Controller) runBlueGreen(canary *flaggerv1.Canary, canaryController can } -func (c *Controller) shouldSkipAnalysis(canary *flaggerv1.Canary, canaryController canary.Controller, meshRouter router.Interface, primaryWeight int, canaryWeight int) bool { +func (c *Controller) runAnalysis(canary *flaggerv1.Canary) bool { + // run external checks + for _, webhook := range canary.GetAnalysis().Webhooks { + if webhook.Type == "" || webhook.Type == flaggerv1.RolloutHook { + err := CallWebhook(canary.Name, canary.Namespace, flaggerv1.CanaryPhaseProgressing, webhook) + if err != nil { + c.recordEventWarningf(canary, "Halt %s.%s advancement external check %s failed %v", + canary.Name, canary.Namespace, webhook.Name, err) + return false + } + } + } + + ok := c.runBuiltinMetricChecks(canary) + if !ok { + return ok + } + + ok = c.runMetricChecks(canary) + if !ok { + return ok + } + + return true +} + +func (c *Controller) shouldSkipAnalysis(canary *flaggerv1.Canary, canaryController canary.Controller, meshRouter router.Interface) bool { if !canary.SkipAnalysis() { return false } // route all traffic to primary - primaryWeight = 100 - canaryWeight = 0 + primaryWeight := 100 + canaryWeight := 0 if err := meshRouter.SetRoutes(canary, primaryWeight, canaryWeight, false); err != nil { c.recordEventWarningf(canary, "%v", err) return false @@ -657,358 +679,6 @@ func (c *Controller) hasCanaryRevisionChanged(canary *flaggerv1.Canary, canaryCo return false } -func (c *Controller) runConfirmRolloutHooks(canary *flaggerv1.Canary, canaryController canary.Controller) bool { - for _, webhook := range canary.GetAnalysis().Webhooks { - if webhook.Type == flaggerv1.ConfirmRolloutHook { - err := CallWebhook(canary.Name, canary.Namespace, flaggerv1.CanaryPhaseProgressing, webhook) - if err != nil { - if canary.Status.Phase != flaggerv1.CanaryPhaseWaiting { - if err := canaryController.SetStatusPhase(canary, flaggerv1.CanaryPhaseWaiting); err != nil { - c.logger.With("canary", fmt.Sprintf("%s.%s", canary.Name, canary.Namespace)).Errorf("%v", err) - } - c.recordEventWarningf(canary, "Halt %s.%s advancement waiting for approval %s", - canary.Name, canary.Namespace, webhook.Name) - c.alert(canary, "Canary is waiting for approval.", false, flaggerv1.SeverityWarn) - } - return false - } else { - if canary.Status.Phase == flaggerv1.CanaryPhaseWaiting { - if err := canaryController.SetStatusPhase(canary, flaggerv1.CanaryPhaseProgressing); err != nil { - c.logger.With("canary", fmt.Sprintf("%s.%s", canary.Name, canary.Namespace)).Errorf("%v", err) - return false - } - c.recordEventInfof(canary, "Confirm-rollout check %s passed", webhook.Name) - return false - } - } - } - } - return true -} - -func (c *Controller) runConfirmPromotionHooks(canary *flaggerv1.Canary) bool { - for _, webhook := range canary.GetAnalysis().Webhooks { - if webhook.Type == flaggerv1.ConfirmPromotionHook { - err := CallWebhook(canary.Name, canary.Namespace, flaggerv1.CanaryPhaseProgressing, webhook) - if err != nil { - c.recordEventWarningf(canary, "Halt %s.%s advancement waiting for promotion approval %s", - canary.Name, canary.Namespace, webhook.Name) - c.alert(canary, "Canary promotion is waiting for approval.", false, flaggerv1.SeverityWarn) - return false - } else { - c.recordEventInfof(canary, "Confirm-promotion check %s passed", webhook.Name) - } - } - } - return true -} - -func (c *Controller) runPreRolloutHooks(canary *flaggerv1.Canary) bool { - for _, webhook := range canary.GetAnalysis().Webhooks { - if webhook.Type == flaggerv1.PreRolloutHook { - err := CallWebhook(canary.Name, canary.Namespace, flaggerv1.CanaryPhaseProgressing, webhook) - if err != nil { - c.recordEventWarningf(canary, "Halt %s.%s advancement pre-rollout check %s failed %v", - canary.Name, canary.Namespace, webhook.Name, err) - return false - } else { - c.recordEventInfof(canary, "Pre-rollout check %s passed", webhook.Name) - } - } - } - return true -} - -func (c *Controller) runPostRolloutHooks(canary *flaggerv1.Canary, phase flaggerv1.CanaryPhase) bool { - for _, webhook := range canary.GetAnalysis().Webhooks { - if webhook.Type == flaggerv1.PostRolloutHook { - err := CallWebhook(canary.Name, canary.Namespace, phase, webhook) - if err != nil { - c.recordEventWarningf(canary, "Post-rollout hook %s failed %v", webhook.Name, err) - return false - } else { - c.recordEventInfof(canary, "Post-rollout check %s passed", webhook.Name) - } - } - } - return true -} - -func (c *Controller) runRollbackHooks(canary *flaggerv1.Canary, phase flaggerv1.CanaryPhase) bool { - for _, webhook := range canary.GetAnalysis().Webhooks { - if webhook.Type == flaggerv1.RollbackHook { - err := CallWebhook(canary.Name, canary.Namespace, phase, webhook) - if err != nil { - c.recordEventInfof(canary, "Rollback hook %s not signaling a rollback", webhook.Name) - } else { - c.recordEventWarningf(canary, "Rollback check %s passed", webhook.Name) - return true - } - } - } - return false -} - -func (c *Controller) runAnalysis(canary *flaggerv1.Canary) bool { - // run external checks - for _, webhook := range canary.GetAnalysis().Webhooks { - if webhook.Type == "" || webhook.Type == flaggerv1.RolloutHook { - err := CallWebhook(canary.Name, canary.Namespace, flaggerv1.CanaryPhaseProgressing, webhook) - if err != nil { - c.recordEventWarningf(canary, "Halt %s.%s advancement external check %s failed %v", - canary.Name, canary.Namespace, webhook.Name, err) - return false - } - } - } - - ok := c.runBuiltinMetricChecks(canary) - if !ok { - return ok - } - - ok = c.runMetricChecks(canary) - if !ok { - return ok - } - - return true -} - -func (c *Controller) runBuiltinMetricChecks(canary *flaggerv1.Canary) bool { - // override the global provider if one is specified in the canary spec - var metricsProvider string - // set the metrics provider to Crossover Prometheus when Crossover is the mesh provider - // For example, `crossover` metrics provider should be used for `smi:crossover` mesh provider - if strings.Contains(c.meshProvider, "crossover") { - metricsProvider = "crossover" - } else { - metricsProvider = c.meshProvider - } - - if canary.Spec.Provider != "" { - metricsProvider = canary.Spec.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 query Prometheus for the canary Kubernetes service if the canary target is Service - if canary.Spec.TargetRef.Kind == "Service" { - metricsProvider = metricsProvider + MetricsProviderServiceSuffix - } - - // create observer based on the mesh provider - observerFactory := c.observerFactory - - // override the global metrics server if one is specified in the canary spec - if canary.Spec.MetricsServer != "" { - var err error - observerFactory, err = observers.NewFactory(canary.Spec.MetricsServer) - if err != nil { - c.recordEventErrorf(canary, "Error building Prometheus client for %s %v", canary.Spec.MetricsServer, err) - return false - } - } - observer := observerFactory.Observer(metricsProvider) - - // run metrics checks - for _, metric := range canary.GetAnalysis().Metrics { - if metric.Interval == "" { - metric.Interval = canary.GetMetricInterval() - } - - if metric.Name == "request-success-rate" { - val, err := observer.GetRequestSuccessRate(toMetricModel(canary, metric.Interval)) - if err != nil { - if errors.Is(err, providers.ErrNoValuesFound) { - c.recordEventWarningf(canary, - "Halt advancement no values found for %s metric %s probably %s.%s is not receiving traffic: %v", - metricsProvider, metric.Name, canary.Spec.TargetRef.Name, canary.Namespace, err) - } else { - c.recordEventErrorf(canary, "Prometheus query failed: %v", err) - } - return false - } - - if metric.ThresholdRange != nil { - tr := *metric.ThresholdRange - if tr.Min != nil && val < *tr.Min { - c.recordEventWarningf(canary, "Halt %s.%s advancement success rate %.2f%% < %v%%", - canary.Name, canary.Namespace, val, *tr.Min) - return false - } - if tr.Max != nil && val > *tr.Max { - c.recordEventWarningf(canary, "Halt %s.%s advancement success rate %.2f%% > %v%%", - canary.Name, canary.Namespace, val, *tr.Max) - return false - } - } else if metric.Threshold > val { - c.recordEventWarningf(canary, "Halt %s.%s advancement success rate %.2f%% < %v%%", - canary.Name, canary.Namespace, val, metric.Threshold) - return false - } - } - - if metric.Name == "request-duration" { - val, err := observer.GetRequestDuration(toMetricModel(canary, metric.Interval)) - if err != nil { - if errors.Is(err, providers.ErrNoValuesFound) { - c.recordEventWarningf(canary, "Halt advancement no values found for %s metric %s probably %s.%s is not receiving traffic", - metricsProvider, metric.Name, canary.Spec.TargetRef.Name, canary.Namespace) - } else { - c.recordEventErrorf(canary, "Prometheus query failed: %v", err) - } - return false - } - if metric.ThresholdRange != nil { - tr := *metric.ThresholdRange - if tr.Min != nil && val < time.Duration(*tr.Min)*time.Millisecond { - c.recordEventWarningf(canary, "Halt %s.%s advancement request duration %v < %v", - canary.Name, canary.Namespace, val, time.Duration(*tr.Min)*time.Millisecond) - return false - } - if tr.Max != nil && val > time.Duration(*tr.Max)*time.Millisecond { - c.recordEventWarningf(canary, "Halt %s.%s advancement request duration %v > %v", - canary.Name, canary.Namespace, val, time.Duration(*tr.Max)*time.Millisecond) - return false - } - } else if val > time.Duration(metric.Threshold)*time.Millisecond { - c.recordEventWarningf(canary, "Halt %s.%s advancement request duration %v > %v", - canary.Name, canary.Namespace, val, time.Duration(metric.Threshold)*time.Millisecond) - return false - } - } - - // in-line PromQL - if metric.Query != "" { - val, err := observerFactory.Client.RunQuery(metric.Query) - if err != nil { - if errors.Is(err, providers.ErrNoValuesFound) { - c.recordEventWarningf(canary, "Halt advancement no values found for metric: %s", - metric.Name) - } else { - c.recordEventErrorf(canary, "Prometheus query failed for %s: %v", metric.Name, err) - } - return false - } - if metric.ThresholdRange != nil { - tr := *metric.ThresholdRange - if tr.Min != nil && val < *tr.Min { - c.recordEventWarningf(canary, "Halt %s.%s advancement %s %.2f < %v", - canary.Name, canary.Namespace, metric.Name, val, *tr.Min) - return false - } - if tr.Max != nil && val > *tr.Max { - c.recordEventWarningf(canary, "Halt %s.%s advancement %s %.2f > %v", - canary.Name, canary.Namespace, metric.Name, val, *tr.Max) - return false - } - } else if val > metric.Threshold { - c.recordEventWarningf(canary, "Halt %s.%s advancement %s %.2f > %v", - canary.Name, canary.Namespace, metric.Name, val, metric.Threshold) - return false - } - } - } - - return true -} - -func (c *Controller) runMetricChecks(canary *flaggerv1.Canary) bool { - for _, metric := range canary.GetAnalysis().Metrics { - if metric.TemplateRef != nil { - namespace := canary.Namespace - if metric.TemplateRef.Namespace != "" { - namespace = metric.TemplateRef.Namespace - } - - template, err := c.flaggerInformers.MetricInformer.Lister().MetricTemplates(namespace).Get(metric.TemplateRef.Name) - if err != nil { - c.recordEventErrorf(canary, "Metric template %s.%s error: %v", metric.TemplateRef.Name, namespace, err) - return false - } - - var credentials map[string][]byte - if template.Spec.Provider.SecretRef != nil { - secret, err := c.kubeClient.CoreV1().Secrets(namespace).Get(template.Spec.Provider.SecretRef.Name, metav1.GetOptions{}) - if err != nil { - c.recordEventErrorf(canary, "Metric template %s.%s secret %s error: %v", - metric.TemplateRef.Name, namespace, template.Spec.Provider.SecretRef.Name, err) - return false - } - credentials = secret.Data - } - - factory := providers.Factory{} - provider, err := factory.Provider(metric.Interval, template.Spec.Provider, credentials) - if err != nil { - c.recordEventErrorf(canary, "Metric template %s.%s provider %s error: %v", - metric.TemplateRef.Name, namespace, template.Spec.Provider.Type, err) - return false - } - - query, err := observers.RenderQuery(template.Spec.Query, toMetricModel(canary, metric.Interval)) - if err != nil { - c.recordEventErrorf(canary, "Metric template %s.%s query render error: %v", - metric.TemplateRef.Name, namespace, err) - return false - } - - val, err := provider.RunQuery(query) - if err != nil { - if errors.Is(err, providers.ErrNoValuesFound) { - c.recordEventWarningf(canary, "Halt advancement no values found for custom metric: %s: %v", - metric.Name, err) - } else { - c.recordEventErrorf(canary, "Metric query failed for %s: %v", metric.Name, err) - } - return false - } - - if metric.ThresholdRange != nil { - tr := *metric.ThresholdRange - if tr.Min != nil && val < *tr.Min { - c.recordEventWarningf(canary, "Halt %s.%s advancement %s %.2f < %v", - canary.Name, canary.Namespace, metric.Name, val, *tr.Min) - return false - } - if tr.Max != nil && val > *tr.Max { - c.recordEventWarningf(canary, "Halt %s.%s advancement %s %.2f > %v", - canary.Name, canary.Namespace, metric.Name, val, *tr.Max) - return false - } - } else if val > metric.Threshold { - c.recordEventWarningf(canary, "Halt %s.%s advancement %s %.2f > %v", - canary.Name, canary.Namespace, metric.Name, val, metric.Threshold) - return false - } - } - } - - return true -} - -func toMetricModel(r *flaggerv1.Canary, interval string) flaggerv1.MetricTemplateModel { - service := r.Spec.TargetRef.Name - if r.Spec.Service.Name != "" { - service = r.Spec.Service.Name - } - ingress := r.Spec.TargetRef.Name - if r.Spec.IngressRef != nil { - ingress = r.Spec.IngressRef.Name - } - return flaggerv1.MetricTemplateModel{ - Name: r.Name, - Namespace: r.Namespace, - Target: r.Spec.TargetRef.Name, - Service: service, - Ingress: ingress, - Interval: interval, - } -} - func (c *Controller) rollback(canary *flaggerv1.Canary, canaryController canary.Controller, meshRouter router.Interface) { if canary.Status.FailedChecks >= canary.GetAnalysisThreshold() { c.recordEventWarningf(canary, "Rolling back %s.%s failed checks threshold reached %v", diff --git a/pkg/controller/scheduler_hooks.go b/pkg/controller/scheduler_hooks.go new file mode 100644 index 00000000..9e72ee37 --- /dev/null +++ b/pkg/controller/scheduler_hooks.go @@ -0,0 +1,100 @@ +package controller + +import ( + "fmt" + + flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1beta1" + "github.com/weaveworks/flagger/pkg/canary" +) + +func (c *Controller) runConfirmRolloutHooks(canary *flaggerv1.Canary, canaryController canary.Controller) bool { + for _, webhook := range canary.GetAnalysis().Webhooks { + if webhook.Type == flaggerv1.ConfirmRolloutHook { + err := CallWebhook(canary.Name, canary.Namespace, flaggerv1.CanaryPhaseProgressing, webhook) + if err != nil { + if canary.Status.Phase != flaggerv1.CanaryPhaseWaiting { + if err := canaryController.SetStatusPhase(canary, flaggerv1.CanaryPhaseWaiting); err != nil { + c.logger.With("canary", fmt.Sprintf("%s.%s", canary.Name, canary.Namespace)).Errorf("%v", err) + } + c.recordEventWarningf(canary, "Halt %s.%s advancement waiting for approval %s", + canary.Name, canary.Namespace, webhook.Name) + c.alert(canary, "Canary is waiting for approval.", false, flaggerv1.SeverityWarn) + } + return false + } else { + if canary.Status.Phase == flaggerv1.CanaryPhaseWaiting { + if err := canaryController.SetStatusPhase(canary, flaggerv1.CanaryPhaseProgressing); err != nil { + c.logger.With("canary", fmt.Sprintf("%s.%s", canary.Name, canary.Namespace)).Errorf("%v", err) + return false + } + c.recordEventInfof(canary, "Confirm-rollout check %s passed", webhook.Name) + return false + } + } + } + } + return true +} + +func (c *Controller) runConfirmPromotionHooks(canary *flaggerv1.Canary) bool { + for _, webhook := range canary.GetAnalysis().Webhooks { + if webhook.Type == flaggerv1.ConfirmPromotionHook { + err := CallWebhook(canary.Name, canary.Namespace, flaggerv1.CanaryPhaseProgressing, webhook) + if err != nil { + c.recordEventWarningf(canary, "Halt %s.%s advancement waiting for promotion approval %s", + canary.Name, canary.Namespace, webhook.Name) + c.alert(canary, "Canary promotion is waiting for approval.", false, flaggerv1.SeverityWarn) + return false + } else { + c.recordEventInfof(canary, "Confirm-promotion check %s passed", webhook.Name) + } + } + } + return true +} + +func (c *Controller) runPreRolloutHooks(canary *flaggerv1.Canary) bool { + for _, webhook := range canary.GetAnalysis().Webhooks { + if webhook.Type == flaggerv1.PreRolloutHook { + err := CallWebhook(canary.Name, canary.Namespace, flaggerv1.CanaryPhaseProgressing, webhook) + if err != nil { + c.recordEventWarningf(canary, "Halt %s.%s advancement pre-rollout check %s failed %v", + canary.Name, canary.Namespace, webhook.Name, err) + return false + } else { + c.recordEventInfof(canary, "Pre-rollout check %s passed", webhook.Name) + } + } + } + return true +} + +func (c *Controller) runPostRolloutHooks(canary *flaggerv1.Canary, phase flaggerv1.CanaryPhase) bool { + for _, webhook := range canary.GetAnalysis().Webhooks { + if webhook.Type == flaggerv1.PostRolloutHook { + err := CallWebhook(canary.Name, canary.Namespace, phase, webhook) + if err != nil { + c.recordEventWarningf(canary, "Post-rollout hook %s failed %v", webhook.Name, err) + return false + } else { + c.recordEventInfof(canary, "Post-rollout check %s passed", webhook.Name) + } + } + } + return true +} + +func (c *Controller) runRollbackHooks(canary *flaggerv1.Canary, phase flaggerv1.CanaryPhase) bool { + for _, webhook := range canary.GetAnalysis().Webhooks { + if webhook.Type == flaggerv1.RollbackHook { + err := CallWebhook(canary.Name, canary.Namespace, phase, webhook) + if err != nil { + c.recordEventInfof(canary, "Rollback hook %s not signaling a rollback", webhook.Name) + } else { + c.recordEventWarningf(canary, "Rollback check %s passed", webhook.Name) + return true + } + } + } + return false +} diff --git a/pkg/controller/scheduler_metrics.go b/pkg/controller/scheduler_metrics.go new file mode 100644 index 00000000..16731f09 --- /dev/null +++ b/pkg/controller/scheduler_metrics.go @@ -0,0 +1,247 @@ +package controller + +import ( + "errors" + "strings" + "time" + + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + + flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1beta1" + "github.com/weaveworks/flagger/pkg/metrics/observers" + "github.com/weaveworks/flagger/pkg/metrics/providers" +) + +func (c *Controller) runBuiltinMetricChecks(canary *flaggerv1.Canary) bool { + // override the global provider if one is specified in the canary spec + var metricsProvider string + // set the metrics provider to Crossover Prometheus when Crossover is the mesh provider + // For example, `crossover` metrics provider should be used for `smi:crossover` mesh provider + if strings.Contains(c.meshProvider, "crossover") { + metricsProvider = "crossover" + } else { + metricsProvider = c.meshProvider + } + + if canary.Spec.Provider != "" { + metricsProvider = canary.Spec.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 query Prometheus for the canary Kubernetes service if the canary target is Service + if canary.Spec.TargetRef.Kind == "Service" { + metricsProvider = metricsProvider + MetricsProviderServiceSuffix + } + + // create observer based on the mesh provider + observerFactory := c.observerFactory + + // override the global metrics server if one is specified in the canary spec + if canary.Spec.MetricsServer != "" { + var err error + observerFactory, err = observers.NewFactory(canary.Spec.MetricsServer) + if err != nil { + c.recordEventErrorf(canary, "Error building Prometheus client for %s %v", canary.Spec.MetricsServer, err) + return false + } + } + observer := observerFactory.Observer(metricsProvider) + + // run metrics checks + for _, metric := range canary.GetAnalysis().Metrics { + if metric.Interval == "" { + metric.Interval = canary.GetMetricInterval() + } + + if metric.Name == "request-success-rate" { + val, err := observer.GetRequestSuccessRate(toMetricModel(canary, metric.Interval)) + if err != nil { + if errors.Is(err, providers.ErrNoValuesFound) { + c.recordEventWarningf(canary, + "Halt advancement no values found for %s metric %s probably %s.%s is not receiving traffic: %v", + metricsProvider, metric.Name, canary.Spec.TargetRef.Name, canary.Namespace, err) + } else { + c.recordEventErrorf(canary, "Prometheus query failed: %v", err) + } + return false + } + + if metric.ThresholdRange != nil { + tr := *metric.ThresholdRange + if tr.Min != nil && val < *tr.Min { + c.recordEventWarningf(canary, "Halt %s.%s advancement success rate %.2f%% < %v%%", + canary.Name, canary.Namespace, val, *tr.Min) + return false + } + if tr.Max != nil && val > *tr.Max { + c.recordEventWarningf(canary, "Halt %s.%s advancement success rate %.2f%% > %v%%", + canary.Name, canary.Namespace, val, *tr.Max) + return false + } + } else if metric.Threshold > val { + c.recordEventWarningf(canary, "Halt %s.%s advancement success rate %.2f%% < %v%%", + canary.Name, canary.Namespace, val, metric.Threshold) + return false + } + } + + if metric.Name == "request-duration" { + val, err := observer.GetRequestDuration(toMetricModel(canary, metric.Interval)) + if err != nil { + if errors.Is(err, providers.ErrNoValuesFound) { + c.recordEventWarningf(canary, "Halt advancement no values found for %s metric %s probably %s.%s is not receiving traffic", + metricsProvider, metric.Name, canary.Spec.TargetRef.Name, canary.Namespace) + } else { + c.recordEventErrorf(canary, "Prometheus query failed: %v", err) + } + return false + } + if metric.ThresholdRange != nil { + tr := *metric.ThresholdRange + if tr.Min != nil && val < time.Duration(*tr.Min)*time.Millisecond { + c.recordEventWarningf(canary, "Halt %s.%s advancement request duration %v < %v", + canary.Name, canary.Namespace, val, time.Duration(*tr.Min)*time.Millisecond) + return false + } + if tr.Max != nil && val > time.Duration(*tr.Max)*time.Millisecond { + c.recordEventWarningf(canary, "Halt %s.%s advancement request duration %v > %v", + canary.Name, canary.Namespace, val, time.Duration(*tr.Max)*time.Millisecond) + return false + } + } else if val > time.Duration(metric.Threshold)*time.Millisecond { + c.recordEventWarningf(canary, "Halt %s.%s advancement request duration %v > %v", + canary.Name, canary.Namespace, val, time.Duration(metric.Threshold)*time.Millisecond) + return false + } + } + + // in-line PromQL + if metric.Query != "" { + val, err := observerFactory.Client.RunQuery(metric.Query) + if err != nil { + if errors.Is(err, providers.ErrNoValuesFound) { + c.recordEventWarningf(canary, "Halt advancement no values found for metric: %s", + metric.Name) + } else { + c.recordEventErrorf(canary, "Prometheus query failed for %s: %v", metric.Name, err) + } + return false + } + if metric.ThresholdRange != nil { + tr := *metric.ThresholdRange + if tr.Min != nil && val < *tr.Min { + c.recordEventWarningf(canary, "Halt %s.%s advancement %s %.2f < %v", + canary.Name, canary.Namespace, metric.Name, val, *tr.Min) + return false + } + if tr.Max != nil && val > *tr.Max { + c.recordEventWarningf(canary, "Halt %s.%s advancement %s %.2f > %v", + canary.Name, canary.Namespace, metric.Name, val, *tr.Max) + return false + } + } else if val > metric.Threshold { + c.recordEventWarningf(canary, "Halt %s.%s advancement %s %.2f > %v", + canary.Name, canary.Namespace, metric.Name, val, metric.Threshold) + return false + } + } + } + + return true +} + +func (c *Controller) runMetricChecks(canary *flaggerv1.Canary) bool { + for _, metric := range canary.GetAnalysis().Metrics { + if metric.TemplateRef != nil { + namespace := canary.Namespace + if metric.TemplateRef.Namespace != "" { + namespace = metric.TemplateRef.Namespace + } + + template, err := c.flaggerInformers.MetricInformer.Lister().MetricTemplates(namespace).Get(metric.TemplateRef.Name) + if err != nil { + c.recordEventErrorf(canary, "Metric template %s.%s error: %v", metric.TemplateRef.Name, namespace, err) + return false + } + + var credentials map[string][]byte + if template.Spec.Provider.SecretRef != nil { + secret, err := c.kubeClient.CoreV1().Secrets(namespace).Get(template.Spec.Provider.SecretRef.Name, metav1.GetOptions{}) + if err != nil { + c.recordEventErrorf(canary, "Metric template %s.%s secret %s error: %v", + metric.TemplateRef.Name, namespace, template.Spec.Provider.SecretRef.Name, err) + return false + } + credentials = secret.Data + } + + factory := providers.Factory{} + provider, err := factory.Provider(metric.Interval, template.Spec.Provider, credentials) + if err != nil { + c.recordEventErrorf(canary, "Metric template %s.%s provider %s error: %v", + metric.TemplateRef.Name, namespace, template.Spec.Provider.Type, err) + return false + } + + query, err := observers.RenderQuery(template.Spec.Query, toMetricModel(canary, metric.Interval)) + if err != nil { + c.recordEventErrorf(canary, "Metric template %s.%s query render error: %v", + metric.TemplateRef.Name, namespace, err) + return false + } + + val, err := provider.RunQuery(query) + if err != nil { + if errors.Is(err, providers.ErrNoValuesFound) { + c.recordEventWarningf(canary, "Halt advancement no values found for custom metric: %s: %v", + metric.Name, err) + } else { + c.recordEventErrorf(canary, "Metric query failed for %s: %v", metric.Name, err) + } + return false + } + + if metric.ThresholdRange != nil { + tr := *metric.ThresholdRange + if tr.Min != nil && val < *tr.Min { + c.recordEventWarningf(canary, "Halt %s.%s advancement %s %.2f < %v", + canary.Name, canary.Namespace, metric.Name, val, *tr.Min) + return false + } + if tr.Max != nil && val > *tr.Max { + c.recordEventWarningf(canary, "Halt %s.%s advancement %s %.2f > %v", + canary.Name, canary.Namespace, metric.Name, val, *tr.Max) + return false + } + } else if val > metric.Threshold { + c.recordEventWarningf(canary, "Halt %s.%s advancement %s %.2f > %v", + canary.Name, canary.Namespace, metric.Name, val, metric.Threshold) + return false + } + } + } + + return true +} + +func toMetricModel(r *flaggerv1.Canary, interval string) flaggerv1.MetricTemplateModel { + service := r.Spec.TargetRef.Name + if r.Spec.Service.Name != "" { + service = r.Spec.Service.Name + } + ingress := r.Spec.TargetRef.Name + if r.Spec.IngressRef != nil { + ingress = r.Spec.IngressRef.Name + } + return flaggerv1.MetricTemplateModel{ + Name: r.Name, + Namespace: r.Namespace, + Target: r.Spec.TargetRef.Name, + Service: service, + Ingress: ingress, + Interval: interval, + } +}