diff --git a/internal/collector/trust_bundle.go b/internal/collector/trust_bundle.go index b68b155ea..14bf1ec15 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 e638136f3..5038a3fee 100644 --- a/internal/reconcile/configmaps.go +++ b/internal/reconcile/configmaps.go @@ -13,6 +13,9 @@ import ( "sigs.k8s.io/controller-runtime/pkg/client" ) +// 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 +// 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{} @@ -23,14 +26,19 @@ 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 err := utils.EnsureCanUpdateOwnedResource(current, configMap.OwnerReferences); err != nil { + return err + } + + 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 000000000..4f929acfe --- /dev/null +++ b/internal/reconcile/configmaps_test.go @@ -0,0 +1,153 @@ +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 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", + }) + 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()) + }) + + 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 792a77e7a..7902f8da4 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 000000000..c94a1a1b2 --- /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 6c094aad1..3a6bb37e1 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 162872d24..cfbc60505 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 7ab9371d5..02afca64f 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 52f39b3c3..85cbb8985 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 c272e33f6..4fec10cb5 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 ef0ba12a8..144ed81e2 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 6c7682142..808019d27 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 000000000..85a3a3964 --- /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 000000000..eec0025ce --- /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()) + }) +})