Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
184 changes: 184 additions & 0 deletions controllers/clusterpromotion_copies_sweeper.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,184 @@
/*
Copyright 2026. projectsveltos.io. All rights reserved.

Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at

http://www.apache.org/licenses/LICENSE-2.0

Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/

package controllers

import (
"context"
"time"

"github.com/go-logr/logr"
corev1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
"sigs.k8s.io/controller-runtime/pkg/client"

configv1beta1 "github.com/projectsveltos/addon-controller/api/v1beta1"
libsveltosv1beta1 "github.com/projectsveltos/libsveltos/api/v1beta1"
logs "github.com/projectsveltos/libsveltos/lib/logsettings"
)

const (
// clusterProfileNameLabel is set by ClusterPromotion (Sveltos Enterprise) on the ConfigMaps/Secrets it
// creates for a stage, with the name of the ClusterProfile consuming them. The same key is used as
// annotation holding the full name (a label value is limited to 63 characters).
clusterProfileNameLabel = "config.projectsveltos.io/clusterprofilename"

// clusterPromotionCopiesSweepInterval is how often ConfigMaps/Secrets created by ClusterPromotion are
// checked
clusterPromotionCopiesSweepInterval = 5 * time.Minute

// clusterPromotionCopiesGracePeriod is how old a ConfigMap/Secret needs to be before it can be removed.
// ClusterPromotion creates them just before the ClusterProfile consuming them.
clusterPromotionCopiesGracePeriod = 5 * time.Minute
)

// ClusterPromotionCopiesSweeper periodically removes the ConfigMaps/Secrets ClusterPromotion created for
// a stage when they are not needed anymore: the ClusterProfile they were created for is gone, or does
// not reference them anymore. It is a safety net. Normally they are removed by the Kubernetes garbage
// collector when the ClusterPromotion is deleted, and by ClusterPromotion itself when it stops
// referencing them. It does not rely on the ClusterPromotion still existing.
type ClusterPromotionCopiesSweeper struct {
client.Client
Interval time.Duration
GracePeriod time.Duration
Logger logr.Logger
}

// NewClusterPromotionCopiesSweeper returns a sweeper using the default interval and grace period
func NewClusterPromotionCopiesSweeper(c client.Client, logger logr.Logger) *ClusterPromotionCopiesSweeper {
return &ClusterPromotionCopiesSweeper{
Client: c,
Interval: clusterPromotionCopiesSweepInterval,
GracePeriod: clusterPromotionCopiesGracePeriod,
Logger: logger,
}
}

// Start implements manager.Runnable. It runs until ctx is done. Added with manager.Add, it only runs
// on the leader and after the caches are started.
func (s *ClusterPromotionCopiesSweeper) Start(ctx context.Context) error {
ticker := time.NewTicker(s.Interval)
defer ticker.Stop()

for {
select {
case <-ctx.Done():
s.Logger.Info("stopping ClusterPromotion copies sweeper")
return nil
case <-ticker.C:
s.sweep(ctx)
}
}
}

func (s *ClusterPromotionCopiesSweeper) sweep(ctx context.Context) {
// ClusterProfiles already looked up during this sweep. A nil entry means not found.
clusterProfiles := map[string]*configv1beta1.ClusterProfile{}

listOptions := []client.ListOption{client.HasLabels{clusterProfileNameLabel}}

configMaps := &corev1.ConfigMapList{}
if err := s.List(ctx, configMaps, listOptions...); err != nil {
s.Logger.V(logs.LogInfo).Info("failed to list ConfigMaps", "error", err)
} else {
for i := range configMaps.Items {
s.removeIfStale(ctx, &configMaps.Items[i], string(libsveltosv1beta1.ConfigMapReferencedResourceKind),
clusterProfiles)
}
}

secrets := &corev1.SecretList{}
if err := s.List(ctx, secrets, listOptions...); err != nil {
s.Logger.V(logs.LogInfo).Info("failed to list Secrets", "error", err)
} else {
for i := range secrets.Items {
s.removeIfStale(ctx, &secrets.Items[i], string(libsveltosv1beta1.SecretReferencedResourceKind),
clusterProfiles)
}
}
}

func (s *ClusterPromotionCopiesSweeper) removeIfStale(ctx context.Context, obj client.Object, kind string,
clusterProfiles map[string]*configv1beta1.ClusterProfile) {

logger := s.Logger.WithValues("kind", kind, "namespace", obj.GetNamespace(), "name", obj.GetName())

stale, err := s.isStale(ctx, obj, kind, clusterProfiles)
if err != nil {
logger.V(logs.LogInfo).Info("failed to verify whether it is stale", "error", err)
return
}
if !stale {
return
}

logger.V(logs.LogInfo).Info("removing resource created by ClusterPromotion: ClusterProfile does not use it anymore")
uid := obj.GetUID()
err = s.Delete(ctx, obj, client.Preconditions{UID: &uid})
if err != nil && !apierrors.IsNotFound(err) {
logger.V(logs.LogInfo).Info("failed to delete", "error", err)
}
}

// isStale returns true if the ClusterProfile the resource was created for does not exist, or does
// not reference the resource anymore. A ClusterProfile being deleted still exists: the resource is
// kept until it is gone.
func (s *ClusterPromotionCopiesSweeper) isStale(ctx context.Context, obj client.Object, kind string,
clusterProfiles map[string]*configv1beta1.ClusterProfile) (bool, error) {

// The annotation holds the full ClusterProfile name. Resources without it are not ours.
clusterProfileName := obj.GetAnnotations()[clusterProfileNameLabel]
if clusterProfileName == "" {
return false, nil
}

if time.Since(obj.GetCreationTimestamp().Time) < s.GracePeriod {
return false, nil
}

clusterProfile, ok := clusterProfiles[clusterProfileName]
if !ok {
clusterProfile = &configv1beta1.ClusterProfile{}
err := s.Get(ctx, client.ObjectKey{Name: clusterProfileName}, clusterProfile)
if err != nil {
if !apierrors.IsNotFound(err) {
return false, err
}
clusterProfile = nil
}
clusterProfiles[clusterProfileName] = clusterProfile
}

if clusterProfile == nil {
return true, nil
}

references := getClusterPromotionReferences(&configv1beta1.ProfileSpec{
PolicyRefs: clusterProfile.Spec.PolicyRefs,
KustomizationRefs: clusterProfile.Spec.KustomizationRefs,
HelmCharts: clusterProfile.Spec.HelmCharts,
PatchesFrom: clusterProfile.Spec.PatchesFrom,
})
for i := range references {
if references[i].Kind == kind && references[i].Namespace == obj.GetNamespace() &&
references[i].Name == obj.GetName() {

return false, nil
}
}

return true, nil
}
133 changes: 133 additions & 0 deletions controllers/clusterpromotion_copies_sweeper_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,133 @@
/*
Copyright 2026. projectsveltos.io. All rights reserved.

Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at

http://www.apache.org/licenses/LICENSE-2.0

Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/

package controllers_test

import (
"context"
"time"

. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"

"github.com/go-logr/logr"
corev1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/client/fake"

configv1beta1 "github.com/projectsveltos/addon-controller/api/v1beta1"
"github.com/projectsveltos/addon-controller/controllers"
libsveltosv1beta1 "github.com/projectsveltos/libsveltos/api/v1beta1"
)

const (
copyClusterProfileLabel = "config.projectsveltos.io/clusterprofilename"
copyNamespace = "default"
)

var _ = Describe("ClusterPromotionCopiesSweeper", func() {
// getCopyConfigMap returns a ConfigMap as created by ClusterPromotion for a stage
getCopyConfigMap := func(clusterProfileName string, age time.Duration) *corev1.ConfigMap {
return &corev1.ConfigMap{
ObjectMeta: metav1.ObjectMeta{
Namespace: copyNamespace,
Name: randomString(),
CreationTimestamp: metav1.NewTime(time.Now().Add(-age)),
Labels: map[string]string{copyClusterProfileLabel: clusterProfileName},
Annotations: map[string]string{copyClusterProfileLabel: clusterProfileName},
},
}
}

getCopySecret := func(clusterProfileName string, age time.Duration) *corev1.Secret {
return &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{
Namespace: copyNamespace,
Name: randomString(),
CreationTimestamp: metav1.NewTime(time.Now().Add(-age)),
Labels: map[string]string{copyClusterProfileLabel: clusterProfileName},
Annotations: map[string]string{copyClusterProfileLabel: clusterProfileName},
},
}
}

exists := func(c client.Client, obj client.Object) bool {
err := c.Get(context.TODO(), client.ObjectKeyFromObject(obj), obj.DeepCopyObject().(client.Object))
if apierrors.IsNotFound(err) {
return false
}
Expect(err).To(BeNil())
return true
}

It("removes resources whose ClusterProfile does not exist, keeping the others", func() {
const oldAge = time.Hour

missingProfile := randomString()
configMapOfMissingProfile := getCopyConfigMap(missingProfile, oldAge)
secretOfMissingProfile := getCopySecret(missingProfile, oldAge)
// Younger than the grace period: ClusterPromotion creates it just before the ClusterProfile
youngConfigMap := getCopyConfigMap(missingProfile, time.Minute)
// Has the label but not the annotation ClusterPromotion sets: not created by ClusterPromotion
notOurs := getCopyConfigMap(missingProfile, oldAge)
notOurs.Annotations = nil

clusterProfile := &configv1beta1.ClusterProfile{
ObjectMeta: metav1.ObjectMeta{Name: randomString()},
}
configMapOfExistingProfile := getCopyConfigMap(clusterProfile.Name, oldAge)
secretOfExistingProfile := getCopySecret(clusterProfile.Name, oldAge)
clusterProfile.Spec.PolicyRefs = []configv1beta1.PolicyRef{
{
Namespace: copyNamespace, Name: configMapOfExistingProfile.Name,
Kind: string(libsveltosv1beta1.ConfigMapReferencedResourceKind),
},
}
clusterProfile.Spec.HelmCharts = []configv1beta1.HelmChart{
{
ValuesFrom: []configv1beta1.ValueFrom{
{
Namespace: copyNamespace, Name: secretOfExistingProfile.Name,
Kind: string(libsveltosv1beta1.SecretReferencedResourceKind),
},
},
},
}

// The ClusterProfile exists but does not reference the ConfigMap anymore
configMapNotReferenced := getCopyConfigMap(clusterProfile.Name, oldAge)

initObjects := []client.Object{
configMapOfMissingProfile, secretOfMissingProfile, youngConfigMap, notOurs,
clusterProfile, configMapOfExistingProfile, secretOfExistingProfile, configMapNotReferenced,
}
c := fake.NewClientBuilder().WithScheme(scheme).WithObjects(initObjects...).Build()

sweeper := controllers.NewClusterPromotionCopiesSweeper(c, logr.Discard())
controllers.SweepClusterPromotionCopies(sweeper, context.TODO())

Expect(exists(c, configMapOfMissingProfile)).To(BeFalse())
Expect(exists(c, secretOfMissingProfile)).To(BeFalse())
Expect(exists(c, configMapNotReferenced)).To(BeFalse())

Expect(exists(c, youngConfigMap)).To(BeTrue())
Expect(exists(c, notOurs)).To(BeTrue())
Expect(exists(c, configMapOfExistingProfile)).To(BeTrue())
Expect(exists(c, secretOfExistingProfile)).To(BeTrue())
})
})
1 change: 1 addition & 0 deletions controllers/export_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,7 @@ var (
ConvertResultStatus = (*ClusterSummaryReconciler).convertResultStatus
RequeueClusterSummaryForReference = (*ClusterSummaryReconciler).requeueClusterSummaryForReference
RequeueClusterPromotionForReference = (*ClusterPromotionReconciler).requeueClusterPromotionForReference
SweepClusterPromotionCopies = (*ClusterPromotionCopiesSweeper).sweep
ClusterPromotionUpdateMaps = (*ClusterPromotionReconciler).updateMaps
ClusterPromotionCleanMaps = (*ClusterPromotionReconciler).cleanMaps
RequeueClusterSummaryForCluster = (*ClusterSummaryReconciler).requeueClusterSummaryForCluster
Expand Down
13 changes: 13 additions & 0 deletions pkg/app/app.go
Original file line number Diff line number Diff line change
Expand Up @@ -199,6 +199,17 @@ func Run() {
}
}

// addClusterPromotionCopiesSweeper adds the runnable removing the ConfigMaps/Secrets ClusterPromotion
// creates per stage when the ClusterProfile consuming them is gone
func addClusterPromotionCopiesSweeper(mgr manager.Manager) {
sweeper := controllers.NewClusterPromotionCopiesSweeper(mgr.GetClient(),
ctrl.Log.WithName("clusterpromotion-copies-sweeper"))
if err := mgr.Add(sweeper); err != nil {
setupLog.Error(err, "unable to add runnable", "runnable", "ClusterPromotionCopiesSweeper")
os.Exit(1)
}
}

func getCacheConfig() (disableFor []client.Object, byObject map[client.Object]cache.ByObject) {
disableFor = []client.Object{}
byObject = map[client.Object]cache.ByObject{}
Expand Down Expand Up @@ -681,6 +692,8 @@ func startControllersAndWatchers(ctx context.Context, mgr manager.Manager) {
os.Exit(1)
}

addClusterPromotionCopiesSweeper(mgr)

// Needs a fleet-wide view of every ClusterSummary to dedup chart keys correctly, so
// this only ever runs on the default (unsharded) deployment, same as the reconcilers
// started above. Running it per-shard would give no benefit (chart-key dedup is
Expand Down
9 changes: 9 additions & 0 deletions test/fv/promotion_configmap_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,9 @@ var _ = Describe("Stage Promotions with referenced ConfigMap", func() {
// every couple of minutes so allow for a whole promotion.
restartTimeout = 10 * time.Minute

// label ClusterPromotion sets on the ConfigMaps/Secrets it creates for a stage
clusterProfileNameLabel = "config.projectsveltos.io/clusterprofilename"

payloadPolicy = `apiVersion: v1
kind: ConfigMap
metadata:
Expand Down Expand Up @@ -163,6 +166,12 @@ data:
content, err := getCopyContent(configMap.Namespace, copyName)
Expect(err).To(BeNil())
Expect(content).To(Equal(configMap.Data["policy0.yaml"]))

Byf("Verify the copy is labeled with the name of the ClusterProfile consuming it")
copyConfigMap := &corev1.ConfigMap{}
Expect(k8sClient.Get(context.TODO(), types.NamespacedName{Namespace: configMap.Namespace, Name: copyName},
copyConfigMap)).To(Succeed())
Expect(copyConfigMap.Labels).To(HaveKeyWithValue(clusterProfileNameLabel, clusterProfile.Name))
}

Byf("Verify the resource is deployed in the cluster matching stage %s", staging.Name)
Expand Down
Loading