diff --git a/migration/examples/cmd/migrate-operators-v0-to-v1/validation_test.go b/migration/examples/cmd/migrate-operators-v0-to-v1/validation_test.go new file mode 100644 index 0000000..414739f --- /dev/null +++ b/migration/examples/cmd/migrate-operators-v0-to-v1/validation_test.go @@ -0,0 +1,60 @@ +package main + +import ( + "context" + "testing" + + "github.com/spf13/cobra" +) + +// These validations run before constructing a Kubernetes client. Keeping them +// tested here makes invalid CLI invocations deterministic and side-effect free. +func TestCommandRejectsAmbiguousAndMissingTargets(t *testing.T) { + t.Cleanup(func() { checkAll, convertAll, cleanupAll, rollbackAll = false, false, false, false }) + cmd := &cobra.Command{} + cmd.SetContext(context.Background()) + tests := []struct { + name string + run func([]string) error + }{ + { + name: "check", + run: func(args []string) error { + return runCheck(cmd, args) + }, + }, + { + name: "convert", + run: func(args []string) error { + return runConvert(cmd, args) + }, + }, + { + name: "cleanup", + run: func(args []string) error { + return runCleanup(cmd, args) + }, + }, + { + name: "rollback", + run: func(args []string) error { + return runRollback(cmd, args) + }, + }, + } + + for _, tt := range tests { + t.Run(tt.name+" missing target", func(t *testing.T) { + checkAll, convertAll, cleanupAll, rollbackAll = false, false, false, false + if err := tt.run(nil); err == nil { + t.Fatal("missing target unexpectedly reached the Kubernetes client") + } + }) + t.Run(tt.name+" ambiguous target", func(t *testing.T) { + checkAll, convertAll, cleanupAll, rollbackAll = true, true, true, true + if err := tt.run([]string{"operator"}); err == nil { + t.Fatal("target combined with --all unexpectedly reached the Kubernetes client") + } + }) + } +} diff --git a/migration/pkg/catalogmigration/catalogmigration.go b/migration/pkg/catalogmigration/catalogmigration.go index d751cce..d9a3b31 100644 --- a/migration/pkg/catalogmigration/catalogmigration.go +++ b/migration/pkg/catalogmigration/catalogmigration.go @@ -20,6 +20,7 @@ import ( const ( // MigratedFromCatalogSourceAnnotation is set on ClusterCatalog when first created or adopted. MigratedFromCatalogSourceAnnotation = "olm.operatorframework.io/migrated-from-catalogsource" + clusterCatalogServingTimeout = 10 * time.Minute ) // CatalogMigratorOptions configures the catalog migration. @@ -346,7 +347,7 @@ func (cm *CatalogMigrator) annotateIfNotPresent(ctx context.Context, cc *ocv1.Cl // waitForServing polls until the ClusterCatalog has Serving=True. func (cm *CatalogMigrator) waitForServing(ctx context.Context, ccName string) error { - return wait.PollUntilContextTimeout(ctx, 5*time.Second, 3*time.Minute, true, func(ctx context.Context) (bool, error) { + return wait.PollUntilContextTimeout(ctx, 5*time.Second, clusterCatalogServingTimeout, true, func(ctx context.Context) (bool, error) { var cc ocv1.ClusterCatalog if err := cm.Client.Get(ctx, client.ObjectKey{Name: ccName}, &cc); err != nil { return false, err diff --git a/migration/pkg/catalogmigration/unit_test.go b/migration/pkg/catalogmigration/unit_test.go index f959b34..709783c 100644 --- a/migration/pkg/catalogmigration/unit_test.go +++ b/migration/pkg/catalogmigration/unit_test.go @@ -144,3 +144,35 @@ func TestWaitForServing(t *testing.T) { t.Fatal("waitForServing(missing) unexpectedly succeeded") } } + +func TestMigrateCatalogsSkipsUnsupportedSourcesWithoutMutation(t *testing.T) { + scheme := runtime.NewScheme() + if err := operatorsv1alpha1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } + if err := ocv1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } + sources := []runtime.Object{ + &operatorsv1alpha1.CatalogSource{ObjectMeta: metav1.ObjectMeta{Name: "configmap", Namespace: "ns"}, Spec: operatorsv1alpha1.CatalogSourceSpec{SourceType: operatorsv1alpha1.SourceTypeConfigmap}}, + &operatorsv1alpha1.CatalogSource{ObjectMeta: metav1.ObjectMeta{Name: "address-only", Namespace: "ns"}, Spec: operatorsv1alpha1.CatalogSourceSpec{SourceType: operatorsv1alpha1.SourceTypeGrpc, Address: "catalog.ns.svc:50051"}}, + &operatorsv1alpha1.CatalogSource{ObjectMeta: metav1.ObjectMeta{Name: "unknown", Namespace: "ns"}, Spec: operatorsv1alpha1.CatalogSourceSpec{SourceType: operatorsv1alpha1.SourceType("unsupported")}}, + } + c := fake.NewClientBuilder().WithScheme(scheme).WithRuntimeObjects(sources...).Build() + results, err := NewCatalogMigrator(c).MigrateCatalogs(context.Background(), CatalogMigratorOptions{}) + if err != nil || len(results) != len(sources) { + t.Fatalf("MigrateCatalogs() = %#v, %v", results, err) + } + for _, result := range results { + if result.Status != "skipped" || result.ClusterCatalogName != "" || result.Reason == "" { + t.Fatalf("unsupported source result = %#v", result) + } + } + var catalogs ocv1.ClusterCatalogList + if err := c.List(context.Background(), &catalogs); err != nil { + t.Fatal(err) + } + if len(catalogs.Items) != 0 { + t.Fatalf("unsupported CatalogSources created ClusterCatalogs: %#v", catalogs.Items) + } +} diff --git a/migration/pkg/migration/labels.go b/migration/pkg/migration/labels.go index b0e9668..aed08c6 100644 --- a/migration/pkg/migration/labels.go +++ b/migration/pkg/migration/labels.go @@ -43,9 +43,6 @@ const ( // MigratedFromCatalogSourceAnnotation is set on ClusterCatalog by the catalog migration tool. MigratedFromCatalogSourceAnnotation = "olm.operatorframework.io/migrated-from-catalogsource" - // fieldManager is the SSA field manager used for all apply operations. - fieldManager = "olm.operatorframework.io/migration" - // cosWaitPollInterval / cosWaitTimeout control how long to wait for a // ClusterObjectSet to reach Succeeded=True. cosWaitPollInterval = 5 * time.Second diff --git a/migration/pkg/migration/migration.go b/migration/pkg/migration/migration.go index 6e33a81..2aabec2 100644 --- a/migration/pkg/migration/migration.go +++ b/migration/pkg/migration/migration.go @@ -3,13 +3,18 @@ package migration import ( "context" "encoding/json" + "errors" "fmt" "strings" + "time" + corev1 "k8s.io/api/core/v1" apiextensionsv1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/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/apimachinery/pkg/types" + "k8s.io/apimachinery/pkg/util/uuid" "k8s.io/apimachinery/pkg/util/wait" "sigs.k8s.io/controller-runtime/pkg/client" @@ -26,6 +31,18 @@ var annotationPrefixesToStrip = []string{ "deployment.kubernetes.io/", } +// createdMigrationResources records only objects created by this invocation. +// Recovery must never infer ownership from predictable names or labels: another +// migration may have created objects with the same values after this one started. +type createdMigrationResources struct { + cos *ocv1.ClusterObjectSet + secrets []corev1.Secret + ce *ocv1.ClusterExtension + ownershipUnknown bool +} + +const migrationInvocationAnnotation = "olm.operatorframework.io/migration-invocation" + // Migrate performs the full migration of an OLMv0-managed operator to OLMv1. // Steps: // 1. Profile the Operator (Subscription/CSV/InstallPlan) @@ -44,6 +61,9 @@ func (m *Migrator) Migrate(ctx context.Context, opts Options) error { if err != nil { return fmt.Errorf("failed to profile operator: %w", err) } + if err := m.ensureClusterExtensionAbsent(ctx, opts.ClusterExtensionName); err != nil { + return err + } info, err := m.GetBundleInfo(ctx, opts, csv, ip) if err != nil { @@ -77,7 +97,6 @@ func (m *Migrator) Migrate(ctx context.Context, opts Options) error { if err != nil { return fmt.Errorf("failed to backup resources: %w", err) } - // Populate CE backup annotations (R2.5) — must happen before PrepareForMigration deletes the Sub. if backup.Subscription != nil { if j, err := json.Marshal(backup.Subscription.Spec); err == nil { @@ -117,15 +136,24 @@ 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.CreateClusterObjectSet(ctx, opts, info); err != nil { - if recoverErr := m.RecoverBeforeCE(ctx, opts, backup); recoverErr != nil { + resources, err := m.createClusterObjectSet(ctx, opts, info) + if err != nil { + if recoverErr := m.recoverCreatedMigrationResources(ctx, opts, backup, resources); recoverErr != nil { return fmt.Errorf("COS creation failed: %w; recovery also failed: %v", err, recoverErr) } return fmt.Errorf("COS creation failed (recovered): %w", err) } - if err := m.CreateClusterExtension(ctx, opts, info); err != nil { - return fmt.Errorf("failed to create ClusterExtension: %w", err) + ce, ownershipUnknown, err := m.createClusterExtension(ctx, opts, info) + if ce != nil { + resources.ce = ce + } + resources.ownershipUnknown = resources.ownershipUnknown || ownershipUnknown + if err != nil { + if recoverErr := m.recoverCreatedMigrationResources(ctx, opts, backup, resources); recoverErr != nil { + return fmt.Errorf("ClusterExtension creation failed: %w; recovery also failed: %v", err, recoverErr) + } + return fmt.Errorf("ClusterExtension creation failed (recovered): %w", err) } m.CleanupOLMv0Resources(ctx, opts, info.PackageName, csv.Name) @@ -242,19 +270,80 @@ func (m *Migrator) RecoverFromBackup(ctx context.Context, opts Options, backup * }) } -// RecoverBeforeCE implements recovery when COS creation fails. -// Deletes the failed COS with orphan cascade, then restores the Subscription. +// RecoverBeforeCE restores the Subscription after a failure before a migration +// resource was created. It deliberately does not delete a predictably named COS +// or Secret because this invocation cannot establish ownership of such objects. func (m *Migrator) RecoverBeforeCE(ctx context.Context, opts Options, backup *Backup) error { - cosName := fmt.Sprintf("%s-1", opts.ClusterExtensionName) - cos := &ocv1.ClusterObjectSet{} - cos.Name = cosName - if err := m.Client.Delete(ctx, cos, client.PropagationPolicy(metav1.DeletePropagationOrphan)); err != nil { - if client.IgnoreNotFound(err) != nil { - return fmt.Errorf("failed to delete COS during recovery: %w", err) + return m.RecoverFromBackup(ctx, opts, backup) +} + +// 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. +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() + if resources == nil { + return m.RecoverBeforeCE(recoveryCtx, opts, backup) + } + if resources.ownershipUnknown { + return fmt.Errorf("migration resource creation outcome is unknown; refusing automatic recovery") + } + if resources.ce != nil { + if err := m.deleteTrackedResource(recoveryCtx, resources.ce); err != nil && client.IgnoreNotFound(err) != nil { + return fmt.Errorf("delete created ClusterExtension during recovery: %w", err) + } + } + if resources.cos != nil { + if err := m.deleteTrackedResource(recoveryCtx, resources.cos, client.PropagationPolicy(metav1.DeletePropagationOrphan)); err != nil && client.IgnoreNotFound(err) != nil { + return fmt.Errorf("delete created ClusterObjectSet during recovery: %w", err) } } + cleanupErr := m.cleanupCreatedSecrets(recoveryCtx, resources.secrets) + recoverErr := m.RecoverFromBackup(recoveryCtx, opts, backup) + return errors.Join(cleanupErr, recoverErr) +} - return m.RecoverFromBackup(ctx, opts, backup) +func (m *Migrator) cleanupCreatedSecrets(ctx context.Context, secrets []corev1.Secret) error { + var errs []error + for i := range secrets { + if err := m.deleteTrackedResource(ctx, &secrets[i]); err != nil && client.IgnoreNotFound(err) != nil { + errs = append(errs, fmt.Errorf("delete COS ref Secret %s: %w", secrets[i].Name, err)) + } + } + return errors.Join(errs...) +} + +func (m *Migrator) deleteTrackedResource(ctx context.Context, obj client.Object, opts ...client.DeleteOption) error { + uid := obj.GetUID() + opts = append(opts, client.Preconditions{UID: &uid}) + return m.Client.Delete(ctx, obj, opts...) +} + +func createdByInvocation(obj client.Object, marker string) bool { + return obj.GetAnnotations()[migrationInvocationAnnotation] == marker +} + +func (m *Migrator) resolveCreatedObject(ctx context.Context, obj client.Object, marker string) bool { + if err := m.Client.Get(ctx, client.ObjectKeyFromObject(obj), obj); err != nil { + return false + } + return createdByInvocation(obj, marker) +} + +// ensureClusterExtensionAbsent rejects a target name before OLMv0 resources +// are removed. CreateClusterExtension remains the race-safe final check. +func (m *Migrator) ensureClusterExtensionAbsent(ctx context.Context, name string) error { + var ce ocv1.ClusterExtension + err := m.Client.Get(ctx, client.ObjectKey{Name: name}, &ce) + if err == nil { + return fmt.Errorf("ClusterExtension %s already exists", name) + } + if client.IgnoreNotFound(err) != nil { + return fmt.Errorf("check ClusterExtension %s: %w", name, err) + } + return nil } // CreateClusterObjectSet builds and creates a COS from the collected resources. @@ -266,6 +355,18 @@ func (m *Migrator) RecoverBeforeCE(ctx context.Context, opts Options, backup *Ba // OLMv1 revision object(s) are appropriate. Track upstream progress at OPRUN-4716 and the // boxcutter ClusterObjectDeployment design. func (m *Migrator) CreateClusterObjectSet(ctx context.Context, opts Options, info *MigrationInfo) error { + resources, err := m.createClusterObjectSet(ctx, opts, info) + if err != nil { + if cleanupErr := m.cleanupCreatedClusterObjectSet(ctx, resources); cleanupErr != nil { + return errors.Join(err, cleanupErr) + } + } + return err +} + +func (m *Migrator) createClusterObjectSet(ctx context.Context, opts Options, info *MigrationInfo) (*createdMigrationResources, error) { + resources := &createdMigrationResources{} + invocationMarker := string(uuid.NewUUID()) cosName := fmt.Sprintf("%s-1", opts.ClusterExtensionName) systemNS := opts.systemNamespace() @@ -291,15 +392,42 @@ func (m *Migrator) CreateClusterObjectSet(ctx context.Context, opts Options, inf } packed, err := packer.pack(phases) if err != nil { - return fmt.Errorf("failed to pack COS objects into Secrets: %w", err) + return resources, fmt.Errorf("failed to pack COS objects into Secrets: %w", err) + } + cleanupSecrets := func() error { + cleanupCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second) + defer cancel() + return m.cleanupCreatedSecrets(cleanupCtx, resources.secrets) + } + failWithSecretCleanup := func(err error) error { + if resources.ownershipUnknown { + return err + } + if cleanupErr := cleanupSecrets(); cleanupErr != nil { + return errors.Join(err, fmt.Errorf("clean up COS ref Secrets: %w", cleanupErr)) + } + return err } // Create ref Secrets before the COS so the COS controller can find them immediately. for i := range packed.Secrets { secret := &packed.Secrets[i] + if secret.Annotations == nil { + secret.Annotations = map[string]string{} + } + secret.Annotations[migrationInvocationAnnotation] = invocationMarker if err := m.Client.Create(ctx, secret); err != nil { - return fmt.Errorf("failed to create COS ref Secret %s: %w", secret.Name, err) + if apierrors.IsAlreadyExists(err) { + return resources, fmt.Errorf("failed to create COS ref Secret %s: %w", secret.Name, err) + } + if m.resolveCreatedObject(context.WithoutCancel(ctx), secret, invocationMarker) { + resources.secrets = append(resources.secrets, *secret) + } else { + resources.ownershipUnknown = true + } + return resources, failWithSecretCleanup(fmt.Errorf("failed to create COS ref Secret %s: %w", secret.Name, err)) } + resources.secrets = append(resources.secrets, *secret) } // Replace inline objects with Secret refs in the phases. @@ -325,6 +453,7 @@ func (m *Migrator) CreateClusterObjectSet(ctx context.Context, opts Options, inf LabelPackageName: info.PackageName, LabelBundleName: info.BundleName, LabelBundleVersion: info.Version, + migrationInvocationAnnotation: invocationMarker, } if info.BundleImage != "" { cosAnnotations[LabelBundleReference] = info.BundleImage @@ -338,20 +467,52 @@ func (m *Migrator) CreateClusterObjectSet(ctx context.Context, opts Options, inf }). WithAnnotations(cosAnnotations) - cosObj := &ocv1.ClusterObjectSet{} - cosObj.Name = cosName - cosData, err := json.Marshal(cos) if err != nil { - return fmt.Errorf("failed to marshal COS: %w", err) + return resources, failWithSecretCleanup(fmt.Errorf("failed to marshal COS: %w", err)) + } + cosObj := &ocv1.ClusterObjectSet{} + if err := json.Unmarshal(cosData, cosObj); err != nil { + return resources, failWithSecretCleanup(fmt.Errorf("failed to decode COS: %w", err)) } - if err := m.Client.Patch(ctx, cosObj, client.RawPatch(types.ApplyPatchType, cosData), - client.ForceOwnership, client.FieldOwner(fieldManager)); err != nil { - return fmt.Errorf("failed to apply ClusterObjectSet: %w", err) + // A migration owns only a newly created revision. Applying an existing COS + // can overwrite its ownership or fail on immutable phases, making recovery + // and Secret cleanup unsafe. + if err := m.Client.Create(ctx, cosObj); err != nil { + if apierrors.IsAlreadyExists(err) { + return resources, failWithSecretCleanup(fmt.Errorf("failed to create ClusterObjectSet: %w", err)) + } + if m.resolveCreatedObject(context.WithoutCancel(ctx), cosObj, invocationMarker) { + resources.cos = cosObj + return resources, fmt.Errorf("failed to create ClusterObjectSet: %w", err) + } + resources.ownershipUnknown = true + return resources, fmt.Errorf("failed to create ClusterObjectSet: %w", err) + } + resources.cos = cosObj + + if err := m.WaitForCOSSucceeded(ctx, cosName); err != nil { + return resources, err } + return resources, nil +} - return m.WaitForCOSSucceeded(ctx, cosName) +func (m *Migrator) cleanupCreatedClusterObjectSet(ctx context.Context, resources *createdMigrationResources) error { + if resources == nil { + return nil + } + if resources.ownershipUnknown { + return fmt.Errorf("ClusterObjectSet creation outcome is unknown; refusing automatic cleanup") + } + cleanupCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second) + defer cancel() + if resources.cos != nil { + if err := m.deleteTrackedResource(cleanupCtx, resources.cos, client.PropagationPolicy(metav1.DeletePropagationOrphan)); err != nil && client.IgnoreNotFound(err) != nil { + return fmt.Errorf("delete created ClusterObjectSet: %w", err) + } + } + return m.cleanupCreatedSecrets(cleanupCtx, resources.secrets) } // WaitForCOSSucceeded waits for the COS to reach Succeeded=True. @@ -381,9 +542,21 @@ func (m *Migrator) WaitForCOSSucceeded(ctx context.Context, cosName string) erro // spec.serviceAccount is NOT set — deprecated and ignored in OLMv1 (R2.5/R7). // Migration annotations (R2.5) are added for AlreadyMigrated/Conflict detection and rollback. func (m *Migrator) CreateClusterExtension(ctx context.Context, opts Options, info *MigrationInfo) error { + ce, _, err := m.createClusterExtension(ctx, opts, info) + if err != nil && ce != nil { + if deleteErr := m.deleteTrackedResource(context.WithoutCancel(ctx), ce); deleteErr != nil && client.IgnoreNotFound(deleteErr) != nil { + return errors.Join(err, fmt.Errorf("delete created ClusterExtension: %w", deleteErr)) + } + } + return err +} + +func (m *Migrator) createClusterExtension(ctx context.Context, opts Options, info *MigrationInfo) (*ocv1.ClusterExtension, bool, error) { + invocationMarker := string(uuid.NewUUID()) // Build annotations (R2.5). annotations := map[string]string{ MigratedFromSubscriptionAnnotation: fmt.Sprintf("%s/%s", opts.SubscriptionNamespace, opts.SubscriptionName), + migrationInvocationAnnotation: invocationMarker, } if info.SubscriptionBackupJSON != "" { annotations[MigrationSubscriptionBackupAnnotation] = info.SubscriptionBackupJSON @@ -465,11 +638,11 @@ func (m *Migrator) CreateClusterExtension(ctx context.Context, opts Options, inf if info.SubscriptionConfig != nil { cfgJSON, err := json.Marshal(info.SubscriptionConfig) if err != nil { - return fmt.Errorf("failed to marshal SubscriptionConfig for CE: %w", err) + return nil, false, fmt.Errorf("failed to marshal SubscriptionConfig for CE: %w", err) } inlineJSON, err := json.Marshal(map[string]json.RawMessage{"deploymentConfig": cfgJSON}) if err != nil { - return fmt.Errorf("failed to marshal CE inline config: %w", err) + return nil, false, fmt.Errorf("failed to marshal CE inline config: %w", err) } ce.Spec.Config = &ocv1.ClusterExtensionConfig{ ConfigType: ocv1.ClusterExtensionConfigTypeInline, @@ -478,10 +651,16 @@ func (m *Migrator) CreateClusterExtension(ctx context.Context, opts Options, inf } if err := m.Client.Create(ctx, ce); err != nil { - return fmt.Errorf("failed to create ClusterExtension: %w", err) + if apierrors.IsAlreadyExists(err) { + return nil, false, fmt.Errorf("failed to create ClusterExtension: %w", err) + } + if m.resolveCreatedObject(context.WithoutCancel(ctx), ce, invocationMarker) { + return ce, false, fmt.Errorf("failed to create ClusterExtension: %w", err) + } + return nil, true, fmt.Errorf("failed to create ClusterExtension: %w", err) } - return m.WaitForClusterExtensionInstalled(ctx, opts.ClusterExtensionName) + return ce, false, m.WaitForClusterExtensionInstalled(ctx, opts.ClusterExtensionName) } // WaitForClusterExtensionInstalled waits for the CE to reach Installed=True. diff --git a/migration/pkg/migration/scan.go b/migration/pkg/migration/scan.go index 9384306..f627310 100644 --- a/migration/pkg/migration/scan.go +++ b/migration/pkg/migration/scan.go @@ -6,6 +6,7 @@ import ( "fmt" "strings" + "k8s.io/apimachinery/pkg/util/validation" "sigs.k8s.io/controller-runtime/pkg/client" operatorsv1alpha1 "github.com/operator-framework/api/pkg/operators/v1alpha1" @@ -343,6 +344,29 @@ func (m *Migrator) RollbackClusterExtension(ctx context.Context, ceName string, } } + // Validate the data needed to restore OLMv0 ownership before deleting any + // OLMv1 resources. A corrupt or incomplete backup must leave the existing + // ClusterExtension and ClusterObjectSet recoverable. + subBackupJSON, ok := ce.Annotations[MigrationSubscriptionBackupAnnotation] + if !ok || subBackupJSON == "" { + return fmt.Errorf("ClusterExtension %s has no migration-subscription-backup annotation; cannot restore Subscription", ceName) + } + subRef := ce.Annotations[MigratedFromSubscriptionAnnotation] + if subRef == "" { + return fmt.Errorf("ClusterExtension %s has no migrated-from-subscription annotation", ceName) + } + var subSpec operatorsv1alpha1.SubscriptionSpec + if err := unmarshalJSON(subBackupJSON, &subSpec); err != nil { + return fmt.Errorf("failed to unmarshal subscription backup: %w", err) + } + if subSpec.Package == "" || subSpec.CatalogSource == "" || subSpec.CatalogSourceNamespace == "" { + return fmt.Errorf("subscription backup is missing required package, source, or sourceNamespace") + } + ns, name, err := splitSubRef(subRef) + if err != nil { + return fmt.Errorf("invalid migrated-from-subscription annotation %q: %w", subRef, err) + } + // Delete CE (orphan cascade — preserves operator workloads) if err := m.Client.Delete(ctx, &ce, client.PropagationPolicy("Orphan")); err != nil { if client.IgnoreNotFound(err) != nil { @@ -361,28 +385,7 @@ func (m *Migrator) RollbackClusterExtension(ctx context.Context, ceName string, } } - // Restore Subscription from backup annotation - subBackupJSON, ok := ce.Annotations["olm.operatorframework.io/migration-subscription-backup"] - if !ok || subBackupJSON == "" { - return fmt.Errorf("ClusterExtension %s has no migration-subscription-backup annotation; cannot restore Subscription", ceName) - } - - subRef := ce.Annotations[MigratedFromSubscriptionAnnotation] - if subRef == "" { - return fmt.Errorf("ClusterExtension %s has no migrated-from-subscription annotation", ceName) - } - - // Restore Subscription - var subSpec operatorsv1alpha1.SubscriptionSpec - if err := unmarshalJSON(subBackupJSON, &subSpec); err != nil { - return fmt.Errorf("failed to unmarshal subscription backup: %w", err) - } - - ns, name, err := splitSubRef(subRef) - if err != nil { - return fmt.Errorf("invalid migrated-from-subscription annotation %q: %w", subRef, err) - } - + // Restore Subscription from the validated backup. restoredSub := &operatorsv1alpha1.Subscription{} restoredSub.Name = name restoredSub.Namespace = ns @@ -443,8 +446,11 @@ func (m *Migrator) CleanupConflict(ctx context.Context, ceName string) error { // splitSubRef splits a "namespace/name" subscription reference into its components. func splitSubRef(ref string) (string, string, error) { + if strings.Count(ref, "/") != 1 { + return "", "", fmt.Errorf("invalid namespace/name ref %q", ref) + } parts := strings.SplitN(ref, "/", 2) - if len(parts) != 2 || parts[0] == "" || parts[1] == "" { + if len(validation.IsDNS1123Label(parts[0])) != 0 || len(validation.IsDNS1123Subdomain(parts[1])) != 0 { return "", "", fmt.Errorf("invalid namespace/name ref %q", ref) } return parts[0], parts[1], nil diff --git a/migration/pkg/migration/unit_test.go b/migration/pkg/migration/unit_test.go index f1a90e3..27558a8 100644 --- a/migration/pkg/migration/unit_test.go +++ b/migration/pkg/migration/unit_test.go @@ -5,6 +5,7 @@ import ( "compress/gzip" "context" "encoding/json" + "errors" "fmt" "net/http" "net/url" @@ -15,6 +16,7 @@ import ( corev1 "k8s.io/api/core/v1" rbacv1 "k8s.io/api/rbac/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/apimachinery/pkg/runtime" @@ -30,6 +32,43 @@ import ( ocv1ac "github.com/operator-framework/operator-controller/applyconfigurations/api/v1" ) +type failingMigrationClient struct { + client.Client + failCOSCreate bool + unknownCOSCreate bool + failCECreate bool + blockCOS bool +} + +func (c failingMigrationClient) Create(ctx context.Context, obj client.Object, opts ...client.CreateOption) error { + if c.failCOSCreate { + if _, ok := obj.(*ocv1.ClusterObjectSet); ok { + return apierrors.NewAlreadyExists(ocv1.GroupVersion.WithResource("clusterobjectsets").GroupResource(), obj.GetName()) + } + } + if c.unknownCOSCreate { + if _, ok := obj.(*ocv1.ClusterObjectSet); ok { + return errors.New("simulated ClusterObjectSet create transport failure") + } + } + if c.failCECreate { + if _, ok := obj.(*ocv1.ClusterExtension); ok { + return apierrors.NewAlreadyExists(ocv1.GroupVersion.WithResource("clusterextensions").GroupResource(), obj.GetName()) + } + } + return c.Client.Create(ctx, obj, opts...) +} + +func (c failingMigrationClient) Get(ctx context.Context, key client.ObjectKey, obj client.Object, opts ...client.GetOption) error { + if c.blockCOS { + if cos, ok := obj.(*ocv1.ClusterObjectSet); ok { + cos.Status.Conditions = []metav1.Condition{{Type: ocv1.ClusterObjectSetTypeSucceeded, Status: metav1.ConditionFalse, Reason: ocv1.ClusterObjectSetReasonBlocked}} + return nil + } + } + return c.Client.Get(ctx, key, obj, opts...) +} + func migrationTestClient(t *testing.T, objects ...runtime.Object) *Migrator { t.Helper() scheme := runtime.NewScheme() @@ -42,6 +81,9 @@ func migrationTestClient(t *testing.T, objects ...runtime.Object) *Migrator { if err := rbacv1.AddToScheme(scheme); err != nil { t.Fatal(err) } + if err := corev1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } if err := ocv1.AddToScheme(scheme); err != nil { t.Fatal(err) } @@ -383,6 +425,35 @@ func TestScanStatesAndPublicHelpers(t *testing.T) { } } +func TestScanAllKeepsMixedUnsafeOperatorsOutOfEligibleResults(t *testing.T) { + ctx := context.Background() + busySub, busyCSV := healthySubscriptionFixtures() + busySub.Name, busyCSV.Name = "busy", "busy.v1" + busySub.Status.InstalledCSV = busyCSV.Name + busySub.Status.State = operatorsv1alpha1.SubscriptionStateUpgradeAvailable + conflictSub := &operatorsv1alpha1.Subscription{ObjectMeta: metav1.ObjectMeta{Name: "conflict", Namespace: "ns"}, Spec: &operatorsv1alpha1.SubscriptionSpec{Package: "conflict"}} + conflictCE := &ocv1.ClusterExtension{ObjectMeta: metav1.ObjectMeta{Name: "conflict", Annotations: map[string]string{MigratedFromSubscriptionAnnotation: "ns/conflict"}}} + already := &ocv1.ClusterExtension{ObjectMeta: metav1.ObjectMeta{Name: "already", Annotations: map[string]string{MigratedFromSubscriptionAnnotation: "ns/gone"}}} + m := migrationTestClient(t, busySub, busyCSV, conflictSub, conflictCE, already) + + results, err := m.ScanAll(ctx) + if err != nil || len(results) != 3 { + t.Fatalf("ScanAll() = %#v, %v", results, err) + } + if eligible := EligibleFromScan(results); len(eligible) != 0 { + t.Fatalf("unsafe mixed scan returned eligible operators: %#v", eligible) + } + got := map[OperatorStatus]int{} + for _, result := range results { + got[result.Status]++ + } + for _, status := range []OperatorStatus{OperatorStatusIneligible, OperatorStatusConflict, OperatorStatusAlreadyMigrated} { + if got[status] != 1 { + t.Fatalf("mixed scan status counts = %#v, want one %s", got, status) + } + } +} + func TestPrerequisitesAndRecoveryErrors(t *testing.T) { ctx := context.Background() sub, csv := healthySubscriptionFixtures() @@ -410,3 +481,313 @@ func TestPrerequisitesAndRecoveryErrors(t *testing.T) { t.Fatal("Cleanup() unexpectedly succeeded for a missing ClusterExtension") } } + +// TestMigrateRejectsUnsafeOperatorsWithoutMutation verifies the safety boundary +// of conversion: failures before preparation must leave OLMv0 resources intact. +func TestMigrateRejectsUnsafeOperatorsWithoutMutation(t *testing.T) { + ctx := context.Background() + tests := []struct { + name string + mutate func(*operatorsv1alpha1.Subscription, *operatorsv1alpha1.ClusterServiceVersion) + }{ + { + name: "not steady", + mutate: func(sub *operatorsv1alpha1.Subscription, _ *operatorsv1alpha1.ClusterServiceVersion) { + sub.Status.State = operatorsv1alpha1.SubscriptionStateUpgradeAvailable + }, + }, + { + name: "incompatible API service", + mutate: func(_ *operatorsv1alpha1.Subscription, csv *operatorsv1alpha1.ClusterServiceVersion) { + csv.Spec.APIServiceDefinitions.Owned = []operatorsv1alpha1.APIServiceDescription{{Name: "v1.widgets"}} + }, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + sub, csv := healthySubscriptionFixtures() + tt.mutate(sub, csv) + m := migrationTestClient(t, sub, csv) + + if err := m.Migrate(ctx, Options{SubscriptionName: sub.Name, SubscriptionNamespace: sub.Namespace}); err == nil { + t.Fatal("Migrate() unexpectedly accepted an unsafe operator") + } + if err := m.Client.Get(ctx, client.ObjectKeyFromObject(sub), &operatorsv1alpha1.Subscription{}); err != nil { + t.Fatalf("unsafe migration deleted Subscription: %v", err) + } + if err := m.Client.Get(ctx, client.ObjectKeyFromObject(csv), &operatorsv1alpha1.ClusterServiceVersion{}); err != nil { + t.Fatalf("unsafe migration deleted CSV: %v", err) + } + if err := m.Client.Get(ctx, client.ObjectKey{Name: sub.Name}, &ocv1.ClusterExtension{}); err == nil { + t.Fatal("unsafe migration created a ClusterExtension") + } + }) + } +} + +func TestRollbackAndCleanupRejectInvalidInputWithoutMutation(t *testing.T) { + ctx := context.Background() + installed := &ocv1.ClusterExtension{ + ObjectMeta: metav1.ObjectMeta{ + Name: "installed", + Annotations: map[string]string{ + MigratedFromSubscriptionAnnotation: "ns/sub", + MigrationSubscriptionBackupAnnotation: `{}`, + }, + }, + Status: ocv1.ClusterExtensionStatus{Conditions: []metav1.Condition{{Type: "Installed", Status: metav1.ConditionTrue}}}, + } + cos := &ocv1.ClusterObjectSet{ObjectMeta: metav1.ObjectMeta{Name: "installed-1"}} + m := migrationTestClient(t, installed, cos) + if err := m.Rollback(ctx, Options{ClusterExtensionName: "installed"}); err == nil { + t.Fatal("rollback of Installed=True ClusterExtension unexpectedly succeeded without acknowledgement") + } + if err := m.Client.Get(ctx, client.ObjectKeyFromObject(installed), &ocv1.ClusterExtension{}); err != nil { + t.Fatalf("unacknowledged rollback deleted ClusterExtension: %v", err) + } + if err := m.Client.Get(ctx, client.ObjectKeyFromObject(cos), &ocv1.ClusterObjectSet{}); err != nil { + t.Fatalf("unacknowledged rollback deleted ClusterObjectSet: %v", err) + } + + invalid := &ocv1.ClusterExtension{ObjectMeta: metav1.ObjectMeta{Name: "invalid", Annotations: map[string]string{MigratedFromSubscriptionAnnotation: "not-a-reference"}}} + sub := &operatorsv1alpha1.Subscription{ObjectMeta: metav1.ObjectMeta{Name: "sub", Namespace: "ns"}} + m = migrationTestClient(t, invalid, sub) + if err := m.CleanupConflict(ctx, "invalid"); err == nil { + t.Fatal("cleanup accepted an invalid migrated-from-subscription annotation") + } + if err := m.Client.Get(ctx, client.ObjectKeyFromObject(sub), &operatorsv1alpha1.Subscription{}); err != nil { + t.Fatalf("invalid conflict cleanup deleted Subscription: %v", err) + } + if err := m.Client.Get(ctx, client.ObjectKeyFromObject(invalid), &ocv1.ClusterExtension{}); err != nil { + t.Fatalf("invalid conflict cleanup deleted ClusterExtension: %v", err) + } +} + +func TestCreateClusterObjectSetCleansTemporarySecretsOnCollision(t *testing.T) { + ctx := context.Background() + m := migrationTestClient(t) + 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"}, + }} + err := m.CreateClusterObjectSet(ctx, Options{SubscriptionName: "sub", SubscriptionNamespace: "ns"}, &MigrationInfo{ + PackageName: "widgets", BundleName: "widgets.v1", Version: "1.0.0", CollectedObjects: []unstructured.Unstructured{object}, + }) + if err == nil || !apierrors.IsAlreadyExists(err) { + t.Fatalf("CreateClusterObjectSet() error = %v, want collision", err) + } + var secrets corev1.SecretList + 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 err := m.Client.Get(ctx, client.ObjectKey{Name: "sub-1"}, &ocv1.ClusterObjectSet{}); err == nil { + t.Fatal("COS collision created a ClusterObjectSet") + } +} + +func TestCreateClusterObjectSetPreservesSecretsAfterUnknownCreateOutcome(t *testing.T) { + ctx := context.Background() + m := migrationTestClient(t) + m.Client = failingMigrationClient{Client: m.Client, unknownCOSCreate: true} + object := unstructured.Unstructured{Object: map[string]interface{}{ + "apiVersion": "v1", "kind": "ConfigMap", "metadata": map[string]interface{}{"name": "operator-config", "namespace": "ns"}, + }} + err := m.CreateClusterObjectSet(ctx, Options{SubscriptionName: "sub", SubscriptionNamespace: "ns"}, &MigrationInfo{ + PackageName: "widgets", BundleName: "widgets.v1", Version: "1.0.0", CollectedObjects: []unstructured.Unstructured{object}, + }) + if err == nil || !strings.Contains(err.Error(), "refusing automatic cleanup") { + t.Fatalf("CreateClusterObjectSet() error = %v, want unknown ownership cleanup refusal", err) + } + var secrets corev1.SecretList + if err := m.Client.List(ctx, &secrets, client.InNamespace("olmv1-system")); err != nil { + t.Fatal(err) + } + if len(secrets.Items) == 0 { + t.Fatal("unknown COS creation outcome removed Secret references") + } +} + +func TestCreateClusterExtensionAlreadyExistsIsKnownOutcome(t *testing.T) { + ctx := context.Background() + m := migrationTestClient(t) + m.Client = failingMigrationClient{Client: m.Client, failCECreate: true} + ce, ownershipUnknown, err := m.createClusterExtension(ctx, Options{SubscriptionName: "sub", SubscriptionNamespace: "ns", ClusterExtensionName: "sub"}, &MigrationInfo{PackageName: "widgets"}) + if err == nil || !apierrors.IsAlreadyExists(err) { + t.Fatalf("createClusterExtension() error = %v, want already exists", err) + } + if ce != nil || ownershipUnknown { + t.Fatalf("createClusterExtension() = ce=%#v, ownershipUnknown=%t, want known non-created outcome", ce, ownershipUnknown) + } +} + +func TestCreateClusterObjectSetCleansSecretsAfterReadinessFailure(t *testing.T) { + ctx := context.Background() + m := migrationTestClient(t) + baseClient := m.Client + m.Client = failingMigrationClient{Client: baseClient, blockCOS: true} + object := unstructured.Unstructured{Object: map[string]interface{}{ + "apiVersion": "v1", "kind": "ConfigMap", "metadata": map[string]interface{}{"name": "operator-config", "namespace": "ns"}, + }} + err := m.CreateClusterObjectSet(ctx, Options{SubscriptionName: "sub", SubscriptionNamespace: "ns"}, &MigrationInfo{ + PackageName: "widgets", BundleName: "widgets.v1", Version: "1.0.0", CollectedObjects: []unstructured.Unstructured{object}, + }) + if err == nil { + t.Fatal("CreateClusterObjectSet() unexpectedly completed without COS status") + } + m.Client = baseClient + var secrets corev1.SecretList + if err := m.Client.List(ctx, &secrets, client.InNamespace("olmv1-system")); err != nil { + t.Fatal(err) + } + if len(secrets.Items) != 0 { + t.Fatalf("COS readiness failure left Secret references: %#v", secrets.Items) + } + if err := m.Client.Get(ctx, client.ObjectKey{Name: "sub-1"}, &ocv1.ClusterObjectSet{}); err == nil { + t.Fatal("COS readiness failure left ClusterObjectSet") + } +} + +func TestRecoverBeforeCEPreservesResourcesItDidNotCreate(t *testing.T) { + ctx := context.Background() + secret := &corev1.Secret{ObjectMeta: metav1.ObjectMeta{ + Name: "sub-1-ref", Namespace: "olmv1-system", + Labels: map[string]string{LabelRevisionName: "sub-1", LabelOwnerName: "sub"}, + }} + cos := &ocv1.ClusterObjectSet{ObjectMeta: metav1.ObjectMeta{Name: "sub-1"}} + m := migrationTestClient(t, secret, cos) + if err := m.RecoverBeforeCE(ctx, Options{ClusterExtensionName: "sub"}, nil); err == nil { + t.Fatal("RecoverBeforeCE() unexpectedly succeeded without a backup") + } + if err := m.Client.Get(ctx, client.ObjectKeyFromObject(secret), &corev1.Secret{}); err != nil { + t.Fatalf("RecoverBeforeCE() deleted a pre-existing COS reference Secret: %v", err) + } + if err := m.Client.Get(ctx, client.ObjectKeyFromObject(cos), &ocv1.ClusterObjectSet{}); err != nil { + t.Fatalf("RecoverBeforeCE() deleted a pre-existing ClusterObjectSet: %v", err) + } +} + +func TestRecoverCreatedMigrationResourcesDeletesOnlyTrackedResources(t *testing.T) { + ctx := context.Background() + createdSecret := corev1.Secret{ObjectMeta: metav1.ObjectMeta{Name: "created-ref", Namespace: "olmv1-system"}} + otherSecret := &corev1.Secret{ObjectMeta: metav1.ObjectMeta{Name: "other-ref", Namespace: "olmv1-system"}} + createdCOS := &ocv1.ClusterObjectSet{ObjectMeta: metav1.ObjectMeta{Name: "created-1"}} + otherCOS := &ocv1.ClusterObjectSet{ObjectMeta: metav1.ObjectMeta{Name: "other-1"}} + createdCE := &ocv1.ClusterExtension{ObjectMeta: metav1.ObjectMeta{Name: "created"}} + m := migrationTestClient(t, &createdSecret, otherSecret, createdCOS, otherCOS, createdCE) + + err := m.recoverCreatedMigrationResources(ctx, Options{}, nil, &createdMigrationResources{ + cos: createdCOS, + secrets: []corev1.Secret{createdSecret}, + ce: createdCE, + }) + if err == nil { + t.Fatal("recovery unexpectedly succeeded without a backup") + } + for _, object := range []client.Object{&createdSecret, createdCOS, createdCE} { + if err := m.Client.Get(ctx, client.ObjectKeyFromObject(object), object); err == nil { + t.Fatalf("recovery left created %T %q", object, object.GetName()) + } + } + for _, object := range []client.Object{otherSecret, otherCOS} { + if err := m.Client.Get(ctx, client.ObjectKeyFromObject(object), object); err != nil { + t.Fatalf("recovery deleted untracked %T %q: %v", object, object.GetName(), err) + } + } +} + +func TestRecoverCreatedMigrationResourcesRefusesUnknownOwnership(t *testing.T) { + ctx := context.Background() + cos := &ocv1.ClusterObjectSet{ObjectMeta: metav1.ObjectMeta{Name: "possibly-foreign-1"}} + m := migrationTestClient(t, cos) + err := m.recoverCreatedMigrationResources(ctx, Options{}, nil, &createdMigrationResources{ + cos: cos, + ownershipUnknown: true, + }) + if err == nil || !strings.Contains(err.Error(), "outcome is unknown") { + t.Fatalf("recovery error = %v, want unknown ownership", err) + } + if err := m.Client.Get(ctx, client.ObjectKeyFromObject(cos), &ocv1.ClusterObjectSet{}); err != nil { + t.Fatalf("recovery deleted a resource with unknown ownership: %v", err) + } +} + +func TestSplitSubRefRejectsMalformedReferences(t *testing.T) { + for _, ref := range []string{"", "ns", "ns/", "/sub", "ns/sub/extra", "invalid_namespace/sub", "ns/invalid_name"} { + if _, _, err := splitSubRef(ref); err == nil { + t.Fatalf("splitSubRef(%q) unexpectedly succeeded", ref) + } + } + if namespace, name, err := splitSubRef("valid-ns/valid.subscription"); err != nil || namespace != "valid-ns" || name != "valid.subscription" { + t.Fatalf("splitSubRef(valid) = %q, %q, %v", namespace, name, err) + } +} + +func TestCreateClusterExtensionRejectsNameCollisionWithoutReplacement(t *testing.T) { + ctx := context.Background() + existing := &ocv1.ClusterExtension{ObjectMeta: metav1.ObjectMeta{Name: "sub", Annotations: map[string]string{"keep": "existing"}}} + m := migrationTestClient(t, existing) + err := m.CreateClusterExtension(ctx, Options{SubscriptionName: "sub", SubscriptionNamespace: "ns"}, &MigrationInfo{PackageName: "widgets"}) + if err == nil { + t.Fatal("CreateClusterExtension() unexpectedly replaced an existing ClusterExtension") + } + var got ocv1.ClusterExtension + if err := m.Client.Get(ctx, client.ObjectKey{Name: "sub"}, &got); err != nil || got.Annotations["keep"] != "existing" { + t.Fatalf("existing ClusterExtension changed after collision: %#v, err=%v", got, err) + } +} + +func TestEnsureClusterExtensionAbsent(t *testing.T) { + ctx := context.Background() + m := migrationTestClient(t, &ocv1.ClusterExtension{ObjectMeta: metav1.ObjectMeta{Name: "sub"}}) + if err := m.ensureClusterExtensionAbsent(ctx, "sub"); err == nil { + t.Fatal("existing ClusterExtension passed pre-migration check") + } + if err := m.ensureClusterExtensionAbsent(ctx, "available"); err != nil { + t.Fatalf("absent ClusterExtension failed pre-migration check: %v", err) + } +} + +func TestMigrateRejectsExistingClusterExtensionWithoutMutation(t *testing.T) { + ctx := context.Background() + sub, csv := healthySubscriptionFixtures() + ce := &ocv1.ClusterExtension{ObjectMeta: metav1.ObjectMeta{Name: "sub"}} + m := migrationTestClient(t, sub, csv, ce) + if err := m.Migrate(ctx, Options{SubscriptionName: sub.Name, SubscriptionNamespace: sub.Namespace}); err == nil { + t.Fatal("Migrate() unexpectedly accepted an existing ClusterExtension") + } + if err := m.Client.Get(ctx, client.ObjectKeyFromObject(sub), &operatorsv1alpha1.Subscription{}); err != nil { + t.Fatalf("existing ClusterExtension migration deleted Subscription: %v", err) + } + if err := m.Client.Get(ctx, client.ObjectKeyFromObject(csv), &operatorsv1alpha1.ClusterServiceVersion{}); err != nil { + t.Fatalf("existing ClusterExtension migration deleted CSV: %v", err) + } +} + +func TestRollbackRejectsMissingAndMalformedBackupsWithoutMutation(t *testing.T) { + ctx := context.Background() + for name, annotations := range map[string]map[string]string{ + "missing backup": {MigratedFromSubscriptionAnnotation: "ns/sub"}, + "malformed backup": {MigratedFromSubscriptionAnnotation: "ns/sub", MigrationSubscriptionBackupAnnotation: "{"}, + "empty backup": {MigratedFromSubscriptionAnnotation: "ns/sub", MigrationSubscriptionBackupAnnotation: `{}`}, + "missing source ref": {MigrationSubscriptionBackupAnnotation: `{}`}, + } { + t.Run(name, func(t *testing.T) { + ce := &ocv1.ClusterExtension{ObjectMeta: metav1.ObjectMeta{Name: "sub", Annotations: annotations}} + cos := &ocv1.ClusterObjectSet{ObjectMeta: metav1.ObjectMeta{Name: "sub-1"}} + m := migrationTestClient(t, ce, cos) + if err := m.Rollback(ctx, Options{ClusterExtensionName: "sub", AcknowledgeInstalled: true}); err == nil { + t.Fatal("Rollback() unexpectedly accepted invalid backup metadata") + } + if err := m.Client.Get(ctx, client.ObjectKeyFromObject(ce), &ocv1.ClusterExtension{}); err != nil { + t.Fatalf("invalid rollback deleted ClusterExtension: %v", err) + } + if err := m.Client.Get(ctx, client.ObjectKeyFromObject(cos), &ocv1.ClusterObjectSet{}); err != nil { + t.Fatalf("invalid rollback deleted ClusterObjectSet: %v", err) + } + }) + } +} diff --git a/test/e2e/migration/e2e_test.go b/test/e2e/migration/e2e_test.go index b807774..3050e35 100644 --- a/test/e2e/migration/e2e_test.go +++ b/test/e2e/migration/e2e_test.go @@ -5,6 +5,7 @@ package e2e import ( "bytes" "encoding/json" + "fmt" "os" "os/exec" "path/filepath" @@ -128,6 +129,60 @@ func TestCatalogSourceMigration(t *testing.T) { run(t, "kubectl", "get", "catalogsource/"+name, "-n", namespace) } +// TestFixtureNegativeGuards verifies that fixture migration refuses unsafe +// inputs without creating OLMv1 resources. It runs before TestMigration, which +// recreates the migrated catalog and performs the successful conversion. +func TestFixtureNegativeGuards(t *testing.T) { + if os.Getenv("E2E_SUITE") != "fixture" { + t.Skip("negative fixture guards run only against the controller-free OLMv0 fixture suite") + } + namespace, subscription := os.Getenv("E2E_NAMESPACE"), os.Getenv("E2E_SUBSCRIPTION") + if namespace == "" || subscription == "" { + t.Fatal("E2E_NAMESPACE and E2E_SUBSCRIPTION are required") + } + + // A non-steady Subscription is unsafe to migrate. Both check and convert + // must reject it before they remove OLMv0 resources or create OLMv1 ones. + t.Cleanup(func() { + if out, err := output("kubectl", "patch", "subscription/"+subscription, "-n", namespace, + "--subresource=status", "--type=merge", "--patch", `{"status":{"state":"AtLatestKnown"}}`); err != nil { + t.Errorf("restore Subscription state: %v\n%s", err, out) + } + }) + run(t, "kubectl", "patch", "subscription/"+subscription, "-n", namespace, + "--subresource=status", "--type=merge", "--patch", `{"status":{"state":"UpgradeAvailable"}}`) + expectCheckFailure(t, "Subscription state", binary(t, "migrate-operators-v0-to-v1"), "check", subscription, "-n", namespace, "--kubeconfig", os.Getenv("KUBECONFIG")) + expectFailure(t, binary(t, "migrate-operators-v0-to-v1"), "convert", subscription, "-n", namespace, "--kubeconfig", os.Getenv("KUBECONFIG")) + assertNoMigrationObjects(t, subscription) + run(t, "kubectl", "patch", "subscription/"+subscription, "-n", namespace, + "--subresource=status", "--type=merge", "--patch", `{"status":{"state":"AtLatestKnown"}}`) + run(t, "kubectl", "get", "subscription/"+subscription, "-n", namespace) + + // Catalog availability is a hard prerequisite. Ask the serving catalog for a + // package that does not exist, then verify the failed resolution is still + // non-mutating. Unlike deleting ClusterCatalogs, this does not race the + // bootstrap catalog reconciler. + packageName, err := output("kubectl", "get", "subscription/"+subscription, "-n", namespace, "-o", "jsonpath={.spec.name}") + if err != nil || strings.TrimSpace(packageName) == "" { + t.Fatalf("get source package: %v (%s)", err, packageName) + } + t.Cleanup(func() { + patch := fmt.Sprintf(`{"spec":{"name":%q}}`, strings.TrimSpace(packageName)) + if out, err := output("kubectl", "patch", "subscription/"+subscription, "-n", namespace, + "--type=merge", "--patch", patch); err != nil { + t.Errorf("restore Subscription package: %v\n%s", err, out) + } + }) + run(t, "kubectl", "patch", "subscription/"+subscription, "-n", namespace, + "--type=merge", "--patch", `{"spec":{"name":"migration-fixture-package-that-does-not-exist"}}`) + expectCheckFailure(t, "No ClusterCatalog found", binary(t, "migrate-operators-v0-to-v1"), "check", subscription, "-n", namespace, "--kubeconfig", os.Getenv("KUBECONFIG")) + expectFailure(t, binary(t, "migrate-operators-v0-to-v1"), "convert", subscription, "-n", namespace, "--kubeconfig", os.Getenv("KUBECONFIG")) + assertNoMigrationObjects(t, subscription) + run(t, "kubectl", "patch", "subscription/"+subscription, "-n", namespace, + "--type=merge", "--patch", fmt.Sprintf(`{"spec":{"name":%q}}`, strings.TrimSpace(packageName))) + run(t, "kubectl", "get", "subscription/"+subscription, "-n", namespace) +} + // TestMigration applies the suite's complete fixture, exercises the two migration // binaries, and observes the API server rather than mocking either OLM controller. // E2E_MANIFEST must create the namespace, a CatalogSource, and the named Subscription. @@ -242,6 +297,41 @@ func run(t *testing.T, command string, args ...string) { } } +// expectFailure requires a command to reject its input and includes its output +// in the test failure to make an accidental success diagnosable. +func expectFailure(t *testing.T, command string, args ...string) { + t.Helper() + if out, err := output(command, args...); err == nil { + t.Fatalf("%s %s unexpectedly succeeded:\n%s", command, strings.Join(args, " "), out) + } +} + +// expectCheckFailure verifies the CLI's deliberate check contract: it returns +// success after reporting failed prerequisites, allowing callers to inspect the +// full report. Conversion itself must still reject those prerequisites. +func expectCheckFailure(t *testing.T, want, command string, args ...string) { + t.Helper() + out, err := output(command, args...) + if err != nil { + t.Fatalf("%s %s returned an error instead of a check report: %v\n%s", command, strings.Join(args, " "), err, out) + } + if !strings.Contains(out, want) { + t.Fatalf("%s %s did not report %q:\n%s", command, strings.Join(args, " "), want, out) + } +} + +// assertNoMigrationObjects proves a rejected conversion did not create either +// OLMv1 resource that would take ownership of the OLMv0 installation. +func assertNoMigrationObjects(t *testing.T, subscription string) { + t.Helper() + if out, err := output("kubectl", "get", "clusterextension/"+subscription); err == nil { + t.Fatalf("rejected conversion created ClusterExtension %s:\n%s", subscription, out) + } + if out, err := output("kubectl", "get", "clusterobjectsets", "-l", "olm.operatorframework.io/owner-name="+subscription, "-o", "name"); err != nil || strings.TrimSpace(out) != "" { + t.Fatalf("rejected conversion created ClusterObjectSet(s): err=%v\n%s", err, out) + } +} + // output runs a command and returns its combined standard output and error. func output(command string, args ...string) (string, error) { cmd := exec.Command(command, args...)