feat: add v2 foundation packages for config and workload abstraction

This commit is contained in:
TheiLLeniumStudios
2025-12-28 08:47:51 +01:00
parent a5d1012570
commit b9a7b8997c
12 changed files with 1757 additions and 1 deletions
+3 -1
View File
@@ -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
+14
View File
@@ -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=
+252
View File
@@ -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
}
+157
View File
@@ -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
}
+148
View File
@@ -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
}
+208
View File
@@ -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
}
+192
View File
@@ -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
}
+194
View File
@@ -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
}
+120
View File
@@ -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
}
+203
View File
@@ -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
}
+74
View File
@@ -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)
}
}
+192
View File
@@ -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
}