feat: Introduce a Generic ResourceReconciler and a generic BaseWorkload to de-duplicate a lot of code

This commit is contained in:
TheiLLeniumStudios
2026-01-05 00:59:57 +01:00
parent c785067a44
commit 5548ce559a
31 changed files with 1491 additions and 1356 deletions
+42 -118
View File
@@ -1,10 +1,6 @@
package controller
import (
"context"
"sync"
"time"
"github.com/go-logr/logr"
"github.com/stakater/Reloader/internal/pkg/alerting"
"github.com/stakater/Reloader/internal/pkg/config"
@@ -14,129 +10,57 @@ import (
"github.com/stakater/Reloader/internal/pkg/webhook"
"github.com/stakater/Reloader/internal/pkg/workload"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/errors"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/predicate"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
)
// ConfigMapReconciler watches ConfigMaps and triggers workload reloads.
type ConfigMapReconciler struct {
client.Client
Log logr.Logger
Config *config.Config
ReloadService *reload.Service
Registry *workload.Registry
Collectors *metrics.Collectors
EventRecorder *events.Recorder
WebhookClient *webhook.Client
Alerter alerting.Alerter
PauseHandler *reload.PauseHandler
type ConfigMapReconciler = ResourceReconciler[*corev1.ConfigMap]
handler *ReloadHandler
initialized bool
initOnce sync.Once
// NewConfigMapReconciler creates a new ConfigMapReconciler with the given dependencies.
func NewConfigMapReconciler(
c client.Client,
log logr.Logger,
cfg *config.Config,
reloadService *reload.Service,
registry *workload.Registry,
collectors *metrics.Collectors,
eventRecorder *events.Recorder,
webhookClient *webhook.Client,
alerter alerting.Alerter,
pauseHandler *reload.PauseHandler,
) *ConfigMapReconciler {
return NewResourceReconciler(
ResourceReconcilerDeps{
Client: c,
Log: log,
Config: cfg,
ReloadService: reloadService,
Registry: registry,
Collectors: collectors,
EventRecorder: eventRecorder,
WebhookClient: webhookClient,
Alerter: alerter,
PauseHandler: pauseHandler,
},
ResourceConfig[*corev1.ConfigMap]{
ResourceType: reload.ResourceTypeConfigMap,
NewResource: func() *corev1.ConfigMap { return &corev1.ConfigMap{} },
CreateChange: func(cm *corev1.ConfigMap, eventType reload.EventType) reload.ResourceChange {
return reload.ConfigMapChange{ConfigMap: cm, EventType: eventType}
},
CreatePredicates: func(cfg *config.Config, hasher *reload.Hasher) predicate.Predicate {
return reload.ConfigMapPredicates(cfg, hasher)
},
},
)
}
// Reconcile handles ConfigMap events and triggers workload reloads as needed.
func (r *ConfigMapReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
startTime := time.Now()
log := r.Log.WithValues("configmap", req.NamespacedName)
r.initOnce.Do(func() {
r.initialized = true
log.Info("ConfigMap controller initialized")
})
r.Collectors.RecordEventReceived("reconcile", "configmap")
var cm corev1.ConfigMap
if err := r.Get(ctx, req.NamespacedName, &cm); err != nil {
if errors.IsNotFound(err) {
if r.Config.ReloadOnDelete {
r.Collectors.RecordEventReceived("delete", "configmap")
result, err := r.handleDelete(ctx, req, log)
if err != nil {
r.Collectors.RecordReconcile("error", time.Since(startTime))
} else {
r.Collectors.RecordReconcile("success", time.Since(startTime))
}
return result, err
}
r.Collectors.RecordSkipped("not_found")
r.Collectors.RecordReconcile("success", time.Since(startTime))
return ctrl.Result{}, nil
}
log.Error(err, "failed to get ConfigMap")
r.Collectors.RecordError("get_configmap")
r.Collectors.RecordReconcile("error", time.Since(startTime))
return ctrl.Result{}, err
}
if r.Config.IsNamespaceIgnored(cm.Namespace) {
log.V(1).Info("skipping ConfigMap in ignored namespace")
r.Collectors.RecordSkipped("ignored_namespace")
r.Collectors.RecordReconcile("success", time.Since(startTime))
return ctrl.Result{}, nil
}
result, err := r.reloadHandler().Process(ctx, cm.Namespace, cm.Name, reload.ResourceTypeConfigMap,
func(workloads []workload.WorkloadAccessor) []reload.ReloadDecision {
return r.ReloadService.Process(reload.ConfigMapChange{
ConfigMap: &cm,
EventType: reload.EventTypeUpdate,
}, workloads)
}, log)
if err != nil {
r.Collectors.RecordReconcile("error", time.Since(startTime))
} else {
r.Collectors.RecordReconcile("success", time.Since(startTime))
}
return result, err
}
func (r *ConfigMapReconciler) handleDelete(ctx context.Context, req ctrl.Request, log logr.Logger) (ctrl.Result, error) {
log.Info("handling ConfigMap deletion")
cm := &corev1.ConfigMap{}
cm.Name = req.Name
cm.Namespace = req.Namespace
return r.reloadHandler().Process(ctx, req.Namespace, req.Name, reload.ResourceTypeConfigMap,
func(workloads []workload.WorkloadAccessor) []reload.ReloadDecision {
return r.ReloadService.Process(reload.ConfigMapChange{
ConfigMap: cm,
EventType: reload.EventTypeDelete,
}, workloads)
}, log)
}
func (r *ConfigMapReconciler) reloadHandler() *ReloadHandler {
if r.handler == nil {
r.handler = &ReloadHandler{
Client: r.Client,
Lister: workload.NewLister(r.Client, r.Registry, r.Config),
ReloadService: r.ReloadService,
WebhookClient: r.WebhookClient,
Collectors: r.Collectors,
EventRecorder: r.EventRecorder,
Alerter: r.Alerter,
PauseHandler: r.PauseHandler,
}
}
return r.handler
}
// SetupWithManager sets up the controller with the Manager.
func (r *ConfigMapReconciler) SetupWithManager(mgr ctrl.Manager) error {
return ctrl.NewControllerManagedBy(mgr).
For(&corev1.ConfigMap{}).
WithEventFilter(BuildEventFilter(
reload.ConfigMapPredicates(r.Config, r.ReloadService.Hasher()),
r.Config, &r.initialized,
)).
Complete(r)
// SetupConfigMapReconciler sets up a ConfigMap reconciler with the manager.
func SetupConfigMapReconciler(mgr ctrl.Manager, r *ConfigMapReconciler) error {
return r.SetupWithManager(mgr, &corev1.ConfigMap{})
}
var _ reconcile.Reconciler = &ConfigMapReconciler{}
+1 -1
View File
@@ -32,7 +32,7 @@ func (h *ReloadHandler) Process(
ctx context.Context,
namespace, resourceName string,
resourceType reload.ResourceType,
getDecisions func([]workload.WorkloadAccessor) []reload.ReloadDecision,
getDecisions func([]workload.Workload) []reload.ReloadDecision,
log logr.Logger,
) (ctrl.Result, error) {
workloads, err := h.Lister.List(ctx, namespace)
+30 -27
View File
@@ -127,11 +127,12 @@ func NewManagerWithRestConfig(opts ManagerOptions, restConfig *rest.Config) (ctr
func SetupReconcilers(mgr ctrl.Manager, cfg *config.Config, log logr.Logger, collectors *metrics.Collectors) error {
registry := workload.NewRegistry(
workload.RegistryOptions{
ArgoRolloutsEnabled: cfg.ArgoRolloutsEnabled,
DeploymentConfigEnabled: cfg.DeploymentConfigEnabled,
ArgoRolloutsEnabled: cfg.ArgoRolloutsEnabled,
DeploymentConfigEnabled: cfg.DeploymentConfigEnabled,
RolloutStrategyAnnotation: cfg.Annotations.RolloutStrategy,
},
)
reloadService := reload.NewService(cfg)
reloadService := reload.NewService(cfg, log.WithName("reload"))
eventRecorder := events.NewRecorder(mgr.GetEventRecorderFor("reloader"))
pauseHandler := reload.NewPauseHandler(cfg)
@@ -150,36 +151,38 @@ func SetupReconcilers(mgr ctrl.Manager, cfg *config.Config, log logr.Logger, col
// Setup ConfigMap reconciler
if !cfg.IsResourceIgnored("configmaps") {
if err := (&ConfigMapReconciler{
Client: mgr.GetClient(),
Log: log.WithName("configmap-reconciler"),
Config: cfg,
ReloadService: reloadService,
Registry: registry,
Collectors: collectors,
EventRecorder: eventRecorder,
WebhookClient: webhookClient,
Alerter: alerter,
PauseHandler: pauseHandler,
}).SetupWithManager(mgr); err != nil {
cmReconciler := NewConfigMapReconciler(
mgr.GetClient(),
log.WithName("configmap-reconciler"),
cfg,
reloadService,
registry,
collectors,
eventRecorder,
webhookClient,
alerter,
pauseHandler,
)
if err := SetupConfigMapReconciler(mgr, cmReconciler); err != nil {
return fmt.Errorf("setting up configmap reconciler: %w", err)
}
}
// Setup Secret reconciler
if !cfg.IsResourceIgnored("secrets") {
if err := (&SecretReconciler{
Client: mgr.GetClient(),
Log: log.WithName("secret-reconciler"),
Config: cfg,
ReloadService: reloadService,
Registry: registry,
Collectors: collectors,
EventRecorder: eventRecorder,
WebhookClient: webhookClient,
Alerter: alerter,
PauseHandler: pauseHandler,
}).SetupWithManager(mgr); err != nil {
secretReconciler := NewSecretReconciler(
mgr.GetClient(),
log.WithName("secret-reconciler"),
cfg,
reloadService,
registry,
collectors,
eventRecorder,
webhookClient,
alerter,
pauseHandler,
)
if err := SetupSecretReconciler(mgr, secretReconciler); err != nil {
return fmt.Errorf("setting up secret reconciler: %w", err)
}
}
@@ -0,0 +1,186 @@
package controller
import (
"context"
"sync"
"time"
"github.com/go-logr/logr"
"github.com/stakater/Reloader/internal/pkg/alerting"
"github.com/stakater/Reloader/internal/pkg/config"
"github.com/stakater/Reloader/internal/pkg/events"
"github.com/stakater/Reloader/internal/pkg/metrics"
"github.com/stakater/Reloader/internal/pkg/reload"
"github.com/stakater/Reloader/internal/pkg/webhook"
"github.com/stakater/Reloader/internal/pkg/workload"
"k8s.io/apimachinery/pkg/api/errors"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/predicate"
)
// ResourceReconcilerDeps holds shared dependencies for resource reconcilers.
type ResourceReconcilerDeps struct {
Client client.Client
Log logr.Logger
Config *config.Config
ReloadService *reload.Service
Registry *workload.Registry
Collectors *metrics.Collectors
EventRecorder *events.Recorder
WebhookClient *webhook.Client
Alerter alerting.Alerter
PauseHandler *reload.PauseHandler
}
// ResourceConfig provides type-specific configuration for a resource reconciler.
type ResourceConfig[T client.Object] struct {
// ResourceType identifies the type of resource (configmap or secret).
ResourceType reload.ResourceType
// NewResource creates a new instance of the resource type.
NewResource func() T
// CreateChange creates a change event for the resource.
CreateChange func(resource T, eventType reload.EventType) reload.ResourceChange
// CreatePredicates creates the predicates for this resource type.
CreatePredicates func(cfg *config.Config, hasher *reload.Hasher) predicate.Predicate
}
// ResourceReconciler is a generic reconciler for ConfigMaps and Secrets.
type ResourceReconciler[T client.Object] struct {
ResourceReconcilerDeps
ResourceConfig[T]
handler *ReloadHandler
initialized bool
initOnce sync.Once
}
// NewResourceReconciler creates a new generic resource reconciler.
func NewResourceReconciler[T client.Object](
deps ResourceReconcilerDeps,
cfg ResourceConfig[T],
) *ResourceReconciler[T] {
return &ResourceReconciler[T]{
ResourceReconcilerDeps: deps,
ResourceConfig: cfg,
}
}
// Reconcile handles resource events and triggers workload reloads as needed.
func (r *ResourceReconciler[T]) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
startTime := time.Now()
resourceType := string(r.ResourceType)
log := r.Log.WithValues(resourceType, req.NamespacedName)
r.initOnce.Do(func() {
r.initialized = true
log.Info(resourceType + " controller initialized")
})
r.Collectors.RecordEventReceived("reconcile", resourceType)
resource := r.NewResource()
if err := r.Client.Get(ctx, req.NamespacedName, resource); err != nil {
if errors.IsNotFound(err) {
return r.handleNotFound(ctx, req, log, startTime)
}
log.Error(err, "failed to get "+resourceType)
r.Collectors.RecordError("get_" + resourceType)
r.Collectors.RecordReconcile("error", time.Since(startTime))
return ctrl.Result{}, err
}
namespace := resource.GetNamespace()
if r.Config.IsNamespaceIgnored(namespace) {
log.V(1).Info("skipping " + resourceType + " in ignored namespace")
r.Collectors.RecordSkipped("ignored_namespace")
r.Collectors.RecordReconcile("success", time.Since(startTime))
return ctrl.Result{}, nil
}
result, err := r.reloadHandler().Process(ctx, req.Namespace, req.Name, r.ResourceType,
func(workloads []workload.Workload) []reload.ReloadDecision {
return r.ReloadService.Process(r.CreateChange(resource, reload.EventTypeUpdate), workloads)
}, log)
r.recordReconcile(startTime, err)
return result, err
}
func (r *ResourceReconciler[T]) handleNotFound(
ctx context.Context,
req ctrl.Request,
log logr.Logger,
startTime time.Time,
) (ctrl.Result, error) {
if r.Config.ReloadOnDelete {
r.Collectors.RecordEventReceived("delete", string(r.ResourceType))
result, err := r.handleDelete(ctx, req, log)
r.recordReconcile(startTime, err)
return result, err
}
r.Collectors.RecordSkipped("not_found")
r.Collectors.RecordReconcile("success", time.Since(startTime))
return ctrl.Result{}, nil
}
func (r *ResourceReconciler[T]) handleDelete(
ctx context.Context,
req ctrl.Request,
log logr.Logger,
) (ctrl.Result, error) {
log.Info("handling " + string(r.ResourceType) + " deletion")
// Create a minimal resource with just name/namespace for the delete event
resource := r.NewResource()
resource.SetName(req.Name)
resource.SetNamespace(req.Namespace)
return r.reloadHandler().Process(ctx, req.Namespace, req.Name, r.ResourceType,
func(workloads []workload.Workload) []reload.ReloadDecision {
return r.ReloadService.Process(r.CreateChange(resource, reload.EventTypeDelete), workloads)
}, log)
}
func (r *ResourceReconciler[T]) recordReconcile(startTime time.Time, err error) {
if err != nil {
r.Collectors.RecordReconcile("error", time.Since(startTime))
} else {
r.Collectors.RecordReconcile("success", time.Since(startTime))
}
}
func (r *ResourceReconciler[T]) reloadHandler() *ReloadHandler {
if r.handler == nil {
r.handler = &ReloadHandler{
Client: r.Client,
Lister: workload.NewLister(r.Client, r.Registry, r.Config),
ReloadService: r.ReloadService,
WebhookClient: r.WebhookClient,
Collectors: r.Collectors,
EventRecorder: r.EventRecorder,
Alerter: r.Alerter,
PauseHandler: r.PauseHandler,
}
}
return r.handler
}
// Initialized returns whether the reconciler has been initialized.
func (r *ResourceReconciler[T]) Initialized() *bool {
return &r.initialized
}
// SetupWithManager sets up the controller with the Manager.
func (r *ResourceReconciler[T]) SetupWithManager(mgr ctrl.Manager, forObject T) error {
return ctrl.NewControllerManagedBy(mgr).
For(forObject).
WithEventFilter(BuildEventFilter(
r.CreatePredicates(r.Config, r.ReloadService.Hasher()),
r.Config, r.Initialized(),
)).
Complete(r)
}
+21 -172
View File
@@ -2,13 +2,10 @@ package controller
import (
"context"
"maps"
"github.com/stakater/Reloader/internal/pkg/reload"
"github.com/stakater/Reloader/internal/pkg/workload"
batchv1 "k8s.io/api/batch/v1"
"k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/util/retry"
"sigs.k8s.io/controller-runtime/pkg/client"
)
@@ -48,33 +45,31 @@ func UpdateObjectWithRetry(
// UpdateWorkloadWithRetry updates a workload with exponential backoff on conflict.
// On conflict, it re-fetches the object, re-applies the reload changes, and retries.
// For Jobs and CronJobs, special handling is applied:
// - Jobs are deleted and recreated with the same spec
// - CronJobs create a new Job from their template
// For Argo Rollouts, special handling is applied based on the rollout strategy annotation.
// Workloads use their UpdateStrategy to determine how they're updated:
// - UpdateStrategyPatch: uses strategic merge patch with retry (most workloads)
// - UpdateStrategyRecreate: deletes and recreates (Jobs)
// - UpdateStrategyCreateNew: creates a new resource from template (CronJobs)
// Deployments have additional pause handling for paused rollouts.
func UpdateWorkloadWithRetry(
ctx context.Context,
c client.Client,
reloadService *reload.Service,
pauseHandler *reload.PauseHandler,
wl workload.WorkloadAccessor,
wl workload.Workload,
resourceName string,
resourceType reload.ResourceType,
namespace string,
hash string,
autoReload bool,
) (bool, error) {
// Handle special workload types
switch wl.Kind() {
case workload.KindJob:
return updateJobWithRecreate(ctx, c, reloadService, wl, resourceName, resourceType, namespace, hash, autoReload)
case workload.KindCronJob:
return updateCronJobWithNewJob(ctx, c, reloadService, wl, resourceName, resourceType, namespace, hash, autoReload)
case workload.KindArgoRollout:
return updateArgoRollout(ctx, c, reloadService, wl, resourceName, resourceType, namespace, hash, autoReload)
case workload.KindDeployment:
return updateDeploymentWithPause(ctx, c, reloadService, pauseHandler, wl, resourceName, resourceType, namespace, hash, autoReload)
switch wl.UpdateStrategy() {
case workload.UpdateStrategyRecreate, workload.UpdateStrategyCreateNew:
return updateWithSpecialStrategy(ctx, c, reloadService, wl, resourceName, resourceType, namespace, hash, autoReload)
default:
// UpdateStrategyPatch: use standard retry logic with special handling for Deployments
if wl.Kind() == workload.KindDeployment {
return updateDeploymentWithPause(ctx, c, reloadService, pauseHandler, wl, resourceName, resourceType, namespace, hash, autoReload)
}
return updateStandardWorkload(ctx, c, reloadService, wl, resourceName, resourceType, namespace, hash, autoReload)
}
}
@@ -85,7 +80,7 @@ func retryWithReload(
ctx context.Context,
c client.Client,
reloadService *reload.Service,
wl workload.WorkloadAccessor,
wl workload.Workload,
resourceName string,
resourceType reload.ResourceType,
namespace string,
@@ -133,7 +128,7 @@ func updateStandardWorkload(
ctx context.Context,
c client.Client,
reloadService *reload.Service,
wl workload.WorkloadAccessor,
wl workload.Workload,
resourceName string,
resourceType reload.ResourceType,
namespace string,
@@ -154,7 +149,7 @@ func updateDeploymentWithPause(
c client.Client,
reloadService *reload.Service,
pauseHandler *reload.PauseHandler,
wl workload.WorkloadAccessor,
wl workload.Workload,
resourceName string,
resourceType reload.ResourceType,
namespace string,
@@ -176,25 +171,19 @@ func updateDeploymentWithPause(
)
}
// updateJobWithRecreate deletes the Job and recreates it with the updated spec.
// Jobs are immutable after creation, so we must delete and recreate.
func updateJobWithRecreate(
// updateWithSpecialStrategy handles workloads that don't use standard patch.
// It applies reload changes, then delegates to the workload's PerformSpecialUpdate.
func updateWithSpecialStrategy(
ctx context.Context,
c client.Client,
reloadService *reload.Service,
wl workload.WorkloadAccessor,
wl workload.Workload,
resourceName string,
resourceType reload.ResourceType,
namespace string,
hash string,
autoReload bool,
) (bool, error) {
jobWl, ok := wl.(*workload.JobWorkload)
if !ok {
return false, nil
}
// Apply reload changes to the workload
updated, err := reloadService.ApplyReload(
ctx,
wl,
@@ -212,145 +201,5 @@ func updateJobWithRecreate(
return false, nil
}
oldJob := jobWl.GetJob()
newJob := oldJob.DeepCopy()
// Delete the old job with background propagation
policy := metav1.DeletePropagationBackground
if err := c.Delete(
ctx, oldJob, &client.DeleteOptions{
PropagationPolicy: &policy,
},
); err != nil {
if !errors.IsNotFound(err) {
return false, err
}
}
// Clear fields that should not be specified when creating a new Job
newJob.ResourceVersion = ""
newJob.UID = ""
newJob.CreationTimestamp = metav1.Time{}
newJob.Status = batchv1.JobStatus{}
// Remove problematic labels that are auto-generated
delete(newJob.Spec.Template.Labels, "controller-uid")
delete(newJob.Spec.Template.Labels, batchv1.ControllerUidLabel)
delete(newJob.Spec.Template.Labels, batchv1.JobNameLabel)
delete(newJob.Spec.Template.Labels, "job-name")
// Remove the selector to allow it to be auto-generated
newJob.Spec.Selector = nil
// Create the new job with same spec
if err := c.Create(ctx, newJob, client.FieldOwner(workload.FieldManager)); err != nil {
return false, err
}
return true, nil
}
// updateCronJobWithNewJob creates a new Job from the CronJob's template.
// CronJobs don't get updated directly; instead, a new Job is triggered.
func updateCronJobWithNewJob(
ctx context.Context,
c client.Client,
reloadService *reload.Service,
wl workload.WorkloadAccessor,
resourceName string,
resourceType reload.ResourceType,
namespace string,
hash string,
autoReload bool,
) (bool, error) {
cronJobWl, ok := wl.(*workload.CronJobWorkload)
if !ok {
return false, nil
}
// Apply reload changes to get the updated spec
updated, err := reloadService.ApplyReload(
ctx,
wl,
resourceName,
resourceType,
namespace,
hash,
autoReload,
)
if err != nil {
return false, err
}
if !updated {
return false, nil
}
cronJob := cronJobWl.GetCronJob()
annotations := make(map[string]string)
annotations["cronjob.kubernetes.io/instantiate"] = "manual"
maps.Copy(annotations, cronJob.Spec.JobTemplate.Annotations)
job := &batchv1.Job{
ObjectMeta: metav1.ObjectMeta{
GenerateName: cronJob.Name + "-",
Namespace: cronJob.Namespace,
Annotations: annotations,
Labels: cronJob.Spec.JobTemplate.Labels,
OwnerReferences: []metav1.OwnerReference{
*metav1.NewControllerRef(cronJob, batchv1.SchemeGroupVersion.WithKind("CronJob")),
},
},
Spec: cronJob.Spec.JobTemplate.Spec,
}
if err := c.Create(ctx, job, client.FieldOwner(workload.FieldManager)); err != nil {
return false, err
}
savedAnnotations := maps.Clone(cronJob.Spec.JobTemplate.Spec.Template.Annotations)
err = UpdateObjectWithRetry(
ctx, c, cronJob, func() (bool, error) {
if cronJob.Spec.JobTemplate.Spec.Template.Annotations == nil {
cronJob.Spec.JobTemplate.Spec.Template.Annotations = make(map[string]string)
}
maps.Copy(cronJob.Spec.JobTemplate.Spec.Template.Annotations, savedAnnotations)
return true, nil
},
)
if err != nil {
return false, err
}
return true, nil
}
// updateArgoRollout updates an Argo Rollout using its custom Update method.
// This handles the rollout strategy annotation to determine whether to do
// a standard rollout or set the restartAt field.
func updateArgoRollout(
ctx context.Context,
c client.Client,
reloadService *reload.Service,
wl workload.WorkloadAccessor,
resourceName string,
resourceType reload.ResourceType,
namespace string,
hash string,
autoReload bool,
) (bool, error) {
rolloutWl, ok := wl.(*workload.RolloutWorkload)
if !ok {
return false, nil
}
return retryWithReload(
ctx, c, reloadService, wl, resourceName, resourceType, namespace, hash, autoReload,
func() error {
return rolloutWl.Update(ctx, c)
},
)
return wl.PerformSpecialUpdate(ctx, c)
}
+15 -14
View File
@@ -4,6 +4,7 @@ import (
"context"
"testing"
"github.com/go-logr/logr/testr"
"github.com/stakater/Reloader/internal/pkg/config"
"github.com/stakater/Reloader/internal/pkg/controller"
"github.com/stakater/Reloader/internal/pkg/reload"
@@ -22,14 +23,14 @@ func TestUpdateWorkloadWithRetry_WorkloadTypes(t *testing.T) {
tests := []struct {
name string
object runtime.Object
workload func(runtime.Object) workload.WorkloadAccessor
workload func(runtime.Object) workload.Workload
resourceType reload.ResourceType
verify func(t *testing.T, c client.Client)
}{
{
name: "Deployment",
object: testutil.NewDeployment("test-deployment", "default", nil),
workload: func(o runtime.Object) workload.WorkloadAccessor {
workload: func(o runtime.Object) workload.Workload {
return workload.NewDeploymentWorkload(o.(*appsv1.Deployment))
},
resourceType: reload.ResourceTypeConfigMap,
@@ -46,7 +47,7 @@ func TestUpdateWorkloadWithRetry_WorkloadTypes(t *testing.T) {
{
name: "DaemonSet",
object: testutil.NewDaemonSet("test-daemonset", "default", nil),
workload: func(o runtime.Object) workload.WorkloadAccessor {
workload: func(o runtime.Object) workload.Workload {
return workload.NewDaemonSetWorkload(o.(*appsv1.DaemonSet))
},
resourceType: reload.ResourceTypeSecret,
@@ -63,7 +64,7 @@ func TestUpdateWorkloadWithRetry_WorkloadTypes(t *testing.T) {
{
name: "StatefulSet",
object: testutil.NewStatefulSet("test-statefulset", "default", nil),
workload: func(o runtime.Object) workload.WorkloadAccessor {
workload: func(o runtime.Object) workload.Workload {
return workload.NewStatefulSetWorkload(o.(*appsv1.StatefulSet))
},
resourceType: reload.ResourceTypeConfigMap,
@@ -80,7 +81,7 @@ func TestUpdateWorkloadWithRetry_WorkloadTypes(t *testing.T) {
{
name: "Job",
object: testutil.NewJob("test-job", "default"),
workload: func(o runtime.Object) workload.WorkloadAccessor {
workload: func(o runtime.Object) workload.Workload {
return workload.NewJobWorkload(o.(*batchv1.Job))
},
resourceType: reload.ResourceTypeConfigMap,
@@ -97,7 +98,7 @@ func TestUpdateWorkloadWithRetry_WorkloadTypes(t *testing.T) {
{
name: "CronJob",
object: testutil.NewCronJob("test-cronjob", "default"),
workload: func(o runtime.Object) workload.WorkloadAccessor {
workload: func(o runtime.Object) workload.Workload {
return workload.NewCronJobWorkload(o.(*batchv1.CronJob))
},
resourceType: reload.ResourceTypeSecret,
@@ -120,7 +121,7 @@ func TestUpdateWorkloadWithRetry_WorkloadTypes(t *testing.T) {
t.Run(
tt.name, func(t *testing.T) {
cfg := config.NewDefault()
reloadService := reload.NewService(cfg)
reloadService := reload.NewService(cfg, testr.New(t))
fakeClient := fake.NewClientBuilder().
WithScheme(testutil.NewScheme()).
@@ -201,7 +202,7 @@ func TestUpdateWorkloadWithRetry_Strategies(t *testing.T) {
tt.name, func(t *testing.T) {
cfg := config.NewDefault()
cfg.ReloadStrategy = tt.strategy
reloadService := reload.NewService(cfg)
reloadService := reload.NewService(cfg, testr.New(t))
deployment := testutil.NewDeployment("test-deployment", "default", nil)
fakeClient := fake.NewClientBuilder().
@@ -246,7 +247,7 @@ func TestUpdateWorkloadWithRetry_Strategies(t *testing.T) {
func TestUpdateWorkloadWithRetry_NoUpdate(t *testing.T) {
cfg := config.NewDefault()
reloadService := reload.NewService(cfg)
reloadService := reload.NewService(cfg, testr.New(t))
deployment := testutil.NewDeployment("test-deployment", "default", nil)
deployment.Spec.Template.Spec.Containers[0].Env = []corev1.EnvVar{
@@ -306,7 +307,7 @@ func TestResourceTypeKind(t *testing.T) {
func TestUpdateWorkloadWithRetry_PauseDeployment(t *testing.T) {
cfg := config.NewDefault()
reloadService := reload.NewService(cfg)
reloadService := reload.NewService(cfg, testr.New(t))
pauseHandler := reload.NewPauseHandler(cfg)
deployment := testutil.NewDeployment(
@@ -367,7 +368,7 @@ func TestUpdateWorkloadWithRetry_PauseDeployment(t *testing.T) {
// TestUpdateWorkloadWithRetry_PauseWithExplicitAnnotation tests pause with explicit configmap annotation (no auto).
func TestUpdateWorkloadWithRetry_PauseWithExplicitAnnotation(t *testing.T) {
cfg := config.NewDefault()
reloadService := reload.NewService(cfg)
reloadService := reload.NewService(cfg, testr.New(t))
pauseHandler := reload.NewPauseHandler(cfg)
deployment := testutil.NewDeployment(
@@ -428,7 +429,7 @@ func TestUpdateWorkloadWithRetry_PauseWithExplicitAnnotation(t *testing.T) {
// TestUpdateWorkloadWithRetry_PauseWithSecretReload tests pause with Secret-triggered reload.
func TestUpdateWorkloadWithRetry_PauseWithSecretReload(t *testing.T) {
cfg := config.NewDefault()
reloadService := reload.NewService(cfg)
reloadService := reload.NewService(cfg, testr.New(t))
pauseHandler := reload.NewPauseHandler(cfg)
deployment := testutil.NewDeployment(
@@ -485,7 +486,7 @@ func TestUpdateWorkloadWithRetry_PauseWithSecretReload(t *testing.T) {
// TestUpdateWorkloadWithRetry_PauseWithAutoSecret tests pause with auto annotation + Secret change.
func TestUpdateWorkloadWithRetry_PauseWithAutoSecret(t *testing.T) {
cfg := config.NewDefault()
reloadService := reload.NewService(cfg)
reloadService := reload.NewService(cfg, testr.New(t))
pauseHandler := reload.NewPauseHandler(cfg)
deployment := testutil.NewDeployment(
@@ -536,7 +537,7 @@ func TestUpdateWorkloadWithRetry_PauseWithAutoSecret(t *testing.T) {
func TestUpdateWorkloadWithRetry_NoPauseWithoutAnnotation(t *testing.T) {
cfg := config.NewDefault()
reloadService := reload.NewService(cfg)
reloadService := reload.NewService(cfg, testr.New(t))
pauseHandler := reload.NewPauseHandler(cfg)
deployment := testutil.NewDeployment(
+42 -118
View File
@@ -1,10 +1,6 @@
package controller
import (
"context"
"sync"
"time"
"github.com/go-logr/logr"
"github.com/stakater/Reloader/internal/pkg/alerting"
"github.com/stakater/Reloader/internal/pkg/config"
@@ -14,129 +10,57 @@ import (
"github.com/stakater/Reloader/internal/pkg/webhook"
"github.com/stakater/Reloader/internal/pkg/workload"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/errors"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/predicate"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
)
// SecretReconciler watches Secrets and triggers workload reloads.
type SecretReconciler struct {
client.Client
Log logr.Logger
Config *config.Config
ReloadService *reload.Service
Registry *workload.Registry
Collectors *metrics.Collectors
EventRecorder *events.Recorder
WebhookClient *webhook.Client
Alerter alerting.Alerter
PauseHandler *reload.PauseHandler
type SecretReconciler = ResourceReconciler[*corev1.Secret]
handler *ReloadHandler
initialized bool
initOnce sync.Once
// NewSecretReconciler creates a new SecretReconciler with the given dependencies.
func NewSecretReconciler(
c client.Client,
log logr.Logger,
cfg *config.Config,
reloadService *reload.Service,
registry *workload.Registry,
collectors *metrics.Collectors,
eventRecorder *events.Recorder,
webhookClient *webhook.Client,
alerter alerting.Alerter,
pauseHandler *reload.PauseHandler,
) *SecretReconciler {
return NewResourceReconciler(
ResourceReconcilerDeps{
Client: c,
Log: log,
Config: cfg,
ReloadService: reloadService,
Registry: registry,
Collectors: collectors,
EventRecorder: eventRecorder,
WebhookClient: webhookClient,
Alerter: alerter,
PauseHandler: pauseHandler,
},
ResourceConfig[*corev1.Secret]{
ResourceType: reload.ResourceTypeSecret,
NewResource: func() *corev1.Secret { return &corev1.Secret{} },
CreateChange: func(s *corev1.Secret, eventType reload.EventType) reload.ResourceChange {
return reload.SecretChange{Secret: s, EventType: eventType}
},
CreatePredicates: func(cfg *config.Config, hasher *reload.Hasher) predicate.Predicate {
return reload.SecretPredicates(cfg, hasher)
},
},
)
}
// Reconcile handles Secret events and triggers workload reloads as needed.
func (r *SecretReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
startTime := time.Now()
log := r.Log.WithValues("secret", req.NamespacedName)
r.initOnce.Do(func() {
r.initialized = true
log.Info("Secret controller initialized")
})
r.Collectors.RecordEventReceived("reconcile", "secret")
var secret corev1.Secret
if err := r.Get(ctx, req.NamespacedName, &secret); err != nil {
if errors.IsNotFound(err) {
if r.Config.ReloadOnDelete {
r.Collectors.RecordEventReceived("delete", "secret")
result, err := r.handleDelete(ctx, req, log)
if err != nil {
r.Collectors.RecordReconcile("error", time.Since(startTime))
} else {
r.Collectors.RecordReconcile("success", time.Since(startTime))
}
return result, err
}
r.Collectors.RecordSkipped("not_found")
r.Collectors.RecordReconcile("success", time.Since(startTime))
return ctrl.Result{}, nil
}
log.Error(err, "failed to get Secret")
r.Collectors.RecordError("get_secret")
r.Collectors.RecordReconcile("error", time.Since(startTime))
return ctrl.Result{}, err
}
if r.Config.IsNamespaceIgnored(secret.Namespace) {
log.V(1).Info("skipping Secret in ignored namespace")
r.Collectors.RecordSkipped("ignored_namespace")
r.Collectors.RecordReconcile("success", time.Since(startTime))
return ctrl.Result{}, nil
}
result, err := r.reloadHandler().Process(ctx, secret.Namespace, secret.Name, reload.ResourceTypeSecret,
func(workloads []workload.WorkloadAccessor) []reload.ReloadDecision {
return r.ReloadService.Process(reload.SecretChange{
Secret: &secret,
EventType: reload.EventTypeUpdate,
}, workloads)
}, log)
if err != nil {
r.Collectors.RecordReconcile("error", time.Since(startTime))
} else {
r.Collectors.RecordReconcile("success", time.Since(startTime))
}
return result, err
}
func (r *SecretReconciler) handleDelete(ctx context.Context, req ctrl.Request, log logr.Logger) (ctrl.Result, error) {
log.Info("handling Secret deletion")
secret := &corev1.Secret{}
secret.Name = req.Name
secret.Namespace = req.Namespace
return r.reloadHandler().Process(ctx, req.Namespace, req.Name, reload.ResourceTypeSecret,
func(workloads []workload.WorkloadAccessor) []reload.ReloadDecision {
return r.ReloadService.Process(reload.SecretChange{
Secret: secret,
EventType: reload.EventTypeDelete,
}, workloads)
}, log)
}
func (r *SecretReconciler) reloadHandler() *ReloadHandler {
if r.handler == nil {
r.handler = &ReloadHandler{
Client: r.Client,
Lister: workload.NewLister(r.Client, r.Registry, r.Config),
ReloadService: r.ReloadService,
WebhookClient: r.WebhookClient,
Collectors: r.Collectors,
EventRecorder: r.EventRecorder,
Alerter: r.Alerter,
PauseHandler: r.PauseHandler,
}
}
return r.handler
}
// SetupWithManager sets up the controller with the Manager.
func (r *SecretReconciler) SetupWithManager(mgr ctrl.Manager) error {
return ctrl.NewControllerManagedBy(mgr).
For(&corev1.Secret{}).
WithEventFilter(BuildEventFilter(
reload.SecretPredicates(r.Config, r.ReloadService.Hasher()),
r.Config, &r.initialized,
)).
Complete(r)
// SetupSecretReconciler sets up a Secret reconciler with the manager.
func SetupSecretReconciler(mgr ctrl.Manager, r *SecretReconciler) error {
return r.SetupWithManager(mgr, &corev1.Secret{})
}
var _ reconcile.Reconciler = &SecretReconciler{}
+64 -42
View File
@@ -4,6 +4,7 @@ import (
"context"
"testing"
"github.com/go-logr/logr"
"github.com/go-logr/logr/testr"
"github.com/stakater/Reloader/internal/pkg/alerting"
"github.com/stakater/Reloader/internal/pkg/config"
@@ -21,56 +22,77 @@ import (
"sigs.k8s.io/controller-runtime/pkg/client/fake"
)
// testDeps holds shared test dependencies.
type testDeps struct {
client *fake.ClientBuilder
log logr.Logger
cfg *config.Config
reloadService *reload.Service
registry *workload.Registry
collectors *metrics.Collectors
eventRecorder *events.Recorder
webhookClient *webhook.Client
alerter alerting.Alerter
}
// newTestDeps creates shared test dependencies for reconciler tests.
func newTestDeps(t *testing.T, cfg *config.Config, objects ...runtime.Object) testDeps {
t.Helper()
log := testr.New(t)
collectors := metrics.NewCollectors()
return testDeps{
client: fake.NewClientBuilder().
WithScheme(testutil.NewScheme()).
WithRuntimeObjects(objects...),
log: log,
cfg: cfg,
reloadService: reload.NewService(cfg, log),
registry: workload.NewRegistry(workload.RegistryOptions{
ArgoRolloutsEnabled: cfg.ArgoRolloutsEnabled,
DeploymentConfigEnabled: cfg.DeploymentConfigEnabled,
RolloutStrategyAnnotation: cfg.Annotations.RolloutStrategy,
}),
collectors: &collectors,
eventRecorder: events.NewRecorder(nil),
webhookClient: webhook.NewClient("", log),
alerter: &alerting.NoOpAlerter{},
}
}
// newConfigMapReconciler creates a ConfigMapReconciler for testing.
func newConfigMapReconciler(t *testing.T, cfg *config.Config, objects ...runtime.Object) *controller.ConfigMapReconciler {
t.Helper()
fakeClient := fake.NewClientBuilder().
WithScheme(testutil.NewScheme()).
WithRuntimeObjects(objects...).
Build()
collectors := metrics.NewCollectors()
return &controller.ConfigMapReconciler{
Client: fakeClient,
Log: testr.New(t),
Config: cfg,
ReloadService: reload.NewService(cfg),
Registry: workload.NewRegistry(workload.RegistryOptions{
ArgoRolloutsEnabled: cfg.ArgoRolloutsEnabled,
DeploymentConfigEnabled: cfg.DeploymentConfigEnabled,
}),
Collectors: &collectors,
EventRecorder: events.NewRecorder(nil),
WebhookClient: webhook.NewClient("", testr.New(t)),
Alerter: &alerting.NoOpAlerter{},
}
deps := newTestDeps(t, cfg, objects...)
return controller.NewConfigMapReconciler(
deps.client.Build(),
deps.log,
deps.cfg,
deps.reloadService,
deps.registry,
deps.collectors,
deps.eventRecorder,
deps.webhookClient,
deps.alerter,
nil,
)
}
// newSecretReconciler creates a SecretReconciler for testing.
func newSecretReconciler(t *testing.T, cfg *config.Config, objects ...runtime.Object) *controller.SecretReconciler {
t.Helper()
fakeClient := fake.NewClientBuilder().
WithScheme(testutil.NewScheme()).
WithRuntimeObjects(objects...).
Build()
collectors := metrics.NewCollectors()
return &controller.SecretReconciler{
Client: fakeClient,
Log: testr.New(t),
Config: cfg,
ReloadService: reload.NewService(cfg),
Registry: workload.NewRegistry(workload.RegistryOptions{
ArgoRolloutsEnabled: cfg.ArgoRolloutsEnabled,
DeploymentConfigEnabled: cfg.DeploymentConfigEnabled,
}),
Collectors: &collectors,
EventRecorder: events.NewRecorder(nil),
WebhookClient: webhook.NewClient("", testr.New(t)),
Alerter: &alerting.NoOpAlerter{},
}
deps := newTestDeps(t, cfg, objects...)
return controller.NewSecretReconciler(
deps.client.Build(),
deps.log,
deps.cfg,
deps.reloadService,
deps.registry,
deps.collectors,
deps.eventRecorder,
deps.webhookClient,
deps.alerter,
nil,
)
}
// newNamespaceReconciler creates a NamespaceReconciler for testing.