mirror of
https://github.com/fluxcd/flagger.git
synced 2026-04-15 06:57:34 +00:00
feat: weighted deployments
This commit is contained in:
@@ -3,6 +3,7 @@ package controller
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
|
||||
@@ -172,6 +173,13 @@ func alertMetadata(canary *flaggerv1.Canary) []notifier.Field {
|
||||
canary.GetAnalysis().StepWeight,
|
||||
canary.GetAnalysis().MaxWeight),
|
||||
})
|
||||
} else if len(canary.GetAnalysis().StepWeights) > 0 {
|
||||
fields = append(fields, notifier.Field{
|
||||
Name: "Traffic routing",
|
||||
Value: fmt.Sprintf("Weight steps: %s max: %v",
|
||||
strings.Trim(strings.Join(strings.Fields(fmt.Sprint(canary.GetAnalysis().StepWeights)), ","), "[]"),
|
||||
canary.GetAnalysis().MaxWeight),
|
||||
})
|
||||
} else if len(canary.GetAnalysis().Match) > 0 {
|
||||
fields = append(fields, notifier.Field{
|
||||
Name: "Traffic routing",
|
||||
|
||||
+53
-24
@@ -14,6 +14,33 @@ import (
|
||||
"github.com/weaveworks/flagger/pkg/router"
|
||||
)
|
||||
|
||||
func (c *Controller) fullWeight(canary *flaggerv1.Canary) int {
|
||||
if canary.GetAnalysis().FullWeight > 0 {
|
||||
return canary.GetAnalysis().FullWeight
|
||||
}
|
||||
// set max weight default value to 100%
|
||||
return 100
|
||||
}
|
||||
|
||||
func (c *Controller) nextStepWeight(canary *flaggerv1.Canary, canaryWeight int) int {
|
||||
var stepWeightsLen = len(canary.GetAnalysis().StepWeights)
|
||||
if canary.GetAnalysis().StepWeight > 0 || stepWeightsLen == 0 {
|
||||
return canary.GetAnalysis().StepWeight
|
||||
}
|
||||
|
||||
if canaryWeight == 0 {
|
||||
return canary.GetAnalysis().StepWeights[0]
|
||||
}
|
||||
|
||||
for i := 0; i < stepWeightsLen-1; i++ {
|
||||
if canary.GetAnalysis().StepWeights[i] == canaryWeight {
|
||||
return canary.GetAnalysis().StepWeights[i+1] - canaryWeight
|
||||
}
|
||||
}
|
||||
|
||||
return c.fullWeight(canary) - canaryWeight
|
||||
}
|
||||
|
||||
// scheduleCanaries synchronises the canary map with the jobs map,
|
||||
// for new canaries new jobs are created and started
|
||||
// for the removed canaries the jobs are stopped and deleted
|
||||
@@ -173,8 +200,7 @@ func (c *Controller) advanceCanary(name string, namespace string) {
|
||||
return
|
||||
}
|
||||
|
||||
// set max weight default value to 100%
|
||||
maxWeight := 100
|
||||
maxWeight := c.fullWeight(cd)
|
||||
if cd.GetAnalysis().MaxWeight > 0 {
|
||||
maxWeight = cd.GetAnalysis().MaxWeight
|
||||
}
|
||||
@@ -207,7 +233,7 @@ func (c *Controller) advanceCanary(name string, namespace string) {
|
||||
cd.Spec.TargetRef.Name, cd.Namespace)
|
||||
|
||||
// route all traffic back to primary
|
||||
primaryWeight = 100
|
||||
primaryWeight = c.fullWeight(cd)
|
||||
canaryWeight = 0
|
||||
if err := meshRouter.SetRoutes(cd, primaryWeight, canaryWeight, false); err != nil {
|
||||
c.recordEventWarningf(cd, "%v", err)
|
||||
@@ -343,7 +369,7 @@ func (c *Controller) advanceCanary(name string, namespace string) {
|
||||
}
|
||||
|
||||
// strategy: Canary progressive traffic increase
|
||||
if cd.GetAnalysis().StepWeight > 0 {
|
||||
if c.nextStepWeight(cd, canaryWeight) > 0 {
|
||||
c.runCanary(cd, canaryController, meshRouter, mirrored, canaryWeight, primaryWeight, maxWeight)
|
||||
}
|
||||
|
||||
@@ -362,22 +388,22 @@ func (c *Controller) runPromotionTrafficShift(canary *flaggerv1.Canary, canaryCo
|
||||
// route all traffic to primary in one go when promotion step wight is not set
|
||||
if canary.Spec.Analysis.StepWeightPromotion == 0 {
|
||||
c.recordEventInfof(canary, "Routing all traffic to primary")
|
||||
if err := meshRouter.SetRoutes(canary, 100, 0, false); err != nil {
|
||||
if err := meshRouter.SetRoutes(canary, c.fullWeight(canary), 0, false); err != nil {
|
||||
c.recordEventWarningf(canary, "%v", err)
|
||||
return
|
||||
}
|
||||
c.recorder.SetWeight(canary, 100, 0)
|
||||
c.recorder.SetWeight(canary, c.fullWeight(canary), 0)
|
||||
if err := canaryController.SetStatusPhase(canary, flaggerv1.CanaryPhaseFinalising); err != nil {
|
||||
c.recordEventWarningf(canary, "%v", err)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
// increment the primary traffic weight until it reaches 100%
|
||||
// increment the primary traffic weight until it reaches 100%/full weight
|
||||
if canaryWeight > 0 {
|
||||
primaryWeight += canary.GetAnalysis().StepWeightPromotion
|
||||
if primaryWeight > 100 {
|
||||
primaryWeight = 100
|
||||
if primaryWeight > c.fullWeight(canary) {
|
||||
primaryWeight = c.fullWeight(canary)
|
||||
}
|
||||
canaryWeight -= canary.GetAnalysis().StepWeightPromotion
|
||||
if canaryWeight < 0 {
|
||||
@@ -391,7 +417,7 @@ func (c *Controller) runPromotionTrafficShift(canary *flaggerv1.Canary, canaryCo
|
||||
c.recordEventInfof(canary, "Advance %s.%s primary weight %v", canary.Name, canary.Namespace, primaryWeight)
|
||||
|
||||
// finalize promotion
|
||||
if primaryWeight == 100 {
|
||||
if primaryWeight == c.fullWeight(canary) {
|
||||
if err := canaryController.SetStatusPhase(canary, flaggerv1.CanaryPhaseFinalising); err != nil {
|
||||
c.recordEventWarningf(canary, "%v", err)
|
||||
}
|
||||
@@ -415,27 +441,30 @@ func (c *Controller) runCanary(canary *flaggerv1.Canary, canaryController canary
|
||||
// If in "mirror" mode, do one step of mirroring before shifting traffic to canary.
|
||||
// When mirroring, all requests go to primary and canary, but only responses from
|
||||
// primary go back to the user.
|
||||
|
||||
var nextStepWeight int
|
||||
nextStepWeight = c.nextStepWeight(canary, canaryWeight)
|
||||
if canary.GetAnalysis().Mirror && canaryWeight == 0 {
|
||||
if !mirrored {
|
||||
mirrored = true
|
||||
primaryWeight = 100
|
||||
primaryWeight = c.fullWeight(canary)
|
||||
canaryWeight = 0
|
||||
} else {
|
||||
mirrored = false
|
||||
primaryWeight = 100 - canary.GetAnalysis().StepWeight
|
||||
canaryWeight = canary.GetAnalysis().StepWeight
|
||||
primaryWeight = c.fullWeight(canary) - nextStepWeight
|
||||
canaryWeight = nextStepWeight
|
||||
}
|
||||
c.logger.With("canary", fmt.Sprintf("%s.%s", canary.Name, canary.Namespace)).
|
||||
Infof("Running mirror step %d/%d/%t", primaryWeight, canaryWeight, mirrored)
|
||||
} else {
|
||||
|
||||
primaryWeight -= canary.GetAnalysis().StepWeight
|
||||
primaryWeight -= nextStepWeight
|
||||
if primaryWeight < 0 {
|
||||
primaryWeight = 0
|
||||
}
|
||||
canaryWeight += canary.GetAnalysis().StepWeight
|
||||
if canaryWeight > 100 {
|
||||
canaryWeight = 100
|
||||
canaryWeight += nextStepWeight
|
||||
if canaryWeight > c.fullWeight(canary) {
|
||||
canaryWeight = c.fullWeight(canary)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -483,11 +512,11 @@ func (c *Controller) runAB(canary *flaggerv1.Canary, canaryController canary.Con
|
||||
|
||||
// route traffic to canary and increment iterations
|
||||
if canary.GetAnalysis().Iterations > canary.Status.Iterations {
|
||||
if err := meshRouter.SetRoutes(canary, 0, 100, false); err != nil {
|
||||
if err := meshRouter.SetRoutes(canary, 0, c.fullWeight(canary), false); err != nil {
|
||||
c.recordEventWarningf(canary, "%v", err)
|
||||
return
|
||||
}
|
||||
c.recorder.SetWeight(canary, 0, 100)
|
||||
c.recorder.SetWeight(canary, 0, c.fullWeight(canary))
|
||||
|
||||
if err := canaryController.SetStatusIterations(canary, canary.Status.Iterations+1); err != nil {
|
||||
c.recordEventWarningf(canary, "%v", err)
|
||||
@@ -529,7 +558,7 @@ func (c *Controller) runBlueGreen(canary *flaggerv1.Canary, canaryController can
|
||||
// If in "mirror" mode, mirror requests during the entire B/G canary test
|
||||
if provider != "kubernetes" &&
|
||||
canary.GetAnalysis().Mirror && !mirrored {
|
||||
if err := meshRouter.SetRoutes(canary, 100, 0, true); err != nil {
|
||||
if err := meshRouter.SetRoutes(canary, c.fullWeight(canary), 0, true); err != nil {
|
||||
c.recordEventWarningf(canary, "%v", err)
|
||||
}
|
||||
c.logger.With("canary", fmt.Sprintf("%s.%s", canary.Name, canary.Namespace)).
|
||||
@@ -557,11 +586,11 @@ func (c *Controller) runBlueGreen(canary *flaggerv1.Canary, canaryController can
|
||||
} else {
|
||||
c.recordEventInfof(canary, "Routing all traffic to canary")
|
||||
}
|
||||
if err := meshRouter.SetRoutes(canary, 0, 100, false); err != nil {
|
||||
if err := meshRouter.SetRoutes(canary, 0, c.fullWeight(canary), false); err != nil {
|
||||
c.recordEventWarningf(canary, "%v", err)
|
||||
return
|
||||
}
|
||||
c.recorder.SetWeight(canary, 0, 100)
|
||||
c.recorder.SetWeight(canary, 0, c.fullWeight(canary))
|
||||
}
|
||||
|
||||
// increment iterations
|
||||
@@ -622,7 +651,7 @@ func (c *Controller) shouldSkipAnalysis(canary *flaggerv1.Canary, canaryControll
|
||||
}
|
||||
|
||||
// route all traffic to primary
|
||||
primaryWeight := 100
|
||||
primaryWeight := c.fullWeight(canary)
|
||||
canaryWeight := 0
|
||||
if err := meshRouter.SetRoutes(canary, primaryWeight, canaryWeight, false); err != nil {
|
||||
c.recordEventWarningf(canary, "%v", err)
|
||||
@@ -756,7 +785,7 @@ func (c *Controller) rollback(canary *flaggerv1.Canary, canaryController canary.
|
||||
}
|
||||
|
||||
// route all traffic back to primary
|
||||
primaryWeight := 100
|
||||
primaryWeight := c.fullWeight(canary)
|
||||
canaryWeight := 0
|
||||
if err := meshRouter.SetRoutes(canary, primaryWeight, canaryWeight, false); err != nil {
|
||||
c.recordEventWarningf(canary, "%v", err)
|
||||
|
||||
Reference in New Issue
Block a user