Files
Reloader/internal/pkg/controller/resource_reconciler.go
T
2026-06-22 20:34:40 +05:00

227 lines
7.6 KiB
Go

package controller
import (
"context"
"time"
"github.com/go-logr/logr"
"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"
"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"
)
// ResourceReconcilerDeps holds shared dependencies for resource reconcilers.
type ResourceReconcilerDeps struct {
Client client.Client
// APIReader is an optional non-cached reader for ResolveChange lookups.
APIReader client.Reader
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
NamespaceCache *NamespaceCache
}
// 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
// ResolveChange derives the change from a second object instead of CreateChange;
// ok=false skips the event (CSI: change comes from the parent SecretProviderClass).
ResolveChange func(ctx context.Context, reader client.Reader, log logr.Logger, resource T) (reload.ResourceChange, bool)
// SkipOnNotFound treats a missing object as a no-op (CSI deletes SPCPS as pods roll).
SkipOnNotFound bool
// BuildFilter overrides the default BuildEventFilter for the watch.
BuildFilter 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
}
// 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.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
}
if r.NamespaceCache != nil && r.NamespaceCache.IsEnabled() && !r.NamespaceCache.Contains(namespace) {
log.V(1).Info("skipping "+resourceType+" in namespace not matching selector", "namespace", namespace)
r.Collectors.RecordSkipped("namespace_selector")
r.Collectors.RecordReconcile("success", time.Since(startTime))
return ctrl.Result{}, nil
}
change, ok := r.buildChange(ctx, log, resource)
if !ok {
r.Collectors.RecordSkipped("resolve_skipped")
r.Collectors.RecordReconcile("success", time.Since(startTime))
return ctrl.Result{}, nil
}
result, err := r.reloadHandler().Process(
ctx, change.GetNamespace(), change.GetName(), r.ResourceType,
func(workloads []workload.Workload) []reload.ReloadDecision {
return r.ReloadService.Process(change, workloads)
}, log,
)
r.recordReconcile(startTime, err)
return result, err
}
// buildChange returns the change via ResolveChange when set, else CreateChange.
func (r *ResourceReconciler[T]) buildChange(ctx context.Context, log logr.Logger, resource T) (reload.ResourceChange, bool) {
if r.ResolveChange != nil {
return r.ResolveChange(ctx, r.APIReader, log, resource)
}
return r.CreateChange(resource, reload.EventTypeUpdate), true
}
func (r *ResourceReconciler[T]) handleNotFound(
ctx context.Context,
req ctrl.Request,
log logr.Logger,
startTime time.Time,
) (ctrl.Result, error) {
if r.SkipOnNotFound {
r.Collectors.RecordSkipped("not_found")
r.Collectors.RecordReconcile("success", time.Since(startTime))
return ctrl.Result{}, nil
}
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
}
// SetupWithManager sets up the controller with the Manager.
func (r *ResourceReconciler[T]) SetupWithManager(mgr ctrl.Manager, forObject T) error {
var filter predicate.Predicate
if r.BuildFilter != nil {
filter = r.BuildFilter(r.Config, r.ReloadService.Hasher())
} else {
// time.Now() lets the create predicate ignore initial-sync replays of
// pre-existing resources (older creation timestamps) while honoring later creates.
filter = BuildEventFilter(
r.CreatePredicates(r.Config, r.ReloadService.Hasher()),
r.Config, time.Now(),
)
}
return ctrl.NewControllerManagedBy(mgr).
For(forObject).
WithEventFilter(filter).
Complete(r)
}