From 561345297bb991889a35b9c9c915b721713cb131 Mon Sep 17 00:00:00 2001 From: Super User Date: Mon, 27 Jul 2026 23:11:28 +0530 Subject: [PATCH 1/2] LOG-9591: do not overwrite ConfigMaps owned by other resources Refuse to update an existing ConfigMap when its ownerReferences do not match the ClusterLogForwarder. This prevents CLO from overwriting LokiStack (or other) ConfigMaps that share the same name, and surfaces the conflict as a deployment error on the forwarder. Co-authored-by: Cursor --- internal/reconcile/configmaps.go | 24 +++-- internal/reconcile/configmaps_test.go | 123 ++++++++++++++++++++++++++ 2 files changed, 141 insertions(+), 6 deletions(-) create mode 100644 internal/reconcile/configmaps_test.go diff --git a/internal/reconcile/configmaps.go b/internal/reconcile/configmaps.go index e638136f32..0d7e09a765 100644 --- a/internal/reconcile/configmaps.go +++ b/internal/reconcile/configmaps.go @@ -13,6 +13,10 @@ import ( "sigs.k8s.io/controller-runtime/pkg/client" ) +// Configmap creates or updates a ConfigMap owned by the ClusterLogForwarder. +// If a ConfigMap with the same name already exists and is not owned by the desired +// owner (for example, a LokiStack ConfigMap), it is left unchanged and an error is +// returned so the operator does not overwrite foreign resources (LOG-9591). func Configmap(k8Client client.Client, reader client.Reader, configMap *corev1.ConfigMap, opts ...comparators.ComparisonOption) error { return retry.RetryOnConflict(retry.DefaultRetry, func() error { current := &corev1.ConfigMap{} @@ -23,14 +27,22 @@ func Configmap(k8Client client.Client, reader client.Reader, configMap *corev1.C } return fmt.Errorf("failed to get %v configmap: %v", key, err) } - if configmaps.AreSame(current, configMap, opts...) && utils.HasSameOwner(current.OwnerReferences, configMap.OwnerReferences) { + + if !utils.HasSameOwner(current.OwnerReferences, configMap.OwnerReferences) { + return fmt.Errorf( + "configmap %s/%s already exists and is not owned by this ClusterLogForwarder; refusing to overwrite", + key.Namespace, key.Name, + ) + } + + if configmaps.AreSame(current, configMap, opts...) { return nil - } else { - current.Data = configMap.Data - current.Labels = configMap.Labels - current.Annotations = configMap.Annotations - current.OwnerReferences = configMap.OwnerReferences } + + current.Data = configMap.Data + current.Labels = configMap.Labels + current.Annotations = configMap.Annotations + current.OwnerReferences = configMap.OwnerReferences return k8Client.Update(context.TODO(), current) }) } diff --git a/internal/reconcile/configmaps_test.go b/internal/reconcile/configmaps_test.go new file mode 100644 index 0000000000..ddc6fb13dd --- /dev/null +++ b/internal/reconcile/configmaps_test.go @@ -0,0 +1,123 @@ +package reconcile_test + +import ( + "context" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" + "github.com/openshift/cluster-logging-operator/internal/reconcile" + "github.com/openshift/cluster-logging-operator/internal/runtime" + "github.com/openshift/cluster-logging-operator/internal/utils" + "github.com/openshift/cluster-logging-operator/internal/utils/comparators" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + k8sruntime "k8s.io/apimachinery/pkg/runtime" + "k8s.io/client-go/kubernetes/scheme" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" +) + +var _ = Describe("reconciling ConfigMap", func() { + + const ( + namespace = "openshift-logging" + name = "lokistack-config" + ) + + var ( + clfOwner = metav1.OwnerReference{ + APIVersion: "observability.openshift.io/v1", + Kind: "ClusterLogForwarder", + Name: "lokistack", + UID: "clf-uid", + Controller: utils.GetPtr(true), + } + lokiOwner = metav1.OwnerReference{ + APIVersion: "loki.grafana.com/v1", + Kind: "LokiStack", + Name: "lokistack", + UID: "loki-uid", + Controller: utils.GetPtr(true), + } + ) + + newClient := func(objs ...client.Object) client.Client { + globalScheme := k8sruntime.NewScheme() + Expect(scheme.AddToScheme(globalScheme)).To(Succeed()) + return fake.NewClientBuilder().WithScheme(globalScheme).WithObjects(objs...).Build() + } + + getConfigMap := func(k8sClient client.Client) *corev1.ConfigMap { + result := &corev1.ConfigMap{} + Expect(k8sClient.Get(context.TODO(), client.ObjectKey{Namespace: namespace, Name: name}, result)).To(Succeed()) + return result + } + + desiredCollectorCM := func() *corev1.ConfigMap { + cm := runtime.NewConfigMap(namespace, name, map[string]string{ + "vector.toml": "data_dir = \"/var/lib/vector\"", + }) + utils.AddOwnerRefToObject(cm, clfOwner) + return cm + } + + It("should create the ConfigMap when it does not exist", func() { + k8sClient := newClient() + desired := desiredCollectorCM() + + Expect(reconcile.Configmap(k8sClient, k8sClient, desired, comparators.CompareLabels)).To(Succeed()) + + result := getConfigMap(k8sClient) + Expect(result.Data).To(HaveKey("vector.toml")) + Expect(result.OwnerReferences).To(Equal([]metav1.OwnerReference{clfOwner})) + }) + + It("should update the ConfigMap when owned by the ClusterLogForwarder", func() { + existing := runtime.NewConfigMap(namespace, name, map[string]string{ + "vector.toml": "old", + }) + utils.AddOwnerRefToObject(existing, clfOwner) + k8sClient := newClient(existing) + + desired := desiredCollectorCM() + Expect(reconcile.Configmap(k8sClient, k8sClient, desired, comparators.CompareLabels)).To(Succeed()) + + result := getConfigMap(k8sClient) + Expect(result.Data["vector.toml"]).To(Equal("data_dir = \"/var/lib/vector\"")) + Expect(result.OwnerReferences).To(Equal([]metav1.OwnerReference{clfOwner})) + }) + + It("should refuse to overwrite a ConfigMap owned by another resource", func() { + existing := runtime.NewConfigMap(namespace, name, map[string]string{ + "config.yaml": "auth_enabled: false", + }) + utils.AddOwnerRefToObject(existing, lokiOwner) + k8sClient := newClient(existing) + + desired := desiredCollectorCM() + err := reconcile.Configmap(k8sClient, k8sClient, desired, comparators.CompareLabels) + Expect(err).To(HaveOccurred()) + Expect(err.Error()).To(ContainSubstring("refusing to overwrite")) + + result := getConfigMap(k8sClient) + Expect(result.Data).To(HaveKey("config.yaml")) + Expect(result.Data).NotTo(HaveKey("vector.toml")) + Expect(result.OwnerReferences).To(Equal([]metav1.OwnerReference{lokiOwner})) + }) + + It("should refuse to take ownership of an existing ConfigMap with no owner", func() { + existing := runtime.NewConfigMap(namespace, name, map[string]string{ + "config.yaml": "auth_enabled: false", + }) + k8sClient := newClient(existing) + + desired := desiredCollectorCM() + err := reconcile.Configmap(k8sClient, k8sClient, desired, comparators.CompareLabels) + Expect(err).To(HaveOccurred()) + Expect(err.Error()).To(ContainSubstring("refusing to overwrite")) + + result := getConfigMap(k8sClient) + Expect(result.Data).To(HaveKey("config.yaml")) + Expect(result.OwnerReferences).To(BeEmpty()) + }) +}) From 4c7fa7c4b28a590296d91e2270f438821630b860 Mon Sep 17 00:00:00 2001 From: Super User Date: Thu, 30 Jul 2026 22:55:26 +0530 Subject: [PATCH 2/2] LOG-9591: refuse overwrite of foreign-owned resources Extend ownership checks beyond ConfigMaps to DaemonSet, Deployment, Service, ServiceAccount, ServiceMonitor, NetworkPolicy, Role/RoleBinding, ClusterRoleBinding, and trusted-CA ConfigMap. Shared helper matches owners by UID and surfaces a clear refuse-to-overwrite error. Co-authored-by: Cursor --- internal/collector/trust_bundle.go | 7 +- internal/reconcile/configmaps.go | 12 +-- internal/reconcile/configmaps_test.go | 30 ++++++++ internal/reconcile/daemon_sets.go | 5 +- .../reconcile/daemonset_ownership_test.go | 73 +++++++++++++++++++ internal/reconcile/deployment.go | 5 +- internal/reconcile/networkpolicy.go | 4 + internal/reconcile/rbac.go | 14 +++- internal/reconcile/rbac_test.go | 16 ++++ internal/reconcile/service.go | 4 + internal/reconcile/service_account.go | 4 + internal/reconcile/service_monitor.go | 4 + internal/utils/ownership.go | 56 ++++++++++++++ internal/utils/ownership_test.go | 56 ++++++++++++++ 14 files changed, 278 insertions(+), 12 deletions(-) create mode 100644 internal/reconcile/daemonset_ownership_test.go create mode 100644 internal/utils/ownership.go create mode 100644 internal/utils/ownership_test.go diff --git a/internal/collector/trust_bundle.go b/internal/collector/trust_bundle.go index b68b155ea7..14bf1ec150 100644 --- a/internal/collector/trust_bundle.go +++ b/internal/collector/trust_bundle.go @@ -8,6 +8,7 @@ import ( log "github.com/ViaQ/logerr/v2/log/static" "github.com/openshift/cluster-logging-operator/internal/constants" "github.com/openshift/cluster-logging-operator/internal/runtime" + "github.com/openshift/cluster-logging-operator/internal/utils" corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/util/wait" @@ -23,10 +24,14 @@ var ( // ReconcileTrustedCABundleConfigMap creates or returns an existing Trusted CA Bundle ConfigMap. // By setting label "config.openshift.io/inject-trusted-cabundle: true", the cert is automatically filled/updated. func ReconcileTrustedCABundleConfigMap(k8sClient client.Client, namespace, name string, owner metav1.OwnerReference) error { + desiredOwners := []metav1.OwnerReference{owner} cm := runtime.NewConfigMap(namespace, name, nil) op, err := controllerutil.CreateOrUpdate(context.TODO(), k8sClient, cm, func() error { + if err := utils.EnsureCanUpdateOwnedResource(cm, desiredOwners); err != nil { + return err + } cm.Labels = map[string]string{constants.InjectTrustedCABundleLabel: "true"} - cm.OwnerReferences = []metav1.OwnerReference{owner} + cm.OwnerReferences = desiredOwners return nil }) diff --git a/internal/reconcile/configmaps.go b/internal/reconcile/configmaps.go index 0d7e09a765..5038a3fee4 100644 --- a/internal/reconcile/configmaps.go +++ b/internal/reconcile/configmaps.go @@ -13,10 +13,9 @@ import ( "sigs.k8s.io/controller-runtime/pkg/client" ) -// Configmap creates or updates a ConfigMap owned by the ClusterLogForwarder. +// Configmap creates or updates a ConfigMap owned by the desired ownerReferences. // If a ConfigMap with the same name already exists and is not owned by the desired -// owner (for example, a LokiStack ConfigMap), it is left unchanged and an error is -// returned so the operator does not overwrite foreign resources (LOG-9591). +// owners, it is left unchanged and an error is returned (LOG-9591). func Configmap(k8Client client.Client, reader client.Reader, configMap *corev1.ConfigMap, opts ...comparators.ComparisonOption) error { return retry.RetryOnConflict(retry.DefaultRetry, func() error { current := &corev1.ConfigMap{} @@ -28,11 +27,8 @@ func Configmap(k8Client client.Client, reader client.Reader, configMap *corev1.C return fmt.Errorf("failed to get %v configmap: %v", key, err) } - if !utils.HasSameOwner(current.OwnerReferences, configMap.OwnerReferences) { - return fmt.Errorf( - "configmap %s/%s already exists and is not owned by this ClusterLogForwarder; refusing to overwrite", - key.Namespace, key.Name, - ) + if err := utils.EnsureCanUpdateOwnedResource(current, configMap.OwnerReferences); err != nil { + return err } if configmaps.AreSame(current, configMap, opts...) { diff --git a/internal/reconcile/configmaps_test.go b/internal/reconcile/configmaps_test.go index ddc6fb13dd..4f929acfed 100644 --- a/internal/reconcile/configmaps_test.go +++ b/internal/reconcile/configmaps_test.go @@ -105,6 +105,18 @@ var _ = Describe("reconciling ConfigMap", func() { Expect(result.OwnerReferences).To(Equal([]metav1.OwnerReference{lokiOwner})) }) + It("should refuse an identical ConfigMap owned by another resource", func() { + existing := runtime.NewConfigMap(namespace, name, map[string]string{ + "vector.toml": "data_dir = \"/var/lib/vector\"", + }) + utils.AddOwnerRefToObject(existing, lokiOwner) + k8sClient := newClient(existing) + + err := reconcile.Configmap(k8sClient, k8sClient, desiredCollectorCM(), comparators.CompareLabels) + Expect(err).To(HaveOccurred()) + Expect(err.Error()).To(ContainSubstring("refusing to overwrite")) + }) + It("should refuse to take ownership of an existing ConfigMap with no owner", func() { existing := runtime.NewConfigMap(namespace, name, map[string]string{ "config.yaml": "auth_enabled: false", @@ -120,4 +132,22 @@ var _ = Describe("reconciling ConfigMap", func() { Expect(result.Data).To(HaveKey("config.yaml")) Expect(result.OwnerReferences).To(BeEmpty()) }) + + It("should allow update when the desired owner UID is present among additional owners", func() { + existing := runtime.NewConfigMap(namespace, name, map[string]string{ + "vector.toml": "old", + }) + utils.AddOwnerRefToObject(existing, clfOwner) + utils.AddOwnerRefToObject(existing, metav1.OwnerReference{ + APIVersion: "v1", + Kind: "ConfigMap", + Name: "extra", + UID: "extra-uid", + }) + k8sClient := newClient(existing) + + Expect(reconcile.Configmap(k8sClient, k8sClient, desiredCollectorCM(), comparators.CompareLabels)).To(Succeed()) + result := getConfigMap(k8sClient) + Expect(result.Data["vector.toml"]).To(Equal("data_dir = \"/var/lib/vector\"")) + }) }) diff --git a/internal/reconcile/daemon_sets.go b/internal/reconcile/daemon_sets.go index 792a77e7a5..7902f8da4e 100644 --- a/internal/reconcile/daemon_sets.go +++ b/internal/reconcile/daemon_sets.go @@ -6,6 +6,7 @@ import ( log "github.com/ViaQ/logerr/v2/log/static" "github.com/openshift/cluster-logging-operator/internal/runtime" + "github.com/openshift/cluster-logging-operator/internal/utils" apps "k8s.io/api/apps/v1" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" @@ -16,7 +17,9 @@ import ( func DaemonSet(k8Client client.Client, desired *apps.DaemonSet) error { ds := runtime.NewDaemonSet(desired.Namespace, desired.Name) op, err := controllerutil.CreateOrUpdate(context.TODO(), k8Client, ds, func() error { - // Update the daemonset with our desired state + if err := utils.EnsureCanUpdateOwnedResource(ds, desired.OwnerReferences); err != nil { + return err + } ds.Labels = desired.Labels ds.Spec = desired.Spec ds.OwnerReferences = desired.OwnerReferences diff --git a/internal/reconcile/daemonset_ownership_test.go b/internal/reconcile/daemonset_ownership_test.go new file mode 100644 index 0000000000..c94a1a1b22 --- /dev/null +++ b/internal/reconcile/daemonset_ownership_test.go @@ -0,0 +1,73 @@ +package reconcile_test + +import ( + "context" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" + "github.com/openshift/cluster-logging-operator/internal/reconcile" + "github.com/openshift/cluster-logging-operator/internal/runtime" + "github.com/openshift/cluster-logging-operator/internal/utils" + appsv1 "k8s.io/api/apps/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + k8sruntime "k8s.io/apimachinery/pkg/runtime" + "k8s.io/client-go/kubernetes/scheme" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" +) + +var _ = Describe("reconciling DaemonSet ownership", func() { + const ( + namespace = "openshift-logging" + name = "lokistack" + ) + + var ( + clfOwner = metav1.OwnerReference{ + APIVersion: "observability.openshift.io/v1", + Kind: "ClusterLogForwarder", + Name: "lokistack", + UID: "clf-uid", + Controller: utils.GetPtr(true), + } + lokiOwner = metav1.OwnerReference{ + APIVersion: "loki.grafana.com/v1", + Kind: "LokiStack", + Name: "lokistack", + UID: "loki-uid", + Controller: utils.GetPtr(true), + } + ) + + newClient := func(objs ...client.Object) client.Client { + globalScheme := k8sruntime.NewScheme() + Expect(scheme.AddToScheme(globalScheme)).To(Succeed()) + Expect(appsv1.AddToScheme(globalScheme)).To(Succeed()) + return fake.NewClientBuilder().WithScheme(globalScheme).WithObjects(objs...).Build() + } + + desiredDS := func() *appsv1.DaemonSet { + ds := runtime.NewDaemonSet(namespace, name) + ds.Spec.Selector = &metav1.LabelSelector{MatchLabels: map[string]string{"app": "collector"}} + ds.Spec.Template.Labels = map[string]string{"app": "collector"} + utils.AddOwnerRefToObject(ds, clfOwner) + return ds + } + + It("should refuse to overwrite a DaemonSet owned by another resource", func() { + existing := runtime.NewDaemonSet(namespace, name) + existing.Spec.Selector = &metav1.LabelSelector{MatchLabels: map[string]string{"app": "loki"}} + existing.Spec.Template.Labels = map[string]string{"app": "loki"} + utils.AddOwnerRefToObject(existing, lokiOwner) + k8sClient := newClient(existing) + + err := reconcile.DaemonSet(k8sClient, desiredDS()) + Expect(err).To(HaveOccurred()) + Expect(err.Error()).To(ContainSubstring("refusing to overwrite")) + + got := &appsv1.DaemonSet{} + Expect(k8sClient.Get(context.TODO(), client.ObjectKey{Namespace: namespace, Name: name}, got)).To(Succeed()) + Expect(got.Spec.Selector.MatchLabels["app"]).To(Equal("loki")) + Expect(got.OwnerReferences).To(Equal([]metav1.OwnerReference{lokiOwner})) + }) +}) diff --git a/internal/reconcile/deployment.go b/internal/reconcile/deployment.go index 6c094aad1d..3a6bb37e1f 100644 --- a/internal/reconcile/deployment.go +++ b/internal/reconcile/deployment.go @@ -6,6 +6,7 @@ import ( log "github.com/ViaQ/logerr/v2/log/static" "github.com/openshift/cluster-logging-operator/internal/runtime" + "github.com/openshift/cluster-logging-operator/internal/utils" apps "k8s.io/api/apps/v1" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" @@ -16,7 +17,9 @@ import ( func Deployment(k8Client client.Client, desired *apps.Deployment) error { dpl := runtime.NewDeployment(desired.Namespace, desired.Name) op, err := controllerutil.CreateOrUpdate(context.TODO(), k8Client, dpl, func() error { - // Update the deployment with our desired state + if err := utils.EnsureCanUpdateOwnedResource(dpl, desired.OwnerReferences); err != nil { + return err + } dpl.Labels = desired.Labels dpl.Spec = desired.Spec dpl.OwnerReferences = desired.OwnerReferences diff --git a/internal/reconcile/networkpolicy.go b/internal/reconcile/networkpolicy.go index 162872d248..cfbc605058 100644 --- a/internal/reconcile/networkpolicy.go +++ b/internal/reconcile/networkpolicy.go @@ -6,6 +6,7 @@ import ( log "github.com/ViaQ/logerr/v2/log/static" "github.com/openshift/cluster-logging-operator/internal/runtime" + "github.com/openshift/cluster-logging-operator/internal/utils" networkingv1 "k8s.io/api/networking/v1" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" @@ -16,6 +17,9 @@ import ( func NetworkPolicy(k8Client client.Client, desired *networkingv1.NetworkPolicy) error { np := runtime.NewNetworkPolicy(desired.Namespace, desired.Name) op, err := controllerutil.CreateOrUpdate(context.TODO(), k8Client, np, func() error { + if err := utils.EnsureCanUpdateOwnedResource(np, desired.OwnerReferences); err != nil { + return err + } np.Labels = desired.Labels np.Spec = desired.Spec np.OwnerReferences = desired.OwnerReferences diff --git a/internal/reconcile/rbac.go b/internal/reconcile/rbac.go index 7ab9371d5f..02afca64f9 100644 --- a/internal/reconcile/rbac.go +++ b/internal/reconcile/rbac.go @@ -6,6 +6,7 @@ import ( log "github.com/ViaQ/logerr/v2/log/static" "github.com/openshift/cluster-logging-operator/internal/runtime" + "github.com/openshift/cluster-logging-operator/internal/utils" rbacv1 "k8s.io/api/rbac/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" "sigs.k8s.io/controller-runtime/pkg/client" @@ -15,7 +16,9 @@ import ( func Role(k8Client client.Client, desired *rbacv1.Role) error { role := runtime.NewRole(desired.Namespace, desired.Name) op, err := controllerutil.CreateOrUpdate(context.TODO(), k8Client, role, func() error { - // Update the role with our desired state + if err := utils.EnsureCanUpdateOwnedResource(role, desired.OwnerReferences); err != nil { + return err + } role.Rules = desired.Rules role.OwnerReferences = desired.OwnerReferences return nil @@ -38,6 +41,10 @@ func RoleBinding(k8Client client.Client, desired *rbacv1.RoleBinding) error { return err } + if err := utils.EnsureCanUpdateOwnedResource(existing, desired.OwnerReferences); err != nil { + return err + } + if existing.RoleRef != desired.RoleRef { log.V(3).Info("Deleting roleBinding due to roleRef change", "name", desired.Name, "namespace", desired.Namespace) if err := k8Client.Delete(context.TODO(), existing); err != nil { @@ -65,6 +72,10 @@ func ClusterRoleBinding(k8sClient client.Client, name string, generator func() * return err } + if err := utils.EnsureCanUpdateOwnedResource(existing, desired.OwnerReferences); err != nil { + return err + } + if existing.RoleRef != desired.RoleRef { log.V(3).Info("Deleting clusterRoleBinding due to roleRef change", "name", name) if err := k8sClient.Delete(context.TODO(), existing); err != nil { @@ -75,6 +86,7 @@ func ClusterRoleBinding(k8sClient client.Client, name string, generator func() * } existing.Subjects = desired.Subjects + existing.OwnerReferences = desired.OwnerReferences log.V(3).Info("Updating clusterRoleBinding", "name", name) return k8sClient.Update(context.TODO(), existing) } diff --git a/internal/reconcile/rbac_test.go b/internal/reconcile/rbac_test.go index 52f39b3c37..85cbb89850 100644 --- a/internal/reconcile/rbac_test.go +++ b/internal/reconcile/rbac_test.go @@ -55,6 +55,7 @@ var _ = Describe("reconciling RoleBinding", func() { It("should update Subjects and OwnerReferences when roleRef is unchanged", func() { existing := runtime.NewRoleBinding(namespace, name, newRef, rbacv1.Subject{Kind: "ServiceAccount", Name: "old-sa", Namespace: namespace}) + existing.OwnerReferences = []metav1.OwnerReference{ownerRef} k8sClient := newClient(existing) desired := runtime.NewRoleBinding(namespace, name, newRef, subject) @@ -71,6 +72,7 @@ var _ = Describe("reconciling RoleBinding", func() { It("should delete and recreate the RoleBinding when roleRef changes", func() { existing := runtime.NewRoleBinding(namespace, name, oldRef, rbacv1.Subject{Kind: "ServiceAccount", Name: "old-sa", Namespace: namespace}) + existing.OwnerReferences = []metav1.OwnerReference{ownerRef} k8sClient := newClient(existing) desired := runtime.NewRoleBinding(namespace, name, newRef, subject) @@ -83,4 +85,18 @@ var _ = Describe("reconciling RoleBinding", func() { Expect(result.Subjects).To(Equal([]rbacv1.Subject{subject})) Expect(result.OwnerReferences).To(Equal([]metav1.OwnerReference{ownerRef})) }) + + It("should refuse to overwrite a RoleBinding owned by another resource", func() { + foreignOwner := metav1.OwnerReference{APIVersion: "v1", Kind: "ConfigMap", Name: "other", UID: "other-uid"} + existing := runtime.NewRoleBinding(namespace, name, newRef, subject) + existing.OwnerReferences = []metav1.OwnerReference{foreignOwner} + k8sClient := newClient(existing) + + desired := runtime.NewRoleBinding(namespace, name, newRef, subject) + desired.OwnerReferences = []metav1.OwnerReference{ownerRef} + + err := reconcile.RoleBinding(k8sClient, desired) + Expect(err).To(HaveOccurred()) + Expect(err.Error()).To(ContainSubstring("refusing to overwrite")) + }) }) diff --git a/internal/reconcile/service.go b/internal/reconcile/service.go index c272e33f68..4fec10cb58 100644 --- a/internal/reconcile/service.go +++ b/internal/reconcile/service.go @@ -6,6 +6,7 @@ import ( log "github.com/ViaQ/logerr/v2/log/static" "github.com/openshift/cluster-logging-operator/internal/runtime" + "github.com/openshift/cluster-logging-operator/internal/utils" corev1 "k8s.io/api/core/v1" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" @@ -16,6 +17,9 @@ import ( func Service(k8Client client.Client, desired *corev1.Service) error { sm := runtime.NewService(desired.Namespace, desired.Name) op, err := controllerutil.CreateOrUpdate(context.TODO(), k8Client, sm, func() error { + if err := utils.EnsureCanUpdateOwnedResource(sm, desired.OwnerReferences); err != nil { + return err + } // Set annotations upon creation if sm.CreationTimestamp.IsZero() { diff --git a/internal/reconcile/service_account.go b/internal/reconcile/service_account.go index ef0ba12a85..144ed81e29 100644 --- a/internal/reconcile/service_account.go +++ b/internal/reconcile/service_account.go @@ -6,6 +6,7 @@ import ( log "github.com/ViaQ/logerr/v2/log/static" "github.com/openshift/cluster-logging-operator/internal/runtime" + "github.com/openshift/cluster-logging-operator/internal/utils" v1 "k8s.io/api/core/v1" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" @@ -14,6 +15,9 @@ import ( func ServiceAccount(k8Client client.Client, desired *v1.ServiceAccount) (*v1.ServiceAccount, error) { sa := runtime.NewServiceAccount(desired.Namespace, desired.Name) op, err := controllerutil.CreateOrUpdate(context.TODO(), k8Client, sa, func() error { + if err := utils.EnsureCanUpdateOwnedResource(sa, desired.OwnerReferences); err != nil { + return err + } sa.Annotations = desired.Annotations sa.OwnerReferences = desired.OwnerReferences sa.Finalizers = desired.Finalizers diff --git a/internal/reconcile/service_monitor.go b/internal/reconcile/service_monitor.go index 6c76821427..808019d27c 100644 --- a/internal/reconcile/service_monitor.go +++ b/internal/reconcile/service_monitor.go @@ -6,6 +6,7 @@ import ( log "github.com/ViaQ/logerr/v2/log/static" "github.com/openshift/cluster-logging-operator/internal/runtime" + "github.com/openshift/cluster-logging-operator/internal/utils" monitoringv1 "github.com/prometheus-operator/prometheus-operator/pkg/apis/monitoring/v1" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" @@ -16,6 +17,9 @@ import ( func ServiceMonitor(k8Client client.Client, desired *monitoringv1.ServiceMonitor) error { sm := runtime.NewServiceMonitor(desired.Namespace, desired.Name) op, err := controllerutil.CreateOrUpdate(context.TODO(), k8Client, sm, func() error { + if err := utils.EnsureCanUpdateOwnedResource(sm, desired.OwnerReferences); err != nil { + return err + } sm.Labels = desired.Labels sm.Spec = desired.Spec sm.Annotations = desired.Annotations diff --git a/internal/utils/ownership.go b/internal/utils/ownership.go new file mode 100644 index 0000000000..85a3a39648 --- /dev/null +++ b/internal/utils/ownership.go @@ -0,0 +1,56 @@ +package utils + +import ( + "fmt" + + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +// EnsureCanUpdateOwnedResource returns nil when the object has not been persisted yet +// (create path) or when it is already owned by all of the desired owners. +// If the object exists and is not owned by the desired owners (foreign-owned or +// unowned), it returns an error and callers must not mutate the object. +// +// Ownership is matched by OwnerReference UID so additional ownerRefs on the +// object do not cause false conflicts (LOG-9591). +func EnsureCanUpdateOwnedResource(obj metav1.Object, desiredOwners []metav1.OwnerReference) error { + if obj.GetResourceVersion() == "" { + return nil + } + if IsOwnedByDesired(obj.GetOwnerReferences(), desiredOwners) { + return nil + } + return ResourceOwnershipConflictError(obj) +} + +// IsOwnedByDesired reports whether every desired owner UID is present among current owners. +// An empty desired owner list only matches when current owners are also empty. +func IsOwnedByDesired(current, desired []metav1.OwnerReference) bool { + if len(desired) == 0 { + return len(current) == 0 + } + currentUIDs := make(map[string]struct{}, len(current)) + for _, ref := range current { + currentUIDs[string(ref.UID)] = struct{}{} + } + for _, want := range desired { + if _, ok := currentUIDs[string(want.UID)]; !ok { + return false + } + } + return true +} + +// ResourceOwnershipConflictError builds a stable error for ownership conflicts. +func ResourceOwnershipConflictError(obj metav1.Object) error { + if ns := obj.GetNamespace(); ns != "" { + return fmt.Errorf( + "resource %s/%s already exists and is not owned by the expected owner; refusing to overwrite", + ns, obj.GetName(), + ) + } + return fmt.Errorf( + "resource %s already exists and is not owned by the expected owner; refusing to overwrite", + obj.GetName(), + ) +} diff --git a/internal/utils/ownership_test.go b/internal/utils/ownership_test.go new file mode 100644 index 0000000000..eec0025ce7 --- /dev/null +++ b/internal/utils/ownership_test.go @@ -0,0 +1,56 @@ +package utils_test + +import ( + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" + "github.com/openshift/cluster-logging-operator/internal/runtime" + "github.com/openshift/cluster-logging-operator/internal/utils" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +var _ = Describe("EnsureCanUpdateOwnedResource", func() { + var ( + clfOwner = metav1.OwnerReference{ + APIVersion: "observability.openshift.io/v1", + Kind: "ClusterLogForwarder", + Name: "lokistack", + UID: "clf-uid", + Controller: utils.GetPtr(true), + } + lokiOwner = metav1.OwnerReference{ + APIVersion: "loki.grafana.com/v1", + Kind: "LokiStack", + Name: "lokistack", + UID: "loki-uid", + Controller: utils.GetPtr(true), + } + ) + + It("should allow create when the object is not persisted", func() { + cm := runtime.NewConfigMap("openshift-logging", "lokistack-config", nil) + Expect(utils.EnsureCanUpdateOwnedResource(cm, []metav1.OwnerReference{clfOwner})).To(Succeed()) + }) + + It("should allow update when owned by the desired owner", func() { + cm := runtime.NewConfigMap("openshift-logging", "lokistack-config", nil) + cm.SetResourceVersion("1") + utils.AddOwnerRefToObject(cm, clfOwner) + Expect(utils.EnsureCanUpdateOwnedResource(cm, []metav1.OwnerReference{clfOwner})).To(Succeed()) + }) + + It("should refuse update when owned by another resource", func() { + cm := runtime.NewConfigMap("openshift-logging", "lokistack-config", nil) + cm.SetResourceVersion("1") + utils.AddOwnerRefToObject(cm, lokiOwner) + err := utils.EnsureCanUpdateOwnedResource(cm, []metav1.OwnerReference{clfOwner}) + Expect(err).To(HaveOccurred()) + Expect(err.Error()).To(ContainSubstring("refusing to overwrite")) + }) + + It("should refuse update when the object has no owner", func() { + cm := runtime.NewConfigMap("openshift-logging", "lokistack-config", nil) + cm.SetResourceVersion("1") + err := utils.EnsureCanUpdateOwnedResource(cm, []metav1.OwnerReference{clfOwner}) + Expect(err).To(HaveOccurred()) + }) +})