Files
capsule/internal/controllers/resources/namespaced.go
T
fed61d2815 feat: replicating resources upon namespace creation (#2080)
* feat: replicating resources upon namespace creation

This feature enhances the (Global)TenantResource replication by
triggering a replication of resources upon a Namespace creation: this
speeds up the replication of resources without waiting for the
resyncPeriod that could be delayed for several reasons.

Signed-off-by: Dario Tranchitella <dario@tranchitella.eu>

* fix: addressing fix from github copilot

Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
Signed-off-by: Dario Tranchitella <dario@tranchitella.eu>

---------

Signed-off-by: Dario Tranchitella <dario@tranchitella.eu>
Co-authored-by: Oliver Bähler <26610571+oliverbaehler@users.noreply.github.com>
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
2026-08-17 10:51:09 +02:00

714 lines
19 KiB
Go

// Copyright 2020-2026 Project Capsule Authors
// SPDX-License-Identifier: Apache-2.0
package resources
import (
"context"
"fmt"
"reflect"
"strconv"
"github.com/go-logr/logr"
gherrors "github.com/pkg/errors"
corev1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/util/sets"
"k8s.io/client-go/util/retry"
"sigs.k8s.io/cluster-api/util/patch"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/builder"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
"sigs.k8s.io/controller-runtime/pkg/handler"
ctrllog "sigs.k8s.io/controller-runtime/pkg/log"
"sigs.k8s.io/controller-runtime/pkg/predicate"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
capsulev1beta2 "github.com/projectcapsule/capsule/api/v1beta2"
"github.com/projectcapsule/capsule/internal/cache"
cutils "github.com/projectcapsule/capsule/internal/controllers/utils"
"github.com/projectcapsule/capsule/internal/metrics"
caperrors "github.com/projectcapsule/capsule/pkg/api/errors"
"github.com/projectcapsule/capsule/pkg/api/meta"
"github.com/projectcapsule/capsule/pkg/api/processor"
"github.com/projectcapsule/capsule/pkg/runtime/configuration"
tenantresourceindexer "github.com/projectcapsule/capsule/pkg/runtime/indexers/tenantresource"
"github.com/projectcapsule/capsule/pkg/runtime/predicates"
"github.com/projectcapsule/capsule/pkg/runtime/sanitize"
tpl "github.com/projectcapsule/capsule/pkg/template"
"github.com/projectcapsule/capsule/pkg/tenant"
)
type namespacedResourceController struct {
client client.Client
reader client.Reader
log logr.Logger
processor processor.Processor
collector Collector
configuration configuration.Configuration
metrics *metrics.TenantResourceRecorder
clients impersonatedClientLoader[*capsulev1beta2.TenantResource]
impersonation *cache.ImpersonationCache
}
func (r *namespacedResourceController) SetupWithManager(mgr ctrl.Manager, ctrlConfig cutils.ControllerOptions) error {
r.client = mgr.GetClient()
r.reader = mgr.GetAPIReader()
r.processor = processor.Processor{
Configuration: r.configuration,
AllowCrossNamespaceSelection: false,
GatherClient: mgr.GetAPIReader(),
Mapper: mgr.GetRESTMapper(),
}
r.collector = NewCollector(
mgr.GetAPIReader(),
mgr.GetRESTMapper(),
)
r.clients = impersonatedClientLoader[*capsulev1beta2.TenantResource]{
client: r.client,
configuration: r.configuration,
impersonation: r.impersonation,
resolve: namespacedServiceAccount,
}
return ctrl.NewControllerManagedBy(mgr).
For(
&capsulev1beta2.TenantResource{},
builder.WithPredicates(
predicate.Or(
predicate.GenerationChangedPredicate{},
predicates.ReconcileRequestedPredicate{},
),
),
).
Watches(
&capsulev1beta2.TenantResource{},
handler.EnqueueRequestsFromMapFunc(r.enqueueDependentTenantResources),
builder.WithPredicates(predicates.DependencyStateChangedPredicate{}),
).
Watches(
&capsulev1beta2.CapsuleConfiguration{},
handler.EnqueueRequestsFromMapFunc(r.enqueueAllResources),
builder.WithPredicates(
predicates.CapsuleConfigSpecImpersonationChangedPredicate{},
predicates.NamesMatchingPredicate{Names: []string{ctrlConfig.ConfigurationName}},
),
).
Watches(
&corev1.Namespace{},
handler.EnqueueRequestsFromMapFunc(r.enqueueTenantResourcesForNamespace),
builder.WithPredicates(
predicates.LabelPresentPredicate{Label: meta.TenantLabel},
),
).
Watches(
&capsulev1beta2.Tenant{},
handler.EnqueueRequestsFromMapFunc(r.enqueueTenantResourcesForTenant),
builder.WithPredicates(predicates.TenantNamespacesChangedPredicate{}),
).
WithOptions(ctrlConfig.Runtime.ToControllerOptions()).
Complete(r)
}
func (r *namespacedResourceController) Reconcile(ctx context.Context, request reconcile.Request) (res reconcile.Result, err error) {
log := ctrllog.FromContext(ctx)
log.V(5).Info("start processing")
tntResource := &capsulev1beta2.TenantResource{}
if err := r.client.Get(ctx, request.NamespacedName, tntResource); err != nil {
if apierrors.IsNotFound(err) {
log.V(5).Info("Request object not found, could have been deleted after reconcile request")
r.metrics.DeleteMetrics(request.Name, request.Namespace)
return reconcile.Result{}, nil
}
return reconcile.Result{}, err
}
requeue := reconcile.Result{
RequeueAfter: jitteredResync(tntResource.Spec.ResyncPeriod.Duration),
}
patchHelper, err := patch.NewHelper(tntResource, r.client)
if err != nil {
return reconcile.Result{}, gherrors.Wrap(err, "failed to init patch helper")
}
var statusErr error
//nolint:dupl
defer func() {
meta.RemoveReconcileTriggerAnnotation(tntResource)
reconcileErr := err
if statusErr != nil {
reconcileErr = statusErr
}
if uerr := r.updateStatus(ctx, tntResource, reconcileErr); uerr != nil {
if caperrors.IgnoreGone(uerr) {
err = nil
return
}
err = fmt.Errorf("cannot update tenantresource status: %w", uerr)
return
}
r.metrics.RecordConditions(tntResource)
if e := patchHelper.Patch(ctx, tntResource); e != nil {
if caperrors.IgnoreGone(e) {
err = nil
return
}
res = reconcile.Result{}
err = gherrors.Wrap(e, "failed to patch TenantResource")
return
}
// Controller-runtime should not receive handled reconciliation errors.
err = nil
}()
// On Deletion these checks are skipped.
//nolint:nestif
if tntResource.DeletionTimestamp.IsZero() {
if tntResource.Spec.IsCordoned() {
log.V(5).Info("tenant resource cordoned")
return reconcile.Result{}, nil
}
for _, dep := range tntResource.Spec.DependsOn {
d := &capsulev1beta2.TenantResource{}
if getErr := r.client.Get(ctx, types.NamespacedName{
Name: dep.Name.String(),
Namespace: tntResource.GetNamespace(),
}, d); getErr != nil {
if apierrors.IsNotFound(getErr) {
statusErr = fmt.Errorf("dependency %s not found", dep.Name)
} else {
statusErr = getErr
}
return requeue, nil
}
stat := d.Status.Conditions.GetConditionByType(meta.ReadyCondition)
if stat == nil || stat.Status != metav1.ConditionTrue {
statusErr = fmt.Errorf("dependency %s not ready", dep.Name)
return requeue, nil
}
}
}
// Load client must be first since it updates the new serviceaccount used which can then be directly
// posted to the status.
c, loadErr := r.loadClient(ctx, log, tntResource)
if loadErr != nil {
statusErr = gherrors.Wrap(loadErr, "failed to load serviceaccount client")
return requeue, nil
}
// Best-Effort for Updating the status
if updateErr := r.updateReconcilingStatus(ctx, tntResource); updateErr != nil {
if caperrors.IgnoreGone(updateErr) {
return reconcile.Result{}, nil
}
log.Error(updateErr, "failed to update status")
}
if c == nil {
statusErr = fmt.Errorf("received empty client for serviceaccount")
return requeue, nil
}
statusErr = r.reconcile(ctx, c, tntResource)
if len(tntResource.Status.ProcessedItems) > 0 {
controllerutil.AddFinalizer(tntResource, meta.ControllerFinalizer)
} else {
controllerutil.RemoveFinalizer(tntResource, meta.ControllerFinalizer)
}
controllerutil.RemoveFinalizer(tntResource, meta.LegacyResourceFinalizer)
return requeue, nil
}
func (r *namespacedResourceController) enqueueTenantResourcesForTenant(ctx context.Context, obj client.Object) []reconcile.Request {
tnt, ok := obj.(*capsulev1beta2.Tenant)
if !ok {
return nil
}
seen := map[types.NamespacedName]struct{}{}
out := make([]reconcile.Request, 0)
for _, ns := range tnt.Status.Namespaces {
list := &capsulev1beta2.TenantResourceList{}
if err := r.client.List(ctx, list, client.InNamespace(ns)); err != nil {
continue
}
for i := range list.Items {
key := types.NamespacedName{
Name: list.Items[i].Name,
Namespace: list.Items[i].Namespace,
}
if _, exists := seen[key]; exists {
continue
}
seen[key] = struct{}{}
out = append(out, reconcile.Request{NamespacedName: key})
}
}
return out
}
func (r *namespacedResourceController) enqueueDependentTenantResources(
ctx context.Context,
obj client.Object,
) []ctrl.Request {
changed, ok := obj.(*capsulev1beta2.TenantResource)
if !ok {
return nil
}
var list capsulev1beta2.TenantResourceList
if err := r.client.List(
ctx,
&list,
client.InNamespace(changed.Namespace),
client.MatchingFields{tenantresourceindexer.NamespacedDependenciesFieldName: changed.Name},
); err != nil {
return nil
}
reqs := make([]ctrl.Request, 0, len(list.Items))
for i := range list.Items {
reqs = append(reqs, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(&list.Items[i])})
}
return reqs
}
// Requeue TenantResources if there is changes to namespaces of the same tenant
// We are not relying on the tenant status, as we might have a terminating lock caused by TenantResources.
func (r *namespacedResourceController) enqueueTenantResourcesForNamespace(
ctx context.Context,
obj client.Object,
) []reconcile.Request {
ns, ok := obj.(*corev1.Namespace)
if !ok {
return nil
}
labelValue, ok := ns.Labels[meta.TenantLabel]
if !ok || labelValue == "" {
return nil
}
var namespaces corev1.NamespaceList
if err := r.client.List(
ctx,
&namespaces,
client.MatchingLabels{meta.TenantLabel: labelValue},
); err != nil {
r.log.Error(err, "failed to list namespaces by label", "label", meta.TenantLabel, "value", labelValue)
return nil
}
requests := make([]reconcile.Request, 0, 16)
seen := make(map[types.NamespacedName]struct{})
for i := range namespaces.Items {
var trList capsulev1beta2.TenantResourceList
if err := r.client.List(
ctx,
&trList,
client.InNamespace(namespaces.Items[i].Name),
); err != nil {
r.log.Error(err, "failed to list TenantResources", "namespace", namespaces.Items[i].Name)
continue
}
for j := range trList.Items {
key := types.NamespacedName{
Namespace: trList.Items[j].Namespace,
Name: trList.Items[j].Name,
}
if _, exists := seen[key]; exists {
continue
}
seen[key] = struct{}{}
requests = append(requests, reconcile.Request{NamespacedName: key})
}
}
return requests
}
//nolint:dupl
func (r *namespacedResourceController) enqueueAllResources(ctx context.Context, _ client.Object) []reconcile.Request {
var list capsulev1beta2.TenantResourceList
if err := r.client.List(ctx, &list); err != nil {
r.log.V(1).Error(err, "unable to list TenantResources for config-triggered reconcile")
return nil
}
reqs := make([]reconcile.Request, 0, len(list.Items))
for i := range list.Items {
reqs = append(reqs, reconcile.Request{
NamespacedName: types.NamespacedName{
Name: list.Items[i].Name,
Namespace: list.Items[i].Namespace,
},
})
}
return reqs
}
func (r *namespacedResourceController) reconcile(
ctx context.Context,
c client.Client,
tntResource *capsulev1beta2.TenantResource,
) error {
log := ctrllog.FromContext(ctx)
// Adding the default value for the status
if tntResource.Status.ProcessedItems == nil {
tntResource.Status.ProcessedItems = make([]meta.ObjectReferenceStatus, 0)
}
// Retrieving the parent of the Tenant Resource:
// can be owned, or being deployed in one of its Namespace.
// we cant resolve via status.namespaces, as when a namespace is deleted it is no longer references by the tenant
// causing a deletion blockade.
ns := &corev1.Namespace{}
if err := r.client.Get(ctx, types.NamespacedName{Name: tntResource.GetNamespace()}, ns); err != nil {
return err
}
tnt, err := tenant.GetTenantByOwnerreferences(ctx, r.client, ns.GetOwnerReferences())
if err != nil {
return err
}
if tnt == nil {
log.Info("skipping sync, the current Namespace is not belonging to any Tenant")
return nil
}
acc := processor.Accumulator{}
// Gather Resources
if tntResource.DeletionTimestamp.IsZero() {
err := r.gatherResources(
ctx,
c,
log,
tntResource,
*tnt,
acc,
)
if err != nil {
return err
}
}
return r.processor.Reconcile(
ctx,
log,
c,
&tntResource.Status.ProcessedItems,
acc,
processor.ProcessorOptions{
FieldOwnerPrefix: getFieldOwner(tntResource.GetName(), tntResource.GetNamespace()),
Prune: *tntResource.Spec.PruningOnDelete,
Adopt: *tntResource.Spec.Settings.Adopt,
Force: *tntResource.Spec.Settings.Force,
Owner: nil,
})
}
func (r *namespacedResourceController) gatherResources(
ctx context.Context,
c client.Client,
log logr.Logger,
tntResource *capsulev1beta2.TenantResource,
tnt capsulev1beta2.Tenant,
acc processor.Accumulator,
) (err error) {
opts := CollectorOptions{
Accumulator: acc,
AllowCrossNamespaceSelection: false,
AllowClusterScopedObjects: false,
ValidatorNamespaces: tpl.NewNamespaceValidator(false, sets.New[string](tnt.Status.Namespaces...)),
}
for resourceIndex, resource := range tntResource.Spec.Resources {
objs, err := r.collector.CollectNamespacedItems(ctx, c, opts, resource, &corev1.Namespace{ObjectMeta: metav1.ObjectMeta{Name: tntResource.GetNamespace()}}, tnt)
if err != nil {
return err
}
for g := range objs {
log.V(5).Info("found replication source object", "name", g.Name, "namespace", g.Namespace, "kind", g.Kind)
}
namespaces, err := r.collector.selectedTenantNamespaces(ctx, log, tnt, resource)
if err != nil {
return err
}
i := 0
for _, innerNs := range namespaces {
opts.Iterator = NewCollectorIteratorOptions(&tnt, innerNs, resource)
for _, obj := range objs {
if obj.GetNamespace() == innerNs.GetName() {
continue
}
target := obj.DeepCopy()
if err := sanitize.SanitizeObject(target, c.Scheme(), r.collector.objectSanitizeOptions); err != nil {
return err
}
target.SetNamespace(innerNs.GetName())
log.V(4).Info("adding replication for namespaced item", "name", target.GetName(), "namespace", target.GetNamespace(), "kind", target.GetKind())
err = r.collector.AddToAccumulation(&tnt, innerNs, opts, resource, target, "replica", false)
if err != nil {
return err
}
}
err = r.collector.Collect(
ctx,
c,
opts,
&tnt,
strconv.Itoa((resourceIndex)),
resource,
innerNs,
)
if err != nil {
return err
}
i++
}
}
return nil
}
func (r *namespacedResourceController) loadClient(
ctx context.Context,
log logr.Logger,
tntResource *capsulev1beta2.TenantResource,
) (client.Client, error) {
c, sa, err := r.clients.Load(ctx, log, tntResource)
// The resolved identity is posted to the status even along a failure, as it states
// which ServiceAccount the replication was attempted with.
tntResource.Status.ServiceAccount = sa
if err != nil {
return nil, err
}
return c, nil
}
func (r *namespacedResourceController) updateReconcilingStatus(ctx context.Context, instance *capsulev1beta2.TenantResource) error {
return retry.RetryOnConflict(retry.DefaultBackoff, func() (err error) {
latest := &capsulev1beta2.TenantResource{}
if err = r.reader.Get(ctx, types.NamespacedName{Name: instance.GetName(), Namespace: instance.GetNamespace()}, latest); err != nil {
return err
}
if latest.Status.ObservedGeneration == instance.GetGeneration() {
return nil
}
latest.Status.ServiceAccount = instance.Status.ServiceAccount
latest.Status.Conditions.UpdateConditionByType(meta.NewReadyConditionReconcilingReason(instance))
if err := r.client.Status().Update(ctx, latest); err != nil {
if apierrors.IsNotFound(err) {
return nil
}
return err
}
// Keep the in-memory object aligned with what we just wrote.
instance.Status = latest.Status
return nil
})
}
func (r *namespacedResourceController) updateStatus(ctx context.Context, instance *capsulev1beta2.TenantResource, reconcileError error) error {
instance.Status.UpdateStats()
return retry.RetryOnConflict(retry.DefaultBackoff, func() (err error) {
latest := &capsulev1beta2.TenantResource{}
if err = r.reader.Get(ctx, types.NamespacedName{Name: instance.GetName(), Namespace: instance.GetNamespace()}, latest); err != nil {
return err
}
originalStatus := latest.Status.DeepCopy()
latest.Status = instance.Status
latest.Status.ObservedGeneration = instance.GetGeneration()
// Set Ready Condition
readyCondition := meta.NewReadyCondition(instance)
if reconcileError != nil {
readyCondition.Message = reconcileError.Error()
readyCondition.Status = metav1.ConditionFalse
readyCondition.Reason = meta.FailedReason
}
latest.Status.Conditions.UpdateConditionByType(readyCondition)
// Set Cordoned Condition
cordonedCondition := meta.NewCordonedCondition(instance)
if *instance.Spec.Cordoned {
cordonedCondition.Reason = meta.CordonedReason
cordonedCondition.Message = "is cordoned"
cordonedCondition.Status = metav1.ConditionTrue
}
latest.Status.Conditions.UpdateConditionByType(cordonedCondition)
if reflect.DeepEqual(*originalStatus, latest.Status) {
return nil
}
if err := r.client.Status().Update(ctx, latest); err != nil {
return err
}
// Keep the in-memory object aligned with what we just wrote.
instance.Status = latest.Status
return nil
})
}
func ForeachNamespace(
ctx context.Context,
controllerClient client.Client,
resourceClient client.Client,
collector Collector,
opts CollectorOptions,
log logr.Logger,
resource capsulev1beta2.ResourceSpec,
resourceIndex int,
tnt capsulev1beta2.Tenant,
acc processor.Accumulator,
) (err error) {
namespaces, err := tenant.CollectTenantNamespaceByLabel(ctx, controllerClient, tnt, resource.NamespaceSelector)
if err != nil {
return err
}
for _, ns := range namespaces {
if ns.DeletionTimestamp != nil {
terminating, err := tenant.NamespaceIsPendingUnmanagedTerminationByStatus(ctx, controllerClient, &ns)
if err != nil {
return err
}
// Skip this namespace so resources are cleaned
if terminating {
continue
}
}
opts.Iterator = NewCollectorIteratorOptions(&tnt, &ns, resource)
objs, err := collector.CollectNamespacedItems(ctx, resourceClient, opts, resource, &ns, tnt)
if err != nil {
return err
}
i := 0
for _, obj := range objs {
for _, innerNs := range namespaces {
if obj.GetNamespace() == innerNs.GetName() {
continue
}
target := obj.DeepCopy()
target.SetNamespace(innerNs.GetName())
log.V(4).Info("adding replication for namespaced item", "name", target.GetName(), "namespace", target.GetNamespace(), "kind", target.GetKind())
err = collector.AddToAccumulation(&tnt, &innerNs, opts, resource, target, strconv.Itoa(resourceIndex)+"/replica-"+strconv.Itoa(i), true)
if err != nil {
return err
}
}
i++
}
err = collector.Collect(
ctx,
resourceClient,
opts,
&tnt,
strconv.Itoa((resourceIndex)),
resource,
&ns,
)
if err != nil {
return err
}
}
return nil
}