Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
53 changes: 49 additions & 4 deletions migration/examples/cmd/migrate-operators-v0-to-v1/convert.go
Original file line number Diff line number Diff line change
@@ -1,11 +1,14 @@
package main

import (
"context"
"encoding/json"
"errors"
"fmt"
"time"

"github.com/spf13/cobra"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"

"github.com/operator-framework/library-olm/migration/pkg/migration"
)
Expand Down Expand Up @@ -70,6 +73,9 @@ func runConvert(cmd *cobra.Command, args []string) error { //nolint:nestif
if !convertAll && len(args) == 0 {
return fmt.Errorf("specify an operator name or --all")
}
if convertAll && convertInstallNs != "" {
return fmt.Errorf("--install-namespace requires a single operator")
}

c, restCfg, err := newClient()
if err != nil {
Expand Down Expand Up @@ -152,7 +158,6 @@ func runConvert(cmd *cobra.Command, args []string) error { //nolint:nestif
AcknowledgeNotSteadyState: convertAckNotSteady,
}
opts.ApplyDefaults()

if convertDryRun {
return runConvertDryRun(cmd, m, opts)
}
Expand Down Expand Up @@ -218,12 +223,20 @@ func runConvert(cmd *cobra.Command, args []string) error { //nolint:nestif
return fmt.Errorf("ClusterObjectSet prerequisite check failed: %w", err)
}
success(fmt.Sprintf("ClusterObjectSet API established; using operator-controller namespace %s", opts.SystemNamespace))
if err := m.PrepareInstallNamespace(ctx, opts); err != nil {
return fmt.Errorf("install namespace preparation failed: %w", err)
}

stepHeader(4, "Collecting operator resources")
objects, err := m.CollectResources(ctx, opts, csv, ip, bundleInfo.PackageName)
if err != nil {
return fmt.Errorf("failed to collect resources: %w", err)
}
sourceObjects := make([]unstructured.Unstructured, len(objects))
for i := range objects {
sourceObjects[i] = *objects[i].DeepCopy()
}
migration.RewriteInstallNamespace(objects, opts.SubscriptionNamespace, opts.InstallNamespace)
bundleInfo.CollectedObjects = objects
kindCounts := make(map[string]int)
for _, obj := range objects {
Expand Down Expand Up @@ -258,22 +271,46 @@ func runConvert(cmd *cobra.Command, args []string) error { //nolint:nestif
success("Resources backed up in memory (CE annotation backup authoritative)")

stepHeader(6, "Preparing operator for migration")
info("Deleting Subscription and CSV (orphan cascade — workloads keep running)...")
info("Deleting Subscription and CSV (orphan cascade)...")
if err := m.PrepareForMigration(ctx, opts, csv); err != nil {
return fmt.Errorf("preparation failed: %w", err)
}
success("OLMv0 management removed")
restoreSourceDeployments, err := m.ScaleSourceDeployments(ctx, sourceObjects, opts)
if err != nil {
recoveryCtx, cancel := migration.NewRecoveryContext(ctx)
defer cancel()
if recoverErr := m.RecoverFromBackup(recoveryCtx, opts, backup); recoverErr != nil {
return fmt.Errorf("scale source Deployments: %w; recovery also failed: %v", err, recoverErr)
}
return fmt.Errorf("scale source Deployments failed (recovered): %w", err)
}
if opts.InstallNamespace != opts.SubscriptionNamespace {
success("Source Deployments scaled to zero before target cutover")
}

stepHeader(7, "Creating OLMv1 migration resources")
info(fmt.Sprintf("Applying COS %s-1 with %d objects and creating its ClusterExtension...", opts.ClusterExtensionName, len(bundleInfo.CollectedObjects)))
startProgress()
if err := m.CreateMigrationResources(ctx, opts, bundleInfo, backup); err != nil {
result, err := m.CreateMigrationResources(ctx, opts, bundleInfo, backup)
if err != nil {
clearProgress()
if result.TargetMayBeActive {
return fmt.Errorf("create migration resources: %w; target ClusterObjectSet may be active, so source Deployments remain scaled to zero", err)
}
restoreCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second)
defer cancel()
if restoreErr := restoreSourceDeployments(restoreCtx); restoreErr != nil {
return fmt.Errorf("create migration resources: %w; restore source Deployments: %v", err, restoreErr)
}
return err
}
clearProgress()
success(fmt.Sprintf("ClusterObjectSet %s-1 reached Succeeded=True", opts.ClusterExtensionName))
success(fmt.Sprintf("ClusterExtension %s is Installed", opts.ClusterExtensionName))
if err := m.DeleteSourceNamespaceResources(ctx, sourceObjects, opts); err != nil {
return fmt.Errorf("delete source install resources: %w", err)
}

stepHeader(8, "Cleaning up OLMv0 resources")
cleanupResult := m.CleanupOLMv0Resources(ctx, opts, bundleInfo.PackageName, csv.Name)
Expand All @@ -287,7 +324,9 @@ func runConvert(cmd *cobra.Command, args []string) error { //nolint:nestif
success(action.Description)
}
}

if err := cleanupResult.Err(); err != nil {
return fmt.Errorf("clean up OLMv0 resources: %w", err)
}
banner(fmt.Sprintf("Migration complete! %s is now managed by OLMv1", bundleInfo.PackageName))
fmt.Println()
return nil
Expand Down Expand Up @@ -368,6 +407,12 @@ func dryRunCleanupPlan(opts migration.Options, info *migration.MigrationInfo) []
} else {
lines = append(lines, "Retain OperatorGroup(s); --delete-operatorgroup was not specified.")
}
if opts.InstallNamespace != opts.SubscriptionNamespace {
lines = append(lines,
fmt.Sprintf("Create or update install namespace %s with PSA/SCC labels copied from %s.", opts.InstallNamespace, opts.SubscriptionNamespace),
fmt.Sprintf("Move collected namespaced operator resources from %s to %s and delete their source copies after ClusterExtension installation.", opts.SubscriptionNamespace, opts.InstallNamespace))
lines = append(lines, fmt.Sprintf("Retain source namespace %s.", opts.SubscriptionNamespace))
}
return lines
}

Expand Down
19 changes: 19 additions & 0 deletions migration/examples/cmd/migrate-operators-v0-to-v1/convert_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -39,3 +39,22 @@ func TestDryRunCleanupPlan(t *testing.T) {
t.Fatalf("dry-run cleanup plan does not describe the default OperatorGroup behavior:\n%s", plan)
}
}

