Make the pod selector configurable

- default labels: app, name and app.kubernetes.io/name
This commit is contained in:
stefanprodan
2019-04-15 12:57:25 +03:00
parent 60f51ad7d5
commit 6ef72e2550
9 changed files with 76 additions and 41 deletions
+9
View File
@@ -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)
+44 -25
View File
@@ -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
}
+8 -8
View File
@@ -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())
}
+5 -4
View File
@@ -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{
+2
View File
@@ -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,
+1
View File
@@ -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,
+3 -2
View File
@@ -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
}
+2 -1
View File
@@ -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,
}
}
+2 -1
View File
@@ -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,