diff --git a/go.mod b/go.mod index 05edeccd..af8cb962 100644 --- a/go.mod +++ b/go.mod @@ -10,12 +10,14 @@ require ( github.com/prometheus/client_golang v1.22.0 github.com/sirupsen/logrus v1.9.3 github.com/spf13/cobra v1.10.1 + github.com/spf13/pflag v1.0.9 github.com/stretchr/testify v1.10.0 k8s.io/api v0.32.3 k8s.io/apimachinery v0.32.3 k8s.io/client-go v0.32.3 k8s.io/kubectl v0.32.3 k8s.io/utils v0.0.0-20251002143259-bc988d571ff4 + sigs.k8s.io/controller-runtime v0.19.4 ) require ( @@ -24,6 +26,7 @@ require ( github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect github.com/elazarl/goproxy v0.0.0-20240726154733-8b0c20506380 // indirect github.com/emicklei/go-restful/v3 v3.12.2 // indirect + github.com/evanphx/json-patch/v5 v5.9.0 // indirect github.com/fxamacker/cbor/v2 v2.8.0 // indirect github.com/go-logr/logr v1.4.2 // indirect github.com/go-openapi/jsonpointer v0.21.1 // indirect @@ -50,7 +53,6 @@ require ( github.com/prometheus/common v0.63.0 // indirect github.com/prometheus/procfs v0.16.0 // indirect github.com/smartystreets/goconvey v1.7.2 // indirect - github.com/spf13/pflag v1.0.9 // indirect github.com/x448/float16 v0.8.4 // indirect golang.org/x/net v0.39.0 // indirect golang.org/x/oauth2 v0.29.0 // indirect diff --git a/go.sum b/go.sum index 59339eaf..dd99ea92 100644 --- a/go.sum +++ b/go.sum @@ -13,10 +13,14 @@ github.com/elazarl/goproxy v0.0.0-20240726154733-8b0c20506380 h1:1NyRx2f4W4WBRyg github.com/elazarl/goproxy v0.0.0-20240726154733-8b0c20506380/go.mod h1:thX175TtLTzLj3p7N/Q9IiKZ7NF+p72cvL91emV0hzo= github.com/emicklei/go-restful/v3 v3.12.2 h1:DhwDP0vY3k8ZzE0RunuJy8GhNpPL6zqLkDf9B/a0/xU= github.com/emicklei/go-restful/v3 v3.12.2/go.mod h1:6n3XBCmQQb25CM2LCACGz8ukIrRry+4bhvbpWn3mrbc= +github.com/evanphx/json-patch/v5 v5.9.0 h1:kcBlZQbplgElYIlo/n1hJbls2z/1awpXxpRi0/FOJfg= +github.com/evanphx/json-patch/v5 v5.9.0/go.mod h1:VNkHZ/282BpEyt/tObQO8s5CMPmYYq14uClGH4abBuQ= github.com/fxamacker/cbor/v2 v2.8.0 h1:fFtUGXUzXPHTIUdne5+zzMPTfffl3RD5qYnkY40vtxU= github.com/fxamacker/cbor/v2 v2.8.0/go.mod h1:vM4b+DJCtHn+zz7h3FFp/hDAI9WNWCsZj23V5ytsSxQ= github.com/go-logr/logr v1.4.2 h1:6pFjapn8bFcIbiKo3XT4j/BhANplGihG6tvd+8rYgrY= github.com/go-logr/logr v1.4.2/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= +github.com/go-logr/zapr v1.3.0 h1:XGdV8XW8zdwFiwOA2Dryh1gj2KRQyOOoNmBy4EplIcQ= +github.com/go-logr/zapr v1.3.0/go.mod h1:YKepepNBd1u/oyhd/yQmtjVXmm9uML4IXUgMOwR8/Gg= github.com/go-openapi/jsonpointer v0.21.1 h1:whnzv/pNXtK2FbX/W9yJfRmE2gsmkfahjMKB0fZvcic= github.com/go-openapi/jsonpointer v0.21.1/go.mod h1:50I1STOfbY1ycR8jGz8DaMeLCdXiI6aDteEdRNNzpdk= github.com/go-openapi/jsonreference v0.21.0 h1:Rs+Y7hSXT83Jacb7kFyjn4ijOuVGSvOdF2+tg1TRrwQ= @@ -119,9 +123,15 @@ github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9de github.com/yuin/goldmark v1.2.1/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= +go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0= +go.uber.org/multierr v1.11.0/go.mod h1:20+QtiLqy0Nd6FdQB9TLXag12DsQkrbs3htMFfDN80Y= +go.uber.org/zap v1.26.0 h1:sI7k6L95XOKS281NhVKOFCUNIvv9e0w4BF8N3u+tCRo= +go.uber.org/zap v1.26.0/go.mod h1:dtElttAiwGvoJ/vj4IwHBS/gXsEu/pZ50mUIRWuG0so= golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= +golang.org/x/exp v0.0.0-20230515195305-f3d0a9c9a5cc h1:mCRnTeVUjcrhlRmO0VK8a6k6Rrf6TF9htwo2pJVSjIU= +golang.org/x/exp v0.0.0-20230515195305-f3d0a9c9a5cc/go.mod h1:V1LtkGg67GoY2N1AnLN78QLrzxkLyJw7RJb1gzOOz9w= golang.org/x/mod v0.2.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= golang.org/x/mod v0.3.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= golang.org/x/net v0.0.0-20190311183353-d8887717615a/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= @@ -175,6 +185,8 @@ gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= k8s.io/api v0.32.3 h1:Hw7KqxRusq+6QSplE3NYG4MBxZw1BZnq4aP4cJVINls= k8s.io/api v0.32.3/go.mod h1:2wEDTXADtm/HA7CCMD8D8bK4yuBUptzaRhYcYEEYA3k= +k8s.io/apiextensions-apiserver v0.31.0 h1:fZgCVhGwsclj3qCw1buVXCV6khjRzKC5eCFt24kyLSk= +k8s.io/apiextensions-apiserver v0.31.0/go.mod h1:b9aMDEYaEe5sdK+1T0KU78ApR/5ZVp4i56VacZYEHxk= k8s.io/apimachinery v0.32.3 h1:JmDuDarhDmA/Li7j3aPrwhpNBA94Nvk5zLeOge9HH1U= k8s.io/apimachinery v0.32.3/go.mod h1:GpHVgxoKlTxClKcteaeuF1Ul/lDVb74KpZcxcmLDElE= k8s.io/client-go v0.32.3 h1:RKPVltzopkSgHS7aS98QdscAgtgah/+zmpAogooIqVU= @@ -187,6 +199,8 @@ k8s.io/kubectl v0.32.3 h1:VMi584rbboso+yjfv0d8uBHwwxbC438LKq+dXd5tOAI= k8s.io/kubectl v0.32.3/go.mod h1:6Euv2aso5GKzo/UVMacV6C7miuyevpfI91SvBvV9Zdg= k8s.io/utils v0.0.0-20251002143259-bc988d571ff4 h1:SjGebBtkBqHFOli+05xYbK8YF1Dzkbzn+gDM4X9T4Ck= k8s.io/utils v0.0.0-20251002143259-bc988d571ff4/go.mod h1:OLgZIPagt7ERELqWJFomSt595RzquPNLL48iOWgYOg0= +sigs.k8s.io/controller-runtime v0.19.4 h1:SUmheabttt0nx8uJtoII4oIP27BVVvAKFvdvGFwV/Qo= +sigs.k8s.io/controller-runtime v0.19.4/go.mod h1:iRmWllt8IlaLjvTTDLhRBXIEtkCK6hwVBJJsYS9Ajf4= sigs.k8s.io/json v0.0.0-20241014173422-cfa47c3a1cc8 h1:gBQPwqORJ8d8/YNZWEjoZs7npUVDpVXUUOFfW6CgAqE= sigs.k8s.io/json v0.0.0-20241014173422-cfa47c3a1cc8/go.mod h1:mdzfpAEoE6DHQEN0uh9ZbOCuHbLK5wOm7dK4ctXE9Tg= sigs.k8s.io/randfill v0.0.0-20250304075658-069ef1bbf016/go.mod h1:XeLlZ/jmk4i1HRopwe7/aU3H5n1zNUcX6TM94b3QxOY= diff --git a/internal/pkg/config/config.go b/internal/pkg/config/config.go new file mode 100644 index 00000000..7a228284 --- /dev/null +++ b/internal/pkg/config/config.go @@ -0,0 +1,252 @@ +// Package config provides configuration management for Reloader. +// It replaces the old global variables pattern with an immutable Config struct. +package config + +import ( + "time" + + "k8s.io/apimachinery/pkg/labels" +) + +// ReloadStrategy defines how Reloader triggers workload restarts. +type ReloadStrategy string + +const ( + // ReloadStrategyEnvVars adds/updates environment variables to trigger restart. + // This is the default and recommended strategy for GitOps compatibility. + ReloadStrategyEnvVars ReloadStrategy = "env-vars" + + // ReloadStrategyAnnotations adds/updates pod template annotations to trigger restart. + ReloadStrategyAnnotations ReloadStrategy = "annotations" +) + +// ArgoRolloutStrategy defines the strategy for Argo Rollout updates. +type ArgoRolloutStrategy string + +const ( + // ArgoRolloutStrategyRestart uses the restart mechanism for Argo Rollouts. + ArgoRolloutStrategyRestart ArgoRolloutStrategy = "restart" + + // ArgoRolloutStrategyRollout uses the rollout mechanism for Argo Rollouts. + ArgoRolloutStrategyRollout ArgoRolloutStrategy = "rollout" +) + +// Config holds all configuration for Reloader. +// This struct is immutable after creation - all fields should be set during initialization. +type Config struct { + // Annotations holds customizable annotation keys. + Annotations AnnotationConfig + + // AutoReloadAll enables automatic reload for all resources without requiring annotations. + AutoReloadAll bool + + // ReloadStrategy determines how workload restarts are triggered. + ReloadStrategy ReloadStrategy + + // ArgoRolloutsEnabled enables support for Argo Rollouts workload type. + ArgoRolloutsEnabled bool + + // ArgoRolloutStrategy determines how Argo Rollouts are updated. + ArgoRolloutStrategy ArgoRolloutStrategy + + // ReloadOnCreate enables watching for resource creation events. + ReloadOnCreate bool + + // ReloadOnDelete enables watching for resource deletion events. + ReloadOnDelete bool + + // SyncAfterRestart triggers a sync operation after a restart is performed. + SyncAfterRestart bool + + // EnableHA enables high-availability mode with leader election. + EnableHA bool + + // WebhookURL is an optional URL to send notifications to instead of triggering reload. + WebhookURL string + + // Filtering configuration + IgnoredResources []string // ConfigMaps/Secrets to ignore (case-insensitive) + IgnoredWorkloads []string // Workload types to ignore + IgnoredNamespaces []string // Namespaces to ignore + NamespaceSelectors []labels.Selector + ResourceSelectors []labels.Selector + + // Logging configuration + LogFormat string // "json" or "" for default + LogLevel string // trace, debug, info, warning, error, fatal, panic + + // Metrics configuration + MetricsAddr string // Address to serve metrics on (default :9090) + + // Profiling configuration + EnablePProf bool + PProfAddr string + + // Alerting configuration + Alerting AlertingConfig + + // Leader election configuration + LeaderElection LeaderElectionConfig + + // WatchedNamespace limits watching to a specific namespace (empty = all namespaces) + WatchedNamespace string + + // SyncPeriod is the period for re-syncing watched resources + SyncPeriod time.Duration +} + +// AnnotationConfig holds all customizable annotation keys. +type AnnotationConfig struct { + // Prefix is the base prefix for all annotations (default: reloader.stakater.com) + Prefix string + + // Auto annotations + Auto string // reloader.stakater.com/auto + ConfigmapAuto string // configmap.reloader.stakater.com/auto + SecretAuto string // secret.reloader.stakater.com/auto + + // Reload annotations (explicit resource names) + ConfigmapReload string // configmap.reloader.stakater.com/reload + SecretReload string // secret.reloader.stakater.com/reload + + // Exclude annotations + ConfigmapExclude string // configmaps.exclude.reloader.stakater.com/reload + SecretExclude string // secrets.exclude.reloader.stakater.com/reload + + // Ignore annotation + Ignore string // reloader.stakater.com/ignore + + // Search/Match annotations + Search string // reloader.stakater.com/search + Match string // reloader.stakater.com/match + + // Rollout strategy annotation + RolloutStrategy string // reloader.stakater.com/rollout-strategy + + // Pause annotations + PausePeriod string // deployment.reloader.stakater.com/pause-period + PausedAt string // deployment.reloader.stakater.com/paused-at + + // Last reloaded from annotation (set by Reloader) + LastReloadedFrom string // reloader.stakater.com/last-reloaded-from +} + +// AlertingConfig holds configuration for alerting integrations. +type AlertingConfig struct { + SlackWebhookURL string + TeamsWebhookURL string + GChatWebhookURL string +} + +// LeaderElectionConfig holds configuration for leader election. +type LeaderElectionConfig struct { + LockName string + Namespace string + Identity string +} + +// NewDefault creates a Config with default values. +func NewDefault() *Config { + return &Config{ + Annotations: DefaultAnnotations(), + AutoReloadAll: false, + ReloadStrategy: ReloadStrategyEnvVars, + ArgoRolloutsEnabled: false, + ArgoRolloutStrategy: ArgoRolloutStrategyRollout, + ReloadOnCreate: false, + ReloadOnDelete: false, + SyncAfterRestart: false, + EnableHA: false, + WebhookURL: "", + IgnoredResources: []string{}, + IgnoredWorkloads: []string{}, + IgnoredNamespaces: []string{}, + NamespaceSelectors: []labels.Selector{}, + ResourceSelectors: []labels.Selector{}, + LogFormat: "", + LogLevel: "info", + MetricsAddr: ":9090", + EnablePProf: false, + PProfAddr: ":6060", + Alerting: AlertingConfig{}, + LeaderElection: LeaderElectionConfig{ + LockName: "stakater-reloader-lock", + }, + WatchedNamespace: "", + SyncPeriod: 0, + } +} + +// DefaultAnnotations returns the default annotation configuration. +func DefaultAnnotations() AnnotationConfig { + return AnnotationConfig{ + Prefix: "reloader.stakater.com", + Auto: "reloader.stakater.com/auto", + ConfigmapAuto: "configmap.reloader.stakater.com/auto", + SecretAuto: "secret.reloader.stakater.com/auto", + ConfigmapReload: "configmap.reloader.stakater.com/reload", + SecretReload: "secret.reloader.stakater.com/reload", + ConfigmapExclude: "configmaps.exclude.reloader.stakater.com/reload", + SecretExclude: "secrets.exclude.reloader.stakater.com/reload", + Ignore: "reloader.stakater.com/ignore", + Search: "reloader.stakater.com/search", + Match: "reloader.stakater.com/match", + RolloutStrategy: "reloader.stakater.com/rollout-strategy", + PausePeriod: "deployment.reloader.stakater.com/pause-period", + PausedAt: "deployment.reloader.stakater.com/paused-at", + LastReloadedFrom: "reloader.stakater.com/last-reloaded-from", + } +} + +// IsResourceIgnored checks if a resource name should be ignored (case-insensitive). +func (c *Config) IsResourceIgnored(name string) bool { + for _, ignored := range c.IgnoredResources { + if equalFold(ignored, name) { + return true + } + } + return false +} + +// IsWorkloadIgnored checks if a workload type should be ignored (case-insensitive). +func (c *Config) IsWorkloadIgnored(workloadType string) bool { + for _, ignored := range c.IgnoredWorkloads { + if equalFold(ignored, workloadType) { + return true + } + } + return false +} + +// IsNamespaceIgnored checks if a namespace should be ignored. +func (c *Config) IsNamespaceIgnored(namespace string) bool { + for _, ignored := range c.IgnoredNamespaces { + if ignored == namespace { + return true + } + } + return false +} + +// equalFold is a simple case-insensitive string comparison. +func equalFold(s, t string) bool { + if len(s) != len(t) { + return false + } + for i := 0; i < len(s); i++ { + c1, c2 := s[i], t[i] + if c1 != c2 { + // Convert to lowercase for comparison + if 'A' <= c1 && c1 <= 'Z' { + c1 += 'a' - 'A' + } + if 'A' <= c2 && c2 <= 'Z' { + c2 += 'a' - 'A' + } + if c1 != c2 { + return false + } + } + } + return true +} diff --git a/internal/pkg/config/flags.go b/internal/pkg/config/flags.go new file mode 100644 index 00000000..e4363764 --- /dev/null +++ b/internal/pkg/config/flags.go @@ -0,0 +1,157 @@ +package config + +import ( + "strings" + + "github.com/spf13/pflag" +) + +// flagValues holds intermediate string values from CLI flags +// that need further parsing into the Config struct. +type flagValues struct { + namespaceSelectors string + resourceSelectors string + ignoredResources string + ignoredWorkloads string + ignoredNamespaces string + isArgoRollouts string + reloadOnCreate string + reloadOnDelete string +} + +var fv flagValues + +// BindFlags binds configuration flags to the provided flag set. +// Call this before parsing flags, then call ApplyFlags after parsing. +func BindFlags(fs *pflag.FlagSet, cfg *Config) { + // Auto reload + fs.BoolVar(&cfg.AutoReloadAll, "auto-reload-all", cfg.AutoReloadAll, + "Automatically reload all resources when their configmaps/secrets are updated, without requiring annotations") + + // Reload strategy + fs.StringVar((*string)(&cfg.ReloadStrategy), "reload-strategy", string(cfg.ReloadStrategy), + "Strategy for triggering workload restart: 'env-vars' (default, GitOps friendly) or 'annotations'") + + // Argo Rollouts + fs.StringVar(&fv.isArgoRollouts, "is-argo-rollouts", "false", + "Enable Argo Rollouts support (true/false)") + + // Event watching + fs.StringVar(&fv.reloadOnCreate, "reload-on-create", "false", + "Reload when configmaps/secrets are created (true/false)") + fs.StringVar(&fv.reloadOnDelete, "reload-on-delete", "false", + "Reload when configmaps/secrets are deleted (true/false)") + + // Sync after restart + fs.BoolVar(&cfg.SyncAfterRestart, "sync-after-restart", cfg.SyncAfterRestart, + "Trigger sync operation after restart") + + // High availability + fs.BoolVar(&cfg.EnableHA, "enable-ha", cfg.EnableHA, + "Enable high-availability mode with leader election") + + // Webhook + fs.StringVar(&cfg.WebhookURL, "webhook-url", cfg.WebhookURL, + "URL to send notification instead of triggering reload") + + // Filtering - resources + fs.StringVar(&fv.ignoredResources, "resources-to-ignore", "", + "Comma-separated list of configmap/secret names to ignore (case-insensitive)") + fs.StringVar(&fv.ignoredWorkloads, "workload-types-to-ignore", "", + "Comma-separated list of workload types to ignore (Deployment, DaemonSet, StatefulSet)") + fs.StringVar(&fv.ignoredNamespaces, "namespaces-to-ignore", "", + "Comma-separated list of namespaces to ignore") + + // Filtering - selectors + fs.StringVar(&fv.namespaceSelectors, "namespace-selector", "", + "Comma-separated list of namespace label selectors") + fs.StringVar(&fv.resourceSelectors, "resource-label-selector", "", + "Comma-separated list of resource label selectors") + + // Logging + fs.StringVar(&cfg.LogFormat, "log-format", cfg.LogFormat, + "Log format: 'json' or empty for default") + fs.StringVar(&cfg.LogLevel, "log-level", cfg.LogLevel, + "Log level: trace, debug, info, warning, error, fatal, panic") + + // Metrics + fs.StringVar(&cfg.MetricsAddr, "metrics-addr", cfg.MetricsAddr, + "Address to serve metrics on") + + // Profiling + fs.BoolVar(&cfg.EnablePProf, "enable-pprof", cfg.EnablePProf, + "Enable pprof profiling server") + fs.StringVar(&cfg.PProfAddr, "pprof-addr", cfg.PProfAddr, + "Address for pprof server") + + // Annotation customization + fs.StringVar(&cfg.Annotations.Auto, "auto-annotation", cfg.Annotations.Auto, + "Custom annotation for auto-reload") + fs.StringVar(&cfg.Annotations.ConfigmapAuto, "configmap-auto-annotation", cfg.Annotations.ConfigmapAuto, + "Custom annotation for configmap auto-reload") + fs.StringVar(&cfg.Annotations.SecretAuto, "secret-auto-annotation", cfg.Annotations.SecretAuto, + "Custom annotation for secret auto-reload") + fs.StringVar(&cfg.Annotations.ConfigmapReload, "configmap-reload-annotation", cfg.Annotations.ConfigmapReload, + "Custom annotation for configmap reload") + fs.StringVar(&cfg.Annotations.SecretReload, "secret-reload-annotation", cfg.Annotations.SecretReload, + "Custom annotation for secret reload") + fs.StringVar(&cfg.Annotations.Ignore, "ignore-annotation", cfg.Annotations.Ignore, + "Custom annotation for ignoring resources") + fs.StringVar(&cfg.Annotations.Search, "search-annotation", cfg.Annotations.Search, + "Custom annotation for search-based matching") + fs.StringVar(&cfg.Annotations.Match, "match-annotation", cfg.Annotations.Match, + "Custom annotation for match-based matching") + + // Watched namespace (for single-namespace mode) + fs.StringVar(&cfg.WatchedNamespace, "watch-namespace", cfg.WatchedNamespace, + "Namespace to watch (empty for all namespaces)") +} + +// ApplyFlags applies flag values that need post-processing. +// Call this after parsing flags. +func ApplyFlags(cfg *Config) error { + // Parse boolean string flags + cfg.ArgoRolloutsEnabled = parseBoolString(fv.isArgoRollouts) + cfg.ReloadOnCreate = parseBoolString(fv.reloadOnCreate) + cfg.ReloadOnDelete = parseBoolString(fv.reloadOnDelete) + + // Parse comma-separated lists + cfg.IgnoredResources = splitAndTrim(fv.ignoredResources) + cfg.IgnoredWorkloads = splitAndTrim(fv.ignoredWorkloads) + cfg.IgnoredNamespaces = splitAndTrim(fv.ignoredNamespaces) + + // Parse selectors + var err error + cfg.NamespaceSelectors, err = ParseSelectors(splitAndTrim(fv.namespaceSelectors)) + if err != nil { + return err + } + cfg.ResourceSelectors, err = ParseSelectors(splitAndTrim(fv.resourceSelectors)) + if err != nil { + return err + } + + return nil +} + +// parseBoolString parses a string as a boolean, defaulting to false. +func parseBoolString(s string) bool { + s = strings.ToLower(strings.TrimSpace(s)) + return s == "true" || s == "1" || s == "yes" +} + +// splitAndTrim splits a comma-separated string and trims whitespace. +func splitAndTrim(s string) []string { + if s == "" { + return nil + } + parts := strings.Split(s, ",") + result := make([]string, 0, len(parts)) + for _, p := range parts { + p = strings.TrimSpace(p) + if p != "" { + result = append(result, p) + } + } + return result +} diff --git a/internal/pkg/config/validation.go b/internal/pkg/config/validation.go new file mode 100644 index 00000000..8a3bbfe5 --- /dev/null +++ b/internal/pkg/config/validation.go @@ -0,0 +1,148 @@ +package config + +import ( + "fmt" + "strings" + + "k8s.io/apimachinery/pkg/labels" +) + +// ValidationError represents a configuration validation error. +type ValidationError struct { + Field string + Message string +} + +func (e ValidationError) Error() string { + return fmt.Sprintf("config.%s: %s", e.Field, e.Message) +} + +// ValidationErrors is a collection of validation errors. +type ValidationErrors []ValidationError + +func (e ValidationErrors) Error() string { + if len(e) == 0 { + return "" + } + if len(e) == 1 { + return e[0].Error() + } + var b strings.Builder + b.WriteString("multiple configuration errors:\n") + for _, err := range e { + b.WriteString(" - ") + b.WriteString(err.Error()) + b.WriteString("\n") + } + return b.String() +} + +// Validate checks the configuration for errors and normalizes values. +func (c *Config) Validate() error { + var errs ValidationErrors + + // Validate ReloadStrategy + switch c.ReloadStrategy { + case ReloadStrategyEnvVars, ReloadStrategyAnnotations: + // valid + case "": + c.ReloadStrategy = ReloadStrategyEnvVars + default: + errs = append(errs, ValidationError{ + Field: "ReloadStrategy", + Message: fmt.Sprintf("invalid value %q, must be %q or %q", c.ReloadStrategy, ReloadStrategyEnvVars, ReloadStrategyAnnotations), + }) + } + + // Validate ArgoRolloutStrategy + switch c.ArgoRolloutStrategy { + case ArgoRolloutStrategyRestart, ArgoRolloutStrategyRollout: + // valid + case "": + c.ArgoRolloutStrategy = ArgoRolloutStrategyRollout + default: + errs = append(errs, ValidationError{ + Field: "ArgoRolloutStrategy", + Message: fmt.Sprintf("invalid value %q, must be %q or %q", c.ArgoRolloutStrategy, ArgoRolloutStrategyRestart, ArgoRolloutStrategyRollout), + }) + } + + // Validate LogLevel + switch strings.ToLower(c.LogLevel) { + case "trace", "debug", "info", "warn", "warning", "error", "fatal", "panic", "": + // valid + default: + errs = append(errs, ValidationError{ + Field: "LogLevel", + Message: fmt.Sprintf("invalid log level %q", c.LogLevel), + }) + } + + // Validate LogFormat + switch strings.ToLower(c.LogFormat) { + case "json", "": + // valid + default: + errs = append(errs, ValidationError{ + Field: "LogFormat", + Message: fmt.Sprintf("invalid log format %q, must be \"json\" or empty", c.LogFormat), + }) + } + + // Normalize IgnoredResources to lowercase for consistent comparison + c.IgnoredResources = normalizeToLower(c.IgnoredResources) + + // Normalize IgnoredWorkloads to lowercase + c.IgnoredWorkloads = normalizeToLower(c.IgnoredWorkloads) + + if len(errs) > 0 { + return errs + } + return nil +} + +// normalizeToLower converts all strings in the slice to lowercase and removes empty strings. +func normalizeToLower(items []string) []string { + if len(items) == 0 { + return items + } + result := make([]string, 0, len(items)) + for _, item := range items { + item = strings.TrimSpace(strings.ToLower(item)) + if item != "" { + result = append(result, item) + } + } + return result +} + +// ParseSelectors parses a slice of selector strings into label selectors. +func ParseSelectors(selectorStrings []string) ([]labels.Selector, error) { + if len(selectorStrings) == 0 { + return nil, nil + } + + selectors := make([]labels.Selector, 0, len(selectorStrings)) + for _, s := range selectorStrings { + s = strings.TrimSpace(s) + if s == "" { + continue + } + selector, err := labels.Parse(s) + if err != nil { + return nil, fmt.Errorf("invalid selector %q: %w", s, err) + } + selectors = append(selectors, selector) + } + return selectors, nil +} + +// MustParseSelectors parses selectors and panics on error. +// Use only when selectors are known to be valid (e.g., from validated config). +func MustParseSelectors(selectorStrings []string) []labels.Selector { + selectors, err := ParseSelectors(selectorStrings) + if err != nil { + panic(err) + } + return selectors +} diff --git a/internal/pkg/workload/cronjob.go b/internal/pkg/workload/cronjob.go new file mode 100644 index 00000000..42df8eca --- /dev/null +++ b/internal/pkg/workload/cronjob.go @@ -0,0 +1,208 @@ +package workload + +import ( + "context" + + batchv1 "k8s.io/api/batch/v1" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +// CronJobWorkload wraps a Kubernetes CronJob. +// Note: CronJobs have a special update mechanism - instead of updating the CronJob itself, +// Reloader creates a new Job from the CronJob's template. +type CronJobWorkload struct { + cronjob *batchv1.CronJob +} + +// NewCronJobWorkload creates a new CronJobWorkload. +func NewCronJobWorkload(c *batchv1.CronJob) *CronJobWorkload { + return &CronJobWorkload{cronjob: c} +} + +// Ensure CronJobWorkload implements WorkloadAccessor. +var _ WorkloadAccessor = (*CronJobWorkload)(nil) + +func (w *CronJobWorkload) Kind() Kind { + return KindCronJob +} + +func (w *CronJobWorkload) GetObject() client.Object { + return w.cronjob +} + +func (w *CronJobWorkload) GetName() string { + return w.cronjob.Name +} + +func (w *CronJobWorkload) GetNamespace() string { + return w.cronjob.Namespace +} + +func (w *CronJobWorkload) GetAnnotations() map[string]string { + return w.cronjob.Annotations +} + +// GetPodTemplateAnnotations returns annotations from the JobTemplate's pod template. +func (w *CronJobWorkload) GetPodTemplateAnnotations() map[string]string { + if w.cronjob.Spec.JobTemplate.Spec.Template.Annotations == nil { + w.cronjob.Spec.JobTemplate.Spec.Template.Annotations = make(map[string]string) + } + return w.cronjob.Spec.JobTemplate.Spec.Template.Annotations +} + +func (w *CronJobWorkload) SetPodTemplateAnnotation(key, value string) { + if w.cronjob.Spec.JobTemplate.Spec.Template.Annotations == nil { + w.cronjob.Spec.JobTemplate.Spec.Template.Annotations = make(map[string]string) + } + w.cronjob.Spec.JobTemplate.Spec.Template.Annotations[key] = value +} + +func (w *CronJobWorkload) GetContainers() []corev1.Container { + return w.cronjob.Spec.JobTemplate.Spec.Template.Spec.Containers +} + +func (w *CronJobWorkload) SetContainers(containers []corev1.Container) { + w.cronjob.Spec.JobTemplate.Spec.Template.Spec.Containers = containers +} + +func (w *CronJobWorkload) GetInitContainers() []corev1.Container { + return w.cronjob.Spec.JobTemplate.Spec.Template.Spec.InitContainers +} + +func (w *CronJobWorkload) SetInitContainers(containers []corev1.Container) { + w.cronjob.Spec.JobTemplate.Spec.Template.Spec.InitContainers = containers +} + +func (w *CronJobWorkload) GetVolumes() []corev1.Volume { + return w.cronjob.Spec.JobTemplate.Spec.Template.Spec.Volumes +} + +// Update for CronJob is a no-op - use CreateJobFromCronJob instead. +// CronJobs trigger reloads by creating a new Job from their template. +func (w *CronJobWorkload) Update(ctx context.Context, c client.Client) error { + // CronJobs don't get updated directly - a new Job is created instead + // This is handled by the reload package's special CronJob logic + return nil +} + +func (w *CronJobWorkload) DeepCopy() Workload { + return &CronJobWorkload{cronjob: w.cronjob.DeepCopy()} +} + +func (w *CronJobWorkload) GetEnvFromSources() []corev1.EnvFromSource { + var sources []corev1.EnvFromSource + for _, container := range w.cronjob.Spec.JobTemplate.Spec.Template.Spec.Containers { + sources = append(sources, container.EnvFrom...) + } + for _, container := range w.cronjob.Spec.JobTemplate.Spec.Template.Spec.InitContainers { + sources = append(sources, container.EnvFrom...) + } + return sources +} + +func (w *CronJobWorkload) UsesConfigMap(name string) bool { + spec := &w.cronjob.Spec.JobTemplate.Spec.Template.Spec + + // Check volumes + for _, vol := range spec.Volumes { + if vol.ConfigMap != nil && vol.ConfigMap.Name == name { + return true + } + if vol.Projected != nil { + for _, source := range vol.Projected.Sources { + if source.ConfigMap != nil && source.ConfigMap.Name == name { + return true + } + } + } + } + + // Check containers + for _, container := range spec.Containers { + for _, envFrom := range container.EnvFrom { + if envFrom.ConfigMapRef != nil && envFrom.ConfigMapRef.Name == name { + return true + } + } + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.ConfigMapKeyRef != nil && env.ValueFrom.ConfigMapKeyRef.Name == name { + return true + } + } + } + + // Check init containers + for _, container := range spec.InitContainers { + for _, envFrom := range container.EnvFrom { + if envFrom.ConfigMapRef != nil && envFrom.ConfigMapRef.Name == name { + return true + } + } + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.ConfigMapKeyRef != nil && env.ValueFrom.ConfigMapKeyRef.Name == name { + return true + } + } + } + + return false +} + +func (w *CronJobWorkload) UsesSecret(name string) bool { + spec := &w.cronjob.Spec.JobTemplate.Spec.Template.Spec + + // Check volumes + for _, vol := range spec.Volumes { + if vol.Secret != nil && vol.Secret.SecretName == name { + return true + } + if vol.Projected != nil { + for _, source := range vol.Projected.Sources { + if source.Secret != nil && source.Secret.Name == name { + return true + } + } + } + } + + // Check containers + for _, container := range spec.Containers { + for _, envFrom := range container.EnvFrom { + if envFrom.SecretRef != nil && envFrom.SecretRef.Name == name { + return true + } + } + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.SecretKeyRef != nil && env.ValueFrom.SecretKeyRef.Name == name { + return true + } + } + } + + // Check init containers + for _, container := range spec.InitContainers { + for _, envFrom := range container.EnvFrom { + if envFrom.SecretRef != nil && envFrom.SecretRef.Name == name { + return true + } + } + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.SecretKeyRef != nil && env.ValueFrom.SecretKeyRef.Name == name { + return true + } + } + } + + return false +} + +func (w *CronJobWorkload) GetOwnerReferences() []metav1.OwnerReference { + return w.cronjob.OwnerReferences +} + +// GetCronJob returns the underlying CronJob for special handling. +func (w *CronJobWorkload) GetCronJob() *batchv1.CronJob { + return w.cronjob +} diff --git a/internal/pkg/workload/daemonset.go b/internal/pkg/workload/daemonset.go new file mode 100644 index 00000000..ca51f4b5 --- /dev/null +++ b/internal/pkg/workload/daemonset.go @@ -0,0 +1,192 @@ +package workload + +import ( + "context" + + appsv1 "k8s.io/api/apps/v1" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +// DaemonSetWorkload wraps a Kubernetes DaemonSet. +type DaemonSetWorkload struct { + daemonset *appsv1.DaemonSet +} + +// NewDaemonSetWorkload creates a new DaemonSetWorkload. +func NewDaemonSetWorkload(d *appsv1.DaemonSet) *DaemonSetWorkload { + return &DaemonSetWorkload{daemonset: d} +} + +// Ensure DaemonSetWorkload implements WorkloadAccessor. +var _ WorkloadAccessor = (*DaemonSetWorkload)(nil) + +func (w *DaemonSetWorkload) Kind() Kind { + return KindDaemonSet +} + +func (w *DaemonSetWorkload) GetObject() client.Object { + return w.daemonset +} + +func (w *DaemonSetWorkload) GetName() string { + return w.daemonset.Name +} + +func (w *DaemonSetWorkload) GetNamespace() string { + return w.daemonset.Namespace +} + +func (w *DaemonSetWorkload) GetAnnotations() map[string]string { + return w.daemonset.Annotations +} + +func (w *DaemonSetWorkload) GetPodTemplateAnnotations() map[string]string { + if w.daemonset.Spec.Template.Annotations == nil { + w.daemonset.Spec.Template.Annotations = make(map[string]string) + } + return w.daemonset.Spec.Template.Annotations +} + +func (w *DaemonSetWorkload) SetPodTemplateAnnotation(key, value string) { + if w.daemonset.Spec.Template.Annotations == nil { + w.daemonset.Spec.Template.Annotations = make(map[string]string) + } + w.daemonset.Spec.Template.Annotations[key] = value +} + +func (w *DaemonSetWorkload) GetContainers() []corev1.Container { + return w.daemonset.Spec.Template.Spec.Containers +} + +func (w *DaemonSetWorkload) SetContainers(containers []corev1.Container) { + w.daemonset.Spec.Template.Spec.Containers = containers +} + +func (w *DaemonSetWorkload) GetInitContainers() []corev1.Container { + return w.daemonset.Spec.Template.Spec.InitContainers +} + +func (w *DaemonSetWorkload) SetInitContainers(containers []corev1.Container) { + w.daemonset.Spec.Template.Spec.InitContainers = containers +} + +func (w *DaemonSetWorkload) GetVolumes() []corev1.Volume { + return w.daemonset.Spec.Template.Spec.Volumes +} + +func (w *DaemonSetWorkload) Update(ctx context.Context, c client.Client) error { + return c.Update(ctx, w.daemonset) +} + +func (w *DaemonSetWorkload) DeepCopy() Workload { + return &DaemonSetWorkload{daemonset: w.daemonset.DeepCopy()} +} + +func (w *DaemonSetWorkload) GetEnvFromSources() []corev1.EnvFromSource { + var sources []corev1.EnvFromSource + for _, container := range w.daemonset.Spec.Template.Spec.Containers { + sources = append(sources, container.EnvFrom...) + } + for _, container := range w.daemonset.Spec.Template.Spec.InitContainers { + sources = append(sources, container.EnvFrom...) + } + return sources +} + +func (w *DaemonSetWorkload) UsesConfigMap(name string) bool { + // Check volumes + for _, vol := range w.daemonset.Spec.Template.Spec.Volumes { + if vol.ConfigMap != nil && vol.ConfigMap.Name == name { + return true + } + if vol.Projected != nil { + for _, source := range vol.Projected.Sources { + if source.ConfigMap != nil && source.ConfigMap.Name == name { + return true + } + } + } + } + + // Check envFrom + for _, container := range w.daemonset.Spec.Template.Spec.Containers { + for _, envFrom := range container.EnvFrom { + if envFrom.ConfigMapRef != nil && envFrom.ConfigMapRef.Name == name { + return true + } + } + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.ConfigMapKeyRef != nil && env.ValueFrom.ConfigMapKeyRef.Name == name { + return true + } + } + } + + // Check init containers + for _, container := range w.daemonset.Spec.Template.Spec.InitContainers { + for _, envFrom := range container.EnvFrom { + if envFrom.ConfigMapRef != nil && envFrom.ConfigMapRef.Name == name { + return true + } + } + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.ConfigMapKeyRef != nil && env.ValueFrom.ConfigMapKeyRef.Name == name { + return true + } + } + } + + return false +} + +func (w *DaemonSetWorkload) UsesSecret(name string) bool { + // Check volumes + for _, vol := range w.daemonset.Spec.Template.Spec.Volumes { + if vol.Secret != nil && vol.Secret.SecretName == name { + return true + } + if vol.Projected != nil { + for _, source := range vol.Projected.Sources { + if source.Secret != nil && source.Secret.Name == name { + return true + } + } + } + } + + // Check envFrom + for _, container := range w.daemonset.Spec.Template.Spec.Containers { + for _, envFrom := range container.EnvFrom { + if envFrom.SecretRef != nil && envFrom.SecretRef.Name == name { + return true + } + } + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.SecretKeyRef != nil && env.ValueFrom.SecretKeyRef.Name == name { + return true + } + } + } + + // Check init containers + for _, container := range w.daemonset.Spec.Template.Spec.InitContainers { + for _, envFrom := range container.EnvFrom { + if envFrom.SecretRef != nil && envFrom.SecretRef.Name == name { + return true + } + } + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.SecretKeyRef != nil && env.ValueFrom.SecretKeyRef.Name == name { + return true + } + } + } + + return false +} + +func (w *DaemonSetWorkload) GetOwnerReferences() []metav1.OwnerReference { + return w.daemonset.OwnerReferences +} diff --git a/internal/pkg/workload/deployment.go b/internal/pkg/workload/deployment.go new file mode 100644 index 00000000..0ddb5b45 --- /dev/null +++ b/internal/pkg/workload/deployment.go @@ -0,0 +1,194 @@ +package workload + +import ( + "context" + + appsv1 "k8s.io/api/apps/v1" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +// DeploymentWorkload wraps a Kubernetes Deployment. +type DeploymentWorkload struct { + deployment *appsv1.Deployment +} + +// NewDeploymentWorkload creates a new DeploymentWorkload. +func NewDeploymentWorkload(d *appsv1.Deployment) *DeploymentWorkload { + return &DeploymentWorkload{deployment: d} +} + +// Ensure DeploymentWorkload implements WorkloadAccessor. +var _ WorkloadAccessor = (*DeploymentWorkload)(nil) + +func (w *DeploymentWorkload) Kind() Kind { + return KindDeployment +} + +func (w *DeploymentWorkload) GetObject() client.Object { + return w.deployment +} + +func (w *DeploymentWorkload) GetName() string { + return w.deployment.Name +} + +func (w *DeploymentWorkload) GetNamespace() string { + return w.deployment.Namespace +} + +func (w *DeploymentWorkload) GetAnnotations() map[string]string { + return w.deployment.Annotations +} + +func (w *DeploymentWorkload) GetPodTemplateAnnotations() map[string]string { + if w.deployment.Spec.Template.Annotations == nil { + w.deployment.Spec.Template.Annotations = make(map[string]string) + } + return w.deployment.Spec.Template.Annotations +} + +func (w *DeploymentWorkload) SetPodTemplateAnnotation(key, value string) { + if w.deployment.Spec.Template.Annotations == nil { + w.deployment.Spec.Template.Annotations = make(map[string]string) + } + w.deployment.Spec.Template.Annotations[key] = value +} + +func (w *DeploymentWorkload) GetContainers() []corev1.Container { + return w.deployment.Spec.Template.Spec.Containers +} + +func (w *DeploymentWorkload) SetContainers(containers []corev1.Container) { + w.deployment.Spec.Template.Spec.Containers = containers +} + +func (w *DeploymentWorkload) GetInitContainers() []corev1.Container { + return w.deployment.Spec.Template.Spec.InitContainers +} + +func (w *DeploymentWorkload) SetInitContainers(containers []corev1.Container) { + w.deployment.Spec.Template.Spec.InitContainers = containers +} + +func (w *DeploymentWorkload) GetVolumes() []corev1.Volume { + return w.deployment.Spec.Template.Spec.Volumes +} + +func (w *DeploymentWorkload) Update(ctx context.Context, c client.Client) error { + return c.Update(ctx, w.deployment) +} + +func (w *DeploymentWorkload) DeepCopy() Workload { + return &DeploymentWorkload{deployment: w.deployment.DeepCopy()} +} + +func (w *DeploymentWorkload) GetEnvFromSources() []corev1.EnvFromSource { + var sources []corev1.EnvFromSource + for _, container := range w.deployment.Spec.Template.Spec.Containers { + sources = append(sources, container.EnvFrom...) + } + for _, container := range w.deployment.Spec.Template.Spec.InitContainers { + sources = append(sources, container.EnvFrom...) + } + return sources +} + +func (w *DeploymentWorkload) UsesConfigMap(name string) bool { + // Check volumes + for _, vol := range w.deployment.Spec.Template.Spec.Volumes { + if vol.ConfigMap != nil && vol.ConfigMap.Name == name { + return true + } + if vol.Projected != nil { + for _, source := range vol.Projected.Sources { + if source.ConfigMap != nil && source.ConfigMap.Name == name { + return true + } + } + } + } + + // Check envFrom + for _, container := range w.deployment.Spec.Template.Spec.Containers { + for _, envFrom := range container.EnvFrom { + if envFrom.ConfigMapRef != nil && envFrom.ConfigMapRef.Name == name { + return true + } + } + // Check individual env vars + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.ConfigMapKeyRef != nil && env.ValueFrom.ConfigMapKeyRef.Name == name { + return true + } + } + } + + // Check init containers + for _, container := range w.deployment.Spec.Template.Spec.InitContainers { + for _, envFrom := range container.EnvFrom { + if envFrom.ConfigMapRef != nil && envFrom.ConfigMapRef.Name == name { + return true + } + } + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.ConfigMapKeyRef != nil && env.ValueFrom.ConfigMapKeyRef.Name == name { + return true + } + } + } + + return false +} + +func (w *DeploymentWorkload) UsesSecret(name string) bool { + // Check volumes + for _, vol := range w.deployment.Spec.Template.Spec.Volumes { + if vol.Secret != nil && vol.Secret.SecretName == name { + return true + } + if vol.Projected != nil { + for _, source := range vol.Projected.Sources { + if source.Secret != nil && source.Secret.Name == name { + return true + } + } + } + } + + // Check envFrom + for _, container := range w.deployment.Spec.Template.Spec.Containers { + for _, envFrom := range container.EnvFrom { + if envFrom.SecretRef != nil && envFrom.SecretRef.Name == name { + return true + } + } + // Check individual env vars + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.SecretKeyRef != nil && env.ValueFrom.SecretKeyRef.Name == name { + return true + } + } + } + + // Check init containers + for _, container := range w.deployment.Spec.Template.Spec.InitContainers { + for _, envFrom := range container.EnvFrom { + if envFrom.SecretRef != nil && envFrom.SecretRef.Name == name { + return true + } + } + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.SecretKeyRef != nil && env.ValueFrom.SecretKeyRef.Name == name { + return true + } + } + } + + return false +} + +func (w *DeploymentWorkload) GetOwnerReferences() []metav1.OwnerReference { + return w.deployment.OwnerReferences +} diff --git a/internal/pkg/workload/interface.go b/internal/pkg/workload/interface.go new file mode 100644 index 00000000..6f805af4 --- /dev/null +++ b/internal/pkg/workload/interface.go @@ -0,0 +1,120 @@ +// Package workload provides an abstraction layer for Kubernetes workload types. +// It allows uniform handling of Deployments, DaemonSets, StatefulSets, Jobs, CronJobs, and Argo Rollouts. +// +// Note: Jobs and CronJobs have special update mechanisms: +// - Job: deleted and recreated with the same spec +// - CronJob: a new Job is created from the CronJob's template +package workload + +import ( + "context" + + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +// Kind represents the type of workload. +type Kind string + +const ( + KindDeployment Kind = "Deployment" + KindDaemonSet Kind = "DaemonSet" + KindStatefulSet Kind = "StatefulSet" + KindArgoRollout Kind = "Rollout" + KindJob Kind = "Job" + KindCronJob Kind = "CronJob" +) + +// Workload provides a uniform interface for managing Kubernetes workloads. +// All implementations must be safe for concurrent use. +type Workload interface { + // Kind returns the workload type. + Kind() Kind + + // GetObject returns the underlying Kubernetes object. + GetObject() client.Object + + // GetName returns the workload name. + GetName() string + + // GetNamespace returns the workload namespace. + GetNamespace() string + + // GetAnnotations returns the workload's annotations. + GetAnnotations() map[string]string + + // GetPodTemplateAnnotations returns annotations from the pod template spec. + GetPodTemplateAnnotations() map[string]string + + // SetPodTemplateAnnotation sets an annotation on the pod template. + SetPodTemplateAnnotation(key, value string) + + // GetContainers returns all containers (including init containers). + GetContainers() []corev1.Container + + // SetContainers updates the containers. + SetContainers(containers []corev1.Container) + + // GetInitContainers returns all init containers. + GetInitContainers() []corev1.Container + + // SetInitContainers updates the init containers. + SetInitContainers(containers []corev1.Container) + + // GetVolumes returns the pod template volumes. + GetVolumes() []corev1.Volume + + // Update persists changes to the workload. + Update(ctx context.Context, c client.Client) error + + // DeepCopy returns a deep copy of the workload. + DeepCopy() Workload +} + +// Accessor provides read-only access to workload configuration. +// Use this interface when you only need to inspect workload state. +type Accessor interface { + // Kind returns the workload type. + Kind() Kind + + // GetName returns the workload name. + GetName() string + + // GetNamespace returns the workload namespace. + GetNamespace() string + + // GetAnnotations returns the workload's annotations. + GetAnnotations() map[string]string + + // GetPodTemplateAnnotations returns annotations from the pod template spec. + GetPodTemplateAnnotations() map[string]string + + // GetContainers returns all containers (including init containers). + GetContainers() []corev1.Container + + // GetInitContainers returns all init containers. + GetInitContainers() []corev1.Container + + // GetVolumes returns the pod template volumes. + GetVolumes() []corev1.Volume + + // GetEnvFromSources returns all envFrom sources from all containers. + GetEnvFromSources() []corev1.EnvFromSource + + // UsesConfigMap checks if the workload uses a specific ConfigMap. + UsesConfigMap(name string) bool + + // UsesSecret checks if the workload uses a specific Secret. + UsesSecret(name string) bool + + // GetOwnerReferences returns the owner references of the workload. + GetOwnerReferences() []metav1.OwnerReference +} + +// WorkloadAccessor provides both Workload and Accessor interfaces. +// This is the primary type returned by the registry. +type WorkloadAccessor interface { + Workload + Accessor +} diff --git a/internal/pkg/workload/job.go b/internal/pkg/workload/job.go new file mode 100644 index 00000000..85b01e9b --- /dev/null +++ b/internal/pkg/workload/job.go @@ -0,0 +1,203 @@ +package workload + +import ( + "context" + + batchv1 "k8s.io/api/batch/v1" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +// JobWorkload wraps a Kubernetes Job. +// Note: Jobs have a special update mechanism - instead of updating the Job, +// Reloader deletes and recreates it with the same spec. +type JobWorkload struct { + job *batchv1.Job +} + +// NewJobWorkload creates a new JobWorkload. +func NewJobWorkload(j *batchv1.Job) *JobWorkload { + return &JobWorkload{job: j} +} + +// Ensure JobWorkload implements WorkloadAccessor. +var _ WorkloadAccessor = (*JobWorkload)(nil) + +func (w *JobWorkload) Kind() Kind { + return KindJob +} + +func (w *JobWorkload) GetObject() client.Object { + return w.job +} + +func (w *JobWorkload) GetName() string { + return w.job.Name +} + +func (w *JobWorkload) GetNamespace() string { + return w.job.Namespace +} + +func (w *JobWorkload) GetAnnotations() map[string]string { + return w.job.Annotations +} + +func (w *JobWorkload) GetPodTemplateAnnotations() map[string]string { + if w.job.Spec.Template.Annotations == nil { + w.job.Spec.Template.Annotations = make(map[string]string) + } + return w.job.Spec.Template.Annotations +} + +func (w *JobWorkload) SetPodTemplateAnnotation(key, value string) { + if w.job.Spec.Template.Annotations == nil { + w.job.Spec.Template.Annotations = make(map[string]string) + } + w.job.Spec.Template.Annotations[key] = value +} + +func (w *JobWorkload) GetContainers() []corev1.Container { + return w.job.Spec.Template.Spec.Containers +} + +func (w *JobWorkload) SetContainers(containers []corev1.Container) { + w.job.Spec.Template.Spec.Containers = containers +} + +func (w *JobWorkload) GetInitContainers() []corev1.Container { + return w.job.Spec.Template.Spec.InitContainers +} + +func (w *JobWorkload) SetInitContainers(containers []corev1.Container) { + w.job.Spec.Template.Spec.InitContainers = containers +} + +func (w *JobWorkload) GetVolumes() []corev1.Volume { + return w.job.Spec.Template.Spec.Volumes +} + +// Update for Job is a no-op - use RecreateJob instead. +// Jobs trigger reloads by being deleted and recreated. +func (w *JobWorkload) Update(ctx context.Context, c client.Client) error { + // Jobs don't get updated directly - they are deleted and recreated + // This is handled by the reload package's special Job logic + return nil +} + +func (w *JobWorkload) DeepCopy() Workload { + return &JobWorkload{job: w.job.DeepCopy()} +} + +func (w *JobWorkload) GetEnvFromSources() []corev1.EnvFromSource { + var sources []corev1.EnvFromSource + for _, container := range w.job.Spec.Template.Spec.Containers { + sources = append(sources, container.EnvFrom...) + } + for _, container := range w.job.Spec.Template.Spec.InitContainers { + sources = append(sources, container.EnvFrom...) + } + return sources +} + +func (w *JobWorkload) UsesConfigMap(name string) bool { + // Check volumes + for _, vol := range w.job.Spec.Template.Spec.Volumes { + if vol.ConfigMap != nil && vol.ConfigMap.Name == name { + return true + } + if vol.Projected != nil { + for _, source := range vol.Projected.Sources { + if source.ConfigMap != nil && source.ConfigMap.Name == name { + return true + } + } + } + } + + // Check containers + for _, container := range w.job.Spec.Template.Spec.Containers { + for _, envFrom := range container.EnvFrom { + if envFrom.ConfigMapRef != nil && envFrom.ConfigMapRef.Name == name { + return true + } + } + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.ConfigMapKeyRef != nil && env.ValueFrom.ConfigMapKeyRef.Name == name { + return true + } + } + } + + // Check init containers + for _, container := range w.job.Spec.Template.Spec.InitContainers { + for _, envFrom := range container.EnvFrom { + if envFrom.ConfigMapRef != nil && envFrom.ConfigMapRef.Name == name { + return true + } + } + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.ConfigMapKeyRef != nil && env.ValueFrom.ConfigMapKeyRef.Name == name { + return true + } + } + } + + return false +} + +func (w *JobWorkload) UsesSecret(name string) bool { + // Check volumes + for _, vol := range w.job.Spec.Template.Spec.Volumes { + if vol.Secret != nil && vol.Secret.SecretName == name { + return true + } + if vol.Projected != nil { + for _, source := range vol.Projected.Sources { + if source.Secret != nil && source.Secret.Name == name { + return true + } + } + } + } + + // Check containers + for _, container := range w.job.Spec.Template.Spec.Containers { + for _, envFrom := range container.EnvFrom { + if envFrom.SecretRef != nil && envFrom.SecretRef.Name == name { + return true + } + } + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.SecretKeyRef != nil && env.ValueFrom.SecretKeyRef.Name == name { + return true + } + } + } + + // Check init containers + for _, container := range w.job.Spec.Template.Spec.InitContainers { + for _, envFrom := range container.EnvFrom { + if envFrom.SecretRef != nil && envFrom.SecretRef.Name == name { + return true + } + } + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.SecretKeyRef != nil && env.ValueFrom.SecretKeyRef.Name == name { + return true + } + } + } + + return false +} + +func (w *JobWorkload) GetOwnerReferences() []metav1.OwnerReference { + return w.job.OwnerReferences +} + +// GetJob returns the underlying Job for special handling. +func (w *JobWorkload) GetJob() *batchv1.Job { + return w.job +} diff --git a/internal/pkg/workload/registry.go b/internal/pkg/workload/registry.go new file mode 100644 index 00000000..55f55b2d --- /dev/null +++ b/internal/pkg/workload/registry.go @@ -0,0 +1,74 @@ +package workload + +import ( + "fmt" + + appsv1 "k8s.io/api/apps/v1" + batchv1 "k8s.io/api/batch/v1" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +// Registry provides factory methods for creating Workload instances. +type Registry struct { + argoRolloutsEnabled bool +} + +// NewRegistry creates a new workload registry. +func NewRegistry(argoRolloutsEnabled bool) *Registry { + return &Registry{ + argoRolloutsEnabled: argoRolloutsEnabled, + } +} + +// SupportedKinds returns all supported workload kinds. +func (r *Registry) SupportedKinds() []Kind { + kinds := []Kind{ + KindDeployment, + KindDaemonSet, + KindStatefulSet, + KindJob, + KindCronJob, + } + if r.argoRolloutsEnabled { + kinds = append(kinds, KindArgoRollout) + } + return kinds +} + +// FromObject creates a WorkloadAccessor from a Kubernetes object. +func (r *Registry) FromObject(obj client.Object) (WorkloadAccessor, error) { + switch o := obj.(type) { + case *appsv1.Deployment: + return NewDeploymentWorkload(o), nil + case *appsv1.DaemonSet: + return NewDaemonSetWorkload(o), nil + case *appsv1.StatefulSet: + return NewStatefulSetWorkload(o), nil + case *batchv1.Job: + return NewJobWorkload(o), nil + case *batchv1.CronJob: + return NewCronJobWorkload(o), nil + default: + return nil, fmt.Errorf("unsupported object type: %T", obj) + } +} + +// KindFromString converts a string to a Kind. +func KindFromString(s string) (Kind, error) { + switch s { + case "Deployment", "deployment", "deployments": + return KindDeployment, nil + case "DaemonSet", "daemonset", "daemonsets": + return KindDaemonSet, nil + case "StatefulSet", "statefulset", "statefulsets": + return KindStatefulSet, nil + case "Rollout", "rollout", "rollouts": + return KindArgoRollout, nil + case "Job", "job", "jobs": + return KindJob, nil + case "CronJob", "cronjob", "cronjobs": + return KindCronJob, nil + default: + return "", fmt.Errorf("unknown workload kind: %s", s) + } +} diff --git a/internal/pkg/workload/statefulset.go b/internal/pkg/workload/statefulset.go new file mode 100644 index 00000000..003cef3d --- /dev/null +++ b/internal/pkg/workload/statefulset.go @@ -0,0 +1,192 @@ +package workload + +import ( + "context" + + appsv1 "k8s.io/api/apps/v1" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +// StatefulSetWorkload wraps a Kubernetes StatefulSet. +type StatefulSetWorkload struct { + statefulset *appsv1.StatefulSet +} + +// NewStatefulSetWorkload creates a new StatefulSetWorkload. +func NewStatefulSetWorkload(s *appsv1.StatefulSet) *StatefulSetWorkload { + return &StatefulSetWorkload{statefulset: s} +} + +// Ensure StatefulSetWorkload implements WorkloadAccessor. +var _ WorkloadAccessor = (*StatefulSetWorkload)(nil) + +func (w *StatefulSetWorkload) Kind() Kind { + return KindStatefulSet +} + +func (w *StatefulSetWorkload) GetObject() client.Object { + return w.statefulset +} + +func (w *StatefulSetWorkload) GetName() string { + return w.statefulset.Name +} + +func (w *StatefulSetWorkload) GetNamespace() string { + return w.statefulset.Namespace +} + +func (w *StatefulSetWorkload) GetAnnotations() map[string]string { + return w.statefulset.Annotations +} + +func (w *StatefulSetWorkload) GetPodTemplateAnnotations() map[string]string { + if w.statefulset.Spec.Template.Annotations == nil { + w.statefulset.Spec.Template.Annotations = make(map[string]string) + } + return w.statefulset.Spec.Template.Annotations +} + +func (w *StatefulSetWorkload) SetPodTemplateAnnotation(key, value string) { + if w.statefulset.Spec.Template.Annotations == nil { + w.statefulset.Spec.Template.Annotations = make(map[string]string) + } + w.statefulset.Spec.Template.Annotations[key] = value +} + +func (w *StatefulSetWorkload) GetContainers() []corev1.Container { + return w.statefulset.Spec.Template.Spec.Containers +} + +func (w *StatefulSetWorkload) SetContainers(containers []corev1.Container) { + w.statefulset.Spec.Template.Spec.Containers = containers +} + +func (w *StatefulSetWorkload) GetInitContainers() []corev1.Container { + return w.statefulset.Spec.Template.Spec.InitContainers +} + +func (w *StatefulSetWorkload) SetInitContainers(containers []corev1.Container) { + w.statefulset.Spec.Template.Spec.InitContainers = containers +} + +func (w *StatefulSetWorkload) GetVolumes() []corev1.Volume { + return w.statefulset.Spec.Template.Spec.Volumes +} + +func (w *StatefulSetWorkload) Update(ctx context.Context, c client.Client) error { + return c.Update(ctx, w.statefulset) +} + +func (w *StatefulSetWorkload) DeepCopy() Workload { + return &StatefulSetWorkload{statefulset: w.statefulset.DeepCopy()} +} + +func (w *StatefulSetWorkload) GetEnvFromSources() []corev1.EnvFromSource { + var sources []corev1.EnvFromSource + for _, container := range w.statefulset.Spec.Template.Spec.Containers { + sources = append(sources, container.EnvFrom...) + } + for _, container := range w.statefulset.Spec.Template.Spec.InitContainers { + sources = append(sources, container.EnvFrom...) + } + return sources +} + +func (w *StatefulSetWorkload) UsesConfigMap(name string) bool { + // Check volumes + for _, vol := range w.statefulset.Spec.Template.Spec.Volumes { + if vol.ConfigMap != nil && vol.ConfigMap.Name == name { + return true + } + if vol.Projected != nil { + for _, source := range vol.Projected.Sources { + if source.ConfigMap != nil && source.ConfigMap.Name == name { + return true + } + } + } + } + + // Check envFrom + for _, container := range w.statefulset.Spec.Template.Spec.Containers { + for _, envFrom := range container.EnvFrom { + if envFrom.ConfigMapRef != nil && envFrom.ConfigMapRef.Name == name { + return true + } + } + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.ConfigMapKeyRef != nil && env.ValueFrom.ConfigMapKeyRef.Name == name { + return true + } + } + } + + // Check init containers + for _, container := range w.statefulset.Spec.Template.Spec.InitContainers { + for _, envFrom := range container.EnvFrom { + if envFrom.ConfigMapRef != nil && envFrom.ConfigMapRef.Name == name { + return true + } + } + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.ConfigMapKeyRef != nil && env.ValueFrom.ConfigMapKeyRef.Name == name { + return true + } + } + } + + return false +} + +func (w *StatefulSetWorkload) UsesSecret(name string) bool { + // Check volumes + for _, vol := range w.statefulset.Spec.Template.Spec.Volumes { + if vol.Secret != nil && vol.Secret.SecretName == name { + return true + } + if vol.Projected != nil { + for _, source := range vol.Projected.Sources { + if source.Secret != nil && source.Secret.Name == name { + return true + } + } + } + } + + // Check envFrom + for _, container := range w.statefulset.Spec.Template.Spec.Containers { + for _, envFrom := range container.EnvFrom { + if envFrom.SecretRef != nil && envFrom.SecretRef.Name == name { + return true + } + } + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.SecretKeyRef != nil && env.ValueFrom.SecretKeyRef.Name == name { + return true + } + } + } + + // Check init containers + for _, container := range w.statefulset.Spec.Template.Spec.InitContainers { + for _, envFrom := range container.EnvFrom { + if envFrom.SecretRef != nil && envFrom.SecretRef.Name == name { + return true + } + } + for _, env := range container.Env { + if env.ValueFrom != nil && env.ValueFrom.SecretKeyRef != nil && env.ValueFrom.SecretKeyRef.Name == name { + return true + } + } + } + + return false +} + +func (w *StatefulSetWorkload) GetOwnerReferences() []metav1.OwnerReference { + return w.statefulset.OwnerReferences +}