func TestDryRunCleanupPlanNamespaceChange(t *testing.T) {
opts := migration.Options{
SubscriptionName: "widget-operator",
SubscriptionNamespace: "operators",
InstallNamespace: "widget-system",
}
info := &migration.MigrationInfo{PackageName: "widgets", BundleName: "widgets.v1.2.3"}
plan := strings.Join(dryRunCleanupPlan(opts, info), "\n")
for _, expected := range []string{
"Create or update install namespace widget-system with PSA/SCC labels copied from operators.",
"Move collected namespaced operator resources from operators to widget-system and delete their source copies after ClusterExtension installation.",
"Retain source namespace operators.",
} {
if !strings.Contains(plan, expected) {
t.Fatalf("dry-run cleanup plan missing %q:\n%s", expected, plan)
}
}
}
95 changes: 78 additions & 17 deletions migration/pkg/migration/migration.go
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,13 @@ type createdMigrationResources struct {
ownershipUnknown bool
}

// MigrationResourcesResult records whether target resources may have started
// reconciling when CreateMigrationResources returns. Callers must not restart
// scaled source controllers while TargetMayBeActive is true.
type MigrationResourcesResult struct {
TargetMayBeActive bool
}

const migrationInvocationAnnotation = "olm.operatorframework.io/migration-invocation"

