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,