From b937c4ea8d9423961033085108ce27387a288154 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sun, 7 Jul 2019 02:21:13 +0300 Subject: [PATCH 1/4] Implement leader election Add enable-leader-election and leader-election-namespace flags --- cmd/flagger/main.go | 120 +++++++++++++++++++++++++++++++++++--------- go.sum | 1 + 2 files changed, 97 insertions(+), 24 deletions(-) diff --git a/cmd/flagger/main.go b/cmd/flagger/main.go index f7a308d0..1aabb4a8 100644 --- a/cmd/flagger/main.go +++ b/cmd/flagger/main.go @@ -1,6 +1,7 @@ package main import ( + "context" "flag" "fmt" "log" @@ -8,7 +9,7 @@ import ( "strings" "time" - semver "github.com/Masterminds/semver" + "github.com/Masterminds/semver" clientset "github.com/weaveworks/flagger/pkg/client/clientset/versioned" informers "github.com/weaveworks/flagger/pkg/client/informers/externalversions" "github.com/weaveworks/flagger/pkg/controller" @@ -20,31 +21,37 @@ import ( "github.com/weaveworks/flagger/pkg/signals" "github.com/weaveworks/flagger/pkg/version" "go.uber.org/zap" + "k8s.io/apimachinery/pkg/util/uuid" "k8s.io/client-go/kubernetes" _ "k8s.io/client-go/plugin/pkg/client/auth/gcp" "k8s.io/client-go/tools/cache" "k8s.io/client-go/tools/clientcmd" + "k8s.io/client-go/tools/leaderelection" + "k8s.io/client-go/tools/leaderelection/resourcelock" + "k8s.io/client-go/transport" _ "k8s.io/code-generator/cmd/client-gen/generators" ) var ( - masterURL string - kubeconfig string - metricsServer string - controlLoopInterval time.Duration - logLevel string - port string - msteamsURL string - slackURL string - slackUser string - slackChannel string - threadiness int - zapReplaceGlobals bool - zapEncoding string - namespace string - meshProvider string - selectorLabels string - ver bool + masterURL string + kubeconfig string + metricsServer string + controlLoopInterval time.Duration + logLevel string + port string + msteamsURL string + slackURL string + slackUser string + slackChannel string + threadiness int + zapReplaceGlobals bool + zapEncoding string + namespace string + meshProvider string + selectorLabels string + enableLeaderElection bool + leaderElectionNamespace string + ver bool ) func init() { @@ -62,8 +69,10 @@ func init() { flag.BoolVar(&zapReplaceGlobals, "zap-replace-globals", false, "Whether to change the logging level of the global zap logger.") 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, appmesh, supergloo, nginx or smi.") + flag.StringVar(&meshProvider, "mesh-provider", "istio", "Service mesh provider, can be istio, linkerd, appmesh, supergloo, nginx or smi.") flag.StringVar(&selectorLabels, "selector-labels", "app,name,app.kubernetes.io/name", "List of pod labels that Flagger uses to create pod selectors.") + flag.BoolVar(&enableLeaderElection, "enable-leader-election", false, "Enable leader election.") + flag.StringVar(&leaderElectionNamespace, "leader-election-namespace", "kube-system", "Namespace used to create the leader election config map.") flag.BoolVar(&ver, "version", false, "Print version") } @@ -194,14 +203,77 @@ func main() { } } - // start controller - go func(ctrl *controller.Controller) { - if err := ctrl.Run(threadiness, stopCh); err != nil { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + cfg.Wrap(transport.ContextCanceller(ctx, fmt.Errorf("the leader is shutting down"))) + + go func() { + <-stopCh + cancel() + }() + + runController := func() { + if err := c.Run(threadiness, stopCh); err != nil { logger.Fatalf("Error running controller: %v", err) } - }(c) + } - <-stopCh + if enableLeaderElection { + ns := leaderElectionNamespace + if namespace != "" { + ns = namespace + } + startLeaderElection(ctx, runController, ns, kubeClient, logger) + } else { + runController() + } +} + +func startLeaderElection(ctx context.Context, run func(), ns string, kubeClient kubernetes.Interface, logger *zap.SugaredLogger) { + configMapName := "flagger-leader-election" + id, err := os.Hostname() + if err != nil { + logger.Fatalf("Error running controller: %v", err) + } + id = id + "_" + string(uuid.NewUUID()) + + lock, err := resourcelock.New( + resourcelock.ConfigMapsResourceLock, + ns, + configMapName, + kubeClient.CoreV1(), + kubeClient.CoordinationV1(), + resourcelock.ResourceLockConfig{ + Identity: id, + }, + ) + if err != nil { + logger.Fatalf("Error running controller: %v", err) + } + + logger.Infof("Starting leader election id: %s configmap: %s namespace: %s", id, configMapName, ns) + leaderelection.RunOrDie(ctx, leaderelection.LeaderElectionConfig{ + Lock: lock, + ReleaseOnCancel: true, + LeaseDuration: 60 * time.Second, + RenewDeadline: 15 * time.Second, + RetryPeriod: 5 * time.Second, + Callbacks: leaderelection.LeaderCallbacks{ + OnStartedLeading: func(ctx context.Context) { + logger.Info("Acting as elected leader") + run() + }, + OnStoppedLeading: func() { + logger.Infof("Leader election lost") + }, + OnNewLeader: func(identity string) { + if identity != id { + logger.Infof("Another instance has been elected as leader: %v", identity) + } + }, + }, + }) } func initNotifier(logger *zap.SugaredLogger) (client notifier.Interface) { diff --git a/go.sum b/go.sum index 3eb69fc9..b000eb21 100644 --- a/go.sum +++ b/go.sum @@ -154,6 +154,7 @@ github.com/google/gofuzz v1.0.0 h1:A8PeW59pxE9IoFRqBp37U+mSNaQoZ46F1f0f863XSXw= github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= github.com/google/martian v2.1.0+incompatible/go.mod h1:9I4somxYTbIHy5NJKHRl3wXiIaQGbYVAs8BPL6v8lEs= github.com/google/pprof v0.0.0-20181206194817-3ea8567a2e57/go.mod h1:zfwlbNMJ+OItoe0UupaVj+oy1omPYYDuagoSzA8v9mc= +github.com/google/uuid v1.0.0 h1:b4Gk+7WdP/d3HZH8EJsZpvV7EtDOgaZLtnaNGIu1adA= github.com/google/uuid v1.0.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= github.com/googleapis/gax-go/v2 v2.0.4/go.mod h1:0Wqv26UfaUD9n4G6kQubkQ+KchISgw+vpHVxEJEs9eg= github.com/googleapis/gnostic v0.0.0-20170426233943-68f4ded48ba9/go.mod h1:sJBsCZ4ayReDTBIg8b9dl28c5xFWyhBTVRp3pOg5EKY= From a7f4b6d2ae1ceb9b6f06e1dd2f599e2f74708210 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sun, 7 Jul 2019 12:08:08 +0300 Subject: [PATCH 2/4] Add leader election and pod anti affinity to chart --- charts/flagger/README.md | 2 ++ charts/flagger/templates/deployment.yaml | 20 +++++++++++++++----- charts/flagger/values.yaml | 6 ++++-- 3 files changed, 21 insertions(+), 7 deletions(-) diff --git a/charts/flagger/README.md b/charts/flagger/README.md index 3641cac2..12e91463 100644 --- a/charts/flagger/README.md +++ b/charts/flagger/README.md @@ -64,6 +64,8 @@ Parameter | Description | Default `slack.channel` | Slack channel | None `slack.user` | Slack username | `flagger` `msteams.url` | Microsoft Teams incoming webhook | None +`leaderElection.enabled` | leader election must be enabled when running more than one replica | `false` +`leaderElection.replicaCount` | number of replicas | `1` `rbac.create` | if `true`, create and use RBAC resources | `true` `rbac.pspEnabled` | If `true`, create and use a restricted pod security policy | `false` `crd.create` | if `true`, create Flagger's CRDs | `true` diff --git a/charts/flagger/templates/deployment.yaml b/charts/flagger/templates/deployment.yaml index 24f59b01..268224c8 100644 --- a/charts/flagger/templates/deployment.yaml +++ b/charts/flagger/templates/deployment.yaml @@ -8,7 +8,7 @@ metadata: app.kubernetes.io/managed-by: {{ .Release.Service }} app.kubernetes.io/instance: {{ .Release.Name }} spec: - replicas: 1 + replicas: {{ .Values.leaderElection.replicaCount }} strategy: type: Recreate selector: @@ -22,6 +22,16 @@ spec: app.kubernetes.io/instance: {{ .Release.Name }} spec: serviceAccountName: {{ template "flagger.serviceAccountName" . }} + affinity: + podAntiAffinity: + preferredDuringSchedulingIgnoredDuringExecution: + - weight: 100 + podAffinityTerm: + labelSelector: + matchLabels: + app.kubernetes.io/name: {{ template "flagger.name" . }} + app.kubernetes.io/instance: {{ .Release.Name }} + topologyKey: kubernetes.io/hostname containers: - name: flagger securityContext: @@ -54,6 +64,10 @@ spec: {{- if .Values.msteams.url }} - -msteams-url={{ .Values.msteams.url }} {{- end }} + {{- if .Values.leaderElection.enabled }} + - -enable-leader-election=true + - -leader-election-namespace={{ .Release.Namespace }} + {{- end }} livenessProbe: exec: command: @@ -78,10 +92,6 @@ spec: {{ toYaml .Values.resources | indent 12 }} {{- with .Values.nodeSelector }} nodeSelector: -{{ toYaml . | indent 8 }} - {{- end }} - {{- with .Values.affinity }} - affinity: {{ toYaml . | indent 8 }} {{- end }} {{- with .Values.tolerations }} diff --git a/charts/flagger/values.yaml b/charts/flagger/values.yaml index d99717dd..185b0a53 100644 --- a/charts/flagger/values.yaml +++ b/charts/flagger/values.yaml @@ -23,6 +23,10 @@ msteams: # MS Teams incoming webhook URL url: +leaderElection: + enabled: false + replicaCount: 1 + serviceAccount: # serviceAccount.create: Whether to create a service account or not create: true @@ -54,8 +58,6 @@ nodeSelector: {} tolerations: [] -affinity: {} - prometheus: # to be used with AppMesh or nginx ingress install: false From b1bb9fa114910035f28cd9e924fc0beff5095c3d Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sun, 7 Jul 2019 12:08:33 +0300 Subject: [PATCH 3/4] Enable leader election for e2e testing --- test/e2e-linkerd.sh | 2 ++ 1 file changed, 2 insertions(+) diff --git a/test/e2e-linkerd.sh b/test/e2e-linkerd.sh index 4545c64c..5dfa94f5 100755 --- a/test/e2e-linkerd.sh +++ b/test/e2e-linkerd.sh @@ -22,6 +22,8 @@ kind load docker-image test/flagger:latest echo '>>> Installing Flagger' helm upgrade -i flagger ${REPO_ROOT}/charts/flagger \ --namespace linkerd \ +--set leaderElection.enabled=true \ +--set leaderElection.replicaCount=2 \ --set metricsServer=http://linkerd-prometheus:9090 \ --set meshProvider=smi:linkerd From 10c61daee45303501e333ff6b2f937e455c4acc9 Mon Sep 17 00:00:00 2001 From: stefanprodan Date: Sun, 7 Jul 2019 12:52:32 +0300 Subject: [PATCH 4/4] Exit when losing leadership --- Makefile | 9 ++++++++- cmd/flagger/main.go | 8 +++++++- pkg/controller/controller.go | 9 --------- 3 files changed, 15 insertions(+), 11 deletions(-) diff --git a/Makefile b/Makefile index fbca7d9c..eb575eb6 100644 --- a/Makefile +++ b/Makefile @@ -8,7 +8,14 @@ TS=$(shell date +%Y-%m-%d_%H-%M-%S) run: GO111MODULE=on go run cmd/flagger/* -kubeconfig=$$HOME/.kube/config -log-level=info -mesh-provider=istio -namespace=test \ - -metrics-server=https://prometheus.istio.weavedx.com + -metrics-server=https://prometheus.istio.weavedx.com \ + -enable-leader-election=true + +run2: + GO111MODULE=on go run cmd/flagger/* -kubeconfig=$$HOME/.kube/config -log-level=info -mesh-provider=istio -namespace=test \ + -metrics-server=https://prometheus.istio.weavedx.com \ + -enable-leader-election=true \ + -port=9092 run-appmesh: GO111MODULE=on go run cmd/flagger/* -kubeconfig=$$HOME/.kube/config -log-level=info -mesh-provider=appmesh \ diff --git a/cmd/flagger/main.go b/cmd/flagger/main.go index 1aabb4a8..9a17aaba 100644 --- a/cmd/flagger/main.go +++ b/cmd/flagger/main.go @@ -203,22 +203,27 @@ func main() { } } + // leader election context ctx, cancel := context.WithCancel(context.Background()) defer cancel() + // prevents new requests when leadership is lost cfg.Wrap(transport.ContextCanceller(ctx, fmt.Errorf("the leader is shutting down"))) + // cancel leader election context on shutdown signals go func() { <-stopCh cancel() }() + // wrap controller run runController := func() { if err := c.Run(threadiness, stopCh); err != nil { logger.Fatalf("Error running controller: %v", err) } } + // run controller when this instance wins the leader election if enableLeaderElection { ns := leaderElectionNamespace if namespace != "" { @@ -265,7 +270,8 @@ func startLeaderElection(ctx context.Context, run func(), ns string, kubeClient run() }, OnStoppedLeading: func() { - logger.Infof("Leader election lost") + logger.Infof("Leadership lost") + os.Exit(1) }, OnNewLeader: func(identity string) { if identity != id { diff --git a/pkg/controller/controller.go b/pkg/controller/controller.go index 80dec4a8..16a0b8a3 100644 --- a/pkg/controller/controller.go +++ b/pkg/controller/controller.go @@ -216,15 +216,6 @@ func (c *Controller) syncHandler(key string) error { } c.canaries.Store(fmt.Sprintf("%s.%s", cd.Name, cd.Namespace), cd) - - //if cd.Spec.TargetRef.Kind == "Deployment" { - // err = c.bootstrapDeployment(cd) - // if err != nil { - // c.logger.Warnf("%s.%s bootstrap error %v", cd.Name, cd.Namespace, err) - // return err - // } - //} - c.logger.Infof("Synced %s", key) return nil