From a7f124fa5fbc738023d6373705ae42f70d9cd982 Mon Sep 17 00:00:00 2001 From: Todd Short Date: Fri, 18 Sep 2026 15:59:36 -0400 Subject: [PATCH 1/6] fix: discover catalogd serving certificate Signed-off-by: Todd Short --- migration/pkg/migration/catalog.go | 71 +++++++++++++++++++++++----- migration/pkg/migration/unit_test.go | 31 ++++++++++-- 2 files changed, 86 insertions(+), 16 deletions(-) diff --git a/migration/pkg/migration/catalog.go b/migration/pkg/migration/catalog.go index 8e449e7..f7ed6c2 100644 --- a/migration/pkg/migration/catalog.go +++ b/migration/pkg/migration/catalog.go @@ -12,6 +12,7 @@ import ( "strings" "time" + corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/util/wait" "k8s.io/client-go/kubernetes" @@ -115,22 +116,26 @@ func catalogEndpoint(ctx context.Context, catalog *ocv1.ClusterCatalog, config * if err != nil { return "", nil, nil, fmt.Errorf("create Kubernetes client for catalog port-forward: %w", err) } - catalogConfig, err := catalogdTLSConfig(ctx, clientset, config) + namespace, serverName, err := catalogdServiceLocation(catalog) if err != nil { return "", nil, nil, err } - if inCluster { - return catalog.Status.URLs.Base + "/api/v1/all", func() {}, catalogConfig, nil + podName, err := catalogdLeader(ctx, clientset, namespace) + if err != nil { + return "", nil, nil, err } - podName, err := catalogdLeader(ctx, clientset) + catalogConfig, err := catalogdTLSConfig(ctx, clientset, config, namespace, podName) if err != nil { return "", nil, nil, err } + if inCluster { + return catalog.Status.URLs.Base + "/api/v1/all", func() {}, catalogConfig, nil + } u, err := url.Parse(config.Host) if err != nil { return "", nil, nil, err } - u.Path = path.Join(u.Path, "api", "v1", "namespaces", "olmv1-system", "pods", podName, "portforward") + u.Path = path.Join(u.Path, "api", "v1", "namespaces", namespace, "pods", podName, "portforward") rt, upgrader, err := spdy.RoundTripperFor(config) if err != nil { return "", nil, nil, fmt.Errorf("create catalogd port-forward: %w", err) @@ -161,13 +166,46 @@ func catalogEndpoint(ctx context.Context, catalog *ocv1.ClusterCatalog, config * close(stop) return "", nil, nil, err } - catalogConfig.ServerName = "localhost" + // The local port-forward address is not the catalogd certificate's identity. + // Verify the service DNS name advertised by ClusterCatalog status instead. + catalogConfig.ServerName = serverName return fmt.Sprintf("https://127.0.0.1:%d/catalogs/%s/api/v1/all", ports[0].Local, catalog.Name), func() { close(stop) }, catalogConfig, nil } -// catalogdTLSConfig replaces the Kubernetes API CA with catalogd's serving CA. -func catalogdTLSConfig(ctx context.Context, clientset kubernetes.Interface, config *rest.Config) (*rest.Config, error) { - secret, err := clientset.CoreV1().Secrets("cert-manager").Get(ctx, "olmv1-ca", metav1.GetOptions{}) +// catalogdServiceLocation derives catalogd's service namespace and TLS server +// name from the in-cluster endpoint published by operator-controller. This +// avoids imposing either the upstream cert-manager layout or OpenShift's +// service-ca layout on migration users. +func catalogdServiceLocation(catalog *ocv1.ClusterCatalog) (string, string, error) { + if catalog.Status.URLs == nil || catalog.Status.URLs.Base == "" { + return "", "", fmt.Errorf("catalog %s has no base URL in status", catalog.Name) + } + u, err := url.Parse(catalog.Status.URLs.Base) + if err != nil { + return "", "", fmt.Errorf("parse catalog %s base URL: %w", catalog.Name, err) + } + host := u.Hostname() + parts := strings.Split(host, ".") + if len(parts) < 3 || parts[0] == "" || parts[1] == "" || parts[2] != "svc" { + return "", "", fmt.Errorf("catalog %s base URL %q does not use a Kubernetes service hostname", catalog.Name, catalog.Status.URLs.Base) + } + return parts[1], host, nil +} + +// catalogdTLSConfig replaces the Kubernetes API CA with the CA carried by the +// serving certificate mounted in the current catalogd leader Pod. The secret +// name is deliberately discovered from the Pod: upstream installs use a +// cert-manager secret while OpenShift uses a service-ca-generated secret. +func catalogdTLSConfig(ctx context.Context, clientset kubernetes.Interface, config *rest.Config, namespace, podName string) (*rest.Config, error) { + pod, err := clientset.CoreV1().Pods(namespace).Get(ctx, podName, metav1.GetOptions{}) + if err != nil { + return nil, fmt.Errorf("get catalogd pod: %w", err) + } + secretName := catalogdServingCertificateSecret(pod) + if secretName == "" { + return nil, fmt.Errorf("catalogd pod %s/%s has no serving certificate Secret", namespace, podName) + } + secret, err := clientset.CoreV1().Secrets(namespace).Get(ctx, secretName, metav1.GetOptions{}) if err != nil { return nil, fmt.Errorf("get catalogd CA: %w", err) } @@ -190,17 +228,26 @@ func catalogdTLSConfig(ctx context.Context, clientset kubernetes.Interface, conf return catalogConfig, nil } +func catalogdServingCertificateSecret(pod *corev1.Pod) string { + for _, volume := range pod.Spec.Volumes { + if volume.Name == "catalogserver-certs" && volume.Secret != nil { + return volume.Secret.SecretName + } + } + return "" +} + // catalogdLeader waits for catalogd's leader Lease to reference a current pod. -func catalogdLeader(ctx context.Context, clientset kubernetes.Interface) (string, error) { +func catalogdLeader(ctx context.Context, clientset kubernetes.Interface, namespace string) (string, error) { var lastErr error var leader string err := wait.PollUntilContextTimeout(ctx, time.Second, 30*time.Second, true, func(context.Context) (bool, error) { - pods, err := clientset.CoreV1().Pods("olmv1-system").List(ctx, metav1.ListOptions{LabelSelector: "app.kubernetes.io/name=catalogd"}) + pods, err := clientset.CoreV1().Pods(namespace).List(ctx, metav1.ListOptions{LabelSelector: "app.kubernetes.io/name=catalogd"}) if err != nil { lastErr = fmt.Errorf("list catalogd pods: %w", err) return false, nil } - lease, err := clientset.CoordinationV1().Leases("olmv1-system").Get(ctx, "catalogd-operator-lock", metav1.GetOptions{}) + lease, err := clientset.CoordinationV1().Leases(namespace).Get(ctx, "catalogd-operator-lock", metav1.GetOptions{}) if err != nil { lastErr = fmt.Errorf("get catalogd leader lease: %w", err) return false, nil diff --git a/migration/pkg/migration/unit_test.go b/migration/pkg/migration/unit_test.go index 27558a8..b1e1d1c 100644 --- a/migration/pkg/migration/unit_test.go +++ b/migration/pkg/migration/unit_test.go @@ -161,7 +161,7 @@ not-json } } -func TestCatalogdTLSConfigUsesCatalogCAWithoutAPIServerTransportSettings(t *testing.T) { +func TestCatalogdTLSConfigUsesServingCertificateWithoutAPIServerTransportSettings(t *testing.T) { ca := []byte("catalogd-ca") config := &rest.Config{ TLSClientConfig: rest.TLSClientConfig{CAData: []byte("api-server-ca"), Insecure: true}, @@ -169,12 +169,18 @@ func TestCatalogdTLSConfigUsesCatalogCAWithoutAPIServerTransportSettings(t *test return &url.URL{Scheme: "https", Host: "proxy.example"}, nil }, } - clientset := k8sfake.NewSimpleClientset(&corev1.Secret{ - ObjectMeta: metav1.ObjectMeta{Name: "olmv1-ca", Namespace: "cert-manager"}, + clientset := k8sfake.NewSimpleClientset(&corev1.Pod{ + ObjectMeta: metav1.ObjectMeta{Name: "catalogd", Namespace: "catalogd-ns"}, + Spec: corev1.PodSpec{Volumes: []corev1.Volume{{ + Name: "catalogserver-certs", + VolumeSource: corev1.VolumeSource{Secret: &corev1.SecretVolumeSource{SecretName: "catalogserver-cert"}}, + }}}, + }, &corev1.Secret{ + ObjectMeta: metav1.ObjectMeta{Name: "catalogserver-cert", Namespace: "catalogd-ns"}, Data: map[string][]byte{"ca.crt": ca}, }) - got, err := catalogdTLSConfig(context.Background(), clientset, config) + got, err := catalogdTLSConfig(context.Background(), clientset, config, "catalogd-ns", "catalogd") if err != nil { t.Fatalf("catalogdTLSConfig() error = %v", err) } @@ -183,6 +189,23 @@ func TestCatalogdTLSConfigUsesCatalogCAWithoutAPIServerTransportSettings(t *test } } +func TestCatalogdServiceLocation(t *testing.T) { + catalog := &ocv1.ClusterCatalog{ObjectMeta: metav1.ObjectMeta{Name: "catalog"}} + catalog.Status.URLs = &ocv1.ClusterCatalogURLs{Base: "https://catalogd-service.openshift-catalogd.svc/catalogs/catalog"} + namespace, serverName, err := catalogdServiceLocation(catalog) + if err != nil { + t.Fatalf("catalogdServiceLocation() error = %v", err) + } + if namespace != "openshift-catalogd" || serverName != "catalogd-service.openshift-catalogd.svc" { + t.Fatalf("catalogdServiceLocation() = %q, %q, want openshift-catalogd and service DNS name", namespace, serverName) + } + + catalog.Status.URLs.Base = "https://catalog.example.test/catalogs/catalog" + if _, _, err := catalogdServiceLocation(catalog); err == nil { + t.Fatal("catalogdServiceLocation() accepted a non-service hostname") + } +} + func TestCompatibilityPureChecks(t *testing.T) { for _, properties := range []string{ `[{"type":"olm.package.required","value":{"packageName":"dep"}}]`, From cf98a162e0030c2a7c7ccc1c2ba9fca9ad608ea6 Mon Sep 17 00:00:00 2001 From: Todd Short Date: Fri, 18 Sep 2026 16:05:00 -0400 Subject: [PATCH 2/6] fix: show cleanup in migration dry run Signed-off-by: Todd Short --- .../cmd/migrate-operators-v0-to-v1/convert.go | 32 +++++++++++++++ .../convert_test.go | 41 +++++++++++++++++++ 2 files changed, 73 insertions(+) create mode 100644 migration/examples/cmd/migrate-operators-v0-to-v1/convert_test.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 6988ff1..b6300c3 100644 --- a/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go +++ b/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go @@ -303,6 +303,8 @@ func runConvertDryRun(cmd *cobra.Command, m *migration.Migrator, opts migration. } success(fmt.Sprintf("Package: %s Version: %s Channel: %s", info.PackageName, info.Version, valueOrDefault(info.Channel, "(default)"))) + fmt.Printf("\n Resources that would be created:\n") + detail("ClusterObjectSet:", fmt.Sprintf("%s-1 (wait for Succeeded=True before creating the ClusterExtension)", opts.ClusterExtensionName)) fmt.Printf("\n Resources that would be placed into ClusterObjectSet %s-1:\n", opts.ClusterExtensionName) kindCounts := make(map[string]int) @@ -325,11 +327,41 @@ func runConvertDryRun(cmd *cobra.Command, m *migration.Migrator, opts migration. detail("Channel:", valueOrDefault(info.Channel, "(none set)")) detail("CollisionProtection:", "IfNoController") + fmt.Printf("\n OLMv0 resources that would be deleted or changed:\n") + for _, line := range dryRunCleanupPlan(opts, info) { + info2(line) + } + + fmt.Printf("\n Backup plan:\n") + info2("Store Subscription and OperatorGroup specifications in ClusterExtension annotations before deletion.") + if opts.BackupDirectory != "" { + info2(fmt.Sprintf("Write Subscription, OperatorGroup, CSV, and InstallPlan YAML to %s before deletion (not written during dry run).", opts.BackupDirectory)) + } + fmt.Println() info2("No cluster resources were modified (dry run).") return nil } +// dryRunCleanupPlan describes all OLMv0 cleanup actions performed by a normal +// conversion. It intentionally calls no API: dry-run must remain non-mutating. +func dryRunCleanupPlan(opts migration.Options, info *migration.MigrationInfo) []string { + lines := []string{ + fmt.Sprintf("Delete Subscription %s/%s with orphan propagation (operator workloads remain).", opts.SubscriptionNamespace, opts.SubscriptionName), + fmt.Sprintf("Delete ClusterServiceVersion %s/%s with orphan propagation (operator workloads remain).", opts.SubscriptionNamespace, info.BundleName), + fmt.Sprintf("Delete Operator CR %s.%s.", info.PackageName, opts.SubscriptionNamespace), + fmt.Sprintf("Delete OperatorCondition %s/%s if present.", opts.SubscriptionNamespace, info.BundleName), + fmt.Sprintf("Delete copied ClusterServiceVersions derived from %s if present, with orphan propagation.", info.BundleName), + "Retain InstallPlan resources; conversion does not delete them.", + } + if opts.DeleteOperatorGroup { + lines = append(lines, "Delete OperatorGroup(s) only when no Subscriptions remain; strip OLM ownership labels from their aggregation ClusterRoles first.") + } else { + lines = append(lines, "Retain OperatorGroup(s); --delete-operatorgroup was not specified.") + } + return lines +} + func info2(msg string) { fmt.Printf(" %s\n", msg) } 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 new file mode 100644 index 0000000..747f86a --- /dev/null +++ b/migration/examples/cmd/migrate-operators-v0-to-v1/convert_test.go @@ -0,0 +1,41 @@ +package main + +import ( + "strings" + "testing" + + "github.com/operator-framework/library-olm/migration/pkg/migration" +) + +func TestDryRunCleanupPlan(t *testing.T) { + opts := migration.Options{ + SubscriptionName: "widget-operator", + SubscriptionNamespace: "operators", + DeleteOperatorGroup: true, + } + info := &migration.MigrationInfo{ + PackageName: "widgets", + BundleName: "widgets.v1.2.3", + } + + plan := strings.Join(dryRunCleanupPlan(opts, info), "\n") + for _, expected := range []string{ + "Delete Subscription operators/widget-operator with orphan propagation", + "Delete ClusterServiceVersion operators/widgets.v1.2.3 with orphan propagation", + "Delete Operator CR widgets.operators", + "Delete OperatorCondition operators/widgets.v1.2.3 if present", + "Delete copied ClusterServiceVersions derived from widgets.v1.2.3 if present", + "Retain InstallPlan resources", + "Delete OperatorGroup(s) only when no Subscriptions remain", + } { + if !strings.Contains(plan, expected) { + t.Errorf("dry-run cleanup plan does not include %q:\n%s", expected, plan) + } + } + + opts.DeleteOperatorGroup = false + plan = strings.Join(dryRunCleanupPlan(opts, info), "\n") + if !strings.Contains(plan, "Retain OperatorGroup(s); --delete-operatorgroup was not specified.") { + t.Fatalf("dry-run cleanup plan does not describe the default OperatorGroup behavior:\n%s", plan) + } +} From 2df88a700e86245370e805bd77ca44cac257a48e Mon Sep 17 00:00:00 2001 From: Todd Short Date: Fri, 18 Sep 2026 16:27:20 -0400 Subject: [PATCH 3/6] fix: preflight migration target resources Signed-off-by: Todd Short --- .../cmd/migrate-operators-v0-to-v1/convert.go | 13 ++- migration/pkg/migration/migration.go | 74 +++++++++++++- migration/pkg/migration/readiness.go | 13 +++ migration/pkg/migration/secretpacker.go | 3 +- migration/pkg/migration/types.go | 28 +++--- migration/pkg/migration/unit_test.go | 97 ++++++++++++++++++- 6 files changed, 206 insertions(+), 22 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 b6300c3..a5cca94 100644 --- a/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go +++ b/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go @@ -211,6 +211,14 @@ func runConvert(cmd *cobra.Command, args []string) error { //nolint:nestif bundleInfo.ResolvedCatalogName = catalogName success(fmt.Sprintf("Selected ClusterCatalog: %s", catalogName)) + // Verify every OLMv1 prerequisite before deleting the Subscription or CSV. + // This also discovers the operator-controller namespace used by SecretPacker. + opts, err = m.PrepareClusterObjectSet(ctx, opts) + if err != nil { + return fmt.Errorf("ClusterObjectSet prerequisite check failed: %w", err) + } + success(fmt.Sprintf("ClusterObjectSet API established; using operator-controller namespace %s", opts.SystemNamespace)) + stepHeader(4, "Collecting operator resources") objects, err := m.CollectResources(ctx, opts, csv, ip, bundleInfo.PackageName) if err != nil { @@ -261,7 +269,10 @@ func runConvert(cmd *cobra.Command, args []string) error { //nolint:nestif startProgress() if err := m.CreateClusterObjectSet(ctx, opts, bundleInfo); err != nil { clearProgress() - return fmt.Errorf("COS creation failed: %w", err) + if recoverErr := m.RecoverBeforeCE(ctx, opts, backup); recoverErr != nil { + return fmt.Errorf("COS creation failed: %w; recovery also failed: %v", err, recoverErr) + } + return fmt.Errorf("COS creation failed (recovered): %w", err) } clearProgress() success(fmt.Sprintf("ClusterObjectSet %s-1 reached Succeeded=True", opts.ClusterExtensionName)) diff --git a/migration/pkg/migration/migration.go b/migration/pkg/migration/migration.go index 2aabec2..a40f1eb 100644 --- a/migration/pkg/migration/migration.go +++ b/migration/pkg/migration/migration.go @@ -8,6 +8,7 @@ import ( "strings" "time" + appsv1 "k8s.io/api/apps/v1" corev1 "k8s.io/api/core/v1" apiextensionsv1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" @@ -24,6 +25,11 @@ import ( ocv1ac "github.com/operator-framework/operator-controller/applyconfigurations/api/v1" ) +const ( + clusterObjectSetCRDName = "clusterobjectsets.olm.operatorframework.io" + operatorControllerDeployName = "operator-controller-controller-manager" +) + // annotationPrefixesToStrip are annotation prefixes that should be removed from migrated resources. var annotationPrefixesToStrip = []string{ "kubectl.kubernetes.io/", @@ -93,6 +99,13 @@ func (m *Migrator) Migrate(ctx context.Context, opts Options) error { } info.ResolvedCatalogName = catalogName + // Fail before taking OLMv0 out of management if the target API or its + // SecretPacker namespace is not available on this cluster. + opts, err = m.PrepareClusterObjectSet(ctx, opts) + if err != nil { + return err + } + backup, err := m.BackupResources(ctx, opts, csv, ip) if err != nil { return fmt.Errorf("failed to backup resources: %w", err) @@ -161,6 +174,59 @@ func (m *Migrator) Migrate(ctx context.Context, opts Options) error { return nil } +// PrepareClusterObjectSet validates that this cluster can accept the OLMv1 +// objects migration creates and resolves where SecretPacker data belongs. It +// must run before PrepareForMigration, which deletes OLMv0 management objects. +func (m *Migrator) PrepareClusterObjectSet(ctx context.Context, opts Options) (Options, error) { + opts.ApplyDefaults() + if err := m.ensureClusterObjectSetCRD(ctx); err != nil { + return opts, err + } + if opts.SystemNamespace != "" { + return opts, nil + } + namespace, err := m.operatorControllerNamespace(ctx) + if err != nil { + return opts, err + } + opts.SystemNamespace = namespace + return opts, nil +} + +func (m *Migrator) ensureClusterObjectSetCRD(ctx context.Context) error { + var crd apiextensionsv1.CustomResourceDefinition + if err := m.Client.Get(ctx, client.ObjectKey{Name: clusterObjectSetCRDName}, &crd); err != nil { + return fmt.Errorf("ClusterObjectSet CRD %q is required before migration: %w", clusterObjectSetCRDName, err) + } + for _, condition := range crd.Status.Conditions { + if condition.Type == apiextensionsv1.Established && condition.Status == apiextensionsv1.ConditionTrue { + return nil + } + } + return fmt.Errorf("ClusterObjectSet CRD %q is not established", clusterObjectSetCRDName) +} + +func (m *Migrator) operatorControllerNamespace(ctx context.Context) (string, error) { + var deployments appsv1.DeploymentList + if err := m.Client.List(ctx, &deployments, client.MatchingLabels{"app.kubernetes.io/name": "operator-controller"}); err != nil { + return "", fmt.Errorf("list operator-controller Deployments: %w", err) + } + var matches []string + for _, deployment := range deployments.Items { + if deployment.Name == operatorControllerDeployName { + matches = append(matches, deployment.Namespace) + } + } + switch len(matches) { + case 1: + return matches[0], nil + case 0: + return "", fmt.Errorf("operator-controller Deployment %q was not found; set Options.SystemNamespace only after installing a compatible operator-controller", operatorControllerDeployName) + default: + return "", fmt.Errorf("found operator-controller Deployment %q in multiple namespaces %v; set Options.SystemNamespace", operatorControllerDeployName, matches) + } +} + // EnsurePrerequisites verifies that all prerequisites for migration are met. func (m *Migrator) EnsurePrerequisites(ctx context.Context, opts Options) (*operatorsv1alpha1.ClusterServiceVersion, *operatorsv1alpha1.InstallPlan, *PreMigrationReport, *PreMigrationReport, error) { readiness, err := m.CheckReadiness(ctx, opts) @@ -355,6 +421,11 @@ func (m *Migrator) ensureClusterExtensionAbsent(ctx context.Context, name string // 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 { + var err error + opts, err = m.PrepareClusterObjectSet(ctx, opts) + if err != nil { + return err + } resources, err := m.createClusterObjectSet(ctx, opts, info) if err != nil { if cleanupErr := m.cleanupCreatedClusterObjectSet(ctx, resources); cleanupErr != nil { @@ -368,7 +439,6 @@ func (m *Migrator) createClusterObjectSet(ctx context.Context, opts Options, inf resources := &createdMigrationResources{} invocationMarker := string(uuid.NewUUID()) cosName := fmt.Sprintf("%s-1", opts.ClusterExtensionName) - systemNS := opts.systemNamespace() cosObjects := make([]ocv1ac.ClusterObjectSetObjectApplyConfiguration, 0, len(info.CollectedObjects)) for _, obj := range info.CollectedObjects { @@ -388,7 +458,7 @@ func (m *Migrator) createClusterObjectSet(ctx context.Context, opts Options, inf packer := &secretPacker{ RevisionName: cosName, OwnerName: opts.ClusterExtensionName, - SystemNamespace: systemNS, + SystemNamespace: opts.SystemNamespace, } packed, err := packer.pack(phases) if err != nil { diff --git a/migration/pkg/migration/readiness.go b/migration/pkg/migration/readiness.go index 9c696d9..43a8277 100644 --- a/migration/pkg/migration/readiness.go +++ b/migration/pkg/migration/readiness.go @@ -13,6 +13,19 @@ import ( // It checks Subscription state, CSV health, uniqueness, and dependency status. func (m *Migrator) CheckReadiness(ctx context.Context, opts Options) (*PreMigrationReport, error) { report := &PreMigrationReport{} + if err := m.ensureClusterObjectSetCRD(ctx); err != nil { + report.Checks = append(report.Checks, CheckResult{ + Name: "ClusterObjectSet API", + Passed: false, + Message: err.Error(), + }) + } else { + report.Checks = append(report.Checks, CheckResult{ + Name: "ClusterObjectSet API", + Passed: true, + Message: "ClusterObjectSet CRD is established", + }) + } var sub operatorsv1alpha1.Subscription if err := m.Client.Get(ctx, types.NamespacedName{ diff --git a/migration/pkg/migration/secretpacker.go b/migration/pkg/migration/secretpacker.go index edf47e8..958be16 100644 --- a/migration/pkg/migration/secretpacker.go +++ b/migration/pkg/migration/secretpacker.go @@ -39,7 +39,8 @@ type secretPacker struct { RevisionName string // OwnerName is the CE name — recorded as a label on each Secret. OwnerName string - // SystemNamespace is where Secrets are created (e.g. "olmv1-system"). + // SystemNamespace is where Secrets are created (the discovered + // operator-controller namespace, unless explicitly overridden). SystemNamespace string } diff --git a/migration/pkg/migration/types.go b/migration/pkg/migration/types.go index b80f55f..cff567f 100644 --- a/migration/pkg/migration/types.go +++ b/migration/pkg/migration/types.go @@ -5,6 +5,7 @@ import ( "os" "path/filepath" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/types" "k8s.io/client-go/rest" @@ -51,18 +52,11 @@ type Options struct { DeleteOperatorGroup bool // SystemNamespace is the namespace where COS ref Secrets are created (R2.4). - // Defaults to "olmv1-system" when empty. + // When empty, migration discovers the operator-controller Deployment namespace. + // It is an override for unusual installations, not a production default. SystemNamespace string } -// systemNamespace returns the effective system namespace. -func (o Options) systemNamespace() string { - if o.SystemNamespace != "" { - return o.SystemNamespace - } - return "olmv1-system" -} - // ApplyDefaults fills in default values for any unset optional fields. func (o *Options) ApplyDefaults() { if o.ClusterExtensionName == "" { @@ -131,15 +125,21 @@ func (b *Backup) SaveToDisk(dir string) error { if err := os.MkdirAll(dir, 0o750); err != nil { return fmt.Errorf("failed to create backup directory: %w", err) } - if err := writeYAMLFile(filepath.Join(dir, "subscription.yaml"), b.Subscription); err != nil { + subscription := b.Subscription.DeepCopy() + subscription.TypeMeta = metav1.TypeMeta{APIVersion: "operators.coreos.com/v1alpha1", Kind: "Subscription"} + if err := writeYAMLFile(filepath.Join(dir, "subscription.yaml"), subscription); err != nil { return fmt.Errorf("failed to write subscription.yaml: %w", err) } if b.OperatorGroup != nil { - if err := writeYAMLFile(filepath.Join(dir, "operatorgroup.yaml"), b.OperatorGroup); err != nil { + operatorGroup := b.OperatorGroup.DeepCopy() + operatorGroup.TypeMeta = metav1.TypeMeta{APIVersion: "operators.coreos.com/v1", Kind: "OperatorGroup"} + if err := writeYAMLFile(filepath.Join(dir, "operatorgroup.yaml"), operatorGroup); err != nil { return fmt.Errorf("failed to write operatorgroup.yaml: %w", err) } } - if err := writeYAMLFile(filepath.Join(dir, "clusterserviceversion.yaml"), b.ClusterServiceVersion); err != nil { + csv := b.ClusterServiceVersion.DeepCopy() + csv.TypeMeta = metav1.TypeMeta{APIVersion: "operators.coreos.com/v1alpha1", Kind: "ClusterServiceVersion"} + if err := writeYAMLFile(filepath.Join(dir, "clusterserviceversion.yaml"), csv); err != nil { return fmt.Errorf("failed to write clusterserviceversion.yaml: %w", err) } if b.InstallPlan != nil { @@ -147,7 +147,9 @@ func (b *Backup) SaveToDisk(dir string) error { if err := os.MkdirAll(ipDir, 0o750); err != nil { return fmt.Errorf("failed to create installplans directory: %w", err) } - if err := writeYAMLFile(filepath.Join(ipDir, b.InstallPlan.Name+".yaml"), b.InstallPlan); err != nil { + installPlan := b.InstallPlan.DeepCopy() + installPlan.TypeMeta = metav1.TypeMeta{APIVersion: "operators.coreos.com/v1alpha1", Kind: "InstallPlan"} + if err := writeYAMLFile(filepath.Join(ipDir, b.InstallPlan.Name+".yaml"), installPlan); err != nil { return fmt.Errorf("failed to write installplan: %w", err) } } diff --git a/migration/pkg/migration/unit_test.go b/migration/pkg/migration/unit_test.go index b1e1d1c..6a98f81 100644 --- a/migration/pkg/migration/unit_test.go +++ b/migration/pkg/migration/unit_test.go @@ -14,8 +14,10 @@ import ( "strings" "testing" + appsv1 "k8s.io/api/apps/v1" corev1 "k8s.io/api/core/v1" rbacv1 "k8s.io/api/rbac/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" @@ -87,6 +89,12 @@ func migrationTestClient(t *testing.T, objects ...runtime.Object) *Migrator { if err := ocv1.AddToScheme(scheme); err != nil { t.Fatal(err) } + if err := appsv1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } + if err := apiextensionsv1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } return NewMigrator(fake.NewClientBuilder().WithScheme(scheme).WithRuntimeObjects(objects...).Build(), nil) } @@ -96,9 +104,19 @@ func healthySubscriptionFixtures() (*operatorsv1alpha1.Subscription, *operatorsv return sub, csv } +func establishedClusterObjectSetCRD() *apiextensionsv1.CustomResourceDefinition { + return &apiextensionsv1.CustomResourceDefinition{ + ObjectMeta: metav1.ObjectMeta{Name: clusterObjectSetCRDName}, + Status: apiextensionsv1.CustomResourceDefinitionStatus{Conditions: []apiextensionsv1.CustomResourceDefinitionCondition{{ + Type: apiextensionsv1.Established, + Status: apiextensionsv1.ConditionTrue, + }}}, + } +} + func TestCheckReadiness(t *testing.T) { sub, csv := healthySubscriptionFixtures() - m := migrationTestClient(t, sub, csv) + m := migrationTestClient(t, sub, csv, establishedClusterObjectSetCRD()) report, err := m.CheckReadiness(context.Background(), Options{SubscriptionName: "sub", SubscriptionNamespace: "ns"}) if err != nil || !report.Passed() { t.Fatalf("healthy report=%#v err=%v", report, err) @@ -107,7 +125,7 @@ func TestCheckReadiness(t *testing.T) { busySub, busyCSV := healthySubscriptionFixtures() busySub.Status.State = operatorsv1alpha1.SubscriptionStateUpgradeAvailable busyCSV.Status.Phase = operatorsv1alpha1.CSVPhaseFailed - m = migrationTestClient(t, busySub, busyCSV) + m = migrationTestClient(t, busySub, busyCSV, establishedClusterObjectSetCRD()) report, err = m.CheckReadiness(context.Background(), Options{SubscriptionName: "sub", SubscriptionNamespace: "ns"}) if err != nil || report.Passed() || len(report.FailedChecks()) != 2 { t.Fatalf("unacknowledged report=%#v err=%v", report, err) @@ -119,13 +137,31 @@ func TestCheckReadiness(t *testing.T) { depSub, depCSV := healthySubscriptionFixtures() depSub.Annotations = map[string]string{"olm.generated-by": "parent"} - m = migrationTestClient(t, depSub, depCSV, &operatorsv1alpha1.Subscription{ObjectMeta: metav1.ObjectMeta{Name: "other", Namespace: "other"}, Spec: &operatorsv1alpha1.SubscriptionSpec{Package: "widgets"}}) + m = migrationTestClient(t, depSub, depCSV, establishedClusterObjectSetCRD(), &operatorsv1alpha1.Subscription{ObjectMeta: metav1.ObjectMeta{Name: "other", Namespace: "other"}, Spec: &operatorsv1alpha1.SubscriptionSpec{Package: "widgets"}}) report, err = m.CheckReadiness(context.Background(), Options{SubscriptionName: "sub", SubscriptionNamespace: "ns"}) if err != nil || report.Passed() || len(report.FailedChecks()) != 2 { t.Fatalf("dependency/duplicate report=%#v err=%v", report, err) } } +func TestCheckReadinessRequiresClusterObjectSetCRD(t *testing.T) { + sub, csv := healthySubscriptionFixtures() + m := migrationTestClient(t, sub, csv) + report, err := m.CheckReadiness(context.Background(), Options{SubscriptionName: "sub", SubscriptionNamespace: "ns"}) + if err != nil { + t.Fatalf("CheckReadiness() error = %v", err) + } + for _, check := range report.Checks { + if check.Name == "ClusterObjectSet API" { + if check.Passed || !strings.Contains(check.Message, clusterObjectSetCRDName) { + t.Fatalf("ClusterObjectSet API check = %#v, want failed missing-CRD check", check) + } + return + } + } + t.Fatal("CheckReadiness() did not report the ClusterObjectSet API check") +} + func TestCompatibilityOperatorGroupAndConditionOverrides(t *testing.T) { og := &operatorsv1.OperatorGroup{ObjectMeta: metav1.ObjectMeta{Name: "og", Namespace: "ns"}, Spec: operatorsv1.OperatorGroupSpec{TargetNamespaces: []string{"target"}, ServiceAccountName: "restricted"}} csv := &operatorsv1alpha1.ClusterServiceVersion{ObjectMeta: metav1.ObjectMeta{Name: "widgets.v1", Namespace: "ns"}} @@ -296,7 +332,7 @@ func TestSecretPacker(t *testing.T) { func TestOptionsReportsAndAnnotationFiltering(t *testing.T) { opts := Options{SubscriptionName: "sub", SubscriptionNamespace: "ns"} opts.ApplyDefaults() - if opts.ClusterExtensionName != "sub" || opts.InstallNamespace != "ns" || opts.systemNamespace() != "olmv1-system" { + if opts.ClusterExtensionName != "sub" || opts.InstallNamespace != "ns" || opts.SystemNamespace != "" { t.Fatalf("defaults: %#v", opts) } report := &PreMigrationReport{Checks: []CheckResult{{Passed: true}, {Name: "bad", Passed: false}}} @@ -309,6 +345,54 @@ func TestOptionsReportsAndAnnotationFiltering(t *testing.T) { } } +func TestPrepareClusterObjectSet(t *testing.T) { + establishedCRD := &apiextensionsv1.CustomResourceDefinition{ + ObjectMeta: metav1.ObjectMeta{Name: clusterObjectSetCRDName}, + Status: apiextensionsv1.CustomResourceDefinitionStatus{Conditions: []apiextensionsv1.CustomResourceDefinitionCondition{{ + Type: apiextensionsv1.Established, + Status: apiextensionsv1.ConditionTrue, + }}}, + } + controller := &appsv1.Deployment{ + ObjectMeta: metav1.ObjectMeta{ + Name: operatorControllerDeployName, + Namespace: "openshift-operator-controller", + Labels: map[string]string{"app.kubernetes.io/name": "operator-controller"}, + }, + } + m := migrationTestClient(t, establishedCRD, controller) + + got, err := m.PrepareClusterObjectSet(context.Background(), Options{SubscriptionName: "sub", SubscriptionNamespace: "operators"}) + if err != nil { + t.Fatalf("PrepareClusterObjectSet() error = %v", err) + } + if got.SystemNamespace != "openshift-operator-controller" { + t.Fatalf("SystemNamespace = %q, want discovered operator-controller namespace", got.SystemNamespace) + } + + got, err = m.PrepareClusterObjectSet(context.Background(), Options{SystemNamespace: "chosen"}) + if err != nil { + t.Fatalf("PrepareClusterObjectSet() with override error = %v", err) + } + if got.SystemNamespace != "chosen" { + t.Fatalf("SystemNamespace override = %q, want chosen", got.SystemNamespace) + } +} + +func TestPrepareClusterObjectSetRequiresEstablishedCRD(t *testing.T) { + m := migrationTestClient(t) + if _, err := m.PrepareClusterObjectSet(context.Background(), Options{}); err == nil || !strings.Contains(err.Error(), clusterObjectSetCRDName) { + t.Fatalf("PrepareClusterObjectSet() error = %v, want missing CRD error", err) + } + + m = migrationTestClient(t, &apiextensionsv1.CustomResourceDefinition{ + ObjectMeta: metav1.ObjectMeta{Name: clusterObjectSetCRDName}, + }) + if _, err := m.PrepareClusterObjectSet(context.Background(), Options{}); err == nil || !strings.Contains(err.Error(), "not established") { + t.Fatalf("PrepareClusterObjectSet() error = %v, want unestablished CRD error", err) + } +} + func TestEmptyLabelSelector(t *testing.T) { for name, selector := range map[string]*metav1.LabelSelector{ "nil": nil, @@ -360,6 +444,9 @@ func TestBackupSaveToDisk(t *testing.T) { if err != nil || !strings.Contains(string(data), name) { t.Errorf("backup %s = %q, err=%v", file, data, err) } + if !strings.Contains(string(data), "apiVersion: operators.coreos.com/") || !strings.Contains(string(data), "kind:") { + t.Errorf("backup %s is not an applyable Kubernetes manifest: %q", file, data) + } } } @@ -480,7 +567,7 @@ func TestScanAllKeepsMixedUnsafeOperatorsOutOfEligibleResults(t *testing.T) { func TestPrerequisitesAndRecoveryErrors(t *testing.T) { ctx := context.Background() sub, csv := healthySubscriptionFixtures() - m := migrationTestClient(t, sub, csv) + m := migrationTestClient(t, sub, csv, establishedClusterObjectSetCRD()) gotCSV, ip, readiness, compatibility, err := m.EnsurePrerequisites(ctx, Options{SubscriptionName: "sub", SubscriptionNamespace: "ns"}) if err != nil || gotCSV.Name != csv.Name || ip != nil || !readiness.Passed() || compatibility == nil { t.Fatalf("EnsurePrerequisites() = csv=%#v ip=%#v readiness=%#v compatibility=%#v err=%v", gotCSV, ip, readiness, compatibility, err) From 5dfc385003498fe58157849a21a1599662a84f87 Mon Sep 17 00:00:00 2001 From: Todd Short Date: Fri, 18 Sep 2026 16:31:43 -0400 Subject: [PATCH 4/6] fix: preflight COS for migration dry runs Signed-off-by: Todd Short --- .../examples/cmd/migrate-operators-v0-to-v1/convert.go | 9 +++++++++ 1 file changed, 9 insertions(+) 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 a5cca94..588f90e 100644 --- a/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go +++ b/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go @@ -308,6 +308,15 @@ func runConvertDryRun(cmd *cobra.Command, m *migration.Migrator, opts migration. ctx := cmd.Context() fmt.Printf("\n%s%s🔍 Dry run: %s/%s%s\n", colorBold, colorCyan, opts.SubscriptionNamespace, opts.SubscriptionName, colorReset) + // Dry-run must reject a target that cannot create a COS, just as a real + // conversion would. This is read-only and runs before gathering the preview. + var err error + opts, err = m.PrepareClusterObjectSet(ctx, opts) + if err != nil { + return fmt.Errorf("ClusterObjectSet prerequisite check failed: %w", err) + } + success(fmt.Sprintf("ClusterObjectSet API established; using operator-controller namespace %s", opts.SystemNamespace)) + info, err := m.GatherMigrationInfo(ctx, opts) if err != nil { return fmt.Errorf("failed to gather migration info: %w", err) From 120f5c0f1690869b7d03c91f58c4d33f95e03f92 Mon Sep 17 00:00:00 2001 From: Todd Short Date: Mon, 21 Sep 2026 14:46:07 -0400 Subject: [PATCH 5/6] test: satisfy COS preflight in recovery tests Signed-off-by: Todd Short --- migration/pkg/migration/unit_test.go | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/migration/pkg/migration/unit_test.go b/migration/pkg/migration/unit_test.go index 6a98f81..91413d8 100644 --- a/migration/pkg/migration/unit_test.go +++ b/migration/pkg/migration/unit_test.go @@ -676,12 +676,12 @@ func TestRollbackAndCleanupRejectInvalidInputWithoutMutation(t *testing.T) { func TestCreateClusterObjectSetCleansTemporarySecretsOnCollision(t *testing.T) { ctx := context.Background() - m := migrationTestClient(t) + 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"}, }} - err := m.CreateClusterObjectSet(ctx, Options{SubscriptionName: "sub", SubscriptionNamespace: "ns"}, &MigrationInfo{ + err := m.CreateClusterObjectSet(ctx, Options{SubscriptionName: "sub", SubscriptionNamespace: "ns", SystemNamespace: "olmv1-system"}, &MigrationInfo{ PackageName: "widgets", BundleName: "widgets.v1", Version: "1.0.0", CollectedObjects: []unstructured.Unstructured{object}, }) if err == nil || !apierrors.IsAlreadyExists(err) { @@ -701,12 +701,12 @@ func TestCreateClusterObjectSetCleansTemporarySecretsOnCollision(t *testing.T) { func TestCreateClusterObjectSetPreservesSecretsAfterUnknownCreateOutcome(t *testing.T) { ctx := context.Background() - m := migrationTestClient(t) + m := migrationTestClient(t, establishedClusterObjectSetCRD()) 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{ + err := m.CreateClusterObjectSet(ctx, Options{SubscriptionName: "sub", SubscriptionNamespace: "ns", SystemNamespace: "olmv1-system"}, &MigrationInfo{ PackageName: "widgets", BundleName: "widgets.v1", Version: "1.0.0", CollectedObjects: []unstructured.Unstructured{object}, }) if err == nil || !strings.Contains(err.Error(), "refusing automatic cleanup") { From e8a7a91a234dbca2945340dabce1c9430e009f3f Mon Sep 17 00:00:00 2001 From: Todd Short Date: Mon, 21 Sep 2026 15:22:46 -0400 Subject: [PATCH 6/6] fix: recover CLI migrations after extension failure Signed-off-by: Todd Short --- .../cmd/migrate-operators-v0-to-v1/convert.go | 21 ++------ migration/pkg/migration/migration.go | 49 ++++++++++++------- migration/pkg/migration/unit_test.go | 38 +++++++++++++- 3 files changed, 72 insertions(+), 36 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 588f90e..d981329 100644 --- a/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go +++ b/migration/examples/cmd/migrate-operators-v0-to-v1/convert.go @@ -264,29 +264,18 @@ func runConvert(cmd *cobra.Command, args []string) error { //nolint:nestif } success("OLMv0 management removed") - stepHeader(7, "Creating ClusterObjectSet") - info(fmt.Sprintf("Applying COS %s-1 with %d objects...", opts.ClusterExtensionName, len(bundleInfo.CollectedObjects))) + 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.CreateClusterObjectSet(ctx, opts, bundleInfo); err != nil { + if err := m.CreateMigrationResources(ctx, opts, bundleInfo, backup); err != nil { clearProgress() - if recoverErr := m.RecoverBeforeCE(ctx, opts, backup); recoverErr != nil { - return fmt.Errorf("COS creation failed: %w; recovery also failed: %v", err, recoverErr) - } - return fmt.Errorf("COS creation failed (recovered): %w", err) + return err } clearProgress() success(fmt.Sprintf("ClusterObjectSet %s-1 reached Succeeded=True", opts.ClusterExtensionName)) - - stepHeader(8, "Creating ClusterExtension") - startProgress() - if err := m.CreateClusterExtension(ctx, opts, bundleInfo); err != nil { - clearProgress() - return fmt.Errorf("failed to create ClusterExtension: %w", err) - } - clearProgress() success(fmt.Sprintf("ClusterExtension %s is Installed", opts.ClusterExtensionName)) - stepHeader(9, "Cleaning up OLMv0 resources") + stepHeader(8, "Cleaning up OLMv0 resources") cleanupResult := m.CleanupOLMv0Resources(ctx, opts, bundleInfo.PackageName, csv.Name) for _, action := range cleanupResult.Actions { switch { diff --git a/migration/pkg/migration/migration.go b/migration/pkg/migration/migration.go index a40f1eb..c2fe00a 100644 --- a/migration/pkg/migration/migration.go +++ b/migration/pkg/migration/migration.go @@ -149,24 +149,8 @@ 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") - 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) - } - - 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) + if err := m.CreateMigrationResources(ctx, opts, info, backup); err != nil { + return err } m.CleanupOLMv0Resources(ctx, opts, info.PackageName, csv.Name) @@ -179,6 +163,9 @@ func (m *Migrator) Migrate(ctx context.Context, opts Options) error { // must run before PrepareForMigration, which deletes OLMv0 management objects. func (m *Migrator) PrepareClusterObjectSet(ctx context.Context, opts Options) (Options, error) { opts.ApplyDefaults() + if err := m.ensureClusterExtensionAbsent(ctx, opts.ClusterExtensionName); err != nil { + return opts, err + } if err := m.ensureClusterObjectSetCRD(ctx); err != nil { return opts, err } @@ -435,6 +422,32 @@ func (m *Migrator) CreateClusterObjectSet(ctx context.Context, opts Options, inf return err } +// 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 { + 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) + } + + 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) + } + return nil +} + func (m *Migrator) createClusterObjectSet(ctx context.Context, opts Options, info *MigrationInfo) (*createdMigrationResources, error) { resources := &createdMigrationResources{} invocationMarker := string(uuid.NewUUID()) diff --git a/migration/pkg/migration/unit_test.go b/migration/pkg/migration/unit_test.go index 91413d8..a5944d1 100644 --- a/migration/pkg/migration/unit_test.go +++ b/migration/pkg/migration/unit_test.go @@ -40,6 +40,7 @@ type failingMigrationClient struct { unknownCOSCreate bool failCECreate bool blockCOS bool + succeedCOS bool } func (c failingMigrationClient) Create(ctx context.Context, obj client.Object, opts ...client.CreateOption) error { @@ -58,6 +59,11 @@ func (c failingMigrationClient) Create(ctx context.Context, obj client.Object, o return apierrors.NewAlreadyExists(ocv1.GroupVersion.WithResource("clusterextensions").GroupResource(), obj.GetName()) } } + if c.succeedCOS { + if cos, ok := obj.(*ocv1.ClusterObjectSet); ok { + cos.Status.Conditions = []metav1.Condition{{Type: ocv1.ClusterObjectSetTypeSucceeded, Status: metav1.ConditionTrue}} + } + } return c.Client.Create(ctx, obj, opts...) } @@ -379,6 +385,14 @@ func TestPrepareClusterObjectSet(t *testing.T) { } } +func TestPrepareClusterObjectSetRejectsExistingClusterExtension(t *testing.T) { + m := migrationTestClient(t, &ocv1.ClusterExtension{ObjectMeta: metav1.ObjectMeta{Name: "sub"}}) + _, err := m.PrepareClusterObjectSet(context.Background(), Options{SubscriptionName: "sub", SubscriptionNamespace: "operators"}) + if err == nil || !strings.Contains(err.Error(), "ClusterExtension sub already exists") { + t.Fatalf("PrepareClusterObjectSet() error = %v, want existing ClusterExtension error", err) + } +} + func TestPrepareClusterObjectSetRequiresEstablishedCRD(t *testing.T) { m := migrationTestClient(t) if _, err := m.PrepareClusterObjectSet(context.Background(), Options{}); err == nil || !strings.Contains(err.Error(), clusterObjectSetCRDName) { @@ -734,15 +748,35 @@ func TestCreateClusterExtensionAlreadyExistsIsKnownOutcome(t *testing.T) { } } -func TestCreateClusterObjectSetCleansSecretsAfterReadinessFailure(t *testing.T) { +func TestCreateMigrationResourcesCleansTrackedCOSAfterClusterExtensionFailure(t *testing.T) { ctx := context.Background() m := migrationTestClient(t) baseClient := m.Client + m.Client = failingMigrationClient{Client: baseClient, succeedCOS: true, failCECreate: true} + 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{ + 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) + } + 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") + } +} + +func TestCreateClusterObjectSetCleansSecretsAfterReadinessFailure(t *testing.T) { + ctx := context.Background() + m := migrationTestClient(t, establishedClusterObjectSetCRD()) + 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{ + err := m.CreateClusterObjectSet(ctx, Options{SubscriptionName: "sub", SubscriptionNamespace: "ns", SystemNamespace: "olmv1-system"}, &MigrationInfo{ PackageName: "widgets", BundleName: "widgets.v1", Version: "1.0.0", CollectedObjects: []unstructured.Unstructured{object}, }) if err == nil {