From 4b17788a77f46e7320b0a9af1f4f1aee3df1525e Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Fri, 12 Apr 2019 16:02:23 +0300 Subject: [PATCH 01/22] Add Envoy request duration P99 query --- pkg/metrics/envoy.go | 54 +++++++++++++++++++++++++++++++++++++++ pkg/metrics/envoy_test.go | 23 +++++++++++++++++ 2 files changed, 77 insertions(+) diff --git a/pkg/metrics/envoy.go b/pkg/metrics/envoy.go index 04002bb7..b4487df3 100644 --- a/pkg/metrics/envoy.go +++ b/pkg/metrics/envoy.go @@ -4,6 +4,7 @@ import ( "fmt" "net/url" "strconv" + "time" ) const envoySuccessRateQuery = ` @@ -63,3 +64,56 @@ func (c *Observer) GetEnvoySuccessRate(name string, namespace string, metric str } return *rate, nil } + +const envoyRequestDurationQuery = ` +histogram_quantile(0.99, sum(rate( +envoy_cluster_upstream_rq_time_bucket{kubernetes_namespace="{{ .Namespace }}", +kubernetes_pod_name=~"{{ .Name }}-[0-9a-zA-Z]+(-[0-9a-zA-Z]+)"} +[{{ .Interval }}])) by (le)) +` + +// GetEnvoyRequestDuration returns the 99P requests delay using envoy_cluster_upstream_rq_time_bucket metrics +func (c *Observer) GetEnvoyRequestDuration(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, envoyRequestDurationQuery) + 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/envoy_test.go b/pkg/metrics/envoy_test.go index 12d93570..0ae9a96c 100644 --- a/pkg/metrics/envoy_test.go +++ b/pkg/metrics/envoy_test.go @@ -26,3 +26,26 @@ func Test_EnvoySuccessRateQueryRender(t *testing.T) { t.Errorf("\nGot %s \nWanted %s", query, expected) } } + +func Test_EnvoyRequestDurationQueryRender(t *testing.T) { + meta := struct { + Name string + Namespace string + Interval string + }{ + "podinfo", + "default", + "1m", + } + + query, err := render(meta, envoyRequestDurationQuery) + if err != nil { + t.Fatal(err) + } + + expected := `histogram_quantile(0.99, sum(rate(envoy_cluster_upstream_rq_time_bucket{kubernetes_namespace="default",kubernetes_pod_name=~"podinfo-[0-9a-zA-Z]+(-[0-9a-zA-Z]+)"}[1m])) by (le))` + + if query != expected { + t.Errorf("\nGot %s \nWanted %s", query, expected) + } +} From c651ef00c9663ac9adb159d81bbfce0e919e9afd Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Fri, 12 Apr 2019 16:09:00 +0300 Subject: [PATCH 02/22] Add Envoy request duration test --- pkg/metrics/observer_test.go | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) diff --git a/pkg/metrics/observer_test.go b/pkg/metrics/observer_test.go index 2bcc0ab5..dab4cd52 100644 --- a/pkg/metrics/observer_test.go +++ b/pkg/metrics/observer_test.go @@ -27,6 +27,25 @@ func TestCanaryObserver_GetEnvoySuccessRate(t *testing.T) { } +func TestCanaryObserver_GetEnvoyRequestDuration(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 := NewObserver(ts.URL) + + val, err := observer.GetEnvoyRequestDuration("podinfo", "default", "envoy_cluster_upstream_rq_time_bucket", "1m") + if err != nil { + t.Fatal(err.Error()) + } + + if val != 200*time.Millisecond { + t.Errorf("Got %v wanted %v", val, 200*time.Millisecond) + } +} + 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"]}]}}` From e091d6a50d1688cbeff886c6c0fa8a0ebf70c1cc Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Fri, 12 Apr 2019 16:25:34 +0300 Subject: [PATCH 03/22] Set Envoy request duration to ms --- pkg/metrics/envoy.go | 2 +- pkg/metrics/observer_test.go | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/pkg/metrics/envoy.go b/pkg/metrics/envoy.go index b4487df3..5373f532 100644 --- a/pkg/metrics/envoy.go +++ b/pkg/metrics/envoy.go @@ -114,6 +114,6 @@ func (c *Observer) GetEnvoyRequestDuration(name string, namespace string, metric if rate == nil { return 0, fmt.Errorf("no values found for metric %s", metric) } - ms := time.Duration(int64(*rate*1000)) * time.Millisecond + ms := time.Duration(int64(*rate)) * time.Millisecond return ms, nil } diff --git a/pkg/metrics/observer_test.go b/pkg/metrics/observer_test.go index dab4cd52..fcd12eb4 100644 --- a/pkg/metrics/observer_test.go +++ b/pkg/metrics/observer_test.go @@ -29,7 +29,7 @@ func TestCanaryObserver_GetEnvoySuccessRate(t *testing.T) { func TestCanaryObserver_GetEnvoyRequestDuration(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"]}]}}` + json := `{"status":"success","data":{"resultType":"vector","result":[{"metric":{},"value":[1545905245.596,"200"]}]}}` w.Write([]byte(json)) })) defer ts.Close() From 352ed898d4a16ad313511a4c5bf8d85298a19370 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Fri, 12 Apr 2019 17:00:04 +0300 Subject: [PATCH 04/22] Add request success rate and duration metrics alias --- pkg/controller/scheduler.go | 101 ++++++++++++++++++++++-------------- 1 file changed, 61 insertions(+), 40 deletions(-) diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index 4c699f73..02c7206b 100644 --- a/pkg/controller/scheduler.go +++ b/pkg/controller/scheduler.go @@ -494,56 +494,77 @@ func (c *Controller) analyseCanary(r *flaggerv1.Canary) bool { metric.Interval = r.GetMetricInterval() } - if metric.Name == "envoy_cluster_upstream_rq" { - val, err := c.observer.GetEnvoySuccessRate(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.GetMetricsServer(), err) + // App Mesh checks + if c.meshProvider == "appmesh" { + if metric.Name == "request-success-rate" || metric.Name == "envoy_cluster_upstream_rq" { + val, err := c.observer.GetEnvoySuccessRate(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.GetMetricsServer(), err) + } + return false } - return false - } - if float64(metric.Threshold) > val { - c.recordEventWarningf(r, "Halt %s.%s advancement success rate %.2f%% < %v%%", - r.Name, r.Namespace, val, metric.Threshold) - return false - } - } - - if metric.Name == "istio_requests_total" { - 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.GetMetricsServer(), err) + if float64(metric.Threshold) > val { + c.recordEventWarningf(r, "Halt %s.%s advancement success rate %.2f%% < %v%%", + r.Name, r.Namespace, val, metric.Threshold) + return false } - return false } - if float64(metric.Threshold) > val { - c.recordEventWarningf(r, "Halt %s.%s advancement success rate %.2f%% < %v%%", - r.Name, r.Namespace, val, metric.Threshold) - return false + + if metric.Name == "request-duration" || metric.Name == "envoy_cluster_upstream_rq_time_bucket" { + val, err := c.observer.GetEnvoyRequestDuration(r.Spec.TargetRef.Name, r.Namespace, metric.Name, metric.Interval) + if err != nil { + c.recordEventErrorf(r, "Metrics server %s query failed: %v", c.observer.GetMetricsServer(), err) + return false + } + t := time.Duration(metric.Threshold) * time.Millisecond + if val > t { + c.recordEventWarningf(r, "Halt %s.%s advancement request duration %v > %v", + r.Name, r.Namespace, val, t) + return false + } } } - if metric.Name == "istio_request_duration_seconds_bucket" { - 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.GetMetricsServer(), err) - return false + // Istio checks + if c.meshProvider == "istio" { + if metric.Name == "request-success-rate" || metric.Name == "istio_requests_total" { + 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.GetMetricsServer(), err) + } + return false + } + if float64(metric.Threshold) > val { + c.recordEventWarningf(r, "Halt %s.%s advancement success rate %.2f%% < %v%%", + r.Name, r.Namespace, val, metric.Threshold) + return false + } } - t := time.Duration(metric.Threshold) * time.Millisecond - if val > t { - c.recordEventWarningf(r, "Halt %s.%s advancement request duration %v > %v", - r.Name, r.Namespace, val, t) - return false + + if metric.Name == "request-duration" || metric.Name == "istio_request_duration_seconds_bucket" { + 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.GetMetricsServer(), err) + return false + } + t := time.Duration(metric.Threshold) * time.Millisecond + if val > t { + c.recordEventWarningf(r, "Halt %s.%s advancement request duration %v > %v", + r.Name, r.Namespace, val, t) + return false + } } } + // custom checks if metric.Query != "" { val, err := c.observer.GetScalar(metric.Query) if err != nil { From 4ab9ceafc16f6a3ad6973a66812170ce3c8d60c3 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Fri, 12 Apr 2019 17:43:47 +0300 Subject: [PATCH 05/22] Use metrics alias in e2e tests --- test/e2e-tests.sh | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/test/e2e-tests.sh b/test/e2e-tests.sh index 3756a10f..eed594b6 100755 --- a/test/e2e-tests.sh +++ b/test/e2e-tests.sh @@ -45,10 +45,10 @@ spec: maxWeight: 50 stepWeight: 10 metrics: - - name: istio_requests_total + - name: request-success-rate threshold: 99 interval: 1m - - name: istio_request_duration_seconds_bucket + - name: request-duration threshold: 500 interval: 30s - name: "404s percentage" From 4ac66299694e88e565f055bceb5d3f104fd12c15 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sat, 13 Apr 2019 15:33:03 +0300 Subject: [PATCH 06/22] Exclude docs branches from CI --- .circleci/config.yml | 2 +- .travis.yml | 11 +++++------ 2 files changed, 6 insertions(+), 7 deletions(-) diff --git a/.circleci/config.yml b/.circleci/config.yml index 6c812451..4a44d24e 100644 --- a/.circleci/config.yml +++ b/.circleci/config.yml @@ -17,6 +17,6 @@ workflows: filters: branches: ignore: - - gh-pages + - /gh-pages.*/ - /docs-.*/ - /release-.*/ diff --git a/.travis.yml b/.travis.yml index b9dc35a3..c1535e23 100644 --- a/.travis.yml +++ b/.travis.yml @@ -1,6 +1,11 @@ sudo: required language: go +branches: + except: + - /gh-pages.*/ + - /docs-.*/ + go: - 1.12.x @@ -12,13 +17,7 @@ addons: packages: - docker-ce -#before_script: -# - go get -u sigs.k8s.io/kind -# - curl https://raw.githubusercontent.com/kubernetes/helm/master/scripts/get | bash -# - 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/ - script: - - set -e - make test-fmt - make test-codegen - go test -race -coverprofile=coverage.txt -covermode=atomic $(go list ./pkg/...) From e0fc5ecb39f0304a2fc8114c4421dcc024307297 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sat, 13 Apr 2019 15:37:41 +0300 Subject: [PATCH 07/22] Add hook type to CRD - pre-rollout execute webhook before routing traffic to canary - rollout execute webhook during the canary analysis on each iteration - post-rollout execute webhook after the canary has been promoted or rolled back Add canary phase to webhook payload --- pkg/apis/flagger/v1alpha3/types.go | 20 +++++++++++++++++--- 1 file changed, 17 insertions(+), 3 deletions(-) diff --git a/pkg/apis/flagger/v1alpha3/types.go b/pkg/apis/flagger/v1alpha3/types.go index 92bc5516..4b0b9d12 100755 --- a/pkg/apis/flagger/v1alpha3/types.go +++ b/pkg/apis/flagger/v1alpha3/types.go @@ -148,11 +148,24 @@ type CanaryMetric struct { Query string `json:"query,omitempty"` } +// HookType can be pre, post or during rollout +type HookType string + +const ( + // RolloutHook execute webhook during the canary analysis + RolloutHook HookType = "rollout" + // PreRolloutHook execute webhook before routing traffic to canary + PreRolloutHook HookType = "pre-rollout" + // PreRolloutHook execute webhook after the canary analysis + PostRolloutHook HookType = "post-rollout" +) + // CanaryWebhook holds the reference to external checks used for canary analysis type CanaryWebhook struct { - Name string `json:"name"` - URL string `json:"url"` - Timeout string `json:"timeout"` + Type HookType `json:"type"` + Name string `json:"name"` + URL string `json:"url"` + Timeout string `json:"timeout"` // +optional Metadata *map[string]string `json:"metadata,omitempty"` } @@ -161,6 +174,7 @@ type CanaryWebhook struct { type CanaryWebhookPayload struct { Name string `json:"name"` Namespace string `json:"namespace"` + Phase CanaryPhase `json:"phase"` Metadata map[string]string `json:"metadata,omitempty"` } From edcff9cd1552dc161b680f63c17ac80dfb78e280 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sat, 13 Apr 2019 15:43:23 +0300 Subject: [PATCH 08/22] Execute pre/post rollout webhooks - halt the canary advancement if pre-rollout hooks are failing - include the canary status (Succeeded/Failed) in the post-rollout webhook payload - ignore post-rollout webhook failures - log pre/post rollout webhook response result --- pkg/controller/scheduler.go | 55 ++++++++++++++++++++++++++++++---- pkg/controller/webhook.go | 3 +- pkg/controller/webhook_test.go | 4 +-- 3 files changed, 54 insertions(+), 8 deletions(-) diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index 02c7206b..9af0b354 100644 --- a/pkg/controller/scheduler.go +++ b/pkg/controller/scheduler.go @@ -240,6 +240,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh } c.recorder.SetStatus(cd, flaggerv1.CanaryFailed) + c.runPostRolloutHooks(cd, flaggerv1.CanaryFailed) return } @@ -247,6 +248,15 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh // skip check if no traffic is routed to canary if canaryWeight == 0 { c.recordEventInfof(cd, "Starting canary analysis for %s.%s", cd.Spec.TargetRef.Name, cd.Namespace) + + // run pre-rollout web hooks + if ok := c.runPreRolloutHooks(cd); !ok { + if err := c.deployer.SetStatusFailedChecks(cd, cd.Status.FailedChecks+1); err != nil { + c.recordEventWarningf(cd, "%v", err) + return + } + return + } } else { if ok := c.analyseCanary(cd); !ok { if err := c.deployer.SetStatusFailedChecks(cd, cd.Status.FailedChecks+1); err != nil { @@ -314,6 +324,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh return } c.recorder.SetStatus(cd, flaggerv1.CanarySucceeded) + c.runPostRolloutHooks(cd, flaggerv1.CanarySucceeded) c.sendNotification(cd, "Canary analysis completed successfully, promotion finished.", false, false) return @@ -380,6 +391,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh return } c.recorder.SetStatus(cd, flaggerv1.CanarySucceeded) + c.runPostRolloutHooks(cd, flaggerv1.CanarySucceeded) c.sendNotification(cd, "Canary analysis completed successfully, promotion finished.", false, false) } @@ -477,14 +489,47 @@ func (c *Controller) hasCanaryRevisionChanged(cd *flaggerv1.Canary) bool { return false } +func (c *Controller) runPreRolloutHooks(canary *flaggerv1.Canary) bool { + for _, webhook := range canary.Spec.CanaryAnalysis.Webhooks { + if webhook.Type == flaggerv1.PreRolloutHook { + err := CallWebhook(canary.Name, canary.Namespace, flaggerv1.CanaryProgressing, 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.Spec.CanaryAnalysis.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) analyseCanary(r *flaggerv1.Canary) bool { // run external checks for _, webhook := range r.Spec.CanaryAnalysis.Webhooks { - err := CallWebhook(r.Name, r.Namespace, webhook) - if err != nil { - c.recordEventWarningf(r, "Halt %s.%s advancement external check %s failed %v", - r.Name, r.Namespace, webhook.Name, err) - return false + if webhook.Type == "" || webhook.Type == flaggerv1.RolloutHook { + err := CallWebhook(r.Name, r.Namespace, flaggerv1.CanaryProgressing, webhook) + if err != nil { + c.recordEventWarningf(r, "Halt %s.%s advancement external check %s failed %v", + r.Name, r.Namespace, webhook.Name, err) + return false + } } } diff --git a/pkg/controller/webhook.go b/pkg/controller/webhook.go index 4bb1d7f0..550921fc 100644 --- a/pkg/controller/webhook.go +++ b/pkg/controller/webhook.go @@ -15,10 +15,11 @@ import ( // CallWebhook does a HTTP POST to an external service and // returns an error if the response status code is non-2xx -func CallWebhook(name string, namespace string, w flaggerv1.CanaryWebhook) error { +func CallWebhook(name string, namespace string, phase flaggerv1.CanaryPhase, w flaggerv1.CanaryWebhook) error { payload := flaggerv1.CanaryWebhookPayload{ Name: name, Namespace: namespace, + Phase: phase, } if w.Metadata != nil { diff --git a/pkg/controller/webhook_test.go b/pkg/controller/webhook_test.go index d9ddb913..528dc516 100644 --- a/pkg/controller/webhook_test.go +++ b/pkg/controller/webhook_test.go @@ -19,7 +19,7 @@ func TestCallWebhook(t *testing.T) { Metadata: &map[string]string{"key1": "val1"}, } - err := CallWebhook("podinfo", "default", hook) + err := CallWebhook("podinfo", "default", flaggerv1.CanaryProgressing, hook) if err != nil { t.Fatal(err.Error()) } @@ -35,7 +35,7 @@ func TestCallWebhook_StatusCode(t *testing.T) { URL: ts.URL, } - err := CallWebhook("podinfo", "default", hook) + err := CallWebhook("podinfo", "default", flaggerv1.CanaryProgressing, hook) if err == nil { t.Errorf("Got no error wanted %v", http.StatusInternalServerError) } From 19e625d38ed0f6a736b38d56da4d73865e80a5b0 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sat, 13 Apr 2019 20:30:19 +0300 Subject: [PATCH 09/22] Add pre/post rollout webhooks to docs --- docs/gitbook/how-it-works.md | 54 ++++++++++++++++++++++++++---------- 1 file changed, 39 insertions(+), 15 deletions(-) diff --git a/docs/gitbook/how-it-works.md b/docs/gitbook/how-it-works.md index 7940d0d5..598e299a 100644 --- a/docs/gitbook/how-it-works.md +++ b/docs/gitbook/how-it-works.md @@ -2,7 +2,7 @@ [Flagger](https://github.com/weaveworks/flagger) takes a Kubernetes deployment and optionally a horizontal pod autoscaler \(HPA\) and creates a series of objects -\(Kubernetes deployments, ClusterIP services and Istio virtual services\) to drive the canary analysis and promotion. +\(Kubernetes deployments, ClusterIP services and Istio or App Mesh virtual services\) to drive the canary analysis and promotion. ![Flagger Canary Process](https://raw.githubusercontent.com/weaveworks/flagger/master/docs/diagrams/flagger-canary-hpa.png) @@ -268,16 +268,22 @@ Gated canary promotion stages: * check primary and canary deployments status * halt advancement if a rolling update is underway * halt advancement if pods are unhealthy +* call pre-rollout webhooks are check results + * halt advancement if any hook returned a non HTTP 2xx result + * increment the failed checks counter * increase canary traffic weight percentage from 0% to 5% (step weight) -* call webhooks and check results +* call rollout webhooks and check results * check canary HTTP request success rate and latency * halt advancement if any metric is under the specified threshold * increment the failed checks counter * check if the number of failed checks reached the threshold * route all traffic to primary * scale to zero the canary deployment and mark it as failed + * call post-rollout webhooks + * post the analysis result to Slack * wait for the canary deployment to be updated and start over * increase canary traffic weight by 5% (step weight) till it reaches 50% (max weight) + * halt advancement if any webhook call fails * halt advancement while canary request success rate is under the threshold * halt advancement while canary request duration P99 is over the threshold * halt advancement if the primary or canary deployment becomes unhealthy @@ -290,6 +296,8 @@ Gated canary promotion stages: * route all traffic to primary * scale to zero the canary deployment * mark rollout as finished +* call post-rollout webhooks +* post the analysis result to Slack * wait for the canary deployment to be updated and start over ### Canary Analysis @@ -524,39 +532,55 @@ rate reaches the 5% threshold, then the canary fails. When specifying a query, Flagger will run the promql query and convert the result to float64. Then it compares the query result value with the metric threshold value. - ### Webhooks -The canary analysis can be extended with webhooks. -Flagger will call each webhook URL and determine from the response status code (HTTP 2xx) if the canary is failing or not. +The canary analysis can be extended with webhooks. Flagger will call each webhook URL and +determine from the response status code (HTTP 2xx) if the canary is failing or not. + +There are three types of hooks: +* Pre-rollout hooks are executed before routing traffic to canary. +The canary advancement is paused if a pre-rollout hook fails and if the number of failures reach the +threshold the canary will be rollback. +* Rollout hooks are executed during the analysis on each iteration before the metric checks. +If a rollout hook call fails the canary advancement is paused and eventfully rolled back. +* Post-rollout hooks are executed after the canary has been promoted or rolled back. +If a post rollout hook fails the error is logged. Spec: ```yaml canaryAnalysis: webhooks: - - name: integration-test - url: http://int-runner.test:8080/ - timeout: 30s - metadata: - test: "all" - token: "16688eb5e9f289f1991c" - - name: db-test + - name: "smoke test" + type: pre-rollout url: http://migration-check.db/query timeout: 30s metadata: key1: "val1" key2: "val2" + - name: "load test" + type: rollout + url: http://flagger-loadtester.test/ + timeout: 15s + metadata: + cmd: "hey -z 1m -q 5 -c 2 http://podinfo-canary.test:9898/" + - name: "notify" + type: post-rollout + url: http://telegram.bot:8080/ + timeout: 5s + metadata: + some: "message" ``` -> **Note** that the sum of all webhooks timeouts should be lower than the control loop interval. +> **Note** that the sum of all rollout webhooks timeouts should be lower than the analysis interval. Webhook payload (HTTP POST): ```json { "name": "podinfo", - "namespace": "test", + "namespace": "test", + "phase": "Progressing", "metadata": { "test": "all", "token": "16688eb5e9f289f1991c" @@ -676,4 +700,4 @@ webhooks: ``` When the canary analysis starts, the load tester will initiate a [clone_and_start request](https://github.com/naver/ngrinder/wiki/REST-API-PerfTest) to the nGrinder server and start a new performance test. the load tester will periodically poll the nGrinder server -for the status of the test, and prevent duplicate requests from being sent in subsequent analysis loops. \ No newline at end of file +for the status of the test, and prevent duplicate requests from being sent in subsequent analysis loops. From 663fa08cc16b4ab42a56983fa3fa63ca199042f3 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sat, 13 Apr 2019 21:21:51 +0300 Subject: [PATCH 10/22] Add hook type and status to CRD schema validation --- artifacts/flagger/crd.yaml | 35 ++++++++++++++++++++++++++++--- charts/flagger/templates/crd.yaml | 27 ++++++++++++++++++++++++ 2 files changed, 59 insertions(+), 3 deletions(-) diff --git a/artifacts/flagger/crd.yaml b/artifacts/flagger/crd.yaml index c8817b80..134020e6 100644 --- a/artifacts/flagger/crd.yaml +++ b/artifacts/flagger/crd.yaml @@ -2,6 +2,8 @@ apiVersion: apiextensions.k8s.io/v1beta1 kind: CustomResourceDefinition metadata: name: canaries.flagger.app + annotations: + helm.sh/resource-policy: keep spec: group: flagger.app version: v1alpha3 @@ -39,9 +41,9 @@ spec: properties: spec: required: - - targetRef - - service - - canaryAnalysis + - targetRef + - service + - canaryAnalysis properties: progressDeadlineSeconds: type: number @@ -119,9 +121,36 @@ spec: properties: name: type: string + type: + type: string + enum: + - "" + - pre-rollout + - rollout + - post-rollout url: type: string format: url timeout: type: string pattern: "^[0-9]+(m|s)" + status: + properties: + phase: + type: string + enum: + - "" + - Initialized + - Progressing + - Succeeded + - Failed + canaryWeight: + type: number + failedChecks: + type: number + iterations: + type: number + lastAppliedSpec: + type: string + lastTransitionTime: + type: string diff --git a/charts/flagger/templates/crd.yaml b/charts/flagger/templates/crd.yaml index 1302acee..aad996b8 100644 --- a/charts/flagger/templates/crd.yaml +++ b/charts/flagger/templates/crd.yaml @@ -122,10 +122,37 @@ spec: properties: name: type: string + type: + type: string + enum: + - "" + - pre-rollout + - rollout + - post-rollout url: type: string format: url timeout: type: string pattern: "^[0-9]+(m|s)" + status: + properties: + phase: + type: string + enum: + - "" + - Initialized + - Progressing + - Succeeded + - Failed + canaryWeight: + type: number + failedChecks: + type: number + iterations: + type: number + lastAppliedSpec: + type: string + lastTransitionTime: + type: string {{- end }} From f46882c778f32baf10da81fa544ed468c8701977 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sun, 14 Apr 2019 12:24:35 +0300 Subject: [PATCH 11/22] Update cert-manager to v0.7 in GKE docs --- .../gitbook/install/flagger-install-on-google-cloud.md | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/docs/gitbook/install/flagger-install-on-google-cloud.md b/docs/gitbook/install/flagger-install-on-google-cloud.md index 94a2f7b7..25dfa373 100644 --- a/docs/gitbook/install/flagger-install-on-google-cloud.md +++ b/docs/gitbook/install/flagger-install-on-google-cloud.md @@ -186,7 +186,7 @@ Install cert-manager's CRDs: ```bash CERT_REPO=https://raw.githubusercontent.com/jetstack/cert-manager -kubectl apply -f ${CERT_REPO}/release-0.6/deploy/manifests/00-crds.yaml +kubectl apply -f ${CERT_REPO}/release-0.7/deploy/manifests/00-crds.yaml ``` Create the cert-manager namespace and disable resource validation: @@ -200,10 +200,12 @@ kubectl label namespace cert-manager certmanager.k8s.io/disable-validation=true Install cert-manager with Helm: ```bash -helm repo update && helm upgrade -i cert-manager \ +helm repo add jetstack https://charts.jetstack.io && \ +helm repo update && \ +helm upgrade -i cert-manager \ --namespace cert-manager \ ---version v0.6.0 \ -stable/cert-manager +--version v0.7.0 \ +jetstack/cert-manager ``` ### Istio Gateway TLS setup From a09dc2cbd8ace3756c262cfaa2f66674e791e2bf Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Mon, 15 Apr 2019 11:25:45 +0300 Subject: [PATCH 12/22] Rename logging package --- cmd/flagger/main.go | 4 ++-- cmd/loadtester/main.go | 4 ++-- pkg/loadtester/runner_test.go | 4 ++-- pkg/loadtester/task_ngrinder_test.go | 4 ++-- pkg/{logging => logger}/logger.go | 14 +------------- pkg/router/router_test.go | 4 ++-- 6 files changed, 11 insertions(+), 23 deletions(-) rename pkg/{logging => logger}/logger.go (88%) diff --git a/cmd/flagger/main.go b/cmd/flagger/main.go index ed2fa659..e0d702f4 100644 --- a/cmd/flagger/main.go +++ b/cmd/flagger/main.go @@ -6,7 +6,7 @@ import ( clientset "github.com/weaveworks/flagger/pkg/client/clientset/versioned" 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/logger" "github.com/weaveworks/flagger/pkg/metrics" "github.com/weaveworks/flagger/pkg/notifier" "github.com/weaveworks/flagger/pkg/server" @@ -58,7 +58,7 @@ func init() { func main() { flag.Parse() - logger, err := logging.NewLoggerWithEncoding(logLevel, zapEncoding) + logger, err := logger.NewLoggerWithEncoding(logLevel, zapEncoding) if err != nil { log.Fatalf("Error creating logger: %v", err) } diff --git a/cmd/loadtester/main.go b/cmd/loadtester/main.go index dfdafdbc..a3e98d3a 100644 --- a/cmd/loadtester/main.go +++ b/cmd/loadtester/main.go @@ -3,7 +3,7 @@ package main import ( "flag" "github.com/weaveworks/flagger/pkg/loadtester" - "github.com/weaveworks/flagger/pkg/logging" + "github.com/weaveworks/flagger/pkg/logger" "github.com/weaveworks/flagger/pkg/signals" "go.uber.org/zap" "log" @@ -30,7 +30,7 @@ func init() { func main() { flag.Parse() - logger, err := logging.NewLoggerWithEncoding(logLevel, zapEncoding) + logger, err := logger.NewLoggerWithEncoding(logLevel, zapEncoding) if err != nil { log.Fatalf("Error creating logger: %v", err) } diff --git a/pkg/loadtester/runner_test.go b/pkg/loadtester/runner_test.go index 1b43509b..1c7ab9c7 100644 --- a/pkg/loadtester/runner_test.go +++ b/pkg/loadtester/runner_test.go @@ -1,14 +1,14 @@ package loadtester import ( - "github.com/weaveworks/flagger/pkg/logging" + "github.com/weaveworks/flagger/pkg/logger" "testing" "time" ) func TestTaskRunner_Start(t *testing.T) { stop := make(chan struct{}) - logger, _ := logging.NewLogger("debug") + logger, _ := logger.NewLogger("debug") tr := NewTaskRunner(logger, time.Hour) go tr.Start(10*time.Millisecond, stop) diff --git a/pkg/loadtester/task_ngrinder_test.go b/pkg/loadtester/task_ngrinder_test.go index 3402ba88..699aea72 100644 --- a/pkg/loadtester/task_ngrinder_test.go +++ b/pkg/loadtester/task_ngrinder_test.go @@ -3,7 +3,7 @@ package loadtester import ( "context" "fmt" - "github.com/weaveworks/flagger/pkg/logging" + "github.com/weaveworks/flagger/pkg/logger" "gopkg.in/h2non/gock.v1" "testing" "time" @@ -12,7 +12,7 @@ import ( func TestTaskNGrinder(t *testing.T) { server := "http://ngrinder:8080" cloneId := "960" - logger, _ := logging.NewLoggerWithEncoding("debug", "console") + logger, _ := logger.NewLoggerWithEncoding("debug", "console") canary := "podinfo.default" taskFactory, ok := GetTaskFactory(TaskTypeNGrinder) if !ok { diff --git a/pkg/logging/logger.go b/pkg/logger/logger.go similarity index 88% rename from pkg/logging/logger.go rename to pkg/logger/logger.go index 8c1838e5..87ba8d76 100644 --- a/pkg/logging/logger.go +++ b/pkg/logger/logger.go @@ -1,9 +1,6 @@ -package logging +package logger import ( - "fmt" - "os" - "go.uber.org/zap" "go.uber.org/zap/zapcore" ) @@ -64,12 +61,3 @@ func NewLoggerWithEncoding(logLevel, zapEncoding string) (*zap.SugaredLogger, er } return logger.Sugar(), nil } - -// Console writes to stdout if the console env var exists -func Console(a ...interface{}) (n int, err error) { - if os.Getenv("console") != "" { - return fmt.Fprintln(os.Stdout, a...) - } - - return 0, nil -} diff --git a/pkg/router/router_test.go b/pkg/router/router_test.go index 08f77cf9..32c5f770 100644 --- a/pkg/router/router_test.go +++ b/pkg/router/router_test.go @@ -6,7 +6,7 @@ import ( istiov1alpha3 "github.com/weaveworks/flagger/pkg/apis/istio/v1alpha3" clientset "github.com/weaveworks/flagger/pkg/client/clientset/versioned" fakeFlagger "github.com/weaveworks/flagger/pkg/client/clientset/versioned/fake" - "github.com/weaveworks/flagger/pkg/logging" + "github.com/weaveworks/flagger/pkg/logger" "go.uber.org/zap" appsv1 "k8s.io/api/apps/v1" hpav1 "k8s.io/api/autoscaling/v1" @@ -35,7 +35,7 @@ func setupfakeClients() fakeClients { kubeClient := fake.NewSimpleClientset(newMockDeployment(), newMockABTestDeployment()) meshClient := fakeFlagger.NewSimpleClientset() - logger, _ := logging.NewLogger("debug") + logger, _ := logger.NewLogger("debug") return fakeClients{ canary: canary, From 60f51ad7d50a0725f2eb02307b52be6497b85487 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Mon, 15 Apr 2019 11:27:08 +0300 Subject: [PATCH 13/22] Move deployer and config tracker to canary package --- pkg/canary/deployer.go | 338 ++++++++++++ pkg/{controller => canary}/deployer_test.go | 37 +- pkg/canary/mock.go | 470 ++++++++++++++++ pkg/canary/ready.go | 112 ++++ pkg/canary/status.go | 104 ++++ pkg/{controller => canary}/tracker.go | 46 +- pkg/controller/controller.go | 19 +- pkg/controller/controller_test.go | 31 +- pkg/controller/deployer.go | 566 -------------------- pkg/controller/scheduler.go | 30 +- 10 files changed, 1117 insertions(+), 636 deletions(-) create mode 100644 pkg/canary/deployer.go rename pkg/{controller => canary}/deployer_test.go (91%) create mode 100644 pkg/canary/mock.go create mode 100644 pkg/canary/ready.go create mode 100644 pkg/canary/status.go rename pkg/{controller => canary}/tracker.go (88%) delete mode 100644 pkg/controller/deployer.go diff --git a/pkg/canary/deployer.go b/pkg/canary/deployer.go new file mode 100644 index 00000000..937dff30 --- /dev/null +++ b/pkg/canary/deployer.go @@ -0,0 +1,338 @@ +package canary + +import ( + "crypto/rand" + "encoding/base64" + "encoding/json" + "fmt" + "github.com/google/go-cmp/cmp" + "github.com/google/go-cmp/cmp/cmpopts" + flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1alpha3" + clientset "github.com/weaveworks/flagger/pkg/client/clientset/versioned" + "go.uber.org/zap" + "io" + appsv1 "k8s.io/api/apps/v1" + hpav1 "k8s.io/api/autoscaling/v2beta1" + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/api/resource" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/client-go/kubernetes" +) + +// Deployer is managing the operations for Kubernetes deployment kind +type Deployer struct { + KubeClient kubernetes.Interface + FlaggerClient clientset.Interface + Logger *zap.SugaredLogger + ConfigTracker ConfigTracker +} + +// Initialize creates the primary deployment and hpa +// and scales to zero the canary deployment +func (c *Deployer) Initialize(cd *flaggerv1.Canary) error { + primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) + if err := c.createPrimaryDeployment(cd); err != nil { + return fmt.Errorf("creating deployment %s.%s failed: %v", primaryName, cd.Namespace, err) + } + + if cd.Status.Phase == "" { + c.Logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("Scaling down %s.%s", cd.Spec.TargetRef.Name, cd.Namespace) + if err := c.Scale(cd, 0); err != nil { + return err + } + } + + if cd.Spec.AutoscalerRef != nil && cd.Spec.AutoscalerRef.Kind == "HorizontalPodAutoscaler" { + if err := c.createPrimaryHpa(cd); err != nil { + return fmt.Errorf("creating hpa %s.%s failed: %v", primaryName, cd.Namespace, err) + } + } + return nil +} + +// Promote copies the pod spec, secrets and config maps from canary to primary +func (c *Deployer) Promote(cd *flaggerv1.Canary) error { + targetName := cd.Spec.TargetRef.Name + primaryName := fmt.Sprintf("%s-primary", targetName) + + canary, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) + } + return fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) + } + + primary, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(primaryName, metav1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return fmt.Errorf("deployment %s.%s not found", primaryName, cd.Namespace) + } + return fmt.Errorf("deployment %s.%s query error %v", primaryName, cd.Namespace, err) + } + + // promote secrets and config maps + configRefs, err := c.ConfigTracker.GetTargetConfigs(cd) + if err != nil { + return err + } + if err := c.ConfigTracker.CreatePrimaryConfigs(cd, configRefs); err != nil { + return err + } + + primaryCopy := primary.DeepCopy() + primaryCopy.Spec.ProgressDeadlineSeconds = canary.Spec.ProgressDeadlineSeconds + primaryCopy.Spec.MinReadySeconds = canary.Spec.MinReadySeconds + primaryCopy.Spec.RevisionHistoryLimit = canary.Spec.RevisionHistoryLimit + primaryCopy.Spec.Strategy = canary.Spec.Strategy + + // update spec with primary secrets and config maps + primaryCopy.Spec.Template.Spec = c.ConfigTracker.ApplyPrimaryConfigs(canary.Spec.Template.Spec, configRefs) + + // update pod annotations to ensure a rolling update + annotations, err := c.makeAnnotations(canary.Spec.Template.Annotations) + if err != nil { + return err + } + primaryCopy.Spec.Template.Annotations = annotations + + primaryCopy.Spec.Template.Labels = makePrimaryLabels(canary.Spec.Template.Labels, primaryName) + + _, err = c.KubeClient.AppsV1().Deployments(cd.Namespace).Update(primaryCopy) + if err != nil { + return fmt.Errorf("updating deployment %s.%s template spec failed: %v", + primaryCopy.GetName(), primaryCopy.Namespace, err) + } + + return nil +} + +// HasDeploymentChanged returns true if the canary deployment pod spec has changed +func (c *Deployer) HasDeploymentChanged(cd *flaggerv1.Canary) (bool, error) { + targetName := cd.Spec.TargetRef.Name + canary, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return false, fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) + } + return false, fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) + } + + if cd.Status.LastAppliedSpec == "" { + return true, nil + } + + newSpec := &canary.Spec.Template.Spec + oldSpecJson, err := base64.StdEncoding.DecodeString(cd.Status.LastAppliedSpec) + if err != nil { + return false, fmt.Errorf("%s.%s decode error %v", cd.Name, cd.Namespace, err) + } + oldSpec := &corev1.PodSpec{} + err = json.Unmarshal(oldSpecJson, oldSpec) + if err != nil { + return false, fmt.Errorf("%s.%s unmarshal error %v", cd.Name, cd.Namespace, err) + } + + if diff := cmp.Diff(*newSpec, *oldSpec, cmpopts.IgnoreUnexported(resource.Quantity{})); diff != "" { + //fmt.Println(diff) + return true, nil + } + + return false, nil +} + +// Scale sets the canary deployment replicas +func (c *Deployer) Scale(cd *flaggerv1.Canary, replicas int32) error { + targetName := cd.Spec.TargetRef.Name + dep, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) + } + return fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) + } + + depCopy := dep.DeepCopy() + depCopy.Spec.Replicas = int32p(replicas) + + _, err = c.KubeClient.AppsV1().Deployments(dep.Namespace).Update(depCopy) + if err != nil { + return fmt.Errorf("scaling %s.%s to %v failed: %v", depCopy.GetName(), depCopy.Namespace, replicas, err) + } + return nil +} + +func (c *Deployer) createPrimaryDeployment(cd *flaggerv1.Canary) error { + targetName := cd.Spec.TargetRef.Name + primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) + + canaryDep, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return fmt.Errorf("deployment %s.%s not found, retrying", targetName, cd.Namespace) + } + return err + } + + if appSel, ok := canaryDep.Spec.Selector.MatchLabels["app"]; !ok || appSel != canaryDep.Name { + return fmt.Errorf("invalid label selector! Deployment %s.%s spec.selector.matchLabels must contain selector 'app: %s'", + targetName, cd.Namespace, targetName) + } + + primaryDep, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(primaryName, metav1.GetOptions{}) + if errors.IsNotFound(err) { + // create primary secrets and config maps + configRefs, err := c.ConfigTracker.GetTargetConfigs(cd) + if err != nil { + return err + } + if err := c.ConfigTracker.CreatePrimaryConfigs(cd, configRefs); err != nil { + return err + } + annotations, err := c.makeAnnotations(canaryDep.Spec.Template.Annotations) + if err != nil { + return err + } + + replicas := int32(1) + if canaryDep.Spec.Replicas != nil && *canaryDep.Spec.Replicas > 0 { + replicas = *canaryDep.Spec.Replicas + } + + // create primary deployment + primaryDep = &appsv1.Deployment{ + ObjectMeta: metav1.ObjectMeta{ + Name: primaryName, + Labels: canaryDep.Labels, + Namespace: cd.Namespace, + OwnerReferences: []metav1.OwnerReference{ + *metav1.NewControllerRef(cd, schema.GroupVersionKind{ + Group: flaggerv1.SchemeGroupVersion.Group, + Version: flaggerv1.SchemeGroupVersion.Version, + Kind: flaggerv1.CanaryKind, + }), + }, + }, + Spec: appsv1.DeploymentSpec{ + ProgressDeadlineSeconds: canaryDep.Spec.ProgressDeadlineSeconds, + MinReadySeconds: canaryDep.Spec.MinReadySeconds, + RevisionHistoryLimit: canaryDep.Spec.RevisionHistoryLimit, + Replicas: int32p(replicas), + Strategy: canaryDep.Spec.Strategy, + Selector: &metav1.LabelSelector{ + MatchLabels: map[string]string{ + "app": primaryName, + }, + }, + Template: corev1.PodTemplateSpec{ + ObjectMeta: metav1.ObjectMeta{ + Labels: makePrimaryLabels(canaryDep.Spec.Template.Labels, primaryName), + Annotations: annotations, + }, + // update spec with the primary secrets and config maps + Spec: c.ConfigTracker.ApplyPrimaryConfigs(canaryDep.Spec.Template.Spec, configRefs), + }, + }, + } + + _, err = c.KubeClient.AppsV1().Deployments(cd.Namespace).Create(primaryDep) + if err != nil { + return err + } + + c.Logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("Deployment %s.%s created", primaryDep.GetName(), cd.Namespace) + } + + return nil +} + +func (c *Deployer) createPrimaryHpa(cd *flaggerv1.Canary) error { + primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) + hpa, err := c.KubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Get(cd.Spec.AutoscalerRef.Name, metav1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return fmt.Errorf("HorizontalPodAutoscaler %s.%s not found, retrying", + cd.Spec.AutoscalerRef.Name, cd.Namespace) + } + return err + } + primaryHpaName := fmt.Sprintf("%s-primary", cd.Spec.AutoscalerRef.Name) + primaryHpa, err := c.KubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Get(primaryHpaName, metav1.GetOptions{}) + + if errors.IsNotFound(err) { + primaryHpa = &hpav1.HorizontalPodAutoscaler{ + ObjectMeta: metav1.ObjectMeta{ + Name: primaryHpaName, + Namespace: cd.Namespace, + Labels: hpa.Labels, + OwnerReferences: []metav1.OwnerReference{ + *metav1.NewControllerRef(cd, schema.GroupVersionKind{ + Group: flaggerv1.SchemeGroupVersion.Group, + Version: flaggerv1.SchemeGroupVersion.Version, + Kind: flaggerv1.CanaryKind, + }), + }, + }, + Spec: hpav1.HorizontalPodAutoscalerSpec{ + ScaleTargetRef: hpav1.CrossVersionObjectReference{ + Name: primaryName, + Kind: hpa.Spec.ScaleTargetRef.Kind, + APIVersion: hpa.Spec.ScaleTargetRef.APIVersion, + }, + MinReplicas: hpa.Spec.MinReplicas, + MaxReplicas: hpa.Spec.MaxReplicas, + Metrics: hpa.Spec.Metrics, + }, + } + + _, err = c.KubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Create(primaryHpa) + if err != nil { + return err + } + c.Logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("HorizontalPodAutoscaler %s.%s created", primaryHpa.GetName(), cd.Namespace) + } + + return nil +} + +// makeAnnotations appends an unique ID to annotations map +func (c *Deployer) makeAnnotations(annotations map[string]string) (map[string]string, error) { + idKey := "flagger-id" + res := make(map[string]string) + uuid := make([]byte, 16) + n, err := io.ReadFull(rand.Reader, uuid) + if n != len(uuid) || err != nil { + return res, err + } + uuid[8] = uuid[8]&^0xc0 | 0x80 + uuid[6] = uuid[6]&^0xf0 | 0x40 + id := fmt.Sprintf("%x-%x-%x-%x-%x", uuid[0:4], uuid[4:6], uuid[6:8], uuid[8:10], uuid[10:]) + + for k, v := range annotations { + if k != idKey { + res[k] = v + } + } + res[idKey] = id + + return res, nil +} + +func makePrimaryLabels(labels map[string]string, primaryName string) map[string]string { + idKey := "app" + res := make(map[string]string) + for k, v := range labels { + if k != idKey { + res[k] = v + } + } + res[idKey] = primaryName + + return res +} + +func int32p(i int32) *int32 { + return &i +} diff --git a/pkg/controller/deployer_test.go b/pkg/canary/deployer_test.go similarity index 91% rename from pkg/controller/deployer_test.go rename to pkg/canary/deployer_test.go index 099e0166..fc53558a 100644 --- a/pkg/controller/deployer_test.go +++ b/pkg/canary/deployer_test.go @@ -1,4 +1,4 @@ -package controller +package canary import ( "testing" @@ -8,8 +8,8 @@ import ( ) func TestCanaryDeployer_Sync(t *testing.T) { - mocks := SetupMocks(false) - err := mocks.deployer.Sync(mocks.canary) + mocks := SetupMocks() + err := mocks.deployer.Initialize(mocks.canary) if err != nil { t.Fatal(err.Error()) } @@ -94,8 +94,8 @@ func TestCanaryDeployer_Sync(t *testing.T) { } func TestCanaryDeployer_IsNewSpec(t *testing.T) { - mocks := SetupMocks(false) - err := mocks.deployer.Sync(mocks.canary) + mocks := SetupMocks() + err := mocks.deployer.Initialize(mocks.canary) if err != nil { t.Fatal(err.Error()) } @@ -106,7 +106,7 @@ func TestCanaryDeployer_IsNewSpec(t *testing.T) { t.Fatal(err.Error()) } - isNew, err := mocks.deployer.IsNewSpec(mocks.canary) + isNew, err := mocks.deployer.HasDeploymentChanged(mocks.canary) if err != nil { t.Fatal(err.Error()) } @@ -117,8 +117,8 @@ func TestCanaryDeployer_IsNewSpec(t *testing.T) { } func TestCanaryDeployer_Promote(t *testing.T) { - mocks := SetupMocks(false) - err := mocks.deployer.Sync(mocks.canary) + mocks := SetupMocks() + err := mocks.deployer.Initialize(mocks.canary) if err != nil { t.Fatal(err.Error()) } @@ -162,8 +162,8 @@ func TestCanaryDeployer_Promote(t *testing.T) { } func TestCanaryDeployer_IsReady(t *testing.T) { - mocks := SetupMocks(false) - err := mocks.deployer.Sync(mocks.canary) + mocks := SetupMocks() + err := mocks.deployer.Initialize(mocks.canary) if err != nil { t.Error("Expected primary readiness check to fail") } @@ -180,8 +180,8 @@ func TestCanaryDeployer_IsReady(t *testing.T) { } func TestCanaryDeployer_SetFailedChecks(t *testing.T) { - mocks := SetupMocks(false) - err := mocks.deployer.Sync(mocks.canary) + mocks := SetupMocks() + err := mocks.deployer.Initialize(mocks.canary) if err != nil { t.Fatal(err.Error()) } @@ -202,8 +202,8 @@ func TestCanaryDeployer_SetFailedChecks(t *testing.T) { } func TestCanaryDeployer_SetState(t *testing.T) { - mocks := SetupMocks(false) - err := mocks.deployer.Sync(mocks.canary) + mocks := SetupMocks() + err := mocks.deployer.Initialize(mocks.canary) if err != nil { t.Fatal(err.Error()) } @@ -224,8 +224,8 @@ func TestCanaryDeployer_SetState(t *testing.T) { } func TestCanaryDeployer_SyncStatus(t *testing.T) { - mocks := SetupMocks(false) - err := mocks.deployer.Sync(mocks.canary) + mocks := SetupMocks() + err := mocks.deployer.Initialize(mocks.canary) if err != nil { t.Fatal(err.Error()) } @@ -263,8 +263,8 @@ func TestCanaryDeployer_SyncStatus(t *testing.T) { } func TestCanaryDeployer_Scale(t *testing.T) { - mocks := SetupMocks(false) - err := mocks.deployer.Sync(mocks.canary) + mocks := SetupMocks() + err := mocks.deployer.Initialize(mocks.canary) if err != nil { t.Fatal(err.Error()) } @@ -279,5 +279,4 @@ func TestCanaryDeployer_Scale(t *testing.T) { if *c.Spec.Replicas != 2 { t.Errorf("Got replicas %v wanted %v", *c.Spec.Replicas, 2) } - } diff --git a/pkg/canary/mock.go b/pkg/canary/mock.go new file mode 100644 index 00000000..17a434af --- /dev/null +++ b/pkg/canary/mock.go @@ -0,0 +1,470 @@ +package canary + +import ( + "github.com/weaveworks/flagger/pkg/apis/flagger/v1alpha3" + clientset "github.com/weaveworks/flagger/pkg/client/clientset/versioned" + fakeFlagger "github.com/weaveworks/flagger/pkg/client/clientset/versioned/fake" + "github.com/weaveworks/flagger/pkg/logger" + "go.uber.org/zap" + appsv1 "k8s.io/api/apps/v1" + hpav1 "k8s.io/api/autoscaling/v1" + hpav2 "k8s.io/api/autoscaling/v2beta1" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/kubernetes/fake" +) + +type Mocks struct { + canary *v1alpha3.Canary + kubeClient kubernetes.Interface + flaggerClient clientset.Interface + deployer Deployer + logger *zap.SugaredLogger +} + +func SetupMocks() Mocks { + // init canary + canary := newTestCanary() + flaggerClient := fakeFlagger.NewSimpleClientset(canary) + + // init kube clientset and register mock objects + kubeClient := fake.NewSimpleClientset( + newTestDeployment(), + newTestHPA(), + NewTestConfigMap(), + NewTestConfigMapEnv(), + NewTestConfigMapVol(), + NewTestSecret(), + NewTestSecretEnv(), + NewTestSecretVol(), + ) + + logger, _ := logger.NewLogger("debug") + + deployer := Deployer{ + FlaggerClient: flaggerClient, + KubeClient: kubeClient, + Logger: logger, + ConfigTracker: ConfigTracker{ + Logger: logger, + KubeClient: kubeClient, + FlaggerClient: flaggerClient, + }, + } + + return Mocks{ + canary: canary, + deployer: deployer, + logger: logger, + flaggerClient: flaggerClient, + kubeClient: kubeClient, + } +} + +func NewTestConfigMap() *corev1.ConfigMap { + return &corev1.ConfigMap{ + TypeMeta: metav1.TypeMeta{APIVersion: corev1.SchemeGroupVersion.String()}, + ObjectMeta: metav1.ObjectMeta{ + Namespace: "default", + Name: "podinfo-config-env", + }, + Data: map[string]string{ + "color": "red", + }, + } +} + +func NewTestConfigMapV2() *corev1.ConfigMap { + return &corev1.ConfigMap{ + TypeMeta: metav1.TypeMeta{APIVersion: corev1.SchemeGroupVersion.String()}, + ObjectMeta: metav1.ObjectMeta{ + Namespace: "default", + Name: "podinfo-config-env", + }, + Data: map[string]string{ + "color": "blue", + "output": "console", + }, + } +} + +func NewTestConfigMapEnv() *corev1.ConfigMap { + return &corev1.ConfigMap{ + TypeMeta: metav1.TypeMeta{APIVersion: corev1.SchemeGroupVersion.String()}, + ObjectMeta: metav1.ObjectMeta{ + Namespace: "default", + Name: "podinfo-config-all-env", + }, + Data: map[string]string{ + "color": "red", + }, + } +} + +func NewTestConfigMapVol() *corev1.ConfigMap { + return &corev1.ConfigMap{ + TypeMeta: metav1.TypeMeta{APIVersion: corev1.SchemeGroupVersion.String()}, + ObjectMeta: metav1.ObjectMeta{ + Namespace: "default", + Name: "podinfo-config-vol", + }, + Data: map[string]string{ + "color": "red", + }, + } +} + +func NewTestSecret() *corev1.Secret { + return &corev1.Secret{ + TypeMeta: metav1.TypeMeta{APIVersion: corev1.SchemeGroupVersion.String()}, + ObjectMeta: metav1.ObjectMeta{ + Namespace: "default", + Name: "podinfo-secret-env", + }, + Type: corev1.SecretTypeOpaque, + Data: map[string][]byte{ + "apiKey": []byte("test"), + }, + } +} + +func NewTestSecretV2() *corev1.Secret { + return &corev1.Secret{ + TypeMeta: metav1.TypeMeta{APIVersion: corev1.SchemeGroupVersion.String()}, + ObjectMeta: metav1.ObjectMeta{ + Namespace: "default", + Name: "podinfo-secret-env", + }, + Type: corev1.SecretTypeOpaque, + Data: map[string][]byte{ + "apiKey": []byte("test2"), + }, + } +} + +func NewTestSecretEnv() *corev1.Secret { + return &corev1.Secret{ + TypeMeta: metav1.TypeMeta{APIVersion: corev1.SchemeGroupVersion.String()}, + ObjectMeta: metav1.ObjectMeta{ + Namespace: "default", + Name: "podinfo-secret-all-env", + }, + Type: corev1.SecretTypeOpaque, + Data: map[string][]byte{ + "apiKey": []byte("test"), + }, + } +} + +func NewTestSecretVol() *corev1.Secret { + return &corev1.Secret{ + TypeMeta: metav1.TypeMeta{APIVersion: corev1.SchemeGroupVersion.String()}, + ObjectMeta: metav1.ObjectMeta{ + Namespace: "default", + Name: "podinfo-secret-vol", + }, + Type: corev1.SecretTypeOpaque, + Data: map[string][]byte{ + "apiKey": []byte("test"), + }, + } +} + +func newTestCanary() *v1alpha3.Canary { + cd := &v1alpha3.Canary{ + TypeMeta: metav1.TypeMeta{APIVersion: v1alpha3.SchemeGroupVersion.String()}, + ObjectMeta: metav1.ObjectMeta{ + Namespace: "default", + Name: "podinfo", + }, + Spec: v1alpha3.CanarySpec{ + TargetRef: hpav1.CrossVersionObjectReference{ + Name: "podinfo", + APIVersion: "apps/v1", + Kind: "Deployment", + }, + AutoscalerRef: &hpav1.CrossVersionObjectReference{ + Name: "podinfo", + APIVersion: "autoscaling/v2beta1", + Kind: "HorizontalPodAutoscaler", + }, Service: v1alpha3.CanaryService{ + Port: 9898, + }, CanaryAnalysis: v1alpha3.CanaryAnalysis{ + Threshold: 10, + StepWeight: 10, + MaxWeight: 50, + Metrics: []v1alpha3.CanaryMetric{ + { + Name: "istio_requests_total", + Threshold: 99, + Interval: "1m", + }, + { + Name: "istio_request_duration_seconds_bucket", + Threshold: 500, + Interval: "1m", + }, + }, + }, + }, + } + return cd +} + +func newTestDeployment() *appsv1.Deployment { + d := &appsv1.Deployment{ + TypeMeta: metav1.TypeMeta{APIVersion: appsv1.SchemeGroupVersion.String()}, + ObjectMeta: metav1.ObjectMeta{ + Namespace: "default", + Name: "podinfo", + }, + Spec: appsv1.DeploymentSpec{ + Selector: &metav1.LabelSelector{ + MatchLabels: map[string]string{ + "app": "podinfo", + }, + }, + Template: corev1.PodTemplateSpec{ + ObjectMeta: metav1.ObjectMeta{ + Labels: map[string]string{ + "app": "podinfo", + }, + }, + Spec: corev1.PodSpec{ + Containers: []corev1.Container{ + { + Name: "podinfo", + Image: "quay.io/stefanprodan/podinfo:1.2.0", + Command: []string{ + "./podinfo", + "--port=9898", + }, + Args: nil, + WorkingDir: "", + Ports: []corev1.ContainerPort{ + { + Name: "http", + ContainerPort: 9898, + Protocol: corev1.ProtocolTCP, + }, + }, + Env: []corev1.EnvVar{ + { + Name: "PODINFO_UI_COLOR", + ValueFrom: &corev1.EnvVarSource{ + ConfigMapKeyRef: &corev1.ConfigMapKeySelector{ + LocalObjectReference: corev1.LocalObjectReference{ + Name: "podinfo-config-env", + }, + Key: "color", + }, + }, + }, + { + Name: "API_KEY", + ValueFrom: &corev1.EnvVarSource{ + SecretKeyRef: &corev1.SecretKeySelector{ + LocalObjectReference: corev1.LocalObjectReference{ + Name: "podinfo-secret-env", + }, + Key: "apiKey", + }, + }, + }, + }, + EnvFrom: []corev1.EnvFromSource{ + { + ConfigMapRef: &corev1.ConfigMapEnvSource{ + LocalObjectReference: corev1.LocalObjectReference{ + Name: "podinfo-config-all-env", + }, + }, + }, + { + SecretRef: &corev1.SecretEnvSource{ + LocalObjectReference: corev1.LocalObjectReference{ + Name: "podinfo-secret-all-env", + }, + }, + }, + }, + VolumeMounts: []corev1.VolumeMount{ + { + Name: "config", + MountPath: "/etc/podinfo/config", + ReadOnly: true, + }, + { + Name: "secret", + MountPath: "/etc/podinfo/secret", + ReadOnly: true, + }, + }, + }, + }, + Volumes: []corev1.Volume{ + { + Name: "config", + VolumeSource: corev1.VolumeSource{ + ConfigMap: &corev1.ConfigMapVolumeSource{ + LocalObjectReference: corev1.LocalObjectReference{ + Name: "podinfo-config-vol", + }, + }, + }, + }, + { + Name: "secret", + VolumeSource: corev1.VolumeSource{ + Secret: &corev1.SecretVolumeSource{ + SecretName: "podinfo-secret-vol", + }, + }, + }, + }, + }, + }, + }, + } + + return d +} + +func newTestDeploymentV2() *appsv1.Deployment { + d := &appsv1.Deployment{ + TypeMeta: metav1.TypeMeta{APIVersion: appsv1.SchemeGroupVersion.String()}, + ObjectMeta: metav1.ObjectMeta{ + Namespace: "default", + Name: "podinfo", + }, + Spec: appsv1.DeploymentSpec{ + Selector: &metav1.LabelSelector{ + MatchLabels: map[string]string{ + "app": "podinfo", + }, + }, + Template: corev1.PodTemplateSpec{ + ObjectMeta: metav1.ObjectMeta{ + Labels: map[string]string{ + "app": "podinfo", + }, + }, + Spec: corev1.PodSpec{ + Containers: []corev1.Container{ + { + Name: "podinfo", + Image: "quay.io/stefanprodan/podinfo:1.2.1", + Ports: []corev1.ContainerPort{ + { + Name: "http", + ContainerPort: 9898, + Protocol: corev1.ProtocolTCP, + }, + }, + Command: []string{ + "./podinfo", + "--port=9898", + }, + Env: []corev1.EnvVar{ + { + Name: "PODINFO_UI_COLOR", + ValueFrom: &corev1.EnvVarSource{ + ConfigMapKeyRef: &corev1.ConfigMapKeySelector{ + LocalObjectReference: corev1.LocalObjectReference{ + Name: "podinfo-config-env", + }, + Key: "color", + }, + }, + }, + { + Name: "API_KEY", + ValueFrom: &corev1.EnvVarSource{ + SecretKeyRef: &corev1.SecretKeySelector{ + LocalObjectReference: corev1.LocalObjectReference{ + Name: "podinfo-secret-env", + }, + Key: "apiKey", + }, + }, + }, + }, + EnvFrom: []corev1.EnvFromSource{ + { + ConfigMapRef: &corev1.ConfigMapEnvSource{ + LocalObjectReference: corev1.LocalObjectReference{ + Name: "podinfo-config-all-env", + }, + }, + }, + }, + VolumeMounts: []corev1.VolumeMount{ + { + Name: "config", + MountPath: "/etc/podinfo/config", + ReadOnly: true, + }, + { + Name: "secret", + MountPath: "/etc/podinfo/secret", + ReadOnly: true, + }, + }, + }, + }, + Volumes: []corev1.Volume{ + { + Name: "config", + VolumeSource: corev1.VolumeSource{ + ConfigMap: &corev1.ConfigMapVolumeSource{ + LocalObjectReference: corev1.LocalObjectReference{ + Name: "podinfo-config-vol", + }, + }, + }, + }, + { + Name: "secret", + VolumeSource: corev1.VolumeSource{ + Secret: &corev1.SecretVolumeSource{ + SecretName: "podinfo-secret-vol", + }, + }, + }, + }, + }, + }, + }, + } + + return d +} + +func newTestHPA() *hpav2.HorizontalPodAutoscaler { + h := &hpav2.HorizontalPodAutoscaler{ + TypeMeta: metav1.TypeMeta{APIVersion: hpav2.SchemeGroupVersion.String()}, + ObjectMeta: metav1.ObjectMeta{ + Namespace: "default", + Name: "podinfo", + }, + Spec: hpav2.HorizontalPodAutoscalerSpec{ + ScaleTargetRef: hpav2.CrossVersionObjectReference{ + Name: "podinfo", + APIVersion: "apps/v1", + Kind: "Deployment", + }, + Metrics: []hpav2.MetricSpec{ + { + Type: "Resource", + Resource: &hpav2.ResourceMetricSource{ + Name: "cpu", + TargetAverageUtilization: int32p(99), + }, + }, + }, + }, + } + + return h +} diff --git a/pkg/canary/ready.go b/pkg/canary/ready.go new file mode 100644 index 00000000..c7884ba4 --- /dev/null +++ b/pkg/canary/ready.go @@ -0,0 +1,112 @@ +package canary + +import ( + "fmt" + "time" + + flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1alpha3" + appsv1 "k8s.io/api/apps/v1" + "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +// IsPrimaryReady checks the primary deployment status and returns an error if +// the deployment is in the middle of a rolling update or if the pods are unhealthy +// it will return a non retriable error if the rolling update is stuck +func (c *Deployer) IsPrimaryReady(cd *flaggerv1.Canary) (bool, error) { + primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) + primary, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(primaryName, metav1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return true, fmt.Errorf("deployment %s.%s not found", primaryName, cd.Namespace) + } + return true, fmt.Errorf("deployment %s.%s query error %v", primaryName, cd.Namespace, err) + } + + retriable, err := c.isDeploymentReady(primary, cd.GetProgressDeadlineSeconds()) + if err != nil { + return retriable, fmt.Errorf("Halt advancement %s.%s %s", primaryName, cd.Namespace, err.Error()) + } + + if primary.Spec.Replicas == int32p(0) { + return true, fmt.Errorf("Halt %s.%s advancement primary deployment is scaled to zero", + cd.Name, cd.Namespace) + } + return true, nil +} + +// IsCanaryReady checks the primary deployment status and returns an error if +// the deployment is in the middle of a rolling update or if the pods are unhealthy +// it will return a non retriable error if the rolling update is stuck +func (c *Deployer) IsCanaryReady(cd *flaggerv1.Canary) (bool, error) { + targetName := cd.Spec.TargetRef.Name + canary, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return true, fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) + } + return true, fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) + } + + retriable, err := c.isDeploymentReady(canary, cd.GetProgressDeadlineSeconds()) + if err != nil { + if retriable { + return retriable, fmt.Errorf("Halt advancement %s.%s %s", targetName, cd.Namespace, err.Error()) + } else { + return retriable, fmt.Errorf("deployment does not have minimum availability for more than %vs", + cd.GetProgressDeadlineSeconds()) + } + } + + return true, nil +} + +// isDeploymentReady determines if a deployment is ready by checking the status conditions +// if a deployment has exceeded the progress deadline it returns a non retriable error +func (c *Deployer) isDeploymentReady(deployment *appsv1.Deployment, deadline int) (bool, error) { + retriable := true + if deployment.Generation <= deployment.Status.ObservedGeneration { + progress := c.getDeploymentCondition(deployment.Status, appsv1.DeploymentProgressing) + if progress != nil { + // Determine if the deployment is stuck by checking if there is a minimum replicas unavailable condition + // and if the last update time exceeds the deadline + available := c.getDeploymentCondition(deployment.Status, appsv1.DeploymentAvailable) + if available != nil && available.Status == "False" && available.Reason == "MinimumReplicasUnavailable" { + from := available.LastUpdateTime + delta := time.Duration(deadline) * time.Second + retriable = !from.Add(delta).Before(time.Now()) + } + } + + if progress != nil && progress.Reason == "ProgressDeadlineExceeded" { + return false, fmt.Errorf("deployment %q exceeded its progress deadline", deployment.GetName()) + } else if deployment.Spec.Replicas != nil && deployment.Status.UpdatedReplicas < *deployment.Spec.Replicas { + return retriable, fmt.Errorf("waiting for rollout to finish: %d out of %d new replicas have been updated", + deployment.Status.UpdatedReplicas, *deployment.Spec.Replicas) + } else if deployment.Status.Replicas > deployment.Status.UpdatedReplicas { + return retriable, fmt.Errorf("waiting for rollout to finish: %d old replicas are pending termination", + deployment.Status.Replicas-deployment.Status.UpdatedReplicas) + } else if deployment.Status.AvailableReplicas < deployment.Status.UpdatedReplicas { + return retriable, fmt.Errorf("waiting for rollout to finish: %d of %d updated replicas are available", + deployment.Status.AvailableReplicas, deployment.Status.UpdatedReplicas) + } + + } else { + return true, fmt.Errorf("waiting for rollout to finish: observed deployment generation less then desired generation") + } + + return true, nil +} + +func (c *Deployer) getDeploymentCondition( + status appsv1.DeploymentStatus, + conditionType appsv1.DeploymentConditionType, +) *appsv1.DeploymentCondition { + for i := range status.Conditions { + c := status.Conditions[i] + if c.Type == conditionType { + return &c + } + } + return nil +} diff --git a/pkg/canary/status.go b/pkg/canary/status.go new file mode 100644 index 00000000..e613e570 --- /dev/null +++ b/pkg/canary/status.go @@ -0,0 +1,104 @@ +package canary + +import ( + "encoding/base64" + "encoding/json" + "fmt" + + flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1alpha3" + "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +// SyncStatus encodes the canary pod spec and updates the canary status +func (c *Deployer) SyncStatus(cd *flaggerv1.Canary, status flaggerv1.CanaryStatus) error { + dep, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(cd.Spec.TargetRef.Name, metav1.GetOptions{}) + if err != nil { + if errors.IsNotFound(err) { + return fmt.Errorf("deployment %s.%s not found", cd.Spec.TargetRef.Name, cd.Namespace) + } + return fmt.Errorf("deployment %s.%s query error %v", cd.Spec.TargetRef.Name, cd.Namespace, err) + } + + specJson, err := json.Marshal(dep.Spec.Template.Spec) + if err != nil { + return fmt.Errorf("deployment %s.%s marshal error %v", cd.Spec.TargetRef.Name, cd.Namespace, err) + } + + configs, err := c.ConfigTracker.GetConfigRefs(cd) + if err != nil { + return fmt.Errorf("configs query error %v", err) + } + + cdCopy := cd.DeepCopy() + cdCopy.Status.Phase = status.Phase + cdCopy.Status.CanaryWeight = status.CanaryWeight + cdCopy.Status.FailedChecks = status.FailedChecks + cdCopy.Status.Iterations = status.Iterations + cdCopy.Status.LastAppliedSpec = base64.StdEncoding.EncodeToString(specJson) + cdCopy.Status.LastTransitionTime = metav1.Now() + cdCopy.Status.TrackedConfigs = configs + + cd, err = c.FlaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) + if err != nil { + return fmt.Errorf("canary %s.%s status update error %v", cdCopy.Name, cdCopy.Namespace, err) + } + return nil +} + +// SetStatusFailedChecks updates the canary failed checks counter +func (c *Deployer) SetStatusFailedChecks(cd *flaggerv1.Canary, val int) error { + cdCopy := cd.DeepCopy() + cdCopy.Status.FailedChecks = val + cdCopy.Status.LastTransitionTime = metav1.Now() + + cd, err := c.FlaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) + if err != nil { + return fmt.Errorf("canary %s.%s status update error %v", cdCopy.Name, cdCopy.Namespace, err) + } + return nil +} + +// SetStatusWeight updates the canary status weight value +func (c *Deployer) SetStatusWeight(cd *flaggerv1.Canary, val int) error { + cdCopy := cd.DeepCopy() + cdCopy.Status.CanaryWeight = val + cdCopy.Status.LastTransitionTime = metav1.Now() + + cd, err := c.FlaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) + if err != nil { + return fmt.Errorf("canary %s.%s status update error %v", cdCopy.Name, cdCopy.Namespace, err) + } + return nil +} + +// SetStatusIterations updates the canary status iterations value +func (c *Deployer) SetStatusIterations(cd *flaggerv1.Canary, val int) error { + cdCopy := cd.DeepCopy() + cdCopy.Status.Iterations = val + cdCopy.Status.LastTransitionTime = metav1.Now() + + cd, err := c.FlaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) + if err != nil { + return fmt.Errorf("canary %s.%s status update error %v", cdCopy.Name, cdCopy.Namespace, err) + } + return nil +} + +// SetStatusPhase updates the canary status phase +func (c *Deployer) SetStatusPhase(cd *flaggerv1.Canary, phase flaggerv1.CanaryPhase) error { + cdCopy := cd.DeepCopy() + cdCopy.Status.Phase = phase + cdCopy.Status.LastTransitionTime = metav1.Now() + + if phase != flaggerv1.CanaryProgressing { + cdCopy.Status.CanaryWeight = 0 + cdCopy.Status.Iterations = 0 + } + + cd, err := c.FlaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) + if err != nil { + return fmt.Errorf("canary %s.%s status update error %v", cdCopy.Name, cdCopy.Namespace, err) + } + return nil +} diff --git a/pkg/controller/tracker.go b/pkg/canary/tracker.go similarity index 88% rename from pkg/controller/tracker.go rename to pkg/canary/tracker.go index 5644b875..f5c41327 100644 --- a/pkg/controller/tracker.go +++ b/pkg/canary/tracker.go @@ -1,4 +1,4 @@ -package controller +package canary import ( "crypto/sha256" @@ -16,9 +16,9 @@ import ( // ConfigTracker is managing the operations for Kubernetes ConfigMaps and Secrets type ConfigTracker struct { - kubeClient kubernetes.Interface - flaggerClient clientset.Interface - logger *zap.SugaredLogger + KubeClient kubernetes.Interface + FlaggerClient clientset.Interface + Logger *zap.SugaredLogger } type ConfigRefType string @@ -50,7 +50,7 @@ func checksum(data interface{}) string { // getRefFromConfigMap transforms a Kubernetes ConfigMap into a ConfigRef // and computes the checksum of the ConfigMap data func (ct *ConfigTracker) getRefFromConfigMap(name string, namespace string) (*ConfigRef, error) { - config, err := ct.kubeClient.CoreV1().ConfigMaps(namespace).Get(name, metav1.GetOptions{}) + config, err := ct.KubeClient.CoreV1().ConfigMaps(namespace).Get(name, metav1.GetOptions{}) if err != nil { return nil, err } @@ -65,7 +65,7 @@ func (ct *ConfigTracker) getRefFromConfigMap(name string, namespace string) (*Co // getRefFromConfigMap transforms a Kubernetes Secret into a ConfigRef // and computes the checksum of the Secret data func (ct *ConfigTracker) getRefFromSecret(name string, namespace string) (*ConfigRef, error) { - secret, err := ct.kubeClient.CoreV1().Secrets(namespace).Get(name, metav1.GetOptions{}) + secret, err := ct.KubeClient.CoreV1().Secrets(namespace).Get(name, metav1.GetOptions{}) if err != nil { return nil, err } @@ -75,7 +75,7 @@ func (ct *ConfigTracker) getRefFromSecret(name string, namespace string) (*Confi secret.Type != corev1.SecretTypeBasicAuth && secret.Type != corev1.SecretTypeSSHAuth && secret.Type != corev1.SecretTypeTLS { - ct.logger.Debugf("ignoring secret %s.%s type not supported %v", name, namespace, secret.Type) + ct.Logger.Debugf("ignoring secret %s.%s type not supported %v", name, namespace, secret.Type) return nil, nil } @@ -91,7 +91,7 @@ func (ct *ConfigTracker) getRefFromSecret(name string, namespace string) (*Confi func (ct *ConfigTracker) GetTargetConfigs(cd *flaggerv1.Canary) (map[string]ConfigRef, error) { res := make(map[string]ConfigRef) targetName := cd.Spec.TargetRef.Name - targetDep, err := ct.kubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) + targetDep, err := ct.KubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { if errors.IsNotFound(err) { return res, fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) @@ -104,7 +104,7 @@ func (ct *ConfigTracker) GetTargetConfigs(cd *flaggerv1.Canary) (map[string]Conf if cmv := volume.ConfigMap; cmv != nil { config, err := ct.getRefFromConfigMap(cmv.Name, cd.Namespace) if err != nil { - ct.logger.Errorf("configMap %s.%s query error %v", cmv.Name, cd.Namespace, err) + ct.Logger.Errorf("configMap %s.%s query error %v", cmv.Name, cd.Namespace, err) continue } if config != nil { @@ -115,7 +115,7 @@ func (ct *ConfigTracker) GetTargetConfigs(cd *flaggerv1.Canary) (map[string]Conf if sv := volume.Secret; sv != nil { secret, err := ct.getRefFromSecret(sv.SecretName, cd.Namespace) if err != nil { - ct.logger.Errorf("secret %s.%s query error %v", sv.SecretName, cd.Namespace, err) + ct.Logger.Errorf("secret %s.%s query error %v", sv.SecretName, cd.Namespace, err) continue } if secret != nil { @@ -133,7 +133,7 @@ func (ct *ConfigTracker) GetTargetConfigs(cd *flaggerv1.Canary) (map[string]Conf name := env.ValueFrom.ConfigMapKeyRef.LocalObjectReference.Name config, err := ct.getRefFromConfigMap(name, cd.Namespace) if err != nil { - ct.logger.Errorf("configMap %s.%s query error %v", name, cd.Namespace, err) + ct.Logger.Errorf("configMap %s.%s query error %v", name, cd.Namespace, err) continue } if config != nil { @@ -143,7 +143,7 @@ func (ct *ConfigTracker) GetTargetConfigs(cd *flaggerv1.Canary) (map[string]Conf name := env.ValueFrom.SecretKeyRef.LocalObjectReference.Name secret, err := ct.getRefFromSecret(name, cd.Namespace) if err != nil { - ct.logger.Errorf("secret %s.%s query error %v", name, cd.Namespace, err) + ct.Logger.Errorf("secret %s.%s query error %v", name, cd.Namespace, err) continue } if secret != nil { @@ -159,7 +159,7 @@ func (ct *ConfigTracker) GetTargetConfigs(cd *flaggerv1.Canary) (map[string]Conf name := envFrom.ConfigMapRef.LocalObjectReference.Name config, err := ct.getRefFromConfigMap(name, cd.Namespace) if err != nil { - ct.logger.Errorf("configMap %s.%s query error %v", name, cd.Namespace, err) + ct.Logger.Errorf("configMap %s.%s query error %v", name, cd.Namespace, err) continue } if config != nil { @@ -169,7 +169,7 @@ func (ct *ConfigTracker) GetTargetConfigs(cd *flaggerv1.Canary) (map[string]Conf name := envFrom.SecretRef.LocalObjectReference.Name secret, err := ct.getRefFromSecret(name, cd.Namespace) if err != nil { - ct.logger.Errorf("secret %s.%s query error %v", name, cd.Namespace, err) + ct.Logger.Errorf("secret %s.%s query error %v", name, cd.Namespace, err) continue } if secret != nil { @@ -221,7 +221,7 @@ func (ct *ConfigTracker) HasConfigChanged(cd *flaggerv1.Canary) (bool, error) { for _, cfg := range configs { if trackedConfigs[cfg.GetName()] != cfg.Checksum { - ct.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)). + ct.Logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)). Infof("%s %s has changed", cfg.Type, cfg.Name) return true, nil } @@ -236,7 +236,7 @@ func (ct *ConfigTracker) CreatePrimaryConfigs(cd *flaggerv1.Canary, refs map[str for _, ref := range refs { switch ref.Type { case ConfigRefMap: - config, err := ct.kubeClient.CoreV1().ConfigMaps(cd.Namespace).Get(ref.Name, metav1.GetOptions{}) + config, err := ct.KubeClient.CoreV1().ConfigMaps(cd.Namespace).Get(ref.Name, metav1.GetOptions{}) if err != nil { return err } @@ -258,10 +258,10 @@ func (ct *ConfigTracker) CreatePrimaryConfigs(cd *flaggerv1.Canary, refs map[str } // update or insert primary ConfigMap - _, err = ct.kubeClient.CoreV1().ConfigMaps(cd.Namespace).Update(primaryConfigMap) + _, err = ct.KubeClient.CoreV1().ConfigMaps(cd.Namespace).Update(primaryConfigMap) if err != nil { if errors.IsNotFound(err) { - _, err = ct.kubeClient.CoreV1().ConfigMaps(cd.Namespace).Create(primaryConfigMap) + _, err = ct.KubeClient.CoreV1().ConfigMaps(cd.Namespace).Create(primaryConfigMap) if err != nil { return err } @@ -270,10 +270,10 @@ func (ct *ConfigTracker) CreatePrimaryConfigs(cd *flaggerv1.Canary, refs map[str } } - ct.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)). + ct.Logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)). Infof("ConfigMap %s synced", primaryConfigMap.GetName()) case ConfigRefSecret: - secret, err := ct.kubeClient.CoreV1().Secrets(cd.Namespace).Get(ref.Name, metav1.GetOptions{}) + secret, err := ct.KubeClient.CoreV1().Secrets(cd.Namespace).Get(ref.Name, metav1.GetOptions{}) if err != nil { return err } @@ -296,10 +296,10 @@ func (ct *ConfigTracker) CreatePrimaryConfigs(cd *flaggerv1.Canary, refs map[str } // update or insert primary Secret - _, err = ct.kubeClient.CoreV1().Secrets(cd.Namespace).Update(primarySecret) + _, err = ct.KubeClient.CoreV1().Secrets(cd.Namespace).Update(primarySecret) if err != nil { if errors.IsNotFound(err) { - _, err = ct.kubeClient.CoreV1().Secrets(cd.Namespace).Create(primarySecret) + _, err = ct.KubeClient.CoreV1().Secrets(cd.Namespace).Create(primarySecret) if err != nil { return err } @@ -308,7 +308,7 @@ func (ct *ConfigTracker) CreatePrimaryConfigs(cd *flaggerv1.Canary, refs map[str } } - ct.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)). + ct.Logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)). Infof("Secret %s synced", primarySecret.GetName()) } } diff --git a/pkg/controller/controller.go b/pkg/controller/controller.go index 6fbc2b91..09388431 100644 --- a/pkg/controller/controller.go +++ b/pkg/controller/controller.go @@ -2,6 +2,7 @@ package controller import ( "fmt" + "github.com/weaveworks/flagger/pkg/canary" "github.com/weaveworks/flagger/pkg/metrics" "sync" "time" @@ -41,7 +42,7 @@ type Controller struct { logger *zap.SugaredLogger canaries *sync.Map jobs map[string]CanaryJob - deployer CanaryDeployer + deployer canary.Deployer observer metrics.Observer recorder metrics.Recorder notifier *notifier.Slack @@ -70,14 +71,14 @@ func NewController( eventRecorder := eventBroadcaster.NewRecorder( scheme.Scheme, corev1.EventSource{Component: controllerAgentName}) - deployer := CanaryDeployer{ - logger: logger, - kubeClient: kubeClient, - flaggerClient: flaggerClient, - configTracker: ConfigTracker{ - logger: logger, - kubeClient: kubeClient, - flaggerClient: flaggerClient, + deployer := canary.Deployer{ + Logger: logger, + KubeClient: kubeClient, + FlaggerClient: flaggerClient, + ConfigTracker: canary.ConfigTracker{ + Logger: logger, + KubeClient: kubeClient, + FlaggerClient: flaggerClient, }, } diff --git a/pkg/controller/controller_test.go b/pkg/controller/controller_test.go index dd605d7c..2954d125 100644 --- a/pkg/controller/controller_test.go +++ b/pkg/controller/controller_test.go @@ -4,10 +4,11 @@ import ( "github.com/weaveworks/flagger/pkg/apis/flagger/v1alpha3" istiov1alpha1 "github.com/weaveworks/flagger/pkg/apis/istio/common/v1alpha1" istiov1alpha3 "github.com/weaveworks/flagger/pkg/apis/istio/v1alpha3" + "github.com/weaveworks/flagger/pkg/canary" clientset "github.com/weaveworks/flagger/pkg/client/clientset/versioned" fakeFlagger "github.com/weaveworks/flagger/pkg/client/clientset/versioned/fake" informers "github.com/weaveworks/flagger/pkg/client/informers/externalversions" - "github.com/weaveworks/flagger/pkg/logging" + "github.com/weaveworks/flagger/pkg/logger" "github.com/weaveworks/flagger/pkg/metrics" "github.com/weaveworks/flagger/pkg/router" "go.uber.org/zap" @@ -34,7 +35,7 @@ type Mocks struct { kubeClient kubernetes.Interface meshClient clientset.Interface flaggerClient clientset.Interface - deployer CanaryDeployer + deployer canary.Deployer observer metrics.Observer ctrl *Controller logger *zap.SugaredLogger @@ -43,11 +44,11 @@ type Mocks struct { func SetupMocks(abtest bool) Mocks { // init canary - canary := newTestCanary() + c := newTestCanary() if abtest { - canary = newTestCanaryAB() + c = newTestCanaryAB() } - flaggerClient := fakeFlagger.NewSimpleClientset(canary) + flaggerClient := fakeFlagger.NewSimpleClientset(c) // init kube clientset and register mock objects kubeClient := fake.NewSimpleClientset( @@ -61,17 +62,17 @@ func SetupMocks(abtest bool) Mocks { NewTestSecretVol(), ) - logger, _ := logging.NewLogger("debug") + logger, _ := logger.NewLogger("debug") // init controller helpers - deployer := CanaryDeployer{ - flaggerClient: flaggerClient, - kubeClient: kubeClient, - logger: logger, - configTracker: ConfigTracker{ - logger: logger, - kubeClient: kubeClient, - flaggerClient: flaggerClient, + deployer := canary.Deployer{ + Logger: logger, + KubeClient: kubeClient, + FlaggerClient: flaggerClient, + ConfigTracker: canary.ConfigTracker{ + Logger: logger, + KubeClient: kubeClient, + FlaggerClient: flaggerClient, }, } observer := metrics.NewObserver("fake") @@ -102,7 +103,7 @@ func SetupMocks(abtest bool) Mocks { meshRouter := rf.MeshRouter("istio") return Mocks{ - canary: canary, + canary: c, observer: observer, deployer: deployer, logger: logger, diff --git a/pkg/controller/deployer.go b/pkg/controller/deployer.go deleted file mode 100644 index 8bd0c80c..00000000 --- a/pkg/controller/deployer.go +++ /dev/null @@ -1,566 +0,0 @@ -package controller - -import ( - "crypto/rand" - "encoding/base64" - "encoding/json" - "fmt" - "io" - "time" - - "github.com/google/go-cmp/cmp" - "github.com/google/go-cmp/cmp/cmpopts" - flaggerv1 "github.com/weaveworks/flagger/pkg/apis/flagger/v1alpha3" - clientset "github.com/weaveworks/flagger/pkg/client/clientset/versioned" - "go.uber.org/zap" - appsv1 "k8s.io/api/apps/v1" - hpav1 "k8s.io/api/autoscaling/v2beta1" - corev1 "k8s.io/api/core/v1" - "k8s.io/apimachinery/pkg/api/errors" - "k8s.io/apimachinery/pkg/api/resource" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/runtime/schema" - "k8s.io/client-go/kubernetes" -) - -// CanaryDeployer is managing the operations for Kubernetes deployment kind -type CanaryDeployer struct { - kubeClient kubernetes.Interface - flaggerClient clientset.Interface - logger *zap.SugaredLogger - configTracker ConfigTracker -} - -// Promote copies the pod spec, secrets and config maps from canary to primary -func (c *CanaryDeployer) Promote(cd *flaggerv1.Canary) error { - targetName := cd.Spec.TargetRef.Name - primaryName := fmt.Sprintf("%s-primary", targetName) - - canary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) - if err != nil { - if errors.IsNotFound(err) { - return fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) - } - return fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) - } - - primary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(primaryName, metav1.GetOptions{}) - if err != nil { - if errors.IsNotFound(err) { - return fmt.Errorf("deployment %s.%s not found", primaryName, cd.Namespace) - } - return fmt.Errorf("deployment %s.%s query error %v", primaryName, cd.Namespace, err) - } - - // promote secrets and config maps - configRefs, err := c.configTracker.GetTargetConfigs(cd) - if err != nil { - return err - } - if err := c.configTracker.CreatePrimaryConfigs(cd, configRefs); err != nil { - return err - } - - primaryCopy := primary.DeepCopy() - primaryCopy.Spec.ProgressDeadlineSeconds = canary.Spec.ProgressDeadlineSeconds - primaryCopy.Spec.MinReadySeconds = canary.Spec.MinReadySeconds - primaryCopy.Spec.RevisionHistoryLimit = canary.Spec.RevisionHistoryLimit - primaryCopy.Spec.Strategy = canary.Spec.Strategy - - // update spec with primary secrets and config maps - primaryCopy.Spec.Template.Spec = c.configTracker.ApplyPrimaryConfigs(canary.Spec.Template.Spec, configRefs) - - // update pod annotations to ensure a rolling update - annotations, err := c.makeAnnotations(canary.Spec.Template.Annotations) - if err != nil { - return err - } - primaryCopy.Spec.Template.Annotations = annotations - - primaryCopy.Spec.Template.Labels = makePrimaryLabels(canary.Spec.Template.Labels, primaryName) - - _, err = c.kubeClient.AppsV1().Deployments(cd.Namespace).Update(primaryCopy) - if err != nil { - return fmt.Errorf("updating deployment %s.%s template spec failed: %v", - primaryCopy.GetName(), primaryCopy.Namespace, err) - } - - return nil -} - -// IsPrimaryReady checks the primary deployment status and returns an error if -// the deployment is in the middle of a rolling update or if the pods are unhealthy -// it will return a non retriable error if the rolling update is stuck -func (c *CanaryDeployer) IsPrimaryReady(cd *flaggerv1.Canary) (bool, error) { - primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) - primary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(primaryName, metav1.GetOptions{}) - if err != nil { - if errors.IsNotFound(err) { - return true, fmt.Errorf("deployment %s.%s not found", primaryName, cd.Namespace) - } - return true, fmt.Errorf("deployment %s.%s query error %v", primaryName, cd.Namespace, err) - } - - retriable, err := c.isDeploymentReady(primary, cd.GetProgressDeadlineSeconds()) - if err != nil { - return retriable, fmt.Errorf("Halt advancement %s.%s %s", primaryName, cd.Namespace, err.Error()) - } - - if primary.Spec.Replicas == int32p(0) { - return true, fmt.Errorf("Halt %s.%s advancement primary deployment is scaled to zero", - cd.Name, cd.Namespace) - } - return true, nil -} - -// IsCanaryReady checks the primary deployment status and returns an error if -// the deployment is in the middle of a rolling update or if the pods are unhealthy -// it will return a non retriable error if the rolling update is stuck -func (c *CanaryDeployer) IsCanaryReady(cd *flaggerv1.Canary) (bool, error) { - targetName := cd.Spec.TargetRef.Name - canary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) - if err != nil { - if errors.IsNotFound(err) { - return true, fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) - } - return true, fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) - } - - retriable, err := c.isDeploymentReady(canary, cd.GetProgressDeadlineSeconds()) - if err != nil { - if retriable { - return retriable, fmt.Errorf("Halt advancement %s.%s %s", targetName, cd.Namespace, err.Error()) - } else { - return retriable, fmt.Errorf("deployment does not have minimum availability for more than %vs", - cd.GetProgressDeadlineSeconds()) - } - } - - return true, nil -} - -// IsNewSpec returns true if the canary deployment pod spec has changed -func (c *CanaryDeployer) IsNewSpec(cd *flaggerv1.Canary) (bool, error) { - targetName := cd.Spec.TargetRef.Name - canary, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) - if err != nil { - if errors.IsNotFound(err) { - return false, fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) - } - return false, fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) - } - - if cd.Status.LastAppliedSpec == "" { - return true, nil - } - - newSpec := &canary.Spec.Template.Spec - oldSpecJson, err := base64.StdEncoding.DecodeString(cd.Status.LastAppliedSpec) - if err != nil { - return false, fmt.Errorf("%s.%s decode error %v", cd.Name, cd.Namespace, err) - } - oldSpec := &corev1.PodSpec{} - err = json.Unmarshal(oldSpecJson, oldSpec) - if err != nil { - return false, fmt.Errorf("%s.%s unmarshal error %v", cd.Name, cd.Namespace, err) - } - - if diff := cmp.Diff(*newSpec, *oldSpec, cmpopts.IgnoreUnexported(resource.Quantity{})); diff != "" { - //fmt.Println(diff) - return true, nil - } - - return false, nil -} - -// ShouldAdvance determines if the canary analysis can proceed -func (c *CanaryDeployer) ShouldAdvance(cd *flaggerv1.Canary) (bool, error) { - if cd.Status.LastAppliedSpec == "" || cd.Status.Phase == flaggerv1.CanaryProgressing { - return true, nil - } - - newDep, err := c.IsNewSpec(cd) - if err != nil { - return false, err - } - if newDep { - return newDep, nil - } - - newCfg, err := c.configTracker.HasConfigChanged(cd) - if err != nil { - return false, err - } - - return newCfg, nil - -} - -// SetStatusFailedChecks updates the canary failed checks counter -func (c *CanaryDeployer) SetStatusFailedChecks(cd *flaggerv1.Canary, val int) error { - cdCopy := cd.DeepCopy() - cdCopy.Status.FailedChecks = val - cdCopy.Status.LastTransitionTime = metav1.Now() - - cd, err := c.flaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) - if err != nil { - return fmt.Errorf("canary %s.%s status update error %v", cdCopy.Name, cdCopy.Namespace, err) - } - return nil -} - -// SetStatusWeight updates the canary status weight value -func (c *CanaryDeployer) SetStatusWeight(cd *flaggerv1.Canary, val int) error { - cdCopy := cd.DeepCopy() - cdCopy.Status.CanaryWeight = val - cdCopy.Status.LastTransitionTime = metav1.Now() - - cd, err := c.flaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) - if err != nil { - return fmt.Errorf("canary %s.%s status update error %v", cdCopy.Name, cdCopy.Namespace, err) - } - return nil -} - -// SetStatusIterations updates the canary status iterations value -func (c *CanaryDeployer) SetStatusIterations(cd *flaggerv1.Canary, val int) error { - cdCopy := cd.DeepCopy() - cdCopy.Status.Iterations = val - cdCopy.Status.LastTransitionTime = metav1.Now() - - cd, err := c.flaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) - if err != nil { - return fmt.Errorf("canary %s.%s status update error %v", cdCopy.Name, cdCopy.Namespace, err) - } - return nil -} - -// SetStatusWeight updates the canary status weight value -func (c *CanaryDeployer) IncrementStatusIterations(cd *flaggerv1.Canary) error { - cdCopy := cd.DeepCopy() - cdCopy.Status.Iterations = cdCopy.Status.Iterations + 1 - cdCopy.Status.LastTransitionTime = metav1.Now() - - cd, err := c.flaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) - if err != nil { - return fmt.Errorf("canary %s.%s status update error %v", cdCopy.Name, cdCopy.Namespace, err) - } - return nil -} - -// SetStatusPhase updates the canary status phase -func (c *CanaryDeployer) SetStatusPhase(cd *flaggerv1.Canary, phase flaggerv1.CanaryPhase) error { - cdCopy := cd.DeepCopy() - cdCopy.Status.Phase = phase - cdCopy.Status.LastTransitionTime = metav1.Now() - - if phase != flaggerv1.CanaryProgressing { - cdCopy.Status.CanaryWeight = 0 - cdCopy.Status.Iterations = 0 - } - - cd, err := c.flaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) - if err != nil { - return fmt.Errorf("canary %s.%s status update error %v", cdCopy.Name, cdCopy.Namespace, err) - } - return nil -} - -// SyncStatus encodes the canary pod spec and updates the canary status -func (c *CanaryDeployer) SyncStatus(cd *flaggerv1.Canary, status flaggerv1.CanaryStatus) error { - dep, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(cd.Spec.TargetRef.Name, metav1.GetOptions{}) - if err != nil { - if errors.IsNotFound(err) { - return fmt.Errorf("deployment %s.%s not found", cd.Spec.TargetRef.Name, cd.Namespace) - } - return fmt.Errorf("deployment %s.%s query error %v", cd.Spec.TargetRef.Name, cd.Namespace, err) - } - - specJson, err := json.Marshal(dep.Spec.Template.Spec) - if err != nil { - return fmt.Errorf("deployment %s.%s marshal error %v", cd.Spec.TargetRef.Name, cd.Namespace, err) - } - - configs, err := c.configTracker.GetConfigRefs(cd) - if err != nil { - return fmt.Errorf("configs query error %v", err) - } - - cdCopy := cd.DeepCopy() - cdCopy.Status.Phase = status.Phase - cdCopy.Status.CanaryWeight = status.CanaryWeight - cdCopy.Status.FailedChecks = status.FailedChecks - cdCopy.Status.Iterations = status.Iterations - cdCopy.Status.LastAppliedSpec = base64.StdEncoding.EncodeToString(specJson) - cdCopy.Status.LastTransitionTime = metav1.Now() - cdCopy.Status.TrackedConfigs = configs - - cd, err = c.flaggerClient.FlaggerV1alpha3().Canaries(cd.Namespace).UpdateStatus(cdCopy) - if err != nil { - return fmt.Errorf("canary %s.%s status update error %v", cdCopy.Name, cdCopy.Namespace, err) - } - return nil -} - -// Scale sets the canary deployment replicas -func (c *CanaryDeployer) Scale(cd *flaggerv1.Canary, replicas int32) error { - targetName := cd.Spec.TargetRef.Name - dep, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) - if err != nil { - if errors.IsNotFound(err) { - return fmt.Errorf("deployment %s.%s not found", targetName, cd.Namespace) - } - return fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) - } - - depCopy := dep.DeepCopy() - depCopy.Spec.Replicas = int32p(replicas) - - _, err = c.kubeClient.AppsV1().Deployments(dep.Namespace).Update(depCopy) - if err != nil { - return fmt.Errorf("scaling %s.%s to %v failed: %v", depCopy.GetName(), depCopy.Namespace, replicas, err) - } - return nil -} - -// Reconcile creates the primary deployment and hpa -// and scales to zero the canary deployment -func (c *CanaryDeployer) Sync(cd *flaggerv1.Canary) error { - primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) - if err := c.createPrimaryDeployment(cd); err != nil { - return fmt.Errorf("creating deployment %s.%s failed: %v", primaryName, cd.Namespace, err) - } - - if cd.Status.Phase == "" { - c.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("Scaling down %s.%s", cd.Spec.TargetRef.Name, cd.Namespace) - if err := c.Scale(cd, 0); err != nil { - return err - } - } - - if cd.Spec.AutoscalerRef != nil && cd.Spec.AutoscalerRef.Kind == "HorizontalPodAutoscaler" { - if err := c.createPrimaryHpa(cd); err != nil { - return fmt.Errorf("creating hpa %s.%s failed: %v", primaryName, cd.Namespace, err) - } - } - return nil -} - -func (c *CanaryDeployer) createPrimaryDeployment(cd *flaggerv1.Canary) error { - targetName := cd.Spec.TargetRef.Name - primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) - - canaryDep, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) - if err != nil { - if errors.IsNotFound(err) { - return fmt.Errorf("deployment %s.%s not found, retrying", targetName, cd.Namespace) - } - return err - } - - if appSel, ok := canaryDep.Spec.Selector.MatchLabels["app"]; !ok || appSel != canaryDep.Name { - return fmt.Errorf("invalid label selector! Deployment %s.%s spec.selector.matchLabels must contain selector 'app: %s'", - targetName, cd.Namespace, targetName) - } - - primaryDep, err := c.kubeClient.AppsV1().Deployments(cd.Namespace).Get(primaryName, metav1.GetOptions{}) - if errors.IsNotFound(err) { - // create primary secrets and config maps - configRefs, err := c.configTracker.GetTargetConfigs(cd) - if err != nil { - return err - } - if err := c.configTracker.CreatePrimaryConfigs(cd, configRefs); err != nil { - return err - } - annotations, err := c.makeAnnotations(canaryDep.Spec.Template.Annotations) - if err != nil { - return err - } - - replicas := int32(1) - if canaryDep.Spec.Replicas != nil && *canaryDep.Spec.Replicas > 0 { - replicas = *canaryDep.Spec.Replicas - } - - // create primary deployment - primaryDep = &appsv1.Deployment{ - ObjectMeta: metav1.ObjectMeta{ - Name: primaryName, - Labels: canaryDep.Labels, - Namespace: cd.Namespace, - OwnerReferences: []metav1.OwnerReference{ - *metav1.NewControllerRef(cd, schema.GroupVersionKind{ - Group: flaggerv1.SchemeGroupVersion.Group, - Version: flaggerv1.SchemeGroupVersion.Version, - Kind: flaggerv1.CanaryKind, - }), - }, - }, - Spec: appsv1.DeploymentSpec{ - ProgressDeadlineSeconds: canaryDep.Spec.ProgressDeadlineSeconds, - MinReadySeconds: canaryDep.Spec.MinReadySeconds, - RevisionHistoryLimit: canaryDep.Spec.RevisionHistoryLimit, - Replicas: int32p(replicas), - Strategy: canaryDep.Spec.Strategy, - Selector: &metav1.LabelSelector{ - MatchLabels: map[string]string{ - "app": primaryName, - }, - }, - Template: corev1.PodTemplateSpec{ - ObjectMeta: metav1.ObjectMeta{ - Labels: makePrimaryLabels(canaryDep.Spec.Template.Labels, primaryName), - Annotations: annotations, - }, - // update spec with the primary secrets and config maps - Spec: c.configTracker.ApplyPrimaryConfigs(canaryDep.Spec.Template.Spec, configRefs), - }, - }, - } - - _, err = c.kubeClient.AppsV1().Deployments(cd.Namespace).Create(primaryDep) - if err != nil { - return err - } - - c.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("Deployment %s.%s created", primaryDep.GetName(), cd.Namespace) - } - - return nil -} - -func (c *CanaryDeployer) createPrimaryHpa(cd *flaggerv1.Canary) error { - primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) - hpa, err := c.kubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Get(cd.Spec.AutoscalerRef.Name, metav1.GetOptions{}) - if err != nil { - if errors.IsNotFound(err) { - return fmt.Errorf("HorizontalPodAutoscaler %s.%s not found, retrying", - cd.Spec.AutoscalerRef.Name, cd.Namespace) - } - return err - } - primaryHpaName := fmt.Sprintf("%s-primary", cd.Spec.AutoscalerRef.Name) - primaryHpa, err := c.kubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Get(primaryHpaName, metav1.GetOptions{}) - - if errors.IsNotFound(err) { - primaryHpa = &hpav1.HorizontalPodAutoscaler{ - ObjectMeta: metav1.ObjectMeta{ - Name: primaryHpaName, - Namespace: cd.Namespace, - Labels: hpa.Labels, - OwnerReferences: []metav1.OwnerReference{ - *metav1.NewControllerRef(cd, schema.GroupVersionKind{ - Group: flaggerv1.SchemeGroupVersion.Group, - Version: flaggerv1.SchemeGroupVersion.Version, - Kind: flaggerv1.CanaryKind, - }), - }, - }, - Spec: hpav1.HorizontalPodAutoscalerSpec{ - ScaleTargetRef: hpav1.CrossVersionObjectReference{ - Name: primaryName, - Kind: hpa.Spec.ScaleTargetRef.Kind, - APIVersion: hpa.Spec.ScaleTargetRef.APIVersion, - }, - MinReplicas: hpa.Spec.MinReplicas, - MaxReplicas: hpa.Spec.MaxReplicas, - Metrics: hpa.Spec.Metrics, - }, - } - - _, err = c.kubeClient.AutoscalingV2beta1().HorizontalPodAutoscalers(cd.Namespace).Create(primaryHpa) - if err != nil { - return err - } - c.logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("HorizontalPodAutoscaler %s.%s created", primaryHpa.GetName(), cd.Namespace) - } - - return nil -} - -// isDeploymentReady determines if a deployment is ready by checking the status conditions -// if a deployment has exceeded the progress deadline it returns a non retriable error -func (c *CanaryDeployer) isDeploymentReady(deployment *appsv1.Deployment, deadline int) (bool, error) { - retriable := true - if deployment.Generation <= deployment.Status.ObservedGeneration { - progress := c.getDeploymentCondition(deployment.Status, appsv1.DeploymentProgressing) - if progress != nil { - // Determine if the deployment is stuck by checking if there is a minimum replicas unavailable condition - // and if the last update time exceeds the deadline - available := c.getDeploymentCondition(deployment.Status, appsv1.DeploymentAvailable) - if available != nil && available.Status == "False" && available.Reason == "MinimumReplicasUnavailable" { - from := available.LastUpdateTime - delta := time.Duration(deadline) * time.Second - retriable = !from.Add(delta).Before(time.Now()) - } - } - - if progress != nil && progress.Reason == "ProgressDeadlineExceeded" { - return false, fmt.Errorf("deployment %q exceeded its progress deadline", deployment.GetName()) - } else if deployment.Spec.Replicas != nil && deployment.Status.UpdatedReplicas < *deployment.Spec.Replicas { - return retriable, fmt.Errorf("waiting for rollout to finish: %d out of %d new replicas have been updated", - deployment.Status.UpdatedReplicas, *deployment.Spec.Replicas) - } else if deployment.Status.Replicas > deployment.Status.UpdatedReplicas { - return retriable, fmt.Errorf("waiting for rollout to finish: %d old replicas are pending termination", - deployment.Status.Replicas-deployment.Status.UpdatedReplicas) - } else if deployment.Status.AvailableReplicas < deployment.Status.UpdatedReplicas { - return retriable, fmt.Errorf("waiting for rollout to finish: %d of %d updated replicas are available", - deployment.Status.AvailableReplicas, deployment.Status.UpdatedReplicas) - } - - } else { - return true, fmt.Errorf("waiting for rollout to finish: observed deployment generation less then desired generation") - } - - return true, nil -} - -func (c *CanaryDeployer) getDeploymentCondition( - status appsv1.DeploymentStatus, - conditionType appsv1.DeploymentConditionType, -) *appsv1.DeploymentCondition { - for i := range status.Conditions { - c := status.Conditions[i] - if c.Type == conditionType { - return &c - } - } - return nil -} - -// makeAnnotations appends an unique ID to annotations map -func (c *CanaryDeployer) makeAnnotations(annotations map[string]string) (map[string]string, error) { - idKey := "flagger-id" - res := make(map[string]string) - uuid := make([]byte, 16) - n, err := io.ReadFull(rand.Reader, uuid) - if n != len(uuid) || err != nil { - return res, err - } - uuid[8] = uuid[8]&^0xc0 | 0x80 - uuid[6] = uuid[6]&^0xf0 | 0x40 - id := fmt.Sprintf("%x-%x-%x-%x-%x", uuid[0:4], uuid[4:6], uuid[6:8], uuid[8:10], uuid[10:]) - - for k, v := range annotations { - if k != idKey { - res[k] = v - } - } - res[idKey] = id - - return res, nil -} - -func makePrimaryLabels(labels map[string]string, primaryName string) map[string]string { - idKey := "app" - res := make(map[string]string) - for k, v := range labels { - if k != idKey { - res[k] = v - } - } - res[idKey] = primaryName - - return res -} diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index 9af0b354..b2602bab 100644 --- a/pkg/controller/scheduler.go +++ b/pkg/controller/scheduler.go @@ -90,7 +90,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) // create primary deployment and hpa if needed - if err := c.deployer.Sync(cd); err != nil { + if err := c.deployer.Initialize(cd); err != nil { c.recordEventWarningf(cd, "%v", err) return } @@ -111,7 +111,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh return } - shouldAdvance, err := c.deployer.ShouldAdvance(cd) + shouldAdvance, err := c.shouldAdvance(cd) if err != nil { c.recordEventWarningf(cd, "%v", err) return @@ -441,6 +441,28 @@ func (c *Controller) shouldSkipAnalysis(cd *flaggerv1.Canary, meshRouter router. return true } +func (c *Controller) shouldAdvance(cd *flaggerv1.Canary) (bool, error) { + if cd.Status.LastAppliedSpec == "" || cd.Status.Phase == flaggerv1.CanaryProgressing { + return true, nil + } + + newDep, err := c.deployer.HasDeploymentChanged(cd) + if err != nil { + return false, err + } + if newDep { + return newDep, nil + } + + newCfg, err := c.deployer.ConfigTracker.HasConfigChanged(cd) + if err != nil { + return false, err + } + + return newCfg, nil + +} + func (c *Controller) checkCanaryStatus(cd *flaggerv1.Canary, shouldAdvance bool) bool { c.recorder.SetStatus(cd, cd.Status.Phase) if cd.Status.Phase == flaggerv1.CanaryProgressing { @@ -479,10 +501,10 @@ func (c *Controller) checkCanaryStatus(cd *flaggerv1.Canary, shouldAdvance bool) func (c *Controller) hasCanaryRevisionChanged(cd *flaggerv1.Canary) bool { if cd.Status.Phase == flaggerv1.CanaryProgressing { - if diff, _ := c.deployer.IsNewSpec(cd); diff { + if diff, _ := c.deployer.HasDeploymentChanged(cd); diff { return true } - if diff, _ := c.deployer.configTracker.HasConfigChanged(cd); diff { + if diff, _ := c.deployer.ConfigTracker.HasConfigChanged(cd); diff { return true } } From 6ef72e25509cfae8ef3eb2ca95343c84bb8e0d49 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Mon, 15 Apr 2019 12:57:25 +0300 Subject: [PATCH 14/22] Make the pod selector configurable - default labels: app, name and app.kubernetes.io/name --- cmd/flagger/main.go | 9 ++++ pkg/canary/deployer.go | 69 ++++++++++++++++++++----------- pkg/canary/deployer_test.go | 16 +++---- pkg/canary/mock.go | 9 ++-- pkg/controller/controller.go | 2 + pkg/controller/controller_test.go | 1 + pkg/controller/scheduler.go | 5 ++- pkg/router/factory.go | 3 +- pkg/router/kubernetes.go | 3 +- 9 files changed, 76 insertions(+), 41 deletions(-) diff --git a/cmd/flagger/main.go b/cmd/flagger/main.go index e0d702f4..3281235f 100644 --- a/cmd/flagger/main.go +++ b/cmd/flagger/main.go @@ -18,6 +18,7 @@ import ( "k8s.io/client-go/tools/cache" "k8s.io/client-go/tools/clientcmd" "log" + "strings" "time" ) @@ -36,6 +37,7 @@ var ( zapEncoding string namespace string meshProvider string + selectorLabels string ) func init() { @@ -53,6 +55,7 @@ func init() { flag.StringVar(&zapEncoding, "zap-encoding", "json", "Zap logger encoding.") flag.StringVar(&namespace, "namespace", "", "Namespace that flagger would watch canary object") flag.StringVar(&meshProvider, "mesh-provider", "istio", "Service mesh provider, can be istio or appmesh") + flag.StringVar(&selectorLabels, "selector-labels", "app,name,app.kubernetes.io/name", "List of pod labels that Flagger uses to create pod selectors") } func main() { @@ -101,6 +104,11 @@ func main() { logger.Fatalf("Error calling Kubernetes API: %v", err) } + labels := strings.Split(selectorLabels, ",") + if len(labels) < 1 { + logger.Fatalf("At least one selector label is required") + } + logger.Infof("Connected to Kubernetes API %s", ver) if namespace != "" { logger.Infof("Watching namespace %s", namespace) @@ -137,6 +145,7 @@ func main() { slack, meshProvider, version.VERSION, + labels, ) flaggerInformerFactory.Start(stopCh) diff --git a/pkg/canary/deployer.go b/pkg/canary/deployer.go index 937dff30..bb8aa121 100644 --- a/pkg/canary/deployer.go +++ b/pkg/canary/deployer.go @@ -27,29 +27,31 @@ type Deployer struct { FlaggerClient clientset.Interface Logger *zap.SugaredLogger ConfigTracker ConfigTracker + Labels []string } -// Initialize creates the primary deployment and hpa -// and scales to zero the canary deployment -func (c *Deployer) Initialize(cd *flaggerv1.Canary) error { +// Initialize creates the primary deployment, hpa, +// scales to zero the canary deployment and returns the pod selector label +func (c *Deployer) Initialize(cd *flaggerv1.Canary) (string, error) { primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) - if err := c.createPrimaryDeployment(cd); err != nil { - return fmt.Errorf("creating deployment %s.%s failed: %v", primaryName, cd.Namespace, err) + label, err := c.createPrimaryDeployment(cd) + if err != nil { + return "", fmt.Errorf("creating deployment %s.%s failed: %v", primaryName, cd.Namespace, err) } if cd.Status.Phase == "" { c.Logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("Scaling down %s.%s", cd.Spec.TargetRef.Name, cd.Namespace) if err := c.Scale(cd, 0); err != nil { - return err + return "", err } } if cd.Spec.AutoscalerRef != nil && cd.Spec.AutoscalerRef.Kind == "HorizontalPodAutoscaler" { if err := c.createPrimaryHpa(cd); err != nil { - return fmt.Errorf("creating hpa %s.%s failed: %v", primaryName, cd.Namespace, err) + return "", fmt.Errorf("creating hpa %s.%s failed: %v", primaryName, cd.Namespace, err) } } - return nil + return label, nil } // Promote copies the pod spec, secrets and config maps from canary to primary @@ -65,6 +67,12 @@ func (c *Deployer) Promote(cd *flaggerv1.Canary) error { return fmt.Errorf("deployment %s.%s query error %v", targetName, cd.Namespace, err) } + label, err := c.getSelectorLabel(canary) + if err != nil { + return fmt.Errorf("invalid label selector! Deployment %s.%s spec.selector.matchLabels must contain selector 'app: %s'", + targetName, cd.Namespace, targetName) + } + primary, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(primaryName, metav1.GetOptions{}) if err != nil { if errors.IsNotFound(err) { @@ -98,7 +106,7 @@ func (c *Deployer) Promote(cd *flaggerv1.Canary) error { } primaryCopy.Spec.Template.Annotations = annotations - primaryCopy.Spec.Template.Labels = makePrimaryLabels(canary.Spec.Template.Labels, primaryName) + primaryCopy.Spec.Template.Labels = makePrimaryLabels(canary.Spec.Template.Labels, primaryName, label) _, err = c.KubeClient.AppsV1().Deployments(cd.Namespace).Update(primaryCopy) if err != nil { @@ -164,20 +172,21 @@ func (c *Deployer) Scale(cd *flaggerv1.Canary, replicas int32) error { return nil } -func (c *Deployer) createPrimaryDeployment(cd *flaggerv1.Canary) error { +func (c *Deployer) createPrimaryDeployment(cd *flaggerv1.Canary) (string, error) { targetName := cd.Spec.TargetRef.Name primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) canaryDep, err := c.KubeClient.AppsV1().Deployments(cd.Namespace).Get(targetName, metav1.GetOptions{}) if err != nil { if errors.IsNotFound(err) { - return fmt.Errorf("deployment %s.%s not found, retrying", targetName, cd.Namespace) + return "", fmt.Errorf("deployment %s.%s not found, retrying", targetName, cd.Namespace) } - return err + return "", err } - if appSel, ok := canaryDep.Spec.Selector.MatchLabels["app"]; !ok || appSel != canaryDep.Name { - return fmt.Errorf("invalid label selector! Deployment %s.%s spec.selector.matchLabels must contain selector 'app: %s'", + label, err := c.getSelectorLabel(canaryDep) + if err != nil { + return "", fmt.Errorf("invalid label selector! Deployment %s.%s spec.selector.matchLabels must contain selector 'app: %s'", targetName, cd.Namespace, targetName) } @@ -186,14 +195,14 @@ func (c *Deployer) createPrimaryDeployment(cd *flaggerv1.Canary) error { // create primary secrets and config maps configRefs, err := c.ConfigTracker.GetTargetConfigs(cd) if err != nil { - return err + return "", err } if err := c.ConfigTracker.CreatePrimaryConfigs(cd, configRefs); err != nil { - return err + return "", err } annotations, err := c.makeAnnotations(canaryDep.Spec.Template.Annotations) if err != nil { - return err + return "", err } replicas := int32(1) @@ -223,12 +232,12 @@ func (c *Deployer) createPrimaryDeployment(cd *flaggerv1.Canary) error { Strategy: canaryDep.Spec.Strategy, Selector: &metav1.LabelSelector{ MatchLabels: map[string]string{ - "app": primaryName, + label: primaryName, }, }, Template: corev1.PodTemplateSpec{ ObjectMeta: metav1.ObjectMeta{ - Labels: makePrimaryLabels(canaryDep.Spec.Template.Labels, primaryName), + Labels: makePrimaryLabels(canaryDep.Spec.Template.Labels, primaryName, label), Annotations: annotations, }, // update spec with the primary secrets and config maps @@ -239,13 +248,13 @@ func (c *Deployer) createPrimaryDeployment(cd *flaggerv1.Canary) error { _, err = c.KubeClient.AppsV1().Deployments(cd.Namespace).Create(primaryDep) if err != nil { - return err + return "", err } c.Logger.With("canary", fmt.Sprintf("%s.%s", cd.Name, cd.Namespace)).Infof("Deployment %s.%s created", primaryDep.GetName(), cd.Namespace) } - return nil + return label, nil } func (c *Deployer) createPrimaryHpa(cd *flaggerv1.Canary) error { @@ -320,15 +329,25 @@ func (c *Deployer) makeAnnotations(annotations map[string]string) (map[string]st return res, nil } -func makePrimaryLabels(labels map[string]string, primaryName string) map[string]string { - idKey := "app" +// getSelectorLabel returns the selector match label +func (c *Deployer) getSelectorLabel(deployment *appsv1.Deployment) (string, error) { + for _, l := range c.Labels { + if _, ok := deployment.Spec.Selector.MatchLabels[l]; ok { + return l, nil + } + } + + return "", fmt.Errorf("selector not found") +} + +func makePrimaryLabels(labels map[string]string, primaryName string, label string) map[string]string { res := make(map[string]string) for k, v := range labels { - if k != idKey { + if k != label { res[k] = v } } - res[idKey] = primaryName + res[label] = primaryName return res } diff --git a/pkg/canary/deployer_test.go b/pkg/canary/deployer_test.go index fc53558a..d978c78c 100644 --- a/pkg/canary/deployer_test.go +++ b/pkg/canary/deployer_test.go @@ -9,7 +9,7 @@ import ( func TestCanaryDeployer_Sync(t *testing.T) { mocks := SetupMocks() - err := mocks.deployer.Initialize(mocks.canary) + _, err := mocks.deployer.Initialize(mocks.canary) if err != nil { t.Fatal(err.Error()) } @@ -95,7 +95,7 @@ func TestCanaryDeployer_Sync(t *testing.T) { func TestCanaryDeployer_IsNewSpec(t *testing.T) { mocks := SetupMocks() - err := mocks.deployer.Initialize(mocks.canary) + _, err := mocks.deployer.Initialize(mocks.canary) if err != nil { t.Fatal(err.Error()) } @@ -118,7 +118,7 @@ func TestCanaryDeployer_IsNewSpec(t *testing.T) { func TestCanaryDeployer_Promote(t *testing.T) { mocks := SetupMocks() - err := mocks.deployer.Initialize(mocks.canary) + _, err := mocks.deployer.Initialize(mocks.canary) if err != nil { t.Fatal(err.Error()) } @@ -163,7 +163,7 @@ func TestCanaryDeployer_Promote(t *testing.T) { func TestCanaryDeployer_IsReady(t *testing.T) { mocks := SetupMocks() - err := mocks.deployer.Initialize(mocks.canary) + _, err := mocks.deployer.Initialize(mocks.canary) if err != nil { t.Error("Expected primary readiness check to fail") } @@ -181,7 +181,7 @@ func TestCanaryDeployer_IsReady(t *testing.T) { func TestCanaryDeployer_SetFailedChecks(t *testing.T) { mocks := SetupMocks() - err := mocks.deployer.Initialize(mocks.canary) + _, err := mocks.deployer.Initialize(mocks.canary) if err != nil { t.Fatal(err.Error()) } @@ -203,7 +203,7 @@ func TestCanaryDeployer_SetFailedChecks(t *testing.T) { func TestCanaryDeployer_SetState(t *testing.T) { mocks := SetupMocks() - err := mocks.deployer.Initialize(mocks.canary) + _, err := mocks.deployer.Initialize(mocks.canary) if err != nil { t.Fatal(err.Error()) } @@ -225,7 +225,7 @@ func TestCanaryDeployer_SetState(t *testing.T) { func TestCanaryDeployer_SyncStatus(t *testing.T) { mocks := SetupMocks() - err := mocks.deployer.Initialize(mocks.canary) + _, err := mocks.deployer.Initialize(mocks.canary) if err != nil { t.Fatal(err.Error()) } @@ -264,7 +264,7 @@ func TestCanaryDeployer_SyncStatus(t *testing.T) { func TestCanaryDeployer_Scale(t *testing.T) { mocks := SetupMocks() - err := mocks.deployer.Initialize(mocks.canary) + _, err := mocks.deployer.Initialize(mocks.canary) if err != nil { t.Fatal(err.Error()) } diff --git a/pkg/canary/mock.go b/pkg/canary/mock.go index 17a434af..f60444e8 100644 --- a/pkg/canary/mock.go +++ b/pkg/canary/mock.go @@ -46,6 +46,7 @@ func SetupMocks() Mocks { FlaggerClient: flaggerClient, KubeClient: kubeClient, Logger: logger, + Labels: []string{"app", "name"}, ConfigTracker: ConfigTracker{ Logger: logger, KubeClient: kubeClient, @@ -222,13 +223,13 @@ func newTestDeployment() *appsv1.Deployment { Spec: appsv1.DeploymentSpec{ Selector: &metav1.LabelSelector{ MatchLabels: map[string]string{ - "app": "podinfo", + "name": "podinfo", }, }, Template: corev1.PodTemplateSpec{ ObjectMeta: metav1.ObjectMeta{ Labels: map[string]string{ - "app": "podinfo", + "name": "podinfo", }, }, Spec: corev1.PodSpec{ @@ -341,13 +342,13 @@ func newTestDeploymentV2() *appsv1.Deployment { Spec: appsv1.DeploymentSpec{ Selector: &metav1.LabelSelector{ MatchLabels: map[string]string{ - "app": "podinfo", + "name": "podinfo", }, }, Template: corev1.PodTemplateSpec{ ObjectMeta: metav1.ObjectMeta{ Labels: map[string]string{ - "app": "podinfo", + "name": "podinfo", }, }, Spec: corev1.PodSpec{ diff --git a/pkg/controller/controller.go b/pkg/controller/controller.go index 09388431..45c3d938 100644 --- a/pkg/controller/controller.go +++ b/pkg/controller/controller.go @@ -60,6 +60,7 @@ func NewController( notifier *notifier.Slack, meshProvider string, version string, + labels []string, ) *Controller { logger.Debug("Creating event broadcaster") flaggerscheme.AddToScheme(scheme.Scheme) @@ -75,6 +76,7 @@ func NewController( Logger: logger, KubeClient: kubeClient, FlaggerClient: flaggerClient, + Labels: labels, ConfigTracker: canary.ConfigTracker{ Logger: logger, KubeClient: kubeClient, diff --git a/pkg/controller/controller_test.go b/pkg/controller/controller_test.go index 2954d125..8946773a 100644 --- a/pkg/controller/controller_test.go +++ b/pkg/controller/controller_test.go @@ -69,6 +69,7 @@ func SetupMocks(abtest bool) Mocks { Logger: logger, KubeClient: kubeClient, FlaggerClient: flaggerClient, + Labels: []string{"app", "name"}, ConfigTracker: canary.ConfigTracker{ Logger: logger, KubeClient: kubeClient, diff --git a/pkg/controller/scheduler.go b/pkg/controller/scheduler.go index b2602bab..85cdb59f 100644 --- a/pkg/controller/scheduler.go +++ b/pkg/controller/scheduler.go @@ -90,7 +90,8 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh primaryName := fmt.Sprintf("%s-primary", cd.Spec.TargetRef.Name) // create primary deployment and hpa if needed - if err := c.deployer.Initialize(cd); err != nil { + label, err := c.deployer.Initialize(cd) + if err != nil { c.recordEventWarningf(cd, "%v", err) return } @@ -100,7 +101,7 @@ func (c *Controller) advanceCanary(name string, namespace string, skipLivenessCh meshRouter := routerFactory.MeshRouter(c.meshProvider) // create or update ClusterIP services - if err := routerFactory.KubernetesRouter().Reconcile(cd); err != nil { + if err := routerFactory.KubernetesRouter(label).Reconcile(cd); err != nil { c.recordEventWarningf(cd, "%v", err) return } diff --git a/pkg/router/factory.go b/pkg/router/factory.go index 4a879392..31f5d608 100644 --- a/pkg/router/factory.go +++ b/pkg/router/factory.go @@ -26,11 +26,12 @@ func NewFactory(kubeClient kubernetes.Interface, } // KubernetesRouter returns a ClusterIP service router -func (factory *Factory) KubernetesRouter() *KubernetesRouter { +func (factory *Factory) KubernetesRouter(label string) *KubernetesRouter { return &KubernetesRouter{ logger: factory.logger, flaggerClient: factory.flaggerClient, kubeClient: factory.kubeClient, + label: label, } } diff --git a/pkg/router/kubernetes.go b/pkg/router/kubernetes.go index 38d6165b..2d38dfe6 100644 --- a/pkg/router/kubernetes.go +++ b/pkg/router/kubernetes.go @@ -19,6 +19,7 @@ type KubernetesRouter struct { kubeClient kubernetes.Interface flaggerClient clientset.Interface logger *zap.SugaredLogger + label string } // Reconcile creates or updates the primary and canary services @@ -64,7 +65,7 @@ func (c *KubernetesRouter) reconcileService(canary *flaggerv1.Canary, name strin svcSpec := corev1.ServiceSpec{ Type: corev1.ServiceTypeClusterIP, - Selector: map[string]string{"app": target}, + Selector: map[string]string{c.label: target}, Ports: []corev1.ServicePort{ { Name: portName, From 65f716182bb8aff68df3d995bdfb066ab67254d3 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Mon, 15 Apr 2019 13:30:58 +0300 Subject: [PATCH 15/22] Add default selectors to docs --- docs/gitbook/how-it-works.md | 3 +++ 1 file changed, 3 insertions(+) diff --git a/docs/gitbook/how-it-works.md b/docs/gitbook/how-it-works.md index 598e299a..2b325e50 100644 --- a/docs/gitbook/how-it-works.md +++ b/docs/gitbook/how-it-works.md @@ -96,6 +96,9 @@ spec: app: podinfo ``` +Besides `app` Flagger supports `name` and `app.kubernetes.io/name` selectors. If you use a different +convention you can specify your label with the `-selector-labels` flag. + The target deployment should expose a TCP port that will be used by Flagger to create the ClusterIP Service and the Istio Virtual Service. The container port from the target deployment should match the `service.port` value. From adb53c63dd975a4df106f60d289ff1e496e333c1 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Tue, 16 Apr 2019 10:26:00 +0300 Subject: [PATCH 16/22] Add e2e tests for A/B testing and pre/post hooks --- test/e2e-tests.sh | 72 ++++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 71 insertions(+), 1 deletion(-) diff --git a/test/e2e-tests.sh b/test/e2e-tests.sh index eed594b6..1b0c52a0 100755 --- a/test/e2e-tests.sh +++ b/test/e2e-tests.sh @@ -83,7 +83,6 @@ spec: type: cmd cmd: "hey -z 10m -q 10 -c 2 http://podinfo.test:9898/" logCmdOutput: "true" - EOF echo '>>> Waiting for primary to be ready' @@ -126,6 +125,77 @@ done echo '✔ Canary promotion test passed' +cat <>> Triggering A/B testing' +kubectl -n test set image deployment/podinfo podinfod=quay.io/stefanprodan/podinfo:1.4.2 + +echo '>>> Waiting for A/B testing promotion' +retries=50 +count=0 +ok=false +until ${ok}; do + kubectl -n test describe deployment/podinfo-primary | grep '1.4.2' && ok=true || ok=false + sleep 10 + kubectl -n istio-system logs deployment/flagger --tail 1 + count=$(($count + 1)) + if [[ ${count} -eq ${retries} ]]; then + kubectl -n test describe deployment/podinfo + kubectl -n test describe deployment/podinfo-primary + kubectl -n istio-system logs deployment/flagger + echo "No more retries left" + exit 1 + fi +done + +echo '✔ A/B testing promotion test passed' + kubectl -n istio-system logs deployment/flagger echo '✔ All tests passed' \ No newline at end of file From 50b7b74480ad177ca7f694a0d2117289e985cd95 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Tue, 16 Apr 2019 10:41:35 +0300 Subject: [PATCH 17/22] Speed up e2e tests by reducing the number of iterations --- test/e2e-tests.sh | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/test/e2e-tests.sh b/test/e2e-tests.sh index 1b0c52a0..2ef1b7bc 100755 --- a/test/e2e-tests.sh +++ b/test/e2e-tests.sh @@ -42,7 +42,7 @@ spec: canaryAnalysis: interval: 15s threshold: 15 - maxWeight: 50 + maxWeight: 30 stepWeight: 10 metrics: - name: request-success-rate @@ -142,7 +142,7 @@ spec: canaryAnalysis: interval: 10s threshold: 5 - iterations: 10 + iterations: 5 match: - headers: cookie: From 15484363d61ab3399cfc483f0f73aed76e4ab62f Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Tue, 16 Apr 2019 11:00:03 +0300 Subject: [PATCH 18/22] Add A/B testing and hooks to e2e readme --- test/README.md | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/test/README.md b/test/README.md index b338d303..3f8d586b 100644 --- a/test/README.md +++ b/test/README.md @@ -5,7 +5,7 @@ The e2e testing infrastructure is powered by CircleCI and [Kubernetes Kind](http CircleCI e2e workflow: * install latest stable kubectl [e2e-kind.sh](e2e-kind.sh) -* build Kubernetes Kind from master [e2e-kind.sh](e2e-kind.sh) +* install Kubernetes Kind [e2e-kind.sh](e2e-kind.sh) * create local Kubernetes cluster with kind [e2e-kind.sh](e2e-kind.sh) * install latest stable Helm CLI [e2e-istio.sh](e2e-istio.sh) * deploy Tiller on the local cluster [e2e-istio.sh](e2e-istio.sh) @@ -18,7 +18,7 @@ CircleCI e2e workflow: * deploy the load tester in the test namespace [e2e-tests.sh](e2e-tests.sh) * deploy a demo workload (podinfo) in the test namespace [e2e-tests.sh](e2e-tests.sh) * test the canary initialization [e2e-tests.sh](e2e-tests.sh) -* test the canary analysis and promotion [e2e-tests.sh](e2e-tests.sh) - +* test the canary analysis and promotion using weighted traffic and the load testing webhook [e2e-tests.sh](e2e-tests.sh) +* test the A/B testing analysis and promotion using cookies filters and pre/post rollout webhooks [e2e-tests.sh](e2e-tests.sh) From 331942a4eddbe8dc450a6256c970c8b8fb96e554 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Wed, 17 Apr 2019 11:03:38 +0300 Subject: [PATCH 19/22] Switch to Docker Hub from Quay --- .travis.yml | 12 ++++++------ Makefile | 4 ++-- 2 files changed, 8 insertions(+), 8 deletions(-) diff --git a/.travis.yml b/.travis.yml index c1535e23..6195cdcc 100644 --- a/.travis.yml +++ b/.travis.yml @@ -28,16 +28,16 @@ after_success: echo "PR build, skipping image push"; else BRANCH_COMMIT=${TRAVIS_BRANCH}-$(echo ${TRAVIS_COMMIT} | head -c7); - docker tag weaveworks/flagger:latest quay.io/weaveworks/flagger:${BRANCH_COMMIT}; - echo $DOCKER_PASS | docker login -u=$DOCKER_USER --password-stdin quay.io; - docker push quay.io/weaveworks/flagger:${BRANCH_COMMIT}; + docker tag weaveworks/flagger:latest weaveworks/flagger:${BRANCH_COMMIT}; + echo $DOCKER_PASS | docker login -u=$DOCKER_USER --password-stdin; + docker push weaveworks/flagger:${BRANCH_COMMIT}; fi - if [ -z "$TRAVIS_TAG" ]; then echo "Not a release, skipping image push"; else - docker tag weaveworks/flagger:latest quay.io/weaveworks/flagger:${TRAVIS_TAG}; - echo $DOCKER_PASS | docker login -u=$DOCKER_USER --password-stdin quay.io; - docker push quay.io/weaveworks/flagger:$TRAVIS_TAG; + docker tag weaveworks/flagger:latest weaveworks/flagger:${TRAVIS_TAG}; + echo $DOCKER_PASS | docker login -u=$DOCKER_USER --password-stdin; + docker push weaveworks/flagger:$TRAVIS_TAG; fi - bash <(curl -s https://codecov.io/bash) - rm coverage.txt diff --git a/Makefile b/Makefile index 109302aa..67b4314e 100644 --- a/Makefile +++ b/Makefile @@ -88,5 +88,5 @@ reset-test: kubectl apply -f ./artifacts/canaries loadtester-push: - docker build -t quay.io/weaveworks/flagger-loadtester:$(LT_VERSION) . -f Dockerfile.loadtester - docker push quay.io/weaveworks/flagger-loadtester:$(LT_VERSION) \ No newline at end of file + docker build -t weaveworks/flagger-loadtester:$(LT_VERSION) . -f Dockerfile.loadtester + docker push weaveworks/flagger-loadtester:$(LT_VERSION) \ No newline at end of file From cd08afcbeb87813ca06074dcf012b8b2e215394c Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Wed, 17 Apr 2019 11:05:10 +0300 Subject: [PATCH 20/22] Add bash bats task runner - run bats tests (blocking requests) --- Dockerfile.loadtester | 26 ++++++-------------------- pkg/loadtester/bats.go | 39 +++++++++++++++++++++++++++++++++++++++ pkg/loadtester/server.go | 21 +++++++++++++++++++++ pkg/loadtester/task.go | 1 + 4 files changed, 67 insertions(+), 20 deletions(-) create mode 100644 pkg/loadtester/bats.go diff --git a/Dockerfile.loadtester b/Dockerfile.loadtester index 238ad5c1..7c6c5e1d 100644 --- a/Dockerfile.loadtester +++ b/Dockerfile.loadtester @@ -1,20 +1,4 @@ -FROM golang:1.12 AS hey-builder - -RUN mkdir -p /go/src/github.com/rakyll/hey/ - -WORKDIR /go/src/github.com/rakyll/hey - -ADD https://github.com/rakyll/hey/archive/v0.1.1.tar.gz . - -RUN tar xzf v0.1.1.tar.gz --strip 1 - -RUN go get ./... - -RUN CGO_ENABLED=0 GOOS=linux GOARCH=amd64 \ - go install -ldflags '-w -extldflags "-static"' \ - /go/src/github.com/rakyll/hey - -FROM golang:1.11 AS builder +FROM golang:1.12 AS builder RUN mkdir -p /go/src/github.com/weaveworks/flagger/ @@ -26,15 +10,17 @@ RUN go test -race ./pkg/loadtester/ RUN CGO_ENABLED=0 GOOS=linux go build -a -installsuffix cgo -o loadtester ./cmd/loadtester/* -FROM alpine:3.9 +FROM bats/bats:v1.1.0 RUN addgroup -S app \ && adduser -S -g app app \ - && apk --no-cache add ca-certificates curl + && apk --no-cache add ca-certificates curl jq WORKDIR /home/app -COPY --from=hey-builder /go/bin/hey /usr/local/bin/hey +RUN curl -sSLo hey "https://storage.googleapis.com/jblabs/dist/hey_linux_v0.1.2" && \ +chmod +x hey && mv hey /usr/local/bin/hey + COPY --from=builder /go/src/github.com/weaveworks/flagger/loadtester . RUN chown -R app:app ./ diff --git a/pkg/loadtester/bats.go b/pkg/loadtester/bats.go new file mode 100644 index 00000000..d27eb4b0 --- /dev/null +++ b/pkg/loadtester/bats.go @@ -0,0 +1,39 @@ +package loadtester + +import ( + "context" + "fmt" + "os/exec" +) + +const TaskTypeBats = "bats" + +type BatsTask struct { + TaskBase + command string + logCmdOutput bool +} + +func (task *BatsTask) Hash() string { + return hash(task.canary + task.command) +} + +func (task *BatsTask) Run(ctx context.Context) (bool, error) { + cmd := exec.CommandContext(ctx, "bash", "-c", task.command) + out, err := cmd.CombinedOutput() + + if err != nil { + task.logger.With("canary", task.canary).Errorf("command failed %s %v %s", task.command, err, out) + return false, fmt.Errorf(" %v %v", err, out) + } else { + if task.logCmdOutput { + fmt.Printf("%s\n", out) + } + task.logger.With("canary", task.canary).Infof("command finished %s", task.command) + } + return true, nil +} + +func (task *BatsTask) String() string { + return task.command +} diff --git a/pkg/loadtester/server.go b/pkg/loadtester/server.go index 533c5554..c5a34aae 100644 --- a/pkg/loadtester/server.go +++ b/pkg/loadtester/server.go @@ -44,6 +44,27 @@ func ListenAndServe(port string, timeout time.Duration, logger *zap.SugaredLogge if !ok { typ = TaskTypeShell } + + // run bats command (blocking task) + if typ == TaskTypeBats { + bats := BatsTask{ + command: payload.Metadata["cmd"], + logCmdOutput: taskRunner.logCmdOutput, + } + + ctx, cancel := context.WithTimeout(context.Background(), taskRunner.timeout) + defer cancel() + + ok, err := bats.Run(ctx) + if !ok { + w.WriteHeader(http.StatusInternalServerError) + w.Write([]byte(err.Error())) + } + + w.WriteHeader(http.StatusOK) + return + } + taskFactory, ok := GetTaskFactory(typ) if !ok { w.WriteHeader(http.StatusBadRequest) diff --git a/pkg/loadtester/task.go b/pkg/loadtester/task.go index 2e30228a..395b34a9 100644 --- a/pkg/loadtester/task.go +++ b/pkg/loadtester/task.go @@ -24,6 +24,7 @@ type TaskBase struct { func (task *TaskBase) Canary() string { return task.canary } + func hash(str string) string { fnvHash := fnv.New32() fnvBytes := fnvHash.Sum([]byte(str)) From a82eb7b01f925d96536e2ac1595d0bfd18763ed8 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Wed, 17 Apr 2019 11:13:31 +0300 Subject: [PATCH 21/22] Release v0.11.0 --- artifacts/flagger/deployment.yaml | 2 +- artifacts/loadtester/deployment.yaml | 2 +- charts/flagger/Chart.yaml | 4 ++-- charts/flagger/values.yaml | 4 ++-- pkg/version/version.go | 2 +- 5 files changed, 7 insertions(+), 7 deletions(-) diff --git a/artifacts/flagger/deployment.yaml b/artifacts/flagger/deployment.yaml index 6597ed70..39f4462b 100644 --- a/artifacts/flagger/deployment.yaml +++ b/artifacts/flagger/deployment.yaml @@ -22,7 +22,7 @@ spec: serviceAccountName: flagger containers: - name: flagger - image: quay.io/weaveworks/flagger:0.10.0 + image: weaveworks/flagger:0.11.0 imagePullPolicy: IfNotPresent ports: - name: http diff --git a/artifacts/loadtester/deployment.yaml b/artifacts/loadtester/deployment.yaml index 02f8fe95..fe77c3a9 100644 --- a/artifacts/loadtester/deployment.yaml +++ b/artifacts/loadtester/deployment.yaml @@ -17,7 +17,7 @@ spec: spec: containers: - name: loadtester - image: quay.io/stefanprodan/flagger-loadtester:0.2.0 + image: weaveworks/flagger-loadtester:0.2.0 imagePullPolicy: IfNotPresent ports: - name: http diff --git a/charts/flagger/Chart.yaml b/charts/flagger/Chart.yaml index 02d6369d..b0f3812d 100644 --- a/charts/flagger/Chart.yaml +++ b/charts/flagger/Chart.yaml @@ -1,7 +1,7 @@ apiVersion: v1 name: flagger -version: 0.10.0 -appVersion: 0.10.0 +version: 0.11.0 +appVersion: 0.11.0 kubeVersion: ">=1.11.0-0" engine: gotpl description: Flagger is a Kubernetes operator that automates the promotion of canary deployments using Istio routing for traffic shifting and Prometheus metrics for canary analysis. diff --git a/charts/flagger/values.yaml b/charts/flagger/values.yaml index 2b5bcbb9..e20d4140 100644 --- a/charts/flagger/values.yaml +++ b/charts/flagger/values.yaml @@ -1,8 +1,8 @@ # Default values for flagger. image: - repository: quay.io/weaveworks/flagger - tag: 0.10.0 + repository: weaveworks/flagger + tag: 0.11.0 pullPolicy: IfNotPresent metricsServer: "http://prometheus:9090" diff --git a/pkg/version/version.go b/pkg/version/version.go index d00cbd89..e2131ffa 100644 --- a/pkg/version/version.go +++ b/pkg/version/version.go @@ -1,4 +1,4 @@ package version -var VERSION = "0.10.0" +var VERSION = "0.11.0" var REVISION = "unknown" From d0b582048f421604b72c3b4997a6c507763a20e9 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Wed, 17 Apr 2019 11:22:02 +0300 Subject: [PATCH 22/22] Add change log for v0.11.0 --- CHANGELOG.md | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index e5328934..7704cd8a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,23 @@ All notable changes to this project are documented in this file. +## 0.11.0 (2019-04-17) + +Adds pre/post rollout [webhooks](https://docs.flagger.app/how-it-works#webhooks) + +#### Features + +- Add `pre-rollout` and `post-rollout` webhook types [#147](https://github.com/weaveworks/flagger/pull/147) + +#### Improvements + +- Unify App Mesh and Istio builtin metric checks [#146](https://github.com/weaveworks/flagger/pull/146) +- Make the pod selector label configurable [#148](https://github.com/weaveworks/flagger/pull/148) + +#### Breaking changes + +- Set default `mesh` Istio gateway only if no gateway is specified [#141](https://github.com/weaveworks/flagger/pull/141) + ## 0.10.0 (2019-03-27) Adds support for App Mesh