From ef0fb14a124ae7dd9efc9a24b01215e1f8f46ba6 Mon Sep 17 00:00:00 2001 From: Todd Short Date: Wed, 23 Sep 2026 16:02:36 -0400 Subject: [PATCH 1/3] feat: support cross-namespace migration Signed-off-by: Todd Short --- .../cmd/migrate-operators-v0-to-v1/convert.go | 48 ++- .../convert_test.go | 19 + migration/pkg/migration/migration.go | 80 +++- migration/pkg/migration/namespace.go | 347 ++++++++++++++++++ migration/pkg/migration/unit_test.go | 237 +++++++++++- 5 files changed, 710 insertions(+), 21 deletions(-) create mode 100644 migration/pkg/migration/namespace.go diff --git a/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go b/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go index a48b128..e16a893 100644 --- a/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go +++ b/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go @@ -1,11 +1,14 @@ package main import ( + "context" "encoding/json" "errors" "fmt" + "time" "github.com/spf13/cobra" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "github.com/operator-framework/library-olm/migration/pkg/migration" ) @@ -70,6 +73,9 @@ func runConvert(cmd *cobra.Command, args []string) error { //nolint:nestif if !convertAll && len(args) == 0 { return fmt.Errorf("specify an operator name or --all") } + if convertAll && convertInstallNs != "" { + return fmt.Errorf("--install-namespace requires a single operator") + } c, restCfg, err := newClient() if err != nil { @@ -152,7 +158,6 @@ func runConvert(cmd *cobra.Command, args []string) error { //nolint:nestif AcknowledgeNotSteadyState: convertAckNotSteady, } opts.ApplyDefaults() - if convertDryRun { return runConvertDryRun(cmd, m, opts) } @@ -218,12 +223,20 @@ func runConvert(cmd *cobra.Command, args []string) error { //nolint:nestif return fmt.Errorf("ClusterObjectSet prerequisite check failed: %w", err) } success(fmt.Sprintf("ClusterObjectSet API established; using operator-controller namespace %s", opts.SystemNamespace)) + if err := m.PrepareInstallNamespace(ctx, opts); err != nil { + return fmt.Errorf("install namespace preparation failed: %w", err) + } stepHeader(4, "Collecting operator resources") objects, err := m.CollectResources(ctx, opts, csv, ip, bundleInfo.PackageName) if err != nil { return fmt.Errorf("failed to collect resources: %w", err) } + sourceObjects := make([]unstructured.Unstructured, len(objects)) + for i := range objects { + sourceObjects[i] = *objects[i].DeepCopy() + } + migration.RewriteInstallNamespace(objects, opts.SubscriptionNamespace, opts.InstallNamespace) bundleInfo.CollectedObjects = objects kindCounts := make(map[string]int) for _, obj := range objects { @@ -258,22 +271,41 @@ func runConvert(cmd *cobra.Command, args []string) error { //nolint:nestif success("Resources backed up in memory (CE annotation backup authoritative)") stepHeader(6, "Preparing operator for migration") - info("Deleting Subscription and CSV (orphan cascade — workloads keep running)...") + info("Deleting Subscription and CSV (orphan cascade)...") if err := m.PrepareForMigration(ctx, opts, csv); err != nil { return fmt.Errorf("preparation failed: %w", err) } success("OLMv0 management removed") + restoreSourceDeployments, err := m.ScaleSourceDeployments(ctx, sourceObjects, opts) + if err != nil { + return err + } + if opts.InstallNamespace != opts.SubscriptionNamespace { + success("Source Deployments scaled to zero before target cutover") + } stepHeader(7, "Creating OLMv1 migration resources") info(fmt.Sprintf("Applying COS %s-1 with %d objects and creating its ClusterExtension...", opts.ClusterExtensionName, len(bundleInfo.CollectedObjects))) startProgress() - if err := m.CreateMigrationResources(ctx, opts, bundleInfo, backup); err != nil { + result, err := m.CreateMigrationResources(ctx, opts, bundleInfo, backup) + if err != nil { clearProgress() + if result.TargetMayBeActive { + return fmt.Errorf("create migration resources: %w; target ClusterObjectSet may be active, so source Deployments remain scaled to zero", err) + } + restoreCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second) + defer cancel() + if restoreErr := restoreSourceDeployments(restoreCtx); restoreErr != nil { + return fmt.Errorf("create migration resources: %w; restore source Deployments: %v", err, restoreErr) + } return err } clearProgress() success(fmt.Sprintf("ClusterObjectSet %s-1 reached Succeeded=True", opts.ClusterExtensionName)) success(fmt.Sprintf("ClusterExtension %s is Installed", opts.ClusterExtensionName)) + if err := m.DeleteSourceNamespaceResources(ctx, sourceObjects, opts); err != nil { + return fmt.Errorf("delete source install resources: %w", err) + } stepHeader(8, "Cleaning up OLMv0 resources") cleanupResult := m.CleanupOLMv0Resources(ctx, opts, bundleInfo.PackageName, csv.Name) @@ -287,7 +319,9 @@ func runConvert(cmd *cobra.Command, args []string) error { //nolint:nestif success(action.Description) } } - + if err := cleanupResult.Err(); err != nil { + return fmt.Errorf("clean up OLMv0 resources: %w", err) + } banner(fmt.Sprintf("Migration complete! %s is now managed by OLMv1", bundleInfo.PackageName)) fmt.Println() return nil @@ -368,6 +402,12 @@ func dryRunCleanupPlan(opts migration.Options, info *migration.MigrationInfo) [] } else { lines = append(lines, "Retain OperatorGroup(s); --delete-operatorgroup was not specified.") } + if opts.InstallNamespace != opts.SubscriptionNamespace { + lines = append(lines, + fmt.Sprintf("Create or update install namespace %s with PSA/SCC labels copied from %s.", opts.InstallNamespace, opts.SubscriptionNamespace), + fmt.Sprintf("Move collected namespaced operator resources from %s to %s and delete their source copies after ClusterExtension installation.", opts.SubscriptionNamespace, opts.InstallNamespace)) + lines = append(lines, fmt.Sprintf("Retain source namespace %s.", opts.SubscriptionNamespace)) + } return lines } diff --git a/migration/examples/cmd/migrate-operators-v0-to-v1/convert_test.go b/migration/examples/cmd/migrate-operators-v0-to-v1/convert_test.go index 747f86a..3c2393d 100644 --- a/migration/examples/cmd/migrate-operators-v0-to-v1/convert_test.go +++ b/migration/examples/cmd/migrate-operators-v0-to-v1/convert_test.go @@ -39,3 +39,22 @@ func TestDryRunCleanupPlan(t *testing.T) { t.Fatalf("dry-run cleanup plan does not describe the default OperatorGroup behavior:\n%s", plan) } } + +func TestDryRunCleanupPlanNamespaceChange(t *testing.T) { + opts := migration.Options{ + SubscriptionName: "widget-operator", + SubscriptionNamespace: "operators", + InstallNamespace: "widget-system", + } + info := &migration.MigrationInfo{PackageName: "widgets", BundleName: "widgets.v1.2.3"} + plan := strings.Join(dryRunCleanupPlan(opts, info), "\n") + for _, expected := range []string{ + "Create or update install namespace widget-system with PSA/SCC labels copied from operators.", + "Move collected namespaced operator resources from operators to widget-system and delete their source copies after ClusterExtension installation.", + "Retain source namespace operators.", + } { + if !strings.Contains(plan, expected) { + t.Fatalf("dry-run cleanup plan missing %q:\n%s", expected, plan) + } + } +} diff --git a/migration/pkg/migration/migration.go b/migration/pkg/migration/migration.go index 1e9b8d9..cc35990 100644 --- a/migration/pkg/migration/migration.go +++ b/migration/pkg/migration/migration.go @@ -47,6 +47,13 @@ type createdMigrationResources struct { ownershipUnknown bool } +// MigrationResourcesResult records whether target resources may have started +// reconciling when CreateMigrationResources returns. Callers must not restart +// scaled source controllers while TargetMayBeActive is true. +type MigrationResourcesResult struct { + TargetMayBeActive bool +} + const migrationInvocationAnnotation = "olm.operatorframework.io/migration-invocation" // Migrate performs the full migration of an OLMv0-managed operator to OLMv1. @@ -105,6 +112,9 @@ func (m *Migrator) Migrate(ctx context.Context, opts Options) error { if err != nil { return err } + if err := m.PrepareInstallNamespace(ctx, opts); err != nil { + return err + } backup, err := m.BackupResources(ctx, opts, csv, ip) if err != nil { @@ -140,6 +150,11 @@ func (m *Migrator) Migrate(ctx context.Context, opts Options) error { if err != nil { return fmt.Errorf("failed to collect resources: %w", err) } + sourceObjects := make([]unstructured.Unstructured, len(objects)) + for i := range objects { + sourceObjects[i] = *objects[i].DeepCopy() + } + RewriteInstallNamespace(objects, opts.SubscriptionNamespace, opts.InstallNamespace) info.CollectedObjects = objects // R9: warn about TLS certificate pivot. OLMv0 manages certs directly via its own @@ -149,12 +164,29 @@ func (m *Migrator) Migrate(ctx context.Context, opts Options) error { m.progress("Note: TLS certificate management will transfer from OLMv0 to cert-manager/service-ca; " + "expect pod restarts while new cert secrets are provisioned") - if err := m.CreateMigrationResources(ctx, opts, info, backup); err != nil { + restoreSourceDeployments, err := m.ScaleSourceDeployments(ctx, sourceObjects, opts) + if err != nil { + return err + } + result, err := m.CreateMigrationResources(ctx, opts, info, backup) + if err != nil { + if result.TargetMayBeActive { + return fmt.Errorf("create migration resources: %w; target ClusterObjectSet may be active, so source Deployments remain scaled to zero", err) + } + restoreCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second) + defer cancel() + if restoreErr := restoreSourceDeployments(restoreCtx); restoreErr != nil { + return fmt.Errorf("create migration resources: %w; restore source Deployments: %v", err, restoreErr) + } + return err + } + if err := m.DeleteSourceNamespaceResources(ctx, sourceObjects, opts); err != nil { return err } - m.CleanupOLMv0Resources(ctx, opts, info.PackageName, csv.Name) - + if err := m.CleanupOLMv0Resources(ctx, opts, info.PackageName, csv.Name).Err(); err != nil { + return fmt.Errorf("clean up OLMv0 resources: %w", err) + } return nil } @@ -263,7 +295,8 @@ func (m *Migrator) BackupResources(ctx context.Context, opts Options, csv *opera } // PrepareForMigration removes OLMv0 management of the operator by deleting -// the Subscription and CSV with orphan cascading (operator workloads keep running). +// the Subscription and CSV with orphan cascading. A later cross-namespace +// cutover scales source Deployments down before creating target resources. func (m *Migrator) PrepareForMigration(ctx context.Context, opts Options, csv *operatorsv1alpha1.ClusterServiceVersion) error { // Delete Subscription with orphan cascading sub := &operatorsv1alpha1.Subscription{} @@ -331,9 +364,9 @@ func (m *Migrator) RecoverBeforeCE(ctx context.Context, opts Options, backup *Ba } // recoverCreatedMigrationResources removes objects created by this invocation -// in reverse dependency order, then restores the OLMv0 Subscription. A failure -// to delete a ClusterExtension or COS leaves the remaining resources intact and -// prevents restoration, avoiding concurrent OLMv0 and OLMv1 ownership. +// in reverse dependency order. Once a COS may have reconciled, it refuses to +// restore the OLMv0 Subscription automatically because orphaned target objects +// can still be running while deletion propagates. func (m *Migrator) recoverCreatedMigrationResources(ctx context.Context, opts Options, backup *Backup, resources *createdMigrationResources) error { recoveryCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), subWaitTimeout+30*time.Second) defer cancel() @@ -354,6 +387,9 @@ func (m *Migrator) recoverCreatedMigrationResources(ctx context.Context, opts Op } } cleanupErr := m.cleanupCreatedSecrets(recoveryCtx, resources.secrets) + if resources.cos != nil { + return errors.Join(cleanupErr, fmt.Errorf("target ClusterObjectSet may have reconciled; refusing automatic OLMv0 recovery")) + } recoverErr := m.RecoverFromBackup(recoveryCtx, opts, backup) return errors.Join(cleanupErr, recoverErr) } @@ -423,15 +459,16 @@ func (m *Migrator) CreateClusterObjectSet(ctx context.Context, opts Options, inf } // CreateMigrationResources creates the ClusterObjectSet and its ClusterExtension -// as one recoverable operation. If either creation fails, it removes only the -// objects created by this invocation before restoring the OLMv0 Subscription. -func (m *Migrator) CreateMigrationResources(ctx context.Context, opts Options, info *MigrationInfo, backup *Backup) error { +// as one recoverable operation. Its result reports whether a target COS may have +// started reconciling, including when cleanup was requested after a failure. +func (m *Migrator) CreateMigrationResources(ctx context.Context, opts Options, info *MigrationInfo, backup *Backup) (MigrationResourcesResult, error) { resources, err := m.createClusterObjectSet(ctx, opts, info) if err != nil { + result := MigrationResourcesResult{TargetMayBeActive: resources.cos != nil || resources.ownershipUnknown} if recoverErr := m.recoverCreatedMigrationResources(ctx, opts, backup, resources); recoverErr != nil { - return fmt.Errorf("COS creation failed: %w; recovery also failed: %v", err, recoverErr) + return result, fmt.Errorf("COS creation failed: %w; recovery also failed: %v", err, recoverErr) } - return fmt.Errorf("COS creation failed (recovered): %w", err) + return result, fmt.Errorf("COS creation failed (recovered): %w", err) } ce, ownershipUnknown, err := m.createClusterExtension(ctx, opts, info) @@ -440,12 +477,13 @@ func (m *Migrator) CreateMigrationResources(ctx context.Context, opts Options, i } resources.ownershipUnknown = resources.ownershipUnknown || ownershipUnknown if err != nil { + result := MigrationResourcesResult{TargetMayBeActive: resources.cos != nil || resources.ownershipUnknown} if recoverErr := m.recoverCreatedMigrationResources(ctx, opts, backup, resources); recoverErr != nil { - return fmt.Errorf("ClusterExtension creation failed: %w; recovery also failed: %v", err, recoverErr) + return result, fmt.Errorf("ClusterExtension creation failed: %w; recovery also failed: %v", err, recoverErr) } - return fmt.Errorf("ClusterExtension creation failed (recovered): %w", err) + return result, fmt.Errorf("ClusterExtension creation failed (recovered): %w", err) } - return nil + return MigrationResourcesResult{}, nil } func (m *Migrator) createClusterObjectSet(ctx context.Context, opts Options, info *MigrationInfo) (*createdMigrationResources, error) { @@ -663,7 +701,6 @@ func (m *Migrator) createClusterExtension(ctx context.Context, opts Options, inf if opts.AcknowledgeNotSteadyState { annotations[AnnotationAcknowledgedPrefix+"not-steady-state"] = "true" } - ce := &ocv1.ClusterExtension{ ObjectMeta: metav1.ObjectMeta{ Name: opts.ClusterExtensionName, @@ -779,6 +816,17 @@ type CleanupResult struct { Actions []CleanupAction } +// Err returns all failures encountered while cleaning up OLMv0 resources. +func (r *CleanupResult) Err() error { + var errs []error + for _, action := range r.Actions { + if action.Error != nil { + errs = append(errs, fmt.Errorf("%s: %w", action.Description, action.Error)) + } + } + return errors.Join(errs...) +} + // CleanupOLMv0Resources removes remaining OLMv0 resources after migration. func (m *Migrator) CleanupOLMv0Resources(ctx context.Context, opts Options, packageName, csvName string) *CleanupResult { result := &CleanupResult{} diff --git a/migration/pkg/migration/namespace.go b/migration/pkg/migration/namespace.go new file mode 100644 index 0000000..d25d90d --- /dev/null +++ b/migration/pkg/migration/namespace.go @@ -0,0 +1,347 @@ +package migration + +import ( + "context" + "errors" + "fmt" + "strings" + + appsv1 "k8s.io/api/apps/v1" + 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/apis/meta/v1/unstructured" + "k8s.io/client-go/util/retry" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +const sccPodSecurityLabelSync = "security.openshift.io/scc.podSecurityLabelSync" + +// sourceDeploymentReplica records a live source Deployment's desired scale for +// restoration if target creation cannot complete. +type sourceDeploymentReplica struct { + name string + namespace string + replicas int32 +} + +// PrepareInstallNamespace creates the requested install namespace when needed +// and copies the namespace labels that control pod security admission. It must +// run before OLMv0 management is removed so a failed preparation leaves the +// source installation untouched. +func (m *Migrator) PrepareInstallNamespace(ctx context.Context, opts Options) error { + if opts.InstallNamespace == opts.SubscriptionNamespace { + return nil + } + + var source corev1.Namespace + if err := m.Client.Get(ctx, client.ObjectKey{Name: opts.SubscriptionNamespace}, &source); err != nil { + return fmt.Errorf("get source namespace %q: %w", opts.SubscriptionNamespace, err) + } + + labels := securityNamespaceLabels(source.Labels) + var target corev1.Namespace + err := m.Client.Get(ctx, client.ObjectKey{Name: opts.InstallNamespace}, &target) + if apierrors.IsNotFound(err) { + target = corev1.Namespace{} + target.Name = opts.InstallNamespace + target.Labels = labels + if err := m.Client.Create(ctx, &target); err != nil { + return fmt.Errorf("create install namespace %q: %w", opts.InstallNamespace, err) + } + m.progress(fmt.Sprintf("Created install namespace %s with copied PSA/SCC labels", opts.InstallNamespace)) + return nil + } + if err != nil { + return fmt.Errorf("get install namespace %q: %w", opts.InstallNamespace, err) + } + if unsafePSAEnforcement(target.Labels, labels) { + return fmt.Errorf("source namespace PSA enforcement may weaken existing target namespace %q; choose a target with an explicit equal or weaker enforcement label", opts.InstallNamespace) + } + + changed := false + if target.Labels == nil && len(labels) > 0 { + target.Labels = map[string]string{} + } + for key, value := range labels { + if target.Labels[key] != value { + target.Labels[key] = value + changed = true + } + } + if changed { + if err := m.Client.Update(ctx, &target); err != nil { + return fmt.Errorf("copy PSA/SCC labels to install namespace %q: %w", opts.InstallNamespace, err) + } + m.progress(fmt.Sprintf("Copied PSA/SCC labels to existing install namespace %s", opts.InstallNamespace)) + } + return nil +} + +// securityNamespaceLabels returns the source labels that affect workload +// admission and must accompany an operator to its new install namespace. +func securityNamespaceLabels(labels map[string]string) map[string]string { + result := map[string]string{} + for key, value := range labels { + if strings.HasPrefix(key, "pod-security.kubernetes.io/") || key == sccPodSecurityLabelSync { + result[key] = value + } + } + return result +} + +// unsafePSAEnforcement reports whether copying source labels would weaken an +// existing target enforce level or replace an unknown cluster default. Kubernetes +// does not expose the effective PSA default through the workload API, so an +// existing target without an enforce label must be treated conservatively. +func unsafePSAEnforcement(targetLabels, sourceLabels map[string]string) bool { + target, targetSet := targetLabels["pod-security.kubernetes.io/enforce"] + source, sourceSet := sourceLabels["pod-security.kubernetes.io/enforce"] + if !sourceSet { + return false + } + if !targetSet { + return true + } + if target == source { + return false + } + levels := map[string]int{"privileged": 0, "baseline": 1, "restricted": 2} + return levels[source] < levels[target] +} + +// RewriteInstallNamespace moves collected namespaced objects from the +// Subscription namespace to the requested install namespace. Cluster-scoped +// resources, and namespaced objects outside the source namespace, are left +// unchanged. +func RewriteInstallNamespace(objects []unstructured.Unstructured, sourceNamespace, installNamespace string) { + if sourceNamespace == installNamespace { + return + } + services, serviceAccounts := migratedNames(objects, sourceNamespace) + for i := range objects { + if objects[i].GetNamespace() == sourceNamespace { + objects[i].SetNamespace(installNamespace) + clearServiceAllocations(&objects[i]) + } + rewriteServiceReferences(&objects[i], sourceNamespace, installNamespace, services, serviceAccounts) + } +} + +// migratedNames indexes source Services and ServiceAccounts whose references +// may need retargeting after their namespace changes. +func migratedNames(objects []unstructured.Unstructured, sourceNamespace string) (map[string]struct{}, map[string]struct{}) { + services := map[string]struct{}{} + serviceAccounts := map[string]struct{}{} + for i := range objects { + if objects[i].GetNamespace() != sourceNamespace { + continue + } + switch objects[i].GetKind() { + case "Service": + services[objects[i].GetName()] = struct{}{} + case "ServiceAccount": + serviceAccounts[objects[i].GetName()] = struct{}{} + } + } + return services, serviceAccounts +} + +// rewriteServiceReferences updates cluster-scoped references to Services that +// move with the operator. These references remain valid through cutover after +// the source Service is deleted. +func rewriteServiceReferences(obj *unstructured.Unstructured, sourceNamespace, installNamespace string, services, serviceAccounts map[string]struct{}) { + if obj.GetKind() == "RoleBinding" || obj.GetKind() == "ClusterRoleBinding" { + rewriteServiceAccountSubjects(obj, sourceNamespace, installNamespace, serviceAccounts) + } + switch obj.GetKind() { + case "ValidatingWebhookConfiguration", "MutatingWebhookConfiguration": + webhooks, found, err := unstructured.NestedSlice(obj.Object, "webhooks") + if err != nil || !found { + return + } + for i := range webhooks { + webhook, ok := webhooks[i].(map[string]interface{}) + if !ok { + continue + } + rewriteNestedServiceNamespace(webhook, sourceNamespace, installNamespace, services, "clientConfig", "service") + } + _ = unstructured.SetNestedSlice(obj.Object, webhooks, "webhooks") + case "CustomResourceDefinition": + rewriteNestedServiceNamespace(obj.Object, sourceNamespace, installNamespace, services, "spec", "conversion", "webhook", "clientConfig", "service") + } +} + +// rewriteServiceAccountSubjects updates explicit ServiceAccount subjects that +// refer to a collected account moved into the target namespace. +func rewriteServiceAccountSubjects(obj *unstructured.Unstructured, sourceNamespace, installNamespace string, serviceAccounts map[string]struct{}) { + subjects, found, err := unstructured.NestedSlice(obj.Object, "subjects") + if err != nil || !found { + return + } + for i := range subjects { + subject, ok := subjects[i].(map[string]interface{}) + if !ok || subject["kind"] != "ServiceAccount" || subject["namespace"] != sourceNamespace { + continue + } + name, _ := subject["name"].(string) + if _, found := serviceAccounts[name]; !found { + continue + } + subject["namespace"] = installNamespace + } + _ = unstructured.SetNestedSlice(obj.Object, subjects, "subjects") +} + +// rewriteNestedServiceNamespace retargets a collected Service reference when +// the reference's namespace and name both identify a moved Service. +func rewriteNestedServiceNamespace(object map[string]interface{}, sourceNamespace, installNamespace string, services map[string]struct{}, serviceFields ...string) { + if len(serviceFields) == 0 { + return + } + nameFields := append(append([]string{}, serviceFields...), "name") + namespaceFields := append(append([]string{}, serviceFields...), "namespace") + name, nameFound, nameErr := unstructured.NestedString(object, nameFields...) + namespace, namespaceFound, namespaceErr := unstructured.NestedString(object, namespaceFields...) + if nameErr == nil && namespaceErr == nil && nameFound && namespaceFound && namespace == sourceNamespace { + if _, found := services[name]; found { + _ = unstructured.SetNestedField(object, installNamespace, namespaceFields...) + } + } +} + +// clearServiceAllocations removes values Kubernetes allocated to a Service in +// its source namespace. clusterIP and nodePort are cluster-wide allocations; +// retaining either when moving a Service to another namespace makes the target +// Service conflict with the still-running source Service before cutover. +func clearServiceAllocations(obj *unstructured.Unstructured) { + if obj.GetAPIVersion() != "v1" || obj.GetKind() != "Service" { + return + } + + clusterIP, _, _ := unstructured.NestedString(obj.Object, "spec", "clusterIP") + if clusterIP != "None" { + unstructured.RemoveNestedField(obj.Object, "spec", "clusterIP") + unstructured.RemoveNestedField(obj.Object, "spec", "clusterIPs") + } + // A health-check NodePort is allocated by Kubernetes just like each Service + // port's nodePort. Keep the desired Service shape but allocate fresh ports. + unstructured.RemoveNestedField(obj.Object, "spec", "healthCheckNodePort") + ports, found, err := unstructured.NestedSlice(obj.Object, "spec", "ports") + if err != nil || !found { + return + } + for i := range ports { + port, ok := ports[i].(map[string]interface{}) + if !ok { + continue + } + delete(port, "nodePort") + } + _ = unstructured.SetNestedSlice(obj.Object, ports, "spec", "ports") +} + +// DeleteSourceNamespaceResources removes the old copies of objects that were +// applied in a different install namespace. It runs only after the target CE +// has reached Installed=True, so a failed handoff retains the OLMv0 workload. +// It intentionally does not delete the Namespace itself; that requires the +// separate AcknowledgeNamespaceDelete opt-in. +func (m *Migrator) DeleteSourceNamespaceResources(ctx context.Context, objects []unstructured.Unstructured, opts Options) error { + if opts.InstallNamespace == opts.SubscriptionNamespace { + return nil + } + deleted := 0 + for i := range objects { + obj := objects[i].DeepCopy() + if obj.GetNamespace() != opts.SubscriptionNamespace { + continue + } + uid := obj.GetUID() + if uid == "" { + return fmt.Errorf("refuse to delete source resource %s %s/%s without collected UID", obj.GetKind(), obj.GetNamespace(), obj.GetName()) + } + if err := m.Client.Delete(ctx, obj, client.Preconditions(metav1.Preconditions{UID: &uid})); err != nil && client.IgnoreNotFound(err) != nil { + if apierrors.IsConflict(err) { + return fmt.Errorf("source resource %s %s/%s changed since collection: %w", obj.GetKind(), obj.GetNamespace(), obj.GetName(), err) + } + return fmt.Errorf("delete source resource %s %s/%s: %w", obj.GetKind(), obj.GetNamespace(), obj.GetName(), err) + } + deleted++ + } + if deleted > 0 { + m.progress(fmt.Sprintf("Deleted %d migrated resource(s) from source namespace %s", deleted, opts.SubscriptionNamespace)) + } + return nil +} + +// ScaleSourceDeployments stops collected source Deployments before the target +// namespace starts them. It returns a restore function for use when target +// creation fails, so a failed cutover does not leave the source operator down. +func (m *Migrator) ScaleSourceDeployments(ctx context.Context, objects []unstructured.Unstructured, opts Options) (func(context.Context) error, error) { + noRestore := func(context.Context) error { return nil } + if opts.InstallNamespace == opts.SubscriptionNamespace { + return noRestore, nil + } + + var originals []sourceDeploymentReplica + for i := range objects { + obj := objects[i] + if obj.GetAPIVersion() != "apps/v1" || obj.GetKind() != "Deployment" || obj.GetNamespace() != opts.SubscriptionNamespace { + continue + } + var original *sourceDeploymentReplica + if err := retry.RetryOnConflict(retry.DefaultBackoff, func() error { + var deployment appsv1.Deployment + if err := m.Client.Get(ctx, client.ObjectKey{Namespace: obj.GetNamespace(), Name: obj.GetName()}, &deployment); err != nil { + return err + } + replicas := int32(1) + if deployment.Spec.Replicas != nil { + replicas = *deployment.Spec.Replicas + } + if replicas == 0 { + return nil + } + zero := int32(0) + deployment.Spec.Replicas = &zero + if err := m.Client.Update(ctx, &deployment); err != nil { + return err + } + original = &sourceDeploymentReplica{name: deployment.Name, namespace: deployment.Namespace, replicas: replicas} + return nil + }); err != nil { + if restoreErr := restoreSourceDeploymentReplicas(ctx, m, originals); restoreErr != nil { + return noRestore, fmt.Errorf("scale source Deployment %s/%s: %w; restore scaled Deployments: %v", obj.GetNamespace(), obj.GetName(), err, restoreErr) + } + return noRestore, fmt.Errorf("scale source Deployment %s/%s: %w", obj.GetNamespace(), obj.GetName(), err) + } + if original != nil { + originals = append(originals, *original) + } + } + + return func(restoreCtx context.Context) error { + return restoreSourceDeploymentReplicas(restoreCtx, m, originals) + }, nil +} + +// restoreSourceDeploymentReplicas returns scaled source Deployments to their +// pre-cutover replica counts. +func restoreSourceDeploymentReplicas(ctx context.Context, m *Migrator, originals []sourceDeploymentReplica) error { + var errs []error + for _, original := range originals { + if err := retry.RetryOnConflict(retry.DefaultBackoff, func() error { + var deployment appsv1.Deployment + if err := m.Client.Get(ctx, client.ObjectKey{Namespace: original.namespace, Name: original.name}, &deployment); err != nil { + return err + } + replicas := original.replicas + deployment.Spec.Replicas = &replicas + return m.Client.Update(ctx, &deployment) + }); err != nil { + errs = append(errs, fmt.Errorf("restore source Deployment %s/%s: %w", original.namespace, original.name, err)) + } + } + return errors.Join(errs...) +} diff --git a/migration/pkg/migration/unit_test.go b/migration/pkg/migration/unit_test.go index 0a89cd8..db054db 100644 --- a/migration/pkg/migration/unit_test.go +++ b/migration/pkg/migration/unit_test.go @@ -351,6 +351,238 @@ func TestOptionsReportsAndAnnotationFiltering(t *testing.T) { } } +func TestPrepareInstallNamespace(t *testing.T) { + ctx := context.Background() + source := &corev1.Namespace{ObjectMeta: metav1.ObjectMeta{ + Name: "source", + Labels: map[string]string{ + "pod-security.kubernetes.io/enforce": "restricted", + "pod-security.kubernetes.io/enforce-version": "latest", + "security.openshift.io/scc.podSecurityLabelSync": "true", + "unrelated.example.test/do-not-copy": "value", + }, + }} + m := migrationTestClient(t, source) + opts := Options{SubscriptionNamespace: "source", InstallNamespace: "target"} + if err := m.PrepareInstallNamespace(ctx, opts); err != nil { + t.Fatalf("PrepareInstallNamespace() error = %v", err) + } + var target corev1.Namespace + if err := m.Client.Get(ctx, client.ObjectKey{Name: "target"}, &target); err != nil { + t.Fatalf("get target namespace: %v", err) + } + if got := target.Labels["pod-security.kubernetes.io/enforce"]; got != "restricted" { + t.Fatalf("PSA label = %q, want restricted", got) + } + if got := target.Labels[sccPodSecurityLabelSync]; got != "true" { + t.Fatalf("SCC label = %q, want true", got) + } + if _, found := target.Labels["unrelated.example.test/do-not-copy"]; found { + t.Fatalf("target labels unexpectedly include unrelated source label: %#v", target.Labels) + } + + target.Labels["pod-security.kubernetes.io/enforce"] = "baseline" + if err := m.Client.Update(ctx, &target); err != nil { + t.Fatalf("update target namespace: %v", err) + } + if err := m.PrepareInstallNamespace(ctx, opts); err != nil { + t.Fatalf("PrepareInstallNamespace(existing) error = %v", err) + } + if err := m.Client.Get(ctx, client.ObjectKey{Name: "target"}, &target); err != nil { + t.Fatalf("get updated target namespace: %v", err) + } + if got := target.Labels["pod-security.kubernetes.io/enforce"]; got != "restricted" { + t.Fatalf("updated PSA label = %q, want restricted", got) + } +} + +func TestInstallNamespaceRewriteAndSourceResourceDeletion(t *testing.T) { + ctx := context.Background() + objects := []unstructured.Unstructured{ + {Object: map[string]interface{}{"apiVersion": "v1", "kind": "ConfigMap", "metadata": map[string]interface{}{"name": "in-source", "namespace": "source"}}}, + {Object: map[string]interface{}{"apiVersion": "v1", "kind": "ConfigMap", "metadata": map[string]interface{}{"name": "elsewhere", "namespace": "elsewhere"}}}, + {Object: map[string]interface{}{"apiVersion": "rbac.authorization.k8s.io/v1", "kind": "ClusterRole", "metadata": map[string]interface{}{"name": "cluster-role"}}}, + } + RewriteInstallNamespace(objects, "source", "target") + if got := objects[0].GetNamespace(); got != "target" { + t.Fatalf("source object namespace = %q, want target", got) + } + if got := objects[1].GetNamespace(); got != "elsewhere" { + t.Fatalf("other object namespace = %q, want unchanged", got) + } + if got := objects[2].GetNamespace(); got != "" { + t.Fatalf("cluster-scoped object namespace = %q, want empty", got) + } + service := unstructured.Unstructured{Object: map[string]interface{}{ + "apiVersion": "v1", "kind": "Service", + "metadata": map[string]interface{}{"name": "operator-metrics", "namespace": "source"}, + "spec": map[string]interface{}{ + "clusterIP": "10.96.0.42", + "clusterIPs": []interface{}{"10.96.0.42"}, + "healthCheckNodePort": int64(30001), + "ports": []interface{}{map[string]interface{}{"port": int64(8443), "nodePort": int64(30002)}}, + }, + }} + services := []unstructured.Unstructured{service} + RewriteInstallNamespace(services, "source", "target") + if _, found, _ := unstructured.NestedFieldNoCopy(services[0].Object, "spec", "clusterIP"); found { + t.Fatal("rewritten Service retained source clusterIP") + } + if _, found, _ := unstructured.NestedFieldNoCopy(services[0].Object, "spec", "clusterIPs"); found { + t.Fatal("rewritten Service retained source clusterIPs") + } + if _, found, _ := unstructured.NestedFieldNoCopy(services[0].Object, "spec", "healthCheckNodePort"); found { + t.Fatal("rewritten Service retained source healthCheckNodePort") + } + ports, found, err := unstructured.NestedSlice(services[0].Object, "spec", "ports") + if err != nil || !found { + t.Fatalf("rewritten Service ports = %#v, found=%t, err=%v", ports, found, err) + } + if _, found := ports[0].(map[string]interface{})["nodePort"]; found { + t.Fatal("rewritten Service retained source nodePort") + } + references := []unstructured.Unstructured{ + {Object: map[string]interface{}{"apiVersion": "v1", "kind": "Service", "metadata": map[string]interface{}{"name": "operator-webhook", "namespace": "source"}}}, + {Object: map[string]interface{}{"apiVersion": "v1", "kind": "ServiceAccount", "metadata": map[string]interface{}{"name": "operator", "namespace": "source"}}}, + {Object: map[string]interface{}{ + "apiVersion": "rbac.authorization.k8s.io/v1", "kind": "ClusterRoleBinding", + "subjects": []interface{}{ + map[string]interface{}{"kind": "ServiceAccount", "name": "operator", "namespace": "source"}, + map[string]interface{}{"kind": "User", "name": "unchanged", "namespace": "source"}, + }, + }}, + {Object: map[string]interface{}{ + "apiVersion": "admissionregistration.k8s.io/v1", "kind": "ValidatingWebhookConfiguration", + "webhooks": []interface{}{map[string]interface{}{"clientConfig": map[string]interface{}{"service": map[string]interface{}{"name": "operator-webhook", "namespace": "source"}}}}, + }}, + {Object: map[string]interface{}{ + "apiVersion": "apiextensions.k8s.io/v1", "kind": "CustomResourceDefinition", + "spec": map[string]interface{}{"conversion": map[string]interface{}{"webhook": map[string]interface{}{"clientConfig": map[string]interface{}{"service": map[string]interface{}{"name": "operator-webhook", "namespace": "source"}}}}}, + }}, + } + RewriteInstallNamespace(references, "source", "target") + subjects, _, _ := unstructured.NestedSlice(references[2].Object, "subjects") + if got := subjects[0].(map[string]interface{})["namespace"]; got != "target" { + t.Fatalf("ServiceAccount subject namespace = %q, want target", got) + } + if got := subjects[1].(map[string]interface{})["namespace"]; got != "source" { + t.Fatalf("non-ServiceAccount subject namespace = %q, want source", got) + } + webhooks, _, _ := unstructured.NestedSlice(references[3].Object, "webhooks") + webhookService := webhooks[0].(map[string]interface{})["clientConfig"].(map[string]interface{})["service"].(map[string]interface{}) + if got := webhookService["namespace"]; got != "target" { + t.Fatalf("webhook Service namespace = %q, want target", got) + } + conversionService, found, err := unstructured.NestedMap(references[4].Object, "spec", "conversion", "webhook", "clientConfig", "service") + if err != nil || !found { + t.Fatalf("CRD conversion Service = %#v, found=%t, err=%v", conversionService, found, err) + } + if got := conversionService["namespace"]; got != "target" { + t.Fatalf("CRD conversion Service namespace = %q, want target", got) + } + + sourceConfigMap := &unstructured.Unstructured{Object: map[string]interface{}{ + "apiVersion": "v1", "kind": "ConfigMap", "metadata": map[string]interface{}{"name": "operator-config", "namespace": "source", "uid": "source-config-uid"}, + }} + otherConfigMap := &unstructured.Unstructured{Object: map[string]interface{}{ + "apiVersion": "v1", "kind": "ConfigMap", "metadata": map[string]interface{}{"name": "unrelated", "namespace": "source"}, + }} + resourceMigrator := migrationTestClient(t, sourceConfigMap, otherConfigMap) + if err := resourceMigrator.DeleteSourceNamespaceResources(ctx, []unstructured.Unstructured{*sourceConfigMap.DeepCopy()}, Options{SubscriptionNamespace: "source", InstallNamespace: "target"}); err != nil { + t.Fatalf("DeleteSourceNamespaceResources() error = %v", err) + } + gotSource := &unstructured.Unstructured{} + gotSource.SetAPIVersion("v1") + gotSource.SetKind("ConfigMap") + if err := resourceMigrator.Client.Get(ctx, client.ObjectKeyFromObject(sourceConfigMap), gotSource); !apierrors.IsNotFound(err) { + t.Fatalf("source operator ConfigMap error = %v, want not found", err) + } + gotOther := &unstructured.Unstructured{} + gotOther.SetAPIVersion("v1") + gotOther.SetKind("ConfigMap") + if err := resourceMigrator.Client.Get(ctx, client.ObjectKeyFromObject(otherConfigMap), gotOther); err != nil { + t.Fatalf("unrelated ConfigMap was deleted: %v", err) + } + if err := resourceMigrator.DeleteSourceNamespaceResources(ctx, []unstructured.Unstructured{{Object: map[string]interface{}{ + "apiVersion": "v1", "kind": "ConfigMap", "metadata": map[string]interface{}{"name": "no-uid", "namespace": "source"}, + }}}, Options{SubscriptionNamespace: "source", InstallNamespace: "target"}); err == nil { + t.Fatal("DeleteSourceNamespaceResources() unexpectedly deleted an object without a UID") + } +} + +func TestPrepareInstallNamespaceRejectsWeakerPSA(t *testing.T) { + ctx := context.Background() + source := &corev1.Namespace{ObjectMeta: metav1.ObjectMeta{Name: "source", Labels: map[string]string{"pod-security.kubernetes.io/enforce": "privileged"}}} + target := &corev1.Namespace{ObjectMeta: metav1.ObjectMeta{Name: "target", Labels: map[string]string{"pod-security.kubernetes.io/enforce": "restricted"}}} + m := migrationTestClient(t, source, target) + err := m.PrepareInstallNamespace(ctx, Options{SubscriptionNamespace: "source", InstallNamespace: "target"}) + if err == nil { + t.Fatal("PrepareInstallNamespace() unexpectedly weakened target PSA enforcement") + } + var unchanged corev1.Namespace + if err := m.Client.Get(ctx, client.ObjectKey{Name: "target"}, &unchanged); err != nil { + t.Fatalf("get target namespace: %v", err) + } + if got := unchanged.Labels["pod-security.kubernetes.io/enforce"]; got != "restricted" { + t.Fatalf("target PSA enforcement = %q, want restricted", got) + } +} + +func TestPrepareInstallNamespaceRejectsUnknownPSADefault(t *testing.T) { + ctx := context.Background() + source := &corev1.Namespace{ObjectMeta: metav1.ObjectMeta{Name: "source", Labels: map[string]string{"pod-security.kubernetes.io/enforce": "baseline"}}} + target := &corev1.Namespace{ObjectMeta: metav1.ObjectMeta{Name: "target"}} + m := migrationTestClient(t, source, target) + err := m.PrepareInstallNamespace(ctx, Options{SubscriptionNamespace: "source", InstallNamespace: "target"}) + if err == nil { + t.Fatal("PrepareInstallNamespace() unexpectedly changed an existing target with unknown PSA enforcement") + } + var unchanged corev1.Namespace + if err := m.Client.Get(ctx, client.ObjectKey{Name: "target"}, &unchanged); err != nil { + t.Fatalf("get target namespace: %v", err) + } + if _, found := unchanged.Labels["pod-security.kubernetes.io/enforce"]; found { + t.Fatalf("target PSA enforcement was changed: %#v", unchanged.Labels) + } +} + +func TestScaleSourceDeployments(t *testing.T) { + ctx := context.Background() + replicas := int32(2) + deployment := &appsv1.Deployment{ObjectMeta: metav1.ObjectMeta{Name: "operator", Namespace: "source"}, Spec: appsv1.DeploymentSpec{Replicas: &replicas}} + m := migrationTestClient(t, deployment) + objects := []unstructured.Unstructured{{Object: map[string]interface{}{ + "apiVersion": "apps/v1", "kind": "Deployment", "metadata": map[string]interface{}{"name": "operator", "namespace": "source"}, + }}} + restore, err := m.ScaleSourceDeployments(ctx, objects, Options{SubscriptionNamespace: "source", InstallNamespace: "target"}) + if err != nil { + t.Fatalf("ScaleSourceDeployments() error = %v", err) + } + var scaled appsv1.Deployment + if err := m.Client.Get(ctx, client.ObjectKey{Namespace: "source", Name: "operator"}, &scaled); err != nil { + t.Fatalf("get scaled Deployment: %v", err) + } + if scaled.Spec.Replicas == nil || *scaled.Spec.Replicas != 0 { + t.Fatalf("scaled Deployment replicas = %v, want 0", scaled.Spec.Replicas) + } + if err := restore(ctx); err != nil { + t.Fatalf("restore source Deployments: %v", err) + } + if err := m.Client.Get(ctx, client.ObjectKey{Namespace: "source", Name: "operator"}, &scaled); err != nil { + t.Fatalf("get restored Deployment: %v", err) + } + if scaled.Spec.Replicas == nil || *scaled.Spec.Replicas != 2 { + t.Fatalf("restored Deployment replicas = %v, want 2", scaled.Spec.Replicas) + } +} + +func TestCleanupResultErr(t *testing.T) { + err := (&CleanupResult{Actions: []CleanupAction{{Description: "Delete OperatorCondition", Error: errors.New("forbidden")}}}).Err() + if err == nil || !strings.Contains(err.Error(), "Delete OperatorCondition: forbidden") { + t.Fatalf("CleanupResult.Err() = %v, want action error", err) + } +} + func TestPrepareClusterObjectSet(t *testing.T) { establishedCRD := &apiextensionsv1.CustomResourceDefinition{ ObjectMeta: metav1.ObjectMeta{Name: clusterObjectSetCRDName}, @@ -756,12 +988,15 @@ func TestCreateMigrationResourcesCleansTrackedCOSAfterClusterExtensionFailure(t object := unstructured.Unstructured{Object: map[string]interface{}{ "apiVersion": "v1", "kind": "ConfigMap", "metadata": map[string]interface{}{"name": "operator-config", "namespace": "ns"}, }} - err := m.CreateMigrationResources(ctx, Options{SubscriptionName: "sub", SubscriptionNamespace: "ns", ClusterExtensionName: "sub", SystemNamespace: "olmv1-system"}, &MigrationInfo{ + result, err := m.CreateMigrationResources(ctx, Options{SubscriptionName: "sub", SubscriptionNamespace: "ns", ClusterExtensionName: "sub", SystemNamespace: "olmv1-system"}, &MigrationInfo{ PackageName: "widgets", BundleName: "widgets.v1", Version: "1.0.0", CollectedObjects: []unstructured.Unstructured{object}, }, nil) if err == nil || !strings.Contains(err.Error(), "ClusterExtension creation failed") { t.Fatalf("CreateMigrationResources() error = %v, want ClusterExtension creation failure", err) } + if !result.TargetMayBeActive { + t.Fatal("CreateMigrationResources() TargetMayBeActive = false, want true after COS success") + } m.Client = baseClient if err := m.Client.Get(ctx, client.ObjectKey{Name: "sub-1"}, &ocv1.ClusterObjectSet{}); err == nil { t.Fatal("ClusterExtension failure left the ClusterObjectSet created by this invocation") From 77081230521b672d6400d33d1622b72120f72dc2 Mon Sep 17 00:00:00 2001 From: Todd Short Date: Fri, 25 Sep 2026 12:03:27 -0400 Subject: [PATCH 2/3] fix: recover safely from namespace cutover failures Signed-off-by: Todd Short --- .../cmd/migrate-operators-v0-to-v1/convert.go | 7 +- migration/pkg/migration/migration.go | 9 ++- migration/pkg/migration/namespace.go | 6 +- migration/pkg/migration/unit_test.go | 69 ++++++++++++++++++- 4 files changed, 85 insertions(+), 6 deletions(-) diff --git a/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go b/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go index e16a893..8c56669 100644 --- a/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go +++ b/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go @@ -278,7 +278,12 @@ func runConvert(cmd *cobra.Command, args []string) error { //nolint:nestif success("OLMv0 management removed") restoreSourceDeployments, err := m.ScaleSourceDeployments(ctx, sourceObjects, opts) if err != nil { - return err + recoveryCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second) + defer cancel() + if recoverErr := m.RecoverFromBackup(recoveryCtx, opts, backup); recoverErr != nil { + return fmt.Errorf("scale source Deployments: %w; recovery also failed: %v", err, recoverErr) + } + return fmt.Errorf("scale source Deployments failed (recovered): %w", err) } if opts.InstallNamespace != opts.SubscriptionNamespace { success("Source Deployments scaled to zero before target cutover") diff --git a/migration/pkg/migration/migration.go b/migration/pkg/migration/migration.go index cc35990..ca1ec0b 100644 --- a/migration/pkg/migration/migration.go +++ b/migration/pkg/migration/migration.go @@ -166,7 +166,12 @@ func (m *Migrator) Migrate(ctx context.Context, opts Options) error { restoreSourceDeployments, err := m.ScaleSourceDeployments(ctx, sourceObjects, opts) if err != nil { - return err + recoveryCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second) + defer cancel() + if recoverErr := m.RecoverFromBackup(recoveryCtx, opts, backup); recoverErr != nil { + return fmt.Errorf("scale source Deployments: %w; recovery also failed: %v", err, recoverErr) + } + return fmt.Errorf("scale source Deployments failed (recovered): %w", err) } result, err := m.CreateMigrationResources(ctx, opts, info, backup) if err != nil { @@ -539,6 +544,7 @@ func (m *Migrator) createClusterObjectSet(ctx context.Context, opts Options, inf secret.Annotations[migrationInvocationAnnotation] = invocationMarker if err := m.Client.Create(ctx, secret); err != nil { if apierrors.IsAlreadyExists(err) { + resources.ownershipUnknown = true return resources, fmt.Errorf("failed to create COS ref Secret %s: %w", secret.Name, err) } if m.resolveCreatedObject(context.WithoutCancel(ctx), secret, invocationMarker) { @@ -602,6 +608,7 @@ func (m *Migrator) createClusterObjectSet(ctx context.Context, opts Options, inf // and Secret cleanup unsafe. if err := m.Client.Create(ctx, cosObj); err != nil { if apierrors.IsAlreadyExists(err) { + resources.ownershipUnknown = true return resources, failWithSecretCleanup(fmt.Errorf("failed to create ClusterObjectSet: %w", err)) } if m.resolveCreatedObject(context.WithoutCancel(ctx), cosObj, invocationMarker) { diff --git a/migration/pkg/migration/namespace.go b/migration/pkg/migration/namespace.go index d25d90d..dba7aa9 100644 --- a/migration/pkg/migration/namespace.go +++ b/migration/pkg/migration/namespace.go @@ -5,6 +5,7 @@ import ( "errors" "fmt" "strings" + "time" appsv1 "k8s.io/api/apps/v1" corev1 "k8s.io/api/core/v1" @@ -311,7 +312,10 @@ func (m *Migrator) ScaleSourceDeployments(ctx context.Context, objects []unstruc original = &sourceDeploymentReplica{name: deployment.Name, namespace: deployment.Namespace, replicas: replicas} return nil }); err != nil { - if restoreErr := restoreSourceDeploymentReplicas(ctx, m, originals); restoreErr != nil { + restoreCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second) + restoreErr := restoreSourceDeploymentReplicas(restoreCtx, m, originals) + cancel() + if restoreErr != nil { return noRestore, fmt.Errorf("scale source Deployment %s/%s: %w; restore scaled Deployments: %v", obj.GetNamespace(), obj.GetName(), err, restoreErr) } return noRestore, fmt.Errorf("scale source Deployment %s/%s: %w", obj.GetNamespace(), obj.GetName(), err) diff --git a/migration/pkg/migration/unit_test.go b/migration/pkg/migration/unit_test.go index db054db..5adf326 100644 --- a/migration/pkg/migration/unit_test.go +++ b/migration/pkg/migration/unit_test.go @@ -43,6 +43,24 @@ type failingMigrationClient struct { succeedCOS bool } +// canceledScaleClient cancels the caller's context while failing one scale-down +// update. It verifies partial scale recovery uses an independent context. +type canceledScaleClient struct { + client.Client + cancel context.CancelFunc +} + +func (c canceledScaleClient) Update(ctx context.Context, obj client.Object, opts ...client.UpdateOption) error { + if deployment, ok := obj.(*appsv1.Deployment); ok && deployment.Name == "second" && deployment.Spec.Replicas != nil && *deployment.Spec.Replicas == 0 { + c.cancel() + return errors.New("simulated scale failure") + } + if err := ctx.Err(); err != nil { + return err + } + return c.Client.Update(ctx, obj, opts...) +} + func (c failingMigrationClient) Create(ctx context.Context, obj client.Object, opts ...client.CreateOption) error { if c.failCOSCreate { if _, ok := obj.(*ocv1.ClusterObjectSet); ok { @@ -576,6 +594,33 @@ func TestScaleSourceDeployments(t *testing.T) { } } +func TestScaleSourceDeploymentsRestoresAfterCanceledContext(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + firstReplicas, secondReplicas := int32(2), int32(1) + first := &appsv1.Deployment{ObjectMeta: metav1.ObjectMeta{Name: "first", Namespace: "source"}, Spec: appsv1.DeploymentSpec{Replicas: &firstReplicas}} + second := &appsv1.Deployment{ObjectMeta: metav1.ObjectMeta{Name: "second", Namespace: "source"}, Spec: appsv1.DeploymentSpec{Replicas: &secondReplicas}} + m := migrationTestClient(t, first, second) + m.Client = canceledScaleClient{Client: m.Client, cancel: cancel} + objects := []unstructured.Unstructured{ + {Object: map[string]interface{}{"apiVersion": "apps/v1", "kind": "Deployment", "metadata": map[string]interface{}{"name": "first", "namespace": "source"}}}, + {Object: map[string]interface{}{"apiVersion": "apps/v1", "kind": "Deployment", "metadata": map[string]interface{}{"name": "second", "namespace": "source"}}}, + } + + _, err := m.ScaleSourceDeployments(ctx, objects, Options{SubscriptionNamespace: "source", InstallNamespace: "target"}) + if err == nil || !strings.Contains(err.Error(), "simulated scale failure") { + t.Fatalf("ScaleSourceDeployments() error = %v, want scale failure", err) + } + + var restored appsv1.Deployment + if err := m.Client.Get(context.Background(), client.ObjectKey{Namespace: "source", Name: "first"}, &restored); err != nil { + t.Fatalf("get restored Deployment: %v", err) + } + if restored.Spec.Replicas == nil || *restored.Spec.Replicas != firstReplicas { + t.Fatalf("restored Deployment replicas = %v, want %d", restored.Spec.Replicas, firstReplicas) + } +} + func TestCleanupResultErr(t *testing.T) { err := (&CleanupResult{Actions: []CleanupAction{{Description: "Delete OperatorCondition", Error: errors.New("forbidden")}}}).Err() if err == nil || !strings.Contains(err.Error(), "Delete OperatorCondition: forbidden") { @@ -920,7 +965,7 @@ func TestRollbackAndCleanupRejectInvalidInputWithoutMutation(t *testing.T) { } } -func TestCreateClusterObjectSetCleansTemporarySecretsOnCollision(t *testing.T) { +func TestCreateClusterObjectSetPreservesTemporarySecretsOnCollision(t *testing.T) { ctx := context.Background() m := migrationTestClient(t, establishedClusterObjectSetCRD()) m.Client = failingMigrationClient{Client: m.Client, failCOSCreate: true} @@ -937,14 +982,32 @@ func TestCreateClusterObjectSetCleansTemporarySecretsOnCollision(t *testing.T) { if err := m.Client.List(ctx, &secrets, client.InNamespace("olmv1-system")); err != nil { t.Fatal(err) } - if len(secrets.Items) != 0 { - t.Fatalf("COS collision left temporary Secret(s): %#v", secrets.Items) + if len(secrets.Items) == 0 { + t.Fatal("COS collision removed temporary Secret(s) despite unknown ownership") } if err := m.Client.Get(ctx, client.ObjectKey{Name: "sub-1"}, &ocv1.ClusterObjectSet{}); err == nil { t.Fatal("COS collision created a ClusterObjectSet") } } +func TestCreateMigrationResourcesTreatsCOSCollisionAsPotentiallyActive(t *testing.T) { + ctx := context.Background() + m := migrationTestClient(t, establishedClusterObjectSetCRD()) + m.Client = failingMigrationClient{Client: m.Client, failCOSCreate: true} + object := unstructured.Unstructured{Object: map[string]interface{}{ + "apiVersion": "v1", "kind": "ConfigMap", "metadata": map[string]interface{}{"name": "operator-config", "namespace": "ns"}, + }} + result, err := m.CreateMigrationResources(ctx, Options{SubscriptionName: "sub", SubscriptionNamespace: "ns", ClusterExtensionName: "sub", SystemNamespace: "olmv1-system"}, &MigrationInfo{ + PackageName: "widgets", BundleName: "widgets.v1", Version: "1.0.0", CollectedObjects: []unstructured.Unstructured{object}, + }, nil) + if err == nil || !apierrors.IsAlreadyExists(err) { + t.Fatalf("CreateMigrationResources() error = %v, want COS collision", err) + } + if !result.TargetMayBeActive { + t.Fatal("CreateMigrationResources() TargetMayBeActive = false, want true after COS collision") + } +} + func TestCreateClusterObjectSetPreservesSecretsAfterUnknownCreateOutcome(t *testing.T) { ctx := context.Background() m := migrationTestClient(t, establishedClusterObjectSetCRD()) From 58c914f98d1d151cff93f103fd97e7604706686a Mon Sep 17 00:00:00 2001 From: Todd Short Date: Fri, 25 Sep 2026 16:18:45 -0400 Subject: [PATCH 3/3] fix: allow full subscription recovery timeout Signed-off-by: Todd Short --- .../examples/cmd/migrate-operators-v0-to-v1/convert.go | 2 +- migration/pkg/migration/migration.go | 10 ++++++++-- 2 files changed, 9 insertions(+), 3 deletions(-) diff --git a/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go b/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go index 8c56669..7689a10 100644 --- a/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go +++ b/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go @@ -278,7 +278,7 @@ func runConvert(cmd *cobra.Command, args []string) error { //nolint:nestif success("OLMv0 management removed") restoreSourceDeployments, err := m.ScaleSourceDeployments(ctx, sourceObjects, opts) if err != nil { - recoveryCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second) + recoveryCtx, cancel := migration.NewRecoveryContext(ctx) defer cancel() if recoverErr := m.RecoverFromBackup(recoveryCtx, opts, backup); recoverErr != nil { return fmt.Errorf("scale source Deployments: %w; recovery also failed: %v", err, recoverErr) diff --git a/migration/pkg/migration/migration.go b/migration/pkg/migration/migration.go index ca1ec0b..0f084f3 100644 --- a/migration/pkg/migration/migration.go +++ b/migration/pkg/migration/migration.go @@ -166,7 +166,7 @@ func (m *Migrator) Migrate(ctx context.Context, opts Options) error { restoreSourceDeployments, err := m.ScaleSourceDeployments(ctx, sourceObjects, opts) if err != nil { - recoveryCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second) + recoveryCtx, cancel := NewRecoveryContext(ctx) defer cancel() if recoverErr := m.RecoverFromBackup(recoveryCtx, opts, backup); recoverErr != nil { return fmt.Errorf("scale source Deployments: %w; recovery also failed: %v", err, recoverErr) @@ -323,6 +323,12 @@ func (m *Migrator) PrepareForMigration(ctx context.Context, opts Options, csv *o return nil } +// NewRecoveryContext returns a cancellation-independent context with enough +// time to recreate and observe a restored Subscription. +func NewRecoveryContext(ctx context.Context) (context.Context, context.CancelFunc) { + return context.WithTimeout(context.WithoutCancel(ctx), subWaitTimeout+30*time.Second) +} + // RecoverFromBackup restores the Subscription from backup after a failed preparation. func (m *Migrator) RecoverFromBackup(ctx context.Context, opts Options, backup *Backup) error { if backup == nil { @@ -373,7 +379,7 @@ func (m *Migrator) RecoverBeforeCE(ctx context.Context, opts Options, backup *Ba // restore the OLMv0 Subscription automatically because orphaned target objects // can still be running while deletion propagates. func (m *Migrator) recoverCreatedMigrationResources(ctx context.Context, opts Options, backup *Backup, resources *createdMigrationResources) error { - recoveryCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), subWaitTimeout+30*time.Second) + recoveryCtx, cancel := NewRecoveryContext(ctx) defer cancel() if resources == nil { return m.RecoverBeforeCE(recoveryCtx, opts, backup)