mirror of
https://github.com/fluxcd/flagger.git
synced 2026-04-15 06:57:34 +00:00
pkg/controller: check metrics server's availability during initalization
This commit is contained in:
@@ -14,10 +14,6 @@ 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
|
||||
@@ -119,6 +115,14 @@ func (c *Controller) advanceCanary(name string, namespace string) {
|
||||
return
|
||||
}
|
||||
|
||||
// check metric servers' availability
|
||||
if !cd.SkipAnalysis() && (cd.Status.Phase == "" || cd.Status.Phase == flaggerv1.CanaryPhaseInitializing) {
|
||||
if err := c.checkMetricProviderAvailability(cd); err != nil {
|
||||
c.recordEventErrorf(cd, "Error checking metric providers: %v", err)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// init mesh router
|
||||
meshRouter := c.routerFactory.MeshRouter(provider, labelSelector)
|
||||
|
||||
|
||||
@@ -1,33 +0,0 @@
|
||||
package controller
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
|
||||
flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1beta1"
|
||||
clientset "github.com/weaveworks/flagger/pkg/client/clientset/versioned"
|
||||
)
|
||||
|
||||
func assertPhase(flaggerClient clientset.Interface, canary string, phase flaggerv1.CanaryPhase) error {
|
||||
c, err := flaggerClient.FlaggerV1beta1().Canaries("default").Get(context.TODO(), canary, metav1.GetOptions{})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if c.Status.Phase != phase {
|
||||
return fmt.Errorf("Got canary state %s wanted %s", c.Status.Phase, phase)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func alwaysReady() bool {
|
||||
return true
|
||||
}
|
||||
|
||||
func toFloatPtr(val int) *float64 {
|
||||
v := float64(val)
|
||||
return &v
|
||||
}
|
||||
@@ -79,7 +79,7 @@ func newDaemonSetFixture(c *flaggerv1.Canary) daemonSetFixture {
|
||||
rf := router.NewFactory(nil, kubeClient, flaggerClient, "annotationsPrefix", "", logger, flaggerClient)
|
||||
|
||||
// init observer
|
||||
observerFactory, _ := observers.NewFactory("fake")
|
||||
observerFactory, _ := observers.NewFactory(testMetricsServerURL)
|
||||
|
||||
// init canary factory
|
||||
configTracker := &canary.ConfigTracker{
|
||||
@@ -616,7 +616,7 @@ func newDaemonSetTestService() *corev1.Service {
|
||||
func newDaemonSetTestMetricTemplate() *flaggerv1.MetricTemplate {
|
||||
provider := flaggerv1.MetricTemplateProvider{
|
||||
Type: "prometheus",
|
||||
Address: "fake",
|
||||
Address: testMetricsServerURL,
|
||||
SecretRef: &corev1.LocalObjectReference{
|
||||
Name: "podinfo-secret-env",
|
||||
},
|
||||
|
||||
@@ -107,7 +107,7 @@ func newDeploymentFixture(c *flaggerv1.Canary) fixture {
|
||||
rf := router.NewFactory(nil, kubeClient, flaggerClient, "annotationsPrefix", "", logger, flaggerClient)
|
||||
|
||||
// init observer
|
||||
observerFactory, _ := observers.NewFactory("fake")
|
||||
observerFactory, _ := observers.NewFactory(testMetricsServerURL)
|
||||
|
||||
// init canary factory
|
||||
configTracker := &canary.ConfigTracker{
|
||||
@@ -708,7 +708,7 @@ func newDeploymentTestHPA() *hpav2.HorizontalPodAutoscaler {
|
||||
func newDeploymentTestMetricTemplate() *flaggerv1.MetricTemplate {
|
||||
provider := flaggerv1.MetricTemplateProvider{
|
||||
Type: "prometheus",
|
||||
Address: "fake",
|
||||
Address: testMetricsServerURL,
|
||||
SecretRef: &corev1.LocalObjectReference{
|
||||
Name: "podinfo-secret-env",
|
||||
},
|
||||
|
||||
@@ -84,11 +84,9 @@ func TestScheduler_DeploymentRollback(t *testing.T) {
|
||||
|
||||
// run metric checks
|
||||
mocks.ctrl.advanceCanary("podinfo", "default")
|
||||
require.NoError(t, err)
|
||||
|
||||
// finalise analysis
|
||||
mocks.ctrl.advanceCanary("podinfo", "default")
|
||||
require.NoError(t, err)
|
||||
|
||||
// check status
|
||||
c, err = mocks.flaggerClient.FlaggerV1beta1().Canaries("default").Get(context.TODO(), "podinfo", metav1.GetOptions{})
|
||||
|
||||
@@ -3,6 +3,7 @@ package controller
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
@@ -13,6 +14,66 @@ import (
|
||||
"github.com/weaveworks/flagger/pkg/metrics/providers"
|
||||
)
|
||||
|
||||
const (
|
||||
MetricsProviderServiceSuffix = ":service"
|
||||
)
|
||||
|
||||
// to be called during canary initialization
|
||||
func (c *Controller) checkMetricProviderAvailability(canary *flaggerv1.Canary) error {
|
||||
for _, metric := range canary.GetAnalysis().Metrics {
|
||||
if metric.Name == "request-success-rate" || metric.Name == "request-duration" {
|
||||
observerFactory := c.observerFactory
|
||||
if canary.Spec.MetricsServer != "" {
|
||||
var err error
|
||||
observerFactory, err = observers.NewFactory(canary.Spec.MetricsServer)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error building Prometheus client for %s %v", canary.Spec.MetricsServer, err)
|
||||
}
|
||||
}
|
||||
if ok, err := observerFactory.Client.IsOnline(); !ok || err != nil {
|
||||
return fmt.Errorf("prometheus not avaiable: %v", err)
|
||||
}
|
||||
continue
|
||||
}
|
||||
|
||||
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 {
|
||||
return fmt.Errorf("metric template %s.%s error: %v", metric.TemplateRef.Name, namespace, err)
|
||||
}
|
||||
|
||||
var credentials map[string][]byte
|
||||
if template.Spec.Provider.SecretRef != nil {
|
||||
secret, err := c.kubeClient.CoreV1().Secrets(namespace).Get(context.TODO(), template.Spec.Provider.SecretRef.Name, metav1.GetOptions{})
|
||||
if err != nil {
|
||||
return fmt.Errorf("metric template %s.%s secret %s error: %v",
|
||||
metric.TemplateRef.Name, namespace, template.Spec.Provider.SecretRef.Name, err)
|
||||
}
|
||||
credentials = secret.Data
|
||||
}
|
||||
|
||||
factory := providers.Factory{}
|
||||
provider, err := factory.Provider(metric.Interval, template.Spec.Provider, credentials)
|
||||
if err != nil {
|
||||
return fmt.Errorf("metric template %s.%s provider %s error: %v",
|
||||
metric.TemplateRef.Name, namespace, template.Spec.Provider.Type, err)
|
||||
}
|
||||
|
||||
if ok, err := provider.IsOnline(); !ok || err != nil {
|
||||
return fmt.Errorf("%v in metric tempalte %s.%s not avaiable: %v", template.Spec.Provider.Type,
|
||||
template.Name, template.Namespace, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
c.recordEventInfof(canary, "all the metrics providers are available!")
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *Controller) runBuiltinMetricChecks(canary *flaggerv1.Canary) bool {
|
||||
// override the global provider if one is specified in the canary spec
|
||||
var metricsProvider string
|
||||
|
||||
@@ -0,0 +1,53 @@
|
||||
package controller
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"testing"
|
||||
|
||||
flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1beta1"
|
||||
clientset "github.com/weaveworks/flagger/pkg/client/clientset/versioned"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
)
|
||||
|
||||
var testMetricsServerURL string
|
||||
|
||||
func TestMain(m *testing.M) {
|
||||
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if r.URL.Query()["query"][0] == "vector(1)" {
|
||||
// for IsOnline invoked during canary initialization
|
||||
w.Write([]byte(`{"status":"success","data":{"resultType":"vector","result":[{"metric":{},"value":[1545905245.458,"1"]}]}}`))
|
||||
return
|
||||
}
|
||||
w.Write([]byte(`{"status":"success","data":{"resultType":"vector","result":[{"metric":{},"value":[1545905245.458,"100"]}]}}`))
|
||||
}))
|
||||
|
||||
testMetricsServerURL = ts.URL
|
||||
defer ts.Close()
|
||||
os.Exit(m.Run())
|
||||
}
|
||||
|
||||
func assertPhase(flaggerClient clientset.Interface, canary string, phase flaggerv1.CanaryPhase) error {
|
||||
c, err := flaggerClient.FlaggerV1beta1().Canaries("default").Get(context.TODO(), canary, metav1.GetOptions{})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if c.Status.Phase != phase {
|
||||
return fmt.Errorf("got canary state %s wanted %s", c.Status.Phase, phase)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func alwaysReady() bool {
|
||||
return true
|
||||
}
|
||||
|
||||
func toFloatPtr(val int) *float64 {
|
||||
v := float64(val)
|
||||
return &v
|
||||
}
|
||||
Reference in New Issue
Block a user