// Migrate performs the full migration of an OLMv0-managed operator to OLMv1.
Expand Down Expand Up @@ -105,6 +112,9 @@ func (m *Migrator) Migrate(ctx context.Context, opts Options) error {
if err != nil {
return err
}
if err := m.PrepareInstallNamespace(ctx, opts); err != nil {
return err
}

backup, err := m.BackupResources(ctx, opts, csv, ip)
if err != nil {
Expand Down Expand Up @@ -140,6 +150,11 @@ func (m *Migrator) Migrate(ctx context.Context, opts Options) error {
if err != nil {
return fmt.Errorf("failed to collect resources: %w", err)
}
sourceObjects := make([]unstructured.Unstructured, len(objects))
for i := range objects {
sourceObjects[i] = *objects[i].DeepCopy()
}
RewriteInstallNamespace(objects, opts.SubscriptionNamespace, opts.InstallNamespace)
info.CollectedObjects = objects

// R9: warn about TLS certificate pivot. OLMv0 manages certs directly via its own
Expand All @@ -149,12 +164,34 @@ func (m *Migrator) Migrate(ctx context.Context, opts Options) error {
m.progress("Note: TLS certificate management will transfer from OLMv0 to cert-manager/service-ca; " +
"expect pod restarts while new cert secrets are provisioned")

if err := m.CreateMigrationResources(ctx, opts, info, backup); err != nil {
restoreSourceDeployments, err := m.ScaleSourceDeployments(ctx, sourceObjects, opts)
if err != nil {
recoveryCtx, cancel := NewRecoveryContext(ctx)
defer cancel()
if recoverErr := m.RecoverFromBackup(recoveryCtx, opts, backup); recoverErr != nil {
return fmt.Errorf("scale source Deployments: %w; recovery also failed: %v", err, recoverErr)
}
return fmt.Errorf("scale source Deployments failed (recovered): %w", err)
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
result, err := m.CreateMigrationResources(ctx, opts, info, backup)
if err != nil {
if result.TargetMayBeActive {
return fmt.Errorf("create migration resources: %w; target ClusterObjectSet may be active, so source Deployments remain scaled to zero", err)
}
restoreCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second)
defer cancel()
if restoreErr := restoreSourceDeployments(restoreCtx); restoreErr != nil {
return fmt.Errorf("create migration resources: %w; restore source Deployments: %v", err, restoreErr)
}
return err
}
if err := m.DeleteSourceNamespaceResources(ctx, sourceObjects, opts); err != nil {
return err
}

m.CleanupOLMv0Resources(ctx, opts, info.PackageName, csv.Name)

if err := m.CleanupOLMv0Resources(ctx, opts, info.PackageName, csv.Name).Err(); err != nil {
return fmt.Errorf("clean up OLMv0 resources: %w", err)
}
return nil
}

Expand Down Expand Up @@ -263,7 +300,8 @@ func (m *Migrator) BackupResources(ctx context.Context, opts Options, csv *opera
}

// PrepareForMigration removes OLMv0 management of the operator by deleting
// the Subscription and CSV with orphan cascading (operator workloads keep running).
// the Subscription and CSV with orphan cascading. A later cross-namespace
// cutover scales source Deployments down before creating target resources.
func (m *Migrator) PrepareForMigration(ctx context.Context, opts Options, csv *operatorsv1alpha1.ClusterServiceVersion) error {
// Delete Subscription with orphan cascading
sub := &operatorsv1alpha1.Subscription{}
Expand All @@ -285,6 +323,12 @@ func (m *Migrator) PrepareForMigration(ctx context.Context, opts Options, csv *o
return nil
}

// NewRecoveryContext returns a cancellation-independent context with enough
// time to recreate and observe a restored Subscription.
func NewRecoveryContext(ctx context.Context) (context.Context, context.CancelFunc) {
return context.WithTimeout(context.WithoutCancel(ctx), subWaitTimeout+30*time.Second)
}

// RecoverFromBackup restores the Subscription from backup after a failed preparation.
func (m *Migrator) RecoverFromBackup(ctx context.Context, opts Options, backup *Backup) error {
if backup == nil {
Expand Down Expand Up @@ -331,11 +375,11 @@ func (m *Migrator) RecoverBeforeCE(ctx context.Context, opts Options, backup *Ba
}

// recoverCreatedMigrationResources removes objects created by this invocation
// in reverse dependency order, then restores the OLMv0 Subscription. A failure
// to delete a ClusterExtension or COS leaves the remaining resources intact and
// prevents restoration, avoiding concurrent OLMv0 and OLMv1 ownership.
// in reverse dependency order. Once a COS may have reconciled, it refuses to
// restore the OLMv0 Subscription automatically because orphaned target objects
// can still be running while deletion propagates.
func (m *Migrator) recoverCreatedMigrationResources(ctx context.Context, opts Options, backup *Backup, resources *createdMigrationResources) error {
recoveryCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), subWaitTimeout+30*time.Second)
recoveryCtx, cancel := NewRecoveryContext(ctx)
defer cancel()
if resources == nil {
return m.RecoverBeforeCE(recoveryCtx, opts, backup)
Expand All @@ -354,6 +398,9 @@ func (m *Migrator) recoverCreatedMigrationResources(ctx context.Context, opts Op
}
}
cleanupErr := m.cleanupCreatedSecrets(recoveryCtx, resources.secrets)
if resources.cos != nil {
return errors.Join(cleanupErr, fmt.Errorf("target ClusterObjectSet may have reconciled; refusing automatic OLMv0 recovery"))
}
recoverErr := m.RecoverFromBackup(recoveryCtx, opts, backup)
return errors.Join(cleanupErr, recoverErr)
}
Expand Down Expand Up @@ -423,15 +470,16 @@ func (m *Migrator) CreateClusterObjectSet(ctx context.Context, opts Options, inf
}

// CreateMigrationResources creates the ClusterObjectSet and its ClusterExtension
// as one recoverable operation. If either creation fails, it removes only the
// objects created by this invocation before restoring the OLMv0 Subscription.
func (m *Migrator) CreateMigrationResources(ctx context.Context, opts Options, info *MigrationInfo, backup *Backup) error {
// as one recoverable operation. Its result reports whether a target COS may have
// started reconciling, including when cleanup was requested after a failure.
func (m *Migrator) CreateMigrationResources(ctx context.Context, opts Options, info *MigrationInfo, backup *Backup) (MigrationResourcesResult, error) {
resources, err := m.createClusterObjectSet(ctx, opts, info)
if err != nil {
result := MigrationResourcesResult{TargetMayBeActive: resources.cos != nil || resources.ownershipUnknown}
if recoverErr := m.recoverCreatedMigrationResources(ctx, opts, backup, resources); recoverErr != nil {
return fmt.Errorf("COS creation failed: %w; recovery also failed: %v", err, recoverErr)
return result, fmt.Errorf("COS creation failed: %w; recovery also failed: %v", err, recoverErr)
}
return fmt.Errorf("COS creation failed (recovered): %w", err)
return result, fmt.Errorf("COS creation failed (recovered): %w", err)
}

ce, ownershipUnknown, err := m.createClusterExtension(ctx, opts, info)
Expand All @@ -440,12 +488,13 @@ func (m *Migrator) CreateMigrationResources(ctx context.Context, opts Options, i
}
resources.ownershipUnknown = resources.ownershipUnknown || ownershipUnknown
if err != nil {
result := MigrationResourcesResult{TargetMayBeActive: resources.cos != nil || resources.ownershipUnknown}
if recoverErr := m.recoverCreatedMigrationResources(ctx, opts, backup, resources); recoverErr != nil {
return fmt.Errorf("ClusterExtension creation failed: %w; recovery also failed: %v", err, recoverErr)
return result, fmt.Errorf("ClusterExtension creation failed: %w; recovery also failed: %v", err, recoverErr)
}
return fmt.Errorf("ClusterExtension creation failed (recovered): %w", err)
return result, fmt.Errorf("ClusterExtension creation failed (recovered): %w", err)
}
return nil
return MigrationResourcesResult{}, nil
}

func (m *Migrator) createClusterObjectSet(ctx context.Context, opts Options, info *MigrationInfo) (*createdMigrationResources, error) {
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Expand Down Expand Up @@ -501,6 +550,7 @@ func (m *Migrator) createClusterObjectSet(ctx context.Context, opts Options, inf
secret.Annotations[migrationInvocationAnnotation] = invocationMarker
if err := m.Client.Create(ctx, secret); err != nil {
if apierrors.IsAlreadyExists(err) {
resources.ownershipUnknown = true
return resources, fmt.Errorf("failed to create COS ref Secret %s: %w", secret.Name, err)
}
if m.resolveCreatedObject(context.WithoutCancel(ctx), secret, invocationMarker) {
Expand Down Expand Up @@ -564,6 +614,7 @@ func (m *Migrator) createClusterObjectSet(ctx context.Context, opts Options, inf
// and Secret cleanup unsafe.
if err := m.Client.Create(ctx, cosObj); err != nil {
if apierrors.IsAlreadyExists(err) {
resources.ownershipUnknown = true
return resources, failWithSecretCleanup(fmt.Errorf("failed to create ClusterObjectSet: %w", err))
}
if m.resolveCreatedObject(context.WithoutCancel(ctx), cosObj, invocationMarker) {
Expand Down Expand Up @@ -663,7 +714,6 @@ func (m *Migrator) createClusterExtension(ctx context.Context, opts Options, inf
if opts.AcknowledgeNotSteadyState {
annotations[AnnotationAcknowledgedPrefix+"not-steady-state"] = "true"
}

ce := &ocv1.ClusterExtension{
ObjectMeta: metav1.ObjectMeta{
Name: opts.ClusterExtensionName,
Expand Down Expand Up @@ -779,6 +829,17 @@ type CleanupResult struct {
Actions []CleanupAction
}

// Err returns all failures encountered while cleaning up OLMv0 resources.
func (r *CleanupResult) Err() error {
var errs []error
for _, action := range r.Actions {
if action.Error != nil {
errs = append(errs, fmt.Errorf("%s: %w", action.Description, action.Error))
}
}
return errors.Join(errs...)
}

// CleanupOLMv0Resources removes remaining OLMv0 resources after migration.
func (m *Migrator) CleanupOLMv0Resources(ctx context.Context, opts Options, packageName, csvName string) *CleanupResult {
result := &CleanupResult{}
Expand Down
Loading
Loading